Files
pikasTech-HWLAB/internal/cloud/server-workbench-realtime-http.test.ts
T

1073 lines
49 KiB
TypeScript

// SPEC: PJ2026-0104010803 Workbench唯一投影 draft-2026-06-20-p2-terminal-outbox-recovery; PJ2026-010403 API契约 draft-2026-06-20-p2-terminal-outbox-recovery.
// Responsibility: Workbench realtime and projection blocker regression tests.
import assert from "node:assert/strict";
import { get as httpGet } from "node:http";
import { test } from "bun:test";
import { createCloudApiServer } from "./server.ts";
import { createBackendPerformanceStore } from "./backend-performance.ts";
import { createCodeAgentTraceStore } from "./code-agent-trace-store.ts";
import { classifyWorkbenchReadModelFailure } from "./server-workbench-http.ts";
import { projectionOutboxRealtimeEvents, workbenchRealtimeAfterSeq } from "./server-workbench-realtime-http.ts";
import { codeAgentTurnStatusPayload, createCodeAgentChatResultStore } from "./server-code-agent-http.ts";
import { createWorkbenchTurnProjection, durableTraceStatus, projectionDiagnostics, traceTerminalEvidence } from "./workbench-turn-projection.ts";
import {
ACTOR,
getJson,
waitForCondition,
createDurableFactsRuntimeStore,
createRuntimeStoreFromFacts,
buildDurableFactsForSession,
normalizeTestMessages,
testTimingProjection,
timestampIso,
elapsedMs,
normalizeTestEvents,
filterFacts,
workbenchTestFactOrder,
matchesFact,
matchesTraceFact,
testFactFamilySet,
mergeFacts,
emptyFacts,
normalizeTestStatus,
getSseEvents,
parseSseBlock
} from "./server-workbench-http-test-helpers.ts";
const TRANSACTIONAL_REALTIME_ENV = Object.freeze({
HWLAB_KAFKA_DIRECT_PUBLISH_ENABLED: "false",
HWLAB_WORKBENCH_LIVE_KAFKA_SSE_ENABLED: "false",
HWLAB_KAFKA_TRANSACTIONAL_PROJECTOR_ENABLED: "false",
HWLAB_KAFKA_PROJECTION_OUTBOX_RELAY_ENABLED: "false",
HWLAB_WORKBENCH_PROJECTION_REALTIME_ENABLED: "true"
});
function projectionRealtimeBridge(capabilities = {}) {
return {
started: true,
capabilities: {
directPublish: false,
liveKafkaSse: false,
transactionalProjector: false,
projectionOutboxRelay: false,
projectionRealtime: true,
...capabilities
},
ready: Promise.resolve(),
subscribeLiveHwlabEvents() { return () => {}; },
subscribeProjectionCommits() { return () => {}; },
async stop() {}
};
}
test("projection outbox emits immutable assistant versions in row order", () => {
const traceId = "trc_outbox_versions";
const sessionId = "ses_outbox_versions";
const messageId = "msg_outbox_versions_agent";
const events = projectionOutboxRealtimeEvents({
facts: { messages: [{ messageId, traceId, sessionId, text: "current state must not replace history", projectedSeq: 99 }] },
events: [
{ outboxSeq: 10, outboxEventId: "outbox-progress-1", entityFamily: "messages", entityId: messageId, projectedSeq: 1, projectionRevision: 1, traceId, sessionId, commitType: "message", payload: { family: "messages", fact: { messageId, traceId, sessionId, text: "first progress", projectedSeq: 1, status: "running" } } },
{ outboxSeq: 11, outboxEventId: "outbox-progress-2", entityFamily: "messages", entityId: messageId, projectedSeq: 2, projectionRevision: 2, traceId, sessionId, commitType: "message", payload: { family: "messages", fact: { messageId, traceId, sessionId, text: "second progress", projectedSeq: 2, status: "running" } } }
]
});
assert.deepEqual(events.map((item) => item.payload.message.text), ["first progress", "second progress"]);
assert.deepEqual(events.map((item) => item.payload.cursor.outboxSeq), [10, 11]);
});
test("live Kafka SSE transparently fans out one envelope without DB, snapshot, cursor, replay, or SSE id", async () => {
const sessionId = "ses_live_kafka_sse";
const traceId = "trc_live_kafka_sse";
const subscribers = new Set();
const otelSpans = [];
let releaseBridgeReady;
const bridgeReady = new Promise((resolve) => { releaseBridgeReady = resolve; });
const bridge = {
capabilities: { directPublish: true, liveKafkaSse: true, transactionalProjector: false, projectionOutboxRelay: false, projectionRealtime: false },
ready: bridgeReady,
subscribeLiveHwlabEvents(listener) {
subscribers.add(listener);
return () => subscribers.delete(listener);
},
async stop() {}
};
let dbCalls = 0;
const workbenchRuntime = new Proxy({}, {
get() {
dbCalls += 1;
throw new Error("live SSE must not access Workbench runtime facts");
}
});
const server = createCloudApiServer({
accessController: realtimeAccessController({
sessions: [{ id: sessionId, ownerUserId: ACTOR.id, lastTraceId: traceId }]
}),
workbenchRuntime,
kafkaEventBridge: bridge,
otelSpanEmitter(name, emittedTraceId, env, spanOptions) { otelSpans.push({ name, emittedTraceId, env, spanOptions }); },
env: {
HWLAB_KAFKA_DIRECT_PUBLISH_ENABLED: "true",
HWLAB_WORKBENCH_LIVE_KAFKA_SSE_ENABLED: "true",
HWLAB_KAFKA_TRANSACTIONAL_PROJECTOR_ENABLED: "false",
HWLAB_KAFKA_PROJECTION_OUTBOX_RELAY_ENABLED: "false",
HWLAB_WORKBENCH_PROJECTION_REALTIME_ENABLED: "false"
}
});
await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve));
const envelope = {
schema: "hwlab.event.v1",
eventType: "hwlab.trace.event.projected",
eventId: "hwlab:evt_live_kafka_sse",
traceId,
hwlabSessionId: sessionId,
sessionId,
runId: "run_live_kafka_sse",
commandId: "cmd_live_kafka_sse",
context: { runId: "run_live_kafka_sse", commandId: "cmd_live_kafka_sse", sourceSeq: 4, valuesRedacted: true },
event: { type: "assistant", eventType: "assistant", status: "running", traceId, sessionId, text: "live increment", sourceSeq: 4, terminal: false, valuesPrinted: false },
valuesPrinted: false
};
try {
const { port } = server.address();
const firstPromise = getSseEvents(port, `/v1/workbench/events?sessionId=${sessionId}&afterSeq=999`, 2);
const secondPromise = getSseEvents(port, `/v1/workbench/events?sessionId=${sessionId}&traceId=${traceId}&afterSeq=888`, 2);
await waitForCondition(() => subscribers.size === 2);
for (const listener of subscribers) listener(envelope, { topic: "hwlab.event.v1", partition: 6, offset: "101" });
releaseBridgeReady();
const [first, second] = await Promise.all([firstPromise, secondPromise]);
for (const events of [first, second]) {
assert.deepEqual(events.map((event) => event.event), ["workbench.connected", "hwlab.event.v1"]);
assert.deepEqual(events[0].data.capabilities, bridge.capabilities);
assert.equal(events[0].data.liveOnly, true);
assert.equal(events[0].data.replay, false);
assert.equal(events[0].data.lossPossible, true);
assert.equal(events[0].data.cursor, undefined);
assert.equal(events[1].id, null);
assert.deepEqual(events[1].data, envelope);
}
assert.equal(otelSpans.length, 2);
for (const span of otelSpans) {
assert.equal(span.name, "hwlab.workbench.live_sse.business_event_write");
assert.equal(span.emittedTraceId, traceId);
assert.deepEqual(span.spanOptions.attributes, {
businessTraceId: traceId,
hwlabSessionId: sessionId,
runId: "run_live_kafka_sse",
commandId: "cmd_live_kafka_sse",
topic: "hwlab.event.v1",
partition: 6,
offset: "101",
sourceTopic: null,
sourcePartition: null,
sourceOffset: null,
eventType: "assistant",
terminal: false,
valuesRedacted: true
});
assert.equal(JSON.stringify(span.spanOptions.attributes).includes("live increment"), false);
}
assert.equal(dbCalls, 0);
} finally {
await new Promise((resolve, reject) => server.close((error) => error ? reject(error) : resolve()));
}
});
test("live Kafka SSE heartbeat keeps the transport open without DB, cursor, snapshot, or replay", async () => {
const sessionId = "ses_live_kafka_heartbeat";
let dbCalls = 0;
const server = createCloudApiServer({
accessController: realtimeAccessController({ sessions: [{ id: sessionId, ownerUserId: ACTOR.id }] }),
workbenchRuntime: new Proxy({}, {
get() {
dbCalls += 1;
throw new Error("live heartbeat must not access Workbench projection state");
}
}),
kafkaEventBridge: {
capabilities: { liveKafkaSse: true },
ready: Promise.resolve(),
subscribeLiveHwlabEvents() { return () => {}; },
async stop() {}
},
env: {
HWLAB_KAFKA_DIRECT_PUBLISH_ENABLED: "true",
HWLAB_WORKBENCH_LIVE_KAFKA_SSE_ENABLED: "true",
HWLAB_KAFKA_TRANSACTIONAL_PROJECTOR_ENABLED: "false",
HWLAB_KAFKA_PROJECTION_OUTBOX_RELAY_ENABLED: "false",
HWLAB_WORKBENCH_PROJECTION_REALTIME_ENABLED: "false",
HWLAB_WORKBENCH_SSE_HEARTBEAT_MS: "5"
}
});
await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve));
try {
const events = await getSseEvents(server.address().port, `/v1/workbench/events?sessionId=${sessionId}`, 2);
assert.deepEqual(events.map((event) => event.event), ["workbench.connected", "workbench.heartbeat"]);
assert.equal(events[1].id, null);
assert.equal(events[1].data.liveOnly, true);
assert.equal(events[1].data.replay, false);
assert.equal(events[1].data.cursor, undefined);
assert.equal(events[1].data.snapshot, undefined);
assert.equal(dbCalls, 0);
} finally {
await new Promise((resolve, reject) => server.close((error) => error ? reject(error) : resolve()));
}
});
test("live Kafka SSE rejects foreign or inconsistent ownership scopes before fanout subscription", async () => {
const ownedSession = { id: "ses_live_kafka_owned", ownerUserId: ACTOR.id, lastTraceId: "trc_live_kafka_owned" };
const foreignSession = { id: "ses_live_kafka_foreign", ownerUserId: "usr_live_kafka_foreign", lastTraceId: "trc_live_kafka_foreign" };
let subscriptions = 0;
let projectionReads = 0;
const server = createCloudApiServer({
accessController: realtimeAccessController({ sessions: [ownedSession, foreignSession] }),
workbenchRuntime: new Proxy({}, {
get() {
projectionReads += 1;
throw new Error("live SSE authorization must not access Workbench projection state");
}
}),
kafkaEventBridge: {
capabilities: { liveKafkaSse: true },
ready: Promise.resolve(),
subscribeLiveHwlabEvents() {
subscriptions += 1;
return () => {};
},
async stop() {}
},
env: {
HWLAB_KAFKA_DIRECT_PUBLISH_ENABLED: "true",
HWLAB_WORKBENCH_LIVE_KAFKA_SSE_ENABLED: "true",
HWLAB_KAFKA_TRANSACTIONAL_PROJECTOR_ENABLED: "false",
HWLAB_KAFKA_PROJECTION_OUTBOX_RELAY_ENABLED: "false",
HWLAB_WORKBENCH_PROJECTION_REALTIME_ENABLED: "false"
}
});
await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve));
try {
const { port } = server.address();
const foreignSessionResponse = await getJson(port, `/v1/workbench/events?sessionId=${foreignSession.id}`);
const foreignTraceResponse = await getJson(port, `/v1/workbench/events?traceId=${foreignSession.lastTraceId}`);
const inconsistentResponse = await getJson(port, `/v1/workbench/events?sessionId=${ownedSession.id}&traceId=${foreignSession.lastTraceId}`);
const missingResponse = await getJson(port, "/v1/workbench/events?sessionId=ses_live_kafka_missing");
for (const response of [foreignSessionResponse, foreignTraceResponse, inconsistentResponse, missingResponse]) {
assert.equal(response.status, 404);
assert.equal(response.body.error.code, "workbench_realtime_scope_not_found");
}
assert.equal(subscriptions, 0);
assert.equal(projectionReads, 0);
} finally {
await new Promise((resolve, reject) => server.close((error) => error ? reject(error) : resolve()));
}
});
test("live Kafka SSE fails closed when ownership lookup is not configured", async () => {
let subscriptions = 0;
const server = createCloudApiServer({
accessController: {
async ensureBootstrap() {},
async authenticate() { return { ok: true, actor: ACTOR, session: { id: "uss_realtime_unconfigured" } }; }
},
kafkaEventBridge: {
capabilities: { liveKafkaSse: true },
ready: Promise.resolve(),
subscribeLiveHwlabEvents() {
subscriptions += 1;
return () => {};
},
async stop() {}
},
env: {
HWLAB_KAFKA_DIRECT_PUBLISH_ENABLED: "true",
HWLAB_WORKBENCH_LIVE_KAFKA_SSE_ENABLED: "true",
HWLAB_KAFKA_TRANSACTIONAL_PROJECTOR_ENABLED: "false",
HWLAB_KAFKA_PROJECTION_OUTBOX_RELAY_ENABLED: "false",
HWLAB_WORKBENCH_PROJECTION_REALTIME_ENABLED: "false"
}
});
await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve));
try {
const response = await getJson(server.address().port, "/v1/workbench/events?sessionId=ses_live_kafka_unconfigured");
assert.equal(response.status, 503);
assert.equal(response.body.error.code, "workbench_realtime_authorization_unconfigured");
assert.equal(subscriptions, 0);
} finally {
await new Promise((resolve, reject) => server.close((error) => error ? reject(error) : resolve()));
}
});
test("live Kafka SSE lets admin subscribe to another owner's consistent session and trace", async () => {
const session = { id: "ses_live_kafka_admin", ownerUserId: "usr_live_kafka_owner", lastTraceId: "trc_live_kafka_admin" };
const admin = { ...ACTOR, id: "usr_live_kafka_admin", role: "admin" };
let subscriptions = 0;
const server = createCloudApiServer({
accessController: realtimeAccessController({ actor: admin, sessions: [session] }),
kafkaEventBridge: {
capabilities: { liveKafkaSse: true },
ready: Promise.resolve(),
subscribeLiveHwlabEvents() {
subscriptions += 1;
return () => {};
},
async stop() {}
},
env: {
HWLAB_KAFKA_DIRECT_PUBLISH_ENABLED: "true",
HWLAB_WORKBENCH_LIVE_KAFKA_SSE_ENABLED: "true",
HWLAB_KAFKA_TRANSACTIONAL_PROJECTOR_ENABLED: "false",
HWLAB_KAFKA_PROJECTION_OUTBOX_RELAY_ENABLED: "false",
HWLAB_WORKBENCH_PROJECTION_REALTIME_ENABLED: "false"
}
});
await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve));
try {
const events = await getSseEvents(server.address().port, `/v1/workbench/events?sessionId=${session.id}&traceId=${session.lastTraceId}`, 1);
assert.equal(events[0].event, "workbench.connected");
assert.deepEqual(events[0].data.filters, { sessionId: session.id, traceId: session.lastTraceId });
assert.equal(subscriptions, 1);
} finally {
await new Promise((resolve, reject) => server.close((error) => error ? reject(error) : resolve()));
}
});
test("workbench realtime initial connection emits current snapshot without replaying historical outbox", async () => {
const sessionId = "ses_realtime_snapshot_only";
const traceId = "trc_realtime_snapshot_only";
const calls = [];
const runtime = {
async readAtomicWorkbenchProjectionSync(params = {}) {
calls.push({ ...params });
return {
facts: {
sessions: [{ sessionId, lastTraceId: traceId }],
messages: [{ messageId: "msg_snapshot_only", sessionId, traceId, role: "agent", text: "current", projectedSeq: 9 }],
turns: [{ turnId: traceId, sessionId, traceId, status: "running", projectedSeq: 9 }]
},
events: [{ outboxSeq: 1, entityFamily: "messages", entityId: "msg_old", payload: { family: "messages", fact: { messageId: "msg_old", sessionId, traceId, text: "historical" } } }],
cutoffOutboxSeq: 90,
cursorOutboxSeq: 90,
hasMore: true
};
}
};
const server = createCloudApiServer({
accessController: realtimeAccessController(),
workbenchRuntime: runtime,
kafkaEventBridge: projectionRealtimeBridge(),
env: { ...TRANSACTIONAL_REALTIME_ENV, HWLAB_WORKBENCH_SSE_HEARTBEAT_MS: "60000", HWLAB_WORKBENCH_OUTBOX_TAIL_BATCH_SIZE: "100" }
});
await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve));
try {
const events = await getSseEvents(server.address().port, `/v1/workbench/events?sessionId=${sessionId}`, 3);
assert.deepEqual(events.map((event) => event.event), ["workbench.connected", "workbench.message.snapshot", "workbench.turn.snapshot"]);
assert.equal(JSON.stringify(events).includes("historical"), false);
assert.equal(events[1].id, "90");
assert.equal(calls.length, 1);
assert.equal(calls[0].snapshotOnly, true);
assert.equal(calls[0].deltaOnly, false);
} finally {
await new Promise((resolve, reject) => server.close((error) => error ? reject(error) : resolve()));
}
});
test("live and projection realtime remain independently reachable when both capabilities are enabled", async () => {
const sessionId = "ses_composable_realtime";
const traceId = "trc_composable_realtime";
let projectionReads = 0;
const runtime = {
async readAtomicWorkbenchProjectionSync() {
projectionReads += 1;
return {
facts: {
sessions: [{ sessionId, lastTraceId: traceId }],
messages: [{ messageId: "msg_composable_realtime", sessionId, traceId, role: "agent", text: "projection snapshot", projectedSeq: 1 }],
turns: []
},
events: [],
cutoffOutboxSeq: 1,
cursorOutboxSeq: 1,
hasMore: false
};
}
};
const server = createCloudApiServer({
accessController: realtimeAccessController(),
workbenchRuntime: runtime,
kafkaEventBridge: projectionRealtimeBridge({ directPublish: true, liveKafkaSse: true }),
env: {
HWLAB_KAFKA_DIRECT_PUBLISH_ENABLED: "true",
HWLAB_WORKBENCH_LIVE_KAFKA_SSE_ENABLED: "true",
HWLAB_KAFKA_TRANSACTIONAL_PROJECTOR_ENABLED: "false",
HWLAB_KAFKA_PROJECTION_OUTBOX_RELAY_ENABLED: "false",
HWLAB_WORKBENCH_PROJECTION_REALTIME_ENABLED: "true",
HWLAB_WORKBENCH_SSE_HEARTBEAT_MS: "60000",
HWLAB_WORKBENCH_OUTBOX_TAIL_BATCH_SIZE: "100"
}
});
await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve));
try {
const events = await getSseEvents(server.address().port, `/v1/workbench/projection-events?sessionId=${sessionId}`, 2);
assert.deepEqual(events.map((event) => event.event), ["workbench.connected", "workbench.message.snapshot"]);
assert.equal(events[0].data.realtimeSource, "projection-outbox");
assert.equal(projectionReads, 1);
} finally {
await new Promise((resolve, reject) => server.close((error) => error ? reject(error) : resolve()));
}
});
test("synchronous projection notification is buffered until the unique initial snapshot completes", async () => {
const sessionId = "ses_realtime_initial_barrier";
const traceId = "trc_realtime_initial_barrier";
const calls = [];
let releaseInitial;
let unsubscribeCount = 0;
const initialGate = new Promise((resolve) => { releaseInitial = resolve; });
const kafkaEventBridge = {
subscribeProjectionCommits(listener) {
listener({ sessionId, traceId });
return () => { unsubscribeCount += 1; };
},
async stop() {}
};
const runtime = {
async readAtomicWorkbenchProjectionSync(params = {}) {
calls.push({ ...params });
if (calls.length === 1) {
await initialGate;
return {
facts: {
sessions: [{ sessionId, lastTraceId: traceId }],
messages: [{ messageId: "msg_initial_barrier", sessionId, traceId, role: "agent", text: "snapshot before delta", projectedSeq: 10 }],
turns: []
},
events: [],
cutoffOutboxSeq: 10,
cursorOutboxSeq: 10,
hasMore: false
};
}
const event = { id: "wte_initial_barrier_delta", sessionId, traceId, projectedSeq: 11, message: "delta after snapshot" };
return {
facts: {},
events: [{ outboxSeq: 11, entityFamily: "traceEvents", entityId: event.id, projectedSeq: 11, projectionRevision: 11, sessionId, traceId, commitType: "event", payload: { family: "traceEvents", fact: event } }],
cutoffOutboxSeq: 11,
cursorOutboxSeq: 11,
hasMore: false
};
}
};
const server = createCloudApiServer({
accessController: realtimeAccessController(),
workbenchRuntime: runtime,
kafkaEventBridge,
env: { ...TRANSACTIONAL_REALTIME_ENV, HWLAB_WORKBENCH_SSE_HEARTBEAT_MS: "60000", HWLAB_WORKBENCH_OUTBOX_TAIL_BATCH_SIZE: "100" }
});
await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve));
try {
const eventsPromise = getSseEvents(server.address().port, `/v1/workbench/events?sessionId=${sessionId}`, 3);
await waitForCondition(() => calls.length === 1);
assert.equal(calls[0].snapshotOnly, true);
assert.equal(calls[0].deltaOnly, false);
assert.equal(calls.length, 1);
releaseInitial();
const events = await eventsPromise;
assert.deepEqual(events.map((event) => event.event), ["workbench.connected", "workbench.message.snapshot", "workbench.trace.event"]);
assert.equal(events[1].data.message.text, "snapshot before delta");
assert.equal(events[2].data.event.message, "delta after snapshot");
assert.equal(calls.length, 2);
assert.equal(calls[1].snapshotOnly, false);
assert.equal(calls[1].deltaOnly, true);
assert.equal(calls[1].afterOutboxSeq, 10);
} finally {
releaseInitial?.();
await new Promise((resolve, reject) => server.close((error) => error ? reject(error) : resolve()));
}
assert.equal(unsubscribeCount, 1);
});
test("initial atomic snapshot failure emits an error and closes SSE for reconnect", async () => {
const sessionId = "ses_realtime_initial_failure";
let unsubscribeCount = 0;
const server = createCloudApiServer({
accessController: realtimeAccessController(),
workbenchRuntime: {
async readAtomicWorkbenchProjectionSync() {
const error = new Error("initial snapshot failed");
error.code = "workbench_initial_snapshot_failed";
throw error;
}
},
kafkaEventBridge: {
subscribeProjectionCommits() { return () => { unsubscribeCount += 1; }; },
async stop() {}
},
env: { ...TRANSACTIONAL_REALTIME_ENV, HWLAB_WORKBENCH_SSE_HEARTBEAT_MS: "60000", HWLAB_WORKBENCH_OUTBOX_TAIL_BATCH_SIZE: "100" }
});
await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve));
const controller = new AbortController();
let timeout = null;
try {
const response = await fetch(`http://127.0.0.1:${server.address().port}/v1/workbench/events?sessionId=${sessionId}`, { signal: controller.signal });
assert.equal(response.status, 200);
const body = await Promise.race([
response.text(),
new Promise((_, reject) => { timeout = setTimeout(() => reject(new Error("initial snapshot failure left SSE open")), 500); })
]);
const events = body.trim().split("\n\n").filter(Boolean).map(parseSseBlock);
assert.deepEqual(events.map((event) => event.event), ["workbench.connected", "workbench.error"]);
assert.equal(events[1].data.error.code, "workbench_initial_snapshot_failed");
assert.equal(body.includes("workbench.heartbeat"), false);
await waitForCondition(() => unsubscribeCount === 1);
} finally {
if (timeout) clearTimeout(timeout);
controller.abort();
await new Promise((resolve, reject) => server.close((error) => error ? reject(error) : resolve()));
}
});
test("delta scan failure after a successful snapshot closes SSE for cursor reconnect", async () => {
const sessionId = "ses_realtime_delta_failure";
const traceId = "trc_realtime_delta_failure";
let unsubscribeCount = 0;
let calls = 0;
const server = createCloudApiServer({
accessController: realtimeAccessController(),
workbenchRuntime: {
async readAtomicWorkbenchProjectionSync() {
calls += 1;
if (calls === 1) {
return {
facts: {
sessions: [{ sessionId, lastTraceId: traceId }],
messages: [{ messageId: "msg_delta_failure", sessionId, traceId, role: "agent", text: "durable snapshot", projectedSeq: 5 }],
turns: []
},
events: [],
cutoffOutboxSeq: 5,
cursorOutboxSeq: 5,
hasMore: false
};
}
const error = new Error("delta scan failed");
error.code = "workbench_delta_scan_failed";
throw error;
}
},
kafkaEventBridge: {
subscribeProjectionCommits(listener) {
listener({ sessionId, traceId });
return () => { unsubscribeCount += 1; };
},
async stop() {}
},
env: { ...TRANSACTIONAL_REALTIME_ENV, HWLAB_WORKBENCH_SSE_HEARTBEAT_MS: "60000", HWLAB_WORKBENCH_OUTBOX_TAIL_BATCH_SIZE: "100" }
});
await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve));
const controller = new AbortController();
let timeout = null;
try {
const response = await fetch(`http://127.0.0.1:${server.address().port}/v1/workbench/events?sessionId=${sessionId}`, { signal: controller.signal });
assert.equal(response.status, 200);
const body = await Promise.race([
response.text(),
new Promise((_, reject) => { timeout = setTimeout(() => reject(new Error("delta failure left SSE open")), 500); })
]);
const events = body.trim().split("\n\n").filter(Boolean).map(parseSseBlock);
assert.deepEqual(events.map((event) => event.event), ["workbench.connected", "workbench.message.snapshot", "workbench.error"]);
assert.equal(events[2].data.error.code, "workbench_delta_scan_failed");
assert.equal(body.includes("workbench.heartbeat"), false);
assert.equal(calls, 2);
await waitForCondition(() => unsubscribeCount === 1);
} finally {
if (timeout) clearTimeout(timeout);
controller.abort();
await new Promise((resolve, reject) => server.close((error) => error ? reject(error) : resolve()));
}
});
test("projector commit during initial snapshot wakes one coalesced delta scan without idle polling", async () => {
const sessionId = "ses_realtime_commit_wakeup";
const traceId = "trc_realtime_commit_wakeup";
const calls = [];
let releaseInitial;
let notify = null;
let unsubscribeCount = 0;
const initialGate = new Promise((resolve) => { releaseInitial = resolve; });
const kafkaEventBridge = {
subscribeProjectionCommits(listener) { notify = listener; return () => { unsubscribeCount += 1; }; },
async stop() {}
};
const runtime = {
async readAtomicWorkbenchProjectionSync(params = {}) {
calls.push({ ...params });
if (calls.length === 1) {
await initialGate;
return { facts: { sessions: [{ sessionId, lastTraceId: traceId }], messages: [], turns: [] }, events: [], cutoffOutboxSeq: 10, cursorOutboxSeq: 10, hasMore: false };
}
const event = { id: "wte_commit_wakeup", sessionId, traceId, projectedSeq: 11, message: "commit wakeup" };
return { facts: {}, events: [{ outboxSeq: 11, entityFamily: "traceEvents", entityId: event.id, projectedSeq: 11, projectionRevision: 11, sessionId, traceId, commitType: "event", payload: { family: "traceEvents", fact: event } }], cutoffOutboxSeq: 11, cursorOutboxSeq: 11, hasMore: false };
}
};
const server = createCloudApiServer({
accessController: realtimeAccessController(),
workbenchRuntime: runtime,
kafkaEventBridge,
env: { ...TRANSACTIONAL_REALTIME_ENV, HWLAB_WORKBENCH_SSE_HEARTBEAT_MS: "60000", HWLAB_WORKBENCH_OUTBOX_TAIL_BATCH_SIZE: "100" }
});
await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve));
try {
const eventsPromise = getSseEvents(server.address().port, `/v1/workbench/events?sessionId=${sessionId}`, 2);
await waitForCondition(() => calls.length === 1 && typeof notify === "function");
notify({ sessionId, traceId });
notify({ sessionId, traceId });
releaseInitial();
const events = await eventsPromise;
assert.deepEqual(events.map((event) => event.event), ["workbench.connected", "workbench.trace.event"]);
assert.equal(events[1].data.event.message, "commit wakeup");
assert.equal(calls.length, 2);
assert.equal(calls[0].snapshotOnly, true);
assert.equal(calls[1].deltaOnly, true);
await new Promise((resolve) => setTimeout(resolve, 30));
assert.equal(calls.length, 2);
} finally {
releaseInitial?.();
await new Promise((resolve, reject) => server.close((error) => error ? reject(error) : resolve()));
}
assert.equal(unsubscribeCount, 1);
});
test("workbench realtime early disconnect releases commit subscription during initial scan", async () => {
const sessionId = "ses_realtime_early_disconnect";
let releaseInitial;
let subscribed = 0;
let unsubscribed = 0;
const gate = new Promise((resolve) => { releaseInitial = resolve; });
const server = createCloudApiServer({
accessController: realtimeAccessController(),
workbenchRuntime: { async readAtomicWorkbenchProjectionSync() { await gate; return { facts: {}, events: [], cutoffOutboxSeq: 0, cursorOutboxSeq: 0, hasMore: false }; } },
kafkaEventBridge: { subscribeProjectionCommits() { subscribed += 1; return () => { unsubscribed += 1; }; }, async stop() {} },
env: { ...TRANSACTIONAL_REALTIME_ENV, HWLAB_WORKBENCH_SSE_HEARTBEAT_MS: "60000", HWLAB_WORKBENCH_OUTBOX_TAIL_BATCH_SIZE: "100" }
});
await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve));
let clientRequest = null;
try {
const response = await new Promise((resolve, reject) => {
clientRequest = httpGet(`http://127.0.0.1:${server.address().port}/v1/workbench/events?sessionId=${sessionId}`, resolve);
clientRequest.once("error", reject);
});
assert.equal(response.statusCode, 200);
await waitForCondition(() => subscribed === 1);
response.destroy();
clientRequest.destroy();
await waitForCondition(() => unsubscribed === 1);
releaseInitial();
} finally {
clientRequest?.destroy();
releaseInitial?.();
await new Promise((resolve, reject) => server.close((error) => error ? reject(error) : resolve()));
}
});
test("workbench realtime prefers automatic Last-Event-ID over stale URL cursor", () => {
const url = new URL("http://localhost/v1/workbench/events?afterSeq=10&afterOutboxSeq=11");
assert.equal(workbenchRealtimeAfterSeq({ headers: { "last-event-id": "42" } }, url), 42);
});
test("workbench realtime resets a cursor ahead of scoped cutoff with one fresh snapshot", async () => {
const sessionId = "ses_realtime_cursor_reset";
const traceId = "trc_realtime_cursor_reset";
const calls = [];
const runtime = {
async readAtomicWorkbenchProjectionSync(params = {}) {
calls.push({ ...params });
if (params.snapshotOnly === true) {
return { facts: { sessions: [{ sessionId, lastTraceId: traceId }], messages: [{ messageId: "msg_cursor_reset", sessionId, traceId, text: "fresh", projectedSeq: 5 }], turns: [] }, events: [], cutoffOutboxSeq: 5, cursorOutboxSeq: 5, hasMore: false };
}
return { facts: {}, events: [], cutoffOutboxSeq: 5, cursorOutboxSeq: 5, hasMore: false };
}
};
const server = createCloudApiServer({ accessController: realtimeAccessController(), workbenchRuntime: runtime, kafkaEventBridge: projectionRealtimeBridge(), env: { ...TRANSACTIONAL_REALTIME_ENV, HWLAB_WORKBENCH_SSE_HEARTBEAT_MS: "60000", HWLAB_WORKBENCH_OUTBOX_TAIL_BATCH_SIZE: "100" } });
await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve));
try {
const events = await getSseEvents(server.address().port, `/v1/workbench/events?sessionId=${sessionId}&afterSeq=42`, 2);
assert.deepEqual(events.map((event) => event.event), ["workbench.connected", "workbench.message.snapshot"]);
assert.equal(events[1].id, "5");
assert.deepEqual(calls.map((call) => [call.deltaOnly, call.snapshotOnly]), [[true, false], [false, true]]);
} finally {
await new Promise((resolve, reject) => server.close((error) => error ? reject(error) : resolve()));
}
});
test("workbench realtime stream surfaces facts blocker instead of legacy trace fallback", async () => {
const traceStore = createCodeAgentTraceStore();
const results = createCodeAgentChatResultStore();
const traceId = "trc_workbench_realtime";
const session = {
id: "ses_workbench_realtime",
projectId: "prj_hwpod_workbench",
agentId: "hwlab-code-agent",
status: "running",
ownerUserId: ACTOR.id,
conversationId: "cnv_workbench_realtime",
threadId: "thread-workbench-realtime",
lastTraceId: traceId,
updatedAt: "2026-06-17T02:00:00.000Z",
session: { sessionStatus: "running", lastTraceId: traceId, messages: [{ role: "user", text: "stream", traceId }] }
};
traceStore.append(traceId, { type: "assistant", status: "completed", label: "assistant:completed", terminal: true, message: "legacy trace answer" });
results.set(traceId, { status: "completed", traceId, ownerUserId: ACTOR.id, sessionId: session.id, threadId: session.threadId, finalResponse: { text: "legacy result answer" } });
const accessController = {
store: {
async getAgentSession(sessionId) { return sessionId === session.id ? session : null; },
async getAgentSessionByTraceId(requestTraceId) { return requestTraceId === traceId ? session : null; }
},
async ensureBootstrap() {},
async authenticate() { return { ok: true, actor: ACTOR, session: { id: "uss_workbench_reader" } }; }
};
const workbenchRuntime = {
async readAtomicWorkbenchProjectionSync() {
const error = new Error("atomic projection unavailable");
error.code = "TEST_ATOMIC_PROJECTION_UNAVAILABLE";
throw error;
}
};
const server = createCloudApiServer({
accessController,
traceStore,
codeAgentChatResults: results,
workbenchRuntime,
kafkaEventBridge: projectionRealtimeBridge(),
env: {
...TRANSACTIONAL_REALTIME_ENV,
HWLAB_WORKBENCH_SSE_HEARTBEAT_MS: "60000",
HWLAB_WORKBENCH_OUTBOX_TAIL_BATCH_SIZE: "100"
}
});
await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve));
try {
const { port } = server.address();
const eventsPromise = getSseEvents(port, `/v1/workbench/events?sessionId=${encodeURIComponent(session.id)}&traceId=${encodeURIComponent(traceId)}`, 2);
const events = await eventsPromise;
assert.deepEqual(events.map((event) => event.event), [
"workbench.connected",
"workbench.error"
]);
assert.equal(events[0].data.filters.sessionId, session.id);
assert.equal(events[1].data.error.code, "TEST_ATOMIC_PROJECTION_UNAVAILABLE");
assert.equal(JSON.stringify(events).includes("legacy result answer"), false);
assert.equal(JSON.stringify(events).includes("legacy trace answer"), false);
} finally {
await new Promise((resolve, reject) => server.close((error) => error ? reject(error) : resolve()));
}
});
test("workbench session realtime follows durable outbox beyond stale lastTraceId", async () => {
const staleTraceId = "trc_workbench_realtime_stale";
const traceId = "trc_workbench_realtime_after_seq";
const session = {
id: "ses_workbench_realtime_after_seq",
projectId: "prj_hwpod_workbench",
agentId: "hwlab-code-agent",
status: "running",
ownerUserId: ACTOR.id,
conversationId: "cnv_workbench_realtime_after_seq",
threadId: "thread-workbench-realtime-after-seq",
lastTraceId: staleTraceId,
updatedAt: "2026-06-24T14:00:00.000Z",
session: { sessionStatus: "running", lastTraceId: staleTraceId }
};
const outboxQueries = [];
const workbenchRuntime = {
async readAtomicWorkbenchProjectionSync(params = {}) {
outboxQueries.push({ ...params });
const after = Number(params.afterOutboxSeq ?? 0);
const outboxSeq = after + 1;
const projectedSeq = outboxSeq === 11 ? 7 : 8;
const event = { id: `wte_workbench_realtime_after_seq_${outboxSeq}`, sourceEventId: `src_workbench_realtime_after_seq_${outboxSeq}`, projectedSeq, sourceSeq: projectedSeq, traceId, sessionId: session.id, turnId: traceId, eventType: "backend", status: "running", message: outboxSeq === 11 ? "durable outbox realtime event" : "second outbox page", terminal: false, sealed: false, updatedAt: "2026-06-24T14:00:01.000Z" };
return {
facts: { sessions: [{ sessionId: session.id, ownerUserId: ACTOR.id, lastTraceId: staleTraceId, status: "running" }], messages: [], parts: [], turns: [], traceEvents: [event], checkpoints: [] },
events: [{ outboxSeq, outboxEventId: `outbox-workbench-realtime-after-seq-${outboxSeq}`, entityFamily: "traceEvents", entityId: event.id, projectionRevision: projectedSeq, projectedSeq, traceId, sessionId: session.id, turnId: traceId, commitType: "event", terminal: false, sealed: false, payload: { family: "traceEvents", fact: event }, createdAt: event.updatedAt }],
cutoffOutboxSeq: 12,
cursorOutboxSeq: outboxSeq,
hasMore: outboxSeq < 12
};
}
};
const accessController = {
store: {
async getAgentSession(sessionId) { return sessionId === session.id ? session : null; },
async getAgentSessionByTraceId(requestTraceId) { return requestTraceId === traceId ? session : null; }
},
async ensureBootstrap() {},
async authenticate() { return { ok: true, actor: ACTOR, session: { id: "uss_workbench_reader" } }; }
};
const serverWithKafka = createCloudApiServer({
accessController,
workbenchRuntime,
kafkaEventBridge: projectionRealtimeBridge(),
env: {
...TRANSACTIONAL_REALTIME_ENV,
HWLAB_WORKBENCH_SSE_HEARTBEAT_MS: "60000",
HWLAB_WORKBENCH_OUTBOX_TAIL_BATCH_SIZE: "100"
}
});
await new Promise((resolve) => serverWithKafka.listen(0, "127.0.0.1", resolve));
try {
const { port } = serverWithKafka.address();
const events = await getSseEvents(port, `/v1/workbench/events?sessionId=${encodeURIComponent(session.id)}&afterSeq=10`, 3);
assert.deepEqual(events.map((event) => event.event), ["workbench.connected", "workbench.trace.event", "workbench.trace.event"]);
assert.equal(events[0].id, "10");
assert.equal(events[1].id, "11");
assert.equal(events[1].data.realtimeSource, "projection-outbox");
assert.equal(events[1].data.realtimeAuthority, "workbench-realtime-authority-v2");
assert.equal(events[1].data.entity.family, "traceEvents");
assert.equal(events[1].data.entity.version, 7);
assert.equal(events[1].data.entity.projectionRevision, "7");
assert.equal(events[1].data.event.message, "durable outbox realtime event");
assert.equal(events[1].data.traceId, traceId);
assert.equal(events[1].data.cursor.traceSeq, 7);
assert.equal(events[2].id, "12");
assert.equal(events[2].data.event.message, "second outbox page");
assert.equal(outboxQueries.length, 2);
assert.deepEqual(outboxQueries.map((query) => query.afterOutboxSeq), [10, 11]);
assert.equal(outboxQueries[0].sessionId, session.id);
assert.deepEqual(outboxQueries[0].actor, { id: ACTOR.id, role: ACTOR.role });
} finally {
await new Promise((resolve, reject) => serverWithKafka.close((error) => error ? reject(error) : resolve()));
}
});
test("workbench read model exposes runtime trace projection query failures as projection blockers", async () => {
const traceStore = createCodeAgentTraceStore();
const results = createCodeAgentChatResultStore();
const traceId = "trc_workbench_projection_store_unavailable";
const session = {
id: "ses_workbench_projection_store_unavailable",
projectId: "prj_hwpod_workbench",
agentId: "hwlab-code-agent",
status: "running",
ownerUserId: ACTOR.id,
conversationId: "cnv_workbench_projection_store_unavailable",
threadId: "thread-workbench-projection-store-unavailable",
lastTraceId: traceId,
updatedAt: "2026-06-18T01:10:00.000Z",
session: {
sessionStatus: "running",
lastTraceId: traceId,
messages: [
{ role: "user", text: "projection store unavailable", traceId, status: "sent", createdAt: "2026-06-18T01:09:58.000Z" },
{ role: "agent", text: "", traceId, status: "running", createdAt: "2026-06-18T01:09:59.000Z" }
],
valuesRedacted: true,
secretMaterialStored: false
}
};
traceStore.append(traceId, { seq: 1, type: "backend", status: "running", label: "runner:created", terminal: false, createdAt: "2026-06-18T01:10:00.000Z" });
results.set(traceId, { status: "running", traceId, ownerUserId: ACTOR.id, sessionId: session.id, threadId: session.threadId });
const accessController = {
store: {
async getAgentSession(sessionId) { return sessionId === session.id ? session : null; },
async getAgentSessionByTraceId(requestTraceId) { return requestTraceId === traceId ? session : null; }
},
async ensureBootstrap() {},
async authenticate() { return { ok: true, actor: ACTOR, session: { id: "uss_workbench_reader" } }; }
};
const runtimeStore = {
async queryWorkbenchFacts() {
const error = new Error("runtime query failed");
error.code = "TEST_RUNTIME_QUERY_FAILED";
throw error;
}
};
const server = createCloudApiServer({ accessController, traceStore, workbenchRuntime: runtimeStore, codeAgentChatResults: results });
await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve));
try {
const { port } = server.address();
const turn = await getJson(port, `/v1/workbench/turns/${encodeURIComponent(traceId)}`);
assert.equal(turn.status, 503);
assert.equal(turn.body.error.code, "projection_store_unavailable");
const trace = await getJson(port, `/v1/workbench/traces/${encodeURIComponent(traceId)}/events?limit=10`);
assert.equal(trace.status, 503);
assert.equal(trace.body.error.code, "projection_store_unavailable");
} finally {
await new Promise((resolve, reject) => server.close((error) => error ? reject(error) : resolve()));
}
});
test("workbench trace events remain visible after session lastTraceId moves", async () => {
const oldTraceId = "trc_workbench_history_old_trace";
const newTraceId = "trc_workbench_history_new_trace";
const session = {
id: "ses_workbench_history_trace",
projectId: "prj_hwpod_workbench",
agentId: "hwlab-code-agent",
status: "running",
ownerUserId: ACTOR.id,
conversationId: "cnv_workbench_history_trace",
threadId: "thread-workbench-history-trace",
lastTraceId: newTraceId,
updatedAt: "2026-06-20T10:10:00.000Z",
session: {
messages: [
{ role: "user", text: "old turn", traceId: oldTraceId, status: "sent", createdAt: "2026-06-20T10:00:00.000Z" },
{ role: "agent", text: "old final", traceId: oldTraceId, status: "completed", createdAt: "2026-06-20T10:00:01.000Z" }
],
valuesRedacted: true
}
};
const oldFacts = buildDurableFactsForSession({
session: { ...session, status: "completed", lastTraceId: oldTraceId, updatedAt: "2026-06-20T10:01:00.000Z" },
status: "completed",
finalText: "old final",
runId: "run_workbench_history_old_trace",
commandId: "cmd_workbench_history_old_trace",
events: [
{ projectedSeq: 1, sourceSeq: 1, type: "backend", status: "running", label: "runner:started", createdAt: "2026-06-20T10:00:00.000Z" },
{ projectedSeq: 2, sourceSeq: 2, type: "result", status: "completed", label: "runner:completed", terminal: true, createdAt: "2026-06-20T10:01:00.000Z" }
],
lastProjectedSeq: 2
});
const newFacts = buildDurableFactsForSession({
session: { ...session, session: { messages: [{ role: "user", text: "new turn", traceId: newTraceId, status: "sent" }] } },
status: "running",
runId: "run_workbench_history_new_trace",
commandId: "cmd_workbench_history_new_trace",
events: [{ projectedSeq: 1, sourceSeq: 1, type: "backend", status: "running", label: "runner:started", createdAt: "2026-06-20T10:10:00.000Z" }],
lastProjectedSeq: 1
});
const facts = emptyFacts();
mergeFacts(facts, oldFacts);
mergeFacts(facts, newFacts);
facts.sessions = newFacts.sessions;
const queries = [];
const runtimeStore = {
async queryWorkbenchFacts(params = {}) {
queries.push(params);
const filtered = filterFacts(facts, params);
return {
facts: filtered,
count: Object.values(filtered).reduce((sum, rows) => sum + rows.length, 0),
persistence: { adapter: "test-durable-workbench-facts", durable: true }
};
}
};
const accessController = {
async ensureBootstrap() {},
async authenticate() { return { ok: true, actor: ACTOR, session: { id: "uss_workbench_reader" } }; }
};
const server = createCloudApiServer({ accessController, workbenchRuntime: runtimeStore });
await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve));
try {
const { port } = server.address();
const response = await getJson(port, `/v1/workbench/traces/${encodeURIComponent(oldTraceId)}/events?limit=10`);
assert.equal(response.status, 200);
assert.equal(response.body.status, "succeeded");
assert.equal(response.body.sessionId, session.id);
assert.equal(response.body.traceId, oldTraceId);
assert.equal(response.body.events.length, 2);
assert.equal(response.body.events.at(-1).status, "completed");
assert.equal(response.body.error, undefined);
assert.equal(queries.some((query) => query.sessionId === session.id && Array.isArray(query.families) && query.families.includes("sessions")), true);
} finally {
await new Promise((resolve, reject) => server.close((error) => error ? reject(error) : resolve()));
}
});
test("workbench trace event page keeps per-trace terminal status after later session cancel", async () => {
const completedTraceId = "trc_workbench_completed_before_cancel";
const canceledTraceId = "trc_workbench_later_cancel";
const session = {
id: "ses_workbench_completed_before_cancel",
projectId: "prj_hwpod_workbench",
agentId: "hwlab-code-agent",
status: "canceled",
ownerUserId: ACTOR.id,
conversationId: "cnv_workbench_completed_before_cancel",
threadId: "thread-workbench-completed-before-cancel",
lastTraceId: canceledTraceId,
updatedAt: "2026-06-20T10:20:30.000Z",
session: {
messages: [
{ role: "user", text: "completed first", traceId: completedTraceId, status: "sent", createdAt: "2026-06-20T10:20:00.000Z" },
{ role: "agent", text: "completed final", traceId: completedTraceId, status: "completed", createdAt: "2026-06-20T10:20:10.000Z" },
{ role: "user", text: "cancel later", traceId: canceledTraceId, status: "sent", createdAt: "2026-06-20T10:20:20.000Z" },
{ role: "agent", text: "hwlab-user-cancel", traceId: canceledTraceId, status: "canceled", createdAt: "2026-06-20T10:20:30.000Z" }
],
valuesRedacted: true
}
};
const facts = emptyFacts();
mergeFacts(facts, buildDurableFactsForSession({
session: { ...session, status: "completed", lastTraceId: completedTraceId, updatedAt: "2026-06-20T10:20:10.000Z" },
status: "completed",
finalText: "completed final",
events: [
{ projectedSeq: 1, sourceSeq: 1, type: "backend", status: "running", label: "runner:started", createdAt: "2026-06-20T10:20:00.000Z" },
{ projectedSeq: 2, sourceSeq: 2, type: "result", status: "completed", label: "runner:completed", terminal: true, createdAt: "2026-06-20T10:20:10.000Z" }
],
lastProjectedSeq: 2
}));
mergeFacts(facts, buildDurableFactsForSession({
session: { ...session, status: "canceled", lastTraceId: canceledTraceId, session: { messages: session.session.messages.slice(2), valuesRedacted: true } },
status: "canceled",
finalText: "hwlab-user-cancel",
events: [
{ projectedSeq: 1, sourceSeq: 1, type: "backend", status: "running", label: "runner:started", createdAt: "2026-06-20T10:20:20.000Z" },
{ projectedSeq: 2, sourceSeq: 2, type: "cancel", status: "canceled", label: "runner:canceled", terminal: true, createdAt: "2026-06-20T10:20:30.000Z" }
],
lastProjectedSeq: 2
}));
facts.sessions = facts.sessions.filter((item) => item.lastTraceId === canceledTraceId);
const runtimeStore = {
async queryWorkbenchFacts(params = {}) {
const filtered = filterFacts(facts, params);
return {
facts: filtered,
count: Object.values(filtered).reduce((sum, rows) => sum + rows.length, 0),
persistence: { adapter: "test-durable-workbench-facts", durable: true }
};
}
};
const accessController = {
async ensureBootstrap() {},
async authenticate() { return { ok: true, actor: ACTOR, session: { id: "uss_workbench_reader" } }; }
};
const server = createCloudApiServer({ accessController, workbenchRuntime: runtimeStore });
await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve));
try {
const { port } = server.address();
const response = await getJson(port, `/v1/workbench/traces/${encodeURIComponent(completedTraceId)}/events?limit=10`);
assert.equal(response.status, 200);
assert.equal(response.body.traceStatus, "completed");
assert.equal(response.body.events.at(-1).status, "completed");
} finally {
await new Promise((resolve, reject) => server.close((error) => error ? reject(error) : resolve()));
}
});
function realtimeAccessController({ actor = ACTOR, sessions = [] } = {}) {
const bySessionId = new Map(sessions.map((session) => [session.id, session]));
const byTraceId = new Map(sessions.filter((session) => session.lastTraceId).map((session) => [session.lastTraceId, session]));
return {
async ensureBootstrap() {},
async authenticate() { return { ok: true, actor, session: { id: "uss_realtime_test" } }; },
async getAgentSession(sessionId) { return bySessionId.get(sessionId) ?? null; },
async getAgentSessionByTraceId(traceId) { return byTraceId.get(traceId) ?? null; }
};
}