返回 CodeWhale
runtime_web_client.test.mjs
根目录 / crates / tui / tests / runtime_web_client.test.mjs
1 import test from "node:test";
2 import assert from "node:assert/strict";
3 import { readFile } from "node:fs/promises";
4
5 import {
6 NO_TARGET,
7 STREAM_EVENT_NAMES,
8 answersForUserInput,
9 applyRuntimeEvent,
10 applySnapshot,
11 buildCreateThreadRequest,
12 claimInFlight,
13 createThreadState,
14 createWebSessionFetch,
15 createStreamConnector,
16 eventStreamUrl,
17 formatRuntimeProvenance,
18 imageInputPresentation,
19 isComposerSubmitKey,
20 groupThreadSummaries,
21 modeLabel,
22 modelOptionLabel,
23 newThreadDefaults,
24 pendingAttentionCount,
25 pendingAttentionLabel,
26 providerOptionLabel,
27 receiptPresentation,
28 recoverSnapshotAndSubscribe,
29 renderRuntimeProvenance,
30 resolveUserInputTarget,
31 restoreDraft,
32 runtimeEventContinuity,
33 saveDraft,
34 sessionTarget,
35 setSafeText,
36 snapshotThenSubscribe,
37 threadTarget,
38 threadProviderLabel,
39 } from "../src/runtime_web/app.mjs";
40
41 function snapshot(threadId = "thread-a", latestSeq = 7) {
42 return {
43 thread: { id: threadId, title: "Test", model: "test", mode: "agent" },
44 turns: [{ id: "turn-1", status: "in_progress" }],
45 items: [
46 {
47 id: "item-1",
48 turn_id: "turn-1",
49 kind: "agent_message",
50 status: "in_progress",
51 summary: "",
52 detail: "Hello",
53 },
54 ],
55 latest_seq: latestSeq,
56 };
57 }
58
59 function runtimeEvent(sequence, event, payload = {}, overrides = {}) {
60 return {
61 schema_version: 1,
62 seq: sequence,
63 event,
64 kind: event,
65 thread_id: "thread-a",
66 turn_id: "turn-1",
67 item_id: null,
68 payload,
69 ...overrides,
70 };
71 }
72
73 function cssDeclarations(styles, selectorPattern) {
74 const match = styles.match(new RegExp(`${selectorPattern}\\s*\\{([^}]*)\\}`));
75 assert.ok(match, `missing CSS rule matching ${selectorPattern}`);
76 return match[1];
77 }
78
79 test("embedded web client uses the Ocean Blue Stage semantic palette", async () => {
80 const [styles, html] = await Promise.all([
81 readFile(new URL("../src/runtime_web/styles.css", import.meta.url), "utf8"),
82 readFile(new URL("../src/runtime_web/index.html", import.meta.url), "utf8"),
83 ]);
84
85 for (const token of [
86 "--bg: #020711",
87 "--sidebar: #050b16",
88 "--surface: #0e1a30",
89 "--surface-raised: #172945",
90 "--stage-surface: #142747",
91 "--text: #f6f2e8",
92 "--action: #6aaef2",
93 "--status-human: #f6c453",
94 "--status-live: #4fd1c5",
95 "--status-warning: #ff7a59",
96 "--status-danger: #ff86b2",
97 "--ok: #9bd66f",
98 "--radius-control: 6px",
99 "--radius-card: 12px",
100 "--radius-composer: 16px",
101 "--rail: 256px",
102 ]) {
103 assert.match(styles, new RegExp(token.replace(/[.*+?^${}()|[\]\\]/g, "\\$&")));
104 }
105 assert.match(
106 cssDeclarations(styles, "\\.primary-button,\\s*\\.send-button"),
107 /background: var\(--action\)/,
108 );
109 assert.match(
110 cssDeclarations(styles, "\\.status-pip\\.running"),
111 /background: var\(--live\)/,
112 );
113 assert.match(
114 cssDeclarations(styles, "\\.message\\.user \\.message-body"),
115 /background: var\(--plate\)/,
116 );
117 assert.match(
118 cssDeclarations(styles, "\\.attention-card"),
119 /border: 1px solid rgba\(246, 196, 83/,
120 );
121 assert.match(
122 cssDeclarations(styles, "\\.status-banner"),
123 /color: var\(--warning\)/,
124 );
125 assert.match(
126 cssDeclarations(styles, "\\.connection-dot\\.ready"),
127 /background: var\(--ok\)/,
128 );
129 assert.match(html, /name="theme-color" content="#020711"/);
130 });
131
132 test("embedded web client keeps the CWC stage, transcript, and receipt hierarchy quiet", async () => {
133 const [styles, html, source] = await Promise.all([
134 readFile(new URL("../src/runtime_web/styles.css", import.meta.url), "utf8"),
135 readFile(new URL("../src/runtime_web/index.html", import.meta.url), "utf8"),
136 readFile(new URL("../src/runtime_web/app.mjs", import.meta.url), "utf8"),
137 ]);
138
139 assert.match(cssDeclarations(styles, "\\.session"), /background: var\(--stage-surface\)/);
140 assert.match(cssDeclarations(styles, "\\.transcript"), /background: var\(--stage-surface\)/);
141 assert.match(cssDeclarations(styles, "\\.receipt"), /display:\s*flex/);
142 assert.match(cssDeclarations(styles, "\\.receipt-dot"), /background: var\(--live\)/);
143 assert.match(
144 cssDeclarations(styles, "\\.message\\.user \\.message-label"),
145 /display:\s*none/,
146 );
147 assert.match(html, /id="transcript" role="log"[^>]+aria-relevant="additions"/);
148 assert.doesNotMatch(html, /id="transcript" role="log"[^>]+aria-relevant="[^"]*text/);
149 assert.match(source, /card\.append\(element\("span", "receipt-dot"\)\)/);
150 });
151
152 test("thread rail groups typed pending requests without disturbing server order", async () => {
153 const summaries = [
154 { id: "newest-recent", pending_attention_count: 0 },
155 { id: "newest-needs-you", pending_attention_count: 2 },
156 { id: "older-needs-you", pending_attention_count: 1 },
157 { id: "older-recent", pending_attention_count: 0 },
158 ];
159 const groups = groupThreadSummaries(summaries);
160
161 assert.deepEqual(groups.needsYou.map(({ id }) => id), ["newest-needs-you", "older-needs-you"]);
162 assert.deepEqual(groups.recent.map(({ id }) => id), ["newest-recent", "older-recent"]);
163 assert.deepEqual(summaries.map(({ id }) => id), [
164 "newest-recent",
165 "newest-needs-you",
166 "older-needs-you",
167 "older-recent",
168 ]);
169 assert.equal(pendingAttentionCount({ pending_attention_count: -1 }), 0);
170 assert.equal(pendingAttentionLabel({ pending_attention_count: 1 }), "1 item needs your attention");
171 assert.equal(pendingAttentionLabel({ pending_attention_count: 2 }), "2 items need your attention");
172
173 const [html, source] = await Promise.all([
174 readFile(new URL("../src/runtime_web/index.html", import.meta.url), "utf8"),
175 readFile(new URL("../src/runtime_web/app.mjs", import.meta.url), "utf8"),
176 ]);
177 assert.match(html, /id="thread-list" aria-label="Live threads"/);
178 assert.match(source, /group\.setAttribute\("aria-labelledby", headingId\)/);
179 assert.match(source, /attention\.setAttribute\("aria-label", pendingAttentionLabel\(summary\)\)/);
180 });
181
182 test("mobile drawer owns focus and background interaction while it is open", async () => {
183 const [html, source] = await Promise.all([
184 readFile(new URL("../src/runtime_web/index.html", import.meta.url), "utf8"),
185 readFile(new URL("../src/runtime_web/app.mjs", import.meta.url), "utf8"),
186 ]);
187
188 assert.match(html, /id="rail-open"[^>]+aria-controls="thread-rail"[^>]+aria-expanded="false"/);
189 assert.match(html, /id="rail-scrim"[^>]+tabindex="-1" hidden/);
190 assert.match(source, /function openRail\(\)[\s\S]*dom\.railClose\.focus/);
191 assert.match(source, /dom\.session\.setAttribute\("aria-hidden", "true"\)[\s\S]*setInert\(dom\.session, true\)/);
192 assert.match(source, /function closeRail[\s\S]*returnTarget\.focus[\s\S]*applyClosedMobileRailAccessibility/);
193 assert.match(source, /function trapFocusWithin[\s\S]*event\.key !== "Tab"[\s\S]*first\.focus/);
194 assert.match(source, /function trapRailFocus[\s\S]*trapFocusWithin\(event, dom\.rail\)/);
195 assert.match(source, /document\.addEventListener\("keydown", \(event\) => \{\s+if \(dom\.newThreadDialog\.open\) return;\s+if \(trapRailFocus\(event\)\) return;/);
196 assert.match(source, /event\.key === "Escape"[\s\S]*closeRail\(\)/);
197 });
198
199 test("mobile viewport and truth controls survive the software keyboard and coarse input", async () => {
200 const [styles, source] = await Promise.all([
201 readFile(new URL("../src/runtime_web/styles.css", import.meta.url), "utf8"),
202 readFile(new URL("../src/runtime_web/app.mjs", import.meta.url), "utf8"),
203 ]);
204
205 assert.match(styles, /height: var\(--visual-viewport-height\)/);
206 assert.match(source, /globalThis\.visualViewport\?\.addEventListener\("resize", syncVisualViewport\)/);
207 assert.match(styles, /@media \(pointer: coarse\)[\s\S]*min-height: 44px/);
208 assert.match(
209 styles,
210 /@media \(max-width: 800px\)[\s\S]*\.composer textarea,[\s\S]*font-size: 16px/,
211 );
212 assert.match(styles, /@media \(max-width: 430px\)[\s\S]*\.session-facts \{[\s\S]*display: flex/);
213 assert.match(styles, /\.session-facts \.fact-chip\[data-fact="workspace"\][\s\S]*display: none/);
214 assert.match(source, /chip\.dataset\.fact = String\(label \|\| ""\)\.toLowerCase\(\)/);
215 });
216
217 test("stream reconciliation preserves live controls, disclosures, and selected transcript text", async () => {
218 const [html, source] = await Promise.all([
219 readFile(new URL("../src/runtime_web/index.html", import.meta.url), "utf8"),
220 readFile(new URL("../src/runtime_web/app.mjs", import.meta.url), "utf8"),
221 ]);
222
223 assert.match(html, /id="attention" role="region"[^>]+aria-live="assertive"[^>]+aria-relevant="additions"/);
224 assert.equal(source.includes("dom.transcript.replaceChildren"), false);
225 assert.equal(source.includes("dom.attention.replaceChildren"), false);
226 assert.match(source, /function reconcileChildren\(/);
227 assert.match(source, /captureTranscriptSelection\(\)[\s\S]*restoreTranscriptSelection\(selection\)/);
228 assert.match(source, /card\.dataset\.attentionKey = key/);
229 assert.equal(source.includes("focusPendingAttention"), false);
230 assert.doesNotMatch(source, /card\.tabIndex = -1/);
231 });
232
233 test("attention requests stay single-flight without stealing the active control", async () => {
234 const requests = new Set();
235 assert.equal(claimInFlight(requests, "approval:one"), true);
236 assert.equal(claimInFlight(requests, "approval:one"), false);
237 requests.delete("approval:one");
238 assert.equal(claimInFlight(requests, "approval:one"), true);
239
240 const source = await readFile(
241 new URL("../src/runtime_web/app.mjs", import.meta.url),
242 "utf8",
243 );
244 assert.match(source, /if \(!claimInFlight\(app\.inFlightActions, action\)\) return;/);
245 assert.match(source, /setAttentionCardBusy\(action, true\)/);
246 assert.match(source, /finally \{\s+app\.inFlightActions\.delete\(action\);\s+setAttentionCardBusy\(action, false\);/);
247 });
248
249 test("degraded workflow receipts surface rejected dispatches as attention", () => {
250 const detail = JSON.stringify({
251 status: "degraded",
252 dispatch_failure_count: 2,
253 dispatch_failures: [{ label: "review", message: "profile unavailable" }],
254 });
255 assert.deepEqual(receiptPresentation({
256 kind: "tool_call",
257 status: "completed",
258 summary: "workflow: degraded",
259 detail,
260 metadata: { status: "degraded", dispatch_failure_count: 2 },
261 }), {
262 label: "Workflow · Needs attention",
263 summary: "2 task dispatches were rejected",
264 raw: `workflow: degraded\n\n${detail}`,
265 failed: true,
266 });
267 });
268
269 test("rail New thread cannot paint over the session fact chips", async () => {
270 const styles = await readFile(
271 new URL("../src/runtime_web/styles.css", import.meta.url),
272 "utf8",
273 );
274 assert.match(cssDeclarations(styles, "\\.rail"), /overflow:\s*hidden/);
275 assert.match(cssDeclarations(styles, "\\.new-thread"), /max-width:\s*100%/);
276 assert.match(cssDeclarations(styles, "\\.session-header"), /overflow:\s*hidden/);
277 assert.match(cssDeclarations(styles, "\\.session-facts"), /flex-wrap:\s*nowrap/);
278 });
279
280 test("production shell keeps readable type, controls, focus, and motion contracts", async () => {
281 const [styles, html] = await Promise.all([
282 readFile(new URL("../src/runtime_web/styles.css", import.meta.url), "utf8"),
283 readFile(new URL("../src/runtime_web/index.html", import.meta.url), "utf8"),
284 ]);
285
286 assert.match(cssDeclarations(styles, "\\.thread-row"), /min-height:\s*62px/);
287 assert.match(cssDeclarations(styles, "\\.thread-title"), /font-size:\s*14px/);
288 assert.match(cssDeclarations(styles, "\\.message-body"), /font-size:\s*15\.5px/);
289 assert.match(cssDeclarations(styles, "\\.composer textarea"), /font-size:\s*15\.5px/);
290 assert.match(cssDeclarations(styles, "\\.composer textarea"), /max-height:\s*220px/);
291 assert.match(
292 styles,
293 /\.primary-button,[\s\S]*?\.icon-button\s*\{[\s\S]*?min-height:\s*36px/,
294 );
295 assert.match(
296 styles,
297 /button:focus-visible,[\s\S]*?outline:\s*3px solid var\(--action\)/,
298 );
299 assert.match(
300 styles,
301 /@media \(prefers-reduced-motion: reduce\)[\s\S]*scroll-behavior: auto !important/,
302 );
303 assert.match(html, /Enter send · Shift\+Enter newline/);
304 });
305
306 test("composer Enter sends without interrupting newlines or IME composition", () => {
307 assert.equal(isComposerSubmitKey({ key: "Enter" }), true);
308 assert.equal(isComposerSubmitKey({ key: "Enter", metaKey: true }), true);
309 assert.equal(isComposerSubmitKey({ key: "Enter", ctrlKey: true }), true);
310 assert.equal(isComposerSubmitKey({ key: "Enter", shiftKey: true }), false);
311 assert.equal(isComposerSubmitKey({ key: "Enter", isComposing: true }), false);
312 assert.equal(isComposerSubmitKey({ key: "a" }), false);
313 });
314
315 test("new thread selection keeps provider and model scoped to the create request", () => {
316 const catalog = {
317 current: "deepseek",
318 providers: [
319 { id: "openai", default_model: "gpt-5.6" },
320 { id: "deepseek", default_model: "deepseek-v4-pro" },
321 ],
322 };
323 assert.deepEqual(newThreadDefaults(catalog), {
324 providerId: "deepseek",
325 modelProviderId: "",
326 model: "deepseek-v4-pro",
327 });
328 assert.deepEqual(
329 buildCreateThreadRequest(" deepseek ", " deepseek-v4-flash-vision-exp "),
330 {
331 model_provider: "deepseek",
332 model: "deepseek-v4-flash-vision-exp",
333 },
334 );
335 assert.throws(() => buildCreateThreadRequest("deepseek", ""), /provider and a model/);
336
337 const namedCustom = {
338 current: "custom",
339 providers: [{
340 id: "custom",
341 model_provider_id: "lm-studio",
342 default_model: "local-vision-model",
343 }],
344 };
345 assert.deepEqual(newThreadDefaults(namedCustom), {
346 providerId: "custom",
347 modelProviderId: "lm-studio",
348 model: "local-vision-model",
349 });
350 assert.deepEqual(
351 buildCreateThreadRequest("custom", "local-vision-model", "lm-studio"),
352 {
353 model_provider: "custom",
354 model: "local-vision-model",
355 model_provider_id: "lm-studio",
356 },
357 );
358 assert.equal(
359 providerOptionLabel({
360 id: "custom",
361 display_name: "Custom",
362 model_provider_id: "lm-studio",
363 }),
364 "Custom · lm-studio",
365 );
366 assert.equal(threadProviderLabel({
367 model_provider: "custom",
368 model_provider_id: "lm-studio",
369 }), "lm-studio");
370 assert.equal(threadProviderLabel({ model_provider: "deepseek" }), "deepseek");
371 });
372
373 test("thread facts keep exact named provider identity visible after creation", async () => {
374 const source = await readFile(
375 new URL("../src/runtime_web/app.mjs", import.meta.url),
376 "utf8",
377 );
378 assert.match(source, /const provider = threadProviderLabel\(thread\)/);
379 assert.match(source, /factChip\("Provider", provider\)/);
380 });
381
382 test("composer send guard survives rerenders until the request settles", async () => {
383 const source = await readFile(
384 new URL("../src/runtime_web/app.mjs", import.meta.url),
385 "utf8",
386 );
387 assert.match(source, /const sending = app\.inFlightActions\.has\(composerSendAction\)/);
388 assert.match(source, /dom\.composerInput\.disabled = sending \|\| !ready/);
389 assert.match(source, /dom\.send\.disabled = sending \|\| !ready/);
390 assert.match(source, /if \(!claimInFlight\(app\.inFlightActions, composerSendAction\)\) return;/);
391 assert.match(source, /finally \{\s+app\.inFlightActions\.delete\(composerSendAction\);\s+renderComposer\(\);/);
392 });
393
394 test("new thread dialog labels exact vision capability without exposing attachments", async () => {
395 const [html, source] = await Promise.all([
396 readFile(new URL("../src/runtime_web/index.html", import.meta.url), "utf8"),
397 readFile(new URL("../src/runtime_web/app.mjs", import.meta.url), "utf8"),
398 ]);
399 const vision = { id: "deepseek-v4-flash-vision-exp", image_input: "supported" };
400 assert.equal(modelOptionLabel(vision), "deepseek-v4-flash-vision-exp · Vision");
401 assert.equal(imageInputPresentation("supported").label, "Vision");
402 assert.equal(imageInputPresentation("unsupported").label, "Text only");
403 assert.equal(imageInputPresentation("unknown").state, "unknown");
404
405 assert.match(html, /id="new-thread-dialog"[^>]+tabindex="-1"[^>]+aria-labelledby="new-thread-title"/);
406 assert.match(html, /id="new-thread-cancel"[^>]+autofocus/);
407 assert.match(html, /id="new-thread-provider" required disabled/);
408 assert.match(html, /id="new-thread-model" required disabled/);
409 assert.match(html, /does not change your Runtime defaults/);
410 assert.doesNotMatch(html, /type="file"/);
411 assert.match(source, /api\("\/v1\/providers"\)/);
412 // The dialog loads the catalog through the bounded, paginated collector
413 // keyed by provider.id; the wire endpoint stays /v1/providers/<id>/models.
414 assert.match(source, /collectProviderModelPages\(provider\.id/);
415 assert.match(source, /\/v1\/providers\/\$\{encodeURIComponent\(provider\)\}\/models\?\$\{query\.toString\(\)\}/);
416 assert.match(source, /body: JSON\.stringify\(request\)/);
417 assert.match(source, /function trapFocusWithin\(event, container\)/);
418 assert.match(source, /dom\.newThreadCancel\.focus\(\{ preventScroll: true \}\)/);
419 assert.match(source, /trapFocusWithin\(event, dom\.newThreadDialog\)/);
420 assert.doesNotMatch(source, /\/v1\/providers\/[^`"']+\/switch/);
421 });
422
423 test("uses the v0.9.6 Work vocabulary for the agent wire mode", () => {
424 assert.equal(modeLabel("agent"), "Work");
425 assert.equal(modeLabel("plan"), "Plan");
426 assert.equal(modeLabel("operate"), "Operate");
427 });
428
429 test("formats and renders exact Runtime build provenance with honest fallbacks", () => {
430 const exactCommit = "abcdef0123456789abcdef0123456789abcdef01";
431 const stamped = {
432 codewhale_version: "0.9.6",
433 codewhale_commit: exactCommit,
434 };
435 assert.equal(formatRuntimeProvenance(stamped), "0.9.6 · abcdef012345");
436
437 const rendered = { textContent: "" };
438 renderRuntimeProvenance(rendered, stamped);
439 assert.equal(rendered.textContent, "0.9.6 · abcdef012345");
440
441 assert.equal(
442 formatRuntimeProvenance({ codewhale_version: "0.9.6", codewhale_commit: "unknown" }),
443 "0.9.6 · source unknown",
444 );
445 assert.equal(
446 formatRuntimeProvenance({ version: "0.9.6", codewhale_commit: "too-short" }),
447 "0.9.6 · source unknown",
448 );
449 assert.equal(formatRuntimeProvenance(null), "version unknown · source unknown");
450 });
451
452 test("loads a consistent snapshot before subscribing from latest_seq", async () => {
453 const state = createThreadState("thread-a");
454 const order = [];
455 const subscribed = await snapshotThenSubscribe({
456 state,
457 threadId: "thread-a",
458 loadSnapshot: async () => {
459 order.push("snapshot");
460 return snapshot("thread-a", 42);
461 },
462 subscribe: (threadId, sequence) => order.push(`subscribe:${threadId}:${sequence}`),
463 });
464
465 assert.equal(subscribed, true);
466 assert.deepEqual(order, ["snapshot", "subscribe:thread-a:42"]);
467 assert.equal(state.latestSeq, 42);
468 });
469
470 test("snapshot recovery waits for the replacement stream to open", async () => {
471 const state = createThreadState("thread-a");
472 let finishOpening;
473 let settled = false;
474 const opening = new Promise((resolve) => {
475 finishOpening = resolve;
476 });
477 const recovery = snapshotThenSubscribe({
478 state,
479 threadId: "thread-a",
480 loadSnapshot: async () => snapshot("thread-a", 43),
481 subscribe: () => opening,
482 }).then((result) => {
483 settled = true;
484 return result;
485 });
486
487 await Promise.resolve();
488 assert.equal(settled, false, "snapshot success alone must not finish recovery");
489 finishOpening();
490 assert.equal(await recovery, true);
491 assert.equal(settled, true);
492 });
493
494 test("a failed replacement stream keeps the gap until a later stream opens", async () => {
495 const state = createThreadState("thread-a");
496 let gap = true;
497 let attempts = 0;
498 const recover = () => recoverSnapshotAndSubscribe({
499 state,
500 threadId: "thread-a",
501 loadSnapshot: async () => snapshot("thread-a", 44 + attempts),
502 subscribe: async () => {
503 attempts += 1;
504 if (attempts === 1) throw new Error("replacement stream did not reopen");
505 },
506 }, () => {
507 gap = false;
508 });
509
510 await assert.rejects(recover(), /did not reopen/);
511 assert.equal(gap, true, "snapshot success must not hide a failed stream handshake");
512 assert.equal(await recover(), true);
513 assert.equal(gap, false, "a later snapshot plus open stream clears the gap");
514 });
515
516 test("drops a stale snapshot selection without opening an event stream", async () => {
517 const state = createThreadState("thread-a");
518 let current = true;
519 let subscribed = false;
520 const result = await snapshotThenSubscribe({
521 state,
522 threadId: "thread-a",
523 loadSnapshot: async () => {
524 current = false;
525 return snapshot();
526 },
527 subscribe: () => {
528 subscribed = true;
529 },
530 isCurrent: () => current,
531 });
532 assert.equal(result, false);
533 assert.equal(subscribed, false);
534 });
535
536 test("reconnect cursor advances monotonically and duplicate or stale-thread events are ignored", () => {
537 const state = createThreadState("thread-a");
538 assert.equal(applySnapshot(state, snapshot("thread-a", 7)), true);
539
540 assert.equal(
541 applyRuntimeEvent(
542 state,
543 runtimeEvent(8, "item.delta", { delta: " world", kind: "agent_message" }, { item_id: "item-1" }),
544 ),
545 true,
546 );
547 assert.equal(
548 applyRuntimeEvent(
549 state,
550 runtimeEvent(8, "item.delta", { delta: " duplicate", kind: "agent_message" }, { item_id: "item-1" }),
551 ),
552 false,
553 );
554 assert.equal(
555 applyRuntimeEvent(state, runtimeEvent(99, "turn.completed", {}, { thread_id: "thread-b" })),
556 false,
557 );
558 assert.equal(state.items.get("item-1").detail, "Hello world");
559 assert.equal(state.latestSeq, 8);
560 assert.equal(eventStreamUrl("thread-a", state.latestSeq), "/v1/threads/thread-a/events?since_seq=8");
561 });
562
563 test("uses the stream predecessor cursor to detect real gaps without assuming global sequences are contiguous", () => {
564 const state = createThreadState("thread-a");
565 applySnapshot(state, snapshot("thread-a", 7));
566
567 const interleaved = runtimeEvent(
568 12,
569 "item.delta",
570 { delta: " after other threads", kind: "agent_message" },
571 { item_id: "item-1", previous_seq: 7 },
572 );
573 assert.equal(runtimeEventContinuity(state, interleaved), "next");
574 assert.equal(applyRuntimeEvent(state, interleaved), true);
575 assert.equal(state.latestSeq, 12);
576
577 const gap = runtimeEvent(
578 15,
579 "approval.required",
580 { approval_id: "approval-missed", tool_name: "exec_shell" },
581 { previous_seq: 14 },
582 );
583 assert.equal(runtimeEventContinuity(state, gap), "gap");
584 assert.equal(applyRuntimeEvent(state, gap), false);
585 assert.equal(state.latestSeq, 12);
586 assert.equal(state.approvals.has("approval-missed"), false);
587 });
588
589 test("registers the full emitted Runtime vocabulary and advances continuity for every event", async () => {
590 const runtimeSource = await readFile(
591 new URL("../src/runtime_threads.rs", import.meta.url),
592 "utf8",
593 );
594 const emittedNames = new Set(
595 [...runtimeSource.matchAll(
596 /"((?:thread|turn|item|approval|user_input|sandbox|agent|tool_call)\.[a-z_]+)"/g,
597 )].map((match) => match[1]),
598 );
599 assert.deepEqual(new Set(STREAM_EVENT_NAMES), emittedNames);
600 assert.equal(STREAM_EVENT_NAMES.includes("thread.created"), false);
601
602 const state = createThreadState("thread-a");
603 applySnapshot(state, snapshot("thread-a", 7));
604 let previousSeq = 7;
605 for (const eventName of STREAM_EVENT_NAMES) {
606 const sequence = previousSeq + 2;
607 const turnBefore = state.turns.get("turn-1");
608 const payload = eventName === "turn.usage"
609 ? { usage: { input_tokens: 100, output_tokens: 20 } }
610 : {};
611 const envelope = runtimeEvent(sequence, eventName, payload, { previous_seq: previousSeq });
612 assert.equal(runtimeEventContinuity(state, envelope), "next", eventName);
613 assert.equal(applyRuntimeEvent(state, envelope), true, eventName);
614 assert.equal(state.latestSeq, sequence, eventName);
615 if (eventName === "turn.usage") {
616 // Request diagnostics advance continuity; only the settled turn owns totals.
617 assert.equal(state.turns.get("turn-1"), turnBefore);
618 assert.equal(applyRuntimeEvent(state, envelope), false);
619 }
620 previousSeq = sequence;
621 }
622 const settledTurn = {
623 id: "turn-1",
624 status: "completed",
625 usage: { input_tokens: 300, output_tokens: 60 },
626 };
627 assert.equal(applyRuntimeEvent(state, runtimeEvent(
628 previousSeq + 1,
629 "turn.completed",
630 { turn: settledTurn },
631 { previous_seq: previousSeq },
632 )), true);
633 assert.deepEqual(state.turns.get("turn-1"), settledTurn);
634 });
635
636 test("gap recovery snapshot restores approval and user-input attention before resubscribing", async () => {
637 const state = createThreadState("thread-a");
638 applySnapshot(state, snapshot("thread-a", 7));
639 const subscriptions = [];
640
641 const recovered = await snapshotThenSubscribe({
642 state,
643 threadId: "thread-a",
644 loadSnapshot: async () => ({
645 ...snapshot("thread-a", 15),
646 pending_approvals: [{
647 id: "approval-recovered",
648 turn_id: "turn-1",
649 tool_name: "exec_command",
650 description: "Run a local check",
651 }],
652 pending_user_inputs: [{
653 id: "input-recovered",
654 turn_id: "turn-1",
655 request: { questions: [{ id: "choice", question: "Continue?", options: [] }] },
656 }],
657 pending_dynamic_tool_calls: [{
658 thread_id: "thread-a",
659 turn_id: "turn-1",
660 call_id: "call-recovered",
661 namespace: "bench",
662 tool: "lookup",
663 arguments: { id: "7" },
664 }],
665 }),
666 subscribe: (threadId, sequence) => subscriptions.push([threadId, sequence]),
667 });
668
669 assert.equal(recovered, true);
670 assert.equal(state.approvals.size, 1);
671 assert.equal(state.approvals.has("approval-recovered"), true);
672 assert.equal(state.userInputs.size, 1);
673 assert.equal(state.userInputs.has("input-recovered"), true);
674 assert.equal(state.dynamicToolCalls.size, 1);
675 assert.equal(state.dynamicToolCalls.get("call-recovered").tool, "lookup");
676 assert.deepEqual(subscriptions, [["thread-a", 15]]);
677
678 const duplicate = runtimeEvent(
679 15,
680 "approval.required",
681 { approval_id: "approval-recovered", tool_name: "exec_command" },
682 { previous_seq: 14 },
683 );
684 assert.equal(applyRuntimeEvent(state, duplicate), false);
685 assert.equal(state.approvals.size, 1);
686 });
687
688 test("browser clears its surfaced gap only after a replacement snapshot subscribes", async () => {
689 const source = await readFile(new URL("../src/runtime_web/app.mjs", import.meta.url), "utf8");
690 assert.match(
691 source,
692 /async function recoverProjection[\s\S]*?connectStream\(id, sequence, generation, true\)/,
693 );
694 assert.match(
695 source,
696 /async function recoverProjection[\s\S]*?recoverSnapshotAndSubscribe\([\s\S]*?app\.streamGap = false;[\s\S]*?if \(!subscribed\) return;\s+renderAll\(\);/,
697 );
698 });
699
700 test("user-input answers stay bound to the selected live thread and pending request", () => {
701 const state = createThreadState("thread-a");
702 state.userInputs.set("input-1", {});
703
704 assert.deepEqual(
705 resolveUserInputTarget("input-1", threadTarget("thread-a"), state),
706 { ok: true, threadId: "thread-a", inputId: "input-1" },
707 );
708 assert.deepEqual(
709 resolveUserInputTarget("input-1", sessionTarget("session-a"), state),
710 { ok: false, reason: "session-not-live" },
711 );
712 assert.deepEqual(
713 resolveUserInputTarget("input-1", NO_TARGET, state),
714 { ok: false, reason: "no-target" },
715 );
716 assert.deepEqual(
717 resolveUserInputTarget("input-2", threadTarget("thread-a"), state),
718 { ok: false, reason: "stale-user-input" },
719 );
720 assert.deepEqual(
721 resolveUserInputTarget("input-1", threadTarget("thread-b"), state),
722 { ok: false, reason: "stale-target" },
723 );
724 });
725
726 test("user-input payloads preserve TUI custom-answer parity and single-select cardinality", () => {
727 const single = {
728 questions: [{
729 id: "path",
730 header: "Path",
731 question: "Which path?",
732 options: [{ label: "A" }, { label: "B" }],
733 allow_free_text: false,
734 multi_select: false,
735 }],
736 };
737 assert.deepEqual(answersForUserInput(single, {}, { path: "A different path" }), {
738 ok: true,
739 answers: [{ id: "path", label: "Other", value: "A different path" }],
740 });
741 assert.equal(
742 answersForUserInput(single, { path: ["A"] }, { path: "also B" }).reason,
743 "multiple-answers",
744 );
745 assert.equal(
746 answersForUserInput(single, { path: ["forged"] }, {}).reason,
747 "invalid-option",
748 );
749 assert.equal(answersForUserInput(single, {}, {}).reason, "missing-answer");
750
751 const multi = {
752 questions: [{
753 ...single.questions[0],
754 id: "checks",
755 multi_select: true,
756 }],
757 };
758 assert.deepEqual(
759 answersForUserInput(multi, { checks: ["A", "B"] }, { checks: "C" }),
760 {
761 ok: true,
762 answers: [
763 { id: "checks", label: "A", value: "A" },
764 { id: "checks", label: "B", value: "B" },
765 { id: "checks", label: "Other", value: "C" },
766 ],
767 },
768 );
769 });
770
771 test("assembles deltas and replaces the live item with its settled receipt", () => {
772 const state = createThreadState("thread-a");
773 applySnapshot(state, { ...snapshot(), items: [], latest_seq: 1 });
774 applyRuntimeEvent(
775 state,
776 runtimeEvent(2, "item.delta", { delta: "one", kind: "agent_message" }, { item_id: "item-new" }),
777 );
778 applyRuntimeEvent(
779 state,
780 runtimeEvent(3, "item.delta", { delta: " two", kind: "agent_message" }, { item_id: "item-new" }),
781 );
782 assert.equal(state.items.get("item-new").detail, "one two");
783
784 applyRuntimeEvent(
785 state,
786 runtimeEvent(
787 4,
788 "item.completed",
789 {
790 item: {
791 id: "item-new",
792 turn_id: "turn-1",
793 kind: "agent_message",
794 status: "completed",
795 summary: "one two",
796 detail: "one two",
797 },
798 },
799 { item_id: "item-new" },
800 ),
801 );
802 assert.equal(state.items.get("item-new").status, "completed");
803 assert.deepEqual(state.itemOrder, ["item-new"]);
804
805 applyRuntimeEvent(
806 state,
807 runtimeEvent(5, "item.delta", { delta: "partial", kind: "tool_call" }, { item_id: "item-stop" }),
808 );
809 applyRuntimeEvent(
810 state,
811 runtimeEvent(
812 6,
813 "item.interrupted",
814 {
815 item: {
816 id: "item-stop",
817 turn_id: "turn-1",
818 kind: "tool_call",
819 status: "interrupted",
820 summary: "Interrupted",
821 detail: "partial",
822 },
823 },
824 { item_id: "item-stop" },
825 ),
826 );
827 assert.equal(state.items.get("item-stop").status, "interrupted");
828
829 applyRuntimeEvent(
830 state,
831 runtimeEvent(
832 7,
833 "item.canceled",
834 {
835 item: {
836 id: "item-compact",
837 turn_id: "turn-1",
838 kind: "compaction",
839 status: "canceled",
840 summary: "Compaction canceled",
841 detail: "Compaction canceled",
842 },
843 },
844 { item_id: "item-compact" },
845 ),
846 );
847 assert.equal(state.items.get("item-compact").status, "canceled");
848 assert.deepEqual(state.itemOrder, ["item-new", "item-stop", "item-compact"]);
849 });
850
851 test("projects agent lifecycle receipts live and settles them without a snapshot reload", () => {
852 const state = createThreadState("thread-a");
853 applySnapshot(state, { ...snapshot(), items: [], latest_seq: 1 });
854
855 const agentItem = (status, summary) => ({
856 id: "item-agent",
857 turn_id: "turn-1",
858 kind: "status",
859 status,
860 summary,
861 detail: summary,
862 });
863 applyRuntimeEvent(
864 state,
865 runtimeEvent(2, "agent.spawned", { item: agentItem("in_progress", "Agent spawned") }),
866 );
867 assert.equal(state.items.get("item-agent").status, "in_progress");
868 assert.deepEqual(state.itemOrder, ["item-agent"]);
869
870 applyRuntimeEvent(
871 state,
872 runtimeEvent(3, "agent.progress", { item: agentItem("in_progress", "Agent checking") }),
873 );
874 applyRuntimeEvent(
875 state,
876 runtimeEvent(4, "agent.completed", { item: agentItem("completed", "Agent completed") }),
877 );
878 assert.equal(state.items.get("item-agent").status, "completed");
879 assert.equal(state.items.get("item-agent").summary, "Agent completed");
880 assert.deepEqual(state.itemOrder, ["item-agent"]);
881
882 applyRuntimeEvent(
883 state,
884 runtimeEvent(5, "agent.list", {
885 item: {
886 id: "item-agent-list",
887 turn_id: "turn-1",
888 kind: "status",
889 status: "completed",
890 summary: "Agent list refreshed",
891 detail: "Agent list refreshed",
892 },
893 }),
894 );
895 assert.equal(state.items.get("item-agent-list").status, "completed");
896 assert.deepEqual(state.itemOrder, ["item-agent", "item-agent-list"]);
897 assert.equal(state.latestSeq, 5);
898 });
899
900 test("tracks approval and user-input attention until each is resolved", () => {
901 const state = createThreadState("thread-a");
902 applySnapshot(state, snapshot());
903 applyRuntimeEvent(
904 state,
905 runtimeEvent(8, "approval.required", { approval_id: "approval-1", tool_name: "exec_shell" }),
906 );
907 applyRuntimeEvent(
908 state,
909 runtimeEvent(9, "user_input.required", {
910 id: "input-1",
911 request: { questions: [{ id: "choice", question: "Choose?", options: [] }] },
912 }),
913 );
914 assert.equal(state.approvals.has("approval-1"), true);
915 assert.equal(state.userInputs.has("input-1"), true);
916
917 applyRuntimeEvent(
918 state,
919 runtimeEvent(10, "approval.decided", { approval_id: "approval-1", decision: "allow" }),
920 );
921 assert.equal(state.approvals.has("approval-1"), false);
922 assert.equal(state.userInputs.has("input-1"), true);
923
924 applyRuntimeEvent(
925 state,
926 runtimeEvent(11, "user_input.answered", { input_id: "input-1" }),
927 );
928 assert.equal(state.userInputs.has("input-1"), false);
929 });
930
931 test("hydrates pending attention from a reload snapshot and clears cancellation events", () => {
932 const state = createThreadState("thread-a");
933 const detail = {
934 ...snapshot(),
935 pending_approvals: [{
936 id: "approval-reload",
937 turn_id: "turn-1",
938 tool_name: "exec_command",
939 description: "Run a local check",
940 }],
941 pending_user_inputs: [{
942 id: "input-reload",
943 turn_id: "turn-1",
944 request: { questions: [{ id: "choice", question: "Continue?", options: [] }] },
945 }],
946 };
947
948 assert.equal(applySnapshot(state, detail), true);
949 assert.equal(state.approvals.get("approval-reload").tool_name, "exec_command");
950 assert.equal(state.userInputs.get("input-reload").turn_id, "turn-1");
951
952 applyRuntimeEvent(
953 state,
954 runtimeEvent(8, "user_input.canceled", { id: "input-reload", terminal: true }),
955 );
956 assert.equal(state.userInputs.has("input-reload"), false);
957 });
958
959 test("turn completion defensively clears attention owned by that turn", () => {
960 const state = createThreadState("thread-a");
961 assert.equal(applySnapshot(state, {
962 ...snapshot(),
963 pending_approvals: [{ id: "approval-terminal", turn_id: "turn-1" }],
964 pending_user_inputs: [{ id: "input-terminal", turn_id: "turn-1", request: { questions: [] } }],
965 pending_dynamic_tool_calls: [{ call_id: "call-terminal", turn_id: "turn-1", tool: "lookup" }],
966 }), true);
967 state.approvals.set("approval-other", { id: "approval-other", turn_id: "turn-other" });
968 state.userInputs.set("input-other", { id: "input-other", turn_id: "turn-other" });
969 state.dynamicToolCalls.set("call-other", { call_id: "call-other", turn_id: "turn-other" });
970
971 assert.equal(applyRuntimeEvent(
972 state,
973 runtimeEvent(8, "turn.completed", { turn: { id: "turn-1", status: "completed" } }),
974 ), true);
975 assert.equal(state.approvals.has("approval-terminal"), false);
976 assert.equal(state.userInputs.has("input-terminal"), false);
977 assert.equal(state.dynamicToolCalls.has("call-terminal"), false);
978 assert.equal(state.approvals.has("approval-other"), true);
979 assert.equal(state.userInputs.has("input-other"), true);
980 assert.equal(state.dynamicToolCalls.has("call-other"), true);
981 });
982
983 test("dynamic tool calls hydrate and disappear exactly once across terminal variants", () => {
984 const state = createThreadState("thread-a");
985 assert.equal(applySnapshot(state, {
986 ...snapshot(),
987 pending_dynamic_tool_calls: [{
988 thread_id: "thread-a",
989 turn_id: "turn-1",
990 call_id: "call-snapshot",
991 tool: "snapshot_lookup",
992 arguments: { id: "snapshot" },
993 }],
994 }), true);
995 assert.equal(state.dynamicToolCalls.get("call-snapshot").tool, "snapshot_lookup");
996
997 assert.equal(applyRuntimeEvent(
998 state,
999 runtimeEvent(8, "tool_call.requested", {
1000 thread_id: "thread-a",
1001 turn_id: "turn-1",
1002 call_id: "call-live",
1003 tool: "live_lookup",
1004 arguments: { id: "live" },
1005 }),
1006 ), true);
1007 assert.equal(state.dynamicToolCalls.size, 2);
1008
1009 assert.equal(applyRuntimeEvent(
1010 state,
1011 runtimeEvent(9, "tool_call.resolved", { call_id: "call-snapshot", status: "resolved" }),
1012 ), true);
1013 assert.equal(state.dynamicToolCalls.has("call-snapshot"), false);
1014 assert.equal(applyRuntimeEvent(
1015 state,
1016 runtimeEvent(9, "tool_call.resolved", { call_id: "call-snapshot", status: "resolved" }),
1017 ), false);
1018 assert.equal(state.dynamicToolCalls.size, 1);
1019
1020 assert.equal(applyRuntimeEvent(
1021 state,
1022 runtimeEvent(10, "tool_call.canceled", { call_id: "call-live", status: "canceled" }),
1023 ), true);
1024 assert.equal(state.dynamicToolCalls.size, 0);
1025
1026 assert.equal(applyRuntimeEvent(
1027 state,
1028 runtimeEvent(11, "tool_call.requested", {
1029 call_id: "call-timeout",
1030 tool: "slow_lookup",
1031 arguments: {},
1032 }),
1033 ), true);
1034 assert.equal(state.dynamicToolCalls.has("call-timeout"), true);
1035 assert.equal(applyRuntimeEvent(
1036 state,
1037 runtimeEvent(12, "tool_call.timeout", { call_id: "call-timeout", status: "timeout" }),
1038 ), true);
1039 assert.equal(state.dynamicToolCalls.size, 0);
1040 });
1041
1042 test("preserves drafts per thread without browser storage", () => {
1043 const drafts = new Map();
1044 saveDraft(drafts, "thread-a", "draft A");
1045 saveDraft(drafts, "thread-b", "draft B");
1046 assert.equal(restoreDraft(drafts, "thread-a"), "draft A");
1047 assert.equal(restoreDraft(drafts, "thread-b"), "draft B");
1048 saveDraft(drafts, "thread-a", "");
1049 assert.equal(restoreDraft(drafts, "thread-a"), "");
1050 });
1051
1052 test("renders hostile Runtime text only through the textContent sink", async () => {
1053 const hostile = `<img src=x onerror=alert(1)><script>alert(2)</script>`;
1054 const fakeElement = { textContent: "" };
1055 setSafeText(fakeElement, hostile);
1056 assert.equal(fakeElement.textContent, hostile);
1057
1058 const source = await readFile(new URL("../src/runtime_web/app.mjs", import.meta.url), "utf8");
1059 assert.equal(source.includes("inner" + "HTML"), false);
1060 assert.equal(source.includes("insertAdjacent" + "HTML"), false);
1061 assert.equal(source.includes("local" + "Storage"), false);
1062 assert.equal(source.includes("codewhale_runtime_token"), false);
1063 });
1064
1065 test("web session fetch keeps request proof outside cookies and strips bootstrap fragments", async () => {
1066 const stored = new Map();
1067 const storage = () => ({ setItem: (key, value) => stored.set(key, value), getItem: (key) => stored.get(key) });
1068 const paths = [];
1069 const calls = [];
1070 const history = { replaceState: (_state, _title, path) => paths.push(path) };
1071 const fetch = async (path, options) => { calls.push({ path, options }); return { ok: true }; };
1072 const request = createWebSessionFetch({
1073 location: { hash: "#p=proof-one", pathname: "/", search: "?view=threads" }, history, storage, fetch,
1074 });
1075 assert.deepEqual(paths, ["/?view=threads"]);
1076 await request("/v1/threads", { method: "POST", body: "{}" });
1077 await request("/__codewhale/web/stream-ticket", { method: "POST" });
1078 const reloaded = createWebSessionFetch({ location: { hash: "" }, history, storage, fetch });
1079 await reloaded("/v1/threads");
1080 for (const { options } of calls) {
1081 assert.equal(options.headers.get("x-codewhale-web-request"), "proof-one");
1082 assert.equal(options.headers.has("cookie"), false);
1083 assert.equal(options.credentials, "same-origin");
1084 assert.equal(options.cache, "no-store");
1085 }
1086 assert.equal(calls[0].options.headers.get("content-type"), "application/json");
1087 assert.equal(eventStreamUrl("thread-a", 4, "ticket-one"), "/v1/threads/thread-a/events?since_seq=4&web_stream_ticket=ticket-one");
1088
1089 const memoryOnly = createWebSessionFetch({
1090 location: { hash: "#p=memory-proof", pathname: "/", search: "" }, history,
1091 storage: () => { throw new Error("storage unavailable"); }, fetch,
1092 });
1093 await memoryOnly("/v1/threads");
1094 assert.equal(calls.at(-1).options.headers.get("x-codewhale-web-request"), "memory-proof");
1095 });
1096
1097 test("web session recovers from the page in an independent tab or without storage", async () => {
1098 for (const storedProof of [null, "old-process-proof", "unavailable"]) {
1099 const stored = new Map([["codewhale_web_request_proof", storedProof]]);
1100 const calls = [];
1101 const request = createWebSessionFetch({
1102 location: { hash: "" },
1103 history: { replaceState: () => assert.fail("no fragment to clear") },
1104 storage: () => {
1105 if (storedProof === "unavailable") throw new Error("storage unavailable");
1106 return { getItem: (key) => stored.get(key), setItem: (key, value) => stored.set(key, value) };
1107 },
1108 pageProof: "current-page-proof",
1109 fetch: async (_path, options) => { calls.push(options); return { ok: true }; },
1110 });
1111 await request("/v1/threads");
1112 await request("/__codewhale/web/stream-ticket", { method: "POST" });
1113 assert.equal(calls.length, 2);
1114 for (const options of calls) {
1115 assert.equal(options.headers.get("x-codewhale-web-request"), "current-page-proof");
1116 }
1117 if (storedProof !== "unavailable") {
1118 assert.equal(stored.get("codewhale_web_request_proof"), "current-page-proof");
1119 }
1120 }
1121 });
1122
1123 function streamHarness(api) {
1124 const app = {
1125 generation: 1, selectedThreadId: "thread-a", threadState: { latestSeq: 4 },
1126 stream: null, streamOpenCancel: null, reconnectTimer: null,
1127 };
1128 const streams = [];
1129 const timers = new Map();
1130 const statuses = [];
1131 let nextTimer = 0;
1132 class EventSource {
1133 constructor(url) { this.url = url; this.closed = false; streams.push(this); }
1134 close() { this.closed = true; }
1135 addEventListener() {}
1136 }
1137 const connector = createStreamConnector({
1138 app, api, EventSource, receive: () => {}, setConnection: () => {},
1139 showStatus: (message) => statuses.push(message),
1140 setTimeout: (callback, delay) => { timers.set(++nextTimer, { callback, delay }); return nextTimer; },
1141 clearTimeout: (id) => timers.delete(id),
1142 });
1143 return {
1144 ...connector, app, streams, timers, statuses,
1145 tick: async () => {
1146 assert.equal(timers.size, 1, "exactly one reconnect should be pending");
1147 const [id, { callback, delay }] = timers.entries().next().value;
1148 timers.delete(id);
1149 callback();
1150 await new Promise((resolve) => setImmediate(resolve));
1151 return delay;
1152 },
1153 };
1154 }
1155
1156 test("stream retries ticket network and 5xx failures with capped backoff and resumes the latest cursor", async () => {
1157 let requests = 0;
1158 const h = streamHarness(async () => {
1159 requests += 1;
1160 if (requests === 2) throw new TypeError("network unavailable");
1161 if (requests > 2 && requests < 10) throw Object.assign(new Error("temporarily unavailable"), { status: 503 });
1162 return { stream_ticket: `ticket-${requests}` };
1163 });
1164 await h.connectStream("thread-a", 4, 1);
1165 h.streams[0].onopen();
1166 h.streams[0].onerror();
1167 assert.equal(h.streams[0].closed, true);
1168 h.app.threadState.latestSeq = 9;
1169 const delays = [];
1170 while (h.streams.length === 1) delays.push(await h.tick());
1171 assert.deepEqual(delays, [900, 1800, 3600, 7200, 14400, 28800, 30000, 30000, 30000]);
1172 assert.match(h.streams[1].url, /since_seq=9&web_stream_ticket=ticket-10$/);
1173 h.streams[1].onopen();
1174 h.streams[1].onerror();
1175 assert.equal(await h.tick(), 900, "opening a stream resets backoff");
1176 h.stopStream();
1177 assert.equal(h.timers.size, 0);
1178 });
1179
1180 test("stream stops ticket retries on 401 and 403", async () => {
1181 for (const status of [401, 403]) {
1182 let requests = 0;
1183 const h = streamHarness(async () => {
1184 if (++requests === 1) return { stream_ticket: "first" };
1185 throw Object.assign(new Error("session ended"), { status });
1186 });
1187 await h.connectStream("thread-a", 4, 1);
1188 h.app.stream.onerror();
1189 await h.tick();
1190 assert.equal(h.timers.size, 0);
1191 assert.equal(h.app.stream, null);
1192 assert.deepEqual(h.statuses, ["session ended"]);
1193 }
1194 });
1195
1196 test("stream ignores superseded ticket completions and failures, including stopped connections", async () => {
1197 const pending = [];
1198 const h = streamHarness(() => new Promise((resolve, reject) => pending.push({ resolve, reject })));
1199 const first = h.connectStream("thread-a", 4, 1);
1200 const second = h.connectStream("thread-a", 4, 1);
1201 pending[1].resolve({ stream_ticket: "current" });
1202 await second;
1203 pending[0].resolve({ stream_ticket: "stale" });
1204 await first;
1205 assert.equal(h.streams.length, 1, "a late ticket cannot create an orphan stream");
1206 assert.match(h.app.stream.url, /current$/);
1207 const third = h.connectStream("thread-a", 4, 1);
1208 assert.equal(h.streams[0].closed, true);
1209 h.stopStream();
1210 pending[2].resolve({ stream_ticket: "stopped" });
1211 await third;
1212 assert.equal(h.streams.length, 1);
1213 assert.equal(h.app.stream, null);
1214 const fourth = h.connectStream("thread-a", 4, 1);
1215 h.app.generation += 1;
1216 pending[3].reject(new Error("old selection failed"));
1217 await fourth;
1218 assert.equal(h.timers.size, 0);
1219 assert.deepEqual(h.statuses, []);
1220 });
1221
1222 test("stream recovery handshake leaves retries to snapshot recovery until open", async () => {
1223 const h = streamHarness(async () => ({ stream_ticket: "recovery" }));
1224 const opened = h.connectStream("thread-a", 4, 1, true);
1225 await new Promise((resolve) => setImmediate(resolve));
1226 h.app.stream.onerror();
1227 await assert.rejects(opened, /did not reopen/);
1228 assert.equal(h.timers.size, 0);
1229 const replacement = h.connectStream("thread-a", 4, 1, true);
1230 await new Promise((resolve) => setImmediate(resolve));
1231 h.app.stream.onopen();
1232 await replacement;
1233 h.app.stream.onerror();
1234 assert.equal(h.timers.size, 1, "normal reconnect resumes after the handshake");
1235 h.stopStream();
1236 });
1237
1238 test("records workspace restore-point receipts on their turn in order", () => {
1239 const state = createThreadState("thread-a");
1240 applySnapshot(state, snapshot("thread-a", 7));
1241 const pre = { kind: "pre_turn", snapshot_id: "c1", tree_id: "t1", session_id: "thread-a" };
1242 const post = { kind: "post_turn", snapshot_id: "c2", tree_id: "t2", session_id: "thread-a" };
1243 assert.equal(applyRuntimeEvent(state, runtimeEvent(8, "turn.workspace_snapshot", pre)), true);
1244 assert.equal(applyRuntimeEvent(state, runtimeEvent(9, "turn.workspace_snapshot", post)), true);
1245 assert.deepEqual(state.turns.get("turn-1").workspace_snapshots, [pre, post]);
1246 // A receipt for a turn the client has not seen still advances continuity.
1247 const orphan = runtimeEvent(10, "turn.workspace_snapshot", pre, { turn_id: "turn-9" });
1248 assert.equal(applyRuntimeEvent(state, orphan), true);
1249 assert.equal(state.latestSeq, 10);
1250 assert.equal(state.turns.has("turn-9"), false);
1251 });
1252
1252 lines Plain Text