From d0d2f46628d60e4ccffa60ba426efe34e485b8af Mon Sep 17 00:00:00 2001 From: root Date: Mon, 13 Jul 2026 21:07:53 +0200 Subject: [PATCH] fix: remove workbench sync repair authority --- internal/cloud/kafka-event-bridge.test.ts | 189 ++------ internal/cloud/kafka-event-bridge.ts | 12 +- internal/cloud/server-workbench-read-http.ts | 20 - .../server-workbench-realtime-http.test.ts | 441 +----------------- internal/cloud/server.ts | 4 +- .../workbench-projection-outbox-events.ts | 2 +- ...kbench-realtime-authority-contract.test.ts | 428 ----------------- .../cloud/workbench-realtime-authority.ts | 162 ------- web/hwlab-cloud-web/scripts/check.ts | 4 +- .../workbench-realtime-runtime.test.ts | 23 +- .../src/api/workbench-events.ts | 41 -- .../src/stores/workbench-colada-keys.ts | 8 - .../src/stores/workbench-colada-mutations.ts | 1 - .../src/stores/workbench-colada-queries.ts | 4 - .../src/stores/workbench-colada.test.ts | 92 +--- ...rkbench-message-projection-runtime.test.ts | 17 +- .../stores/workbench-realtime-authority.ts | 172 +------ .../src/stores/workbench-realtime-plan.ts | 12 +- web/hwlab-cloud-web/src/stores/workbench.ts | 56 +-- .../src/utils/workbench-stream-transport.ts | 11 +- 20 files changed, 88 insertions(+), 1611 deletions(-) delete mode 100644 internal/cloud/workbench-realtime-authority-contract.test.ts delete mode 100644 internal/cloud/workbench-realtime-authority.ts diff --git a/internal/cloud/kafka-event-bridge.test.ts b/internal/cloud/kafka-event-bridge.test.ts index 58644c23..89512113 100644 --- a/internal/cloud/kafka-event-bridge.test.ts +++ b/internal/cloud/kafka-event-bridge.test.ts @@ -322,105 +322,6 @@ test("projects AgentRun terminal_status Kafka event into terminal HWLAB event", assert.equal(projected.event.message, "AgentRun completed"); }); -test("live Kafka mode publishes canonical AgentRun input directly and fans one broker event to all SSE subscribers", async () => { - const consumerConfigs = []; - const subscriptions = []; - const runOptions = []; - const runOrder = []; - const produced = []; - const otelSpans = []; - const consumers = [0, 1].map((index) => ({ - async connect() {}, - async subscribe(input) { subscriptions[index] = input; }, - async run(input) { runOptions[index] = input; runOrder.push(index); }, - async stop() {}, - async disconnect() {} - })); - const producer = { - async connect() {}, - async send(input) { produced.push(input); return [{ topicName: input.topic, partition: 4, baseOffset: "91" }]; }, - async disconnect() {} - }; - let consumerIndex = 0; - const kafkaFactory = () => ({ - consumer(config) { - consumerConfigs.push(config); - return consumers[consumerIndex++]; - }, - producer() { return producer; } - }); - const forbiddenRuntimeStore = new Proxy({}, { - get() { throw new Error("live mode must not touch the runtime store"); } - }); - const bridge = startHwlabKafkaEventBridge({ - env: LIVE_ENV, - runtimeStore: forbiddenRuntimeStore, - kafkaFactory, - logger: null, - otelSpanEmitter(name, traceId, env, options) { otelSpans.push({ name, traceId, env, options }); } - }); - await bridge.ready; - - assert.deepEqual(consumerConfigs.map((item) => item.groupId), ["hwlab-live-bridge-fixed", "hwlab-live-fanout-fixed"]); - assert.deepEqual(runOrder, [1, 0], "fanout group must be ready before direct publish can create a new HWLAB event"); - assert.deepEqual(subscriptions, [ - { topic: "agentrun.event.v1", fromBeginning: false }, - { topic: "hwlab.event.v1", fromBeginning: false } - ]); - const message = canonicalMessage({ offset: "70" }); - await runOptions[0].eachMessage({ topic: "agentrun.event.v1", partition: 0, message }); - assert.equal(produced.length, 1); - assert.equal(produced[0].topic, "hwlab.event.v1"); - const envelope = JSON.parse(produced[0].messages[0].value); - assert.equal(envelope.schema, "hwlab.event.v1"); - assert.equal(envelope.traceId, "trc_projector_test"); - assert.equal(envelope.hwlabSessionId, "ses_projector_test"); - assert.equal(envelope.runId, "run_projector_test"); - assert.equal(envelope.commandId, "cmd_projector_test"); - assert.equal(envelope.event.assistantText, "projected answer"); - assert.equal(otelSpans[0].name, "hwlab.kafka.live.direct_publish"); - assert.equal(otelSpans[0].traceId, "trc_projector_test"); - assert.deepEqual(otelSpans[0].options.attributes, { - businessTraceId: "trc_projector_test", - hwlabSessionId: "ses_projector_test", - runId: "run_projector_test", - commandId: "cmd_projector_test", - topic: "hwlab.event.v1", - partition: 4, - offset: "91", - sourceTopic: "agentrun.event.v1", - sourcePartition: 0, - sourceOffset: "70", - eventType: "assistant_message", - terminal: false, - valuesRedacted: true - }); - - const first = []; - const second = []; - const firstTransports = []; - bridge.subscribeLiveHwlabEvents((value, transport) => { first.push(value); firstTransports.push(transport); }); - bridge.subscribeLiveHwlabEvents((value) => second.push(value)); - await runOptions[1].eachBatch({ - batch: { topic: "hwlab.event.v1", partition: 0, highWatermark: "1", messages: [{ offset: "0", value: Buffer.from(JSON.stringify(envelope)) }] }, - resolveOffset() {}, - async heartbeat() {}, - isRunning: () => true, - isStale: () => false - }); - assert.equal(bridge.liveSubscriberCount(), 2); - assert.deepEqual(first, [envelope]); - assert.deepEqual(second, [envelope]); - assert.deepEqual(firstTransports, [{ topic: "hwlab.event.v1", partition: 0, offset: "0" }]); - assert.equal(otelSpans[1].name, "hwlab.kafka.live.fanout_receive"); - assert.equal(otelSpans[1].options.attributes.topic, "hwlab.event.v1"); - assert.equal(otelSpans[1].options.attributes.partition, 0); - assert.equal(otelSpans[1].options.attributes.offset, "0"); - assert.equal(consumerConfigs.length, 2, "browser subscribers must not allocate Kafka consumer groups"); - assert.equal((await bridge.status()).consumerLag.total, 0); - await bridge.stop(); -}); - test("live Kafka publisher has no PG inbox, facts, checkpoint, or outbox dependency", async () => { const sent = []; const runtimeStore = new Proxy({}, { get() { throw new Error("unexpected DB access"); } }); @@ -490,77 +391,41 @@ test("composable Kafka capabilities reject consumer group collisions", () => { ); }); -test("Workbench bridge rejects direct and transactional dual authority", () => { - assert.throws( - () => startHwlabKafkaEventBridge({ +test("Workbench bridge fails closed unless transactional projection authority is complete", () => { + const invalidAuthorities = [ + { name: "direct/live-only", env: LIVE_ENV }, + { + name: "projection realtime only", + env: { + HWLAB_KAFKA_DIRECT_PUBLISH_ENABLED: "false", + HWLAB_WORKBENCH_LIVE_KAFKA_SSE_ENABLED: "false", + HWLAB_WORKBENCH_KAFKA_REFRESH_REPLAY_ENABLED: "false", + HWLAB_KAFKA_TRANSACTIONAL_PROJECTOR_ENABLED: "false", + HWLAB_KAFKA_PROJECTION_OUTBOX_RELAY_ENABLED: "false", + HWLAB_WORKBENCH_PROJECTION_REALTIME_ENABLED: "true" + } + }, + { name: "projector without outbox relay", env: { ...PROJECTOR_ENV, HWLAB_KAFKA_PROJECTION_OUTBOX_RELAY_ENABLED: "false" } }, + { name: "projector without projection realtime", env: { ...PROJECTOR_ENV, HWLAB_WORKBENCH_PROJECTION_REALTIME_ENABLED: "false" } }, + { + name: "direct and transactional dual authority", env: { ...PROJECTOR_ENV, HWLAB_KAFKA_DIRECT_PUBLISH_ENABLED: "true", HWLAB_WORKBENCH_LIVE_KAFKA_SSE_ENABLED: "true", HWLAB_KAFKA_AGENTRUN_EVENT_GROUP_ID: "hwlab-direct-authority-test", HWLAB_KAFKA_HWLAB_EVENT_GROUP_ID: "hwlab-live-authority-test" - }, - runtimeStore: {} - }), - (error: any) => error?.code === "hwlab_workbench_realtime_dual_authority" - ); -}); - -test("Kafka refresh replay requires explicit YAML budgets and composes with live SSE", () => { - const refreshEnv = { - ...LIVE_ENV, - HWLAB_WORKBENCH_KAFKA_REFRESH_REPLAY_ENABLED: "true", - HWLAB_WORKBENCH_KAFKA_REFRESH_REPLAY_GROUP_PREFIX: "hwlab-v03-cloud-api-sse", - HWLAB_WORKBENCH_KAFKA_REFRESH_REPLAY_TIMEOUT_MS: "30000", - HWLAB_WORKBENCH_KAFKA_REFRESH_REPLAY_SCAN_LIMIT: "1000000", - HWLAB_WORKBENCH_KAFKA_REFRESH_REPLAY_MATCHED_EVENT_LIMIT: "2000", - HWLAB_WORKBENCH_KAFKA_REFRESH_REPLAY_LIVE_BUFFER_LIMIT: "2000" - }; - const config = kafkaEventBridgeConfig(refreshEnv); - assert.deepEqual(config.refreshReplay, { - groupIdPrefix: "hwlab-v03-cloud-api-sse", - timeoutMs: 30000, - scanLimit: 1000000, - matchedEventLimit: 2000, - liveBufferLimit: 2000 - }); - assert.throws( - () => kafkaEventBridgeConfig({ ...refreshEnv, HWLAB_WORKBENCH_KAFKA_REFRESH_REPLAY_SCAN_LIMIT: undefined }), - (error: any) => error?.code === "hwlab_kafka_realtime_config_incomplete" && error?.missing?.includes("HWLAB_WORKBENCH_KAFKA_REFRESH_REPLAY_SCAN_LIMIT") - ); - assert.throws( - () => kafkaEventBridgeConfig({ ...refreshEnv, HWLAB_WORKBENCH_LIVE_KAFKA_SSE_ENABLED: "false" }), - (error: any) => error?.code === "hwlab_kafka_refresh_live_dependency_missing" - ); -}); - -test("projection realtime subscribes to projection commits without enabling Kafka ingest or relay", async () => { - let subscribeCount = 0; - let unsubscribeCount = 0; - const bridge = startHwlabKafkaEventBridge({ - env: { - HWLAB_KAFKA_DIRECT_PUBLISH_ENABLED: "false", - HWLAB_WORKBENCH_LIVE_KAFKA_SSE_ENABLED: "false", - HWLAB_WORKBENCH_KAFKA_REFRESH_REPLAY_ENABLED: "false", - HWLAB_KAFKA_TRANSACTIONAL_PROJECTOR_ENABLED: "false", - HWLAB_KAFKA_PROJECTION_OUTBOX_RELAY_ENABLED: "false", - HWLAB_WORKBENCH_PROJECTION_REALTIME_ENABLED: "true" - }, - runtimeStore: { - async assertWorkbenchTransactionalRealtimeReady() { return { ready: true }; }, - async subscribeWorkbenchProjectionCommits() { - subscribeCount += 1; - return () => { unsubscribeCount += 1; }; } - }, - kafkaFactory() { throw new Error("projection realtime must not create a Kafka client"); }, - logger: null - }); + } + ]; - await bridge.ready; - assert.equal(subscribeCount, 1); - await bridge.stop(); - assert.equal(unsubscribeCount, 1); + for (const authority of invalidAuthorities) { + assert.throws( + () => startHwlabKafkaEventBridge({ env: authority.env, runtimeStore: {} }), + (error: any) => error?.code === "hwlab_workbench_transactional_authority_required", + authority.name + ); + } }); test("projects AgentRun run-created envelope even when no source event is present", () => { diff --git a/internal/cloud/kafka-event-bridge.ts b/internal/cloud/kafka-event-bridge.ts index dc2ddf73..225856f0 100644 --- a/internal/cloud/kafka-event-bridge.ts +++ b/internal/cloud/kafka-event-bridge.ts @@ -119,7 +119,7 @@ export function kafkaEventBridgeConfig(env = process.env) { export function startHwlabKafkaEventBridge({ env = process.env, logger = console, kafkaFactory = defaultKafkaFactory, runtimeStore = null, now = () => new Date().toISOString(), otelSpanEmitter = emitCodeAgentOtelSpan } = {}) { const config = kafkaEventBridgeConfig(env); if (!config) return { started: false, reason: "disabled_or_unconfigured", stop() {}, subscribeProjectionCommits() { return () => {}; }, valuesPrinted: false }; - assertSingleWorkbenchRealtimeAuthority(config.capabilities); + assertTransactionalWorkbenchRealtimeAuthority(config.capabilities); const components = []; if (config.capabilities.directPublish || config.capabilities.liveKafkaSse) { components.push(startLiveHwlabKafkaEventBridge({ config, env, logger, kafkaFactory, otelSpanEmitter })); @@ -130,13 +130,13 @@ export function startHwlabKafkaEventBridge({ env = process.env, logger = console return components.length === 1 ? components[0] : combineKafkaEventBridgeComponents(config, components); } -function assertSingleWorkbenchRealtimeAuthority(capabilities = {}) { +function assertTransactionalWorkbenchRealtimeAuthority(capabilities = {}) { const directAuthority = capabilities.directPublish || capabilities.liveKafkaSse || capabilities.kafkaRefreshReplay; - const projectionAuthority = capabilities.transactionalProjector || capabilities.projectionOutboxRelay || capabilities.projectionRealtime; - if (!directAuthority || !projectionAuthority) return; + const completeProjectionAuthority = capabilities.transactionalProjector && capabilities.projectionOutboxRelay && capabilities.projectionRealtime; + if (!directAuthority && completeProjectionAuthority) return; throw contractError( - "hwlab_workbench_realtime_dual_authority", - "Workbench direct/live Kafka authority cannot run beside the transactional projector authority. Disable direct publish, live-only SSE, and Kafka retention replay before enabling the projector chain." + "hwlab_workbench_transactional_authority_required", + "Workbench Kafka bridge requires transactionalProjector, projectionOutboxRelay, and projectionRealtime together, with directPublish, liveKafkaSse, and kafkaRefreshReplay disabled." ); } diff --git a/internal/cloud/server-workbench-read-http.ts b/internal/cloud/server-workbench-read-http.ts index 78fb4c86..fa644aca 100644 --- a/internal/cloud/server-workbench-read-http.ts +++ b/internal/cloud/server-workbench-read-http.ts @@ -16,8 +16,6 @@ import { import { createWorkbenchReadModel } from "./workbench-read-model.ts"; import { createWorkbenchRuntimeClient } from "./workbench-runtime-client.ts"; import { buildWorkbenchSessionDetail, compactLaunchContext, includeMessagesForSessionDetail } from "./workbench-session-detail-response.ts"; -import { handleWorkbenchSyncHttp } from "./workbench-realtime-authority.ts"; -import { workbenchRealtimeCapabilities } from "./workbench-realtime-capabilities.ts"; import { durableTraceStatus, RUNNING_STATUSES, terminalFinalResponse, TERMINAL_STATUSES } from "./workbench-turn-projection.ts"; import { emitCodeAgentOtelSpan, emitHttpServerRequestSpan } from "./otel-trace.ts"; import * as workbenchFacts from "./server-workbench-facts.ts"; @@ -156,24 +154,6 @@ export async function handleWorkbenchReadModelHttp(request, response, url, optio const auth = perf ? await perf.measure("workbench_auth", () => authenticateWorkbenchRead(request, response, options)) : await authenticateWorkbenchRead(request, response, options); if (!auth) return; - if (url.pathname === "/v1/workbench/sync") { - const capabilities = workbenchRealtimeCapabilities(options.env ?? process.env); - if (!capabilities.projectionRealtime) { - sendJson(response, 503, { - ok: false, - error: { - code: "workbench_projection_realtime_disabled", - message: "Workbench projection sync/replay capability is disabled.", - capabilities, - valuesRedacted: true - } - }); - return; - } - await (perf ? perf.measure("workbench_sync", () => handleWorkbenchSyncHttp(request, response, url, options, auth.actor)) : handleWorkbenchSyncHttp(request, response, url, options, auth.actor)); - return; - } - if (url.pathname === "/v1/workbench/sessions") { if (request.method !== "GET") return methodNotAllowed(response, "GET"); await (perf ? perf.measure("workbench_session_list", () => handleWorkbenchSessionList(request, response, url, options, auth.actor)) : handleWorkbenchSessionList(request, response, url, options, auth.actor)); diff --git a/internal/cloud/server-workbench-realtime-http.test.ts b/internal/cloud/server-workbench-realtime-http.test.ts index be4fe52d..e904be6c 100644 --- a/internal/cloud/server-workbench-realtime-http.test.ts +++ b/internal/cloud/server-workbench-realtime-http.test.ts @@ -90,444 +90,17 @@ test("projection outbox emits immutable assistant versions in row order", () => assert.deepEqual(events.map((item) => item.payload.cursor.outboxSeq), [10, 11]); }); -test.skip("removed live-only product authority transparently fans out one envelope", 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, kafkaRefreshReplay: false, 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"); - } - }); +test("workbench sync repair endpoint is removed from the product API", async () => { 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_WORKBENCH_KAFKA_REFRESH_REPLAY_ENABLED: "false", - HWLAB_KAFKA_TRANSACTIONAL_PROJECTOR_ENABLED: "false", - HWLAB_KAFKA_PROJECTION_OUTBOX_RELAY_ENABLED: "false", - HWLAB_WORKBENCH_PROJECTION_REALTIME_ENABLED: "false" - } + accessController: realtimeAccessController(), + workbenchRuntime: {}, + kafkaEventBridge: projectionRealtimeBridge(), + env: { ...TRANSACTIONAL_REALTIME_ENV } }); 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.skip("removed Kafka retention replay product authority replays retained envelopes", async () => { - const sessionId = "ses_kafka_refresh_replay"; - const traceId = "trc_kafka_refresh_replay"; - const records = [ - refreshRecord(0, sessionId, traceId, "user", { userMessageId: "msg_kafka_refresh_user", messageId: "msg_kafka_refresh_user", text: "retained input" }), - refreshRecord(1, sessionId, traceId, "backend"), - refreshRecord(2, sessionId, traceId, "tool"), - refreshRecord(3, sessionId, traceId, "assistant", { assistantText: "retained answer", text: "retained answer" }), - refreshRecord(4, sessionId, traceId, "result", { terminal: true, status: "completed" }) - ]; - let retentionQuery = null; - let projectionRuntimeCalls = 0; - const bridge = refreshReplayBridge({ - records, - endOffset: "5", - onQuery(params) { retentionQuery = params; }, - ready: new Promise(() => undefined), - liveReady: Promise.resolve() - }); - const server = createCloudApiServer({ - accessController: realtimeAccessController({ sessions: [{ id: sessionId, ownerUserId: ACTOR.id, lastTraceId: traceId }] }), - workbenchRuntime: new Proxy({}, { get() { projectionRuntimeCalls += 1; throw new Error("Kafka refresh must not hydrate projection history"); } }), - kafkaEventBridge: bridge, - env: { ...REFRESH_REALTIME_ENV, HWLAB_WORKBENCH_SSE_HEARTBEAT_MS: "60000" } - }); - await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve)); - - try { - const events = await getSseEvents(server.address().port, `/v1/workbench/events?traceId=${traceId}`, 6); - assert.deepEqual(events.map((event) => event.event), [ - "hwlab.event.v1", - "hwlab.event.v1", - "hwlab.event.v1", - "hwlab.event.v1", - "hwlab.event.v1", - "workbench.connected" - ]); - assert.equal(events[0].data.event.type, "user"); - assert.equal(events[0].data.event.userMessageId, "msg_kafka_refresh_user"); - assert.equal(events[4].data.event.terminal, true); - assert.equal(events[5].data.deliverySemantics, "kafka-retention-then-live"); - assert.equal(events[5].data.liveOnly, false); - assert.equal(events[5].data.replay, true); - assert.equal(events[5].data.lossPossible, false); - assert.deepEqual(events[5].data.filters, { sessionId, traceId }); - assert.equal(events[5].data.refreshReplay.counts.replayed, 5); - assert.equal(retentionQuery.sessionId, sessionId); - assert.equal(retentionQuery.traceId, traceId); - assert.equal(retentionQuery.limit, 50); - assert.equal(retentionQuery.scanLimit, 500); - assert.equal(retentionQuery.timeoutMs, 2500); - assert.equal(retentionQuery.groupIdPrefix, "hwlab-v03-refresh-test"); - assert.equal(retentionQuery.fromBeginning, true); - assert.equal(retentionQuery.signal instanceof AbortSignal, true); - assert.equal(projectionRuntimeCalls, 0); - } finally { - await new Promise((resolve, reject) => server.close((error) => error ? reject(error) : resolve())); - } -}); - -test.skip("removed Kafka retention replay product authority handles client abort", async () => { - const sessionId = "ses_kafka_refresh_ready_abort"; - let releaseLiveReady: (() => void) | null = null; - const liveReady = new Promise((resolve) => { releaseLiveReady = resolve; }); - const bridge = refreshReplayBridge({ liveReady }); - let liveReadyObserved = 0; - let queryCalls = 0; - let subscribeCalls = 0; - const originalQuery = bridge.queryHwlabEventRetention.bind(bridge); - const originalSubscribe = bridge.subscribeLiveHwlabEvents.bind(bridge); - Object.defineProperty(bridge, "liveReady", { - configurable: true, - get() { - liveReadyObserved += 1; - return liveReady; - } - }); - bridge.queryHwlabEventRetention = async (params) => { - queryCalls += 1; - return await originalQuery(params); - }; - bridge.subscribeLiveHwlabEvents = (listener) => { - subscribeCalls += 1; - return originalSubscribe(listener); - }; - const server = createCloudApiServer({ - accessController: realtimeAccessController({ sessions: [{ id: sessionId, ownerUserId: ACTOR.id }] }), - kafkaEventBridge: bridge, - env: { ...REFRESH_REALTIME_ENV, HWLAB_WORKBENCH_SSE_HEARTBEAT_MS: "60000" } - }); - await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve)); - - try { - const controller = new AbortController(); - const response = fetch(`http://127.0.0.1:${server.address().port}/v1/workbench/events?sessionId=${sessionId}`, { signal: controller.signal }).catch((error) => error); - await waitForCondition(() => liveReadyObserved > 0); - controller.abort(); - await response; - releaseLiveReady?.(); - await new Promise((resolve) => setTimeout(resolve, 10)); - assert.equal(queryCalls, 0); - assert.equal(subscribeCalls, 0); - } finally { - releaseLiveReady?.(); - await new Promise((resolve, reject) => server.close((error) => error ? reject(error) : resolve())); - } -}); - -test.skip("removed Kafka retention replay product authority binds trace ownership", async () => { - const sessionId = "ses_kafka_refresh_trace_owner"; - const foreignSessionId = "ses_kafka_refresh_trace_foreign"; - const traceId = "trc_kafka_refresh_trace_owner"; - let retentionQuery = null; - const bridge = refreshReplayBridge({ - records: [refreshRecord(0, foreignSessionId, traceId, "assistant", { text: "foreign session" })], - endOffset: "1", - onQuery(params) { retentionQuery = params; } - }); - const server = createCloudApiServer({ - accessController: realtimeAccessController({ sessions: [{ id: sessionId, ownerUserId: ACTOR.id, lastTraceId: traceId }] }), - kafkaEventBridge: bridge, - env: { ...REFRESH_REALTIME_ENV, HWLAB_WORKBENCH_SSE_HEARTBEAT_MS: "60000" } - }); - await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve)); - - try { - const response = await fetch(`http://127.0.0.1:${server.address().port}/v1/workbench/events?traceId=${traceId}`); - assert.equal(response.status, 200); - const events = (await response.text()).trim().split("\n\n").filter(Boolean).map(parseSseBlock); - assert.deepEqual(events.map((event) => event.event), ["workbench.error"]); - assert.equal(events[0].data.sessionId, sessionId); - assert.equal(events[0].data.traceId, traceId); - assert.equal(events[0].data.error.code, "workbench_kafka_refresh_scope_mismatch"); - assert.equal(retentionQuery.sessionId, sessionId); - assert.equal(retentionQuery.traceId, traceId); - } finally { - await new Promise((resolve, reject) => server.close((error) => error ? reject(error) : resolve())); - } -}); - -test.skip("removed Kafka retention replay product authority reports retention gaps", async () => { - const sessionId = "ses_kafka_refresh_gap"; - const traceId = "trc_kafka_refresh_gap"; - const retained = refreshRecord(0, sessionId, traceId, "user", { userMessageId: "msg_kafka_refresh_gap", messageId: "msg_kafka_refresh_gap", text: "retained input" }); - const missing = refreshRecord(1, sessionId, traceId, "assistant", { assistantText: "missing from scan" }); - const bridge = refreshReplayBridge({ - records: [retained], - endOffset: "2", - beforeQueryResult(listener) { listener(missing.value, { topic: missing.topic, partition: missing.partition, offset: missing.offset }); } - }); - const server = createCloudApiServer({ - accessController: realtimeAccessController({ sessions: [{ id: sessionId, ownerUserId: ACTOR.id, lastTraceId: traceId }] }), - kafkaEventBridge: bridge, - env: { ...REFRESH_REALTIME_ENV, HWLAB_WORKBENCH_SSE_HEARTBEAT_MS: "60000" } - }); - await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve)); - - try { - const response = await fetch(`http://127.0.0.1:${server.address().port}/v1/workbench/events?sessionId=${sessionId}`); - assert.equal(response.status, 200); - const body = await response.text(); - const events = body.trim().split("\n\n").filter(Boolean).map(parseSseBlock); - assert.deepEqual(events.map((event) => event.event), ["hwlab.event.v1", "workbench.error"]); - assert.equal(events[1].data.sessionId, sessionId); - assert.equal(events[1].data.traceId, null); - assert.equal(events[1].data.phase, "flushing"); - assert.equal(events[1].data.error.code, "workbench_kafka_refresh_barrier_gap"); - assert.deepEqual(events[1].data.error.details, { partition: 0 }); - assert.equal(events[1].data.fallback, false); - assert.equal(body.includes("workbench.connected"), false); - assert.equal(body.includes("workbench.trace.snapshot"), false); - assert.equal(body.includes("workbench.message.snapshot"), false); - } finally { - await new Promise((resolve, reject) => server.close((error) => error ? reject(error) : resolve())); - } -}); - -test.skip("removed live-only product authority keeps its transport open", 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_WORKBENCH_KAFKA_REFRESH_REPLAY_ENABLED: "false", - 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.skip("removed live-only product authority rejects foreign ownership", 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_WORKBENCH_KAFKA_REFRESH_REPLAY_ENABLED: "false", - 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.skip("removed live-only product authority fails closed without ownership lookup", 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_WORKBENCH_KAFKA_REFRESH_REPLAY_ENABLED: "false", - 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.skip("removed live-only product authority permits admin ownership scope", 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_WORKBENCH_KAFKA_REFRESH_REPLAY_ENABLED: "false", - 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); + const response = await getJson(server.address().port, "/v1/workbench/sync?sessionId=ses_removed_sync_repair"); + assert.equal(response.status, 404); } finally { await new Promise((resolve, reject) => server.close((error) => error ? reject(error) : resolve())); } diff --git a/internal/cloud/server.ts b/internal/cloud/server.ts index 0965f300..148283d8 100644 --- a/internal/cloud/server.ts +++ b/internal/cloud/server.ts @@ -773,7 +773,7 @@ async function handleRestAdapter(request, response, url, options) { return; } - if (url.pathname === "/v1/workbench/sync" || url.pathname === "/v1/workbench/sessions" || url.pathname.startsWith("/v1/workbench/sessions/") || url.pathname.startsWith("/v1/workbench/turns/") || url.pathname.startsWith("/v1/workbench/traces/")) { + if (url.pathname === "/v1/workbench/sessions" || url.pathname.startsWith("/v1/workbench/sessions/") || url.pathname.startsWith("/v1/workbench/turns/") || url.pathname.startsWith("/v1/workbench/traces/")) { await handleWorkbenchReadModelHttp(request, response, url, options); return; } @@ -1011,7 +1011,7 @@ function navIdForRestPath(pathname, method = "GET") { if (pathname === "/v1/api-keys" || pathname === "/v1/api-keys/default" || pathname.startsWith("/v1/api-keys/")) return "user.apiKeys"; if (pathname === "/v1/users/me/profile" || pathname === "/v1/users/me/password") return "system.settings"; if (pathname === "/v1/workbench/debug/fake-sse" || pathname.startsWith("/v1/workbench/debug/fake-sse/") || pathname === "/v1/workbench/debug/kafka-sse" || pathname.startsWith("/v1/workbench/debug/kafka-sse/")) return "workbench.debug"; - if (pathname === "/v1/workbench/events" || pathname === "/v1/workbench/projection-events" || pathname === "/v1/workbench/sync" || pathname === "/v1/workbench/launches" || pathname === "/v1/workbench/sessions" || pathname.startsWith("/v1/workbench/sessions/") || pathname.startsWith("/v1/workbench/turns/") || pathname.startsWith("/v1/workbench/traces/")) return "workbench.code"; + if (pathname === "/v1/workbench/events" || pathname === "/v1/workbench/projection-events" || pathname === "/v1/workbench/launches" || pathname === "/v1/workbench/sessions" || pathname.startsWith("/v1/workbench/sessions/") || pathname.startsWith("/v1/workbench/turns/") || pathname.startsWith("/v1/workbench/traces/")) return "workbench.code"; if (pathname === "/v1/agent/chat" || pathname === "/v1/agent/sessions" || pathname.startsWith("/v1/agent/sessions/") || pathname === "/v1/agent/chat/inspect" || pathname.startsWith("/v1/agent/chat/result/") || pathname.startsWith("/v1/agent/turns/") || pathname.startsWith("/v1/agent/traces/") || pathname === "/v1/agent/chat/cancel" || pathname === "/v1/agent/chat/steer") return "workbench.code"; if (pathname === "/v1/admin/provider-profiles" || pathname.startsWith("/v1/admin/provider-profiles/")) return "admin.providerProfiles"; if (pathname === "/v1/admin/secrets" || pathname.startsWith("/v1/admin/secrets/")) return "admin.secrets"; diff --git a/internal/cloud/workbench-projection-outbox-events.ts b/internal/cloud/workbench-projection-outbox-events.ts index 7e5fed25..1bcc780a 100644 --- a/internal/cloud/workbench-projection-outbox-events.ts +++ b/internal/cloud/workbench-projection-outbox-events.ts @@ -1,5 +1,5 @@ /* - * Immutable Workbench projection outbox rows are the common SSE and sync-replay authority. + * Immutable Workbench projection outbox rows are the common replay and realtime SSE authority. */ export const WORKBENCH_REALTIME_AUTHORITY_VERSION = "workbench-realtime-authority-v2"; diff --git a/internal/cloud/workbench-realtime-authority-contract.test.ts b/internal/cloud/workbench-realtime-authority-contract.test.ts deleted file mode 100644 index 6a8123bb..00000000 --- a/internal/cloud/workbench-realtime-authority-contract.test.ts +++ /dev/null @@ -1,428 +0,0 @@ -// SPEC: PJ2026-010401080313 Workbench实时权威 draft-2026-07-08-p0-workbench-realtime-authority-v2. -// Responsibility: P1 backend events/sync contract tests; no browser, provider, Kubernetes, or live DB dependency. - -import assert from "node:assert/strict"; -import { test } from "bun:test"; - -import { createCloudApiServer } from "./server.ts"; - -const ACTOR = { id: "usr_workbench_realtime_authority", username: "reader", displayName: "Reader", role: "user", status: "active" }; -const PROJECTION_REALTIME_ENV = Object.freeze({ - HWLAB_KAFKA_DIRECT_PUBLISH_ENABLED: "false", - HWLAB_WORKBENCH_LIVE_KAFKA_SSE_ENABLED: "false", - HWLAB_WORKBENCH_KAFKA_REFRESH_REPLAY_ENABLED: "false", - HWLAB_KAFKA_TRANSACTIONAL_PROJECTOR_ENABLED: "false", - HWLAB_KAFKA_PROJECTION_OUTBOX_RELAY_ENABLED: "false", - HWLAB_WORKBENCH_PROJECTION_REALTIME_ENABLED: "true" -}); -const LIVE_REALTIME_ENV = Object.freeze({ - HWLAB_KAFKA_DIRECT_PUBLISH_ENABLED: "true", - HWLAB_WORKBENCH_LIVE_KAFKA_SSE_ENABLED: "true", - HWLAB_WORKBENCH_KAFKA_REFRESH_REPLAY_ENABLED: "false", - HWLAB_KAFKA_TRANSACTIONAL_PROJECTOR_ENABLED: "false", - HWLAB_KAFKA_PROJECTION_OUTBOX_RELAY_ENABLED: "false", - HWLAB_WORKBENCH_PROJECTION_REALTIME_ENABLED: "false" -}); - -function projectionRealtimeBridge() { - return { - started: true, - capabilities: { - directPublish: false, - liveKafkaSse: false, - transactionalProjector: false, - projectionOutboxRelay: false, - projectionRealtime: true - }, - ready: Promise.resolve(), - subscribeProjectionCommits() { return () => {}; }, - async stop() {} - }; -} - -test("workbench sync fails closed before projection storage in live Kafka capability", async () => { - let syncReads = 0; - const server = createCloudApiServer({ - accessController: createAccessController(), - workbenchRuntime: { - async readAtomicWorkbenchProjectionSync() { - syncReads += 1; - throw new Error("live Kafka sync must not reach projection storage"); - } - }, - kafkaEventBridge: { - started: true, - capabilities: { - directPublish: true, - liveKafkaSse: true, - transactionalProjector: false, - projectionOutboxRelay: false, - projectionRealtime: false - }, - ready: Promise.resolve(), - subscribeLiveHwlabEvents() { return () => {}; }, - async stop() {} - }, - env: { ...LIVE_REALTIME_ENV } - }); - await listen(server); - - try { - const { port } = server.address(); - const response = await getJson(port, "/v1/workbench/sync?sessionId=ses_live_only&since=0"); - assert.equal(response.status, 503); - assert.equal(response.body.error.code, "workbench_projection_realtime_disabled"); - assert.equal(response.body.error.capabilities.liveKafkaSse, true); - assert.equal(syncReads, 0); - } finally { - await close(server); - } -}); - -test("workbench sync and SSE read the same atomic projection outbox", async () => { - const sessionId = "ses_realtime_authority_p1"; - const traceId = "trc_realtime_authority_p1"; - const outboxRows = [{ - outboxSeq: 12, - eventSeq: 44, - aggregateId: sessionId, - aggregateSeq: 5, - projectionRevision: 42, - traceId, - sessionId, - turnId: traceId, - messageId: "msg_realtime_authority_agent", - projectedSeq: 8, - sourceSeq: 8, - sourceEventId: "src_realtime_authority_terminal", - commitType: "terminal", - terminal: true, - sealed: true, - createdAt: "2026-07-08T12:20:00.000Z", - entityFamily: "turns", - entityId: traceId, - payload: { - family: "turns", - fact: { - turnId: traceId, - sessionId, - traceId, - messageId: "msg_realtime_authority_agent", - status: "completed", - projectedSeq: 8, - terminal: true, - sealed: true, - finalResponse: { text: "redacted final", status: "completed", traceId } - }, - valuesRedacted: true - } - }]; - const outboxQueries = []; - const runtime = createRuntime({ facts: durableFacts({ sessionId, traceId }), outboxRows, outboxQueries }); - const server = createCloudApiServer({ - accessController: createAccessController(), - workbenchRuntime: runtime, - kafkaEventBridge: projectionRealtimeBridge(), - env: { - ...PROJECTION_REALTIME_ENV, - HWLAB_WORKBENCH_SSE_HEARTBEAT_MS: "60000", - HWLAB_WORKBENCH_OUTBOX_TAIL_BATCH_SIZE: "100" - } - }); - await listen(server); - - try { - const { port } = server.address(); - const sync = await getJson(port, `/v1/workbench/sync?sessionId=${encodeURIComponent(sessionId)}&since=10&limit=10`); - assert.equal(sync.status, 200); - assert.equal(sync.body.contractVersion, "workbench-sync-v1"); - assert.equal(sync.body.realtimeAuthority, "workbench-realtime-authority-v2"); - assert.equal(sync.body.cursor.outboxSeq, 12); - assert.equal(sync.body.events.length, 1); - assert.equal(sync.body.events[0].family, "turns"); - assert.equal(sync.body.events[0].entity.version, 42); - assert.equal(sync.body.events[0].terminal, true); - assert.equal(sync.body.events[0].sealed, true); - - const events = await getSseEvents(port, `/v1/workbench/events?sessionId=${encodeURIComponent(sessionId)}&traceId=${encodeURIComponent(traceId)}&afterSeq=10`, 2); - assert.deepEqual(events.map((event) => event.event), ["workbench.connected", "workbench.turn.snapshot"]); - 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, "turns"); - assert.equal(events[1].data.entity.version, 42); - assert.equal(events[1].data.turn.finalResponse.text, "redacted final"); - assert.deepEqual(outboxQueries.map((query) => query.afterOutboxSeq), [10, 10]); - assert.equal(outboxQueries.every((query) => query.sessionId === sessionId), true); - assert.deepEqual(outboxQueries.at(-1).actor, { id: ACTOR.id, role: ACTOR.role }); - } finally { - await close(server); - } -}); - -test("workbench sync delta includes all authority families and marks trace rows detail-only", async () => { - const sessionId = "ses_realtime_authority_delta"; - const traceId = "trc_realtime_authority_delta"; - const runtime = createRuntime({ - facts: durableFacts({ - sessionId, - traceId, - finalText: "sealed backend final", - traceEventText: "trace detail progress" - }), - outboxRows: [] - }); - const server = createCloudApiServer({ accessController: createAccessController(), workbenchRuntime: runtime, kafkaEventBridge: projectionRealtimeBridge(), env: { ...PROJECTION_REALTIME_ENV } }); - await listen(server); - - try { - const { port } = server.address(); - const sync = await getJson(port, `/v1/workbench/sync?sessionId=${encodeURIComponent(sessionId)}&since=0`); - assert.equal(sync.status, 200); - assert.equal(sync.body.authority.syncReplay, true); - assert.equal(sync.body.authority.sseEquivalent, true); - assert.equal(sync.body.authority.detailRoutes, "detail-history-only"); - assert.equal(sync.body.families.sessions, 1); - assert.equal(sync.body.families.messages, 2); - assert.equal(sync.body.families.parts, 2); - assert.equal(sync.body.families.turns, 1); - assert.equal(sync.body.families.traceEvents, 1); - assert.equal(sync.body.families.checkpoints, 1); - assert.equal(sync.body.delta.turns[0].finalResponse.text, "sealed backend final"); - assert.equal(sync.body.delta.traceEvents[0].detailProjection, true); - assert.equal(sync.body.delta.traceEvents[0].authority, "trace-detail-only"); - assert.equal(JSON.stringify(sync.body.delta.traceEvents[0]).includes("sealed backend final"), false); - } finally { - await close(server); - } -}); - -test("workbench sync preserves session and trace scope for narrow replay", async () => { - const sessionId = "ses_realtime_authority_narrow"; - const traceId = "trc_realtime_authority_narrow"; - const outboxQueries = []; - const runtime = createRuntime({ facts: durableFacts({ sessionId, traceId }), outboxQueries }); - const server = createCloudApiServer({ accessController: createAccessController(), workbenchRuntime: runtime, kafkaEventBridge: projectionRealtimeBridge(), env: { ...PROJECTION_REALTIME_ENV } }); - await listen(server); - - try { - const { port } = server.address(); - const sync = await getJson(port, `/v1/workbench/sync?sessionId=${encodeURIComponent(sessionId)}&traceId=${encodeURIComponent(traceId)}&since=7`); - assert.equal(sync.status, 200); - assert.deepEqual(sync.body.scope, { kind: "trace", sessionId, traceId, since: 7 }); - assert.equal(outboxQueries.length, 1); - assert.equal(outboxQueries[0].sessionId, sessionId); - assert.equal(outboxQueries[0].traceId, traceId); - } finally { - await close(server); - } -}); - -test("workbench trace event detail route declares detail-only authority", async () => { - const sessionId = "ses_realtime_authority_detail"; - const traceId = "trc_realtime_authority_detail"; - const runtime = createRuntime({ - facts: durableFacts({ sessionId, traceId, finalText: "sealed final", traceEventText: "detail row" }), - outboxRows: [] - }); - const server = createCloudApiServer({ accessController: createAccessController(), workbenchRuntime: runtime, kafkaEventBridge: projectionRealtimeBridge(), env: { ...PROJECTION_REALTIME_ENV } }); - await listen(server); - - try { - const { port } = server.address(); - const detail = await getJson(port, `/v1/workbench/traces/${encodeURIComponent(traceId)}/events?limit=10`); - assert.equal(detail.status, 200); - assert.equal(detail.body.detailProjection, true); - assert.equal(detail.body.authority, "trace-detail-only"); - assert.equal(detail.body.realtimeAuthority, "workbench-realtime-authority-v2"); - assert.equal(detail.body.events.length, 1); - assert.equal(detail.body.events[0].authority, "trace-detail-only"); - assert.equal(detail.body.finalResponse, undefined); - } finally { - await close(server); - } -}); - -test("workbench sync rejects unscoped automatic repair requests", async () => { - const server = createCloudApiServer({ accessController: createAccessController(), workbenchRuntime: createRuntime({ facts: durableFacts({}) }), kafkaEventBridge: projectionRealtimeBridge(), env: { ...PROJECTION_REALTIME_ENV } }); - await listen(server); - - try { - const { port } = server.address(); - const response = await getJson(port, "/v1/workbench/sync?since=0"); - assert.equal(response.status, 400); - assert.equal(response.body.error.code, "workbench_sync_scope_required"); - } finally { - await close(server); - } -}); - -function createAccessController() { - return { - async ensureBootstrap() {}, - async authenticate() { - return { ok: true, actor: ACTOR, session: { id: "uss_workbench_realtime_authority" } }; - } - }; -} - -function createRuntime({ facts, outboxRows = [], outboxQueries = [] }) { - return { - async readAtomicWorkbenchProjectionSync(params = {}) { - outboxQueries.push({ ...params }); - const after = Number(params.afterOutboxSeq ?? params.afterSeq ?? 0); - const events = outboxRows.filter((row) => Number(row.outboxSeq) > after); - const cutoffOutboxSeq = outboxRows.reduce((max, row) => Math.max(max, Number(row.outboxSeq) || 0), 0); - return { - facts: filterFacts(facts, params), - events, - cutoffOutboxSeq, - cursorOutboxSeq: cutoffOutboxSeq, - hasMore: false, - valuesRedacted: true - }; - }, - async queryWorkbenchFacts(params = {}) { - return { - facts: filterFacts(facts, params), - count: Object.values(facts).reduce((sum, rows) => sum + rows.length, 0), - persistence: { adapter: "test-realtime-authority", durable: true } - }; - } - }; -} - -function durableFacts({ sessionId = "ses_realtime_authority", traceId = "trc_realtime_authority", finalText = "sealed final", traceEventText = "detail progress" } = {}) { - const now = "2026-07-08T12:20:00.000Z"; - return { - sessions: [{ - sessionId, - ownerUserId: ACTOR.id, - projectId: "prj_realtime_authority", - conversationId: "cnv_realtime_authority", - threadId: "thread-realtime-authority", - status: "completed", - lastTraceId: traceId, - projectedSeq: 42, - terminal: true, - sealed: true, - updatedAt: now, - valuesRedacted: true - }], - messages: [ - { messageId: "msg_realtime_authority_user", sessionId, turnId: traceId, traceId, role: "user", status: "sent", projectedSeq: 1, text: "hi", updatedAt: now, valuesRedacted: true }, - { messageId: "msg_realtime_authority_agent", sessionId, turnId: traceId, traceId, role: "agent", status: "completed", projectedSeq: 42, terminal: true, sealed: true, text: finalText, updatedAt: now, valuesRedacted: true } - ], - parts: [ - { partId: "prt_realtime_authority_user", messageId: "msg_realtime_authority_user", sessionId, turnId: traceId, traceId, partIndex: 0, partType: "text", status: "sent", text: "hi", projectedSeq: 1, updatedAt: now, valuesRedacted: true }, - { partId: "prt_realtime_authority_final", messageId: "msg_realtime_authority_agent", sessionId, turnId: traceId, traceId, partIndex: 0, partType: "final_response", status: "completed", text: finalText, projectedSeq: 42, terminal: true, sealed: true, updatedAt: now, valuesRedacted: true } - ], - turns: [{ - turnId: traceId, - sessionId, - traceId, - messageId: "msg_realtime_authority_agent", - status: "completed", - projectedSeq: 42, - terminal: true, - sealed: true, - finalResponse: { text: finalText, status: "completed", traceId, valuesPrinted: false }, - timing: { startedAt: now, lastEventAt: now, finishedAt: now, durationMs: 0, valuesRedacted: true }, - updatedAt: now, - valuesRedacted: true - }], - traceEvents: [{ - id: "wte_realtime_authority_detail", - traceId, - sessionId, - turnId: traceId, - messageId: "msg_realtime_authority_agent", - projectedSeq: 42, - sourceSeq: 42, - sourceEventId: "src_realtime_authority_detail", - eventType: "assistant_message", - status: "completed", - message: traceEventText, - terminal: true, - sealed: true, - occurredAt: now, - updatedAt: now, - valuesRedacted: true - }], - checkpoints: [{ - traceId, - sessionId, - turnId: traceId, - projectedSeq: 42, - sourceSeq: 42, - projectionStatus: "caught_up", - projectionHealth: "healthy", - terminal: true, - sealed: true, - updatedAt: now, - valuesRedacted: true - }] - }; -} - -function filterFacts(facts, params = {}) { - const families = new Set(Array.isArray(params.families) ? params.families : ["sessions", "messages", "parts", "turns", "traceEvents", "checkpoints"]); - const sessionId = params.sessionId; - const traceId = params.traceId; - return { - sessions: families.has("sessions") ? facts.sessions.filter((record) => (!sessionId || record.sessionId === sessionId) && (!traceId || record.lastTraceId === traceId)) : [], - messages: families.has("messages") ? facts.messages.filter((record) => (!sessionId || record.sessionId === sessionId) && (!traceId || record.traceId === traceId)) : [], - parts: families.has("parts") ? facts.parts.filter((record) => (!sessionId || record.sessionId === sessionId) && (!traceId || record.traceId === traceId)) : [], - turns: families.has("turns") ? facts.turns.filter((record) => (!sessionId || record.sessionId === sessionId) && (!traceId || record.traceId === traceId)) : [], - traceEvents: families.has("traceEvents") ? facts.traceEvents.filter((record) => (!sessionId || record.sessionId === sessionId) && (!traceId || record.traceId === traceId)) : [], - checkpoints: families.has("checkpoints") ? facts.checkpoints.filter((record) => (!sessionId || record.sessionId === sessionId) && (!traceId || record.traceId === traceId)) : [] - }; -} - -async function getJson(port, path) { - const response = await fetch(`http://127.0.0.1:${port}${path}`); - return { status: response.status, body: await response.json() }; -} - -async function getSseEvents(port, path, count) { - const controller = new AbortController(); - const response = await fetch(`http://127.0.0.1:${port}${path}`, { signal: controller.signal }); - assert.equal(response.status, 200); - const reader = response.body.getReader(); - const decoder = new TextDecoder(); - let buffer = ""; - const events = []; - try { - while (events.length < count) { - const { value, done } = await reader.read(); - if (done) break; - buffer += decoder.decode(value, { stream: true }); - for (;;) { - const index = buffer.indexOf("\n\n"); - if (index < 0) break; - const block = buffer.slice(0, index); - buffer = buffer.slice(index + 2); - events.push(parseSseBlock(block)); - if (events.length >= count) break; - } - } - } finally { - controller.abort(); - } - return events; -} - -function parseSseBlock(block) { - const lines = block.split(/\n/u); - const event = lines.find((line) => line.startsWith("event:"))?.slice("event:".length).trim() ?? "message"; - const id = lines.find((line) => line.startsWith("id:"))?.slice("id:".length).trim() ?? null; - const dataLine = lines.find((line) => line.startsWith("data:"))?.slice("data:".length).trim() ?? "{}"; - return { event, id, data: JSON.parse(dataLine) }; -} - -function listen(server) { - return new Promise((resolve) => server.listen(0, "127.0.0.1", resolve)); -} - -function close(server) { - return new Promise((resolve, reject) => server.close((error) => error ? reject(error) : resolve())); -} diff --git a/internal/cloud/workbench-realtime-authority.ts b/internal/cloud/workbench-realtime-authority.ts deleted file mode 100644 index 1a23d3da..00000000 --- a/internal/cloud/workbench-realtime-authority.ts +++ /dev/null @@ -1,162 +0,0 @@ -/* - * SPEC: PJ2026-010401080313 Workbench实时权威 draft-2026-07-08-p0-workbench-realtime-authority-v2. - * Responsibility: Workbench realtime authority sync/replay contract helpers. REST detail routes remain detail/history only. - */ -import { safeSessionId, safeTraceId, sendJson } from "./server-http-utils.ts"; -import { createWorkbenchRuntimeClient } from "./workbench-runtime-client.ts"; -import { projectionOutboxRealtimeEvents } from "./workbench-projection-outbox-events.ts"; -import { workbenchRealtimeCapabilities } from "./workbench-realtime-capabilities.ts"; - -const SYNC_CONTRACT_VERSION = "workbench-sync-v1"; -const REALTIME_AUTHORITY_VERSION = "workbench-realtime-authority-v2"; -const DEFAULT_SYNC_LIMIT = 100; -const MAX_SYNC_LIMIT = 500; - -export async function handleWorkbenchSyncHttp(request, response, url, options = {}, actor = null) { - if (request.method !== "GET") { - sendJson(response, 405, { ok: false, error: { code: "method_not_allowed", message: "GET required." } }); - return; - } - const capabilities = workbenchRealtimeCapabilities(options.env ?? process.env); - if (!capabilities.projectionRealtime) { - sendJson(response, 503, { - ok: false, - error: { - code: "workbench_projection_realtime_disabled", - message: "Workbench projection sync/replay capability is disabled.", - valuesRedacted: true - } - }); - return; - } - const sessionId = safeSessionId(url.searchParams.get("sessionId") ?? url.searchParams.get("includeSessionId")); - const traceId = safeTraceId(url.searchParams.get("traceId")); - if (!sessionId && !traceId) { - sendJson(response, 400, { - ok: false, - error: { - code: "workbench_sync_scope_required", - message: "Workbench sync requires sessionId or traceId.", - valuesRedacted: true - } - }); - return; - } - const runtime = workbenchRealtimeRuntime(request, options); - if (typeof runtime?.readAtomicWorkbenchProjectionSync !== "function") { - sendJson(response, 503, { - ok: false, - error: { - code: "workbench_sync_runtime_unconfigured", - message: "Workbench sync requires one atomic durable projection snapshot reader.", - valuesRedacted: true - } - }); - return; - } - - try { - const since = realtimeSinceCursor(request, url); - const limit = boundedSyncLimit(url.searchParams.get("limit")); - const snapshot = await runtime.readAtomicWorkbenchProjectionSync({ - afterOutboxSeq: since, - afterSeq: since, - limit, - ...(sessionId ? { sessionId } : {}), - ...(traceId ? { traceId } : {}), - actor: actor ? { id: actor.id, role: actor.role ?? "user" } : undefined - }); - const rows = Array.isArray(snapshot?.events) ? snapshot.events : []; - const events = projectionOutboxRealtimeEvents({ events: rows, facts: {} }).map((item) => item.payload); - const latestOutboxSeq = nonNegativeInteger(snapshot?.cursorOutboxSeq); - const delta = normalizeSyncFacts(snapshot?.facts); - sendJson(response, 200, { - ok: true, - status: "succeeded", - contractVersion: SYNC_CONTRACT_VERSION, - realtimeAuthority: REALTIME_AUTHORITY_VERSION, - scope: { - kind: traceId ? "trace" : "session", - sessionId, - traceId, - since - }, - cursor: { - outboxSeq: latestOutboxSeq, - since, - cutoffOutboxSeq: nonNegativeInteger(snapshot?.cutoffOutboxSeq), - hasMore: snapshot?.hasMore === true - }, - events, - delta, - families: Object.fromEntries(Object.entries(delta).map(([family, values]) => [family, values.length])), - authority: { - syncReplay: true, - sseEquivalent: true, - detailRoutes: "detail-history-only", - automaticRepairEndpoint: "/v1/workbench/sync", - valuesRedacted: true - }, - valuesRedacted: true, - secretMaterialStored: false - }); - } catch (error) { - sendJson(response, 503, { - ok: false, - error: { - code: error?.code ?? "workbench_sync_failed", - message: error?.message ?? "Workbench sync failed.", - retryable: error?.data?.retryable !== false, - valuesRedacted: true - } - }); - } -} - -function workbenchRealtimeRuntime(request, options = {}) { - return options.workbenchRuntime - ?? createWorkbenchRuntimeClient({ env: options.env ?? process.env, fetch: options.fetch, logger: options.logger, traceparent: request?.hwlabHttpRequestContext?.traceparent }); -} - -function realtimeSinceCursor(request, url) { - for (const value of [url.searchParams.get("since"), url.searchParams.get("afterSeq"), url.searchParams.get("afterOutboxSeq"), request?.headers?.["last-event-id"]]) { - const parsed = nonNegativeInteger(value); - if (parsed > 0) return parsed; - } - return 0; -} - -function boundedSyncLimit(value) { - const parsed = Number.parseInt(String(value ?? ""), 10); - if (!Number.isInteger(parsed) || parsed <= 0) return DEFAULT_SYNC_LIMIT; - return Math.min(Math.max(parsed, 1), MAX_SYNC_LIMIT); -} - -function normalizeSyncFacts(facts = {}) { - return { - sessions: factArray(facts.sessions), - messages: factArray(facts.messages), - parts: factArray(facts.parts), - turns: factArray(facts.turns), - traceEvents: factArray(facts.traceEvents).map((event) => ({ - ...event, - detailProjection: true, - authority: "trace-detail-only" - })), - checkpoints: factArray(facts.checkpoints) - }; -} - -function factArray(value) { - return Array.isArray(value) ? value : []; -} - -function textValue(value) { - const text = typeof value === "string" ? value.trim() : value === null || value === undefined ? "" : String(value).trim(); - return text || ""; -} - -function nonNegativeInteger(value) { - const parsed = Number.parseInt(String(value ?? ""), 10); - return Number.isInteger(parsed) && parsed >= 0 ? parsed : 0; -} diff --git a/web/hwlab-cloud-web/scripts/check.ts b/web/hwlab-cloud-web/scripts/check.ts index 4a5965cc..2c7f9d1b 100644 --- a/web/hwlab-cloud-web/scripts/check.ts +++ b/web/hwlab-cloud-web/scripts/check.ts @@ -193,9 +193,9 @@ assertIncludes(workbenchColadaSource, "staleTime", "Workbench query min-interval assertIncludes(workbenchPerformanceSource, "recordWorkbenchRuntimeDiagnostic", "Workbench performance probe must record runtime diagnostics for monitor root cause visibility"); assertIncludes(workbenchPerformanceSource, "clearResourceTimings", "Workbench performance probe must bound browser ResourceTiming retention after API enrichment"); assertIncludes(workbenchStoreSource, "recordWorkbenchRuntimeDiagnostic", "Workbench store must surface SSE recovery diagnostics to the performance probe"); -assertIncludes(workbenchRealtimePlanSource, "recovery.actions.includes(\"sync-replay\")", "Realtime recovery planner must consume transport-owned actions explicitly"); +assertIncludes(workbenchRealtimePlanSource, "recovery.actions.includes(\"events-reconnect\")", "Realtime recovery planner must reconnect the projection SSE explicitly"); assert.doesNotMatch(workbenchRealtimePlanSource, /schedule-session-list/u, "Realtime recovery planner must not restore legacy session-list repair actions"); -assertIncludes(workbenchRealtimePlanSource, "authority: \"automatic-recovery\"", "Realtime recovery planner must classify transport recovery as automatic recovery authority"); +assertIncludes(workbenchRealtimePlanSource, "afterOutboxSeq: finiteNumber(recovery.outboxSeq)", "Realtime recovery planner must resume projection SSE from the durable outbox cursor"); assert.doesNotMatch(workbenchRealtimePlanSource, /force:\s*true/u, "Realtime recovery planner must not turn transport recovery into force-refresh work"); assertIncludes(workbenchColadaSource, "const state = await queryCache.refresh(entry);", "Workbench reads must preserve Colada staleTime/min-interval governance"); assert.doesNotMatch(workbenchColadaSource, /queryCache\.fetch\(entry/u, "Workbench reads must not call queryCache.fetch(entry), which bypasses Colada freshness governance"); diff --git a/web/hwlab-cloud-web/scripts/workbench-realtime-runtime.test.ts b/web/hwlab-cloud-web/scripts/workbench-realtime-runtime.test.ts index fa12dba2..3de3a4a5 100644 --- a/web/hwlab-cloud-web/scripts/workbench-realtime-runtime.test.ts +++ b/web/hwlab-cloud-web/scripts/workbench-realtime-runtime.test.ts @@ -23,7 +23,7 @@ import { WORKBENCH_TIMELINE_OPENCODE_PARITY, buildWorkbenchTimelineRows, normali 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, workbenchSyncReplayEvents } from "../src/stores/workbench-realtime-authority.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"; @@ -593,35 +593,30 @@ test("realtime apply planner turns reducer actions into store steps", () => { assert.deepEqual(ignored.steps, []); }); -test("realtime recovery planner uses sync replay instead of legacy repair fan-out", () => { - const recovery = recoveryEvent(["sync-replay"], { outboxSeq: 42 }); +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), ["sync-replay"]); - assert.deepEqual(authorized.steps.map((step) => step.authority), ["automatic-recovery"]); - assert.equal(authorized.steps[0]?.sinceOutboxSeq, 42); + 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), ["sync-replay"]); + 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), ["sync-replay"]); + 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), []); + assert.deepEqual(terminalSealed.steps.map((step) => step.type), ["events-reconnect"]); }); -test("realtime authority accepts events and sync replay through the same entity contract", () => { +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); - - const replay = workbenchSyncReplayEvents({ contractVersion: "workbench-sync-v1", realtimeAuthority: WORKBENCH_REALTIME_AUTHORITY_VERSION, events: [event], delta: [event] }); - assert.equal(replay.length, 1); - assert.equal(replay[0]?.entity?.id, "msg_1"); }); test("realtime authority rejects trace detail-only and incomplete contract payloads", () => { diff --git a/web/hwlab-cloud-web/src/api/workbench-events.ts b/web/hwlab-cloud-web/src/api/workbench-events.ts index bb8db867..38700e39 100644 --- a/web/hwlab-cloud-web/src/api/workbench-events.ts +++ b/web/hwlab-cloud-web/src/api/workbench-events.ts @@ -75,32 +75,6 @@ export interface WorkbenchRealtimeEvent { [key: string]: unknown; } -export interface WorkbenchSyncReplayRequest { - sessionId?: string | null; - traceId?: string | null; - since?: number | null; -} - -export interface WorkbenchSyncReplayResponse { - contractVersion?: string | null; - realtimeAuthority?: string | null; - scope?: Record | null; - cursor?: Record | null; - events?: WorkbenchRealtimeEvent[]; - delta?: WorkbenchRealtimeEvent[] | { - sessions?: Record[]; - messages?: Record[]; - parts?: Record[]; - turns?: Record[]; - traceEvents?: Record[]; - checkpoints?: Record[]; - [key: string]: unknown; - }; - families?: Record | null; - authority?: Record | null; - [key: string]: unknown; -} - export interface WorkbenchEventStreamOptions { realtimeCapabilities: WorkbenchRealtimeCapabilities; sessionId?: string | null; @@ -225,21 +199,6 @@ export function workbenchProjectionEventStreamPath(options: Pick> { - return fetchJson(workbenchSyncReplayPath(input), { - ...options, - timeoutName: options.timeoutName ?? "workbench sync replay" - }); -} - -export function workbenchSyncReplayPath(input: WorkbenchSyncReplayRequest): string { - const params = new URLSearchParams(); - appendParam(params, "sessionId", input.sessionId); - appendParam(params, "traceId", input.traceId); - appendNumberParam(params, "since", input.since); - return `/v1/workbench/sync?${params.toString()}`; -} - function appendParam(params: URLSearchParams, key: string, value: string | null | undefined): void { const text = typeof value === "string" ? value.trim() : ""; if (text) params.set(key, text); diff --git a/web/hwlab-cloud-web/src/stores/workbench-colada-keys.ts b/web/hwlab-cloud-web/src/stores/workbench-colada-keys.ts index a6f75444..7767983f 100644 --- a/web/hwlab-cloud-web/src/stores/workbench-colada-keys.ts +++ b/web/hwlab-cloud-web/src/stores/workbench-colada-keys.ts @@ -19,12 +19,6 @@ export interface WorkbenchTraceEventsKeyInput { limit?: number | null; } -export interface WorkbenchSyncReplayKeyInput { - sessionId?: string | null; - traceId?: string | null; - since?: number | null; -} - export const workbenchColadaKeys = { root: (): EntryKey => ["workbench"], serverState: (): EntryKey => ["workbench", "projection", "server-state"], @@ -33,8 +27,6 @@ export const workbenchColadaKeys = { snapshotSessions: (input: WorkbenchSessionsKeyInput = {}): EntryKey => ["workbench", "snapshot", "sessions", normalizeKeyObject({ includeSessionId: input.includeSessionId ?? null, cursor: input.cursor ?? null, limit: finiteNumber(input.limit) })], snapshotSessionRoot: (): EntryKey => ["workbench", "snapshot", "session"], snapshotSession: (sessionId: string): EntryKey => ["workbench", "snapshot", "session", sessionId], - syncRoot: (): EntryKey => ["workbench", "sync"], - syncReplay: (input: WorkbenchSyncReplayKeyInput = {}): EntryKey => ["workbench", "sync", "replay", normalizeKeyObject({ sessionId: input.sessionId ?? null, traceId: input.traceId ?? null, since: finiteNumber(input.since) })], historyRoot: (): EntryKey => ["workbench", "history"], historySessionMessagesRoot: (): EntryKey => ["workbench", "history", "session-messages"], historySessionMessages: (sessionId: string, input: WorkbenchSessionMessagesKeyInput = {}): EntryKey => ["workbench", "history", "session-messages", sessionId, normalizeKeyObject({ cursor: input.cursor ?? null, limit: finiteNumber(input.limit) })], diff --git a/web/hwlab-cloud-web/src/stores/workbench-colada-mutations.ts b/web/hwlab-cloud-web/src/stores/workbench-colada-mutations.ts index 5ee2ce9d..d6340f82 100644 --- a/web/hwlab-cloud-web/src/stores/workbench-colada-mutations.ts +++ b/web/hwlab-cloud-web/src/stores/workbench-colada-mutations.ts @@ -69,7 +69,6 @@ export function useWorkbenchColadaMutations(): WorkbenchColadaMutations { async function invalidateMutationScope(queryCache: ReturnType, sessionId: string | null | undefined, traceId: string | null | undefined): Promise { await Promise.all([ - queryCache.invalidateQueries({ key: workbenchColadaKeys.syncRoot() }), queryCache.invalidateQueries({ key: workbenchColadaKeys.snapshotRoot() }), sessionId ? queryCache.invalidateQueries({ key: workbenchColadaKeys.historySessionMessagesRoot().concat(sessionId) }) : Promise.resolve(), traceId ? queryCache.invalidateQueries({ key: workbenchColadaKeys.detailRoot() }) : Promise.resolve() diff --git a/web/hwlab-cloud-web/src/stores/workbench-colada-queries.ts b/web/hwlab-cloud-web/src/stores/workbench-colada-queries.ts index a00e56b2..5b992c37 100644 --- a/web/hwlab-cloud-web/src/stores/workbench-colada-queries.ts +++ b/web/hwlab-cloud-web/src/stores/workbench-colada-queries.ts @@ -4,7 +4,6 @@ import { useQueryCache, type EntryKey, type QueryCache } from "@pinia/colada"; import { api } from "@/api"; import type { ApiRequestOptions } from "@/api/client"; -import { fetchWorkbenchSyncReplay, type WorkbenchSyncReplayRequest, type WorkbenchSyncReplayResponse } from "@/api/workbench-events"; import type { AgentChatResultResponse, ApiResult } from "@/types"; import type { SessionListOptions, SessionMessageOptions, WorkbenchMessagePageResponse, WorkbenchSessionDetailResponse, WorkbenchSessionListResponse, WorkbenchTraceRequestOptions } from "@/api/workbench"; import { workbenchColadaKeys } from "./workbench-colada-keys"; @@ -33,7 +32,6 @@ export interface WorkbenchColadaQueries { queryCache: QueryCache; fetchSessions: (options?: SessionListOptions & QueryRunOptions) => Promise>; fetchSession: (sessionId: string, options?: WorkbenchSessionDetailQueryOptions) => Promise>; - fetchSyncReplay: (input: WorkbenchSyncReplayRequest, options?: ApiRequestOptions & QueryRunOptions) => Promise>; fetchSessionMessages: (sessionId: string, options?: SessionMessageOptions & QueryRunOptions) => Promise>; fetchTurn: (traceId: string, options?: WorkbenchTurnQueryOptions) => Promise>; fetchTraceEvents: (traceId: string, options?: WorkbenchTraceEventsQueryOptions) => Promise>; @@ -53,7 +51,6 @@ export function useWorkbenchColadaQueries(): WorkbenchColadaQueries { queryCache, fetchSessions: (options = {}) => runWorkbenchQuery(queryCache, workbenchColadaKeys.sessions(options), () => api.workbench.sessions(options), options), fetchSession: (sessionId, options = {}) => runWorkbenchQuery(queryCache, workbenchColadaKeys.session(sessionId), () => api.workbench.session(sessionId, { timeoutMs: options.timeoutMs ?? null, includeMessages: false }), options), - fetchSyncReplay: (input, options = {}) => runWorkbenchQuery(queryCache, workbenchColadaKeys.syncReplay(input), () => fetchWorkbenchSyncReplay(input, options), options), fetchSessionMessages: (sessionId, options = {}) => runWorkbenchQuery(queryCache, workbenchColadaKeys.sessionMessages(sessionId, options), () => api.workbench.sessionMessages(sessionId, options), options), fetchTurn: (traceId, options = {}) => runWorkbenchQuery(queryCache, workbenchColadaKeys.turn(traceId), () => api.workbench.turn(traceId, options.timeoutMs ?? 8000, options.activityRef), options), fetchTraceEvents: (traceId, options = {}) => runWorkbenchQuery(queryCache, workbenchColadaKeys.traceEvents(traceId, options), () => api.workbench.traceEvents(traceId, options.timeoutMs ?? 8000, options.activityRef, { afterProjectedSeq: options.afterProjectedSeq, limit: options.limit }), options), @@ -64,7 +61,6 @@ export function useWorkbenchColadaQueries(): WorkbenchColadaQueries { invalidateTurn: (traceId) => traceId ? queryCache.invalidateQueries({ key: workbenchColadaKeys.turn(traceId) }) : Promise.resolve(), invalidateTraceEvents: (traceId) => traceId ? queryCache.invalidateQueries({ key: workbenchColadaKeys.traceEventsRoot().concat(traceId) }) : Promise.resolve(), invalidateTraceScope: async (input) => Promise.all([ - queryCache.invalidateQueries({ key: workbenchColadaKeys.syncRoot() }), queryCache.invalidateQueries({ key: workbenchColadaKeys.snapshotRoot() }), input.sessionId ? queryCache.invalidateQueries({ key: workbenchColadaKeys.historySessionMessagesRoot().concat(input.sessionId) }) : Promise.resolve(), input.traceId ? queryCache.invalidateQueries({ key: workbenchColadaKeys.detailRoot() }) : Promise.resolve() diff --git a/web/hwlab-cloud-web/src/stores/workbench-colada.test.ts b/web/hwlab-cloud-web/src/stores/workbench-colada.test.ts index 185eb3d2..2c77d786 100644 --- a/web/hwlab-cloud-web/src/stores/workbench-colada.test.ts +++ b/web/hwlab-cloud-web/src/stores/workbench-colada.test.ts @@ -7,12 +7,9 @@ import { createApp } from "vue"; import { createPinia, setActivePinia } from "pinia"; import { PiniaColada, useQueryCache } from "@pinia/colada"; -import { createWorkbenchServerState, reduceWorkbenchServerState, selectActiveMessages, selectTurnStatusAuthority, type WorkbenchServerState } from "./workbench-server-state"; +import { createWorkbenchServerState, type WorkbenchServerState } from "./workbench-server-state"; import { workbenchColadaKeys } from "./workbench-colada-keys"; import { useWorkbenchColadaReducer } from "./workbench-colada-reducer"; -import { workbenchSyncReplayDiagnostic, workbenchSyncReplayEvents } from "./workbench-realtime-authority"; -import { reduceWorkbenchRealtimeEvent } from "./workbench-event-reducer"; -import type { TurnStatusAuthority } from "./workbench-session"; const storeDir = path.dirname(fileURLToPath(import.meta.url)); @@ -27,7 +24,6 @@ function installColadaContext(): ReturnType { test("workbench colada keys keep serializable hierarchical cache scopes", () => { assert.deepEqual(workbenchColadaKeys.sessions({ limit: 50, includeSessionId: "ses_colada", cursor: "cur_2" }), ["workbench", "snapshot", "sessions", { cursor: "cur_2", includeSessionId: "ses_colada", limit: 50 }]); - assert.deepEqual(workbenchColadaKeys.syncReplay({ sessionId: "ses_colada", traceId: "trc_colada", since: 42 }), ["workbench", "sync", "replay", { sessionId: "ses_colada", since: 42, traceId: "trc_colada" }]); assert.deepEqual(workbenchColadaKeys.sessionMessages("ses_colada", { limit: 20 }), ["workbench", "history", "session-messages", "ses_colada", { cursor: null, limit: 20 }]); assert.deepEqual(workbenchColadaKeys.traceEvents("trc_colada", { afterProjectedSeq: 4, limit: 25 }), ["workbench", "detail", "trace-events", "trc_colada", { afterProjectedSeq: 4, limit: 25 }]); }); @@ -55,95 +51,17 @@ test("workbench colada read path preserves freshness governance", () => { assert.doesNotMatch(source, /queryCache\.fetch\(entry/u); }); -test("workbench colada sync replay is the only live repair query scope", () => { +test("workbench colada has no automatic HTTP live repair query scope", () => { const querySource = fs.readFileSync(path.join(storeDir, "workbench-colada-queries.ts"), "utf8"); const mutationSource = fs.readFileSync(path.join(storeDir, "workbench-colada-mutations.ts"), "utf8"); const storeSource = fs.readFileSync(path.join(storeDir, "workbench.ts"), "utf8"); - assert.match(querySource, /fetchSyncReplay: \(input, options = \{\}\) => runWorkbenchQuery\(queryCache, workbenchColadaKeys\.syncReplay\(input\),/u); - assert.match(storeSource, /workbenchColadaQueries\.fetchSyncReplay\(\{ sessionId, traceId, since: cursor \}/u); - assert.match(storeSource, /responseCursor\?\.hasMore !== true/u); - assert.doesNotMatch(storeSource, /import \{ fetchWorkbenchSyncReplay \} from "@\/api\/workbench-events"/u); + assert.doesNotMatch(querySource, /fetchSyncReplay|syncReplay|\/v1\/workbench\/sync/u); + assert.doesNotMatch(storeSource, /fetchSyncReplay|refreshWorkbenchSyncReplay|\/v1\/workbench\/sync/u); - assert.match(mutationSource, /workbenchColadaKeys\.syncRoot\(\)/u); + assert.doesNotMatch(mutationSource, /workbenchColadaKeys\.syncRoot\(\)/u); assert.match(mutationSource, /workbenchColadaKeys\.snapshotRoot\(\)/u); assert.match(mutationSource, /workbenchColadaKeys\.historySessionMessagesRoot\(\)/u); assert.match(mutationSource, /workbenchColadaKeys\.detailRoot\(\)/u); assert.doesNotMatch(mutationSource, /workbenchColadaKeys\.traceEventsRoot\(\)\.concat/u); }); - -test("workbench sync replay diagnostic exposes authority and terminal seal metadata", () => { - const payload = { - contractVersion: "workbench-sync-v1", - realtimeAuthority: "workbench-realtime-authority-v2", - cursor: { outboxSeq: 12 }, - events: [{ - type: "workbench.turn.snapshot", - realtimeAuthority: "workbench-realtime-authority-v2", - entity: { family: "turn", id: "trc_sync", version: 3, outboxSeq: 12, projectionRevision: "rev_sync" }, - turn: { traceId: "trc_sync", status: "completed", terminal: true } - }] - }; - - const events = workbenchSyncReplayEvents(payload); - const diagnostic = workbenchSyncReplayDiagnostic(payload, events, { reason: "sse-gap", sinceOutboxSeq: 7 }); - - assert.deepEqual(diagnostic, { - code: "workbench_sync_replay_applied", - reason: "sse-gap", - eventCount: 1, - sinceOutboxSeq: 7, - syncCursorOutboxSeq: 12, - contractVersion: "workbench-sync-v1", - realtimeAuthority: "workbench-realtime-authority-v2", - entityFamilyCount: 1, - entityFamilies: "turn", - maxEntityVersion: 3, - projectionRevision: "rev_sync", - terminalSeal: true, - detailProjection: false, - valuesRedacted: true - }); -}); - -test("workbench sync replay projects durable message and turn facts into observer state", () => { - const sessionId = "ses_sync_cross_page"; - const traceId = "trc_sync_cross_page"; - const payload = { - contractVersion: "workbench-sync-v1", - realtimeAuthority: "workbench-realtime-authority-v2", - cursor: { outboxSeq: 31 }, - events: [], - delta: { - messages: [ - { messageId: "msg_sync_cross_page_user", sessionId, traceId, turnId: traceId, role: "user", status: "sent", text: "run it", projectedSeq: 10, updatedAt: "2026-07-08T12:40:00.000Z", valuesRedacted: true }, - { messageId: "msg_sync_cross_page_agent", sessionId, traceId, turnId: traceId, role: "agent", status: "running", text: "", projectedSeq: 11, updatedAt: "2026-07-08T12:40:01.000Z", valuesRedacted: true } - ], - turns: [ - { turnId: traceId, sessionId, traceId, messageId: "msg_sync_cross_page_agent", status: "running", running: true, terminal: false, projectedSeq: 11, updatedAt: "2026-07-08T12:40:01.000Z", valuesRedacted: true } - ] - } - }; - - const events = workbenchSyncReplayEvents(payload); - assert.deepEqual(events.map((event) => event.type), ["message.snapshot", "message.snapshot", "turn.snapshot"]); - assert.equal(events.every((event) => reduceWorkbenchRealtimeEvent(event, `workbench.${event.type}`).action.type !== "ignore"), true); - - let state = createWorkbenchServerState(); - state = reduceWorkbenchServerState(state, { type: "session.detail", session: { sessionId, status: "running", messages: [] } }); - for (const event of events) { - const action = reduceWorkbenchRealtimeEvent(event, `workbench.${event.type}`).action; - if (action.type === "message.snapshot") { - state = reduceWorkbenchServerState(state, { type: "message.snapshot", sessionId: action.realtimeEvent.sessionId, message: action.realtimeEvent.message }); - } else if (action.type === "turn.snapshot") { - state = reduceWorkbenchServerState(state, { type: "turn.status", turn: action.turn as unknown as TurnStatusAuthority }); - } - } - - const messages = selectActiveMessages(state, sessionId); - assert.equal(messages.length, 2); - assert.equal(messages[0]?.role, "user"); - assert.equal(messages[1]?.role, "agent"); - assert.equal(messages[1]?.traceId, traceId); - assert.equal(selectTurnStatusAuthority(state)[traceId]?.status, "running"); -}); diff --git a/web/hwlab-cloud-web/src/stores/workbench-message-projection-runtime.test.ts b/web/hwlab-cloud-web/src/stores/workbench-message-projection-runtime.test.ts index bb1918a8..df1eea60 100644 --- a/web/hwlab-cloud-web/src/stores/workbench-message-projection-runtime.test.ts +++ b/web/hwlab-cloud-web/src/stores/workbench-message-projection-runtime.test.ts @@ -158,9 +158,8 @@ test("workbench active terminal paths seal final response from turn authority", assert.match(completeBlock, /projectTurnAuthorityToMessages\(traceId, result, "complete-trace"\)/u); assert.match(completeBlock, /options\.forceRead[\s\S]*refreshMessageProjectionForTrace\(ownerSessionId, traceId, \{ force: true \}\)/u); assert.doesNotMatch(completeBlock, /scheduleActiveTraceSyncReplay|refreshWorkbenchSyncReplay/u); - assert.match(crossTabSyncBlock, /traceIdFromRealtimeRefreshReason\(reason\)/u); - assert.match(crossTabSyncBlock, /refreshWorkbenchSyncReplay\(id, traceId, null, `cross-tab-sync:\$\{reason\}`\)/u); - assert.doesNotMatch(source, /scheduleActiveTraceSyncReplay|refreshActiveTraceFromSyncReplay|active-sync-replay:repeat|complete-trace-sync-replay/u); + assert.match(source.slice(source.indexOf("function handleWorkbenchProjectionSignal"), source.indexOf("function applyTraceSnapshot")), /restartRealtime\("cross-tab-session-projection", true\)/u); + assert.doesNotMatch(source, /scheduleActiveTraceSyncReplay|refreshActiveTraceFromSyncReplay/u); assert.match(traceDetailReadBlock, /traceAuthorityById\.value\[traceId\] \?\? message\.runnerTrace/u); assert.match(traceDetailReadBlock, /traceEventsDetailReadDecision\(traceId, afterProjectedSeq, message, options\)/u); assert.match(traceDetailReadBlock, /trace_events_auto_read_skip/u); @@ -176,7 +175,7 @@ test("workbench active terminal paths seal final response from turn authority", assert.doesNotMatch(restoreSealBlock, /force:\s*true/u); }); -test("workbench active SSE path never starts periodic sync replay", () => { +test("workbench active SSE path only recovers by reconnecting the events cursor", () => { const source = fs.readFileSync(path.join(storeDir, "workbench.ts"), "utf8"); const activeBlocks = [ source.slice(source.indexOf("async function submitMessage"), source.indexOf("async function cancelAgentMessage")), @@ -187,15 +186,15 @@ test("workbench active SSE path never starts periodic sync replay", () => { ].join("\n"); assert.doesNotMatch(activeBlocks, /refreshWorkbenchSyncReplay|scheduleActiveTraceSyncReplay|setTimeout|setInterval/u); - assert.match(source.slice(source.indexOf("function executeRealtimeRecoveryStep"), source.indexOf("function stopRealtime")), /case "sync-replay"[\s\S]*refreshWorkbenchSyncReplay/u); - assert.match(source.slice(source.indexOf("async function refreshCrossTabSessionFromSyncReplay"), source.indexOf("function completeTrace")), /refreshWorkbenchSyncReplay/u); + assert.match(source.slice(source.indexOf("function executeRealtimeRecoveryStep"), source.indexOf("function stopRealtime")), /case "events-reconnect"[\s\S]*restartRealtime\(step\.reason, true\)/u); + assert.doesNotMatch(source, /refreshWorkbenchSyncReplay|fetchSyncReplay|\/v1\/workbench\/sync/u); }); -test("cross page projection signal always replays durable sync", () => { +test("cross page projection signal reconnects the durable projection SSE", () => { const source = fs.readFileSync(path.join(storeDir, "workbench.ts"), "utf8"); const signalBlock = source.slice(source.indexOf("function handleWorkbenchProjectionSignal"), source.indexOf("function applyTraceSnapshot")); - assert.match(signalBlock, /refreshCrossTabSessionFromSyncReplay\(sessionId, `cross-tab-session-projection:/u); + assert.match(signalBlock, /restartRealtime\("cross-tab-session-projection", true\)/u); assert.doesNotMatch(signalBlock, /messages\.value\.some/u); - assert.doesNotMatch(signalBlock, /return;\s*\n\s*void refreshCrossTabSessionFromSyncReplay/u); + assert.doesNotMatch(signalBlock, /refreshCrossTabSessionFromSyncReplay|refreshWorkbenchSyncReplay/u); }); diff --git a/web/hwlab-cloud-web/src/stores/workbench-realtime-authority.ts b/web/hwlab-cloud-web/src/stores/workbench-realtime-authority.ts index 075931b7..10cfac9f 100644 --- a/web/hwlab-cloud-web/src/stores/workbench-realtime-authority.ts +++ b/web/hwlab-cloud-web/src/stores/workbench-realtime-authority.ts @@ -1,11 +1,9 @@ // SPEC: PJ2026-010401080313 Workbench实时权威 draft-2026-07-08-p0-workbench-realtime-authority-v2. -// Responsibility: Frontend authority gates for Workbench realtime events and sync replay batches. +// Responsibility: Frontend authority gates for Workbench projection outbox realtime and replay SSE events. -import type { WorkbenchRealtimeEvent, WorkbenchSyncReplayResponse } from "@/api/workbench-events"; -import type { ChatMessage } from "@/types"; +import type { WorkbenchRealtimeEvent } from "@/api/workbench-events"; export const WORKBENCH_REALTIME_AUTHORITY_VERSION = "workbench-realtime-authority-v2"; -export const WORKBENCH_SYNC_CONTRACT_VERSION = "workbench-sync-v1"; export interface WorkbenchRealtimeEntityAuthority { family: string; @@ -82,59 +80,6 @@ export function workbenchRealtimeEntityAuthority(event: WorkbenchRealtimeEvent): }; } -export function workbenchSyncReplayEvents(payload: WorkbenchSyncReplayResponse | null | undefined): WorkbenchRealtimeEvent[] { - const events = arrayOfRecords(payload?.events); - const delta = workbenchSyncReplayDeltaEvents(payload); - const seen = new Set(); - const output: WorkbenchRealtimeEvent[] = []; - for (const value of [...events, ...delta]) { - const event = value as WorkbenchRealtimeEvent; - const key = syncEventKey(event); - if (seen.has(key)) continue; - seen.add(key); - output.push(event); - } - return output; -} - -function workbenchSyncReplayDeltaEvents(payload: WorkbenchSyncReplayResponse | null | undefined): WorkbenchRealtimeEvent[] { - const explicitEvents = arrayOfRecords(payload?.delta); - const delta = recordValue(payload?.delta); - const replayEvents = explicitEvents as WorkbenchRealtimeEvent[]; - if (!delta || Array.isArray(payload?.delta)) return replayEvents; - const messages = arrayOfRecords(delta.messages).map((message) => syncReplayMessageSnapshotEvent(message, payload)).filter((event): event is WorkbenchRealtimeEvent => Boolean(event)); - const turns = arrayOfRecords(delta.turns).map((turn) => syncReplayTurnSnapshotEvent(turn, payload)).filter((event): event is WorkbenchRealtimeEvent => Boolean(event)); - return [...replayEvents, ...messages, ...turns]; -} - -export function workbenchSyncReplayDiagnostic(payload: WorkbenchSyncReplayResponse | null | undefined, events: WorkbenchRealtimeEvent[], input: { reason: string; sinceOutboxSeq?: number | null }): Record { - const eventRows = Array.isArray(events) ? events : []; - const cursor = recordValue(payload?.cursor); - const entities = eventRows.map(workbenchRealtimeEntityAuthority).filter((item): item is WorkbenchRealtimeEntityAuthority => item !== null); - const families = Array.from(new Set(entities.map((item) => item.family).filter(Boolean))).sort(); - const versions = entities.map((item) => item.version).filter((value) => Number.isFinite(value)); - const projectionRevision = firstString( - ...entities.map((item) => item.projectionRevision), - ...eventRows.map((event) => stringValue(event.projectionRevision ?? event.entity?.projectionRevision)) - ); - return { - code: "workbench_sync_replay_applied", - reason: input.reason, - eventCount: eventRows.length, - sinceOutboxSeq: finiteNumber(input.sinceOutboxSeq), - syncCursorOutboxSeq: finiteNumber(cursor?.outboxSeq ?? payload?.cursor?.outboxSeq), - contractVersion: stringValue(payload?.contractVersion), - realtimeAuthority: stringValue(payload?.realtimeAuthority), - entityFamilyCount: families.length, - entityFamilies: families.slice(0, 8).join(","), - maxEntityVersion: versions.length > 0 ? Math.max(...versions) : null, - projectionRevision, - terminalSeal: eventRows.some(workbenchRealtimeEventHasTerminalSeal), - detailProjection: eventRows.some((event) => workbenchRealtimeDetailOnly(event, workbenchRealtimeEntityAuthority(event))), - valuesRedacted: true - }; -} - function authorityFailureCode(event: WorkbenchRealtimeEvent, entity: WorkbenchRealtimeEntityAuthority | null, detailProjection: boolean): string | null { if (detailProjection) return "workbench_realtime_detail_only_rejected"; if (stringValue(event.realtimeAuthority) !== WORKBENCH_REALTIME_AUTHORITY_VERSION) return "workbench_realtime_authority_missing"; @@ -147,114 +92,6 @@ function workbenchRealtimeDetailOnly(event: WorkbenchRealtimeEvent, entity: Work return event.detailProjection === true || entity?.detailProjection === true || stringValue(event.authority) === "trace-detail-only" || entity?.authority === "trace-detail-only"; } -function syncEventKey(event: WorkbenchRealtimeEvent): string { - const entity = workbenchRealtimeEntityAuthority(event); - if (entity) return [entity.family, entity.id, entity.version, entity.outboxSeq ?? "~"].join("|"); - return [event.type ?? "~", event.traceId ?? "~", event.sessionId ?? "~", event.cursor?.outboxSeq ?? event.outboxSeq ?? "~", event.cursor?.traceSeq ?? event.traceSeq ?? "~"].join("|"); -} - -function syncReplayMessageSnapshotEvent(message: Record, payload: WorkbenchSyncReplayResponse | null | undefined): WorkbenchRealtimeEvent | null { - const sessionId = textValue(message.sessionId); - if (!sessionId) return null; - const traceId = textValue(message.traceId); - const role = textValue(message.role) ?? "agent"; - const id = textValue(message.messageId ?? message.id) ?? [traceId, role].filter(Boolean).join(":"); - if (!id) return null; - const entity = syncReplayEntity("messages", id, message, payload); - return { - type: "message.snapshot", - contractVersion: stringValue(payload?.contractVersion), - realtimeAuthority: stringValue(payload?.realtimeAuthority), - sessionId, - traceId, - reason: "sync-delta", - cursor: syncReplayCursor(message, payload), - entity, - projectionRevision: entity.projectionRevision, - message: message as ChatMessage, - valuesRedacted: true - }; -} - -function syncReplayTurnSnapshotEvent(turn: Record, payload: WorkbenchSyncReplayResponse | null | undefined): WorkbenchRealtimeEvent | null { - const traceId = textValue(turn.traceId); - const turnId = textValue(turn.turnId ?? traceId); - if (!traceId && !turnId) return null; - const entity = syncReplayEntity("turns", turnId ?? traceId ?? "turn", turn, payload); - return { - type: "turn.snapshot", - contractVersion: stringValue(payload?.contractVersion), - realtimeAuthority: stringValue(payload?.realtimeAuthority), - sessionId: textValue(turn.sessionId), - traceId, - reason: "sync-delta", - cursor: syncReplayCursor(turn, payload), - entity, - projectionRevision: entity.projectionRevision, - turn, - valuesRedacted: true - }; -} - -function syncReplayEntity(family: string, id: string, fact: Record, payload: WorkbenchSyncReplayResponse | null | undefined): NonNullable { - const traceSeq = finiteNumber(fact.projectedSeq ?? fact.sourceSeq ?? fact.seq); - const outboxSeq = finiteNumber(payload?.cursor?.outboxSeq ?? fact.outboxSeq); - const version = finiteNumber(fact.projectionRevision ?? fact.projectedSeq ?? fact.sourceSeq ?? fact.seq ?? outboxSeq) ?? 0; - const projectionRevision = textValue(fact.projectionRevision ?? fact.projectedSeq ?? fact.sourceSeq ?? fact.seq ?? version) ?? "0"; - return { - family, - id, - version, - entityVersion: version, - outboxSeq, - traceSeq, - projectionRevision, - serverCommittedAt: textValue(fact.updatedAt ?? fact.createdAt ?? fact.occurredAt), - authority: "workbench-sync-delta", - detailProjection: false - }; -} - -function syncReplayCursor(fact: Record, payload: WorkbenchSyncReplayResponse | null | undefined): NonNullable { - return { - traceSeq: finiteNumber(fact.projectedSeq ?? fact.sourceSeq ?? fact.seq), - outboxSeq: finiteNumber(payload?.cursor?.outboxSeq ?? fact.outboxSeq) - }; -} - -function workbenchRealtimeEventHasTerminalSeal(event: WorkbenchRealtimeEvent): boolean { - const turn = recordValue(event.turn); - if (turn?.terminal === true) return true; - const message = recordValue(event.message); - const snapshot = recordValue(event.snapshot); - const traceEvent = recordValue(event.event); - return [ - event.status, - turn?.status, - message?.status, - snapshot?.status, - traceEvent?.status, - traceEvent?.type, - ].some(isTerminalStatusText); -} - -function isTerminalStatusText(value: unknown): boolean { - const text = String(value ?? "").trim().toLowerCase().replace(/_/gu, "-"); - return /^(completed|complete|succeeded|success|failed|failure|error|canceled|cancelled|done|terminal|sealed|thread-resume-failed)$/u.test(text); -} - -function firstString(...values: unknown[]): string | null { - for (const value of values) { - const text = stringValue(value); - if (text) return text; - } - return null; -} - -function arrayOfRecords(value: unknown): Record[] { - return Array.isArray(value) ? value.filter((item): item is Record => Boolean(recordValue(item))) : []; -} - function recordValue(value: unknown): Record | null { return value && typeof value === "object" ? value as Record : null; } @@ -264,11 +101,6 @@ function stringValue(value: unknown): string | null { return text || null; } -function textValue(value: unknown): string | null { - const text = typeof value === "string" ? value.trim() : value === null || value === undefined ? "" : String(value).trim(); - return text || null; -} - function finiteNumber(value: unknown): number | null { const number = Number(value); return Number.isFinite(number) && number >= 0 ? Math.trunc(number) : null; diff --git a/web/hwlab-cloud-web/src/stores/workbench-realtime-plan.ts b/web/hwlab-cloud-web/src/stores/workbench-realtime-plan.ts index 03f00d50..8c70fb56 100644 --- a/web/hwlab-cloud-web/src/stores/workbench-realtime-plan.ts +++ b/web/hwlab-cloud-web/src/stores/workbench-realtime-plan.ts @@ -35,12 +35,8 @@ export interface WorkbenchRealtimeRecoveryContext { terminalTraceSealed?: boolean; } -export type WorkbenchRealtimeRecoveryAuthority = "automatic-recovery"; - export type WorkbenchRealtimeRecoveryStep = - | { type: "sync-replay"; sessionId: string | null; traceId: string | null; sinceOutboxSeq: number | null; reason: string; authority: WorkbenchRealtimeRecoveryAuthority }; - -type WorkbenchRealtimeRecoveryStepInput = T extends WorkbenchRealtimeRecoveryStep ? Omit : never; + | { type: "events-reconnect"; sessionId: string | null; traceId: string | null; afterOutboxSeq: number | null; reason: string }; export interface WorkbenchRealtimeRecoveryPlan { sessionId: string | null; @@ -73,7 +69,7 @@ export function planWorkbenchRealtimeRecovery(recovery: WorkbenchStreamTransport const sessionId = normalizeWorkbenchSessionId(recovery.sessionId ?? context.selectedSessionId); const traceId = firstNonEmptyString(recovery.traceId, context.fallbackTraceId); const steps: WorkbenchRealtimeRecoveryStep[] = []; - if (context.terminalTraceSealed !== true && recovery.actions.includes("sync-replay") && (sessionId || traceId)) steps.push(automaticRecoveryStep({ type: "sync-replay", sessionId: sessionId ?? null, traceId: traceId ?? null, sinceOutboxSeq: finiteNumber(recovery.outboxSeq), reason: "realtime-error:sync-replay" })); + if (recovery.actions.includes("events-reconnect") && (sessionId || traceId)) steps.push({ type: "events-reconnect", sessionId: sessionId ?? null, traceId: traceId ?? null, afterOutboxSeq: finiteNumber(recovery.outboxSeq), reason: "realtime-error:events-reconnect" }); return { sessionId: sessionId ?? null, traceId: traceId ?? null, @@ -88,10 +84,6 @@ export function planWorkbenchRealtimeRecovery(recovery: WorkbenchStreamTransport }; } -function automaticRecoveryStep(step: WorkbenchRealtimeRecoveryStepInput): WorkbenchRealtimeRecoveryStep { - return { ...step, authority: "automatic-recovery" } as WorkbenchRealtimeRecoveryStep; -} - function finiteNumber(value: unknown): number | null { const number = Number(value); return Number.isFinite(number) && number >= 0 ? Math.trunc(number) : null; diff --git a/web/hwlab-cloud-web/src/stores/workbench.ts b/web/hwlab-cloud-web/src/stores/workbench.ts index cd25c89b..582724da 100644 --- a/web/hwlab-cloud-web/src/stores/workbench.ts +++ b/web/hwlab-cloud-web/src/stores/workbench.ts @@ -70,7 +70,6 @@ import { traceResultHasTerminalEvidence, } from "./workbench-message-projection-runtime"; import { planWorkbenchRealtimeApply, planWorkbenchRealtimeRecovery, type WorkbenchRealtimeApplyStep, type WorkbenchRealtimeRecoveryStep } from "./workbench-realtime-plan"; -import { workbenchSyncReplayDiagnostic, workbenchSyncReplayEvents } from "./workbench-realtime-authority"; import { useWorkbenchColadaMutations } from "./workbench-colada-mutations"; import { useWorkbenchColadaQueries } from "./workbench-colada-queries"; import { useWorkbenchColadaReducer } from "./workbench-colada-reducer"; @@ -634,7 +633,7 @@ export const useWorkbenchStore = defineStore("workbench", () => { module: "workbench-turn-status", traceId: options.traceId ?? traceIds[0] ?? null, outcome: "ok", - diagnostic: { code: "workbench_turn_status_auto_read_disabled", reason: options.reason ?? "auto-turn-status-hydrate", source: "sync-replay-authority", requestedCount: traceIds.length, valuesRedacted: true } + diagnostic: { code: "workbench_turn_status_auto_read_disabled", reason: options.reason ?? "auto-turn-status-hydrate", source: "projection-sse-authority", requestedCount: traceIds.length, valuesRedacted: true } }); } return; @@ -667,7 +666,7 @@ export const useWorkbenchStore = defineStore("workbench", () => { const id = firstNonEmptyString(traceId); if (!id) return; if (options.force !== true) { - recordWorkbenchRuntimeDiagnostic({ module: "workbench-turn-status", traceId: id, outcome: "ok", diagnostic: { code: "workbench_turn_status_auto_read_disabled", source: "sync-replay-authority", valuesRedacted: true } }); + recordWorkbenchRuntimeDiagnostic({ module: "workbench-turn-status", traceId: id, outcome: "ok", diagnostic: { code: "workbench_turn_status_auto_read_disabled", source: "projection-sse-authority", valuesRedacted: true } }); return; } if (traceTerminalBodyIsVisible(id)) { @@ -890,7 +889,7 @@ export const useWorkbenchStore = defineStore("workbench", () => { const traceId = message.traceId ?? message.runnerTrace?.traceId; if (!traceId) return; if (options.force !== true) { - recordWorkbenchRuntimeDiagnostic({ module: "workbench-trace-events-read", sessionId: message.sessionId ?? message.runnerTrace?.sessionId ?? null, traceId, outcome: "ok", diagnostic: { code: "trace_events_auto_read_disabled", reason: "sync-replay-authority", source: "trace-detail-read", valuesRedacted: true } }); + recordWorkbenchRuntimeDiagnostic({ module: "workbench-trace-events-read", sessionId: message.sessionId ?? message.runnerTrace?.sessionId ?? null, traceId, outcome: "ok", diagnostic: { code: "trace_events_auto_read_disabled", reason: "projection-sse-authority", source: "trace-detail-read", valuesRedacted: true } }); return; } await readTraceEventsForMessageNow(message, options); @@ -979,7 +978,7 @@ export const useWorkbenchStore = defineStore("workbench", () => { if (options.force !== true) { if (candidates.length > 0) { const traceId = firstNonEmptyString(candidates[0]?.traceId, candidates[0]?.runnerTrace?.traceId); - recordWorkbenchRuntimeDiagnostic({ module: "workbench-trace-events-read", sessionId: candidates[0]?.sessionId ?? candidates[0]?.runnerTrace?.sessionId ?? null, traceId: traceId ?? null, outcome: "ok", diagnostic: { code: "trace_events_auto_read_disabled", reason: "sync-replay-authority", source: "trace-detail-read", candidateCount: candidates.length, valuesRedacted: true } }); + recordWorkbenchRuntimeDiagnostic({ module: "workbench-trace-events-read", sessionId: candidates[0]?.sessionId ?? candidates[0]?.runnerTrace?.sessionId ?? null, traceId: traceId ?? null, outcome: "ok", diagnostic: { code: "trace_events_auto_read_disabled", reason: "projection-sse-authority", source: "trace-detail-read", candidateCount: candidates.length, valuesRedacted: true } }); } return; } @@ -1123,7 +1122,7 @@ export const useWorkbenchStore = defineStore("workbench", () => { restartRealtime("reattach"); } - function restartRealtime(reason: string): void { + function restartRealtime(reason: string, forceReconnect = false): void { const traceId = realtimeTraceId(); const sessionId = selectedSessionId.value ?? null; void reason; @@ -1135,6 +1134,7 @@ export const useWorkbenchStore = defineStore("workbench", () => { realtimeCapabilities, sessionId, traceId, + forceReconnect, errorRecoveryMinMs: runtimePolicy.workbenchRealtimeErrorSyncReplayMinMs, flushMaxItemsPerChunk: runtimePolicy.workbenchRealtimeFlushMaxItemsPerChunk, flushMaxChunkMs: runtimePolicy.workbenchRealtimeFlushMaxChunkMs, @@ -1186,35 +1186,12 @@ export const useWorkbenchStore = defineStore("workbench", () => { function executeRealtimeRecoveryStep(step: WorkbenchRealtimeRecoveryStep): void { switch (step.type) { - case "sync-replay": - void refreshWorkbenchSyncReplay(step.sessionId, step.traceId, step.sinceOutboxSeq, step.reason); + case "events-reconnect": + restartRealtime(step.reason, true); return; } } - async function refreshWorkbenchSyncReplay(sessionId: string | null, traceId: string | null, sinceOutboxSeq: number | null, reason: string): Promise { - if (!realtimeCapabilities.projectionRealtime) return; - let cursor = firstFiniteNumber(sinceOutboxSeq) ?? 0; - for (;;) { - const result = await workbenchColadaQueries.fetchSyncReplay({ sessionId, traceId, since: cursor }, { timeoutMs: 8000, activityRef: () => activityRef.value }); - if (!result.ok || !result.data) { - recordWorkbenchRuntimeDiagnostic({ module: "workbench-sync-replay", sessionId, traceId, outcome: "network", diagnostic: { code: "workbench_sync_replay_failed", reason, status: result.status, apiError: result.apiError, valuesRedacted: true } }); - return; - } - const events = workbenchSyncReplayEvents(result.data); - for (const event of events) applyRealtimeEvent(event, realtimeEventName(event)); - recordWorkbenchRuntimeDiagnostic({ module: "workbench-sync-replay", sessionId, traceId, outcome: "ok", diagnostic: workbenchSyncReplayDiagnostic(result.data, events, { reason, sinceOutboxSeq: cursor }) }); - const responseCursor = recordValue(result.data.cursor); - if (responseCursor?.hasMore !== true) return; - const nextCursor = firstFiniteNumber(responseCursor.outboxSeq); - if (nextCursor === undefined || nextCursor <= cursor) { - recordWorkbenchRuntimeDiagnostic({ module: "workbench-sync-replay", sessionId, traceId, outcome: "error", diagnostic: { code: "workbench_sync_replay_cursor_stalled", reason, cursor, nextCursor, valuesRedacted: true } }); - return; - } - cursor = nextCursor; - } - } - function stopRealtime(): void { liveRealtimeReadySessionId.value = null; realtimeTransport.stop(); @@ -1553,7 +1530,8 @@ export const useWorkbenchStore = defineStore("workbench", () => { const sessionId = normalizeWorkbenchSessionId(record.sessionId); if (!sessionId || sessionId !== activeSessionId.value) return; const traceId = firstNonEmptyString(record.traceId); - void refreshCrossTabSessionFromSyncReplay(sessionId, `cross-tab-session-projection:${firstNonEmptyString(record.reason, traceId, "session") ?? "session"}`); + recordActivity(`cross-tab-session-projection:${firstNonEmptyString(record.reason, traceId, "session") ?? "session"}`); + restartRealtime("cross-tab-session-projection", true); } function applyTraceSnapshot(traceId: string, snapshot: TraceSnapshot, canonicalSession = false): void { @@ -1582,18 +1560,6 @@ export const useWorkbenchStore = defineStore("workbench", () => { if (!traceProjectionIsTerminalSealed(traceId, serverState.value.messagesBySessionId[ownerSessionId] ?? [])) scheduleSessionListRefresh(ownerSessionId, runtimePolicy.sessionListRealtimeRefreshDelayMs); } - async function refreshCrossTabSessionFromSyncReplay(sessionId: string, reason: string): Promise { - const id = normalizeWorkbenchSessionId(sessionId); - if (!id || id !== activeSessionId.value) return; - const traceId = traceIdFromRealtimeRefreshReason(reason); - if (traceTerminalBodyIsVisible(traceId, id)) { - recordWorkbenchRuntimeDiagnostic({ module: "workbench-terminal-priority", sessionId: id, traceId, outcome: "ok", diagnostic: { code: "terminal_low_priority_session_sync_skip", reason, source: "realtime-session-sync", valuesRedacted: true } }); - return; - } - recordActivity(reason); - await refreshWorkbenchSyncReplay(id, traceId, null, `cross-tab-sync:${reason}`); - } - function completeTrace(traceId: string, result: AgentChatResultResponse, options: { forceRead?: boolean } = {}): void { const authoritySessionId = traceResultSessionId(result); const ownerSessionId = traceOwnerSessionId(traceId, authoritySessionId); @@ -1886,7 +1852,7 @@ export const useWorkbenchStore = defineStore("workbench", () => { sessionId: message.sessionId ?? message.runnerTrace?.sessionId ?? null, traceId, outcome: "ok", - diagnostic: { code: "workbench_restored_turn_auto_read_disabled", reason: "load-session:restored-turn", source: "sync-replay-authority", valuesRedacted: true } + diagnostic: { code: "workbench_restored_turn_auto_read_disabled", reason: "load-session:restored-turn", source: "projection-sse-authority", valuesRedacted: true } }); } return source; diff --git a/web/hwlab-cloud-web/src/utils/workbench-stream-transport.ts b/web/hwlab-cloud-web/src/utils/workbench-stream-transport.ts index 8fe49afa..39b5ea1d 100644 --- a/web/hwlab-cloud-web/src/utils/workbench-stream-transport.ts +++ b/web/hwlab-cloud-web/src/utils/workbench-stream-transport.ts @@ -18,6 +18,7 @@ export interface WorkbenchStreamTransportRestartInput { sessionId?: string | null; traceId?: string | null; afterSeq?: number | null; + forceReconnect?: boolean; errorRecoveryMinMs?: number | null; flushMaxItemsPerChunk?: number | null; flushMaxChunkMs?: number | null; @@ -38,7 +39,7 @@ export interface WorkbenchStreamTransportRestartResult { export type WorkbenchStreamTransportPhase = "idle" | "connecting" | "open" | "event" | "error" | "closed" | "blocked"; -export type WorkbenchStreamTransportRecoveryAction = "sync-replay"; +export type WorkbenchStreamTransportRecoveryAction = "events-reconnect"; export interface WorkbenchStreamTransportState { key: string; @@ -102,7 +103,7 @@ export class WorkbenchStreamTransportRuntime { restart(input: WorkbenchStreamTransportRestartInput): WorkbenchStreamTransportRestartResult { const key = workbenchRealtimeScopeKey(input.sessionId ?? null, input.traceId ?? null); - if (key === this.key && this.stream) return { key, changed: false, streamStarted: true }; + if (key === this.key && this.stream && input.forceReconnect !== true) return { key, changed: false, streamStarted: true }; this.stop(); this.key = key; this.armWait(key); @@ -113,7 +114,7 @@ export class WorkbenchStreamTransportRuntime { realtimeCapabilities: input.realtimeCapabilities, sessionId: input.sessionId ?? null, traceId: input.traceId ?? null, - afterSeq: !input.realtimeCapabilities.liveKafkaSse && input.realtimeCapabilities.projectionRealtime ? input.afterSeq ?? this.cursorByKey.get(key)?.outboxSeq ?? null : null, + afterSeq: input.afterSeq ?? this.cursorByKey.get(key)?.outboxSeq ?? null, flushMaxItemsPerChunk: input.flushMaxItemsPerChunk, flushMaxChunkMs: input.flushMaxChunkMs, flushYieldMs: input.flushYieldMs, @@ -130,7 +131,7 @@ export class WorkbenchStreamTransportRuntime { this.requestRecovery(input, event.type || "eventsource-error"); }, onEvent: (event, eventName) => { - if (!input.realtimeCapabilities.liveKafkaSse && input.realtimeCapabilities.projectionRealtime) this.rememberCursor(key, event); + this.rememberCursor(key, event); this.markLive(key); this.emitState(input, "event", eventName); input.onEvent(event, eventName); @@ -201,7 +202,7 @@ export class WorkbenchStreamTransportRuntime { const sessionId = input.sessionId ?? null; const traceId = input.traceId ?? null; const cursor = this.currentCursor(key); - const actions: WorkbenchStreamTransportRecoveryAction[] = !input.realtimeCapabilities.liveKafkaSse && input.realtimeCapabilities.projectionRealtime && (sessionId || traceId) ? ["sync-replay"] : []; + const actions: WorkbenchStreamTransportRecoveryAction[] = input.realtimeCapabilities.projectionRealtime && (sessionId || traceId) ? ["events-reconnect"] : []; const diagnostic = this.diagnosticEnvelope("workbench_sse_recovery", reason, key, actions); input.onRecovery?.({ key, sessionId, traceId, tick: this.wait?.tick ?? this.tick, reason, actions, outboxSeq: cursor.outboxSeq, traceSeq: cursor.traceSeq, diagnostic }); }