Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
72 changes: 47 additions & 25 deletions cmd/agentloop/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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")
Expand All @@ -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) {
Expand All @@ -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
Expand All @@ -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)
}
Expand Down Expand Up @@ -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})
Expand All @@ -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})
Expand Down
123 changes: 109 additions & 14 deletions cmd/agentloop/main_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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
}
}
Expand Down Expand Up @@ -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
}
Expand Down
1 change: 1 addition & 0 deletions docs/PRD.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.)
Expand Down
49 changes: 47 additions & 2 deletions internal/loop/runner.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
}

Expand Down Expand Up @@ -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
Expand All @@ -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
Expand Down Expand Up @@ -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:
Expand Down
Loading