From 9687abf674f530c61da4e6dbc0c8ac47d9977f94 Mon Sep 17 00:00:00 2001 From: juliaye Date: Fri, 4 Sep 2026 00:15:02 +0300 Subject: [PATCH 1/3] =?UTF-8?q?agents=20proxy:=20opencode=20M3=20=E2=80=94?= =?UTF-8?q?=20history=20on=20re-attach?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit GET /session/{id}/message is served from a one-shot replay_only StreamSession pass, reconstructing each hosted run as its user message (run.started's prompt text — the SSE wire delivers it as data.agent, not the proto's user_input; both are accepted) and assistant message (token deltas accumulated; reasoning as its own part). opencode's attach replays history client-side by fetching this endpoint (?limit=100 observed live), so re-attaching a fresh proxy shows prior turns; the session list now includes the bridged session whenever it has replayable history, which is what --continue resumes through. Ids are minted from turn event times with the real time-encoding, with monotonic T allocation across turns — two turns sharing a millisecond would otherwise interleave (caught by test). time.completed comes from the run.completed event, not the run's max event time — late stragglers (run.log after completion) otherwise inflate a 2s turn into a 20m badge (found live). Also stubs GET /session/{id}/diff and /todo (empty arrays, the TUI's post-turn refresh lookups). Live-verified: a fresh proxy + `opencode attach --continue` renders the session's six prior turns with correct prompts, answers, order, and durations. Co-Authored-By: Claude Fable 5 --- internal/agentproxy/opencode/bridge.go | 13 +- internal/agentproxy/opencode/facade.go | 11 +- internal/agentproxy/opencode/facade_test.go | 20 +- internal/agentproxy/opencode/history.go | 238 +++++++++++++++++++ internal/agentproxy/opencode/history_test.go | 132 ++++++++++ 5 files changed, 389 insertions(+), 25 deletions(-) create mode 100644 internal/agentproxy/opencode/history.go create mode 100644 internal/agentproxy/opencode/history_test.go diff --git a/internal/agentproxy/opencode/bridge.go b/internal/agentproxy/opencode/bridge.go index 60facd68d..00f4f7fe4 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 } diff --git a/internal/agentproxy/opencode/facade.go b/internal/agentproxy/opencode/facade.go index 1496ab84c..20615b281 100644 --- a/internal/agentproxy/opencode/facade.go +++ b/internal/agentproxy/opencode/facade.go @@ -163,15 +163,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..28d553ebf --- /dev/null +++ b/internal/agentproxy/opencode/history.go @@ -0,0 +1,238 @@ +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). No cache: history is +// fetched per request, and the TUI asks once per attach. + +// 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 +} + +// fetchHistory replays the durable event history and reconstructs completed +// turns. Only text is reconstructed in M3 — tool-call parts are M4. +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.fetchHistory(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.fetchHistory(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..50c5d29c5 --- /dev/null +++ b/internal/agentproxy/opencode/history_test.go @@ -0,0 +1,132 @@ +package opencode + +import ( + "encoding/json" + "net/http" + "testing" + + "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"]) +} + +// 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") +} From 5899192651bcfeb21f77cd4399ccb6ebf8dd6ea0 Mon Sep 17 00:00:00 2001 From: juliaye Date: Fri, 4 Sep 2026 00:30:07 +0300 Subject: [PATCH 2/3] =?UTF-8?q?agents=20proxy:=20opencode=20=E2=80=94=20ca?= =?UTF-8?q?che=20history=20(replay=20linger=20made=20resume=20~16s)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit harness-api's replay_only stream lingers ~8s after catching up (measured: 68 events, close at exactly 8.0s both runs), and an attach fetches history twice (session-list gating + the message list), so resume rendered ~16s late. History is now reconstructed once and cached — the mutex doubles as single-flight — warmed in the background off the facade's first request (the TUI's health preflight precedes the attach burst), and invalidated when a live turn completes. Warm-path timings: 8s -> 11ms. The 8s replay-close linger itself is a harness-api behavior worth fixing server-side (it also slows doctl agents logs); tracked outside this PR. Co-Authored-By: Claude Fable 5 --- internal/agentproxy/opencode/bridge.go | 3 ++ internal/agentproxy/opencode/facade.go | 18 ++++++++- internal/agentproxy/opencode/history.go | 40 ++++++++++++++++++-- internal/agentproxy/opencode/history_test.go | 39 +++++++++++++++++++ 4 files changed, 96 insertions(+), 4 deletions(-) diff --git a/internal/agentproxy/opencode/bridge.go b/internal/agentproxy/opencode/bridge.go index 00f4f7fe4..867f33f92 100644 --- a/internal/agentproxy/opencode/bridge.go +++ b/internal/agentproxy/opencode/bridge.go @@ -378,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 @@ -404,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 20615b281..7e57699af 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,26 @@ type Facade struct { // session, so the session list grows its single entry (see // handleSessionList). sessionCreated bool + + // History cache (see history() in history.go): histMu also serves as + // single-flight for the slow replay_only fetch. + histMu sync.Mutex + hist []historyMessage + histValid bool } // 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) } diff --git a/internal/agentproxy/opencode/history.go b/internal/agentproxy/opencode/history.go index 28d553ebf..95be19eed 100644 --- a/internal/agentproxy/opencode/history.go +++ b/internal/agentproxy/opencode/history.go @@ -49,8 +49,42 @@ func (ht *historyTurn) endMs() int64 { 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 mutex doubles as single-flight: concurrent +// callers wait for the one in-progress replay instead of starting their own. +// +// The cache is warmed at the first request the facade sees (the TUI's +// /global/health preflight fires well before the attach burst) and +// invalidated when a live turn completes (invalidateHistory). +func (f *Facade) history(ctx context.Context) ([]historyMessage, error) { + f.histMu.Lock() + defer f.histMu.Unlock() + if f.histValid { + return f.hist, nil + } + msgs, err := f.fetchHistory(ctx) + if err != nil { + return nil, err + } + f.hist, f.histValid = msgs, true + return msgs, nil +} + +// invalidateHistory drops the cache; the next history() call replays fresh. +func (f *Facade) invalidateHistory() { + f.histMu.Lock() + 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. +// 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, @@ -213,7 +247,7 @@ func (f *Facade) turnMessages(ht *historyTurn, t uint64) []historyMessage { // 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.fetchHistory(r.Context()) + msgs, err := f.history(r.Context()) if err != nil { http.Error(w, fmt.Sprintf("replaying session history failed: %v", err), http.StatusBadGateway) return @@ -233,6 +267,6 @@ func (f *Facade) handleMessageList(w http.ResponseWriter, r *http.Request) { // 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.fetchHistory(ctx) + 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 index 50c5d29c5..442aa7aa6 100644 --- a/internal/agentproxy/opencode/history_test.go +++ b/internal/agentproxy/opencode/history_test.go @@ -3,8 +3,10 @@ package opencode import ( "encoding/json" "net/http" + "net/http/httptest" "testing" + "github.com/digitalocean/doctl/do" "github.com/digitalocean/doctl/internal/agentproxy/agentproxytest" "github.com/digitalocean/godo" "github.com/stretchr/testify/assert" @@ -115,6 +117,43 @@ func TestSessionListIncludesSessionWithHistory(t *testing.T) { 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) { From 83ede4b9000fedd548fb6f323c69bf9c7e4516bf Mon Sep 17 00:00:00 2001 From: juliaye Date: Tue, 15 Sep 2026 14:30:39 -0700 Subject: [PATCH 3/3] =?UTF-8?q?agents=20proxy:=20opencode=20M3=20review=20?= =?UTF-8?q?fixes=20=E2=80=94=20fetch=20history=20outside=20histMu?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit history() held histMu across the whole ~8s replay, and invalidateHistory() needs that same lock — so a turn completing while history was still loading stalled the event loop (and every frame behind it) until the replay finished. The fetch now runs outside the lock: single-flight rides a histFetching flag + done channel (waiters also honor their ctx), and a histGen counter detects an invalidation that overlapped a fetch — the overlapped result is returned to its caller but not cached. Regression test holds a replay mid-stream via the harness's WaitForHITL gate and asserts invalidateHistory returns immediately. Also refreshes the stale file header that still said "No cache: history is fetched per request" from before the cache commit. Co-Authored-By: Claude Fable 5 --- internal/agentproxy/opencode/facade.go | 15 +++-- internal/agentproxy/opencode/history.go | 65 +++++++++++++++----- internal/agentproxy/opencode/history_test.go | 59 ++++++++++++++++++ 3 files changed, 118 insertions(+), 21 deletions(-) diff --git a/internal/agentproxy/opencode/facade.go b/internal/agentproxy/opencode/facade.go index 7e57699af..2e91301b8 100644 --- a/internal/agentproxy/opencode/facade.go +++ b/internal/agentproxy/opencode/facade.go @@ -74,11 +74,16 @@ type Facade struct { // handleSessionList). sessionCreated bool - // History cache (see history() in history.go): histMu also serves as - // single-flight for the slow replay_only fetch. - histMu sync.Mutex - hist []historyMessage - histValid 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. diff --git a/internal/agentproxy/opencode/history.go b/internal/agentproxy/opencode/history.go index 95be19eed..151b4d116 100644 --- a/internal/agentproxy/opencode/history.go +++ b/internal/agentproxy/opencode/history.go @@ -15,8 +15,9 @@ import ( // 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). No cache: history is -// fetched per request, and the TUI asks once per attach. +// 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. @@ -54,29 +55,61 @@ func (ht *historyTurn) endMs() int64 { // 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 mutex doubles as single-flight: concurrent -// callers wait for the one in-progress replay instead of starting their own. +// 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) and -// invalidated when a live turn completes (invalidateHistory). +// /global/health preflight fires well before the attach burst). func (f *Facade) history(ctx context.Context) ([]historyMessage, error) { - f.histMu.Lock() - defer f.histMu.Unlock() - if f.histValid { - return f.hist, nil - } - msgs, err := f.fetchHistory(ctx) - if err != nil { - return nil, err + 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 } - f.hist, f.histValid = msgs, true - return msgs, nil } // 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() diff --git a/internal/agentproxy/opencode/history_test.go b/internal/agentproxy/opencode/history_test.go index 442aa7aa6..3a081fe6b 100644 --- a/internal/agentproxy/opencode/history_test.go +++ b/internal/agentproxy/opencode/history_test.go @@ -1,10 +1,13 @@ 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" @@ -169,3 +172,59 @@ func TestHistoryFailedTurn(t *testing.T) { 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") +}