diff --git a/docs/superpowers/specs/2026-07-20-workflow-log-streaming-design.md b/docs/superpowers/specs/2026-07-20-workflow-log-streaming-design.md new file mode 100644 index 0000000..1d1f252 --- /dev/null +++ b/docs/superpowers/specs/2026-07-20-workflow-log-streaming-design.md @@ -0,0 +1,193 @@ +# Workflow Log Streaming — Design + +**Date:** 2026-07-20 +**Status:** Approved (design) — ready for implementation planning +**Scope:** Stream step stdout/stderr live from agent to server-side log files, tail them live in the UI, and auto-expire them on a retention period. Enhancement to the already-merged Server Workflows feature. No auth/orgs, no inventory. + +--- + +## 1. Summary + +Today a workflow step buffers all stdout/stderr in agent RAM, ships it in one terminal `StepResult`, and the server persists the whole body into the `workflow_runs` Mongo document. Long/chatty steps risk: agent memory blow-up, the gRPC 4MB message ceiling, and the Mongo 16MB document cap. + +Change to **live streaming**: + +1. Agent streams output chunks over the existing `CommandStream` as the process runs. +2. Server appends chunks (secret-masked) to a **per-server-run log file** on disk — not Mongo. +3. UI tails the file live via **SSE** while a server-run is running; slices per-step by byte offset after completion. +4. A **retention sweeper** deletes old run-log directories on a configurable period (default 30 days, set in Settings). + +`workflow_runs` documents shrink: they no longer carry `stdout`/`stderr` bodies, only status/exit/attempts/output_env/timestamps plus a per-step `log_offset`. + +--- + +## 2. Locked decisions + +| Topic | Decision | +|-------|----------| +| Transport | Reuse bidirectional `CommandStream`. New `AgentMessage.StepOutput` chunk message. | +| Chunk shape | `{command_id, seq, data, eof}`. Interleaved stdout+stderr in execution order. | +| Terminal result | `StepResult` still sent at step end, now carries only `exit_code` + `output_env` (no stdout/stderr). | +| Log granularity | **One file per server-run**: `//.log`, with a marker line before each step. | +| Streams | **Interleaved** — single synchronized writer on the agent, terminal-order output. | +| Masking | **Server-side** (agent can't tell secret env from normal env). Per-stream carry buffer of `maxSecretLen-1` bytes so a secret split across a chunk boundary still masks; flushed on EOF. | +| Live tail | **SSE** at per-server-run granularity: `GET /api/runs/:runId/servers/:serverId/logs/stream`. Post-run whole-file fetch + per-step offset slice. | +| Retention | `settings.workflow_log_retention_days`, default **30**, editable in `/settings`. Hourly sweeper deletes `//` dirs older than retention by run `finished_at`. | +| Log dir | Env `VANTAGE_WORKFLOW_LOG_DIR`, default `/workflow-logs`. Created `0700`. | +| Mongo | No log bodies in `workflow_runs`. Disk is the source of truth for output. | + +--- + +## 3. gRPC protocol (`proto/vantage/v1/vantage.proto` + both `pb.go` files) + +Add to the `AgentMessage` oneof: `StepOutputChunk step_output = 6;` + +```protobuf +message StepOutputChunk { + string command_id = 1; + uint64 seq = 2; // monotonic per command_id, 0-based + bytes data = 3; // raw interleaved stdout+stderr bytes + bool eof = 4; // true on the final (empty) chunk +} +``` + +`StepResult` is unchanged in shape but `stdout`/`stderr` are now left empty by the agent (kept in the message for backward-compat / error notes only — server ignores them for log content). The server still reads `exit_code` and `output_env` from `StepResult`. + +Hand-written JSON-codec struct added to **both** `server/internal/grpc/pb/vantage.pb.go` and `agent/internal/grpc/pb/vantage.pb.go`, identical. `AgentMessage` gains `StepOutput *StepOutputChunk` in both. + +`data` is `[]byte` in the Go structs (JSON-codec base64-encodes it, which is fine). + +--- + +## 4. Agent (`agent/internal/exec/exec.go`) + +`RunStep` signature gains a chunk sink: + +```go +func RunStep(cmd *pb.RunStepCmd, emit func(seq uint64, data []byte)) *pb.StepResult +``` + +- Replace the two `bytes.Buffer`s with a single `streamWriter` set as **both** `c.Stdout` and `c.Stderr`. Its `Write` takes a mutex (so stdout+stderr interleave without interleaving *within* a write), assigns the next `seq`, and calls `emit(seq, copyOfBytes)`. Chunks are whatever the OS pipe delivers (typically ≤64KB); no extra buffering/line-assembly. +- `StepResult` returns with `Stdout`/`Stderr` empty; `ExitCode` and `OutputEnv` populated as today (env parsing unchanged). +- On timeout/exec error, put the short note in `StepResult.Stderr` (terminal, not streamed) so the runner can still surface a failure reason even if nothing streamed. + +Agent loop (`agent/internal/sync/sync.go`, the `cmd.RunStep != nil` goroutine): pass an `emit` closure that sends `AgentMessage{ServerId, AgentToken, StepOutput: &pb.StepOutputChunk{CommandId, Seq, Data}}` through the existing mutex-guarded `send()`. After `RunStep` returns, send a final `StepOutput{eof:true, seq:last+1}` then the terminal `StepResult` (both via `send()`). Ordering: all chunks, then eof, then StepResult. + +--- + +## 5. Server write path + +### 5.1 Log writer registry (`server/internal/services/steplogs.go`) + +Parallel to `StepResults`. Keyed by `command_id`: + +```go +type stepLogWriter struct { + f *os.File + mu sync.Mutex + carry []byte // held-back tail for boundary-safe masking + secrets []string // secret literals to mask + maxSecret int +} +var StepLogs = &stepLogRegistry{ ... } +func (r *stepLogRegistry) Open(commandID, path string, secrets []string) (*stepLogWriter, error) +func (r *stepLogRegistry) Append(commandID string, data []byte) // masked write +func (r *stepLogRegistry) Close(commandID string) // flush carry, close file +``` + +- `Append` masking: concatenate `carry+data`, mask all secret literals (`ReplaceAll(v,"***")`), then write everything except the last `maxSecret-1` bytes; keep those as the new `carry`. `Close` masks+writes the remaining carry. If `secrets` empty, write straight through (no carry). +- The file handle is opened append-only (`O_APPEND|O_CREATE|O_WRONLY`, `0600`); dir `0700`. + +### 5.2 Stream delivery (`server/internal/grpc/server.go`) + +In the receive loop, after the `m.StepResult` block, add: + +```go +if m.StepOutput != nil { + if m.StepOutput.Eof { + services.StepLogs.Close(m.StepOutput.CommandId) + } else { + services.StepLogs.Append(m.StepOutput.CommandId, m.StepOutput.Data) + } +} +``` + +### 5.3 Runner changes (`server/internal/services/workflow_runner.go`) + +- Resolve the log dir + run/server file path once per server-run; ensure `//` exists. +- Before dispatching each step: write the step marker line to the file (`\n===== step : =====\n`), record the current file byte offset as the step's `log_offset` (persisted on the `StepRun`), and `StepLogs.Open(commandID, path, secretVals)` **before** `DispatchRunStep` (same ordering rule as `StepResults.Await`). +- `dispatchAndWait` no longer expects stdout/stderr in the result. On terminal `StepResult`, `StepLogs.Close(commandID)` is driven by the agent's eof; the runner also calls `Close` defensively on timeout/dispatch-failure (idempotent). +- **Drop** `stdout`/`stderr` from `finishStep` persistence. Masking of the streamed body is done in `Append`; `run_env`/`output_env` masking (existing, from the merged fix) stays. +- The shared server-run file is written by two writers that never overlap in time (steps are serial, and the runner writes each step marker *before* `StepLogs.Open`): (a) the runner writes markers directly to the path, serially between steps; (b) `StepLogs` writes chunks during a step. **Resolved approach:** `Open(commandID, path, secrets)` opens the path fresh with `O_APPEND|O_CREATE|O_WRONLY` for that step and `Close(commandID)` closes it on eof. One handle live at a time per server-run (serial steps guarantee this), so there is no shared-handle race and no ref-counting. The runner's marker write is a separate short `O_APPEND` open/write/close on the same path. + +### 5.4 Data model (`server/internal/models/workflow.go`) + +`StepRun`: +- **Remove** `Stdout`, `Stderr` string fields. +- **Add** `LogOffset int64 `bson:"log_offset" json:"log_offset"`` — byte offset in the server-run file where this step's marker begins. + +`ServerRun` gains nothing structural (its file path is derivable: `//.log`). + +--- + +## 6. REST API (`server/internal/api/workflows.go`) + +- `GET /api/runs/:runId/servers/:serverId/logs` — returns the whole server-run log file (`text/plain`). 404 if absent. Used post-run and as SSE fallback. +- `GET /api/runs/:runId/servers/:serverId/logs/stream` — **SSE**. Opens the file, streams existing content as `data:` events, then polls for appends (~500ms) emitting new bytes, until the server-run status is terminal (success/failed/skipped/cancelled) AND no more bytes, then sends a final `event: done` and closes. Sets `Content-Type: text/event-stream`, disables gin's buffering. Guards against path traversal (runId/serverId are used as literal path segments — validate they are UUIDs / contain no separators). + +Log content served by these endpoints is already masked (masking happens at write time), so no masking needed on read. + +--- + +## 7. Retention + +### 7.1 Setting + +`settings` collection gains `workflow_log_retention_days int` (default 30 when unset). Read/write via the existing settings service + surfaced in `/settings` UI as a number input. `0` or negative disables sweeping (keep forever) — document this. + +### 7.2 Sweeper (`server/internal/services/steplogs.go` or `logsweeper.go`) + +- `StartLogSweeper()` launched at server startup (next to index setup): hourly `time.Ticker`. +- Each tick: read retention setting; if ≤0 skip. Compute cutoff = `now - retentionDays`. For each `//` dir, look up the run's `finished_at` (query `workflow_runs` by run_id); if finished and older than cutoff, `os.RemoveAll` the dir. Fallback to dir mtime if the run doc is gone. +- Also run once at startup. + +--- + +## 8. Frontend + +### 8.1 API client (`web/lib/api.ts`) + +- `StepRun`: remove `stdout`/`stderr`; add `log_offset: number`. +- Add `getServerRunLog(runId, serverId): Promise` (GET .../logs). +- SSE consumed directly via `EventSource` in the component (not through the `request` helper), URL built from the same base. +- Settings type gains `workflow_log_retention_days`. + +### 8.2 Run detail (`web/app/workflows/[id]/runs/[runId]/page.tsx`) + +- Per-server card: while the server-run is `running`, open an `EventSource` to the stream endpoint and render a live `
` terminal that appends incoming chunks (auto-scroll). Close the source on `event: done`, unmount, or terminal status.
+- After completion: fetch the whole file once and render it; step `
` still list status/exit/attempts pills. (Per-step slicing by `log_offset` is optional polish — v1 may show the whole server log under the card and keep step pills as the status summary.) +- Remove reliance on `st.stdout`/`st.stderr` (fields gone). + +### 8.3 Settings (`web/app/settings/page.tsx`) + +- Add a "Workflow log retention (days)" number input bound to `workflow_log_retention_days`, saved via the existing settings mutation. Note that `0` = keep forever. + +--- + +## 9. Security + +- Secret masking moves to the streaming write path but remains server-side and boundary-safe (carry buffer). Same `***` replacement. +- Log files `0600`, dirs `0700`, under a dedicated log dir. +- SSE/read endpoints validate `runId`/`serverId` as UUID-shaped path segments to prevent traversal; they are session-authed (same `apiGroup`). +- Terminal `StepResult.Stderr` (error notes only) is still masked before any persistence (it is no longer persisted as log body; if surfaced, mask against secretVals). + +--- + +## 10. Out of scope + +- Per-step (rather than per-server) live SSE channels. +- Log compression / rotation within a run, remote log storage (S3), download-as-zip. +- Full-text search over logs. +- Backfilling/migrating already-existing `workflow_runs` stdout/stderr into files (pre-existing runs keep whatever they had; new field just won't be set — acceptable, feature is new). +- Tests (skipped, consistent with the Workflows iteration). +```