diff --git a/cmd/synapbus/main.go b/cmd/synapbus/main.go index 920cf98..a6052a3 100644 --- a/cmd/synapbus/main.go +++ b/cmd/synapbus/main.go @@ -513,6 +513,7 @@ func runServe(cmd *cobra.Command, args []string) error { harnessRunsStore := runs.New(db.DB, slog.Default()) harnessRegistry.Observer = harnessRunsStore reactorEngine.SetHarnessRegistry(harnessRegistry) + reactorEngine.SetReactionNotifier(&reactorReactionAdapter{svc: reactionService}) slog.Info("harness registry configured", "backends", harnessRegistry.Names(), ) @@ -713,6 +714,7 @@ func runServe(cmd *cobra.Command, args []string) error { TrustService: trustService, ReactorStore: reactorStore, ReactorEngine: reactorEngine, + HarnessRunsStore: harnessRunsStore, BaseURL: baseURL, WikiService: wikiService, }) @@ -953,6 +955,18 @@ func (a *attachmentLinkerAdapter) GetByMessageID(ctx context.Context, messageID return results, nil } +// reactorReactionAdapter adapts reactions.Service to the reactor's +// ReactionNotifier interface. It wraps Toggle so the reactor only +// sees one simple AddReaction call. +type reactorReactionAdapter struct { + svc *reactions.Service +} + +func (a *reactorReactionAdapter) AddReaction(ctx context.Context, messageID int64, agentName, reactionType string) error { + _, err := a.svc.Toggle(ctx, messageID, agentName, reactionType, nil) + return err +} + // reactionEnricherAdapter adapts reactions.Service to messaging.ReactionEnricher. type reactionEnricherAdapter struct { svc *reactions.Service diff --git a/examples/cold-topic-explainer/wrapper.sh b/examples/cold-topic-explainer/wrapper.sh index d4d4a33..8602cab 100755 --- a/examples/cold-topic-explainer/wrapper.sh +++ b/examples/cold-topic-explainer/wrapper.sh @@ -44,8 +44,11 @@ PROMPT="$(cat GEMINI.md) Incoming DM from @${FROM}: ${BODY}" -# Preserve the exact prompt and any stderr noise for forensics. -printf '%s' "$PROMPT" > gemini.prompt.txt +# Preserve the exact prompt the model received — the subprocess +# harness reads prompt.txt after the run completes and stores it in +# harness_runs.prompt so the Web UI can show "what the model saw". +printf '%s' "$PROMPT" > prompt.txt + set +e RAW=$(gemini -m "$GEMINI_MODEL" --approval-mode yolo -p "$PROMPT" 2>gemini.stderr.log) GEMINI_EXIT=$? @@ -55,8 +58,10 @@ set -e # stdout when its MCP config can't reach a server. Strip it. RESPONSE=$(printf '%s' "$RAW" | sed 's|^MCP issues detected\. Run /mcp list for status\.||') -# Save the cleaned response for forensics before we decide next steps. -printf '%s' "$RESPONSE" > gemini.stdout.txt +# Save both the raw and the cleaned response. `response.txt` is the +# one the harness persists into harness_runs.response. +printf '%s' "$RAW" > gemini.stdout.raw +printf '%s' "$RESPONSE" > response.txt if [ -z "$RESPONSE" ]; then log "empty gemini response (exit=$GEMINI_EXIT); last stderr:" diff --git a/internal/api/router.go b/internal/api/router.go index 8db312f..15ddf3b 100644 --- a/internal/api/router.go +++ b/internal/api/router.go @@ -10,6 +10,7 @@ import ( "github.com/synapbus/synapbus/internal/apikeys" "github.com/synapbus/synapbus/internal/attachments" "github.com/synapbus/synapbus/internal/channels" + "github.com/synapbus/synapbus/internal/harness/runs" "github.com/synapbus/synapbus/internal/k8s" "github.com/synapbus/synapbus/internal/messaging" "github.com/synapbus/synapbus/internal/reactor" @@ -41,6 +42,7 @@ type RouterConfig struct { TrustService *trust.Service ReactorStore *reactor.Store ReactorEngine *reactor.Reactor + HarnessRunsStore *runs.Store WikiService *wiki.Service SSEHub *SSEHub Broadcaster *SSEBroadcaster @@ -246,7 +248,14 @@ func NewRouterWithConfig(cfg RouterConfig) chi.Router { // Reactive Runs if cfg.ReactorStore != nil && cfg.ReactorEngine != nil && cfg.AgentService != nil { - runsHandler := NewRunsHandler(cfg.ReactorStore, cfg.ReactorEngine, agents.NewSQLiteAgentStore(cfg.DB)) + runsHandler := NewRunsHandler( + cfg.ReactorStore, + cfg.ReactorEngine, + agents.NewSQLiteAgentStore(cfg.DB), + cfg.HarnessRunsStore, + cfg.MsgService, + cfg.DB, + ) r.Group(func(r chi.Router) { r.Use(authMiddleware) diff --git a/internal/api/runs_handler.go b/internal/api/runs_handler.go index 47b3e8e..9144d3b 100644 --- a/internal/api/runs_handler.go +++ b/internal/api/runs_handler.go @@ -1,6 +1,7 @@ package api import ( + "database/sql" "net/http" "strconv" "time" @@ -8,22 +9,37 @@ import ( "github.com/go-chi/chi/v5" "github.com/synapbus/synapbus/internal/agents" + "github.com/synapbus/synapbus/internal/harness/runs" + "github.com/synapbus/synapbus/internal/messaging" "github.com/synapbus/synapbus/internal/reactor" ) // RunsHandler handles REST API requests for reactive runs. type RunsHandler struct { - store *reactor.Store - reactor *reactor.Reactor - agentStore agents.AgentStore + store *reactor.Store + reactor *reactor.Reactor + agentStore agents.AgentStore + harnessRuns *runs.Store + msgService *messaging.MessagingService + db *sql.DB } // NewRunsHandler creates a new runs handler. -func NewRunsHandler(store *reactor.Store, r *reactor.Reactor, agentStore agents.AgentStore) *RunsHandler { +func NewRunsHandler( + store *reactor.Store, + r *reactor.Reactor, + agentStore agents.AgentStore, + harnessRuns *runs.Store, + msgService *messaging.MessagingService, + db *sql.DB, +) *RunsHandler { return &RunsHandler{ - store: store, - reactor: r, - agentStore: agentStore, + store: store, + reactor: r, + agentStore: agentStore, + harnessRuns: harnessRuns, + msgService: msgService, + db: db, } } @@ -57,7 +73,12 @@ func (h *RunsHandler) ListRuns(w http.ResponseWriter, r *http.Request) { }) } -// GetRun returns a single run by ID. +// GetRun returns a composite view of a reactive run: the reactive_runs +// row itself, the linked harness_runs row (with captured prompt / +// response / usage), the triggering message, the outgoing message the +// agent produced (if any), and a snapshot of the agent's current +// harness config. Everything the Web UI needs to render "what happened +// on this run" in a single request. func (h *RunsHandler) GetRun(w http.ResponseWriter, r *http.Request) { idStr := chi.URLParam(r, "id") id, err := strconv.ParseInt(idStr, 10, 64) @@ -72,7 +93,80 @@ func (h *RunsHandler) GetRun(w http.ResponseWriter, r *http.Request) { return } - writeJSON(w, http.StatusOK, run) + resp := map[string]any{ + "run": run, + } + + // Linked harness_run (may be nil for K8s path which still uses + // the legacy reactive_runs-only flow). + if h.harnessRuns != nil { + hr, _ := h.harnessRuns.GetByReactiveRunID(r.Context(), id) + if hr != nil { + resp["harness_run"] = hr + } + } + + // Triggering message body (what the sender wrote). + if run.TriggerMessageID != nil && h.msgService != nil { + if msg, err := h.msgService.GetMessageByID(r.Context(), *run.TriggerMessageID); err == nil && msg != nil { + resp["trigger_message"] = msg + } + } + + // Agent snapshot — current harness config so the UI can show + // the gemini_md / claude_md the agent is currently running with. + if agent, err := h.agentStore.GetAgentByName(r.Context(), run.AgentName); err == nil && agent != nil { + resp["agent"] = map[string]any{ + "name": agent.Name, + "display_name": agent.DisplayName, + "type": agent.Type, + "harness_name": agent.HarnessName, + "local_command": agent.LocalCommand, + "harness_config_json": agent.HarnessConfigJSON, + "trigger_mode": agent.TriggerMode, + "cooldown_seconds": agent.CooldownSeconds, + "daily_trigger_budget": agent.DailyTriggerBudget, + "max_trigger_depth": agent.MaxTriggerDepth, + } + } + + // Outgoing message — the first DM this agent produced after + // the run started. We find it by querying messages where + // from_agent = this run's agent AND created_at >= run.StartedAt, + // ordered by id. Works for both success and failure cases. + if h.db != nil && run.StartedAt != nil { + var ( + msgID int64 + toAgent sql.NullString + body string + status string + createdAt string + ) + // Wrap both sides in datetime() so SQLite parses and compares + // canonically — the messages table stores created_at as + // 'YYYY-MM-DD HH:MM:SS' (space separator) while Go emits + // RFC3339 with 'T'. A raw string comparison fails silently. + err := h.db.QueryRowContext(r.Context(), + `SELECT id, to_agent, body, status, created_at + FROM messages + WHERE from_agent = ? + AND datetime(created_at) >= datetime(?) + ORDER BY id ASC LIMIT 1`, + run.AgentName, + run.StartedAt.UTC().Format(time.RFC3339), + ).Scan(&msgID, &toAgent, &body, &status, &createdAt) + if err == nil { + resp["outgoing_message"] = map[string]any{ + "id": msgID, + "to_agent": toAgent.String, + "body": body, + "status": status, + "created_at": createdAt, + } + } + } + + writeJSON(w, http.StatusOK, resp) } // RetryRun retries a failed run. diff --git a/internal/harness/harness.go b/internal/harness/harness.go index 6f77078..27088ca 100644 --- a/internal/harness/harness.go +++ b/internal/harness/harness.go @@ -96,6 +96,12 @@ type ExecRequest struct { // asks the backend to resume a prior conversation. SessionID string + // ReactiveRunID, when > 0, is the id of the reactive_runs row + // that triggered this dispatch. Stored alongside the harness_runs + // row so the Web UI can JOIN reactive_runs ↔ harness_runs and + // show operators which subprocess run produced which agent reply. + ReactiveRunID int64 + // Budget bounds wall-clock, tokens, and cost. Budget Budget @@ -125,6 +131,18 @@ type ExecResult struct { // in the response body. ResultJSON json.RawMessage + // Prompt is the rendered prompt the child actually received — + // system instructions + user message + whatever context the + // backend assembled. Subprocess backends populate this by reading + // a `prompt.txt` the wrapper wrote into the workdir. Bounded in + // size by the observer when it persists. + Prompt string + + // Response is the raw text the backend produced before any + // post-processing (routing decisions, JSON parsing). Subprocess + // backends read this from `response.txt`. + Response string + // Usage captures token / cost accounting when the backend can // report it. Zero values mean "not reported". Usage Usage diff --git a/internal/harness/runs/store.go b/internal/harness/runs/store.go index 12f236e..355c35d 100644 --- a/internal/harness/runs/store.go +++ b/internal/harness/runs/store.go @@ -16,27 +16,32 @@ import ( "github.com/synapbus/synapbus/internal/observability" ) -// Run is one row of the harness_runs table. +// Run is one row of the harness_runs table. JSON tags match the +// snake_case convention used everywhere else in the Web UI TypeScript +// layer; fields map 1:1 to column names. type Run struct { - ID int64 - RunID string - AgentName string - Backend string - MessageID *int64 - Status string - ExitCode *int - TraceID string - SpanID string - SessionID string - TokensIn int64 - TokensOut int64 - TokensCached int64 - CostUSD float64 - DurationMs *int64 - ResultJSON string - LogsExcerpt string - CreatedAt time.Time - FinishedAt *time.Time + ID int64 `json:"id"` + RunID string `json:"run_id"` + AgentName string `json:"agent_name"` + Backend string `json:"backend"` + MessageID *int64 `json:"message_id,omitempty"` + ReactiveRunID *int64 `json:"reactive_run_id,omitempty"` + Status string `json:"status"` + ExitCode *int `json:"exit_code,omitempty"` + TraceID string `json:"trace_id,omitempty"` + SpanID string `json:"span_id,omitempty"` + SessionID string `json:"session_id,omitempty"` + TokensIn int64 `json:"tokens_in"` + TokensOut int64 `json:"tokens_out"` + TokensCached int64 `json:"tokens_cached"` + CostUSD float64 `json:"cost_usd"` + DurationMs *int64 `json:"duration_ms,omitempty"` + ResultJSON string `json:"result_json,omitempty"` + LogsExcerpt string `json:"logs_excerpt,omitempty"` + Prompt string `json:"prompt,omitempty"` + Response string `json:"response,omitempty"` + CreatedAt time.Time `json:"created_at"` + FinishedAt *time.Time `json:"finished_at,omitempty"` } // Status constants match the harness_runs.status column domain. @@ -87,14 +92,20 @@ func (s *Store) OnStart(ctx context.Context, agent *agents.Agent, harnessName st id := req.Message.ID msgID = &id } + var reactiveRunID *int64 + if req.ReactiveRunID > 0 { + id := req.ReactiveRunID + reactiveRunID = &id + } _, err := s.db.ExecContext(ctx, - `INSERT INTO harness_runs (run_id, agent_name, backend, message_id, status, trace_id, session_id, created_at) - VALUES (?, ?, ?, ?, ?, ?, ?, CURRENT_TIMESTAMP)`, + `INSERT INTO harness_runs (run_id, agent_name, backend, message_id, reactive_run_id, status, trace_id, session_id, created_at) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, CURRENT_TIMESTAMP)`, req.RunID, agentNameOf(agent), harnessName, msgID, + reactiveRunID, StatusRunning, observability.TraceIDFromContext(ctx), req.SessionID, @@ -125,6 +136,7 @@ func (s *Store) OnFinish(ctx context.Context, agent *agents.Agent, harnessName s var exitCode *int logsExcerpt := "" var resultJSON string + var promptText, responseText string var tokensIn, tokensOut, tokensCached int64 var costUSD float64 sessionID := req.SessionID @@ -136,6 +148,8 @@ func (s *Store) OnFinish(ctx context.Context, agent *agents.Agent, harnessName s if len(res.ResultJSON) > 0 { resultJSON = string(res.ResultJSON) } + promptText = res.Prompt + responseText = res.Response tokensIn = res.Usage.TokensIn tokensOut = res.Usage.TokensOut tokensCached = res.Usage.TokensCached @@ -150,11 +164,21 @@ func (s *Store) OnFinish(ctx context.Context, agent *agents.Agent, harnessName s if execErr != nil { status = StatusFailed } - // Cap logs excerpt to keep row sizes sane (matches design: bounded). - const logsCap = 16 * 1024 + // Cap each large text field so harness_runs rows stay reasonable. + const ( + logsCap = 16 * 1024 + promptCap = 32 * 1024 + responseCap = 32 * 1024 + ) if len(logsExcerpt) > logsCap { logsExcerpt = "... [truncated] ...\n" + logsExcerpt[len(logsExcerpt)-logsCap:] } + if len(promptText) > promptCap { + promptText = "... [truncated " + fmt.Sprintf("%d", len(promptText)-promptCap) + " bytes] ...\n" + promptText[len(promptText)-promptCap:] + } + if len(responseText) > responseCap { + responseText = "... [truncated " + fmt.Sprintf("%d", len(responseText)-responseCap) + " bytes] ...\n" + responseText[len(responseText)-responseCap:] + } traceID := "" if res != nil && res.TraceID != "" { @@ -170,6 +194,7 @@ func (s *Store) OnFinish(ctx context.Context, agent *agents.Agent, harnessName s status = ?, exit_code = ?, trace_id = ?, session_id = ?, tokens_in = ?, tokens_out = ?, tokens_cached = ?, cost_usd = ?, duration_ms = ?, result_json = ?, logs_excerpt = ?, + prompt = ?, response = ?, finished_at = CURRENT_TIMESTAMP WHERE run_id = ?` @@ -177,19 +202,26 @@ func (s *Store) OnFinish(ctx context.Context, agent *agents.Agent, harnessName s status, exitCode, traceID, sessionID, tokensIn, tokensOut, tokensCached, costUSD, durationMs, nullableString(resultJSON), nullableString(logsExcerpt), + nullableString(promptText), nullableString(responseText), req.RunID, ) if err == nil { if n, _ := result.RowsAffected(); n == 0 { // Row didn't exist — insert it fresh. + var reactiveRunID *int64 + if req.ReactiveRunID > 0 { + id := req.ReactiveRunID + reactiveRunID = &id + } _, err = s.db.ExecContext(ctx, - `INSERT INTO harness_runs (run_id, agent_name, backend, status, exit_code, trace_id, session_id, - tokens_in, tokens_out, tokens_cached, cost_usd, duration_ms, result_json, logs_excerpt, created_at, finished_at) - VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP)`, - req.RunID, agentNameOf(agent), harnessName, + `INSERT INTO harness_runs (run_id, agent_name, backend, reactive_run_id, status, exit_code, trace_id, session_id, + tokens_in, tokens_out, tokens_cached, cost_usd, duration_ms, result_json, logs_excerpt, prompt, response, created_at, finished_at) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP)`, + req.RunID, agentNameOf(agent), harnessName, reactiveRunID, status, exitCode, traceID, sessionID, tokensIn, tokensOut, tokensCached, costUSD, durationMs, nullableString(resultJSON), nullableString(logsExcerpt), + nullableString(promptText), nullableString(responseText), ) } } @@ -206,6 +238,21 @@ func (s *Store) GetByRunID(ctx context.Context, runID string) (*Run, error) { return scanRun(row) } +// GetByReactiveRunID returns the most recent harness_run linked to a +// reactive_run. Used by the Web UI to JOIN the two tables without +// leaking SQL into the API layer. Returns nil, nil when no row exists. +func (s *Store) GetByReactiveRunID(ctx context.Context, reactiveRunID int64) (*Run, error) { + row := s.db.QueryRowContext(ctx, + selectSQL()+` WHERE reactive_run_id = ? ORDER BY id DESC LIMIT 1`, + reactiveRunID, + ) + run, err := scanRun(row) + if err == sql.ErrNoRows { + return nil, nil + } + return run, err +} + // ListByAgent returns recent runs for an agent, newest first. func (s *Store) ListByAgent(ctx context.Context, agentName string, limit int) ([]*Run, error) { if limit <= 0 { @@ -239,23 +286,23 @@ func nullableString(s string) any { } func selectSQL() string { - return `SELECT id, run_id, agent_name, backend, message_id, status, exit_code, + return `SELECT id, run_id, agent_name, backend, message_id, reactive_run_id, status, exit_code, trace_id, span_id, session_id, tokens_in, tokens_out, tokens_cached, cost_usd, - duration_ms, result_json, logs_excerpt, created_at, finished_at + duration_ms, result_json, logs_excerpt, prompt, response, created_at, finished_at FROM harness_runs` } func scanRun(row *sql.Row) (*Run, error) { var r Run - var msgID sql.NullInt64 + var msgID, reactiveRunID sql.NullInt64 var exitCode sql.NullInt64 - var traceID, spanID, sessionID, resultJSON, logsExcerpt sql.NullString + var traceID, spanID, sessionID, resultJSON, logsExcerpt, prompt, response sql.NullString var durationMs sql.NullInt64 var createdAt, finishedAt sql.NullTime if err := row.Scan( - &r.ID, &r.RunID, &r.AgentName, &r.Backend, &msgID, &r.Status, &exitCode, + &r.ID, &r.RunID, &r.AgentName, &r.Backend, &msgID, &reactiveRunID, &r.Status, &exitCode, &traceID, &spanID, &sessionID, &r.TokensIn, &r.TokensOut, &r.TokensCached, &r.CostUSD, - &durationMs, &resultJSON, &logsExcerpt, &createdAt, &finishedAt, + &durationMs, &resultJSON, &logsExcerpt, &prompt, &response, &createdAt, &finishedAt, ); err != nil { return nil, err } @@ -263,6 +310,10 @@ func scanRun(row *sql.Row) (*Run, error) { v := msgID.Int64 r.MessageID = &v } + if reactiveRunID.Valid { + v := reactiveRunID.Int64 + r.ReactiveRunID = &v + } if exitCode.Valid { v := int(exitCode.Int64) r.ExitCode = &v @@ -272,6 +323,8 @@ func scanRun(row *sql.Row) (*Run, error) { r.SessionID = sessionID.String r.ResultJSON = resultJSON.String r.LogsExcerpt = logsExcerpt.String + r.Prompt = prompt.String + r.Response = response.String if durationMs.Valid { v := durationMs.Int64 r.DurationMs = &v @@ -290,15 +343,15 @@ func scanRuns(rows *sql.Rows) ([]*Run, error) { var out []*Run for rows.Next() { var r Run - var msgID sql.NullInt64 + var msgID, reactiveRunID sql.NullInt64 var exitCode sql.NullInt64 - var traceID, spanID, sessionID, resultJSON, logsExcerpt sql.NullString + var traceID, spanID, sessionID, resultJSON, logsExcerpt, prompt, response sql.NullString var durationMs sql.NullInt64 var createdAt, finishedAt sql.NullTime if err := rows.Scan( - &r.ID, &r.RunID, &r.AgentName, &r.Backend, &msgID, &r.Status, &exitCode, + &r.ID, &r.RunID, &r.AgentName, &r.Backend, &msgID, &reactiveRunID, &r.Status, &exitCode, &traceID, &spanID, &sessionID, &r.TokensIn, &r.TokensOut, &r.TokensCached, &r.CostUSD, - &durationMs, &resultJSON, &logsExcerpt, &createdAt, &finishedAt, + &durationMs, &resultJSON, &logsExcerpt, &prompt, &response, &createdAt, &finishedAt, ); err != nil { return nil, err } @@ -306,6 +359,10 @@ func scanRuns(rows *sql.Rows) ([]*Run, error) { v := msgID.Int64 r.MessageID = &v } + if reactiveRunID.Valid { + v := reactiveRunID.Int64 + r.ReactiveRunID = &v + } if exitCode.Valid { v := int(exitCode.Int64) r.ExitCode = &v @@ -315,6 +372,8 @@ func scanRuns(rows *sql.Rows) ([]*Run, error) { r.SessionID = sessionID.String r.ResultJSON = resultJSON.String r.LogsExcerpt = logsExcerpt.String + r.Prompt = prompt.String + r.Response = response.String if durationMs.Valid { v := durationMs.Int64 r.DurationMs = &v diff --git a/internal/harness/runs/store_test.go b/internal/harness/runs/store_test.go index a941933..a8e41c9 100644 --- a/internal/harness/runs/store_test.go +++ b/internal/harness/runs/store_test.go @@ -25,25 +25,28 @@ func setupDB(t *testing.T) *sql.DB { } schema := ` CREATE TABLE harness_runs ( - id INTEGER PRIMARY KEY AUTOINCREMENT, - run_id TEXT NOT NULL UNIQUE, - agent_name TEXT NOT NULL, - backend TEXT NOT NULL, - message_id INTEGER, - status TEXT NOT NULL, - exit_code INTEGER, - trace_id TEXT, - span_id TEXT, - session_id TEXT, - tokens_in INTEGER NOT NULL DEFAULT 0, - tokens_out INTEGER NOT NULL DEFAULT 0, - tokens_cached INTEGER NOT NULL DEFAULT 0, - cost_usd REAL NOT NULL DEFAULT 0, - duration_ms INTEGER, - result_json TEXT, - logs_excerpt TEXT, - created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, - finished_at DATETIME + id INTEGER PRIMARY KEY AUTOINCREMENT, + run_id TEXT NOT NULL UNIQUE, + agent_name TEXT NOT NULL, + backend TEXT NOT NULL, + message_id INTEGER, + reactive_run_id INTEGER, + status TEXT NOT NULL, + exit_code INTEGER, + trace_id TEXT, + span_id TEXT, + session_id TEXT, + tokens_in INTEGER NOT NULL DEFAULT 0, + tokens_out INTEGER NOT NULL DEFAULT 0, + tokens_cached INTEGER NOT NULL DEFAULT 0, + cost_usd REAL NOT NULL DEFAULT 0, + duration_ms INTEGER, + result_json TEXT, + logs_excerpt TEXT, + prompt TEXT, + response TEXT, + created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, + finished_at DATETIME );` if _, err := db.Exec(schema); err != nil { t.Fatalf("schema: %v", err) diff --git a/internal/harness/subprocess/subprocess.go b/internal/harness/subprocess/subprocess.go index d921f8e..e646df1 100644 --- a/internal/harness/subprocess/subprocess.go +++ b/internal/harness/subprocess/subprocess.go @@ -202,6 +202,17 @@ func (h *Harness) Execute(ctx context.Context, req *harness.ExecRequest) (*harne } } + // Load optional prompt.txt / response.txt that wrappers write as + // a post-hoc audit trail for the Web UI run detail view. + promptText := "" + if raw, readErr := os.ReadFile(filepath.Join(workdir, "prompt.txt")); readErr == nil { + promptText = string(raw) + } + responseText := "" + if raw, readErr := os.ReadFile(filepath.Join(workdir, "response.txt")); readErr == nil { + responseText = string(raw) + } + logs := mergeLogs(&stdout, &stderr, h.cfg.LogsCap) // Cleanup policy: remove workdir on success unless configured to @@ -222,6 +233,8 @@ func (h *Harness) Execute(ctx context.Context, req *harness.ExecRequest) (*harne ExitCode: exitCode, Logs: logs, ResultJSON: resultJSON, + Prompt: promptText, + Response: responseText, } // Distinguish context timeout from plain failures so the caller diff --git a/internal/reactor/reactor.go b/internal/reactor/reactor.go index 12fa47b..4f53a7f 100644 --- a/internal/reactor/reactor.go +++ b/internal/reactor/reactor.go @@ -33,12 +33,13 @@ const ( // Reactor is the reactive agent triggering engine. type Reactor struct { - store *Store - agentStore agents.AgentStore - runner k8spkg.JobRunner - registry *harness.Registry - notifier FailureNotifier - logger *slog.Logger + store *Store + agentStore agents.AgentStore + runner k8spkg.JobRunner + registry *harness.Registry + notifier FailureNotifier + reactions ReactionNotifier + logger *slog.Logger } // FailureNotifier sends system DMs on job failure. @@ -46,6 +47,15 @@ type FailureNotifier interface { NotifyFailure(ctx context.Context, ownerAgentName, agentName, triggerFrom, triggerEvent string, durationMs int64, errorSummary string) error } +// ReactionNotifier lets the reactor mark the triggering DM with a +// workflow reaction (in_progress / done / reject) so the sender can +// see at a glance whether their message is being processed. Errors +// are logged, never returned — reactions are cosmetic and must not +// block the main dispatch flow. +type ReactionNotifier interface { + AddReaction(ctx context.Context, messageID int64, agentName, reactionType string) error +} + // New creates a new Reactor. func New(store *Store, agentStore agents.AgentStore, runner k8spkg.JobRunner, logger *slog.Logger) *Reactor { return &Reactor{ @@ -68,6 +78,29 @@ func (r *Reactor) SetHarnessRegistry(reg *harness.Registry) { r.registry = reg } +// SetReactionNotifier wires the component that marks triggering DMs +// with in_progress / done / reject reactions. Optional — a nil +// notifier simply skips the reaction step. +func (r *Reactor) SetReactionNotifier(n ReactionNotifier) { + r.reactions = n +} + +// react is a best-effort wrapper that swallows and logs errors so the +// dispatch flow is never blocked by a reaction-store hiccup. +func (r *Reactor) react(ctx context.Context, messageID int64, agentName, reactionType string) { + if r.reactions == nil || messageID <= 0 { + return + } + if err := r.reactions.AddReaction(ctx, messageID, agentName, reactionType); err != nil { + r.logger.Debug("reaction notifier failed", + "message_id", messageID, + "agent", agentName, + "reaction", reactionType, + "error", err, + ) + } +} + // Dispatch implements dispatcher.EventDispatcher. Called by MultiDispatcher // when a message event occurs. func (r *Reactor) Dispatch(ctx context.Context, event dispatcher.MessageEvent) error { @@ -273,6 +306,10 @@ func (r *Reactor) dispatchHarness(ctx context.Context, agent *agents.Agent, even metrics.ReactiveTriggersTotal.WithLabelValues(agent.Name, StatusRunning).Inc() metrics.ReactiveAgentState.WithLabelValues(agent.Name).Set(1) + // Eyes-on: mark the triggering DM in_progress as the targeted + // agent so the sender's Web UI shows visual progress. + r.react(ctx, event.MessageID, agent.Name, "in_progress") + r.logger.Info("reactive harness dispatch", "agent", agent.Name, "backend", r.agentBackendKind(agent), @@ -283,9 +320,10 @@ func (r *Reactor) dispatchHarness(ctx context.Context, agent *agents.Agent, even ) req := &harness.ExecRequest{ - RunID: uuid.NewString(), - AgentName: agent.Name, - Agent: agent, + RunID: uuid.NewString(), + AgentName: agent.Name, + Agent: agent, + ReactiveRunID: runID, Message: &messaging.Message{ ID: event.MessageID, FromAgent: event.FromAgent, @@ -343,7 +381,11 @@ func (r *Reactor) runHarness(runID int64, agent *agents.Agent, event dispatcher. todayCount, _ := r.store.CountTodayRuns(ctx, agent.Name) metrics.ReactiveBudgetUsed.WithLabelValues(agent.Name).Set(float64(todayCount)) + // Mark the triggering DM with the terminal reaction. Done wins + // over in_progress via priority (see reactions.reactionPriority), + // so we don't need to remove the in_progress reaction first. if status == StatusFailed { + r.react(context.Background(), event.MessageID, agent.Name, "reject") r.logger.Warn("reactive harness run failed", "agent", agent.Name, "run_id", runID, @@ -351,6 +393,7 @@ func (r *Reactor) runHarness(runID int64, agent *agents.Agent, event dispatcher. ) r.notifyFailure(context.Background(), agent, event, durationMs, errorLog) } else { + r.react(context.Background(), event.MessageID, agent.Name, "done") r.logger.Info("reactive harness run succeeded", "agent", agent.Name, "run_id", runID, diff --git a/internal/storage/schema/020_harness_run_detail.sql b/internal/storage/schema/020_harness_run_detail.sql new file mode 100644 index 0000000..57c2c20 --- /dev/null +++ b/internal/storage/schema/020_harness_run_detail.sql @@ -0,0 +1,9 @@ +-- 020: Harness run detail — link reactive_runs to harness_runs and +-- capture the rendered prompt + raw response so the Web UI can show +-- operators exactly what the agent saw and what it replied. + +ALTER TABLE harness_runs ADD COLUMN reactive_run_id INTEGER; +ALTER TABLE harness_runs ADD COLUMN prompt TEXT; +ALTER TABLE harness_runs ADD COLUMN response TEXT; + +CREATE INDEX idx_harness_runs_reactive ON harness_runs(reactive_run_id); diff --git a/web/src/lib/components/MessageList.svelte b/web/src/lib/components/MessageList.svelte index bc3387e..0a60938 100644 --- a/web/src/lib/components/MessageList.svelte +++ b/web/src/lib/components/MessageList.svelte @@ -2,6 +2,7 @@ import { openThread } from '$lib/stores/thread'; import MessageBody from '$lib/components/MessageBody.svelte'; import AttachmentPreview from '$lib/components/AttachmentPreview.svelte'; + import ReactionPills from '$lib/components/ReactionPills.svelte'; type Attachment = { hash: string; @@ -11,6 +12,12 @@ is_image: boolean; }; + type ReactionEntry = { + agent_name: string; + reaction: string; + metadata?: Record; + }; + type Message = { id: number; conversation_id: number; @@ -22,6 +29,7 @@ created_at: string; reply_count?: number; attachments?: Attachment[]; + reactions?: ReactionEntry[]; }; let { messages = [], showConversationLink = false, agentTypes = {} }: { messages: Message[]; showConversationLink?: boolean; agentTypes?: Record } = $props(); @@ -118,6 +126,15 @@
+ + + + {#if msg.reactions && msg.reactions.length > 0} +
e.stopPropagation()} role="presentation"> + +
+ {/if} + {#if msg.attachments && msg.attachments.length > 0}
diff --git a/web/src/routes/runs/+page.svelte b/web/src/routes/runs/+page.svelte index 1b275f7..f9449b9 100644 --- a/web/src/routes/runs/+page.svelte +++ b/web/src/routes/runs/+page.svelte @@ -197,11 +197,14 @@
{run.error_log}
{/if} - {#if run.status === 'failed'} - - {/if} +
+ View full details → + {#if run.status === 'failed'} + + {/if} +
{/if} @@ -451,4 +454,11 @@ .retry-btn:hover { opacity: 0.9; } + + .detail-actions { + margin-top: 0.75rem; + display: flex; + align-items: center; + gap: 1rem; + } diff --git a/web/src/routes/runs/[id]/+page.svelte b/web/src/routes/runs/[id]/+page.svelte new file mode 100644 index 0000000..c013f7a --- /dev/null +++ b/web/src/routes/runs/[id]/+page.svelte @@ -0,0 +1,393 @@ + + +
+ + + + + Back to runs + + + {#if loading} +
Loading run #{runId}…
+ {:else if error} +
+

Failed to load run

+

{error}

+
+ {:else if run} + +
+
+
+ +
+

+ {run.agent_name} + #{run.id} +

+

+ triggered by @{run.trigger_from} + · {run.trigger_event} + · depth {run.trigger_depth} +

+
+
+
+ {run.status} + {#if harnessRun} + {harnessRun.backend} + {/if} +
+
+ + +
+
+

Duration

+

{formatDuration(run.duration_ms)}

+
+
+

Tokens in → out

+

+ {harnessRun ? `${formatTokens(harnessRun.tokens_in)} → ${formatTokens(harnessRun.tokens_out)}` : '—'} +

+
+
+

Cost

+

+ {harnessRun ? formatCost(harnessRun.cost_usd) : '—'} +

+
+
+

Exit

+

+ {harnessRun?.exit_code != null ? harnessRun.exit_code : '—'} +

+
+ {#if run.started_at} +
+

Started

+

{new Date(run.started_at).toLocaleString()}

+
+ {/if} + {#if run.completed_at} +
+

Completed

+

{new Date(run.completed_at).toLocaleString()}

+
+ {/if} + {#if harnessRun?.trace_id} +
+

Trace ID

+ +
+ {/if} +
+
+ + + {#if triggerMessage} +
+
+ + + +

Triggering message

+ #{triggerMessage.id} +
+
+

+ from @{triggerMessage.from_agent} + {#if triggerMessage.to_agent} → @{triggerMessage.to_agent}{/if} + · priority {triggerMessage.priority} +

+
{triggerMessage.body}
+
+
+ {/if} + + +
+
+ + + + +

What the model saw

+
+
+ {#if harnessConfigParsed?.gemini_md} +
+ + GEMINI.md (agent's current system instructions) + {harnessConfigParsed.gemini_md.length} bytes + +
{harnessConfigParsed.gemini_md}
+
+ {/if} + + {#if harnessConfigParsed?.claude_md} +
+ + CLAUDE.md (agent's current system instructions) + {harnessConfigParsed.claude_md.length} bytes + +
{harnessConfigParsed.claude_md}
+
+ {/if} + + {#if harnessRun?.prompt} +
+ + Rendered prompt (captured from subprocess workdir) + {harnessRun.prompt.length} bytes + +
{harnessRun.prompt}
+
+ {:else} +

No captured prompt — this run's backend does not record prompts, or the wrapper did not write prompt.txt.

+ {/if} +
+
+ + +
+
+ + + +

What the model said

+
+
+ {#if harnessRun?.response} +
{harnessRun.response}
+ {:else if harnessRun?.logs_excerpt} +

No captured response field — showing bounded logs excerpt instead.

+
{harnessRun.logs_excerpt}
+ {:else if run.error_log} +

Error log

+
{run.error_log}
+ {:else} +

Nothing captured.

+ {/if} +
+
+ + + {#if outgoingMessage} +
+
+ + + +

Outgoing message

+ #{outgoingMessage.id} +
+
+

+ @{run.agent_name} + {#if outgoingMessage.to_agent} → @{outgoingMessage.to_agent}{/if} + · {outgoingMessage.status} +

+
{outgoingMessage.body}
+
+
+ {/if} + + +
+
+

Metadata

+
+
+
reactive_run.id {run.id}
+ {#if harnessRun} +
harness_run.run_id {harnessRun.run_id}
+
backend {harnessRun.backend}
+ {#if harnessRun.session_id} +
session_id {harnessRun.session_id}
+ {/if} + {#if harnessRun.tokens_cached > 0} +
tokens_cached {formatTokens(harnessRun.tokens_cached)}
+ {/if} + {/if} + {#if run.k8s_job_name} +
k8s_job {run.k8s_job_name}
+ {/if} + {#if agent} +
agent.harness_name {agent.harness_name || '—'}
+
max_trigger_depth {agent.max_trigger_depth}
+
daily_budget {agent.daily_trigger_budget}
+
cooldown_s {agent.cooldown_seconds}
+ {/if} +
+
+ {/if} +