Files
pikasTech-HWLAB/web/hwlab-cloud-web/scripts/workbench-realtime-runtime.test.ts
T

725 lines
39 KiB
TypeScript

// SPEC: PJ2026-0106050514 Workbench实时运行面 draft-2026-06-30-p0-1297-spec-first.
// Responsibility: Regression tests for pure Workbench realtime runtime helpers before Vue/store integration.
import assert from "node:assert/strict";
import test from "node:test";
import type { ChatMessage } from "../src/types/index.ts";
import { workbenchSessionDetailPathForTest, workbenchSessionMessagesPathForTest } from "../src/api/workbench.ts";
import { workbenchEventStreamPath, workbenchProjectionEventStreamPath } from "../src/api/workbench-events.ts";
import type { WorkbenchStreamTransportRecovery } from "../src/utils/workbench-realtime-runtime.ts";
import { workbenchRealtimeTraceIdForCapabilities, workbenchRealtimeTransportEnabled } from "../src/utils/workbench-stream-transport.ts";
import { workbenchRuntimePolicy } from "../src/config/workbench-runtime-policy.ts";
import { AsyncQueue, work } from "../src/utils/scheduler/async-queue.ts";
import { createCoalescedEventQueue } from "../src/utils/scheduler/coalesced-event-queue.ts";
import { createKeyedSingleflight } from "../src/utils/scheduler/keyed-singleflight.ts";
import { createScopedCache } from "../src/utils/scoped-cache.ts";
import { composeWorkbenchScopedKey, splitWorkbenchScopedKey, workbenchPathKey, workbenchRealtimeScopeKey } from "../src/utils/workbench-key.ts";
import { createSafeStorageRuntime, isStorageQuotaError, migrateLegacyStorage, normalizePersistedValue, readJsonStorage, removePersistedTarget, removeStorageKey, writeJsonStorage, type StorageLike } from "../src/utils/safe-storage.ts";
import { checkWorkbenchHealth, createWorkbenchHealthProbeCache } from "../src/utils/workbench-health.ts";
import { messageDiagnosticView } from "../src/utils/workbench-error-runtime.ts";
import { projectRejectedWorkbenchAdmission } from "../src/stores/workbench-admission-failure.ts";
import { WORKBENCH_TIMELINE_OPENCODE_PARITY, buildWorkbenchTimelineRows, normalizeWorkbenchTimelineMessages, workbenchTimelineSignature } from "../src/stores/workbench-timeline-model.ts";
import { reduceWorkbenchRealtimeEvent, workbenchRealtimeEventIsBusinessActivity } from "../src/stores/workbench-event-reducer.ts";
import { reduceWorkbenchLiveKafkaMessageState, workbenchAgentMessageIdForTrace } from "../src/stores/workbench-live-kafka-event.ts";
import { planWorkbenchRealtimeApply, planWorkbenchRealtimeRecovery } from "../src/stores/workbench-realtime-plan.ts";
import { WORKBENCH_REALTIME_AUTHORITY_VERSION, workbenchRealtimePrimaryAuthorityDecision } from "../src/stores/workbench-realtime-authority.ts";
import { cleanupWorkbenchServerStateDroppedSessions, cleanupWorkbenchServerStateSessions, createWorkbenchServerState, reduceWorkbenchServerState } from "../src/stores/workbench-server-state.ts";
import { cleanupDroppedWorkbenchSessionCaches, trimWorkbenchSessionCache } from "../src/stores/workbench-session-cache.ts";
import { selectActiveTurnStatusRefreshTraceIds } from "../src/stores/workbench-session.ts";
test("Workbench scoped keys encode delimiter characters", () => {
assert.deepEqual(splitWorkbenchScopedKey(workbenchRealtimeScopeKey("ses|one", "trc/two")), ["workbench.realtime", "ses|one", "trc/two"]);
assert.deepEqual(splitWorkbenchScopedKey(composeWorkbenchScopedKey("queue", null, "", "a b")), ["queue", "~", "~", "a b"]);
assert.throws(() => composeWorkbenchScopedKey("queue", "bad\0part"), /null bytes/u);
assert.equal(workbenchPathKey("C:\\Users\\foo\\"), "C:/Users/foo");
});
test("Workbench runtime policy reads injected config while preserving defaults", () => {
const policy = workbenchRuntimePolicy({
sessionListPageLimit: 7,
traceDetailMaxPages: 2,
workbenchSessionDetailMinRefreshMs: 1234,
workbenchSessionMessagesWindowLimit: 9,
workbenchTraceMessagesWindowLimit: 4,
workbenchRealtimeErrorSyncReplayMinMs: 0,
workbenchRealtimeFlushMaxItemsPerChunk: 2,
workbenchRealtimeFlushMaxChunkMs: 6,
workbenchRealtimeFlushYieldMs: 5,
defaultGatewayTimeoutMs: "bad"
});
assert.equal(policy.sessionListPageLimit, 7);
assert.equal(policy.traceDetailMaxPages, 2);
assert.equal(policy.workbenchSessionDetailMinRefreshMs, 1234);
assert.equal(policy.workbenchSessionMessagesWindowLimit, 9);
assert.equal(policy.workbenchTraceMessagesWindowLimit, 4);
assert.equal(policy.workbenchRealtimeErrorSyncReplayMinMs, 0);
assert.equal(policy.workbenchRealtimeFlushMaxItemsPerChunk, 2);
assert.equal(policy.workbenchRealtimeFlushMaxChunkMs, 6);
assert.equal(policy.workbenchRealtimeFlushYieldMs, 5);
assert.equal(policy.defaultGatewayTimeoutMs, 120_000);
});
test("Workbench runtime policy keeps deprecated trace hydration aliases compatible", () => {
const policy = workbenchRuntimePolicy({
traceHydrationPageLimit: 13,
traceHydrationMaxPages: 2,
traceHydrationAutoQueueLimit: 3,
workbenchReadHydrationConcurrency: 4
});
assert.equal(policy.traceDetailPageLimit, 13);
assert.equal(policy.traceDetailMaxPages, 2);
assert.equal(policy.traceDetailAutoQueueLimit, 3);
assert.equal(policy.workbenchDetailReadConcurrency, 4);
});
test("Workbench API uses metadata-only session detail and bounded messages paths independently", () => {
assert.equal(workbenchSessionDetailPathForTest("ses_metadata"), "/v1/workbench/sessions/ses_metadata?includeMessages=false");
assert.equal(workbenchSessionDetailPathForTest("ses_metadata", { timeoutMs: 8000 }), "/v1/workbench/sessions/ses_metadata?includeMessages=false");
assert.equal(workbenchSessionDetailPathForTest("ses_metadata", 8000), "/v1/workbench/sessions/ses_metadata?includeMessages=false");
assert.equal(workbenchSessionDetailPathForTest("ses_metadata", { includeMessages: true }), "/v1/workbench/sessions/ses_metadata?includeMessages=true");
assert.equal(workbenchSessionMessagesPathForTest("ses_metadata", { limit: 9 }), "/v1/workbench/sessions/ses_metadata/messages?limit=9");
});
test("only the projection authority sends the durable afterSeq cursor", () => {
assert.equal(workbenchEventStreamPath({ realtimeCapabilities: { liveKafkaSse: true, kafkaRefreshReplay: false, projectionRealtime: false }, sessionId: "ses_live", traceId: "trc_live", afterSeq: 42 }), "/v1/workbench/events?sessionId=ses_live&traceId=trc_live");
assert.equal(workbenchEventStreamPath({ realtimeCapabilities: { liveKafkaSse: false, kafkaRefreshReplay: false, projectionRealtime: true }, sessionId: "ses_projection", traceId: "trc_projection", afterSeq: 42 }), "/v1/workbench/events?sessionId=ses_projection&traceId=trc_projection&afterSeq=42");
assert.equal(workbenchEventStreamPath({ realtimeCapabilities: { liveKafkaSse: true, kafkaRefreshReplay: true, projectionRealtime: true }, sessionId: "ses_both", traceId: "trc_both", afterSeq: 42 }), "/v1/workbench/events?sessionId=ses_both&traceId=trc_both&afterSeq=42");
assert.equal(workbenchProjectionEventStreamPath({ sessionId: "ses_both", traceId: "trc_both", afterSeq: 42 }), "/v1/workbench/projection-events?sessionId=ses_both&traceId=trc_both&afterSeq=42");
assert.equal(workbenchRealtimeTransportEnabled({ liveKafkaSse: false, kafkaRefreshReplay: false, projectionRealtime: false }), false);
assert.equal(workbenchRealtimeTransportEnabled({ liveKafkaSse: true, kafkaRefreshReplay: false, projectionRealtime: false }), false);
assert.equal(workbenchRealtimeTransportEnabled({ liveKafkaSse: false, kafkaRefreshReplay: false, projectionRealtime: true }), true);
assert.equal(workbenchRealtimeTransportEnabled({ liveKafkaSse: true, kafkaRefreshReplay: true, projectionRealtime: true }), true);
});
test("obsolete live flags cannot erase the active projection trace scope", () => {
const capabilities = { liveKafkaSse: true, kafkaRefreshReplay: true, projectionRealtime: false };
const beforeSubmit = workbenchRealtimeScopeKey("ses_live", workbenchRealtimeTraceIdForCapabilities(capabilities, null, null));
const afterSubmit = workbenchRealtimeScopeKey("ses_live", workbenchRealtimeTraceIdForCapabilities(capabilities, "trc_current_request", "trc_message"));
assert.notEqual(afterSubmit, beforeSubmit);
assert.equal(afterSubmit, workbenchRealtimeScopeKey("ses_live", "trc_current_request"));
assert.equal(workbenchRealtimeTraceIdForCapabilities({ liveKafkaSse: false, kafkaRefreshReplay: false, projectionRealtime: true }, "trc_current_request", "trc_message"), "trc_current_request");
});
test("Error runtime owns Workbench message diagnostic view model", () => {
const degraded = messageDiagnosticView(agentMessage({
status: "running",
text: "",
projection: {
projectionHealth: "degraded",
blocker: { code: "workbench_facts_missing", userMessage: "read model unavailable", traceId: "trc_1" }
}
}));
assert.equal(degraded.visible, true);
assert.equal(degraded.text, "read model unavailable");
assert.equal(degraded.apiError?.code, "workbench_facts_missing");
assert.equal(degraded.diagnostic?.traceId, "trc_1");
const sealed = messageDiagnosticView(agentMessage({
status: "completed",
text: "done",
projection: { projectionHealth: "degraded", blocker: { code: "stale" } }
}));
assert.equal(sealed.visible, false);
});
test("fixed pre-admission diagnostic projects one assistant failure and no local trace event", () => {
const traceId = "trc_fixed_dispatch_failure";
const userMessage: ChatMessage = { id: "msg_fixed_user", messageId: "msg_fixed_user", role: "user", title: "用户", text: "hi", status: "sent", createdAt: "2026-07-12T08:33:19.000Z", sessionId: "ses_fixed", traceId };
const pending = agentMessage({
id: "msg_fixed_agent",
messageId: "msg_fixed_agent",
status: "running",
text: "",
traceId,
sessionId: "ses_fixed",
createdAt: "2026-07-12T08:33:19.000Z"
});
const projection = {
projectionStatus: "blocked",
projectionHealth: "degraded",
blocker: {
code: "schema-invalid",
failureKind: "schema-invalid",
reason: "dispatcher-runtime-misconfigured",
retryable: false,
userMessage: "AgentRun 调度契约或运行面配置不匹配,本次请求未进入实时事件流;请展开诊断修复配置。"
}
} as NonNullable<ChatMessage["projection"]>;
const failed = projectRejectedWorkbenchAdmission(pending, {
message: projection.blocker?.userMessage ?? "failed",
error: {
code: "schema-invalid",
failureKind: "schema-invalid",
reason: "dispatcher-runtime-misconfigured",
retryable: false,
recoveryAction: "repair-agentrun-dispatch-contract",
failureIdentity: "admission:msg_fixed_agent",
primaryPresentation: "assistant-message",
admissionState: "rejected-before-durable-dispatch",
traceId,
runId: null,
commandId: null,
diagnostic: {
contractVersion: "hwlab-error-diagnostic-v1",
traceId: "ae0181ba9a18bad2a561157d312c3692",
code: "schema-invalid",
layer: "agentrun",
retryable: false,
valuesPrinted: false
}
},
projection,
submittedAt: "2026-07-12T08:33:19.000Z",
finishedAt: "2026-07-12T08:33:19.057Z"
});
assert.equal(failed.runnerTrace, null);
assert.deepEqual(failed.traceEvents, []);
assert.equal((failed.admission as { state?: string }).state, "rejected-before-durable-dispatch");
assert.equal(failed.error?.failureKind, "schema-invalid");
assert.equal(failed.error?.reason, "dispatcher-runtime-misconfigured");
assert.equal(failed.error?.recoveryAction, "repair-agentrun-dispatch-contract");
assert.equal(failed.error?.runId, null);
assert.equal(failed.error?.commandId, null);
const rows = buildWorkbenchTimelineRows([userMessage, failed]);
assert.deepEqual(rows.map((row) => row.type), ["UserMessage", "AssistantPart"]);
assert.equal(rows.filter((row) => row.type === "Error").length, 0);
const diagnostic = messageDiagnosticView(failed);
assert.equal(diagnostic.visible, true);
assert.equal(diagnostic.showMessage, false);
assert.equal(diagnostic.apiError?.failureKind, "schema-invalid");
assert.equal(diagnostic.apiError?.reason, "dispatcher-runtime-misconfigured");
assert.equal(diagnostic.apiError?.recoveryAction, "repair-agentrun-dispatch-contract");
assert.equal(diagnostic.diagnostic?.traceId, "ae0181ba9a18bad2a561157d312c3692");
});
test("formal Kafka failure keeps raw events and a distinct error row", () => {
const events = [
{ sourceEventId: "evt_formal_1", sourceSeq: 1, type: "error", failureKind: "provider-unavailable", message: "provider unavailable" },
{ sourceEventId: "evt_formal_2", sourceSeq: 2, type: "result", terminal: true, failureKind: "provider-unavailable" }
];
const formal = agentMessage({
id: "msg_formal_agent",
messageId: "msg_formal_agent",
status: "failed",
text: "",
traceId: "trc_formal_failure",
error: { failureKind: "provider-unavailable", message: "provider unavailable", retryable: true },
runnerTrace: { traceId: "trc_formal_failure", status: "failed", events, eventCount: 2 }
});
const rows = buildWorkbenchTimelineRows([formal]);
assert.equal(rows.filter((row) => row.type === "Error").length, 1);
assert.equal(formal.runnerTrace?.events?.length, 2);
assert.deepEqual(formal.runnerTrace?.events, events);
});
test("scoped cache keeps OpenCode LRU and TTL semantics", () => {
const disposed: string[] = [];
let clock = 0;
let count = 0;
const cache = createScopedCache((key) => ({ key, count: ++count }), {
maxEntries: 2,
ttlMs: 10,
now: () => clock,
dispose: (value) => disposed.push(`${value.key}:${value.count}`)
});
assert.equal(cache.get("a").count, 1);
assert.equal(cache.get("b").count, 2);
cache.get("a");
assert.equal(cache.get("c").count, 3);
assert.equal(cache.peek("b"), undefined);
assert.deepEqual(disposed, ["b:2"]);
clock = 11;
assert.equal(cache.peek("a"), undefined);
assert.deepEqual(disposed, ["b:2", "a:1", "c:3"]);
});
test("async queue and bounded work preserve OpenCode utility behavior", async () => {
const queue = new AsyncQueue<number>();
const first = queue.next();
queue.push(1);
assert.equal(await first, 1);
queue.push(2);
assert.equal(await queue.next(), 2);
const visited: number[] = [];
await work(2, [1, 2, 3], async (item) => {
visited.push(item);
});
assert.deepEqual(visited.sort(), [1, 2, 3]);
});
test("coalesced event queue replaces keyed snapshots and preserves unkeyed events", () => {
const flushed: string[][] = [];
const queue = createCoalescedEventQueue<{ key?: string; value: string }>({
keyOf: (item) => item.key,
schedule: () => undefined,
onFlush: (items) => flushed.push(items.map((item) => item.value))
});
queue.push({ key: "snapshot:trc", value: "old" });
queue.push({ value: "event-1" });
queue.push({ key: "snapshot:trc", value: "new" });
queue.push({ value: "event-2" });
assert.equal(queue.size, 3);
queue.flush();
assert.deepEqual(flushed, [["new", "event-1", "event-2"]]);
});
test("coalesced event queue chunks scheduled flushes and yields between chunks", () => {
const flushed: string[] = [];
const drain: unknown[] = [];
const scheduled: (() => void)[] = [];
const yielded: (() => void)[] = [];
let clock = 0;
const queue = createCoalescedEventQueue<{ key?: string; value: string }>({
keyOf: (item) => item.key,
maxItemsPerChunk: 10,
maxChunkMs: 4,
now: () => clock,
schedule: (flush) => {
scheduled.push(flush);
return () => undefined;
},
yieldSchedule: (flush) => {
yielded.push(flush);
return () => undefined;
},
onFlush: (items) => {
clock += 3;
flushed.push(...items.map((item) => item.value));
},
onDrain: (info) => drain.push(info)
});
queue.push({ key: "snapshot:trc", value: "old" });
queue.push({ value: "event-1" });
queue.push({ key: "snapshot:trc", value: "new" });
queue.push({ value: "event-2" });
queue.push({ value: "event-3" });
assert.equal(queue.size, 4);
assert.equal(scheduled.length, 1);
scheduled.shift()?.();
assert.deepEqual(flushed, ["new", "event-1"]);
assert.equal(queue.size, 2);
assert.equal(yielded.length, 1);
yielded.shift()?.();
assert.deepEqual(flushed, ["new", "event-1", "event-2", "event-3"]);
assert.equal(queue.size, 0);
assert.deepEqual(drain, [{ eventCount: 5, deliveredCount: 4, droppedCount: 1, chunkCount: 2, replacedByKey: 1, flushDurationMs: 12, maxItemsPerChunk: 10, maxChunkMs: 4, reason: "yield" }]);
});
test("keyed single-flight coalesces duplicate work and supports replacement", async () => {
const flight = createKeyedSingleflight<number>();
let calls = 0;
const first = flight.run("sessions", () => {
calls += 1;
return 7;
});
const second = flight.run("sessions", () => {
calls += 1;
return 9;
});
assert.equal(await first, 7);
assert.equal(await second, 7);
assert.equal(calls, 1);
assert.equal(await flight.run("sessions", () => 11, { replace: true, reason: "manual-refresh" }), 11);
assert.equal(calls, 1);
});
test("safe storage reports failures instead of throwing", () => {
const storage = new MemoryStorage();
assert.equal(writeJsonStorage(storage, "k", { a: 1 }).ok, true);
assert.deepEqual(readJsonStorage(storage, "k", { a: 0 }).value, { a: 1 });
assert.equal(removeStorageKey(storage, "k").ok, true);
const broken: StorageLike = {
getItem: () => "{not-json",
setItem: () => { throw new Error("quota"); },
removeItem: () => { throw new Error("security"); }
};
assert.equal(readJsonStorage(broken, "bad", { safe: true }).status, "parse_error");
assert.equal(writeJsonStorage(broken, "bad", { safe: true }).status, "write_error");
assert.equal(removeStorageKey(broken, "bad").status, "remove_error");
});
test("safe storage preserves OpenCode quota fallback, prefix isolation and legacy migration", () => {
const storage = new MemoryStorage();
const runtime = createSafeStorageRuntime({ cacheMaxEntries: 3, cacheMaxBytes: 200, evictPrefix: "hwlab." });
const bad = runtime.localStorageWithPrefix("hwlab.throw.scope", storage);
bad.setItem("value", '{"value":1}');
const before = storage.calls.set;
bad.setItem("value", '{"value":2}');
assert.equal(storage.calls.set, before);
const safe = runtime.localStorageWithPrefix("hwlab.safe.scope", storage);
safe.setItem("value", '{"value":3}');
assert.equal(storage.getItem("hwlab.safe.scope:value"), '{"value":3}');
assert.equal(isStorageQuotaError(new DOMException("quota", "QuotaExceededError")), true);
assert.equal(normalizePersistedValue({ value: 1, nested: { a: true } }, '{"nested":{"b":2}}'), '{"value":1,"nested":{"a":true,"b":2}}');
assert.equal(normalizePersistedValue({ value: "ok" }, '{"value":"\\x"}'), undefined);
storage.setItem("legacy.workspace", '{"value":2}');
const current = runtime.localStorageWithPrefix("hwlab.current", storage);
const migrated = migrateLegacyStorage({ current, legacyStore: runtime.localStorageDirect(storage), stores: [], keys: ["legacy.workspace"], key: "workspace:demo", defaults: { value: 1 } });
assert.equal(migrated, '{"value":2}');
assert.equal(storage.getItem("hwlab.current:workspace:demo"), '{"value":2}');
assert.equal(storage.getItem("legacy.workspace"), null);
storage.setItem("hwlab.current:workspace:demo", '{"value":2}');
storage.setItem("hwlab.legacy:workspace:demo", '{"value":3}');
removePersistedTarget({ storage: "hwlab.current", legacyStorageNames: ["hwlab.legacy"], key: "workspace:demo" }, runtime, storage);
assert.equal(storage.getItem("hwlab.current:workspace:demo"), null);
assert.equal(storage.getItem("hwlab.legacy:workspace:demo"), null);
});
test("timeline model deduplicates authoritative messages with stable row keys", () => {
const rows = buildWorkbenchTimelineRows([
{ id: "msg_user", messageId: "msg_user", role: "user", title: "用户", text: "hi", status: "sent", createdAt: "2026-06-30T00:00:00.000Z", sessionId: "ses_a", traceId: "trc_a" },
{ id: "msg_agent_old", messageId: "msg_agent", role: "agent", title: "Code Agent", text: "", status: "running", createdAt: "2026-06-30T00:00:01.000Z", sessionId: "ses_a", traceId: "trc_a" },
{ id: "msg_agent_new", messageId: "msg_agent", role: "agent", title: "Code Agent", text: "done", status: "completed", createdAt: "2026-06-30T00:00:01.000Z", sessionId: "ses_a", traceId: "trc_a" }
]);
assert.equal(rows.length, 2);
const messageRows = rows.flatMap((row) => row.message ? [row] : []);
assert.deepEqual(messageRows.map((row) => `${row.role}:${row.message.text}`), ["user:hi", "agent:done"]);
assert.match(workbenchTimelineSignature(rows), /completed:4/u);
assert.deepEqual(normalizeWorkbenchTimelineMessages(messageRows.map((row) => row.message)).map((message) => message.messageId), ["msg_user", "msg_agent"]);
});
test("timeline model preserves OpenCode turn gap, thinking, retry and error row parity", () => {
const rows = buildWorkbenchTimelineRows([
{ id: "msg_user_1", messageId: "msg_user_1", role: "user", title: "用户", text: "hi", status: "sent", createdAt: "2026-06-30T00:00:00.000Z", sessionId: "ses_a", traceId: "trc_a" },
{ id: "msg_agent_1", messageId: "msg_agent_1", role: "agent", title: "Code Agent", text: "", status: "running", createdAt: "2026-06-30T00:00:01.000Z", sessionId: "ses_a", traceId: "trc_a" },
{ id: "msg_user_2", messageId: "msg_user_2", role: "user", title: "用户", text: "again", status: "sent", createdAt: "2026-06-30T00:00:02.000Z", sessionId: "ses_a", traceId: "trc_b" },
{ id: "msg_agent_2", messageId: "msg_agent_2", role: "agent", title: "Code Agent", text: "", status: "retry" as never, createdAt: "2026-06-30T00:00:03.000Z", sessionId: "ses_a", traceId: "trc_b", error: { message: "Error: {\"error\":{\"type\":\"Provider\",\"message\":\"offline\"}}" } }
]);
assert.ok(rows.some((row) => row.type === "TurnGap"));
assert.ok(rows.some((row) => row.type === "Thinking"));
assert.ok(rows.some((row) => row.type === "Retry"));
assert.ok(rows.some((row) => row.type === "Error" && row.text === "Provider: offline"));
});
test("timeline model documents OpenCode row parity and preserves comment/diff rows", () => {
assert.ok(WORKBENCH_TIMELINE_OPENCODE_PARITY.some((row) => row.openCodeRows === "Permission" && row.status === "not-applicable"));
assert.ok(WORKBENCH_TIMELINE_OPENCODE_PARITY.some((row) => row.openCodeRows === "Tool" && row.status === "adapted"));
const rows = buildWorkbenchTimelineRows([
{ id: "msg_user", messageId: "msg_user", role: "user", title: "用户", text: "review", status: "sent", createdAt: "2026-06-30T00:00:00.000Z", sessionId: "ses_a", traceId: "trc_a", comments: [{ text: "note" }] } as ChatMessage,
{ id: "msg_agent", messageId: "msg_agent", role: "agent", title: "Code Agent", text: "done", status: "completed", createdAt: "2026-06-30T00:00:01.000Z", sessionId: "ses_a", traceId: "trc_a", summary: { diffs: [{ file: "a.ts" }, { file: "a.ts" }, { file: "b.ts" }] } } as ChatMessage
]);
assert.ok(rows.some((row) => row.type === "CommentStrip"));
const diff = rows.find((row) => row.type === "DiffSummary");
if (!diff || diff.type !== "DiffSummary") throw new Error("missing DiffSummary row");
assert.equal(diff?.diffs?.length, 2);
});
test("session cache trim keeps retained sessions and trims messages by explicit policy", () => {
const result = trimWorkbenchSessionCache({
sessionOrder: ["ses_1", "ses_2", "ses_3"],
sessionsById: {
ses_1: { sessionId: "ses_1" },
ses_2: { sessionId: "ses_2" },
ses_3: { sessionId: "ses_3" }
},
messagesBySessionId: {
ses_1: [{ id: "m1", role: "user", title: "用户", text: "1", status: "sent", createdAt: "2026-06-30T00:00:00.000Z" }],
ses_2: [
{ id: "m2", role: "user", title: "用户", text: "2", status: "sent", createdAt: "2026-06-30T00:00:00.000Z" },
{ id: "m3", role: "agent", title: "Code Agent", text: "3", status: "completed", createdAt: "2026-06-30T00:00:01.000Z" }
],
ses_3: [{ id: "m4", role: "user", title: "用户", text: "4", status: "sent", createdAt: "2026-06-30T00:00:00.000Z" }]
}
}, { retainSessionIds: ["ses_3"], maxSessions: 1, maxMessagesPerSession: 1 });
assert.deepEqual(result.state.sessionOrder, ["ses_1", "ses_3"]);
assert.deepEqual(result.evictedSessionIds, ["ses_2"]);
assert.deepEqual(Object.keys(result.state.sessionsById).sort(), ["ses_1", "ses_3"]);
});
test("session cache cleanup drops stale session/message caches with OpenCode cleanup parity", () => {
const result = cleanupDroppedWorkbenchSessionCaches({
sessionOrder: ["ses_1", "ses_2"],
sessionsById: { ses_1: { sessionId: "ses_1" }, ses_2: { sessionId: "ses_2" } },
messagesBySessionId: { ses_1: [], ses_2: [] }
}, ["ses_2"]);
assert.deepEqual(result.droppedSessionIds, ["ses_1"]);
assert.deepEqual(result.state.sessionOrder, ["ses_2"]);
assert.deepEqual(Object.keys(result.state.sessionsById), ["ses_2"]);
});
test("server-state cleanup removes dropped session trace and turn authority", () => {
let state = createWorkbenchServerState();
state = reduceWorkbenchServerState(state, { type: "session.detail", session: { sessionId: "ses_1", messages: [{ id: "m1", role: "agent", title: "Code Agent", text: "", status: "running", createdAt: "2026-06-30T00:00:00.000Z", traceId: "trc_1", runnerTrace: { traceId: "trc_1", sessionId: "ses_1", events: [] } }] } });
state = reduceWorkbenchServerState(state, { type: "turn.status", turn: { traceId: "trc_1", status: "running", running: true, terminal: false, sessionId: "ses_1" } });
state = reduceWorkbenchServerState(state, { type: "trace.snapshot", traceId: "trc_1", trace: { traceId: "trc_1", sessionId: "ses_1", events: [] } });
const cleaned = cleanupWorkbenchServerStateDroppedSessions(state, []);
assert.deepEqual(Object.keys(cleaned.messagesBySessionId), []);
assert.deepEqual(Object.keys(cleaned.turnStatusByTraceId), []);
assert.deepEqual(Object.keys(cleaned.traceById), []);
const cleanedById = cleanupWorkbenchServerStateSessions(state, ["ses_1"]);
assert.deepEqual(Object.keys(cleanedById.messagesBySessionId), []);
assert.deepEqual(Object.keys(cleanedById.turnStatusByTraceId), []);
assert.deepEqual(Object.keys(cleanedById.traceById), []);
});
test("passive message projection selects active turn status refresh without local request", () => {
const traceIds = selectActiveTurnStatusRefreshTraceIds({
messages: [
agentMessage({ id: "m_old", traceId: "trc_old", status: "completed", text: "done" }),
agentMessage({ id: "m_new", traceId: "trc_new", status: "running", text: "thinking" })
],
currentRequestTraceId: null,
turnStatusAuthority: {}
});
assert.deepEqual(traceIds, ["trc_new"]);
});
test("passive message projection skips sealed terminal turn status refresh", () => {
const traceIds = selectActiveTurnStatusRefreshTraceIds({
messages: [agentMessage({ id: "m_done", traceId: "trc_done", status: "completed", text: "done" })],
turnStatusAuthority: { trc_done: { traceId: "trc_done", status: "completed", running: false, terminal: true, sessionId: "ses_1" } }
});
assert.deepEqual(traceIds, []);
});
test("turn status refresh keeps local request trace priority", () => {
const traceIds = selectActiveTurnStatusRefreshTraceIds({
messages: [
agentMessage({ id: "m_request", traceId: "trc_request", status: "pending", text: "" }),
agentMessage({ id: "m_observed", traceId: "trc_observed", status: "running", text: "thinking" })
],
currentRequestTraceId: "trc_request",
turnStatusAuthority: {},
limit: 2
});
assert.deepEqual(traceIds, ["trc_request", "trc_observed"]);
});
test("realtime event reducer classifies SSE payloads before store side effects", () => {
const trace = reduceWorkbenchRealtimeEvent(realtimeEvent({ type: "trace.event", traceId: "trc_1", event: { traceId: "trc_1", label: "delta" }, snapshot: { traceId: "trc_1", status: "running" }, entity: { family: "traceEvents", id: "trc_1:1", version: 1, projectionRevision: "prj_1" } }), "workbench.trace.event");
assert.equal(trace.activityLabel, "realtime:trace.event");
assert.equal(trace.action.type, "trace.event");
assert.equal(trace.diagnostic.module, "workbench-event-reducer");
const missingAuthority = reduceWorkbenchRealtimeEvent({ type: "message.snapshot", message: agentMessage({ status: "running", sessionId: "ses_1" }) }, "workbench.message.snapshot");
assert.deepEqual(missingAuthority.action, { type: "ignore", reason: "workbench_realtime_authority_missing" });
const detailOnly = reduceWorkbenchRealtimeEvent(realtimeEvent({ type: "trace.event", traceId: "trc_1", detailProjection: true, authority: "trace-detail-only", event: { traceId: "trc_1", label: "detail" }, entity: { family: "traceEvents", id: "trc_1:2", version: 2, projectionRevision: "prj_1", authority: "trace-detail-only" } }), "workbench.trace.event");
assert.deepEqual(detailOnly.action, { type: "ignore", reason: "workbench_realtime_detail_only_rejected" });
const error = reduceWorkbenchRealtimeEvent({ type: "error", traceId: "trc_2", error: { message: "offline" } }, "workbench.error");
assert.equal(error.action.type, "projection.error");
assert.equal(error.diagnostic.traceId, "trc_2");
});
test("live transport connected and heartbeat frames do not extend business activity timeouts", () => {
assert.equal(workbenchRealtimeEventIsBusinessActivity({ type: "connected" }, "workbench.connected", true), false);
assert.equal(workbenchRealtimeEventIsBusinessActivity({ type: "heartbeat" }, "workbench.heartbeat", true), false);
assert.equal(workbenchRealtimeEventIsBusinessActivity({ schema: "hwlab.event.v1", event: { type: "assistant" } }, "hwlab.event.v1", true), true);
assert.equal(workbenchRealtimeEventIsBusinessActivity({ type: "trace.event" }, "workbench.trace.event", false), true);
});
test("live hwlab envelope projects assistant, tool/output, and terminal without replay or finalizer", () => {
const envelope = (event: Record<string, unknown>) => ({
schema: "hwlab.event.v1",
eventType: "hwlab.trace.event.projected",
eventId: `hwlab:${String(event.sourceEventId ?? event.type)}`,
hwlabSessionId: "ses_live_web",
sessionId: "ses_live_web",
traceId: "trc_live_web",
runId: "run_live_web",
commandId: "cmd_live_web",
context: { runId: "run_live_web", commandId: "cmd_live_web", valuesRedacted: true },
event
});
const assistantEnvelope = envelope({ type: "assistant", eventType: "assistant", sourceEventId: "evt_assistant", traceId: "trc_live_web", sessionId: "ses_live_web", status: "running", assistantText: "running increment", terminal: false });
const assistant = reduceWorkbenchRealtimeEvent(assistantEnvelope, "hwlab.event.v1");
assert.equal(assistant.action.type, "trace.event");
let visible = reduceWorkbenchLiveKafkaMessageState({ text: "", status: "running", terminal: false }, assistantEnvelope.event);
assert.deepEqual(visible, { text: "running increment", status: "running", terminal: false });
const toolEnvelope = envelope({ type: "tool", eventType: "tool", sourceEventId: "evt_tool", traceId: "trc_live_web", sessionId: "ses_live_web", status: "completed", toolName: "commandExecution", outputSummary: "ok", terminal: false });
assert.equal(reduceWorkbenchRealtimeEvent(toolEnvelope, "hwlab.event.v1").action.type, "trace.event");
visible = reduceWorkbenchLiveKafkaMessageState(visible, toolEnvelope.event);
assert.deepEqual(visible, { text: "running increment", status: "running", terminal: false });
const outputEnvelope = envelope({ type: "output", eventType: "status", sourceEventId: "evt_output", traceId: "trc_live_web", sessionId: "ses_live_web", status: "running", message: "stdout increment", terminal: false });
visible = reduceWorkbenchLiveKafkaMessageState(visible, outputEnvelope.event);
assert.deepEqual(visible, { text: "running increment", status: "running", terminal: false });
const terminalEnvelope = envelope({ type: "result", eventType: "terminal", sourceEventId: "evt_terminal", traceId: "trc_live_web", sessionId: "ses_live_web", status: "completed", terminal: true });
visible = reduceWorkbenchLiveKafkaMessageState(visible, terminalEnvelope.event);
assert.deepEqual(visible, { text: "running increment", status: "completed", terminal: true });
assert.equal(workbenchAgentMessageIdForTrace("trc_live_web"), "msg_live_web_agent");
const failed = reduceWorkbenchLiveKafkaMessageState(
{ text: "partial response", status: "running", terminal: false },
{ type: "result", eventType: "terminal", status: "failed", message: "provider stream disconnected", terminal: true }
);
assert.deepEqual(failed, { text: "provider stream disconnected", status: "failed", terminal: true });
});
test("realtime apply planner turns reducer actions into store steps", () => {
const reduced = reduceWorkbenchRealtimeEvent(realtimeEvent({ type: "trace.event", traceId: "trc_1", event: { traceId: "trc_1", label: "delta" }, snapshot: { traceId: "trc_1", status: "running" }, entity: { family: "traceEvents", id: "trc_1:1", version: 1, projectionRevision: "prj_1" } }), "workbench.trace.event");
const tracePlan = planWorkbenchRealtimeApply(reduced.action);
assert.deepEqual(tracePlan.steps.map((step) => step.type), ["apply-trace-event"]);
assert.equal(tracePlan.diagnostic.module, "workbench-realtime-plan");
const unavailable = planWorkbenchRealtimeApply({ type: "trace.unavailable", traceId: "trc_2", reason: "lost" });
assert.deepEqual(unavailable.steps, [{ type: "clear-active-trace", traceId: "trc_2", reason: "lost" }]);
const ignored = planWorkbenchRealtimeApply({ type: "ignore", reason: "unsupported" });
assert.deepEqual(ignored.steps, []);
});
test("realtime recovery planner reconnects the projection SSE from its outbox cursor", () => {
const recovery = recoveryEvent(["events-reconnect"], { outboxSeq: 42 });
const authorized = planWorkbenchRealtimeRecovery(recovery, { selectedSessionId: "ses_1", activeSessionId: "ses_1", fallbackTraceId: "trc_1", activeTraceAuthorized: true });
assert.deepEqual(authorized.steps.map((step) => step.type), ["events-reconnect"]);
assert.equal(authorized.steps[0]?.afterOutboxSeq, 42);
assert.equal(authorized.sessionId, "ses_1");
assert.equal(authorized.traceId, "trc_1");
const unauthorized = planWorkbenchRealtimeRecovery(recovery, { selectedSessionId: "ses_1", activeSessionId: "ses_1", fallbackTraceId: "trc_1", activeTraceAuthorized: false });
assert.deepEqual(unauthorized.steps.map((step) => step.type), ["events-reconnect"]);
const inactive = planWorkbenchRealtimeRecovery(recovery, { selectedSessionId: "ses_1", activeSessionId: "ses_other", fallbackTraceId: "trc_1", activeTraceAuthorized: true });
assert.deepEqual(inactive.steps.map((step) => step.type), ["events-reconnect"]);
const terminalSealed = planWorkbenchRealtimeRecovery(recovery, { selectedSessionId: "ses_1", activeSessionId: "ses_1", fallbackTraceId: "trc_1", activeTraceAuthorized: true, terminalTraceSealed: true });
assert.deepEqual(terminalSealed.steps.map((step) => step.type), ["events-reconnect"]);
});
test("realtime authority accepts projection outbox replay and live events through one entity contract", () => {
const event = realtimeEvent({ type: "message.snapshot", sessionId: "ses_1", message: agentMessage({ id: "msg_1", sessionId: "ses_1", status: "completed", text: "done" }), entity: { family: "messages", id: "msg_1", version: 7, outboxSeq: 12, projectionRevision: "prj_7" } });
const decision = workbenchRealtimePrimaryAuthorityDecision(event);
assert.equal(decision.accepted, true);
assert.equal(decision.entity?.family, "messages");
assert.equal(decision.entity?.version, 7);
});
test("realtime authority rejects trace detail-only and incomplete contract payloads", () => {
const detail = realtimeEvent({ type: "trace.event", traceId: "trc_1", detailProjection: true, authority: "trace-detail-only", event: { traceId: "trc_1" }, entity: { family: "traceEvents", id: "trc_1:1", version: 1, projectionRevision: "prj_1", authority: "trace-detail-only" } });
assert.equal(workbenchRealtimePrimaryAuthorityDecision(detail).reason, "workbench_realtime_detail_only_rejected");
const missingEntity = realtimeEvent({ type: "message.snapshot", message: agentMessage({ status: "running", sessionId: "ses_1" }) });
assert.equal(workbenchRealtimePrimaryAuthorityDecision(missingEntity).reason, "workbench_realtime_entity_missing");
const missingProjection = realtimeEvent({ type: "message.snapshot", message: agentMessage({ status: "running", sessionId: "ses_1" }), entity: { family: "messages", id: "msg_1", version: 1 } });
assert.equal(workbenchRealtimePrimaryAuthorityDecision(missingProjection).reason, "workbench_realtime_projection_revision_missing");
});
test("health probe cache records ok and unavailable states", async () => {
const cache = createWorkbenchHealthProbeCache({ cacheMs: 100 });
const ok = await cache.probe({ key: "workbench", fetcher: async () => ({ ready: true }), classify: (value) => value.ready ? "ok" : "degraded" });
assert.equal(ok.state, "ok");
assert.deepEqual(cache.get("workbench")?.value, { ready: true });
cache.clear("workbench");
const failed = await cache.probe({ key: "workbench", fetcher: async () => { throw new Error("offline"); } });
assert.equal(failed.state, "unavailable");
assert.ok(failed.error instanceof Error);
});
test("health probe retries transient failures with caller supplied policy", async () => {
let calls = 0;
const result = await checkWorkbenchHealth({
key: "workbench",
retryCount: 2,
retryDelayMs: 1,
fetcher: async () => {
calls += 1;
if (calls < 3) throw new TypeError("network");
return { ready: true, version: "v" };
},
classify: (value) => value.ready ? "ok" : "degraded"
});
assert.equal(calls, 3);
assert.equal(result.state, "ok");
assert.deepEqual(result.value, { ready: true, version: "v" });
});
class MemoryStorage implements StorageLike {
private readonly values = new Map<string, string>();
readonly calls = { get: 0, set: 0, remove: 0 };
get length(): number {
return this.values.size;
}
key(index: number): string | null {
return Array.from(this.values.keys())[index] ?? null;
}
getItem(key: string): string | null {
this.calls.get += 1;
if (key.startsWith("hwlab.throw")) throw new Error("storage get failed");
return this.values.get(key) ?? null;
}
setItem(key: string, value: string): void {
this.calls.set += 1;
if (key.startsWith("hwlab.quota")) throw new DOMException("quota", "QuotaExceededError");
if (key.startsWith("hwlab.throw")) throw new Error("storage set failed");
this.values.set(key, value);
}
removeItem(key: string): void {
this.calls.remove += 1;
if (key.startsWith("hwlab.throw")) throw new Error("storage remove failed");
this.values.delete(key);
}
}
function agentMessage(overrides: Partial<ChatMessage>): ChatMessage {
return {
id: "msg_agent",
role: "agent",
title: "Code Agent",
text: "",
status: "running",
createdAt: "2026-06-30T00:00:00.000Z",
traceId: "trc_1",
...overrides
} as ChatMessage;
}
function recoveryEvent(actions: WorkbenchStreamTransportRecovery["actions"], cursor: { outboxSeq?: number | null; traceSeq?: number | null } = {}): WorkbenchStreamTransportRecovery {
return {
key: "workbench.realtime|ses_1|trc_1",
sessionId: "ses_1",
traceId: "trc_1",
tick: 1,
reason: "eventsource-error",
actions,
outboxSeq: cursor.outboxSeq ?? null,
traceSeq: cursor.traceSeq ?? null,
diagnostic: { code: "workbench_sse_recovery", valuesRedacted: true }
};
}
function realtimeEvent(input: Record<string, unknown>) {
return { realtimeAuthority: WORKBENCH_REALTIME_AUTHORITY_VERSION, contractVersion: "workbench-sync-v1", ...input };
}