diff --git a/internal/agentproxy/opencode/bridge.go b/internal/agentproxy/opencode/bridge.go index 60facd68d..867f33f92 100644 --- a/internal/agentproxy/opencode/bridge.go +++ b/internal/agentproxy/opencode/bridge.go @@ -161,16 +161,17 @@ func (f *Facade) handleSessionCreate(w http.ResponseWriter, r *http.Request) { f.writeJSON(w, f.sessionObject()) } -// handleSessionList answers GET /session: empty until the client has -// created/used the session, then a single-entry list. An empty list is the -// "fresh server" state the TUI expects on first attach (verified in the M0 -// capture); returning the session unconditionally would make every attach -// look like a resume. +// handleSessionList answers GET /session: the bridged session appears when +// this proxy has used it (sessionCreated) or when the hosted session has +// replayable history — that's what lets `opencode attach --continue` and the +// session picker resume prior turns through a freshly started proxy. A truly +// fresh session lists empty, the "fresh server" state the TUI expects on +// first attach (verified in the M0 capture). func (f *Facade) handleSessionList(w http.ResponseWriter, r *http.Request) { f.mu.Lock() created := f.sessionCreated f.mu.Unlock() - if !created { + if !created && !f.hasHistory(r.Context()) { f.writeJSON(w, []any{}) return } @@ -377,6 +378,8 @@ func (f *Facade) translateEvent(ev godo.HostedAgentEvent, ts *turnState, ew *eve case godo.HostedAgentEventKindRunCompleted: defer f.dropTurn(ev.RunID) + // The finished turn is durable history now; the cache predates it. + f.invalidateHistory() // Finalize the text part with its full content before idling — the // real server does (deltas stream, then the part's final state // re-carries the whole text; plano's adapter must drop that as a @@ -403,6 +406,7 @@ func (f *Facade) translateEvent(ev godo.HostedAgentEvent, ts *turnState, ew *eve case godo.HostedAgentEventKindRunFailed: defer f.dropTurn(ev.RunID) + f.invalidateHistory() var payload struct { Message string `json:"message"` } diff --git a/internal/agentproxy/opencode/facade.go b/internal/agentproxy/opencode/facade.go index 1496ab84c..2e91301b8 100644 --- a/internal/agentproxy/opencode/facade.go +++ b/internal/agentproxy/opencode/facade.go @@ -12,6 +12,7 @@ package opencode import ( + "context" "encoding/json" "fmt" "log" @@ -72,11 +73,31 @@ type Facade struct { // session, so the session list grows its single entry (see // handleSessionList). sessionCreated bool + + // History cache (see history() in history.go). histMu guards only the + // fields — the slow replay_only fetch runs outside it so invalidation + // never blocks behind it; histFetching/histDone single-flight concurrent + // fetchers, and histGen detects an invalidation that overlapped a fetch. + histMu sync.Mutex + hist []historyMessage + histValid bool + histGen int + histFetching bool + histDone chan struct{} } // ServeHTTP implements http.Handler. func (f *Facade) ServeHTTP(w http.ResponseWriter, r *http.Request) { - f.handlerOnce.Do(f.buildMux) + f.handlerOnce.Do(func() { + f.buildMux() + // Warm the history cache off the first request — the TUI's + // /global/health preflight lands well before the attach burst, so + // the slow replay (see history()) usually finishes before the burst + // asks for the session list or messages. + if f.Sessions != nil { + go func() { _, _ = f.history(context.Background()) }() + } + }) f.mux.ServeHTTP(w, r) } @@ -163,15 +184,18 @@ func (f *Facade) buildMux() { mux.HandleFunc("GET /experimental/workspace", f.json([]any{})) mux.HandleFunc("GET /experimental/workspace/status", f.json([]any{})) - // The bridged session: list/create/get plus the prompt bridge (M2). - // History (GET .../message) returns empty until M3 serves it from a - // replay_only stream pass. + // The bridged session: list/create/get plus the prompt bridge (M2) and + // history (M3, served from a replay_only stream pass). diff/todo are the + // two lookups the TUI fires after every turn (post-idle refresh) and on + // resume — real server returns empty arrays for a no-edits session. mux.HandleFunc("GET /session", f.handleSessionList) mux.HandleFunc("POST /session", f.handleSessionCreate) mux.HandleFunc("GET /session/{id}", func(w http.ResponseWriter, r *http.Request) { f.writeJSON(w, f.sessionObject()) }) - mux.HandleFunc("GET /session/{id}/message", f.json([]any{})) + mux.HandleFunc("GET /session/{id}/message", f.handleMessageList) + mux.HandleFunc("GET /session/{id}/diff", f.json([]any{})) + mux.HandleFunc("GET /session/{id}/todo", f.json([]any{})) mux.HandleFunc("POST /session/{id}/prompt_async", f.handlePromptAsync) // The synchronous variant officially awaits the reply; bridging that // faithfully would block a request goroutine for a whole turn. Current diff --git a/internal/agentproxy/opencode/facade_test.go b/internal/agentproxy/opencode/facade_test.go index bd9405e7d..676a804b3 100644 --- a/internal/agentproxy/opencode/facade_test.go +++ b/internal/agentproxy/opencode/facade_test.go @@ -4,7 +4,6 @@ import ( "bufio" "encoding/json" "net/http" - "net/http/httptest" "strings" "testing" @@ -12,16 +11,11 @@ import ( "github.com/stretchr/testify/require" ) -func newTestFacade() *Facade { - return &Facade{SessionID: "sess-1", Dir: "/tmp/ws"} -} - // The attach burst captured from a real `opencode attach` (TestedVersion). // Every route the TUI hits at startup must answer 200 with valid JSON — // a 404 here is exactly the "new client version wants more" drift signal. func TestAttachBurstRoutesAnswer(t *testing.T) { - srv := httptest.NewServer(newTestFacade()) - defer srv.Close() + srv, _ := newBridgedFacade(t) routes := []string{ "/global/health", @@ -65,8 +59,7 @@ func TestAttachBurstRoutesAnswer(t *testing.T) { } func TestHealthReportsTestedVersion(t *testing.T) { - srv := httptest.NewServer(newTestFacade()) - defer srv.Close() + srv, _ := newBridgedFacade(t) resp, err := http.Get(srv.URL + "/global/health") require.NoError(t, err) @@ -131,8 +124,7 @@ func TestSecondEventStreamConsumerConflicts(t *testing.T) { } func TestShareIsRefused(t *testing.T) { - srv := httptest.NewServer(newTestFacade()) - defer srv.Close() + srv, _ := newBridgedFacade(t) resp, err := http.Post(srv.URL+"/session/sess-1/share", "application/json", nil) require.NoError(t, err) @@ -141,8 +133,7 @@ func TestShareIsRefused(t *testing.T) { } func TestUnknownRouteIs404(t *testing.T) { - srv := httptest.NewServer(newTestFacade()) - defer srv.Close() + srv, _ := newBridgedFacade(t) resp, err := http.Get(srv.URL + "/no/such/route") require.NoError(t, err) @@ -154,8 +145,7 @@ func TestUnknownRouteIs404(t *testing.T) { // /config/providers points at the one provider/model pair every other // catalog route advertises. func TestSyntheticCatalogIsConsistent(t *testing.T) { - srv := httptest.NewServer(newTestFacade()) - defer srv.Close() + srv, _ := newBridgedFacade(t) resp, err := http.Get(srv.URL + "/config/providers") require.NoError(t, err) diff --git a/internal/agentproxy/opencode/history.go b/internal/agentproxy/opencode/history.go new file mode 100644 index 000000000..151b4d116 --- /dev/null +++ b/internal/agentproxy/opencode/history.go @@ -0,0 +1,305 @@ +package opencode + +import ( + "context" + "encoding/json" + "fmt" + "net/http" + "strconv" + + "github.com/digitalocean/godo" +) + +// M3: history. `opencode attach` replays a session by fetching +// GET /session/{id}/message?limit=N itself (client-driven — the server never +// pushes history), so re-attaching shows prior turns. The facade serves that +// endpoint from a one-shot replay_only StreamSession pass over the session's +// durable event history, reconstructed into the [{info, parts}] shape +// captured from a real server (see the capture doc). Replays are cached for +// the proxy's lifetime (see history()): fetched once, single-flighted, warmed +// at the facade's first request, and invalidated when a live turn completes. + +// historyMessage is one reconstructed message: the {info, parts} pair the +// history endpoint returns. +type historyMessage struct { + Info map[string]any `json:"info"` + Parts []map[string]any `json:"parts"` +} + +// historyTurn accumulates one hosted run while replaying. +type historyTurn struct { + runID string + startMs int64 + // doneMs is the run.completed/run.failed event time — the turn's real + // duration. lastMs (any event) is only the fallback for turns that never + // completed: stragglers like late run.log events arrive long after + // completion and would inflate the rendered duration (a "20m" turn + // badge on a 2s turn, found live). + doneMs int64 + lastMs int64 + userText string + text []byte + reasoning []byte + failed bool +} + +func (ht *historyTurn) endMs() int64 { + if ht.doneMs != 0 { + return ht.doneMs + } + return ht.lastMs +} + +// history returns the reconstructed message history, cached. The harness's +// replay_only stream takes ~8s to close after it has caught up (a +// server-side linger, measured against the dev stack — 68 events transfer +// instantly and the close arrives at exactly 8.0s), and an attach fetches +// history twice (list gating + the message list), so uncached history made +// re-attach feel ~16s slow. +// +// The slow fetch runs OUTSIDE histMu: the event loop calls +// invalidateHistory() when a turn completes, and holding the lock across the +// ~8s replay would stall event translation behind it. Single-flight rides +// histFetching/histDone instead — concurrent callers wait for the in-flight +// replay (or their ctx), then re-check the cache. A fetch that an +// invalidation overlapped (histGen moved) still returns its snapshot to the +// caller — it was a valid point-in-time read — but is not cached, so the +// next call replays fresh. +// +// The cache is warmed at the first request the facade sees (the TUI's +// /global/health preflight fires well before the attach burst). +func (f *Facade) history(ctx context.Context) ([]historyMessage, error) { + for { + f.histMu.Lock() + if f.histValid { + msgs := f.hist + f.histMu.Unlock() + return msgs, nil + } + if f.histFetching { + done := f.histDone + f.histMu.Unlock() + select { + case <-done: + continue + case <-ctx.Done(): + return nil, ctx.Err() + } + } + f.histFetching = true + f.histDone = make(chan struct{}) + gen := f.histGen + f.histMu.Unlock() + + msgs, err := f.fetchHistory(ctx) + + f.histMu.Lock() + f.histFetching = false + close(f.histDone) + if err == nil && f.histGen == gen { + f.hist, f.histValid = msgs, true + } + f.histMu.Unlock() + return msgs, err + } +} + +// invalidateHistory drops the cache; the next history() call replays fresh. +// Only ever a fast lock — never blocked behind an in-flight fetch (the event +// loop calls this on turn completion and must not stall). +func (f *Facade) invalidateHistory() { + f.histMu.Lock() + f.histGen++ + f.histValid = false + f.hist = nil + f.histMu.Unlock() +} + +// fetchHistory replays the durable event history and reconstructs completed +// turns. Only text is reconstructed in M3 — tool-call parts are M4. Callers +// go through history() for the cache; this always hits the harness. +func (f *Facade) fetchHistory(ctx context.Context) ([]historyMessage, error) { + stream, err := f.Sessions.StreamSession(ctx, f.SessionID, &godo.HostedAgentSessionStreamOptions{ + ReplayOnly: true, + }) + if err != nil { + return nil, err + } + defer stream.Close() + + var order []string + turns := map[string]*historyTurn{} + for stream.Next() { + ev := stream.Current() + if ev.RunID == "" { + continue + } + ht, ok := turns[ev.RunID] + if !ok { + ht = &historyTurn{runID: ev.RunID} + turns[ev.RunID] = ht + order = append(order, ev.RunID) + } + at := eventTimeMs(ev) + if ht.startMs == 0 { + ht.startMs = at + } + if at > ht.lastMs { + ht.lastMs = at + } + switch ev.Kind { + case godo.HostedAgentEventKindRunStarted: + // The proto names this field user_input, but the SSE wire's SPI + // envelope delivers it as "agent" (observed live against the dev + // stack: data:{"agent":""}). Accept both so a wire + // rename toward the proto name doesn't silently drop history. + var payload struct { + UserInput string `json:"user_input"` + Agent string `json:"agent"` + } + _ = json.Unmarshal(ev.Payload, &payload) + ht.userText = payload.UserInput + if ht.userText == "" { + ht.userText = payload.Agent + } + case godo.HostedAgentEventKindTokenChunk: + var payload struct { + Text string `json:"text"` + IsReasoning bool `json:"is_reasoning"` + } + if err := json.Unmarshal(ev.Payload, &payload); err != nil { + continue + } + if payload.IsReasoning { + ht.reasoning = append(ht.reasoning, payload.Text...) + } else { + ht.text = append(ht.text, payload.Text...) + } + case godo.HostedAgentEventKindRunCompleted: + ht.doneMs = at + case godo.HostedAgentEventKindRunFailed: + ht.failed = true + ht.doneMs = at + } + } + if err := stream.Err(); err != nil { + return nil, err + } + + // T allocation is monotonic across turns: two turns whose events share a + // millisecond would otherwise interleave (turn N+1's counter-0 user id + // sorts below turn N's counter-2 assistant id — caught by test, and real + // replays can burst events into one ms). Each turn gets a fresh ms slot + // in T-space at minimum. + var msgs []historyMessage + var lastT uint64 + for _, runID := range order { + ht := turns[runID] + t := ocTimeVal(ht.startMs) + if t <= lastT { + t = lastT + 0x1000 + } + lastT = t + msgs = append(msgs, f.turnMessages(ht, t)...) + } + return msgs, nil +} + +// turnMessages renders one replayed turn as its user and assistant messages. +// Ids are minted from the caller-allocated T (turn event time, forced +// monotonic across turns) with the real id encoding — history ids sort +// against live-minted ids in the same conversation, so a random id here +// reintroduces the shuffled-rendering bug the live path fixed. +func (f *Facade) turnMessages(ht *historyTurn, t uint64) []historyMessage { + sid := f.ocSessionID() + userMsgID := ocIDWithT("msg_", t, 0, runTail(ht.runID, "hu")) + asstMsgID := ocIDWithT("msg_", t, 2, runTail(ht.runID, "ha")) + + var msgs []historyMessage + if ht.userText != "" { + msgs = append(msgs, historyMessage{ + Info: map[string]any{ + "id": userMsgID, "sessionID": sid, + "role": "user", + "time": map[string]any{"created": ht.startMs}, + "agent": "build", + "model": map[string]any{"providerID": providerID, "modelID": modelID}, + "summary": map[string]any{"diffs": []any{}}, + }, + Parts: []map[string]any{{ + "id": ocIDWithT("prt_", t, 1, runTail(ht.runID, "hu")), "messageID": userMsgID, "sessionID": sid, + "type": "text", "text": ht.userText, + }}, + }) + } + + if len(ht.text) == 0 && len(ht.reasoning) == 0 && !ht.failed { + return msgs + } + info := map[string]any{ + "id": asstMsgID, "sessionID": sid, + "role": "assistant", + "time": map[string]any{"created": ht.startMs, "completed": ht.endMs()}, + "mode": "build", + "agent": "build", + "model": map[string]any{"providerID": providerID, "modelID": modelID}, + "path": map[string]any{"cwd": f.Dir, "root": f.Dir}, + "cost": 0, + "tokens": map[string]any{ + "input": 0, "output": 0, "reasoning": 0, + "cache": map[string]any{"read": 0, "write": 0}, + }, + } + if ht.userText != "" { + info["parentID"] = userMsgID + } + if !ht.failed { + info["finish"] = "stop" + } + var parts []map[string]any + partT := t + if len(ht.reasoning) > 0 { + parts = append(parts, map[string]any{ + "id": ocIDWithT("prt_", partT, 3, runTail(ht.runID, "hr")), "messageID": asstMsgID, "sessionID": sid, + "type": "reasoning", "text": string(ht.reasoning), + "time": map[string]any{"start": ht.startMs, "end": ht.endMs()}, + }) + } + if len(ht.text) > 0 { + parts = append(parts, map[string]any{ + "id": ocIDWithT("prt_", partT, 4, runTail(ht.runID, "ht")), "messageID": asstMsgID, "sessionID": sid, + "type": "text", "text": string(ht.text), + "time": map[string]any{"start": ht.startMs, "end": ht.endMs()}, + }) + } + msgs = append(msgs, historyMessage{Info: info, Parts: parts}) + return msgs +} + +// handleMessageList serves GET /session/{id}/message from a fresh replay. +// The TUI sends ?limit=N (100 on resume); the newest N messages win, order +// preserved (oldest first), matching a real server's response. +func (f *Facade) handleMessageList(w http.ResponseWriter, r *http.Request) { + msgs, err := f.history(r.Context()) + if err != nil { + http.Error(w, fmt.Sprintf("replaying session history failed: %v", err), http.StatusBadGateway) + return + } + if s := r.URL.Query().Get("limit"); s != "" { + if limit, err := strconv.Atoi(s); err == nil && limit >= 0 && len(msgs) > limit { + msgs = msgs[len(msgs)-limit:] + } + } + if msgs == nil { + msgs = []historyMessage{} + } + f.writeJSON(w, msgs) +} + +// hasHistory reports whether the session has any replayable turns — the +// session list includes the bridged session when it does, so `--continue` +// and the session picker can find it on a fresh proxy. +func (f *Facade) hasHistory(ctx context.Context) bool { + msgs, err := f.history(ctx) + return err == nil && len(msgs) > 0 +} diff --git a/internal/agentproxy/opencode/history_test.go b/internal/agentproxy/opencode/history_test.go new file mode 100644 index 000000000..3a081fe6b --- /dev/null +++ b/internal/agentproxy/opencode/history_test.go @@ -0,0 +1,230 @@ +package opencode + +import ( + "context" + "encoding/json" + "net/http" + "net/http/httptest" + "strings" + "testing" + "time" + + "github.com/digitalocean/doctl/do" + "github.com/digitalocean/doctl/internal/agentproxy/agentproxytest" + "github.com/digitalocean/godo" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// queueTwoTurnHistory fills the harness's replay queue with two completed +// turns (the second one with reasoning), the durable-history shape M3 +// reconstructs from. +func queueTwoTurnHistory(h *agentproxytest.Harness) { + h.QueueReplayHistory( + // Turn 1 uses the observed wire key ("agent"); turn 2 uses the proto + // name ("user_input") — both must reconstruct. + agentproxytest.Event{RunID: "run-1", Type: string(godo.HostedAgentEventKindRunStarted), Data: json.RawMessage(`{"agent":"first question"}`)}, + agentproxytest.Event{RunID: "run-1", Type: string(godo.HostedAgentEventKindTokenChunk), Data: json.RawMessage(`{"text":"first "}`)}, + agentproxytest.Event{RunID: "run-1", Type: string(godo.HostedAgentEventKindTokenChunk), Data: json.RawMessage(`{"text":"answer"}`)}, + agentproxytest.Event{RunID: "run-1", Type: string(godo.HostedAgentEventKindRunCompleted)}, + agentproxytest.Event{RunID: "run-2", Type: string(godo.HostedAgentEventKindRunStarted), Data: json.RawMessage(`{"user_input":"second question"}`)}, + agentproxytest.Event{RunID: "run-2", Type: string(godo.HostedAgentEventKindTokenChunk), Data: json.RawMessage(`{"text":"hmm","is_reasoning":true}`)}, + agentproxytest.Event{RunID: "run-2", Type: string(godo.HostedAgentEventKindTokenChunk), Data: json.RawMessage(`{"text":"second answer"}`)}, + agentproxytest.Event{RunID: "run-2", Type: string(godo.HostedAgentEventKindRunCompleted)}, + ) +} + +func getMessages(t *testing.T, url string) []historyMessage { + t.Helper() + resp, err := http.Get(url) + require.NoError(t, err) + defer resp.Body.Close() + require.Equal(t, http.StatusOK, resp.StatusCode) + var msgs []historyMessage + require.NoError(t, json.NewDecoder(resp.Body).Decode(&msgs)) + return msgs +} + +func TestHistoryReconstructsTurns(t *testing.T) { + srv, h := newBridgedFacade(t) + queueTwoTurnHistory(h) + + msgs := getMessages(t, srv.URL+"/session/ses_x/message") + require.Len(t, msgs, 4) // user+assistant per turn + + // Turn 1: user question then assistant answer, linked by parentID. + assert.Equal(t, "user", msgs[0].Info["role"]) + require.Len(t, msgs[0].Parts, 1) + assert.Equal(t, "first question", msgs[0].Parts[0]["text"]) + assert.Equal(t, "assistant", msgs[1].Info["role"]) + assert.Equal(t, msgs[0].Info["id"], msgs[1].Info["parentID"]) + require.Len(t, msgs[1].Parts, 1) + assert.Equal(t, "text", msgs[1].Parts[0]["type"]) + assert.Equal(t, "first answer", msgs[1].Parts[0]["text"]) + // Completed turns carry finish + time.completed (the TUI treats a + // message without them as still in flight). + assert.Equal(t, "stop", msgs[1].Info["finish"]) + tm, ok := msgs[1].Info["time"].(map[string]any) + require.True(t, ok) + assert.Contains(t, tm, "completed") + // The crash-critical assistant fields ride along in history too. + assert.Contains(t, msgs[1].Info, "tokens") + assert.Contains(t, msgs[1].Info, "cost") + + // Turn 2: reasoning reconstructs as its own part, before the text part. + require.Len(t, msgs[3].Parts, 2) + assert.Equal(t, "reasoning", msgs[3].Parts[0]["type"]) + assert.Equal(t, "hmm", msgs[3].Parts[0]["text"]) + assert.Equal(t, "second answer", msgs[3].Parts[1]["text"]) + + // Ids sort chronologically across the whole conversation — history ids + // must interleave correctly with live-minted ones (the TUI orders by id). + var prev string + for i, m := range msgs { + id, ok := m.Info["id"].(string) + require.True(t, ok) + assert.Greater(t, id, prev, "message %d id must sort after its predecessor", i) + prev = id + } +} + +func TestHistoryHonorsLimit(t *testing.T) { + srv, h := newBridgedFacade(t) + queueTwoTurnHistory(h) + + msgs := getMessages(t, srv.URL+"/session/ses_x/message?limit=2") + require.Len(t, msgs, 2) + // The newest messages win: turn 2's user+assistant. + assert.Equal(t, "second question", msgs[0].Parts[0]["text"]) + assert.Equal(t, "assistant", msgs[1].Info["role"]) +} + +func TestHistoryEmptyForFreshSession(t *testing.T) { + srv, _ := newBridgedFacade(t) + msgs := getMessages(t, srv.URL+"/session/ses_x/message") + assert.Empty(t, msgs) +} + +// A fresh proxy on a session with prior turns must list the session, or +// `opencode attach --continue` has nothing to resume. +func TestSessionListIncludesSessionWithHistory(t *testing.T) { + srv, h := newBridgedFacade(t) + queueTwoTurnHistory(h) + + resp, err := http.Get(srv.URL + "/session") + require.NoError(t, err) + defer resp.Body.Close() + var list []map[string]any + require.NoError(t, json.NewDecoder(resp.Body).Decode(&list)) + require.Len(t, list, 1) + assert.Equal(t, "ses_"+ocID("", testSessionID), list[0]["id"]) +} + +// History is fetched once and cached (the harness replay stream lingers ~8s +// before closing, and an attach asks twice) until a completed turn +// invalidates it. +func TestHistoryIsCachedUntilInvalidated(t *testing.T) { + h := agentproxytest.New(t, testSessionID) + client, err := godo.New(nil, godo.SetBaseURL(h.Server.URL+"/")) + require.NoError(t, err) + f := &Facade{SessionID: testSessionID, Sessions: do.NewHostedAgentsService(client), Dir: "/tmp/ws"} + srv := httptest.NewServer(f) + t.Cleanup(srv.Close) + + h.QueueReplayHistory( + agentproxytest.Event{RunID: "run-1", Type: string(godo.HostedAgentEventKindRunStarted), Data: json.RawMessage(`{"user_input":"old"}`)}, + agentproxytest.Event{RunID: "run-1", Type: string(godo.HostedAgentEventKindTokenChunk), Data: json.RawMessage(`{"text":"old answer"}`)}, + agentproxytest.Event{RunID: "run-1", Type: string(godo.HostedAgentEventKindRunCompleted)}, + ) + first := getMessages(t, srv.URL+"/session/ses_x/message") + require.Len(t, first, 2) + + // The harness's durable history changes, but the cache still serves the + // old view... + h.QueueReplayHistory( + agentproxytest.Event{RunID: "run-2", Type: string(godo.HostedAgentEventKindRunStarted), Data: json.RawMessage(`{"user_input":"new"}`)}, + agentproxytest.Event{RunID: "run-2", Type: string(godo.HostedAgentEventKindRunCompleted)}, + ) + cached := getMessages(t, srv.URL+"/session/ses_x/message") + require.Len(t, cached, 2) + assert.Equal(t, "old", cached[0].Parts[0]["text"]) + + // ...until a completed live turn invalidates it (translateEvent does + // this; exercised directly here). + f.invalidateHistory() + fresh := getMessages(t, srv.URL+"/session/ses_x/message") + require.Len(t, fresh, 1) // run-2 has no answer text: user message only + assert.Equal(t, "new", fresh[0].Parts[0]["text"]) +} + +// Failed turns keep the user's message visible but don't fabricate a +// completed answer. +func TestHistoryFailedTurn(t *testing.T) { + srv, h := newBridgedFacade(t) + h.QueueReplayHistory( + agentproxytest.Event{RunID: "run-f", Type: string(godo.HostedAgentEventKindRunStarted), Data: json.RawMessage(`{"user_input":"doomed"}`)}, + agentproxytest.Event{RunID: "run-f", Type: string(godo.HostedAgentEventKindRunFailed), Data: json.RawMessage(`{"message":"boom"}`)}, + ) + + msgs := getMessages(t, srv.URL+"/session/ses_x/message") + require.Len(t, msgs, 2) + assert.Equal(t, "user", msgs[0].Info["role"]) + assert.Equal(t, "assistant", msgs[1].Info["role"]) + assert.NotContains(t, msgs[1].Info, "finish") +} + +// invalidateHistory must never wait on an in-flight replay: the event loop +// calls it when a turn completes, and the replay stream lingers ~8s server- +// side — holding histMu across the fetch would stall event translation for +// that long (review finding on the cache commit). +func TestInvalidateHistoryDoesNotBlockOnInFlightFetch(t *testing.T) { + h := agentproxytest.New(t, testSessionID) + client, err := godo.New(nil, godo.SetBaseURL(h.Server.URL+"/")) + require.NoError(t, err) + f := &Facade{SessionID: testSessionID, Sessions: do.NewHostedAgentsService(client), Dir: "/tmp/ws"} + + // Hold the replay mid-stream: the WaitForHITL gate blocks delivery of the + // second event until the test resolves "gate-1", modeling the server-side + // linger for exactly as long as the test needs. (HangStreamAfterEvents + // can't do this — the harness exempts replay_only streams from it.) + h.QueueReplayHistory( + agentproxytest.Event{RunID: "run-1", Type: string(godo.HostedAgentEventKindRunStarted), Data: json.RawMessage(`{"user_input":"held"}`)}, + agentproxytest.Event{RunID: "run-1", Type: string(godo.HostedAgentEventKindRunCompleted), WaitForHITL: "gate-1"}, + ) + + fetchReturned := make(chan struct{}) + go func() { + defer close(fetchReturned) + _, _ = f.history(context.Background()) + }() + + // Wait until the fetch is actually in flight. + require.Eventually(t, func() bool { + f.histMu.Lock() + defer f.histMu.Unlock() + return f.histFetching + }, 5*time.Second, 5*time.Millisecond, "the history fetch never started") + + invalidated := make(chan struct{}) + go func() { + f.invalidateHistory() + close(invalidated) + }() + select { + case <-invalidated: + case <-time.After(2 * time.Second): + t.Fatal("invalidateHistory blocked behind the in-flight replay") + } + + // Release the gate so the fetch completes; its result must NOT be cached + // (the invalidation superseded it), so the next read replays fresh. + resp, err := http.Post(h.Server.URL+"/v2/agents/sessions/"+testSessionID+"/hitl/gate-1", + "application/json", strings.NewReader(`{"outcome":"HITL_OUTCOME_APPROVE"}`)) + require.NoError(t, err) + resp.Body.Close() + <-fetchReturned + f.histMu.Lock() + valid := f.histValid + f.histMu.Unlock() + assert.False(t, valid, "a fetch overlapped by an invalidation must not populate the cache") +}