Trebejo.Stream (Trebejo v2.0.0)

Copy Markdown View Source

Streaming wrapper around Port.open/2 for long-running commands.

Unlike Trebejo.Util.run_cmd/3 which waits for the command to finish and returns the full output, Trebejo.Stream returns a lazy Stream.t() that yields stdout chunks as they arrive, suitable for docker logs -f, tail -f, build logs, etc.

Usage

{:ok, stream} =
  Trebejo.Stream.open("docker", ["logs", "-f", "my-container"],
    stderr: :separate, timeout: 60_000
  )

stream
|> Stream.each(&IO.write/1)
|> Stream.run()

Or consume line by line:

{:ok, lines} = Trebejo.Stream.lines("tail", ["-f", "/var/log/syslog"])
for line <- lines, do: IO.puts(line)

Stderr handling

  • :merge (default) — stderr is redirected to stdout via :stderr_to_stdout. The stream emits binary chunks containing the merged output.
  • :separate — deprecated in 2.0, reserved for 2.1. The underlying Port.open/2 does not yet tag stderr vs stdout reliably across all OTP versions, so for 2.0 we accept the option but silently fall back to :merge behaviour. Calls passing :separate will log a single warning via :logger.warning/1 so callers can find and update them.

Cleanup

Every call to open/3 registers the port with the calling process. When the calling process exits, the port is closed automatically by the BEAM. For short-lived streams you can also call close/1.

Summary

Functions

Force-close a stream early. Safe to call on a stream that has already finished — the port will already be closed by the resource teardown.

Like open/3 but yields complete lines (split on \n) instead of raw chunks. Stderr chunks are tagged with %{kind: :stderr, data: line}.

Open a streaming port to cmd_name with args.

Types

line_chunk()

@type line_chunk() :: String.t() | %{kind: :stdout | :stderr, data: String.t()}

options()

@type options() :: [
  timeout: pos_integer(),
  stderr: :merge | :separate,
  chunk_size: pos_integer(),
  into: pid(),
  cd: String.t(),
  env: [{String.t(), String.t()}]
]

stream_chunk()

@type stream_chunk() :: binary() | %{kind: :stdout | :stderr, data: binary()}

Functions

close(stream)

@spec close(Enumerable.t()) :: :ok

Force-close a stream early. Safe to call on a stream that has already finished — the port will already be closed by the resource teardown.

lines(cmd_name, args, opts \\ [])

@spec lines(binary(), [binary()], options()) ::
  {:ok, Enumerable.t()} | {:error, Trebejo.Error.t()}

Like open/3 but yields complete lines (split on \n) instead of raw chunks. Stderr chunks are tagged with %{kind: :stderr, data: line}.

open(cmd_name, args, opts \\ [])

@spec open(binary(), [binary()], options()) ::
  {:ok, Enumerable.t()} | {:error, Trebejo.Error.t()}

Open a streaming port to cmd_name with args.

Returns {:ok, stream} (a Stream.t() of stream_chunk/0) or {:error, %Trebejo.Error{}}.

The stream is lazy: nothing executes until you start consuming it (via Stream.each/2, Enum.take/2, Stream.run/1, etc.).