fix(server,web): live in-progress output survives refresh and reconnect (#77)

Co-authored-by: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
Yaowei Zheng
2026-07-27 22:31:57 +08:00
committed by GitHub
parent 3d3495528c
commit 094ebd6118
18 changed files with 1214 additions and 40 deletions
+27 -4
View File
@@ -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 (`<epoch>-<seq>`):
// 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 `<epoch>-<seq>`;
- 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
+26 -4
View File
@@ -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(`<epoch>-<seq>`):
// 截至该 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 按通道单调递增,形如 `<epoch>-<seq>`;
- 每通道维护有界重放缓冲(最近 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. 转入实时消费。
## 类型导入
+33
View File
@@ -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
* (`<epoch>-<seq>`); 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;
}
// ---------------------------------------------------------------------------
+19 -1
View File
@@ -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<AppEnv> {
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) => {
+11
View File
@@ -73,6 +73,17 @@ export class Channel {
return this.listeners.size;
}
/**
* Id of the most recently assigned event (`<epoch>-<seq>`; 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);
+237
View File
@@ -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<string, Map<string, OpenFragment>>();
/** 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);
}
}
@@ -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<string>();
/** Per-Agent config generation (key = agentKey), bumped by invalidateAgentRuntimes on vault updates. */
private readonly agentGenerations = new Map<string, number>();
/** 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;
+12
View File
@@ -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[] = [];
+174
View File
@@ -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");
});
});
+136
View File
@@ -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<MessagesResponse> => {
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();
});
});
+44
View File
@@ -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": ' },
+158
View File
@@ -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 <pre> 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 <pre>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 <pre>. */
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);
});
+12 -6
View File
@@ -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 (`<epoch>-<seq>`; 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<string>) => {
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<string>) => {
try {
handlers.onServerEvent(JSON.parse(e.data) as ServerEvent);
handlers.onServerEvent(JSON.parse(e.data) as ServerEvent, e.lastEventId || null);
} catch {
// Same as above.
}
@@ -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,
+77 -18
View File
@@ -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 (`<epoch>-<seq>`, 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 (`<epoch>-<seq>`); 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<OmniMessage[]>;
/** 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<void>;
/** Retry entry point after a history load failure (keeps the buffer, refetches history). */
retry: () => Promise<void>;
/** 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<void> => {
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);
@@ -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
+118 -6
View File
@@ -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<M extends OmniMessage>(msg: M, ts: string): M {
@@ -39,13 +41,13 @@ interface Harness {
errors: Array<string | null>;
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<OmniMessage[]>((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();
+104
View File
@@ -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();