docs: add workflow log streaming spec
This commit is contained in:
@@ -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**: `<logdir>/<run_id>/<server_id>.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 `<logdir>/<run_id>/` dirs older than retention by run `finished_at`. |
|
||||
| Log dir | Env `VANTAGE_WORKFLOW_LOG_DIR`, default `<data>/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 `<logdir>/<run_id>/` exists.
|
||||
- Before dispatching each step: write the step marker line to the file (`\n===== step <order>: <name> =====\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: `<logdir>/<run_id>/<server_id>.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 `<logdir>/<run_id>/` 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<string>` (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 `<pre>` 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 `<details>` 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).
|
||||
```
|
||||
Reference in New Issue
Block a user