From d8903bc92777b6d67e6a6ecf769f5a30203cd1ba Mon Sep 17 00:00:00 2001 From: Rohit C Prasad Date: Wed, 22 Jul 2026 10:40:01 -0700 Subject: [PATCH] Keep interrupted partial streams in the transcript GUI flushes the streaming buffer into a durable item on interrupted/error. Engine persists partial text on the provider-error path like the stop path. e2e red-green verified; full suites pass. --- coworker/engine.py | 6 ++++ surfaces/gui/e2e/fixtures.ts | 15 ++++++++-- surfaces/gui/e2e/interrupt-partial.spec.ts | 32 +++++++++++++++++++++ surfaces/gui/src/App.tsx | 21 +++++++++++++- tests/test_engine_stop.py | 33 ++++++++++++++++++++++ 5 files changed, 104 insertions(+), 3 deletions(-) create mode 100644 surfaces/gui/e2e/interrupt-partial.spec.ts diff --git a/coworker/engine.py b/coworker/engine.py index 2e24b8bb..3c14b179 100644 --- a/coworker/engine.py +++ b/coworker/engine.py @@ -244,6 +244,12 @@ class TurnEngine: if chunk.turn is not None: turn = chunk.turn except Exception as exc: # provider failure + # Same contract as the stop path below: the partial text the user + # watched arrive survives the failure (no tool calls — none finalized). + if streamed: + self.messages.append( + _assistant_message(AssistantTurn(text="".join(streamed))) + ) friendly = friendly_model_error(self.model, exc) payload = { "error": friendly or str(exc), diff --git a/surfaces/gui/e2e/fixtures.ts b/surfaces/gui/e2e/fixtures.ts index 67b11123..3bf66025 100644 --- a/surfaces/gui/e2e/fixtures.ts +++ b/surfaces/gui/e2e/fixtures.ts @@ -567,6 +567,7 @@ export async function mockApi(page: import("@playwright/test").Page) { ws.send(JSON.stringify({ type, data })); send("ready"); let pendingTool = "run_shell"; // which proposal the next approval decision resolves + let epicTimer: ReturnType | null = null; // the slow stream, stoppable via interrupt ws.onMessage((raw) => { const msg = JSON.parse(String(raw)); if (msg.type === "user_message") { @@ -651,11 +652,12 @@ export async function mockApi(page: import("@playwright/test").Page) { if (/stream the epic/i.test(msg.text)) { let ticks = 0; const line = "The epic scrolls ever onward, line upon line upon line. "; - const timer = setInterval(() => { + epicTimer = setInterval(() => { ticks += 1; send("assistant_delta", { text: line.repeat(3) + "\n\n" }); if (ticks >= 40) { - clearInterval(timer); + clearInterval(epicTimer!); + epicTimer = null; send("assistant_message", { text: ("The epic concludes. " + line).repeat(20) }); send("turn_done"); } @@ -686,6 +688,15 @@ export async function mockApi(page: import("@playwright/test").Page) { send("assistant_message", { text: `Done via ${pendingTool} [decision=${msg.decision}]` }); } send("turn_done"); + } else if (msg.type === "interrupt") { + // Stop mid-stream: like the real engine, end the turn with `interrupted` and + // NO assistant_message — the client owns promoting the partial into the transcript. + if (epicTimer) { + clearInterval(epicTimer); + epicTimer = null; + } + send("interrupted", {}); + send("turn_done"); } }); }); diff --git a/surfaces/gui/e2e/interrupt-partial.spec.ts b/surfaces/gui/e2e/interrupt-partial.spec.ts new file mode 100644 index 00000000..d3f0283c --- /dev/null +++ b/surfaces/gui/e2e/interrupt-partial.spec.ts @@ -0,0 +1,32 @@ +// Owner-hit 2026-07-22: Stop mid-stream kept the partial visible — until the NEXT message's +// turn_start wiped it, because the partial only ever lived in the ephemeral streaming buffer +// (assistant_message is what promotes text into the transcript, and an interrupted turn never +// emits one). The fix flushes the buffer into a durable assistant item on interrupted/error. +import { expect } from "@playwright/test"; +import { test } from "./fixtures"; + +test("interrupted partial stream survives the next turn", async ({ page }) => { + await page.goto("/"); + await page.getByText("Draft the launch note").first().click(); + const box = page.getByPlaceholder(/Ask the coworker/); + await box.fill("stream the epic"); + await box.press("Enter"); + + // Let a few deltas land, then stop the turn. + await expect(page.getByText("The epic scrolls ever onward").first()).toBeVisible({ + timeout: 10_000, + }); + await page.getByRole("button", { name: /Stop/ }).click(); + await expect(page.getByText("Interrupted.").first()).toBeVisible({ timeout: 5_000 }); + + // The partial is still on screen after the stop… + await expect(page.getByText("The epic scrolls ever onward").first()).toBeVisible(); + + // …and — the regression — still there after the next turn starts and completes. + await box.fill("continue please"); + await box.press("Enter"); + await expect(page.getByText("Echo: continue please", { exact: false }).first()).toBeVisible({ + timeout: 10_000, + }); + await expect(page.getByText("The epic scrolls ever onward").first()).toBeVisible(); +}); diff --git a/surfaces/gui/src/App.tsx b/surfaces/gui/src/App.tsx index 07728fea..64223c89 100644 --- a/surfaces/gui/src/App.tsx +++ b/surfaces/gui/src/App.tsx @@ -154,7 +154,14 @@ export function App() { const [connected, setConnected] = useState(false); const [running, setRunning] = useState(false); const [items, setItems] = useState([]); - const [streaming, setStreaming] = useState(""); + const [streaming, setStreamingState] = useState(""); + // Ref mirror of `streaming`: the WS handler closure is built once per socket and can't read + // fresh state — the interrupted/error flush below needs the live buffer at event time. + const streamingRef = useRef(""); + const setStreaming = (value: string | ((s: string) => string)) => { + streamingRef.current = typeof value === "function" ? value(streamingRef.current) : value; + setStreamingState(streamingRef.current); + }; const [todo, setTodo] = useState([]); const [sessions, setSessions] = useState([]); const [projects, setProjects] = useState([]); @@ -514,6 +521,16 @@ export function App() { if (gatesWorkspace(agent) && !workspace) return; // Code needs a folder (gate handles it) const handleEvent = (ev: WsEvent) => { const d = ev.data || {}; + // An interrupted/errored turn never emits assistant_message, so its streamed partial + // would otherwise live only in the ephemeral buffer until the next turn_start wipes it + // (owner-hit 2026-07-22). Promote it to a durable transcript item — the engine persists + // the same text server-side, so the live view and a session reload now agree. + const flushPartialStream = () => { + const partial = streamingRef.current; + if (!partial) return; + setStreaming(""); + setItems((p) => [...p, { kind: "assistant", text: partial, ts: Date.now() / 1000 }]); + }; switch (ev.type) { case "ready": setConnected(true); @@ -622,9 +639,11 @@ export function App() { setItems((p) => [...p, { kind: "notice", tone: "warn", text: "Stopped: max iterations reached." }]); break; case "interrupted": + flushPartialStream(); setItems((p) => [...p, { kind: "notice", tone: "warn", text: "Interrupted." }]); break; case "error": + flushPartialStream(); setItems((p) => [...p, { kind: "notice", tone: "warn", text: "Error: " + (d.error || "unknown") }]); break; case "turn_done": diff --git a/tests/test_engine_stop.py b/tests/test_engine_stop.py index 85b24043..790fc08b 100644 --- a/tests/test_engine_stop.py +++ b/tests/test_engine_stop.py @@ -92,6 +92,39 @@ def test_stop_mid_stream_keeps_partial_text(tmp_path): assert "tool_calls" not in last +class FailingStreamProvider(ProviderClient): + """Streams a few deltas, then dies — a provider outage mid-answer.""" + + def complete(self, **kwargs): # pragma: no cover + raise NotImplementedError + + def capabilities(self, model): + return ModelCapabilities() + + def stream(self, *, model, messages, tools=None, **settings): + yield StreamChunk(text_delta="partial ") + yield StreamChunk(text_delta="answer") + raise RuntimeError("provider went away") + + +def test_provider_error_mid_stream_keeps_partial_text(tmp_path): + engine = TurnEngine( + provider=FailingStreamProvider(), + registry=ToolRegistry(), + permissions=PermissionEngine(workspace_root=tmp_path), + model="gpt-5.5", + ) + + async def run(): + return [ev async for ev in engine.run("go")] + + events = asyncio.run(run()) + assert events[-1].type == EventType.ERROR + last = engine.messages[-1] + assert last["role"] == "assistant" and last["content"] == "partial answer" + assert "tool_calls" not in last + + def test_stop_while_awaiting_approval(tmp_path): async def never_answers(_req): await asyncio.Event().wait() # a pending approval card nobody answers