diff --git a/packages/docs/content/server-api.en.md b/packages/docs/content/server-api.en.md index 391c21d..3fd5254 100644 --- a/packages/docs/content/server-api.en.md +++ b/packages/docs/content/server-api.en.md @@ -144,7 +144,7 @@ The paths below omit the `/api/sessions/:sessionId` prefix. For the storage mode | GET | / | Session info (the single-session GET additionally carries `tracePath`, the absolute path of the latest Trace file; list rows omit it) | | PATCH | / | Update: `{approvalMode?, archived?, title?}` | | DELETE | / | Delete the Session (along with its Traces and scratch files) | -| GET | /messages | Full OmniMessage history | +| GET | /messages | Full OmniMessage history; while a Task runs the response also carries `live` (the in-progress stream tail, see below) | | GET | /stream | SSE event stream (next section) | | POST | /tasks | Start a Task: `{input: TaskInputPart[], thinkingLevel?, queueIfBusy?}` → 202. With `queueIfBusy`, a busy session holds the input as a follow-up (`queued: true`) and auto-starts it as an ordinary next task once idle; `task_state` events report the queued count | | POST | /steer | Mid-run steering: `{text}` queues a message for the running Task (delivered between turns as a standalone `[user_steering]` user message) → 202; 409 `not_running` when no Task is in progress | @@ -164,6 +164,28 @@ The paths below omit the `/api/sessions/:sessionId` prefix. For the storage mode General conventions: Sessions the user cannot access always return 404 — their existence is never leaked; only one Task or compaction runs per Session at a time, and conflicts return 409 (`task_in_progress` / `compacting`). +#### The `live` field on GET /messages + +The Trace stores only complete messages (streaming `partial_*` never reaches disk), so history alone cannot show a message that is still streaming. While the Session is running or compacting, the messages response therefore also carries the in-progress stream tail: + +```ts +interface MessagesResponse { + messages: OmniMessage[]; + live?: { + // The Session channel's most recently assigned SSE event id (`-`): + // every event published up to and including this id is already reflected in `fragments`. + cursor: string; + // One synthetic `partial_* start` OmniMessage per open streaming fragment, whose + // payload carries the full accumulated content so far (text/thinking prefix, + // tool-call name + accumulated arguments, tool-output prefix + images), with the + // original `origin` chain preserved (subagent fragments included). + fragments: OmniMessage[]; + }; +} +``` + +`cursor` and `fragments` are captured atomically before the trace read starts. A client using the connect-first pattern (below) applies them after history: when the cursor's epoch matches the epoch of the SSE events it has buffered, it drops every buffered **partial** event with seq ≤ cursor (their content is already accumulated inside `fragments`), feeds `fragments` through its normal reducer, then replays the rest of the buffer. Buffered **complete** messages are never dropped by the cursor — the regular overlap dedup decides for them. `live` is omitted while idle. + Workspace files may be Agent-generated, so `GET /files/content` treats them as untrusted: every response carries `X-Content-Type-Options: nosniff`, and the rest of the headers depend on the two flags (`download=1` wins over `preview=1`): | Query | Content-Type | Content-Disposition | Content-Security-Policy | @@ -255,7 +277,7 @@ export type ServerEvent = ### Delivery Guarantees - Event ids are monotonic per channel, shaped `-`; -- Each channel keeps a bounded replay buffer (most recent 1000 events or 2MB); +- Each channel keeps a bounded replay buffer (most recent 10,000 events or 8MB); - Reconnecting with `Last-Event-ID` replays the gap on a buffer hit; on a miss the server first sends `resync_required`, and the client refetches `/messages` before continuing; - A heartbeat comment line is written every 20 seconds; - Event order: on a reconnect carrying `Last-Event-ID`, **the replayed gap (or `resync_required`) arrives first**, then the initial events — the authoritative `task_state` snapshot and still-pending approval_requests — then the live stream. A fresh connection (no `Last-Event-ID`) skips replay, so its first event is the `task_state` snapshot. @@ -266,8 +288,9 @@ The order the bundled Web App uses: 1. Connect `/stream` first and buffer incoming events; 2. GET `/messages` for the full history; -3. Replay the buffer, deduplicating the overlap; -4. Go live. +3. If the response carries `live` (a Task is running), drop the buffered partials the cursor already covers and seed the `live.fragments` on top of history — the in-progress message reappears with its streamed prefix intact; +4. Replay the buffer, deduplicating the overlap; +5. Go live. ## Type Imports diff --git a/packages/docs/content/server-api.zh.md b/packages/docs/content/server-api.zh.md index 2fdb861..8fb8ffa 100644 --- a/packages/docs/content/server-api.zh.md +++ b/packages/docs/content/server-api.zh.md @@ -144,7 +144,7 @@ Schedule 写操作仅限 Owner。新建 Session 模式的任务,`modelId` 与 | GET | / | Session 信息(单会话 GET 额外携带 `tracePath`:最新 Trace 文件的绝对路径;列表行不含) | | PATCH | / | 更新:`{approvalMode?, archived?, title?}` | | DELETE | / | 删除 Session(连同 Trace 与暂存文件) | -| GET | /messages | 完整 OmniMessage 历史 | +| GET | /messages | 完整 OmniMessage 历史;Task 运行期间响应额外携带 `live`(进行中的流式尾部,见下) | | GET | /stream | SSE 事件流(见下节) | | POST | /tasks | 发起 Task:`{input: TaskInputPart[], thinkingLevel?, queueIfBusy?}` → 202。带 `queueIfBusy` 时,运行中的 Session 会把输入暂存为跟进消息(`queued: true`),空闲后按序自动作为普通 Task 发出;`task_state` 事件携带排队数 | | POST | /steer | 运行中插话:`{text}` 为运行中的 Task 排队一条消息(作为独立的 `[user_steering]` 用户消息随下一轮送达)→ 202;无 Task 运行返回 409 `not_running` | @@ -164,6 +164,27 @@ Schedule 写操作仅限 Owner。新建 Session 模式的任务,`modelId` 与 通用约定:无权访问的 Session 一律返回 404,不泄露其存在性;每个 Session 同时只允许一个 Task 或压缩在运行,冲突时返回 409(`task_in_progress` / `compacting`)。 +#### GET /messages 的 `live` 字段 + +Trace 只存完整消息(流式 `partial_*` 永远不落盘),所以仅靠历史无法呈现一条正在流式输出的消息。因此当 Session 处于运行/压缩状态时,messages 响应额外携带进行中的流式尾部: + +```ts +interface MessagesResponse { + messages: OmniMessage[]; + live?: { + // Session 通道最近分配的 SSE 事件 id(`-`): + // 截至该 id(含)发布的所有事件都已累积进 `fragments`。 + cursor: string; + // 每个未闭合流式片段对应一条合成的 `partial_* start` OmniMessage,其 payload 携带 + // 迄今累积的全部内容(文本/思考前缀、工具调用名 + 已累积参数、工具输出前缀 + 图片), + // 并保留原始 `origin` 链(子智能体片段同样覆盖)。 + fragments: OmniMessage[]; + }; +} +``` + +`cursor` 与 `fragments` 在 Trace 读取开始前原子采集。使用先连接模式(见下)的客户端在应用完历史后处理它们:当 cursor 的 epoch 与本连接已缓冲事件的 epoch 一致时,丢弃 seq ≤ cursor 的已缓冲 **partial** 事件(其内容已累积在 `fragments` 里),把 `fragments` 按正常归约路径喂入,再重放剩余缓冲。已缓冲的**完整**消息从不按 cursor 丢弃 —— 仍由常规重叠去重裁决。空闲时不携带 `live`。 + Workspace 文件可能由 Agent 生成,`GET /files/content` 一律按不可信内容处理:所有响应都带 `X-Content-Type-Options: nosniff`,其余响应头取决于两个开关(`download=1` 优先于 `preview=1`): | 查询参数 | Content-Type | Content-Disposition | Content-Security-Policy | @@ -254,7 +275,7 @@ export type ServerEvent = ### 投递保证 - 事件 id 按通道单调递增,形如 `-`; -- 每通道维护有界重放缓冲(最近 1000 条事件或 2MB); +- 每通道维护有界重放缓冲(最近 10,000 条事件或 8MB); - 携带 `Last-Event-ID` 重连时,命中缓冲则补发缺口;未命中则先发 `resync_required`,客户端重新拉取 `/messages` 后继续消费; - 每 20 秒写一条心跳注释行; - 事件次序:带 `Last-Event-ID` 重连时,**补发的缺口(或 `resync_required`)最先送达**,随后才是初始事件——权威的 `task_state` 快照与未决的 approval_request,再进入实时流;全新连接(无 `Last-Event-ID`)不重放缓冲,首个事件即为 `task_state` 快照。 @@ -265,8 +286,9 @@ export type ServerEvent = 1. 先连接 `/stream` 并缓冲收到的事件; 2. 再 GET `/messages` 拉取完整历史; -3. 回放缓冲区并对重叠消息去重; -4. 转入实时消费。 +3. 若响应携带 `live`(有 Task 在运行),丢弃 cursor 已覆盖的缓冲 partial 事件,并把 `live.fragments` 播种到历史之上 —— 进行中的消息连同已流式输出的前缀一起回到画面; +4. 回放缓冲区并对重叠消息去重; +5. 转入实时消费。 ## 类型导入 diff --git a/packages/server/src/api/types.ts b/packages/server/src/api/types.ts index 1fb725e..20b4923 100644 --- a/packages/server/src/api/types.ts +++ b/packages/server/src/api/types.ts @@ -561,9 +561,42 @@ export interface SessionPatchRequest { title?: string; } +/** + * Live in-progress tail of a running Session, carried by `MessagesResponse.live`. + * + * Contract (see runtime/live-tail.ts and the GET /messages route): the server captures + * `cursor` and `fragments` atomically — in one synchronous tick, before starting the + * trace read — while the Session is running/compacting. + * - `cursor`: the Session channel's most recently assigned SSE event id + * (`-`); every event published up to and including this id is already + * reflected in `fragments`. + * - `fragments`: one synthetic `partial_* start` OmniMessage per open streaming + * fragment, whose payload carries the full accumulated content so far (text/thinking + * prefix, tool-call name + accumulated arguments, tool-output prefix + images), with + * the original `origin` chain preserved. + * + * Client usage (the bundled Web App's connect-first flow): after applying `messages`, + * when the cursor's epoch matches the epoch of the SSE events seen on the current + * connection, drop every buffered **partial** event with seq <= cursor (its content is + * already inside `fragments`), feed `fragments` through the normal reducer path, then + * replay the rest of the buffer. Buffered **complete** messages are never dropped by the + * cursor — the regular overlap dedup decides for them — so nothing is lost even when a + * complete message's trace append is still in flight at read time. + */ +export interface MessagesLiveTail { + cursor: string; + fragments: OmniMessage[]; +} + /** Message history: the full messages and events from concatenating all of this Session's Trace files in order (excludes partial_*). */ export interface MessagesResponse { messages: OmniMessage[]; + /** + * Present only while the Session is running/compacting: the in-progress stream tail + * (open streaming fragments + the channel cursor they cover), so a client joining + * mid-stream can render the currently streaming message. Omitted when idle. + */ + live?: MessagesLiveTail; } // --------------------------------------------------------------------------- diff --git a/packages/server/src/http/routes/sessions.ts b/packages/server/src/http/routes/sessions.ts index f87eeb9..e058adc 100644 --- a/packages/server/src/http/routes/sessions.ts +++ b/packages/server/src/http/routes/sessions.ts @@ -15,6 +15,7 @@ import type { OmniMessage, ThinkingLevelName } from "@prismshadow/penguin-core"; import type { ApprovalMode, FilesStatResponse, + MessagesLiveTail, MessagesResponse, ServerEvent, SessionCategory, @@ -337,12 +338,29 @@ export function sessionsRoutes(deps: AppDeps): Hono { app.get("/:sessionId/messages", async (c) => { const row = resolveSession(c); + // Live tail (running/compacting sessions only): capture the channel cursor and the + // open-fragment snapshot together, synchronously — no await between the two, and both + // BEFORE the trace read starts. That ordering is what makes the client contract safe + // (see MessagesLiveTail in api/types.ts): every published event with id <= cursor is + // already reflected in `fragments`, and partial_* messages never reach the Trace, so + // the client may drop its buffered partials at/or before the cursor and seed from + // `fragments` without loss or duplication. Complete messages are never dropped by the + // cursor — the client's overlap dedup against `messages` decides for them — so a + // complete message whose trace append is still in flight when the read starts is not + // lost either. + let live: MessagesLiveTail | undefined; + if (deps.manager.statusOf(row.sessionId) !== "idle") { + live = { + cursor: deps.channels.get(row.sessionId).lastEventId, + fragments: deps.manager.liveFragments(row.sessionId), + }; + } const messages = await deps.traceService.readMessages( row.projectId, row.agentId, row.sessionId, ); - return c.json({ messages } satisfies MessagesResponse); + return c.json({ messages, ...(live !== undefined ? { live } : {}) } satisfies MessagesResponse); }); app.get("/:sessionId/stream", (c) => { diff --git a/packages/server/src/runtime/channel.ts b/packages/server/src/runtime/channel.ts index d9465d8..252d6a5 100644 --- a/packages/server/src/runtime/channel.ts +++ b/packages/server/src/runtime/channel.ts @@ -73,6 +73,17 @@ export class Channel { return this.listeners.size; } + /** + * Id of the most recently assigned event (`-`; seq 0 when none was assigned + * yet). Unicast (sendTo) seqs count too — the value is a position marker, not a buffer + * lookup key. Used as the live-tail cursor on GET /messages: captured in the same + * synchronous tick as the fragment snapshot, it tells the client which buffered events + * the snapshot already covers. + */ + get lastEventId(): string { + return `${this.epoch}-${this.nextSeq - 1}`; + } + /** Broadcast an event: number it, buffer it (evicting the oldest), notify all subscribers. */ publish(data: unknown, event?: string): ChannelEvent { const entry = this.makeEvent(data, event); diff --git a/packages/server/src/runtime/live-tail.ts b/packages/server/src/runtime/live-tail.ts new file mode 100644 index 0000000..c6c7717 --- /dev/null +++ b/packages/server/src/runtime/live-tail.ts @@ -0,0 +1,237 @@ +/** + * Live tail of a running Session's stream: the accumulated state of every OPEN + * streaming fragment (partial_text / partial_thinking / partial_tool_call / + * partial_tool_call_output), kept per origin chain. + * + * Why: `partial_*` messages never reach the Trace (core's Writer filters them), so a + * client that joins mid-stream — a refresh during a long tool call — cannot rebuild the + * in-progress message from `GET /messages` alone, and its fresh EventSource carries no + * `Last-Event-ID` for the channel buffer to replay. This tracker lets the messages + * endpoint attach, alongside history, one synthetic `partial_* start` OmniMessage per + * open fragment whose payload carries the full accumulated content so far (text/thinking + * prefix, tool-call name + accumulated arguments, tool-output prefix + images); the + * client seeds these into its stream model and live deltas keep appending on top. See + * `MessagesResponse.live` in api/types.ts for the client-facing contract. + * + * Fed by SessionManager.drive in the same synchronous tick as each channel publish, so a + * "channel cursor + fragments" capture between two publishes is always a consistent + * snapshot. Mirrors core PartialAggregator's merge semantics (fragment key = payload type + * + tool_call_id; start reopens, delta accumulates, stop closes) with the origin chain + * added to the key — the aggregator collapses fragments into complete messages, while + * this keeps the running prefix instead. A complete model message with the same identity + * also closes the fragment (covers stop-less closures, e.g. interruption cleanup); the + * whole session entry is dropped when the run ends (SessionManager.drive's finally). + */ +import { isPartialPayload } from "@prismshadow/penguin-core"; +import type { + CompleteModelPayload, + OmniMessage, + PartialModelPayload, +} from "@prismshadow/penguin-core"; + +type PartialKind = PartialModelPayload["type"]; + +/** + * Per-fragment accumulation cap. Environment already front-truncates tool output online + * (default 16k chars), so tool fragments stay small by construction; text/thinking have + * no upstream cap and a fast large-code reply can reach a few hundred KB (see the channel + * buffer sizing note) — 512KB covers that comfortably while bounding a runaway fragment. + * When the cap trips, the TAIL is kept (the prefix is dropped): the seeded content then + * joins seamlessly with the live deltas that follow the capture, and the complete message + * reconciles the full content at the end anyway. Trimming happens with slack so it costs + * one copy per `FRAGMENT_CAP_SLACK` chars of growth, not one per delta. Tail-keeping is a + * deliberate one-rule-for-all: it is seamless for text/thinking/tool output, and for + * `partial_tool_call` arguments it yields a syntactically broken JSON prefix — accepted, + * since arguments approaching 512KB are not reachable in practice and the complete + * message reconciles either way. + */ +const FRAGMENT_CAP = 512 * 1024; +const FRAGMENT_CAP_SLACK = 64 * 1024; + +interface OpenFragment { + kind: PartialKind; + /** Origin chain copied from the opening message (absent = main session). */ + origin?: string[]; + /** Timestamp of the fragment's original start (reused on the synthetic start). */ + timestamp: string; + /** Accumulated text / thinking / tool-call arguments / tool output (tail-capped). */ + buffer: string; + name?: string; + toolCallId?: string; + /** Tool-output images: not incremental — a single delta carries the whole set; a later one overwrites. */ + images?: string[]; +} + +/** The partial kind a complete payload closes out; null when it has no streamed counterpart. */ +function partialKindFor(p: CompleteModelPayload): PartialKind | null { + switch (p.type) { + case "text": + // Only assistant text streams; a user text mid-run (steering) must not close the model's open fragment. + return p.role === "assistant" ? "partial_text" : null; + case "thinking": + return "partial_thinking"; + case "tool_call": + return "partial_tool_call"; + case "tool_call_output": + return "partial_tool_call_output"; + default: + return null; + } +} + +/** Fragment key: origin chain + payload type + tool_call_id (same merge rule as core's PartialAggregator, origin added). */ +function fragmentKey(origin: string[] | undefined, kind: PartialKind, toolCallId: string): string { + // "\0" (the escape, not a raw byte -- a raw NUL makes git classify the file as binary) + // separates the origin chain from the kind: session ids are [A-Za-z0-9_-], so the + // separator can never occur inside a chain segment and keys cannot collide. + return `${(origin ?? []).join("/")}\0${kind}::${toolCallId}`; +} + +function toolCallIdOf(p: object): string { + const id = (p as { tool_call_id?: unknown }).tool_call_id; + return typeof id === "string" ? id : ""; +} + +function appendDelta(frag: OpenFragment, p: PartialModelPayload): void { + switch (p.type) { + case "partial_text": + frag.buffer += p.text; + break; + case "partial_thinking": + frag.buffer += p.thinking; + break; + case "partial_tool_call": + frag.buffer += p.arguments; + if (p.name) frag.name = p.name; + frag.toolCallId = p.tool_call_id; + break; + case "partial_tool_call_output": + frag.buffer += p.output; + if (p.images && p.images.length > 0) frag.images = p.images; + frag.toolCallId = p.tool_call_id; + break; + } + if (frag.buffer.length > FRAGMENT_CAP + FRAGMENT_CAP_SLACK) { + frag.buffer = frag.buffer.slice(frag.buffer.length - FRAGMENT_CAP); + } +} + +/** The synthetic `partial_* start` payload carrying the fragment's full accumulated content. */ +function startPayload(frag: OpenFragment): PartialModelPayload { + switch (frag.kind) { + case "partial_text": + return { type: "partial_text", role: "assistant", event_type: "start", text: frag.buffer }; + case "partial_thinking": + return { + type: "partial_thinking", + role: "assistant", + event_type: "start", + thinking: frag.buffer, + }; + case "partial_tool_call": + return { + type: "partial_tool_call", + role: "assistant", + event_type: "start", + name: frag.name ?? "", + arguments: frag.buffer, + tool_call_id: frag.toolCallId ?? "", + }; + case "partial_tool_call_output": + return { + type: "partial_tool_call_output", + role: "user", + event_type: "start", + output: frag.buffer, + ...(frag.images !== undefined && frag.images.length > 0 ? { images: frag.images } : {}), + tool_call_id: frag.toolCallId ?? "", + }; + } +} + +export class LiveTailTracker { + /** sessionId → open fragments keyed by fragmentKey, in the order they were opened. */ + private readonly sessions = new Map>(); + + /** Feed one published message (call in the same synchronous tick as the channel publish). */ + observe(sessionId: string, msg: OmniMessage): void { + if (msg.type !== "model_msg") return; + const p = msg.payload; + if (!isPartialPayload(p)) { + // A complete message closes the matching open fragment (normally the stop already + // did; this also covers stop-less closures such as interruption cleanup). + const kind = partialKindFor(p as CompleteModelPayload); + if (kind === null) return; + const open = this.sessions.get(sessionId); + if (!open) return; + open.delete(fragmentKey(msg.origin, kind, toolCallIdOf(p))); + if (open.size === 0) this.sessions.delete(sessionId); + return; + } + const key = fragmentKey(msg.origin, p.type, toolCallIdOf(p)); + let open = this.sessions.get(sessionId); + if (p.event_type === "start") { + if (!open) { + open = new Map(); + this.sessions.set(sessionId, open); + } + // start reopens the fragment: a previous same-key fragment (out-of-order) is replaced. + const frag: OpenFragment = { + kind: p.type, + timestamp: msg.timestamp, + buffer: "", + ...(msg.origin && msg.origin.length > 0 ? { origin: [...msg.origin] } : {}), + }; + open.delete(key); // re-insert so the emit order tracks the reopen + open.set(key, frag); + appendDelta(frag, p); + return; + } + let frag = open?.get(key); + if (!frag) { + // delta/stop without a start: lenient, same as core's PartialAggregator. + if (p.event_type === "stop") return; // nothing was open; nothing to keep or clear + if (!open) { + open = new Map(); + this.sessions.set(sessionId, open); + } + frag = { + kind: p.type, + timestamp: msg.timestamp, + buffer: "", + ...(msg.origin && msg.origin.length > 0 ? { origin: [...msg.origin] } : {}), + }; + open.set(key, frag); + } + appendDelta(frag, p); + if (p.event_type === "stop") { + open!.delete(key); + if (open!.size === 0) this.sessions.delete(sessionId); + } + } + + /** + * Synthetic `partial_* start` messages for every open fragment (in open order), each + * carrying the full accumulated content, the original origin chain, and the original + * start timestamp. Empty when the session has no open fragments. + */ + fragments(sessionId: string): OmniMessage[] { + const open = this.sessions.get(sessionId); + if (!open) return []; + const out: OmniMessage[] = []; + for (const frag of open.values()) { + out.push({ + timestamp: frag.timestamp, + type: "model_msg", + payload: startPayload(frag), + ...(frag.origin !== undefined ? { origin: [...frag.origin] } : {}), + }); + } + return out; + } + + /** Drop all fragment state for a session (the run ended; nothing will continue these fragments). */ + clear(sessionId: string): void { + this.sessions.delete(sessionId); + } +} diff --git a/packages/server/src/runtime/session-manager.ts b/packages/server/src/runtime/session-manager.ts index d323b81..36bb830 100644 --- a/packages/server/src/runtime/session-manager.ts +++ b/packages/server/src/runtime/session-manager.ts @@ -49,6 +49,7 @@ import { ApprovalRegistry, makeApprove } from "./approvals.js"; import type { PendingApproval } from "./approvals.js"; import type { ChannelHub } from "./channel.js"; import type { ErrorSink } from "./error-recorder.js"; +import { LiveTailTracker } from "./live-tail.js"; import { asSessionSource } from "./session-sources.js"; import type { SessionSources } from "./session-sources.js"; import { StreamErrorWatcher } from "./stream-error-watcher.js"; @@ -329,6 +330,8 @@ export class SessionManager { private readonly deletingSessions = new Set(); /** Per-Agent config generation (key = agentKey), bumped by invalidateAgentRuntimes on vault updates. */ private readonly agentGenerations = new Map(); + /** Open streaming fragments of running sessions (fed by drive, served to GET /messages; see live-tail.ts). */ + private readonly liveTail = new LiveTailTracker(); private readonly sweepTimer: NodeJS.Timeout; constructor(private readonly deps: SessionManagerDeps) { @@ -356,6 +359,16 @@ export class SessionManager { return this.entries.get(sessionId)?.followUps.length ?? 0; } + /** + * Live tail of a running session: one synthetic `partial_* start` OmniMessage per open + * streaming fragment, carrying the full accumulated content so far (see live-tail.ts). + * Empty when idle or when nothing is streaming. GET /messages attaches this (with a + * channel cursor) so a client joining mid-stream can render the in-progress message. + */ + liveFragments(sessionId: string): OmniMessage[] { + return this.liveTail.fragments(sessionId); + } + /** Number of Sessions for this Agent that are currently running / compacting. */ activeCountForAgent(projectId: string, agentId: string): number { let n = 0; @@ -944,6 +957,10 @@ export class SessionManager { } } } + // Live-tail bookkeeping in the same synchronous tick as the publish below: the + // messages endpoint captures "channel cursor + open fragments" between two + // publishes, so the pair is always a consistent snapshot (see live-tail.ts). + this.liveTail.observe(entry.sessionId, msg); // Re-fetch the channel before every publish (matches publishEvent): the channel // may have been recycled and recreated during a long wait on approval, and // holding a stale reference would send output to an orphaned, detached channel. @@ -966,6 +983,9 @@ export class SessionManager { } finally { // Wrap-up: persist any still-pending LLM failure and clear the tool-name cache (the watcher's state doesn't carry across runs). watcher?.close(); + // The run is over: no fragment will ever continue, so drop the live tail before the + // idle flip (GET /messages stops attaching `live` the moment status reads idle). + this.liveTail.clear(entry.sessionId); entry.approvals.denyAll(); entry.status = "idle"; entry.abort = null; diff --git a/packages/server/test/channel.test.ts b/packages/server/test/channel.test.ts index a64855e..7c28ced 100644 --- a/packages/server/test/channel.test.ts +++ b/packages/server/test/channel.test.ts @@ -26,6 +26,18 @@ describe("channel", () => { expect(JSON.parse(seen[0]!.data)).toEqual({ a: 1 }); }); + it("lastEventId tracks the most recently assigned seq (unicast included; seq 0 when none)", () => { + const ch = new Channel(); + expect(ch.lastEventId).toBe(`${ch.epoch}-0`); + ch.publish("a"); + expect(ch.lastEventId).toBe(`${ch.epoch}-1`); + // Unicast (sendTo) consumes a seq too: lastEventId is a position marker, not a buffer key. + ch.sendTo(() => {}, "hello", "server_event"); + expect(ch.lastEventId).toBe(`${ch.epoch}-2`); + ch.publish("b"); + expect(ch.lastEventId).toBe(`${ch.epoch}-3`); + }); + it("no longer receives after unsubscribe", () => { const ch = new Channel(); const seen: ChannelEvent[] = []; diff --git a/packages/server/test/live-tail.test.ts b/packages/server/test/live-tail.test.ts new file mode 100644 index 0000000..ecafbad --- /dev/null +++ b/packages/server/test/live-tail.test.ts @@ -0,0 +1,174 @@ +/** + * LiveTailTracker unit tests: open-fragment accumulation into synthetic `partial_* start` + * messages, closure on stop / matching complete message / clear, origin preservation for + * subagent fragments, and the tail-keeping cap. + */ +import { describe, expect, it } from "vitest"; +import { + assistantText, + partialText, + partialThinking, + partialToolCall, + partialToolCallOutput, + requestBegin, + thinkingMessage, + toolCall, + toolCallOutput, + userText, + withOrigin, +} from "@prismshadow/penguin-core"; +import type { + PartialTextPayload, + PartialToolCallOutputPayload, + PartialToolCallPayload, +} from "@prismshadow/penguin-core"; +import { LiveTailTracker } from "../src/runtime/live-tail.js"; + +const SID = "s1"; + +describe("live-tail", () => { + it("accumulates text/thinking/tool-call/tool-output fragments into synthetic starts with full prefixes", () => { + const t = new LiveTailTracker(); + t.observe(SID, partialThinking("start")); + t.observe(SID, partialThinking("delta", "let me ")); + t.observe(SID, partialThinking("delta", "think")); + t.observe(SID, partialText("start", "")); + t.observe(SID, partialText("delta", "Hel")); + t.observe(SID, partialText("delta", "lo")); + t.observe(SID, partialToolCall({ eventType: "start", name: "exec_command", toolCallId: "t1" })); + t.observe( + SID, + partialToolCall({ + eventType: "delta", + name: "", + arguments: '{"cmd":"ls"}', + toolCallId: "t1", + }), + ); + + const frags = t.fragments(SID); + expect(frags.map((f) => (f.payload as { type: string }).type)).toEqual([ + "partial_thinking", + "partial_text", + "partial_tool_call", + ]); + // Every synthetic message is a start carrying the accumulated content. + for (const f of frags) { + expect((f.payload as { event_type: string }).event_type).toBe("start"); + } + expect((frags[0]!.payload as { thinking: string }).thinking).toBe("let me think"); + expect((frags[1]!.payload as PartialTextPayload).text).toBe("Hello"); + const call = frags[2]!.payload as PartialToolCallPayload; + expect(call.name).toBe("exec_command"); + expect(call.arguments).toBe('{"cmd":"ls"}'); + expect(call.tool_call_id).toBe("t1"); + }); + + it("stop closes the fragment; a matching complete message also closes it (stop-less closure)", () => { + const t = new LiveTailTracker(); + t.observe(SID, partialText("start", "a")); + t.observe(SID, partialThinking("start", "b")); + expect(t.fragments(SID)).toHaveLength(2); + t.observe(SID, partialText("stop")); + expect(t.fragments(SID)).toHaveLength(1); + // Complete thinking closes the open thinking fragment even without a stop. + t.observe(SID, thinkingMessage("b (complete)")); + expect(t.fragments(SID)).toEqual([]); + }); + + it("a complete tool_call/tool_call_output closes only the fragment with the same tool_call_id", () => { + const t = new LiveTailTracker(); + t.observe(SID, partialToolCallOutput({ eventType: "start", toolCallId: "t1" })); + t.observe( + SID, + partialToolCallOutput({ eventType: "delta", output: "line 1\n", toolCallId: "t1" }), + ); + t.observe(SID, partialToolCall({ eventType: "start", name: "x", toolCallId: "t2" })); + t.observe(SID, toolCallOutput({ output: "other", toolCallId: "t9" })); + expect(t.fragments(SID)).toHaveLength(2); + t.observe(SID, toolCallOutput({ output: "line 1\nline 2\n", toolCallId: "t1" })); + t.observe(SID, toolCall({ name: "x", arguments: "{}", toolCallId: "t2" })); + expect(t.fragments(SID)).toEqual([]); + }); + + it("a user text (steering) does not close the model's open text fragment; events are ignored", () => { + const t = new LiveTailTracker(); + t.observe(SID, partialText("start", "prefix")); + t.observe(SID, userText("steer this way")); + t.observe(SID, requestBegin()); + const frags = t.fragments(SID); + expect(frags).toHaveLength(1); + expect((frags[0]!.payload as PartialTextPayload).text).toBe("prefix"); + // The assistant's complete text does close it. + t.observe(SID, assistantText("prefix done")); + expect(t.fragments(SID)).toEqual([]); + }); + + it("preserves origin: subagent fragments are keyed and emitted with their origin chain", () => { + const t = new LiveTailTracker(); + t.observe(SID, partialText("start", "parent ")); + t.observe(SID, withOrigin(partialText("start", "child "), "c1")); + t.observe(SID, withOrigin(partialText("delta", "progress"), "c1")); + // The parent's complete text closes only the parent fragment (same type, different origin). + t.observe(SID, assistantText("parent done")); + const frags = t.fragments(SID); + expect(frags).toHaveLength(1); + expect(frags[0]!.origin).toEqual(["c1"]); + expect((frags[0]!.payload as PartialTextPayload).text).toBe("child progress"); + }); + + it("tool-output images ride the synthetic start (whole-set semantics)", () => { + const dataUrl = "data:image/png;base64,AAAA"; + const t = new LiveTailTracker(); + t.observe(SID, partialToolCallOutput({ eventType: "start", toolCallId: "t1" })); + t.observe( + SID, + partialToolCallOutput({ + eventType: "delta", + output: "image/png", + toolCallId: "t1", + images: [dataUrl], + }), + ); + const frag = t.fragments(SID)[0]!.payload as PartialToolCallOutputPayload; + expect(frag.output).toBe("image/png"); + expect(frag.images).toEqual([dataUrl]); + }); + + it("caps a runaway fragment by keeping the TAIL (the seed then joins the live deltas seamlessly)", () => { + const t = new LiveTailTracker(); + t.observe(SID, partialText("start", "")); + const chunk = "x".repeat(64 * 1024); + for (let i = 0; i < 12; i += 1) t.observe(SID, partialText("delta", chunk)); // 768KB total + t.observe(SID, partialText("delta", "THE-TAIL")); + const text = (t.fragments(SID)[0]!.payload as PartialTextPayload).text; + expect(text.length).toBeLessThanOrEqual(512 * 1024 + 64 * 1024); + expect(text.endsWith("THE-TAIL")).toBe(true); + }); + + it("clear drops every fragment of the session; other sessions are unaffected", () => { + const t = new LiveTailTracker(); + t.observe(SID, partialText("start", "a")); + t.observe("s2", partialText("start", "b")); + t.clear(SID); + expect(t.fragments(SID)).toEqual([]); + expect(t.fragments("s2")).toHaveLength(1); + }); + + it("a lenient delta without a start still opens a fragment (mirrors core PartialAggregator); a bare stop is a no-op", () => { + const t = new LiveTailTracker(); + t.observe(SID, partialText("stop")); + expect(t.fragments(SID)).toEqual([]); + t.observe(SID, partialText("delta", "late")); + expect((t.fragments(SID)[0]!.payload as PartialTextPayload).text).toBe("late"); + }); + + it("a reopening start replaces the previous same-key fragment", () => { + const t = new LiveTailTracker(); + t.observe(SID, partialText("start", "first")); + t.observe(SID, partialText("start", "second")); + const frags = t.fragments(SID); + expect(frags).toHaveLength(1); + expect((frags[0]!.payload as PartialTextPayload).text).toBe("second"); + }); +}); diff --git a/packages/server/test/messages-live.test.ts b/packages/server/test/messages-live.test.ts new file mode 100644 index 0000000..1db39fb --- /dev/null +++ b/packages/server/test/messages-live.test.ts @@ -0,0 +1,136 @@ +/** + * GET /api/sessions/:id/messages live-tail integration: while a Task runs, the response + * carries `live` — a channel cursor plus one synthetic `partial_* start` per open + * streaming fragment (origin preserved for subagent fragments); once the run ends the + * field disappears and the tail is cleared. + */ +import { afterEach, beforeEach, describe, expect, it } from "vitest"; +import { + approvalDecision, + assistantText, + partialText, + partialThinking, + partialToolCallOutput, + toolCall, + toolCallOutput, + userText, + withOrigin, +} from "@prismshadow/penguin-core"; +import type { ApproveFn, OmniMessage } from "@prismshadow/penguin-core"; +import type { MessagesResponse } from "../src/api/types.js"; +import type { SessionRow } from "../src/db/repos/sessions.js"; +import type { RuntimeSession } from "../src/runtime/session-manager.js"; +import { apiClient, createTestApp, provisionUser, waitFor } from "./helpers.js"; +import type { TestApp } from "./helpers.js"; + +const SID = "session-2026-07-25-10-00-00-cafe0001"; + +/** + * Fake session frozen mid-stream: it streams a thinking fragment, a closed text segment, + * a tool call with a partially streamed output, and a subagent text fragment, then blocks + * on approval (always-ask) so the Task stays running while the test inspects /messages. + */ +function midStreamFakeSession(sessionId: string): RuntimeSession { + return { + sessionId, + toolPermission: () => "rw", + generateTitle: async () => ({ title: null, usage: null }), + compactability: () => "ok" as const, + steer: () => false, + async *run(_input: OmniMessage[], opts: { approve: ApproveFn; signal: AbortSignal }) { + yield partialThinking("start"); + yield partialThinking("delta", "let me think"); + // A closed segment: start+delta+stop+complete leaves no open fragment behind. + yield partialText("start", ""); + yield partialText("delta", "Hello"); + yield partialText("stop"); + yield assistantText("Hello"); + const tc = toolCall({ name: "exec_command", arguments: '{"cmd":"x"}', toolCallId: "tc-lv" }); + yield tc; + yield withOrigin(partialText("start", "child says "), "child-1"); + yield partialToolCallOutput({ eventType: "start", toolCallId: "tc-lv" }); + yield partialToolCallOutput({ eventType: "delta", output: "line 1\n", toolCallId: "tc-lv" }); + const decision = await opts.approve(tc); // blocks until the test decides + yield approvalDecision(decision, "tc-lv"); + yield toolCallOutput({ output: "line 1\nline 2\n", toolCallId: "tc-lv" }); + yield assistantText("done"); + }, + async *compact() {}, + }; +} + +describe("messages live tail", () => { + let t: TestApp; + let cookie: string; + let row: SessionRow; + + beforeEach(async () => { + t = await createTestApp(); + ({ cookie } = await provisionUser(t.app, "livetailer")); + row = { + sessionId: SID, + projectId: "livetailer-default_project", + agentId: "default_agent", + modelId: "m1", + provider: "custom", + workspace: "/tmp/w", + approvalMode: "always-ask", + title: null, + createdAt: new Date().toISOString(), + }; + t.deps.sessionsRepo.insert(row); + }); + afterEach(async () => { + await t.cleanup(); + }); + + const getMessages = async (): Promise => { + const res = await apiClient(t.app, cookie).get(`/api/sessions/${SID}/messages`); + expect(res.status).toBe(200); + return (await res.json()) as MessagesResponse; + }; + + it("no `live` field while idle", async () => { + const body = await getMessages(); + expect(body.messages).toEqual([]); + expect(body.live).toBeUndefined(); + }); + + it("while running: `live` carries the cursor and one synthetic start per open fragment (origin preserved); it disappears once the run ends", async () => { + t.deps.manager.adopt(row, midStreamFakeSession(SID)); + await t.deps.manager.startTask(SID, [userText("go")]); + await waitFor(() => t.deps.manager.pendingApprovalCount(SID) === 1); + + const body = await getMessages(); + expect(body.live).toBeDefined(); + expect(body.live!.cursor).toMatch(/^[0-9a-f]{8}-\d+$/); + const frags = body.live!.fragments.map((f) => ({ + type: (f.payload as { type: string }).type, + origin: f.origin, + })); + // Closed segments (the text) leave nothing; the thinking, the subagent text and the + // tool output are still open. All are `start` events with accumulated content. + expect(frags).toEqual([ + { type: "partial_thinking", origin: undefined }, + { type: "partial_text", origin: ["child-1"] }, + { type: "partial_tool_call_output", origin: undefined }, + ]); + for (const f of body.live!.fragments) { + expect((f.payload as { event_type: string }).event_type).toBe("start"); + } + expect((body.live!.fragments[0]!.payload as { thinking: string }).thinking).toBe( + "let me think", + ); + expect((body.live!.fragments[1]!.payload as { text: string }).text).toBe("child says "); + const out = body.live!.fragments[2]!.payload as { output: string; tool_call_id: string }; + expect(out.output).toBe("line 1\n"); + expect(out.tool_call_id).toBe("tc-lv"); + + // Let the run finish: the tail is cleared and `live` disappears. + expect(t.deps.manager.decideApproval(SID, "tc-lv", "deny")).toBe(true); + await waitFor(() => t.deps.manager.statusOf(SID) === "idle"); + expect(t.deps.manager.liveFragments(SID)).toEqual([]); + const after = await getMessages(); + expect(after.live).toBeUndefined(); + }); +}); diff --git a/packages/web/e2e/mock-llm.mjs b/packages/web/e2e/mock-llm.mjs index 0a070ae..5f5fc8e 100644 --- a/packages/web/e2e/mock-llm.mjs +++ b/packages/web/e2e/mock-llm.mjs @@ -5,6 +5,10 @@ * - files-card probe ("files card test") -> text with two backtick paths (one real, one missing) * - subagent's own turn (its prompt is the only user text) -> final text * - parent asked to delegate ("run a subagent") -> tool_use(run_subagent) + * - "slow stream test" -> tool_use(exec_command) with a command that prints one line + * every 200ms for ~8s (reload-midstream.spec reloads while its output streams) + * - "slow text test" -> a long text streamed one delta every 200ms for ~8s + * (reload-midstream.spec reloads while the TEXT streams) * - last message has tool_result -> final text (turn 2) * - otherwise (first user turn) -> thinking + text + tool_use(exec_command) */ @@ -250,6 +254,46 @@ const server = http.createServer((req, res) => { return; } + // Slow tool-output test case (reload-midstream.spec): a real exec_command whose output + // streams one line every 200ms for ~8s — long enough to reload the page mid-stream and + // watch the output keep growing afterwards. + if (flat.includes("slow stream test") && !hasToolResult) { + block(res, 0, { type: "tool_use", id: "toolu_slow_1", name: "exec_command", input: {} }, [ + { type: "input_json_delta", partial_json: '{"cmd": "for i in $(seq 1 40); do' }, + { type: "input_json_delta", partial_json: ' echo line $i; sleep 0.2; done"}' }, + ]); + messageStop(res, "tool_use", 14); + return; + } + + // Slow TEXT test case (reload-midstream.spec): a single text block streamed one delta + // every 200ms (~8s total) so the page can be reloaded while the assistant text is + // still streaming. Single turn: ends with end_turn, no tool call. + if (flat.includes("slow text test")) { + sse(res, "content_block_start", { + type: "content_block_start", + index: 0, + content_block: { type: "text", text: "" }, + }); + let n = 0; + const tick = () => { + n += 1; + sse(res, "content_block_delta", { + type: "content_block_delta", + index: 0, + delta: { type: "text_delta", text: `chunk-${n} ` }, + }); + if (n < 40) { + setTimeout(tick, 200); + } else { + sse(res, "content_block_stop", { type: "content_block_stop", index: 0 }); + messageStop(res, "end_turn", 80); + } + }; + tick(); + return; + } + if (wantsSubagent && !hasToolResult) { block(res, 0, { type: "tool_use", id: "toolu_mock_sub", name: "run_subagent", input: {} }, [ { type: "input_json_delta", partial_json: '{"prompt": ' }, diff --git a/packages/web/e2e/reload-midstream.spec.mjs b/packages/web/e2e/reload-midstream.spec.mjs new file mode 100644 index 0000000..a633cc9 --- /dev/null +++ b/packages/web/e2e/reload-midstream.spec.mjs @@ -0,0 +1,158 @@ +/** + * Reload mid-stream: while a Task streams a long tool output (or a long assistant text), + * refreshing the page must bring the in-progress message straight back — the already + * streamed prefix visible promptly (before the stream finishes), the content still + * growing live, and the final state identical to a run that was never reloaded. + * + * Mechanics under test: GET /messages returns `live` ({cursor, fragments}) while the + * Task runs; the frontend seeds the synthetic `partial_* start` fragments on top of + * history and drops the buffered partials the snapshot already covers. + * + * The LLM is mock-llm.mjs: "slow stream test" makes it call exec_command with a command + * that prints one line every 200ms for ~8s (real tool execution, really streamed); + * "slow text test" streams a 40-chunk text one delta every 200ms. + */ +import { test, expect } from "@playwright/test"; +import { provisionAndLogin } from "./auth.mjs"; + +const BASE = process.env.BASE_URL; +const MOCK = process.env.MOCK_URL; +const U = "reloaduser"; +const P = "password123"; + +/** Create a session for the user's auto-provisioned project (models PUT is idempotent). */ +async function createSession(page, approvalMode) { + const projects = await (await page.request.get(`${BASE}/api/projects`)).json(); + const projectId = projects.projects[0].projectId; + const put = await page.request.put(`${BASE}/api/projects/${projectId}/models`, { + data: { + defaultModel: { provider: "custom", modelId: "claude-4-8" }, + models: [ + { + provider: "custom", + modelId: "claude-4-8", + apiKey: "sk-mock", + baseUrl: MOCK, + contextWindow: 200000, + }, + ], + }, + }); + expect(put.ok(), "put models").toBeTruthy(); + const res = await page.request.post( + `${BASE}/api/projects/${projectId}/agents/default_agent/sessions`, + { data: { provider: "custom", modelId: "claude-4-8", approvalMode } }, + ); + expect(res.ok(), `create session: ${await res.text()}`).toBeTruthy(); + return (await res.json()).session.sessionId; +} + +/** + * Keep the running work group open past the end of the turn (it auto-collapses at turn + * end unless the user toggled it — same trick as chat.spec), then expand the + * exec_command tool card and return the output
 locator.
+ */
+async function openToolOutput(page) {
+  const group = page
+    .locator("button[aria-expanded]")
+    .filter({ hasText: /运行中|运行完毕/ })
+    .first();
+  await expect(group).toBeVisible();
+  if ((await group.textContent())?.includes("运行中")) {
+    await group.click(); // toggle → marks the group user-toggled
+    await group.click(); // toggle back → deliberately kept open, survives turn end
+  } else if ((await group.getAttribute("aria-expanded")) !== "true") {
+    await group.click(); // finished and collapsed: open it to reach the card
+  }
+  const toolCard = page
+    .locator("button[aria-expanded]")
+    .filter({ hasText: "exec_command" })
+    .first();
+  await expect(toolCard).toBeVisible();
+  if ((await toolCard.getAttribute("aria-expanded")) !== "true") await toolCard.click();
+  await expect(toolCard).toHaveAttribute("aria-expanded", "true");
+  // Two 
s live in the expanded card (arguments, then output); the output one is the
+  // only one containing a literal "line 1".
+  return page.locator("pre", { hasText: /line 1\b/ }).first();
+}
+
+/** Send the "slow stream test" prompt, approve exec_command, and return the output 
. */
+async function startSlowToolRun(page, sessionId) {
+  await page.goto(`${BASE}/chat/${sessionId}`);
+  const ta = page.getByPlaceholder(/输入消息/);
+  await ta.waitFor();
+  await ta.fill("slow stream test");
+  await page.getByRole("button", { name: "发送" }).click();
+  await expect(page.getByText("exec_command").first()).toBeVisible();
+  await page.getByRole("button", { name: "允许" }).click();
+  const outputPre = await openToolOutput(page);
+  await expect(outputPre).toContainText("line 3");
+  return outputPre;
+}
+
+test("in-progress tool output survives a reload: prefix back promptly, still growing, final state matches a never-reloaded run", async ({
+  page,
+}) => {
+  await provisionAndLogin(page.request, U, P);
+
+  // --- Control run (never reloaded): capture the final tool output to compare against ---
+  const controlId = await createSession(page, "always-ask");
+  await startSlowToolRun(page, controlId);
+  await expect(page.getByText("Command finished; the result looks as expected.")).toBeVisible({
+    timeout: 30_000,
+  });
+  const controlPre = await openToolOutput(page);
+  const controlOutput = await controlPre.textContent();
+  expect(controlOutput).toContain("line 40");
+
+  // --- The run under test: reload while the output is streaming ---
+  const sessionId = await createSession(page, "always-ask");
+  await startSlowToolRun(page, sessionId);
+  await page.reload();
+
+  // (a) The already-streamed prefix is back promptly — well before the ~8s command ends.
+  const outputPre = await openToolOutput(page);
+  await expect(outputPre).toContainText("line 2", { timeout: 4000 });
+  // Not finished yet: what we see is genuinely the in-progress stream, not the final state.
+  await expect(outputPre).not.toContainText("line 40");
+
+  // (b) The output keeps growing live on the same connection.
+  await expect(outputPre).toContainText("line 35", { timeout: 15_000 });
+
+  // (c) After completion the final state matches the never-reloaded run exactly
+  //     (every line exactly once — no lost prefix, no duplicated overlap).
+  await expect(page.getByText("Command finished; the result looks as expected.")).toBeVisible({
+    timeout: 30_000,
+  });
+  const reloadedOutput = await (await openToolOutput(page)).textContent();
+  expect(reloadedOutput).toBe(controlOutput);
+});
+
+test("in-progress assistant TEXT survives a reload and keeps streaming", async ({ page }) => {
+  await provisionAndLogin(page.request, U, P);
+  const sessionId = await createSession(page, "allow-all");
+
+  await page.goto(`${BASE}/chat/${sessionId}`);
+  const ta = page.getByPlaceholder(/输入消息/);
+  await ta.waitFor();
+  await ta.fill("slow text test");
+  await page.getByRole("button", { name: "发送" }).click();
+
+  const reply = page.locator(".md-body", { hasText: /chunk-1\s/ }).first();
+  await expect(reply).toContainText("chunk-3 ");
+  await page.reload();
+
+  // (a) The streamed prefix is back promptly, while the ~8s stream is still going.
+  const reloaded = page.locator(".md-body", { hasText: /chunk-1\s/ }).first();
+  await expect(reloaded).toContainText("chunk-2 ", { timeout: 4000 });
+  await expect(reloaded).not.toContainText("chunk-40");
+
+  // (b) It keeps streaming live, up to the full reply.
+  await expect(reloaded).toContainText("chunk-40", { timeout: 15_000 });
+
+  // (c) No duplicated prefix: the buffered partials the snapshot covered were dropped,
+  //     so every chunk appears exactly once.
+  const text = await reloaded.textContent();
+  expect(text.match(/chunk-5 /g)).toHaveLength(1);
+  expect(text.match(/chunk-1 /g)).toHaveLength(1);
+});
diff --git a/packages/web/src/api/sse.ts b/packages/web/src/api/sse.ts
index 2c20654..c05bdbf 100644
--- a/packages/web/src/api/sse.ts
+++ b/packages/web/src/api/sse.ts
@@ -15,10 +15,15 @@ import type { OmniMessage } from "@prismshadow/penguin-core/omnimessage";
 import type { ServerEvent } from "@prismshadow/penguin-server/api";
 
 export interface StreamHandlers {
-  /** A single OmniMessage (full/streaming/event, envelope as-is). */
-  onOmniMessage: (msg: OmniMessage) => void;
-  /** A single server event. */
-  onServerEvent: (event: ServerEvent) => void;
+  /**
+   * A single OmniMessage (full/streaming/event, envelope as-is). `eventId` is the SSE
+   * event id assigned by the server channel (`-`; null if the event carried
+   * none) — stream-controller uses it to align buffered events with the live-tail cursor
+   * that GET /messages returns.
+   */
+  onOmniMessage: (msg: OmniMessage, eventId: string | null) => void;
+  /** A single server event (`eventId`: same as onOmniMessage). */
+  onServerEvent: (event: ServerEvent, eventId: string | null) => void;
   /** Connection established (including a successful auto-reconnect). */
   onOpen?: () => void;
   /**
@@ -37,14 +42,15 @@ function subscribe(url: string, handlers: StreamHandlers): StreamConnection {
   const source = new EventSource(url);
   source.onmessage = (e: MessageEvent) => {
     try {
-      handlers.onOmniMessage(JSON.parse(e.data) as OmniMessage);
+      // Every server event carries an `id:` line; lastEventId is "" only if none did.
+      handlers.onOmniMessage(JSON.parse(e.data) as OmniMessage, e.lastEventId || null);
     } catch {
       // Ignore lines that fail to parse (the protocol guarantees single-line JSON data, so this shouldn't normally happen).
     }
   };
   source.addEventListener("server_event", (e: MessageEvent) => {
     try {
-      handlers.onServerEvent(JSON.parse(e.data) as ServerEvent);
+      handlers.onServerEvent(JSON.parse(e.data) as ServerEvent, e.lastEventId || null);
     } catch {
       // Same as above.
     }
diff --git a/packages/web/src/features/chat/use-session-stream.ts b/packages/web/src/features/chat/use-session-stream.ts
index 8d82550..870f265 100644
--- a/packages/web/src/features/chat/use-session-stream.ts
+++ b/packages/web/src/features/chat/use-session-stream.ts
@@ -146,7 +146,9 @@ export function useSessionStream(
     setPendingTick((t) => t + 1);
 
     const controller = createStreamController({
-      loadMessages: async () => (await getMessages(sessionId)).messages,
+      // The whole response rides through: `live` (in-progress stream tail) lets the
+      // controller seed the currently streaming message after a reload (see stream-controller).
+      loadMessages: () => getMessages(sessionId),
       onTaskState: setTaskState,
       onQueuedFollowUps: setQueuedFollowUps,
       onLoading: setLoading,
diff --git a/packages/web/src/lib/omni/stream-controller.ts b/packages/web/src/lib/omni/stream-controller.ts
index 5a1b305..6a3cbe5 100644
--- a/packages/web/src/lib/omni/stream-controller.ts
+++ b/packages/web/src/lib/omni/stream-controller.ts
@@ -20,11 +20,18 @@
  *   when a resent approval_request can't find its tool card (sub-session
  *   messages aren't written to the parent Trace, so the card can be missing
  *   after a reload), the toolCall carried by the event (with origin) is fed
- *   to the reducer to rebuild the nested card, making the sub-session's approval visible and decidable.
+ *   to the reducer to rebuild the nested card, making the sub-session's approval visible and decidable;
+ * - Live-tail seeding: while a Task runs, `/messages` also returns `live`
+ *   ({cursor, fragments} — see MessagesLiveTail): buffered partial events at
+ *   or before the cursor are dropped (their content is already accumulated
+ *   inside the fragments), the synthetic `partial_* start` fragments are fed
+ *   through the normal reducer path at the cursor's position, and the rest of
+ *   the buffer replays on top — so the in-progress message survives a refresh
+ *   with its streamed prefix intact and keeps streaming.
  */
 import { isEventMessage, isPartialPayload } from "@prismshadow/penguin-core/omnimessage";
 import type { OmniMessage, ToolCallPayload } from "@prismshadow/penguin-core/omnimessage";
-import type { ServerEvent, SessionStatus } from "@prismshadow/penguin-server/api";
+import type { MessagesLiveTail, ServerEvent, SessionStatus } from "@prismshadow/penguin-server/api";
 import {
   approvalKey,
   buildDedupIndex,
@@ -46,11 +53,23 @@ export interface PendingApproval {
   origin?: string[];
 }
 
-type BufferedEvent = { kind: "omni"; msg: OmniMessage } | { kind: "server"; ev: ServerEvent };
+/** Buffered stream event; `id` is the SSE event id (`-`, null when unknown); `seeded` marks a live-tail fragment injected at replay time (never carried over into a later round's buffer — the new round refetches fresh fragments). */
+type BufferedEvent =
+  | { kind: "omni"; msg: OmniMessage; id: string | null; seeded?: boolean }
+  | { kind: "server"; ev: ServerEvent; id: string | null };
+
+/** Parse a channel event id (`-`); null when malformed (same split rule as the server's Channel.replayAfter). */
+function parseEventId(id: string): { epoch: string; seq: number } | null {
+  const sep = id.lastIndexOf("-");
+  if (sep <= 0) return null;
+  const seq = Number.parseInt(id.slice(sep + 1), 10);
+  if (!Number.isInteger(seq) || seq < 0) return null;
+  return { epoch: id.slice(0, sep), seq };
+}
 
 export interface StreamControllerDeps {
-  /** Fetch history messages (GET /api/sessions/:id/messages). */
-  loadMessages: () => Promise;
+  /** Fetch history messages (GET /api/sessions/:id/messages), including the live tail while running. */
+  loadMessages: () => Promise<{ messages: OmniMessage[]; live?: MessagesLiveTail }>;
   /** Authoritative running state from the stream (covers both the subscription snapshot and transition events). */
   onTaskState: (state: SessionStatus) => void;
   /** Queued follow-up count carried on task_state events (absent on old servers -> 0). */
@@ -78,10 +97,10 @@ export interface StreamController {
   load: () => Promise;
   /** Retry entry point after a history load failure (keeps the buffer, refetches history). */
   retry: () => Promise;
-  /** SSE OmniMessage entry point. */
-  handleOmni: (msg: OmniMessage) => void;
-  /** SSE server-event entry point. */
-  handleServer: (ev: ServerEvent) => void;
+  /** SSE OmniMessage entry point (`eventId`: the SSE event id, used for live-tail cursor alignment). */
+  handleOmni: (msg: OmniMessage, eventId?: string | null) => void;
+  /** SSE server-event entry point (`eventId`: same as handleOmni). */
+  handleServer: (ev: ServerEvent, eventId?: string | null) => void;
   /** Register that this end clicked an approval ("manual" label, persists across resync rebuilds). */
   markLocalDecision: (toolCallId: string) => void;
   /** Remove a pending approval by its composite key (optimistic update; also removed as a fallback when the event arrives). */
@@ -101,6 +120,8 @@ export function createStreamController(deps: StreamControllerDeps): StreamContro
   let buffer: BufferedEvent[] = [];
   /** The most recent in-stream task_state (null = the snapshot hasn't arrived yet; history finalization trusts only this value). */
   let streamStatus: SessionStatus | null = null;
+  /** The most recent SSE event id seen on this connection (null = none yet): identifies the channel epoch the live-tail cursor must match. */
+  let lastEventId: string | null = null;
   /** Load epoch: incremented on rebuild/retry; any replay or finalization from an older epoch is discarded. */
   let epoch = 0;
   /** Whether the most recent load failed (retry only takes effect after a failure, to avoid mistakenly replaying history). */
@@ -200,9 +221,42 @@ export function createStreamController(deps: StreamControllerDeps): StreamContro
     await load(epoch, createStreamModel(localDecisions));
   };
 
+  /**
+   * Weave the live tail into the buffered replay (see MessagesLiveTail in the server API):
+   * entries at/or before the cursor come first with their partials dropped (the fragment
+   * snapshot already accumulates them; completes still replay — the dedup gate decides),
+   * then the synthetic fragment starts at the cursor's position, then everything after the
+   * cursor. Skipped entirely when the cursor's epoch doesn't match the epoch of the events
+   * seen on this connection (channel recycled/server restarted between the two requests —
+   * seq comparison would be meaningless; the resync flow covers that path).
+   */
+  const weaveLiveTail = (replay: BufferedEvent[], live: MessagesLiveTail): BufferedEvent[] => {
+    const cursor = parseEventId(live.cursor);
+    if (!cursor) return replay;
+    const seen = lastEventId === null ? null : parseEventId(lastEventId);
+    if (seen !== null && seen.epoch !== cursor.epoch) return replay;
+    const pre: BufferedEvent[] = [];
+    const post: BufferedEvent[] = [];
+    for (const e of replay) {
+      const eid = e.id === null ? null : parseEventId(e.id);
+      if (eid !== null && eid.epoch === cursor.epoch && eid.seq <= cursor.seq) {
+        if (e.kind !== "omni" || !isPartialPayload(e.msg.payload)) pre.push(e);
+      } else {
+        post.push(e);
+      }
+    }
+    const seeds: BufferedEvent[] = live.fragments.map((msg) => ({
+      kind: "omni",
+      msg,
+      id: null,
+      seeded: true,
+    }));
+    return [...pre, ...seeds, ...post];
+  };
+
   const load = async (currentEpoch: number, freshModel?: StreamModel): Promise => {
     try {
-      const messages = await deps.loadMessages();
+      const { messages, live } = await deps.loadMessages();
       if (disposed || currentEpoch !== epoch) return;
       // Rebuild path: make the freshly-built model visible only now, atomically — the old model
       // stayed on screen throughout the refetch above (see rebuild). Initial load / retry pass no
@@ -211,8 +265,9 @@ export function createStreamController(deps: StreamControllerDeps): StreamContro
       const target = model;
       pushMessages(target, messages, now());
       const dedup = buildDedupIndex(messages, 100);
-      // Replay the buffer (events that arrived while fetching history), with dedup.
-      const replay = buffer;
+      // Replay the buffer (events that arrived while fetching history), with dedup; while a
+      // Task runs, the live tail is woven in so the in-progress message is seeded too.
+      const replay = live !== undefined ? weaveLiveTail(buffer, live) : buffer;
       buffer = [];
       for (let i = 0; i < replay.length; i += 1) {
         const e = replay[i]!;
@@ -222,8 +277,10 @@ export function createStreamController(deps: StreamControllerDeps): StreamContro
         if (currentEpoch !== epoch) {
           // A rebuild was triggered mid-replay (e.g. resync_required): this
           // round is discarded, and the remaining events are handed off to
-          // the new round's buffer — phase is left unchanged and the old buffer is never fed to the new model.
-          buffer.push(...replay.slice(i + 1));
+          // the new round's buffer — phase is left unchanged and the old buffer is never
+          // fed to the new model. Seeded fragments are snapshot-bound to THIS round and
+          // are dropped instead of carried over (the new round refetches fresh ones).
+          buffer.push(...replay.slice(i + 1).filter((r) => r.kind !== "omni" || !r.seeded));
           return;
         }
       }
@@ -269,17 +326,19 @@ export function createStreamController(deps: StreamControllerDeps): StreamContro
       epoch += 1;
       await load(epoch, createStreamModel(localDecisions));
     },
-    handleOmni: (msg) => {
+    handleOmni: (msg, eventId = null) => {
       if (disposed) return;
+      if (eventId !== null) lastEventId = eventId;
       if (phase === "buffering") {
-        buffer.push({ kind: "omni", msg });
+        buffer.push({ kind: "omni", msg, id: eventId });
         return;
       }
       feedOmni(msg, null);
       deps.onModelChange();
     },
-    handleServer: (ev) => {
+    handleServer: (ev, eventId = null) => {
       if (disposed) return;
+      if (eventId !== null) lastEventId = eventId;
       // session_title / session_created only affect list display (unrelated
       // to the view model/history): forwarded immediately at any phase, never buffered.
       if (ev.type === "session_title") {
@@ -299,7 +358,7 @@ export function createStreamController(deps: StreamControllerDeps): StreamContro
           deps.onTaskState(ev.state);
           deps.onQueuedFollowUps?.(ev.queued ?? 0);
         }
-        buffer.push({ kind: "server", ev });
+        buffer.push({ kind: "server", ev, id: eventId });
         return;
       }
       handleServer(ev);
diff --git a/packages/web/src/lib/omni/stream-model.ts b/packages/web/src/lib/omni/stream-model.ts
index 2dad16f..e9b732d 100644
--- a/packages/web/src/lib/omni/stream-model.ts
+++ b/packages/web/src/lib/omni/stream-model.ts
@@ -739,6 +739,9 @@ function handlePartial(model: StreamModel, p: PartialModelPayload, tsMs?: number
       if (p.event_type === "start") {
         card.outputStreaming = true;
         if (p.output) card.output += p.output;
+        // A live-tail synthetic start (mid-stream join seed) may already carry the image
+        // set; same whole-set semantics as the delta branch below.
+        if (p.images && p.images.length > 0) card.images = p.images;
         return;
       }
       if (!card.outputStreaming) return; // orphan delta/stop
diff --git a/packages/web/test/stream-controller.test.ts b/packages/web/test/stream-controller.test.ts
index 49db10f..8d30211 100644
--- a/packages/web/test/stream-controller.test.ts
+++ b/packages/web/test/stream-controller.test.ts
@@ -9,17 +9,19 @@ import { describe, expect, it } from "vitest";
 import {
   approvalDecision,
   assistantText,
+  partialText,
+  partialToolCallOutput,
   tokenUsage,
   toolCall,
   userText,
   withOrigin,
 } from "@prismshadow/penguin-core/omnimessage";
 import type { OmniMessage, TokenCounts } from "@prismshadow/penguin-core/omnimessage";
-import type { ServerEvent, SessionStatus } from "@prismshadow/penguin-server/api";
+import type { MessagesLiveTail, ServerEvent, SessionStatus } from "@prismshadow/penguin-server/api";
 import { createStreamController } from "../src/lib/omni/stream-controller";
 import type { StreamController } from "../src/lib/omni/stream-controller";
 import { approvalKey, findToolCard } from "../src/lib/omni/stream-model";
-import type { ToolCallItem } from "../src/lib/omni/stream-model";
+import type { AssistantTextItem, ToolCallItem } from "../src/lib/omni/stream-model";
 
 /** Override a message timestamp (constructor defaults to the current time). */
 function at(msg: M, ts: string): M {
@@ -39,13 +41,13 @@ interface Harness {
   errors: Array;
   loadings: boolean[];
   loadCalls: () => number;
-  resolveLoad: (messages: OmniMessage[]) => void;
+  resolveLoad: (messages: OmniMessage[], live?: MessagesLiveTail) => void;
   rejectLoad: (err: Error) => void;
 }
 
 function createHarness(): Harness {
   const pendingLoads: Array<{
-    resolve: (m: OmniMessage[]) => void;
+    resolve: (m: { messages: OmniMessage[]; live?: MessagesLiveTail }) => void;
     reject: (e: unknown) => void;
   }> = [];
   const states: SessionStatus[] = [];
@@ -54,7 +56,7 @@ function createHarness(): Harness {
   let calls = 0;
   const controller = createStreamController({
     loadMessages: () =>
-      new Promise((resolve, reject) => {
+      new Promise<{ messages: OmniMessage[]; live?: MessagesLiveTail }>((resolve, reject) => {
         calls += 1;
         pendingLoads.push({ resolve, reject });
       }),
@@ -71,7 +73,8 @@ function createHarness(): Harness {
     errors,
     loadings,
     loadCalls: () => calls,
-    resolveLoad: (messages) => pendingLoads.shift()!.resolve(messages),
+    resolveLoad: (messages, live) =>
+      pendingLoads.shift()!.resolve({ messages, ...(live !== undefined ? { live } : {}) }),
     rejectLoad: (err) => pendingLoads.shift()!.reject(err),
   };
 }
@@ -255,6 +258,115 @@ describe("resync rebuild", () => {
   });
 });
 
+describe("live-tail seeding (reload mid-stream)", () => {
+  it("drops buffered partials at/before the cursor and seeds fragments; later partials continue on top", async () => {
+    const h = createHarness();
+    const p = h.controller.load();
+    h.controller.handleServer({ type: "task_state", state: "running" }, "e1-3");
+    // Already covered by the fragment snapshot (seq <= cursor): must be dropped, or the
+    // seeded prefix would double.
+    h.controller.handleOmni(partialText("start", "Hel"), "e1-4");
+    h.controller.handleOmni(partialText("delta", "lo"), "e1-5");
+    // After the cursor: continues the seeded fragment.
+    h.controller.handleOmni(partialText("delta", " world"), "e1-6");
+    h.resolveLoad([at(userText("question"), "2026-07-05T00:00:00.000Z")], {
+      cursor: "e1-5",
+      fragments: [at(partialText("start", "Hello"), "2026-07-05T00:00:01.000Z")],
+    });
+    await p;
+    const texts = h.controller.model.items.filter((i) => i.kind === "assistant_text");
+    expect(texts).toHaveLength(1);
+    expect((texts[0] as AssistantTextItem).text).toBe("Hello world");
+    expect((texts[0] as AssistantTextItem).streaming).toBe(true);
+  });
+
+  it("tool-output fragment seeds onto the history-built card and keeps streaming", async () => {
+    const h = createHarness();
+    const p = h.controller.load();
+    h.controller.handleServer({ type: "task_state", state: "running" }, "e1-7");
+    h.controller.handleOmni(
+      partialToolCallOutput({ eventType: "delta", output: "line 2\n", toolCallId: "t1" }),
+      "e1-9",
+    );
+    h.resolveLoad(
+      [
+        at(userText("run it"), "2026-07-05T00:00:00.000Z"),
+        at(
+          toolCall({ name: "exec_command", arguments: '{"cmd":"x"}', toolCallId: "t1" }),
+          "2026-07-05T00:00:01.000Z",
+        ),
+      ],
+      {
+        cursor: "e1-8",
+        fragments: [
+          at(
+            partialToolCallOutput({ eventType: "start", output: "line 1\n", toolCallId: "t1" }),
+            "2026-07-05T00:00:02.000Z",
+          ),
+        ],
+      },
+    );
+    await p;
+    const cards = h.controller.model.items.filter((i) => i.kind === "tool_call");
+    expect(cards).toHaveLength(1);
+    const card = cards[0] as ToolCallItem;
+    expect(card.output).toBe("line 1\nline 2\n");
+    expect(card.outputStreaming).toBe(true);
+    expect(card.outputComplete).toBe(false);
+  });
+
+  it("epoch mismatch: neither drops nor seeds (buffered partials replay as-is)", async () => {
+    const h = createHarness();
+    const p = h.controller.load();
+    h.controller.handleOmni(partialText("start", "He"), "old-4");
+    h.controller.handleOmni(partialText("delta", "llo"), "old-5");
+    h.resolveLoad([], {
+      cursor: "new-9",
+      fragments: [partialText("start", "Hello")],
+    });
+    await p;
+    // The buffered start+delta applied normally; the fragment was NOT seeded (a seed on
+    // top of the applied start would produce a second item).
+    const texts = h.controller.model.items.filter((i) => i.kind === "assistant_text");
+    expect(texts).toHaveLength(1);
+    expect((texts[0] as AssistantTextItem).text).toBe("Hello");
+  });
+
+  it("complete messages at/before the cursor are never cursor-dropped: overlap dedup decides", async () => {
+    const h = createHarness();
+    const p = h.controller.load();
+    const inHistory = at(assistantText("done"), "2026-07-05T00:00:01.000Z");
+    // Trace append still in flight at read time: the buffered copy is the only copy.
+    const diskLagged = at(assistantText("lagged"), "2026-07-05T00:00:02.000Z");
+    h.controller.handleOmni(inHistory, "e1-4");
+    h.controller.handleOmni(diskLagged, "e1-5");
+    h.resolveLoad([at(userText("q"), "2026-07-05T00:00:00.000Z"), inHistory], {
+      cursor: "e1-6",
+      fragments: [],
+    });
+    await p;
+    const texts = h.controller.model.items.filter((i) => i.kind === "assistant_text");
+    expect(texts.map((i) => (i as AssistantTextItem).text)).toEqual(["done", "lagged"]);
+  });
+
+  it("subagent fragments preserve origin and land on the nested model", async () => {
+    const h = createHarness();
+    const p = h.controller.load();
+    h.controller.handleServer({ type: "task_state", state: "running" }, "e1-2");
+    h.resolveLoad([], {
+      cursor: "e1-2",
+      fragments: [withOrigin(partialText("start", "sub progress"), "c1")],
+    });
+    await p;
+    const sub = h.controller.model.subagents.get("c1");
+    expect(sub).toBeDefined();
+    const texts = sub!.items.filter((i) => i.kind === "assistant_text");
+    expect(texts).toHaveLength(1);
+    expect((texts[0] as AssistantTextItem).text).toBe("sub progress");
+    expect((texts[0] as AssistantTextItem).streaming).toBe(true);
+  });
+});
+
 describe("history load failure and retry (#6)", () => {
   it("failure surfaces the error and stops loading; retry keeps the buffer (snapshot and initial events are not lost)", async () => {
     const h = createHarness();
diff --git a/packages/web/test/stream-model.test.ts b/packages/web/test/stream-model.test.ts
index c55cb07..10868d7 100644
--- a/packages/web/test/stream-model.test.ts
+++ b/packages/web/test/stream-model.test.ts
@@ -224,6 +224,110 @@ describe("partial aggregation and full-message convergence", () => {
   });
 });
 
+describe("live-tail synthetic starts (mid-stream join seeding)", () => {
+  it("a text start carrying the accumulated prefix opens a streaming item on top of history; deltas continue and the full message replaces", () => {
+    const m = createStreamModel();
+    pushMessage(m, userText("question"));
+    pushMessage(m, partialText("start", "Already streamed prefix"));
+    const item = items(m)[1] as AssistantTextItem;
+    expect(item.kind).toBe("assistant_text");
+    expect(item.text).toBe("Already streamed prefix");
+    expect(item.streaming).toBe(true);
+    pushMessage(m, partialText("delta", " + live tail"));
+    expect(item.text).toBe("Already streamed prefix + live tail");
+    pushMessage(m, partialText("stop"));
+    pushMessage(m, assistantText("Already streamed prefix + live tail."));
+    expect(items(m).filter((i) => i.kind === "assistant_text")).toHaveLength(1);
+    expect(item.text).toBe("Already streamed prefix + live tail.");
+  });
+
+  it("a thinking start carrying the accumulated prefix seeds a streaming thinking item with its start time", () => {
+    const m = createStreamModel();
+    pushMessage(m, at(partialThinking("start", "half a thought"), "2026-07-05T00:00:01.000Z"));
+    const item = items(m)[0] as ThinkingItem;
+    expect(item.thinking).toBe("half a thought");
+    expect(item.streaming).toBe(true);
+    expect(item.startedAtMs).toBe(Date.parse("2026-07-05T00:00:01.000Z"));
+  });
+
+  it("an output start seeds the prefix (and images) onto a call-complete card; a start on an outputComplete card is ignored", () => {
+    const dataUrl = "data:image/png;base64,AAAA";
+    const m = createStreamModel();
+    pushMessage(m, toolCall({ name: "exec_command", arguments: '{"cmd":"x"}', toolCallId: "t1" }));
+    const card = items(m)[0] as ToolCallItem;
+    // Synthetic start carries the accumulated prefix + the whole image set.
+    pushMessage(
+      m,
+      partialToolCallOutput({
+        eventType: "start",
+        output: "line 1\nline 2\n",
+        toolCallId: "t1",
+        images: [dataUrl],
+      }),
+    );
+    expect(card.output).toBe("line 1\nline 2\n");
+    expect(card.outputStreaming).toBe(true);
+    expect(card.images).toEqual([dataUrl]);
+    pushMessage(
+      m,
+      partialToolCallOutput({ eventType: "delta", output: "line 3\n", toolCallId: "t1" }),
+    );
+    expect(card.output).toBe("line 1\nline 2\nline 3\n");
+    // Once the output is complete, a stray synthetic start must not reopen or append.
+    pushMessage(m, toolCallOutput({ output: "final", toolCallId: "t1" }));
+    pushMessage(
+      m,
+      partialToolCallOutput({ eventType: "start", output: "stale", toolCallId: "t1" }),
+    );
+    expect(card.output).toBe("final");
+    expect(card.outputStreaming).toBe(false);
+  });
+
+  it("an arguments start for an id whose call is already complete is ignored (no duplicate card, no reset)", () => {
+    const m = createStreamModel();
+    pushMessage(m, toolCall({ name: "exec_command", arguments: '{"cmd":"x"}', toolCallId: "t1" }));
+    pushMessage(
+      m,
+      partialToolCall({
+        eventType: "start",
+        name: "exec_command",
+        arguments: '{"cmd":"x"}',
+        toolCallId: "t1",
+      }),
+    );
+    expect(items(m)).toHaveLength(1);
+    expect((items(m)[0] as ToolCallItem).argumentsText).toBe('{"cmd":"x"}');
+  });
+
+  it("an arguments start carrying accumulated arguments seeds a card that the full message then completes", () => {
+    const m = createStreamModel();
+    pushMessage(
+      m,
+      partialToolCall({
+        eventType: "start",
+        name: "exec_command",
+        arguments: '{"cmd":"seq 1',
+        toolCallId: "t1",
+      }),
+    );
+    const card = items(m)[0] as ToolCallItem;
+    expect(card.argumentsText).toBe('{"cmd":"seq 1');
+    expect(card.callStreaming).toBe(true);
+    pushMessage(
+      m,
+      partialToolCall({ eventType: "delta", name: "", arguments: ' 40"}', toolCallId: "t1" }),
+    );
+    pushMessage(m, partialToolCall({ eventType: "stop", name: "", toolCallId: "t1" }));
+    pushMessage(
+      m,
+      toolCall({ name: "exec_command", arguments: '{"cmd":"seq 1 40"}', toolCallId: "t1" }),
+    );
+    expect(items(m)).toHaveLength(1);
+    expect(card.argumentsText).toBe('{"cmd":"seq 1 40"}');
+    expect(card.callComplete).toBe(true);
+  });
+});
+
 describe("approvals and events", () => {
   it("approval_decision annotates the matching tool card; locally registered ones are manual, the rest remote", () => {
     const m = createStreamModel();