diff --git a/cmd/agentloop/main.go b/cmd/agentloop/main.go index da7059f..5ca91df 100644 --- a/cmd/agentloop/main.go +++ b/cmd/agentloop/main.go @@ -273,6 +273,15 @@ func (s *Server) submitRun(w http.ResponseWriter, r *http.Request) { return } runID := generateRunID() + // Register the run the instant the caller has a run id — before any + // path can exit. Every exit writes its outcome through + // storeRunIfPresent, which refuses to create a row, so a run that is + // never registered loses its result: a blocked goal would 201 with a + // run id that resolves to nothing. + s.mu.Lock() + s.runs[runID] = loop.RunResult{RunID: runID, State: loop.StateQueued} + s.mu.Unlock() + maxSteps := loop.MaxSteps if body.MaxSteps != nil { maxSteps = *body.MaxSteps @@ -315,9 +324,7 @@ func (s *Server) submitRun(w http.ResponseWriter, r *http.Request) { ScreenErrors: []string{err.Error()}, } blocked.Success = boolPtr(false) - s.mu.Lock() - s.runs[runID] = blocked - s.mu.Unlock() + s.storeRunIfPresent(runID, blocked) // The run id must reach the caller even when the goal is blocked: // a 201 with no body is a client that cannot look up what happened. w.Header().Set("Content-Type", "application/json") @@ -342,30 +349,53 @@ func (s *Server) submitRun(w http.ResponseWriter, r *http.Request) { }}, } held.Success = boolPtr(false) - s.mu.Lock() - s.runs[runID] = held - s.mu.Unlock() + s.storeRunIfPresent(runID, held) w.Header().Set("Content-Type", "application/json") w.WriteHeader(http.StatusCreated) writeJSON(w, map[string]any{"run_id": runID, "state": held.State, "exit_reason": held.ExitReason}) return } + // The first of three forms of visibility: the run is registered, so a + // poll resolves and the P75 kill endpoint has something to signal for + // the whole run rather than only once it ends. s.mu.Lock() s.runners[runID] = runner s.gates[runID] = gate s.mu.Unlock() + // The second form: every step boundary refreshes what a watcher + // reads, so a poll shows the steps that have run so far rather than + // nothing at all. Attached before Run starts, or the first boundary + // races the goroutine. + runner.WithProgress(func(partial loop.RunResult) { + s.storeRunIfPresent(runID, partial) + }) go func() { result, err := runner.Run(context.Background()) if err != nil { result.State = loop.StateFailed } - s.mu.Lock() - s.runs[runID] = result - s.mu.Unlock() + // The third form: the finished result. Without this the run would + // stay at its last in-flight snapshot forever, so a watcher would + // never see the terminal state, the exit reason, or the synthesis. + s.storeRunIfPresent(runID, result) }() w.Header().Set("Content-Type", "application/json") w.WriteHeader(http.StatusCreated) - writeJSON(w, map[string]any{"run_id": runID, "state": loop.StateThinking}) + writeJSON(w, map[string]any{"run_id": runID, "state": loop.StateQueued}) +} + +// storeRunIfPresent records a run's result only while the run is still +// known. DELETE removes the key, and a run's final write lands *after* +// that — so an unconditional store let a deleted run reappear, on the +// terminal write and on every progress report in between. A run an +// operator removed stays removed. +func (s *Server) storeRunIfPresent(runID string, result loop.RunResult) { + s.mu.Lock() + defer s.mu.Unlock() + if _, ok := s.runs[runID]; !ok { + return + } + s.runs[runID] = result } func (s *Server) getRun(w http.ResponseWriter, r *http.Request) { @@ -384,8 +414,8 @@ func (s *Server) getRun(w http.ResponseWriter, r *http.Request) { func (s *Server) killRun(w http.ResponseWriter, r *http.Request) { runID := extractRunID(r.URL.Path) s.mu.Lock() + defer s.mu.Unlock() result, ok := s.runs[runID] - s.mu.Unlock() if !ok { http.Error(w, `{"error":"run not found"}`, http.StatusNotFound) return @@ -395,18 +425,14 @@ func (s *Server) killRun(w http.ResponseWriter, r *http.Request) { if result.PartialSynthesis == "" { result.PartialSynthesis = "killed by operator" } - // Signal the live runner's kill channel so the running - // goroutine exits within one step (P75). Without this the - // stored state is stamped but the loop keeps going. - s.mu.Lock() + // Signal the live runner's kill channel so the running goroutine exits + // within one step (P75). The runner is read under the same lock as the + // run state: if the run is already known, the runner should be too. runner := s.runners[runID] - s.mu.Unlock() + s.runs[runID] = result if runner != nil { runner.Kill() } - s.mu.Lock() - s.runs[runID] = result - s.mu.Unlock() w.Header().Set("Content-Type", "application/json") writeJSON(w, result) } @@ -484,9 +510,7 @@ func (s *Server) submitApproval(w http.ResponseWriter, r *http.Request) { if err != nil { result.State = loop.StateFailed } - s.mu.Lock() - s.runs[runID] = result - s.mu.Unlock() + s.storeRunIfPresent(runID, result) } w.Header().Set("Content-Type", "application/json") writeJSON(w, map[string]any{"approved": approved}) @@ -509,9 +533,7 @@ func (s *Server) submitApprovalByID(w http.ResponseWriter, r *http.Request) { if err != nil { result.State = loop.StateFailed } - s.mu.Lock() - s.runs[runID] = result - s.mu.Unlock() + s.storeRunIfPresent(runID, result) } w.Header().Set("Content-Type", "application/json") writeJSON(w, map[string]any{"approved": approved}) diff --git a/cmd/agentloop/main_test.go b/cmd/agentloop/main_test.go index c4efce0..3f10331 100644 --- a/cmd/agentloop/main_test.go +++ b/cmd/agentloop/main_test.go @@ -26,6 +26,87 @@ func newTestServer(t *testing.T) (*httptest.Server, *Server) { return srv, s } +// TestRunVisibleWhileInFlight is the check for the run-visibility bug: +// submit used to return a run id that no endpoint could resolve until the +// loop finished, so a poll 404'd for the whole run and the P75 kill was +// unreachable. The run must be readable the moment the 201 lands. +func TestRunVisibleWhileInFlight(t *testing.T) { + srv, _ := newTestServer(t) + defer srv.Close() + + resp, err := http.Post(srv.URL+"/v1/runs", "application/json", + strings.NewReader(`{"goal":"visibility","context":"test"}`)) + if err != nil { + t.Fatalf("POST: %v", err) + } + var submit map[string]any + json.NewDecoder(resp.Body).Decode(&submit) + resp.Body.Close() + runID, _ := submit["run_id"].(string) + + // No polling, no sleep: the id the caller was just handed must resolve. + getResp, err := http.Get(srv.URL + "/v1/runs/" + runID) + if err != nil { + t.Fatalf("GET immediately after submit: %v", err) + } + defer getResp.Body.Close() + if getResp.StatusCode != http.StatusOK { + t.Fatalf("GET immediately after submit = %d, want 200 (run in flight must be readable)", getResp.StatusCode) + } + var got RunResultResponse + if err := json.NewDecoder(getResp.Body).Decode(&got); err != nil { + t.Fatalf("decode: %v", err) + } + if got.RunID != runID { + t.Errorf("run_id = %q, want %q", got.RunID, runID) + } + + // The kill endpoint must also reach it while it is still running. + killResp, err := http.Post(srv.URL+"/v1/runs/"+runID+"/kill", "application/json", nil) + if err != nil { + t.Fatalf("POST kill: %v", err) + } + defer killResp.Body.Close() + if killResp.StatusCode != http.StatusOK { + t.Fatalf("kill while in flight = %d, want 200 (P75 must reach a live run)", killResp.StatusCode) + } + var killed map[string]any + json.NewDecoder(killResp.Body).Decode(&killed) + if killed["state"] != string(loop.StateKilled) { + t.Errorf("state = %v, want %v", killed["state"], loop.StateKilled) + } + // And deleting a finished run stays deleted: the loop goroutine writes + // its own result after the run ends, and that write must not resurrect + // a run an operator removed (it used to, by re-inserting the map key). + delReq, _ := http.NewRequest("DELETE", srv.URL+"/v1/runs/"+runID, nil) + delResp, err := http.DefaultClient.Do(delReq) + if err != nil { + t.Fatalf("DELETE: %v", err) + } + delResp.Body.Close() + if delResp.StatusCode != http.StatusNoContent { + t.Fatalf("DELETE = %d, want 204", delResp.StatusCode) + } + waitForRunGone(t, srv, runID) + waitForRunGone(t, srv, runID) +} + +// waitForRunGone polls until the run id stops resolving. +func waitForRunGone(t *testing.T, srv *httptest.Server, runID string) { + t.Helper() + for i := 0; i < 200; i++ { + resp, err := http.Get(srv.URL + "/v1/runs/" + runID) + if err == nil { + resp.Body.Close() + if resp.StatusCode == http.StatusNotFound { + return + } + } + time.Sleep(10 * time.Millisecond) + } + t.Fatal("run still resolves after delete") +} + // routes registers all HTTP routes on a ServeMux — shared between // production main() and tests so handlers stay in sync. func routes(s *Server) *http.ServeMux { @@ -149,23 +230,34 @@ func TestLive_KillRun(t *testing.T) { } } -// waitForRun polls GET /v1/runs/{id} until the run result -// is stored in the server. Used by kill/delete/console tests -// to avoid a race with the background goroutine. +func stateFor(t *testing.T, srv *httptest.Server, runID string, i int) string { + t.Helper() + resp, err := http.Get(srv.URL + "/v1/runs/" + runID) + if err != nil { + t.Fatalf("GET run (i=%d): %v", i, err) + } + defer resp.Body.Close() + var result RunResultResponse + json.NewDecoder(resp.Body).Decode(&result) + return result.State +} + +// waitForRun polls GET /v1/runs/{id} until the run reaches a terminal state. +// A run in flight is readable now, so "still thinking" is not "not stored": +// the delete endpoints already gate on the run id, not the terminal state. func waitForRun(t *testing.T, srv *httptest.Server, runID string) { t.Helper() - for i := 0; i < 50; i++ { - getResp, gerr := http.Get(srv.URL + "/v1/runs/" + runID) - if gerr == nil && getResp.StatusCode == http.StatusOK { - getResp.Body.Close() + for i := 0; i < 300; i++ { + switch stateFor(t, srv, runID, i) { + case "", string(loop.StateThinking), string(loop.StateQueued), + string(loop.StateActing), string(loop.StateEvaluating): + // in flight, keep waiting + default: return } - if getResp != nil { - getResp.Body.Close() - } - time.Sleep(10 * time.Millisecond) + time.Sleep(20 * time.Millisecond) } - t.Fatal("run never stored") + t.Fatal("run never reached a terminal state") } // TestDeleteRun submits a run, deletes it, and confirms 404 after. @@ -454,7 +546,10 @@ func TestM6_Demo_HITLWithEvals(t *testing.T) { var result RunResultResponse json.NewDecoder(getResp.Body).Decode(&result) getResp.Body.Close() - if result.State != "" { + // "thinking" means the run is in flight — since the visibility fix + // that is a real answer, not the absence of one, so wait for a + // state the run has actually settled into. + if result.State != "" && result.State != string(loop.StateThinking) { break } } @@ -558,7 +653,7 @@ func TestM5_GateRegisteredForNewRun(t *testing.T) { var result RunResultResponse der := json.NewDecoder(getResp.Body).Decode(&result) getResp.Body.Close() - if der == nil && result.State != "" { + if der == nil && result.State != "" && result.State != string(loop.StateThinking) { state = result.State break } diff --git a/docs/PRD.md b/docs/PRD.md index fc0f229..4c7dd24 100644 --- a/docs/PRD.md +++ b/docs/PRD.md @@ -748,6 +748,7 @@ Written the way an unfriendly reviewer would write it, then answered. Every find **Read next.** §13.1 (scope → milestones), §17 (defaults), §18 (where to discount the source), §22 (this document's own weaknesses). +* Last updated: 2026-09-21 (**A submitted run is visible while it is in flight.** `POST /v1/runs` registered nothing until the loop returned, so the run id it handed back resolved to `404` for the whole run: the documented submit-then-poll flow could not poll, and `POST /v1/runs/{id}/kill` (P75, NFR-1) had no run to signal. The endpoint now registers the run before any exit, the runner publishes a partial at every step boundary (`LoopRunner.WithProgress`), and every store goes through `storeRunIfPresent`, which refuses to create a row — a `DELETE`d run stays deleted instead of being resurrected by the loop's own final write. Found while fixing the first: a blocked goal's `201` pointed at a run that was never stored at all, so the goal-screen verdict was unreadable too. §6's reconciliation note is unchanged; the reviewable part is `cmd/agentloop/main_test.go` (`TestRunVisibleWhileInFlight`) and `TestGuardrailScreenFiresInProduction`, both of which fail if the registration line is removed.) * Last updated: 2026-09-21 (**The guardrail screen is wired.** §7.2 and NFR-1b claimed every message was screened; nothing in the service ever set a screen — only tests did, so the containment claim was false in production. `internal/guardrail` speaks `POST /v1/systemone`, the goal is judged before the first step and every tool result before it reaches the model, a screen that cannot run is recorded rather than passed or fail-closed, and a blocked goal returns its `run_id` instead of an empty 201. The reply-out direction and the ~740ms/call budget entry are still open.) * Last updated: 2026-09-21 (Tier routing reaches the wire. `ChatTier(ctx, tier, msgs…)` replaced the tier-less `Chat` on the model interface, `onegw.Client.WithTiers` maps tier→combo, and the three tiers are now planning/execution/synthesis instead of planning/execution/tiny — the old middle names were combos onegw does not ship. `Reply` records the combo asked for *and* the answering leg, so a fallback is visible. Verified against a stub gateway: three calls to `exec-combo`, one to `synth-combo`.) * 2026-09-21 — (The loop **decides its own steps**. `internal/loop/reason.go` makes one model call per step with the previous step's verbatim result, and the answer (`{tool,args,why,done}`) is what runs — so the loop observes before it reasons, which is the half it was missing. `done` exits `success`/`goal_met`, a new exit reason for the goal predicate firing rather than a bound. A resume replays the approved decision instead of re-asking (a second call can choose a different tool, so the operator would have approved one action and a different one would run). `StepRecord` now carries its `Args` and `Why`. No model client still means the deterministic rotation, unchanged.) diff --git a/internal/loop/runner.go b/internal/loop/runner.go index 1c6c416..6e6943b 100644 --- a/internal/loop/runner.go +++ b/internal/loop/runner.go @@ -159,6 +159,12 @@ type LoopRunner struct { // a DIFFERENT tool, which would run something the operator never saw, // under an approval for something else. Cost is the smaller reason. heldChoice *StepChoice + // onProgress, when set, receives a copy of the partial result at every + // step boundary. It exists because the result reaches the caller only + // at the END of the loop: without it a run in flight is invisible to + // everyone outside it, so a poll 404s and the P75 kill is unreachable + // for the whole run. + onProgress func(RunResult) lastResult RunResult // partial result at pause (M5 resume) } @@ -349,9 +355,35 @@ func (r *LoopRunner) Kill() { } } +// WithProgress returns the runner with fn called at every step boundary +// (and once before the first step), handing the live partial result to +// whoever is watching the run. Without it a run in flight has no result +// to read: the store is empty until the loop returns, so a poll 404s and +// a kill has nothing to stamp. +func (r *LoopRunner) WithProgress(fn func(RunResult)) *LoopRunner { + r.onProgress = fn + return r +} + +// reportProgress hands the current result to the progress sink, if one is +// set. Steps is copied so a later append in the loop cannot race a reader +// that is already encoding the snapshot. +func (r *LoopRunner) reportProgress(result RunResult) { + if r.onProgress == nil { + return + } + snapshot := result + snapshot.Steps = append([]StepRecord(nil), result.Steps...) + r.onProgress(snapshot) +} + // Run executes the bounded loop and returns the RunResult. func (r *LoopRunner) Run(ctx context.Context) (RunResult, error) { - return r.runWith(ctx, true) + result, err := r.runWith(ctx, true) + // One last publish after the loop ends, so a watcher's final view is + // the finished result and not the second-to-last step boundary. + r.reportProgress(result) + return result, err } // Resume re-enters a run that paused for operator approval @@ -364,7 +396,9 @@ func (r *LoopRunner) Resume(ctx context.Context) (RunResult, error) { if r.pausedStep < 0 { return r.lastResult, nil } - return r.runWith(ctx, false) + result, err := r.runWith(ctx, false) + r.reportProgress(result) + return result, err } // runWith builds the run state and delegates the loop to @@ -421,7 +455,18 @@ func (r *LoopRunner) runLoop(ctx context.Context, result RunResult) (RunResult, if r.pausedStep >= 0 { startStep = r.pausedStep } + // Publish the partial before the first ceiling check. A run in flight + // must be readable: the server stores the result only when the loop + // returns, so without this a poll 404s and a kill has no state to + // stamp for the whole run (the documented "submit, then poll" flow, + // and the P75 kill switch). + r.reportProgress(result) + defer r.reportProgress(result) for step := startStep; step < r.cfg.MaxSteps; step++ { + // Publish at every step boundary, so a watcher sees the step that + // just ran rather than a view lagging a step behind. The deferred + // report fires last, so the final publish is the finished result. + r.reportProgress(result) // --- P75 kill switch: check at every iteration boundary --- select { case <-r.killCh: