// Package execlog persists incremental command output for an execution and // announces it on the event stream. // // It exists as its own package because both SSH execution paths need it — // internal/mcp (the agent's auto-run windows) and internal/httpapi (the // post-approval actuator). Those two already carry near-identical copies of // sshExec, and every bug found in this area so far has been a case of the two // copies drifting apart; one shared sink is the cheap way not to repeat that. package execlog import ( "context" "log/slog" "sync" "time" "github.com/dtoro/oikos/internal/db" "github.com/dtoro/oikos/internal/db/sqlcgen" "github.com/dtoro/oikos/internal/observability" "github.com/google/uuid" ) // eventInterval throttles execution.output events. Chunks are persisted as // they arrive, but a chatty command (apt, a long build) can produce hundreds // per second and the SSE broker drops events for slow subscribers — flooding // it would push out the signal.* and approval.* events that actually need to // arrive. The event is only a "there is more output" ping; subscribers re-read // the rows. const eventInterval = time.Second // Sink receives output chunks as they arrive from a remote command. type Sink func(stream string, chunk []byte) // New returns a Sink that writes chunks to execution_logs and emits a // throttled execution.output event, plus a Flush to call when the command // finishes. // // The returned Sink is safe for concurrent use: stdout and stderr are written // from separate goroutines. func New(ctx context.Context, pool *db.Pool, execID uuid.UUID, correlationID string) (Sink, func()) { var ( mu sync.Mutex seq int lastEvent time.Time pending bool ) emit := func() { if err := observability.Event(ctx, sqlcgen.New(pool), "execution.output", &execID, "info", "actuator", correlationID, map[string]any{"execution_id": execID.String()}); err != nil { slog.Debug("execlog: emit output event", "error", err, "execution_id", execID) } } sink := func(stream string, chunk []byte) { if len(chunk) == 0 { return } mu.Lock() seq++ n := seq mu.Unlock() // A failed log write must never fail the command: this is observability, // and the authoritative output still lands in executions.result at the // end. Log and carry on. if _, err := pool.Exec(ctx, `INSERT INTO execution_logs (execution_id, seq, stream, chunk) VALUES ($1, $2, $3, $4)`, execID, n, stream, string(chunk)); err != nil { slog.Debug("execlog: persist chunk", "error", err, "execution_id", execID) return } mu.Lock() due := time.Since(lastEvent) >= eventInterval if due { lastEvent = time.Now() pending = false } else { pending = true } mu.Unlock() if due { emit() } } // Flush emits a final event when output arrived inside the throttle window, // so the last few lines of a short command are not left unannounced. flush := func() { mu.Lock() due := pending pending = false mu.Unlock() if due { emit() } } return sink, flush } // Read returns an execution's persisted output in order. func Read(ctx context.Context, pool *db.Pool, execID uuid.UUID, limit int) ([]Chunk, error) { if limit <= 0 { limit = 1000 } rows, err := pool.Query(ctx, `SELECT seq, stream, chunk, ts FROM execution_logs WHERE execution_id = $1 ORDER BY seq LIMIT $2`, execID, limit) if err != nil { return nil, err } defer rows.Close() var out []Chunk for rows.Next() { var c Chunk if err := rows.Scan(&c.Seq, &c.Stream, &c.Chunk, &c.TS); err != nil { return nil, err } out = append(out, c) } return out, rows.Err() } // Chunk is one persisted slice of command output. type Chunk struct { Seq int `json:"seq"` Stream string `json:"stream"` Chunk string `json:"chunk"` TS time.Time `json:"ts"` }