diff --git a/internal/cloud/kafka-event-bridge.test.ts b/internal/cloud/kafka-event-bridge.test.ts index dc0b0753..1077c4f0 100644 --- a/internal/cloud/kafka-event-bridge.test.ts +++ b/internal/cloud/kafka-event-bridge.test.ts @@ -6,7 +6,22 @@ import { test } from "bun:test"; import { createCloudApiBunServer } from "./bun-server.ts"; import { buildCloudApiReadiness } from "./health-contract.ts"; -import { decodeCanonicalAgentRunKafkaMessage, kafkaDnsLookup, kafkaEventBridgeConfig, projectAgentRunKafkaEventToHwlabEvent, publishAgentRunKafkaMessageLive, relayHwlabKafkaOutboxOnce, startHwlabKafkaEventBridge } from "./kafka-event-bridge.ts"; +import { decodeCanonicalAgentRunKafkaMessage, kafkaDnsLookup, kafkaEventBridgeConfig, kafkaPartitionForKey, projectAgentRunKafkaEventToHwlabEvent, publishAgentRunKafkaMessageLive, relayHwlabKafkaOutboxOnce, startHwlabKafkaEventBridge } from "./kafka-event-bridge.ts"; + +test("Workbench refresh replay resolves the same Kafka partition as the session-keyed producer", () => { + const offsets = [ + { partition: 0, startOffset: "0", endOffset: "100" }, + { partition: 1, startOffset: "0", endOffset: "200" }, + { partition: 2, startOffset: "0", endOffset: "300" } + ]; + const sessionId = "ses_refresh_partition_scope"; + const first = kafkaPartitionForKey(sessionId, "hwlab.event.v1", offsets); + const repeated = kafkaPartitionForKey(sessionId, "hwlab.event.v1", [...offsets].reverse()); + + assert.ok([0, 1, 2].includes(first)); + assert.equal(repeated, first); + assert.equal(kafkaPartitionForKey(null, "hwlab.event.v1", offsets), null); +}); const PROJECTOR_ENV = Object.freeze({ HWLAB_KAFKA_DIRECT_PUBLISH_ENABLED: "false", diff --git a/internal/cloud/kafka-event-bridge.ts b/internal/cloud/kafka-event-bridge.ts index 3b360673..0979ea1f 100644 --- a/internal/cloud/kafka-event-bridge.ts +++ b/internal/cloud/kafka-event-bridge.ts @@ -7,7 +7,7 @@ import { isIP } from "node:net"; import net from "node:net"; import tls from "node:tls"; -import { Kafka, logLevel } from "kafkajs"; +import { Kafka, logLevel, Partitioners } from "kafkajs"; import { emitCodeAgentOtelSpan } from "./otel-trace.ts"; import { buildWorkbenchProjectionEventFacts } from "./workbench-projection-writer.ts"; import { @@ -958,7 +958,16 @@ function requireKafkaProjectorStore(runtimeStore, capabilities = {}) { if (missing.length > 0) throw contractError("hwlab_kafka_projector_store_invalid", `Kafka durable capabilities require a runtime store: ${missing.join(", ")}`); } -export async function queryKafkaEventStream({ env = process.env, stream = "hwlab", topic = null, traceId = null, sessionId = null, runId = null, commandId = null, limit = DEFAULT_QUERY_LIMIT, scanLimit = null, timeoutMs = DEFAULT_QUERY_TIMEOUT_MS, fromBeginning = true, groupIdPrefix = null, signal = null, kafkaFactory = defaultKafkaFactory } = {}) { +export async function queryKafkaEventStream({ env = process.env, stream = "hwlab", topic = null, traceId = null, sessionId = null, runId = null, commandId = null, partitionKey = null, limit = DEFAULT_QUERY_LIMIT, scanLimit = null, timeoutMs = DEFAULT_QUERY_TIMEOUT_MS, fromBeginning = true, groupIdPrefix = null, signal = null, kafkaFactory = defaultKafkaFactory } = {}) { + const queryStartedAtMs = Date.now(); + const timing = { + endOffsetSnapshotMs: 0, + groupOffsetsMs: 0, + consumerConnectMs: 0, + consumerSubscribeMs: 0, + scanMs: 0, + cleanupMs: 0 + }; const brokers = csv(env.HWLAB_KAFKA_BOOTSTRAP_SERVERS); if (brokers.length === 0) throw new Error("HWLAB_KAFKA_BOOTSTRAP_SERVERS is required for Kafka event queries."); const clientId = stringValue(env.HWLAB_KAFKA_CLIENT_ID) || DEFAULT_CLIENT_ID; @@ -975,14 +984,24 @@ export async function queryKafkaEventStream({ env = process.env, stream = "hwlab dnsServers: csv(env.HWLAB_KAFKA_DNS_SERVERS), dnsSearchDomains: csv(env.HWLAB_KAFKA_DNS_SEARCH_DOMAINS) }); + let stageStartedAtMs = Date.now(); const endOffsetSnapshot = await kafkaTopicEndOffsetSnapshot(kafka, resolvedTopic, deadlineMs, signal); + timing.endOffsetSnapshotMs = Date.now() - stageStartedAtMs; const endOffsetByPartition = new Map(endOffsetSnapshot.offsets.map((entry) => [entry.partition, entry.endOffset])); const startOffsetByPartition = new Map(endOffsetSnapshot.offsets.map((entry) => [entry.partition, entry.startOffset])); - const reachedEndPartitions = new Set(endOffsetSnapshot.offsets.filter(kafkaOffsetPartitionIsEmpty).map((entry) => entry.partition)); + const targetPartition = kafkaPartitionForKey(partitionKey, resolvedTopic, endOffsetSnapshot.offsets); + const skippedPartitions = targetPartition === null ? [] : endOffsetSnapshot.offsets.filter((entry) => entry.partition !== targetPartition).map((entry) => entry.partition); + const requiresScopedGroupOffsets = skippedPartitions.length > 0; + const reachedEndPartitions = new Set([...skippedPartitions, ...endOffsetSnapshot.offsets.filter(kafkaOffsetPartitionIsEmpty).map((entry) => entry.partition)]); const verifiedStartPartitions = new Set(reachedEndPartitions); const resolvedGroupIdPrefix = stringValue(groupIdPrefix) || `${clientId}-query`; const groupId = `${resolvedGroupIdPrefix}-${Date.now()}-${randomUUID().slice(0, 8)}`; - const consumer = kafka.consumer({ groupId, allowAutoTopicCreation: false }); + const consumer = kafka.consumer({ + groupId, + allowAutoTopicCreation: false, + minBytes: 1, + maxWaitTimeInMs: 100 + }); const events = []; const firstScannedOffsetByPartition = new Map(); const lastScannedOffsetByPartition = new Map(); @@ -1005,11 +1024,28 @@ export async function queryKafkaEventStream({ env = process.env, stream = "hwlab completionReason = "retention-start-unavailable"; return kafkaQueryResult(); } + if (requiresScopedGroupOffsets) { + stageStartedAtMs = Date.now(); + await prepareKafkaQueryGroupOffsets(kafka, { + groupId, + topic: resolvedTopic, + targetPartition, + offsets: endOffsetSnapshot.offsets, + deadlineMs, + signal + }); + timing.groupOffsetsMs = Date.now() - stageStartedAtMs; + } + stageStartedAtMs = Date.now(); await withinKafkaQueryBudget(consumer.connect(), deadlineMs, "consumer-connect", signal); - await withinKafkaQueryBudget(consumer.subscribe({ topic: resolvedTopic, fromBeginning: fromBeginning !== false }), deadlineMs, "consumer-subscribe", signal); + timing.consumerConnectMs = Date.now() - stageStartedAtMs; + stageStartedAtMs = Date.now(); + await withinKafkaQueryBudget(consumer.subscribe({ topic: resolvedTopic, fromBeginning: !requiresScopedGroupOffsets && fromBeginning !== false }), deadlineMs, "consumer-subscribe", signal); + timing.consumerSubscribeMs = Date.now() - stageStartedAtMs; const alreadyAtBoundedEnd = fromBeginning !== false && endOffsetSnapshot.available && reachedEndPartitions.size === endOffsetByPartition.size; if (alreadyAtBoundedEnd) completionReason = "end-offset"; if (alreadyAtBoundedEnd) return kafkaQueryResult(); + stageStartedAtMs = Date.now(); await new Promise((resolve, reject) => { let finished = false; const finish = (reason = "timeout") => { @@ -1027,64 +1063,79 @@ export async function queryKafkaEventStream({ env = process.env, stream = "hwlab } running = true; consumer.run({ - eachMessage: async ({ topic: messageTopic, partition, message }) => { - if (finished) return; - const endOffset = endOffsetByPartition.get(partition); - if (fromBeginning !== false && endOffsetSnapshot.available && !firstScannedOffsetByPartition.has(partition)) { - const actualStartOffset = stringValue(message.offset); - const expectedStartOffset = startOffsetByPartition.get(partition); - firstScannedOffsetByPartition.set(partition, actualStartOffset); - if (!endOffsetByPartition.has(partition) || !expectedStartOffset || actualStartOffset !== expectedStartOffset) { - finish("retention-start-moved"); - return; + eachBatchAutoResolve: true, + eachBatch: async ({ batch }) => { + for (const message of batch.messages) { + if (finished) return; + const messageTopic = batch.topic; + const partition = batch.partition; + const endOffset = endOffsetByPartition.get(partition); + if (fromBeginning !== false && endOffsetSnapshot.available && !firstScannedOffsetByPartition.has(partition)) { + const actualStartOffset = stringValue(message.offset); + const expectedStartOffset = startOffsetByPartition.get(partition); + firstScannedOffsetByPartition.set(partition, actualStartOffset); + if (!endOffsetByPartition.has(partition) || !expectedStartOffset || actualStartOffset !== expectedStartOffset) { + finish("retention-start-moved"); + return; + } + verifiedStartPartitions.add(partition); } - verifiedStartPartitions.add(partition); - } - if (fromBeginning !== false && endOffset && kafkaOffsetAtOrBeyond(message.offset, endOffset)) { - excludedPostBarrierCount += 1; - reachedEndPartitions.add(partition); + if (fromBeginning !== false && endOffset && kafkaOffsetAtOrBeyond(message.offset, endOffset)) { + excludedPostBarrierCount += 1; + reachedEndPartitions.add(partition); + if (endOffsetSnapshot.available && reachedEndPartitions.size === endOffsetByPartition.size) { + finish(verifiedStartPartitions.size === endOffsetByPartition.size ? "end-offset" : "retention-start-unverified"); + } + continue; + } + scannedCount += 1; + lastScannedOffsetByPartition.set(partition, stringValue(message.offset)); + const valueText = message.value ? Buffer.from(message.value).toString("utf8") : ""; + const value = parseJson(valueText); + if (!value || typeof value !== "object" || Array.isArray(value)) invalidJsonCount += 1; + else { + parsedCount += 1; + if (!eventMatchesFilters(value, { traceId, sessionId, runId, commandId })) filterRejectedCount += 1; + else { + events.push({ + topic: messageTopic, + partition, + offset: message.offset, + key: message.key ? Buffer.from(message.key).toString("utf8") : null, + timestamp: message.timestamp ?? null, + valueSha256: sha256(valueText), + value + }); + } + } + if (fromBeginning !== false && endOffset && kafkaOffsetReached(message.offset, endOffset)) reachedEndPartitions.add(partition); if (endOffsetSnapshot.available && reachedEndPartitions.size === endOffsetByPartition.size) { finish(verifiedStartPartitions.size === endOffsetByPartition.size ? "end-offset" : "retention-start-unverified"); + return; } - return; - } - scannedCount += 1; - lastScannedOffsetByPartition.set(partition, stringValue(message.offset)); - const valueText = message.value ? Buffer.from(message.value).toString("utf8") : ""; - const value = parseJson(valueText); - if (!value || typeof value !== "object" || Array.isArray(value)) invalidJsonCount += 1; - else { - parsedCount += 1; - if (!eventMatchesFilters(value, { traceId, sessionId, runId, commandId })) filterRejectedCount += 1; - else { - events.push({ - topic: messageTopic, - partition, - offset: message.offset, - key: message.key ? Buffer.from(message.key).toString("utf8") : null, - timestamp: message.timestamp ?? null, - valueSha256: sha256(valueText), - value - }); + if (events.length >= maxEvents) { + finish("limit"); + return; + } + if (maxScannedRecords !== null && scannedCount >= maxScannedRecords) { + finish("scan-limit"); + return; } } - if (fromBeginning !== false && endOffset && kafkaOffsetReached(message.offset, endOffset)) reachedEndPartitions.add(partition); - if (endOffsetSnapshot.available && reachedEndPartitions.size === endOffsetByPartition.size) { - return finish(verifiedStartPartitions.size === endOffsetByPartition.size ? "end-offset" : "retention-start-unverified"); - } - if (events.length >= maxEvents) finish("limit"); - else if (maxScannedRecords !== null && scannedCount >= maxScannedRecords) finish("scan-limit"); } }).catch(reject); }); + timing.scanMs = Date.now() - stageStartedAtMs; } catch (error) { if (error?.code === "kafka_query_aborted") completionReason = "aborted"; else throw error; } finally { + stageStartedAtMs = Date.now(); if (timer) clearTimeout(timer); removeAbortListener?.(); if (running) await boundedKafkaPromise(consumer.stop(), remainingKafkaQueryCleanupBudget(deadlineMs)).catch(() => undefined); await boundedDisconnect(consumer, remainingKafkaQueryCleanupBudget(deadlineMs)); + timing.cleanupMs = Date.now() - stageStartedAtMs; } return kafkaQueryResult(); @@ -1126,12 +1177,51 @@ export async function queryKafkaEventStream({ env = process.env, stream = "hwlab scanLimit: maxScannedRecords, timeoutMs: budgetMs, filters: compactObject({ traceId, sessionId, runId, commandId }), + partitionKeyScoped: targetPartition !== null, + targetPartition, + timing: { + ...timing, + totalMs: Date.now() - queryStartedAtMs, + valuesPrinted: false + }, events, valuesPrinted: false }; } } +export function kafkaPartitionForKey(partitionKey, topic, offsets = []) { + const key = stringValue(partitionKey); + const partitionMetadata = offsets + .map((entry) => integerValue(entry?.partition)) + .filter((partition) => partition !== null) + .map((partitionId) => ({ partitionId })) + .sort((left, right) => left.partitionId - right.partitionId); + if (!key || partitionMetadata.length === 0) return null; + return Partitioners.DefaultPartitioner()({ + topic, + partitionMetadata, + message: { key: Buffer.from(key), value: null, headers: {} } + }); +} + +async function prepareKafkaQueryGroupOffsets(kafka, { groupId, topic, targetPartition, offsets, deadlineMs, signal }) { + const admin = kafka.admin(); + try { + await withinKafkaQueryBudget(admin.connect(), deadlineMs, "query-offset-admin-connect", signal); + await withinKafkaQueryBudget(admin.setOffsets({ + groupId, + topic, + partitions: offsets.map((entry) => ({ + partition: entry.partition, + offset: entry.partition === targetPartition ? entry.startOffset : entry.endOffset + })) + }), deadlineMs, "query-offset-admin-set", signal); + } finally { + await boundedDisconnect(admin, remainingKafkaQueryCleanupBudget(deadlineMs)); + } +} + async function kafkaTopicEndOffsetSnapshot(kafka, topic, deadlineMs, signal = null) { if (typeof kafka?.admin !== "function") return { available: false, offsets: [], error: { code: "kafka_admin_unavailable", valuesRedacted: true } }; const admin = kafka.admin(); diff --git a/internal/cloud/server-workbench-realtime-http.test.ts b/internal/cloud/server-workbench-realtime-http.test.ts index 835cceef..91adf653 100644 --- a/internal/cloud/server-workbench-realtime-http.test.ts +++ b/internal/cloud/server-workbench-realtime-http.test.ts @@ -146,6 +146,61 @@ test("product realtime uses live Kafka SSE without projection reads", async () = } }); +test("refresh replay scopes the retained Kafka query to the session producer key", async () => { + const sessionId = "ses_refresh_partition_key"; + const queryCalls = []; + const subscribers = new Set(); + const capabilities = { directPublish: true, liveKafkaSse: true, kafkaRefreshReplay: true, transactionalProjector: false, projectionOutboxRelay: false, projectionRealtime: false }; + const server = createCloudApiServer({ + accessController: realtimeAccessController({ sessions: [{ id: sessionId, ownerUserId: ACTOR.id }] }), + workbenchRuntime: {}, + kafkaEventBridge: { + capabilities, + refreshReplay: { groupIdPrefix: "hwlab-test-refresh", timeoutMs: 5000, scanLimit: 100, matchedEventLimit: 20, liveBufferLimit: 20 }, + ready: Promise.resolve(), + subscribeLiveHwlabEvents(listener) { subscribers.add(listener); return () => subscribers.delete(listener); }, + async queryHwlabEventRetention(options) { + queryCalls.push(options); + return { + topic: "hwlab.event.v1", + events: [], + completionReason: "end-offset", + completion: { reason: "end-offset", complete: true, barrierReached: true, retentionStartVerified: true }, + reachedEndOffsets: true, + endOffsetsAvailable: true, + endOffsets: [{ partition: 2, startOffset: "0", endOffset: "10" }], + scannedCount: 10, + matchedCount: 0 + }; + }, + async stop() {} + }, + env: { + HWLAB_KAFKA_DIRECT_PUBLISH_ENABLED: "true", + HWLAB_WORKBENCH_LIVE_KAFKA_SSE_ENABLED: "true", + HWLAB_WORKBENCH_KAFKA_REFRESH_REPLAY_ENABLED: "true", + HWLAB_KAFKA_TRANSACTIONAL_PROJECTOR_ENABLED: "false", + HWLAB_KAFKA_PROJECTION_OUTBOX_RELAY_ENABLED: "false", + HWLAB_WORKBENCH_PROJECTION_REALTIME_ENABLED: "false", + HWLAB_WORKBENCH_KAFKA_REFRESH_REPLAY_GROUP_PREFIX: "hwlab-test-refresh", + HWLAB_WORKBENCH_KAFKA_REFRESH_REPLAY_TIMEOUT_MS: "5000", + HWLAB_WORKBENCH_KAFKA_REFRESH_REPLAY_SCAN_LIMIT: "100", + HWLAB_WORKBENCH_KAFKA_REFRESH_REPLAY_MATCHED_EVENT_LIMIT: "20", + HWLAB_WORKBENCH_KAFKA_REFRESH_REPLAY_LIVE_BUFFER_LIMIT: "20" + } + }); + 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}`, 1); + assert.equal(events[0].event, "workbench.connected"); + assert.equal(queryCalls.length, 1); + assert.equal(queryCalls[0].sessionId, sessionId); + assert.equal(queryCalls[0].partitionKey, sessionId); + } finally { + await new Promise((resolve, reject) => server.close((error) => error ? reject(error) : resolve())); + } +}); + test("synchronous projection notification is buffered until the unique initial snapshot completes", async () => { const sessionId = "ses_realtime_initial_barrier"; const traceId = "trc_realtime_initial_barrier"; diff --git a/internal/cloud/server-workbench-realtime-http.ts b/internal/cloud/server-workbench-realtime-http.ts index 9cb764bf..867c1229 100644 --- a/internal/cloud/server-workbench-realtime-http.ts +++ b/internal/cloud/server-workbench-realtime-http.ts @@ -765,6 +765,7 @@ async function handleKafkaRefreshReplayWorkbenchRealtimeHttp(request, response, scanLimit: refreshReplay.scanLimit, timeoutMs: refreshReplay.timeoutMs, groupIdPrefix: refreshReplay.groupIdPrefix, + partitionKey: requestedSessionId, fromBeginning: true, signal }), diff --git a/web/hwlab-cloud-web/scripts/check.ts b/web/hwlab-cloud-web/scripts/check.ts index c8c2eab8..741131e1 100644 --- a/web/hwlab-cloud-web/scripts/check.ts +++ b/web/hwlab-cloud-web/scripts/check.ts @@ -199,8 +199,7 @@ assertIncludes(workbenchRealtimePlanSource, "afterOutboxSeq: finiteNumber(recove 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"); -assertIncludes(workbenchStoreSource, "workbenchColadaQueries.fetchSession", "Realtime session detail recovery must enter Colada query facade"); -assertIncludes(workbenchStoreSource, "runtimePolicy.workbenchSessionDetailMinRefreshMs", "Realtime session detail recovery budget must come from runtime policy"); +assert.doesNotMatch(workbenchStoreSource, /workbenchColadaQueries\.fetchSession(?:Detail|Messages)?\(/u, "HTTP session detail and messages must not participate in Workbench business projection"); assert.doesNotMatch(workbenchStoreSource, /refreshRealtimeSessionMessages[\s\S]{0,900}refreshSessionMessageProjectionPage\(id, \{ force: true \}\)/u, "Realtime session message recovery must not force-bypass the message projection refresh budget"); assert.doesNotMatch(workbenchStoreSource, /handleRealtimeStreamError[\s\S]{0,1200}refreshSessions\([^;]+force:\s*true/u, "Realtime stream errors must not force-refresh the full session list"); assert.doesNotMatch(workbenchStoreSource, /refreshActiveTraceFromRest[\s\S]{0,1200}refreshSessions\([^;]+force:\s*true/u, "Active trace sync replay must not force-refresh the full session list"); 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 3590caec..e3878a49 100644 --- a/web/hwlab-cloud-web/scripts/workbench-realtime-runtime.test.ts +++ b/web/hwlab-cloud-web/scripts/workbench-realtime-runtime.test.ts @@ -22,7 +22,6 @@ import { projectRejectedWorkbenchAdmission } from "../src/stores/workbench-admis import { WORKBENCH_TIMELINE_OPENCODE_PARITY, buildWorkbenchTimelineRows, normalizeWorkbenchTimelineMessages, workbenchTimelineSignature } from "../src/stores/workbench-timeline-model.ts"; import { reduceWorkbenchRealtimeEvent } from "../src/stores/workbench-event-reducer.ts"; import { assessWorkbenchRealtimeConnected, planWorkbenchRealtimeApply, planWorkbenchRealtimeRecovery } from "../src/stores/workbench-realtime-plan.ts"; -import { WORKBENCH_REALTIME_AUTHORITY_VERSION, workbenchRealtimePrimaryAuthorityDecision } from "../src/stores/workbench-realtime-authority.ts"; import { cleanupWorkbenchServerStateDroppedSessions, cleanupWorkbenchServerStateSessions, createWorkbenchServerState, reduceWorkbenchServerState } from "../src/stores/workbench-server-state.ts"; import { cleanupDroppedWorkbenchSessionCaches, trimWorkbenchSessionCache } from "../src/stores/workbench-session-cache.ts"; @@ -505,16 +504,16 @@ test("server-state cleanup removes dropped session trace and turn authority", () }); test("realtime event reducer classifies SSE payloads before store side effects", () => { - const trace = reduceWorkbenchRealtimeEvent(realtimeEvent({ type: "trace.event", traceId: "trc_1", event: { traceId: "trc_1", label: "delta" }, snapshot: { traceId: "trc_1", status: "running" }, entity: { family: "traceEvents", id: "trc_1:1", version: 1, projectionRevision: "prj_1" } }), "workbench.trace.event"); + const trace = reduceWorkbenchRealtimeEvent(realtimeEvent({ schema: "hwlab.event.v1", type: "trace.event", traceId: "trc_1", sessionId: "ses_1", event: { traceId: "trc_1", sessionId: "ses_1", label: "delta" } }), "hwlab.event.v1"); assert.equal(trace.activityLabel, "realtime:trace.event"); assert.equal(trace.action.type, "trace.event"); assert.equal(trace.diagnostic.module, "workbench-event-reducer"); - const missingAuthority = reduceWorkbenchRealtimeEvent({ type: "message.snapshot", message: agentMessage({ status: "running", sessionId: "ses_1" }) }, "workbench.message.snapshot"); - assert.deepEqual(missingAuthority.action, { type: "ignore", reason: "workbench_realtime_authority_missing" }); + const legacyMessage = reduceWorkbenchRealtimeEvent({ type: "message.snapshot", message: agentMessage({ status: "running", sessionId: "ses_1" }) }, "workbench.message.snapshot"); + assert.deepEqual(legacyMessage.action, { type: "ignore", reason: "non-kafka-business-projection-removed" }); - const detailOnly = reduceWorkbenchRealtimeEvent(realtimeEvent({ type: "trace.event", traceId: "trc_1", detailProjection: true, authority: "trace-detail-only", event: { traceId: "trc_1", label: "detail" }, entity: { family: "traceEvents", id: "trc_1:2", version: 2, projectionRevision: "prj_1", authority: "trace-detail-only" } }), "workbench.trace.event"); - assert.deepEqual(detailOnly.action, { type: "ignore", reason: "workbench_realtime_detail_only_rejected" }); + const legacyTrace = reduceWorkbenchRealtimeEvent(realtimeEvent({ type: "trace.event", traceId: "trc_1", event: { traceId: "trc_1", label: "detail" } }), "workbench.trace.event"); + assert.deepEqual(legacyTrace.action, { type: "ignore", reason: "non-kafka-business-projection-removed" }); const error = reduceWorkbenchRealtimeEvent({ type: "error", traceId: "trc_2", error: { message: "offline" } }, "workbench.error"); assert.equal(error.action.type, "projection.error"); @@ -522,7 +521,7 @@ test("realtime event reducer classifies SSE payloads before store side effects", }); test("realtime apply planner turns reducer actions into store steps", () => { - const reduced = reduceWorkbenchRealtimeEvent(realtimeEvent({ type: "trace.event", traceId: "trc_1", event: { traceId: "trc_1", label: "delta" }, snapshot: { traceId: "trc_1", status: "running" }, entity: { family: "traceEvents", id: "trc_1:1", version: 1, projectionRevision: "prj_1" } }), "workbench.trace.event"); + const reduced = reduceWorkbenchRealtimeEvent(realtimeEvent({ schema: "hwlab.event.v1", type: "trace.event", traceId: "trc_1", sessionId: "ses_1", event: { traceId: "trc_1", sessionId: "ses_1", label: "delta" } }), "hwlab.event.v1"); const tracePlan = planWorkbenchRealtimeApply(reduced.action); assert.deepEqual(tracePlan.steps.map((step) => step.type), ["apply-trace-event"]); assert.equal(tracePlan.diagnostic.module, "workbench-realtime-plan"); @@ -552,25 +551,6 @@ test("realtime recovery planner reconnects the projection SSE from its outbox cu assert.deepEqual(terminalSealed.steps.map((step) => step.type), ["events-reconnect"]); }); -test("realtime authority accepts projection outbox replay and live events through one entity contract", () => { - const event = realtimeEvent({ type: "message.snapshot", sessionId: "ses_1", message: agentMessage({ id: "msg_1", sessionId: "ses_1", status: "completed", text: "done" }), entity: { family: "messages", id: "msg_1", version: 7, outboxSeq: 12, projectionRevision: "prj_7" } }); - const decision = workbenchRealtimePrimaryAuthorityDecision(event); - assert.equal(decision.accepted, true); - assert.equal(decision.entity?.family, "messages"); - assert.equal(decision.entity?.version, 7); -}); - -test("realtime authority rejects trace detail-only and incomplete contract payloads", () => { - const detail = realtimeEvent({ type: "trace.event", traceId: "trc_1", detailProjection: true, authority: "trace-detail-only", event: { traceId: "trc_1" }, entity: { family: "traceEvents", id: "trc_1:1", version: 1, projectionRevision: "prj_1", authority: "trace-detail-only" } }); - assert.equal(workbenchRealtimePrimaryAuthorityDecision(detail).reason, "workbench_realtime_detail_only_rejected"); - - const missingEntity = realtimeEvent({ type: "message.snapshot", message: agentMessage({ status: "running", sessionId: "ses_1" }) }); - assert.equal(workbenchRealtimePrimaryAuthorityDecision(missingEntity).reason, "workbench_realtime_entity_missing"); - - const missingProjection = realtimeEvent({ type: "message.snapshot", message: agentMessage({ status: "running", sessionId: "ses_1" }), entity: { family: "messages", id: "msg_1", version: 1 } }); - assert.equal(workbenchRealtimePrimaryAuthorityDecision(missingProjection).reason, "workbench_realtime_projection_revision_missing"); -}); - test("health probe cache records ok and unavailable states", async () => { const cache = createWorkbenchHealthProbeCache({ cacheMs: 100 }); const ok = await cache.probe({ key: "workbench", fetcher: async () => ({ ready: true }), classify: (value) => value.ready ? "ok" : "degraded" }); @@ -661,5 +641,5 @@ function recoveryEvent(actions: WorkbenchStreamTransportRecovery["actions"], cur } function realtimeEvent(input: Record) { - return { realtimeAuthority: WORKBENCH_REALTIME_AUTHORITY_VERSION, contractVersion: "workbench-sync-v1", ...input }; + return input; } diff --git a/web/hwlab-cloud-web/src/stores/workbench-event-reducer.ts b/web/hwlab-cloud-web/src/stores/workbench-event-reducer.ts index e9984b5c..46d66de9 100644 --- a/web/hwlab-cloud-web/src/stores/workbench-event-reducer.ts +++ b/web/hwlab-cloud-web/src/stores/workbench-event-reducer.ts @@ -5,13 +5,8 @@ import type { WorkbenchRealtimeEvent } from "../api/workbench-events"; import { firstNonEmptyString } from "../utils"; -import { workbenchRealtimePrimaryAuthorityDecision } from "./workbench-realtime-authority"; - export type WorkbenchRealtimeAction = - | { type: "trace.snapshot"; traceId: string | null; snapshot: WorkbenchRealtimeEvent["snapshot"] } | { type: "trace.event"; traceId: string | null; event: WorkbenchRealtimeEvent["event"]; snapshot: WorkbenchRealtimeEvent["snapshot"]; realtimeEvent: WorkbenchRealtimeEvent } - | { type: "message.snapshot"; realtimeEvent: WorkbenchRealtimeEvent } - | { type: "turn.snapshot"; turn: NonNullable } | { type: "projection.error"; realtimeEvent: WorkbenchRealtimeEvent } | { type: "trace.unavailable"; traceId: string | null; reason: string } | { type: "ignore"; reason: string }; @@ -57,17 +52,12 @@ function reduceRealtimeAction(event: WorkbenchRealtimeEvent, eventName: string): if (event.schema === "hwlab.event.v1" && event.event) { return { type: "trace.event", traceId: realtimeTraceId(event), event: event.event, snapshot: null, realtimeEvent: event }; } - const authority = primaryAuthority(event); - if (authority) return authority; switch (event.type) { - case "trace.snapshot": - return { type: "trace.snapshot", traceId: realtimeTraceId(event), snapshot: event.snapshot ?? null }; case "trace.event": - return { type: "trace.event", traceId: realtimeTraceId(event), event: event.event ?? null, snapshot: event.snapshot ?? null, realtimeEvent: event }; case "message.snapshot": - return event.message ? { type: "message.snapshot", realtimeEvent: event } : { type: "ignore", reason: "message.snapshot.missing-message" }; case "turn.snapshot": - return event.turn ? { type: "turn.snapshot", turn: event.turn } : { type: "ignore", reason: "turn.snapshot.missing-turn" }; + case "trace.snapshot": + return { type: "ignore", reason: "non-kafka-business-projection-removed" }; case "trace.unavailable": return { type: "trace.unavailable", traceId: realtimeTraceId(event), reason: firstNonEmptyString(event.reason, "realtime-trace-unavailable") ?? "realtime-trace-unavailable" }; case "error": @@ -77,13 +67,6 @@ function reduceRealtimeAction(event: WorkbenchRealtimeEvent, eventName: string): return { type: "ignore", reason: firstNonEmptyString(event.type, eventName, "unsupported") ?? "unsupported" }; } -function primaryAuthority(event: WorkbenchRealtimeEvent): WorkbenchRealtimeAction | null { - if (!["trace.snapshot", "trace.event", "message.snapshot", "turn.snapshot"].includes(firstNonEmptyString(event.type) ?? "")) return null; - const decision = workbenchRealtimePrimaryAuthorityDecision(event); - if (decision.accepted) return null; - return { type: "ignore", reason: decision.reason }; -} - function realtimeTraceId(event: WorkbenchRealtimeEvent): string | null { return firstNonEmptyString(event.traceId, event.snapshot?.traceId, event.event?.traceId, event.message?.traceId) ?? null; } diff --git a/web/hwlab-cloud-web/src/stores/workbench-live-kafka-event.test.ts b/web/hwlab-cloud-web/src/stores/workbench-live-kafka-event.test.ts index 5f8b8c6a..32b6dad9 100644 --- a/web/hwlab-cloud-web/src/stores/workbench-live-kafka-event.test.ts +++ b/web/hwlab-cloud-web/src/stores/workbench-live-kafka-event.test.ts @@ -190,6 +190,46 @@ test("session replay keeps steer user input separate from the target agent lifec assert.equal([...agentMessages.values()].filter((message) => message.finalResponse).length, 1); }); +test("refresh rebuilds the completed turn from an empty store using only Kafka SSE events", () => { + const sessionId = "ses_refresh_kafka_only"; + const traceId = "trc_refresh_kafka_only"; + const threadId = "thread_refresh_kafka_only"; + const events = [ + traceEvent(1, "user", { traceId, sessionId, userMessageId: "msg_refresh_user", messageId: "msg_refresh_user", text: "hi", threadId }), + traceEvent(2, "backend", { traceId, sessionId, threadId }), + traceEvent(3, "assistant", { traceId, sessionId, assistantText: "refresh survives", text: "refresh survives", threadId }), + traceEvent(4, "result", { traceId, sessionId, terminal: true, status: "completed", threadId }) + ]; + + const replay = (): ChatMessage[] => { + let state = createWorkbenchServerState(); + state = reduceWorkbenchServerState(state, { type: "session.detail", session: { sessionId, messages: [] } }); + for (const [index, event] of events.entries()) { + const receivedAt = `2026-07-10T10:02:0${index + 1}.000Z`; + if (workbenchLiveKafkaProjectionTarget(event) === "user") { + const previous = selectActiveMessages(state, sessionId).find((message) => message.role === "user" && message.traceId === traceId) ?? null; + const message = projectWorkbenchLiveKafkaUserMessage({ previous, traceId, sessionId, event, receivedAt, threadId }); + if (message) state = reduceWorkbenchServerState(state, { type: "message.upsert", sessionId, message }); + continue; + } + const previous = selectActiveMessages(state, sessionId).find((message) => message.role === "agent" && message.traceId === traceId) ?? null; + const message = projectWorkbenchLiveKafkaMessage({ previous, traceId, sessionId, event, receivedAt, threadId }); + state = reduceWorkbenchServerState(state, { type: "message.upsert", sessionId, message }); + } + return selectActiveMessages(state, sessionId); + }; + + const beforeRefresh = replay(); + const afterRefresh = replay(); + assert.deepEqual(afterRefresh, beforeRefresh); + assert.equal(afterRefresh.length, 2); + assert.equal(afterRefresh[0]?.text, "hi"); + assert.equal(afterRefresh[1]?.status, "completed"); + assert.equal(afterRefresh[1]?.text, "refresh survives"); + assert.equal((afterRefresh[1]?.finalResponse as { text?: string } | null)?.text, "refresh survives"); + assert.equal(afterRefresh[1]?.threadId, threadId); +}); + test("live Kafka projection refreshes lastEventAt from each ingress receipt", () => { const traceId = "trc_live_timing"; const sessionId = "ses_live_timing"; diff --git a/web/hwlab-cloud-web/src/stores/workbench-live-kafka-event.ts b/web/hwlab-cloud-web/src/stores/workbench-live-kafka-event.ts index abee3c58..e74786c3 100644 --- a/web/hwlab-cloud-web/src/stores/workbench-live-kafka-event.ts +++ b/web/hwlab-cloud-web/src/stores/workbench-live-kafka-event.ts @@ -19,6 +19,7 @@ export interface WorkbenchLiveKafkaProjectionInput { event: TraceEvent; receivedAt: string; title?: string; + threadId?: string | null; } export interface WorkbenchLiveKafkaUserProjectionInput { @@ -27,6 +28,7 @@ export interface WorkbenchLiveKafkaUserProjectionInput { sessionId: string; event: TraceEvent; receivedAt: string; + threadId?: string | null; } export function reduceWorkbenchLiveKafkaMessageState(previous: WorkbenchLiveMessageState, event: TraceEvent, previousEvents: TraceEvent[] = []): WorkbenchLiveMessageState { @@ -92,6 +94,7 @@ export function projectWorkbenchLiveKafkaUserMessage(input: WorkbenchLiveKafkaUs targetTraceId: firstNonEmptyString(input.event.targetTraceId), turnId: input.previous?.turnId ?? input.traceId, sessionId: input.sessionId, + threadId: firstNonEmptyString(input.threadId, input.previous?.threadId), createdAt: input.previous?.createdAt ?? createdAt, updatedAt: receivedAt } as ChatMessage; @@ -185,6 +188,7 @@ export function projectWorkbenchLiveKafkaMessage(input: WorkbenchLiveKafkaProjec traceId: input.traceId, turnId: previous?.turnId ?? input.traceId, sessionId: input.sessionId, + threadId: firstNonEmptyString(input.threadId, previous?.threadId), createdAt: previous?.createdAt ?? startedAt, updatedAt: receivedAt, startedAt, 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 2dfb8f36..80ed1aa0 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 @@ -132,47 +132,27 @@ test("terminal response helpers treat failed body as sealed authority", () => { assert.equal(traceHasTerminalResponse(traceId, [unsealed]), false); }); -test("workbench active terminal paths seal final response from turn authority", () => { +test("workbench active state uses Kafka SSE without HTTP session or terminal projection", () => { const source = fs.readFileSync(path.join(storeDir, "workbench.ts"), "utf8"); - const projectBlock = source.slice(source.indexOf("function projectTurnAuthorityToMessages"), source.indexOf("async function submitMessage")); - const realtimeTurnBlock = source.slice(source.indexOf("function applyRealtimeTurnSnapshot"), source.indexOf("function scheduleRealtimeTurnProjection")); - const realtimeTurnProjectionBlock = source.slice(source.indexOf("function flushRealtimeTurnProjection"), source.indexOf("function installRealtimeVisibilityHandler")); - const completeBlock = source.slice(source.indexOf("function completeTrace"), source.indexOf("function applyTerminalResultDiagnostics")); - const terminalDiagnosticsBlock = source.slice(source.indexOf("function applyTerminalResultDiagnostics"), source.indexOf("function applyRealtimeProjectionError")); - const crossTabSyncBlock = source.slice(source.indexOf("async function refreshCrossTabSessionFromSyncReplay"), source.indexOf("function completeTrace")); - const sessionDetailReadBlock = source.slice(source.indexOf("function fetchSessionDetailPage"), source.indexOf("function sessionMessageProjectionWindowLimit")); + const submitBlock = source.slice(source.indexOf("async function submitMessage"), source.indexOf("async function cancelAgentMessage")); + const traceDetailApplyBlock = source.slice(source.indexOf("function applyTraceDetailEventsResult"), source.indexOf("function setProviderProfile")); const traceDetailReadBlock = source.slice(source.indexOf("async function readTraceEventsForExplicitDetailPages"), source.indexOf("async function fetchTraceDetailEventsPage")); - const loadBlock = source.slice(source.indexOf("async function loadWorkbenchSession"), source.indexOf("function reattachRestoredActiveTrace")); + const loadBlock = source.slice(source.indexOf("function loadWorkbenchSession"), source.indexOf("function reattachRestoredActiveTrace")); - assert.match(projectBlock, /terminalAuthorityMessageFromTurnResult\(result\)/u); - assert.match(projectBlock, /type:\s*"message\.upsert"/u); - assert.match(projectBlock, /\.\.\.terminalPatch/u); assert.match(source, /const traceDetailReadSingleflight = createKeyedSingleflight\(\)/u); assert.match(source, /traceDetailReadSingleflight\.run/u); assert.doesNotMatch(source, new RegExp(["refresh", "RealtimeSessionMessages"].join(""), "u")); - assert.match(realtimeTurnBlock, /rememberTurnStatus\(traceId, result\)[\s\S]*scheduleRealtimeTurnProjection\(\{ traceId, result, terminalTurn \}\)/u); - assert.doesNotMatch(realtimeTurnBlock, new RegExp(`applyTurnStatusSnapshot\\(|${["refresh", "TerminalTraceFromRest"].join("")}\\(`, "u")); - assert.match(realtimeTurnProjectionBlock, /syncTurnStatusToMessage\(next\.traceId, next\.result\)/u); - assert.match(realtimeTurnProjectionBlock, /workbench_realtime_turn_projection_budget/u); - assert.doesNotMatch(realtimeTurnProjectionBlock, /refreshWorkbenchSyncReplay|refreshTerminalTraceFromSyncReplay/u); - assert.match(completeBlock, /projectTurnAuthorityToMessages\(traceId, result, "complete-trace"\)/u); - assert.doesNotMatch(completeBlock, /forceRead|refreshMessageProjectionForTrace|fetchWorkbenchTurnStatus|refreshTurnStatusByTraceId/u); - assert.doesNotMatch(terminalDiagnosticsBlock, /readTraceEventsForMessages|fetchSessionMessagesPage|fetchWorkbenchTurnStatus/u); - assert.doesNotMatch(completeBlock, /scheduleActiveTraceSyncReplay|refreshWorkbenchSyncReplay/u); + assert.doesNotMatch(source, /function applyRealtimeTurnSnapshot|function applyRealtimeMessageSnapshot|function applyRealtimeTraceSnapshot/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); - assert.match(sessionDetailReadBlock, /workbenchSessionDetailReadKey\(\{ sessionId, force: options\.force \}\)/u); - assert.match(sessionDetailReadBlock, /fetchSession\(sessionId, \{ includeMessages: false,/u); - assert.match(loadBlock, /const messageLimit = sessionMessageProjectionWindowLimit\(\);/u); - assert.match(loadBlock, /fetchSessionDetailPage\(requestId, \{ reason: "load-session:detail", force: true \}\)/u); - assert.match(loadBlock, /sessionFromWorkbenchSession\(detail\.data\?\.session, \{ includeMessages: false \}\)/u); - assert.match(loadBlock, /const fallbackMessages = seed\?\.sessionId === id \? seed\.messages \?\? \[\] : \[\];/u); - assert.match(loadBlock, /fetchSessionMessagesPage\(normalizedRequestId, \{ limit: messageLimit, reason: "load-session:messages", force: true \}\)/u); - assert.doesNotMatch(loadBlock, /messages:\s*\[\]/u); - assert.doesNotMatch(loadBlock, /limit:\s*100/u); + assert.doesNotMatch(source, /fetchSessionDetailPage|fetchSessionMessagesPage/u); + assert.doesNotMatch(loadBlock, /fetchSession\(|fetchSessionMessages\(/u); + assert.match(loadBlock, /workbenchSessionNavigationSeed\(id/u); + assert.doesNotMatch(submitBlock, /applyTurnStatusSnapshot|completeTrace/u); + assert.doesNotMatch(traceDetailApplyBlock, /rememberTurnStatus|projectTurnAuthorityToMessages|chatPending\.value\s*=\s*false/u); assert.doesNotMatch(source, /sealRestoredActiveTurnMessages|messageNeedsRestoredTurnSeal|refreshSessionStatusAuthority|readTerminalTraceDetailGaps/u); }); @@ -180,9 +160,7 @@ 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")), - source.slice(source.indexOf("function reattachTrace"), source.indexOf("function restartRealtime")), - source.slice(source.indexOf("function flushRealtimeTurnProjection"), source.indexOf("function installRealtimeVisibilityHandler")), - source.slice(source.indexOf("function completeTrace"), source.indexOf("function applyTerminalResultDiagnostics")) + source.slice(source.indexOf("function reattachTrace"), source.indexOf("function restartRealtime")) ].join("\n"); assert.doesNotMatch(activeBlocks, /refreshWorkbenchSyncReplay|scheduleActiveTraceSyncReplay|setTimeout|setInterval/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 deleted file mode 100644 index 10cfac9f..00000000 --- a/web/hwlab-cloud-web/src/stores/workbench-realtime-authority.ts +++ /dev/null @@ -1,107 +0,0 @@ -// SPEC: PJ2026-010401080313 Workbench实时权威 draft-2026-07-08-p0-workbench-realtime-authority-v2. -// Responsibility: Frontend authority gates for Workbench projection outbox realtime and replay SSE events. - -import type { WorkbenchRealtimeEvent } from "@/api/workbench-events"; - -export const WORKBENCH_REALTIME_AUTHORITY_VERSION = "workbench-realtime-authority-v2"; - -export interface WorkbenchRealtimeEntityAuthority { - family: string; - id: string; - version: number; - outboxSeq: number | null; - traceSeq: number | null; - projectionRevision: string | null; - committedAt: string | null; - authority: string | null; - detailProjection: boolean; -} - -export interface WorkbenchRealtimeAuthorityDecision { - accepted: boolean; - reason: string; - entity: WorkbenchRealtimeEntityAuthority | null; - diagnostic: { - module: "workbench-realtime-authority"; - code: string; - reason: string; - realtimeAuthority: string | null; - contractVersion: string | null; - entityFamily: string | null; - entityId: string | null; - entityVersion: number | null; - detailProjection: boolean; - valuesRedacted: true; - }; -} - -export function workbenchRealtimePrimaryAuthorityDecision(event: WorkbenchRealtimeEvent): WorkbenchRealtimeAuthorityDecision { - const entity = workbenchRealtimeEntityAuthority(event); - const detailProjection = workbenchRealtimeDetailOnly(event, entity); - const code = authorityFailureCode(event, entity, detailProjection); - const accepted = code === null; - const reason = accepted ? "accepted" : code; - return { - accepted, - reason, - entity, - diagnostic: { - module: "workbench-realtime-authority", - code: accepted ? "workbench_realtime_authority_accept" : code, - reason, - realtimeAuthority: stringValue(event.realtimeAuthority), - contractVersion: stringValue(event.contractVersion), - entityFamily: entity?.family ?? null, - entityId: entity?.id ?? null, - entityVersion: entity?.version ?? null, - detailProjection, - valuesRedacted: true - } - }; -} - -export function workbenchRealtimeEntityAuthority(event: WorkbenchRealtimeEvent): WorkbenchRealtimeEntityAuthority | null { - const source = recordValue(event.entity); - if (!source) return null; - const family = stringValue(source.family); - const id = stringValue(source.id ?? source.entityId); - const version = finiteNumber(source.version ?? source.entityVersion); - if (!family || !id || version === null) return null; - return { - family, - id, - version, - outboxSeq: finiteNumber(source.outboxSeq ?? event.cursor?.outboxSeq ?? event.outboxSeq), - traceSeq: finiteNumber(source.traceSeq ?? event.cursor?.traceSeq ?? event.traceSeq), - projectionRevision: stringValue(source.projectionRevision ?? event.projectionRevision), - committedAt: stringValue(source.committedAt ?? source.serverCommittedAt ?? event.eventCreatedAt ?? event.serverSentAt), - authority: stringValue(source.authority ?? event.authority), - detailProjection: source.detailProjection === true || event.detailProjection === 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"; - if (!entity) return "workbench_realtime_entity_missing"; - if (!entity.projectionRevision) return "workbench_realtime_projection_revision_missing"; - return null; -} - -function workbenchRealtimeDetailOnly(event: WorkbenchRealtimeEvent, entity: WorkbenchRealtimeEntityAuthority | null): boolean { - return event.detailProjection === true || entity?.detailProjection === true || stringValue(event.authority) === "trace-detail-only" || entity?.authority === "trace-detail-only"; -} - -function recordValue(value: unknown): Record | null { - return value && typeof value === "object" ? value as Record : null; -} - -function stringValue(value: unknown): string | null { - const text = typeof value === "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 b0ae2163..a272d87a 100644 --- a/web/hwlab-cloud-web/src/stores/workbench-realtime-plan.ts +++ b/web/hwlab-cloud-web/src/stores/workbench-realtime-plan.ts @@ -9,10 +9,7 @@ import type { WorkbenchStreamTransportRecovery } from "../utils/workbench-realti import type { WorkbenchRealtimeAction } from "./workbench-event-reducer"; export type WorkbenchRealtimeApplyStep = - | { type: "apply-trace-snapshot"; traceId: string | null; snapshot: WorkbenchRealtimeEvent["snapshot"] } | { type: "apply-trace-event"; traceId: string | null; event: WorkbenchRealtimeEvent["event"]; snapshot: WorkbenchRealtimeEvent["snapshot"]; realtimeEvent: WorkbenchRealtimeEvent } - | { type: "apply-message-snapshot"; realtimeEvent: WorkbenchRealtimeEvent } - | { type: "apply-turn-snapshot"; turn: NonNullable } | { type: "apply-projection-error"; realtimeEvent: WorkbenchRealtimeEvent } | { type: "clear-active-trace"; traceId: string; reason: string }; @@ -142,14 +139,8 @@ function finiteNumber(value: unknown): number | null { function applySteps(action: WorkbenchRealtimeAction): WorkbenchRealtimeApplyStep[] { switch (action.type) { - case "trace.snapshot": - return [{ type: "apply-trace-snapshot", traceId: action.traceId, snapshot: action.snapshot }]; case "trace.event": return [{ type: "apply-trace-event", traceId: action.traceId, event: action.event, snapshot: action.snapshot, realtimeEvent: action.realtimeEvent }]; - case "message.snapshot": - return [{ type: "apply-message-snapshot", realtimeEvent: action.realtimeEvent }]; - case "turn.snapshot": - return [{ type: "apply-turn-snapshot", turn: action.turn }]; case "projection.error": return [{ type: "apply-projection-error", realtimeEvent: action.realtimeEvent }]; case "trace.unavailable": diff --git a/web/hwlab-cloud-web/src/stores/workbench.ts b/web/hwlab-cloud-web/src/stores/workbench.ts index f1c2df6e..1ce81a45 100644 --- a/web/hwlab-cloud-web/src/stores/workbench.ts +++ b/web/hwlab-cloud-web/src/stores/workbench.ts @@ -1,4 +1,4 @@ -// SPEC: PJ2026-0106050514 Workbench实时运行面 draft-2026-06-30-p0-1297-spec-first; PJ2026-010403 API契约 draft-2026-06-18-r1; PJ2026-010401 Web工作台 draft-2026-06-18-r1; PJ2026-0104010803 唯一投影 draft-2026-06-18-p0-unique-projection; draft-2026-06-28-p0-d518-session-timeline-consistency; PJ2026-01060505 Workbench Performance draft-2026-06-17-p0; PJ2026-010401080313 Workbench实时权威 draft-2026-07-08-p0-workbench-realtime-authority-v2. +// SPEC: PJ2026-0106050514 Workbench实时运行面 draft-2026-06-30-p0-1297-spec-first; PJ2026-010403 API契约 draft-2026-06-18-r1; PJ2026-010401 Web工作台 draft-2026-06-18-r1; PJ2026-0104010803 唯一投影 draft-2026-06-18-p0-unique-projection; draft-2026-06-28-p0-d518-session-timeline-consistency; PJ2026-01060505 Workbench Performance draft-2026-06-17-p0; PJ2026-010401080313 Workbench实时权威 draft-2026-07-20-p0-kafka-sse-single-projection. // Responsibility: Session-first Workbench state orchestration for selection, turn admission, and trace lifecycle rendering. import { computed, nextTick, ref } from "vue"; @@ -11,9 +11,8 @@ import { createWorkbenchHealthProbeCache } from "@/utils/workbench-health"; import { agentErrorFromProjection, normalizeApiErrorRecord, normalizeErrorDiagnostic, normalizeProjectionDiagnostic, projectionDiagnosticFromApiFailure, projectionDiagnosticFromFailure } from "@/utils/workbench-error-runtime"; import { readWorkbenchJson, readWorkbenchNumber, readWorkbenchString, removeWorkbenchStorageKey, writeWorkbenchJson, writeWorkbenchString } from "@/utils/workbench-storage-runtime"; import { createWorkbenchStreamTransportRuntime, workbenchRealtimeTraceIdForCapabilities, type WorkbenchRealtimeEvent, type WorkbenchStreamTransportRecovery } from "@/utils/workbench-realtime-runtime"; -import { mergeRunnerTrace, snapshotToRunnerTrace, type TraceSnapshot } from "@/composables/workbench-trace-snapshot"; -import type { WorkbenchMessagePageResponse, WorkbenchSessionDetailResponse } from "@/api/workbench"; -import type { AgentChatResponse, AgentChatResultResponse, AgentRunProvenance, ApiError, ApiResult, ChatMessage, ErrorDiagnostic, LiveSurface, ProjectionBlocker, ProjectionDiagnostic, ProviderProfile, TraceEvent, WorkbenchSessionRecord, WorkbenchTurnTimingProjection } from "@/types"; +import { mergeRunnerTrace, type TraceSnapshot } from "@/composables/workbench-trace-snapshot"; +import type { AgentChatResponse, AgentChatResultResponse, ApiError, ApiResult, ChatMessage, ErrorDiagnostic, LiveSurface, ProjectionBlocker, ProjectionDiagnostic, ProviderProfile, TraceEvent, WorkbenchSessionRecord, WorkbenchTurnTimingProjection } from "@/types"; import { firstNonEmptyString, nextProtocolId, normalizeWorkbenchSessionId, normalizeWorkbenchSessionRouteId } from "@/utils"; import { composeWorkbenchScopedKey, workbenchRealtimeScopeKey } from "@/utils/workbench-key"; import { failWorkbenchSessionSwitch, failWorkbenchSubmitJourney, finishWorkbenchSessionSwitchFullLoad, markWorkbenchSubmitApiAccepted, markWorkbenchTraceEventsReceived, markWorkbenchTraceProjected, recordWorkbenchLoadingState, recordWorkbenchRuntimeDiagnostic, startWorkbenchSessionSwitch, startWorkbenchSubmitJourney } from "@/utils/workbench-performance"; @@ -25,29 +24,20 @@ import { reduceWorkbenchRealtimeEvent, workbenchRealtimeEventIsBusinessActivity, import { projectWorkbenchLiveKafkaMessage, projectWorkbenchLiveKafkaUserMessage, workbenchLiveKafkaEnvelope, workbenchLiveKafkaProjectionTarget } from "./workbench-live-kafka-event"; import { projectRejectedWorkbenchAdmission } from "./workbench-admission-failure"; import { messageHasSealedTerminalResult, messageIsSealedTerminal, traceAuthorityIsSealed } from "./workbench-terminal-authority"; -import { boundedProjectionMessageLimit } from "./workbench-message-projection-budget"; import { agentErrorDisplayText, agentErrorFromApiFailure, agentRunFromMessage, - agentRunFromResult, asAgentRun, - clearRunnerTraceTransientDiagnostics, finalResponseText, firstFiniteNumber, isTerminalMessageStatus, isTraceActiveStatus, messageHasTerminalResponse, messageNeedsTraceDetailRead, - messageStatusPatchForTerminalMerge, - terminalAuthorityMessageFromTurnResult, - terminalMessagePatchFromTurnResult, messageText, messageTimingPatch, - messageTimingPatchForMerge, messageTimingPatchFromProjection, - mergeTerminalResultTrace, - nonBlockingProjection, normalizeAgentError, normalizeMessageRunnerTrace, normalizedStatusText, @@ -57,20 +47,18 @@ import { projectionFromMessage, projectionFromResult, recordValue, - shouldClearCompletedTurnDiagnostics, shouldSuppressTransientWorkbenchReadFailure, terminalMessageTimingPatchForNormalize, turnResultIsTerminalForMerge, turnResultStatusForMerge, traceHasEvents, - traceResultHasTerminalEvidence, } from "./workbench-message-projection-runtime"; import { assessWorkbenchRealtimeConnected, planWorkbenchRealtimeApply, planWorkbenchRealtimeRecovery, type WorkbenchRealtimeApplyStep, type WorkbenchRealtimeRecoveryStep } from "./workbench-realtime-plan"; import { useWorkbenchColadaMutations } from "./workbench-colada-mutations"; import { useWorkbenchColadaQueries } from "./workbench-colada-queries"; import { useWorkbenchColadaReducer } from "./workbench-colada-reducer"; -import { traceEventsAutoReadDecision, workbenchSessionDetailReadKey, workbenchSessionMessagesReadKey, workbenchTraceEventsReadKey, type TraceEventsReadRangeRecord } from "./workbench-session-messages-read-budget"; -import { realtimeSnapshotToTraceSnapshot, terminalSealResultWithoutTraceEvents, traceDetailProjectedSeq, traceNextProjectedSeq } from "./workbench-trace-detail"; +import { traceEventsAutoReadDecision, workbenchTraceEventsReadKey, type TraceEventsReadRangeRecord } from "./workbench-session-messages-read-budget"; +import { traceDetailProjectedSeq, traceNextProjectedSeq } from "./workbench-trace-detail"; import { appendRawHwlabIngressFrame, createRawHwlabIngressState } from "./workbench-raw-hwlab-ingress"; const WORKBENCH_SESSION_PROJECTION_SIGNAL_CHANNEL = "hwlab.workbench.sessionProjection.v1"; @@ -87,23 +75,6 @@ interface SelectSessionOptions { source?: SessionSelectionSource; } -interface SessionMessagesReadOptions { - limit: number; - reason: string; - force?: boolean; -} - -interface SessionDetailReadOptions { - reason: string; - force?: boolean; -} - -interface RealtimeTurnProjectionItem { - traceId: string; - result: AgentChatResultResponse; - terminalTurn: boolean; -} - function workbenchMessageIdForTrace(traceId: string, role: "user" | "agent"): string { const suffix = firstNonEmptyString(traceId) ?.replace(/^trc_/u, "") @@ -121,11 +92,7 @@ export const useWorkbenchStore = defineStore("workbench", () => { const workbenchColadaQueries = useWorkbenchColadaQueries(); const workbenchColadaMutations = useWorkbenchColadaMutations(); const traceDetailReadSingleflight = createKeyedSingleflight(); - const sessionMessagesReadSingleflight = createKeyedSingleflight>(); - const sessionDetailReadSingleflight = createKeyedSingleflight>(); const traceEventsReadRanges = new Map(); - const realtimeTurnProjectionQueue = new Map(); - let realtimeTurnProjectionScheduled = false; const providerProfile = ref(readString("hwlab.workbench.providerProfile.v1", "codex")); const providerOptions = ref(defaultProviderProfileOptions(providerProfile.value)); const recentDrafts = ref(readRecentDrafts()); @@ -194,6 +161,8 @@ export const useWorkbenchStore = defineStore("workbench", () => { sessionDetailLoadingId.value = routeSessionId; recordWorkbenchLoadingState({ scope: "session_detail", active: true, reason: "hydrate", sessionId: routeSessionId }); setActiveSessionSelection(routeSessionId, "route"); + rememberSessionDetail(workbenchSessionNavigationSeed(routeSessionId)); + restartRealtime("hydrate-route"); } const sessionsResult = await workbenchColadaQueries.fetchSessions({ includeSessionId, limit: runtimePolicy.sessionListPageLimit, minIntervalMs: runtimePolicy.sessionListMinRefreshIntervalMs }); const listedSessions = sessionsResult.ok ? workbenchSessionsFromPayload(sessionsResult.data) : []; @@ -218,7 +187,7 @@ export const useWorkbenchStore = defineStore("workbench", () => { sessionDetailLoadingId.value = targetSessionId; recordWorkbenchLoadingState({ scope: "session_detail", active: true, reason: "hydrate", sessionId: targetSessionId }); } - const selected = targetSessionId ? await loadWorkbenchSession(targetSessionId, listedSessions.find((item) => item.sessionId === targetSessionId) ?? null) : null; + const selected = targetSessionId ? loadWorkbenchSession(targetSessionId, listedSessions.find((item) => item.sessionId === targetSessionId) ?? null) : null; if (selected && !isArchivedSession(selected) && routeRequestId && !routeSessionId && requestEpoch === selectionEpoch.value) setActiveSessionSelection(selected.sessionId, "route"); const selectedIsCurrent = Boolean(selected && !isArchivedSession(selected) && isCurrentSessionSelection(requestEpoch, selected.sessionId)); if (sessionsResult.ok || (selectedIsCurrent && selected)) error.value = null; @@ -272,7 +241,8 @@ export const useWorkbenchStore = defineStore("workbench", () => { const response = await workbenchColadaMutations.createAgentSession({ providerProfile: providerProfile.value }); loading.value = false; recordWorkbenchLoadingState({ scope: "workbench", active: false, reason: "create_session", sessionId: activeSessionId.value }); - const created = response.ok ? sessionFromWorkbenchSession(response.data?.session) : null; + const createdResponse = response.ok ? sessionFromWorkbenchSession(response.data?.session, { includeMessages: false }) : null; + const created = createdResponse ? workbenchSessionNavigationSeed(createdResponse.sessionId, createdResponse) : null; if (!response.ok || !created) { error.value = response.error ?? "session create failed"; return; @@ -308,7 +278,7 @@ export const useWorkbenchStore = defineStore("workbench", () => { rememberSessionList(mergeSessionIntoList(sessions.value, existing)); sessionsReady.value = true; } - const selected = await loadWorkbenchSession(requestId, existing); + const selected = loadWorkbenchSession(requestId, existing); if (selected && !isArchivedSession(selected) && !normalized && requestEpoch === selectionEpoch.value) setActiveSessionSelection(selected.sessionId, source); if (!isCurrentSessionRequest(requestEpoch, requestId, selected?.sessionId ?? normalized)) { clearSessionDetailLoading(requestId); @@ -517,28 +487,6 @@ export const useWorkbenchStore = defineStore("workbench", () => { return isTraceActiveStatus(turn?.status) || isTraceActiveStatus(message?.status); } - function fetchSessionMessagesPage(sessionId: string, options: SessionMessagesReadOptions): Promise> { - const key = workbenchSessionMessagesReadKey({ sessionId, limit: options.limit, force: options.force }); - return sessionMessagesReadSingleflight.run(key, () => workbenchColadaQueries.fetchSessionMessages(sessionId, { limit: options.limit, minIntervalMs: runtimePolicy.workbenchSessionMessagesMinRefreshMs, force: options.force }), { reason: options.reason }); - } - - function fetchSessionDetailPage(sessionId: string, options: SessionDetailReadOptions): Promise> { - const key = workbenchSessionDetailReadKey({ sessionId, force: options.force }); - return sessionDetailReadSingleflight.run(key, () => workbenchColadaQueries.fetchSession(sessionId, { includeMessages: false, minIntervalMs: runtimePolicy.workbenchSessionDetailMinRefreshMs, force: options.force }), { reason: options.reason }); - } - - function sessionMessageProjectionWindowLimit(): number { - return boundedProjectionMessageLimit(runtimePolicy.workbenchSessionMessagesWindowLimit, runtimePolicy.sessionListPageLimit); - } - - function applyTurnStatusSnapshot(traceId: string, result: AgentChatResultResponse | TraceSnapshot): void { - const ownerSessionId = traceOwnerSessionId(traceId, traceResultSessionId(result)); - if (!ownerSessionId) return; - if (ownerSessionId === activeSessionId.value) recordActivity(`turn:${firstNonEmptyString(result.lastEventLabel, result.status, result.waitingFor, "status") ?? "status"}`); - rememberTurnStatus(traceId, result); - syncTurnStatusToMessage(traceId, result); - } - function rememberTurnStatus(traceId: string, result: AgentChatResultResponse | TraceSnapshot): void { const id = firstNonEmptyString(result.traceId, traceId); if (!id) return; @@ -560,44 +508,6 @@ export const useWorkbenchStore = defineStore("workbench", () => { }); } - function syncTurnStatusToMessage(traceId: string, result: AgentChatResultResponse | TraceSnapshot): void { - projectTurnAuthorityToMessages(traceId, result as AgentChatResultResponse, "turn-status"); - } - - function projectTurnAuthorityToMessages(traceId: string, result: AgentChatResultResponse, reason: string): boolean { - const authoritySessionId = traceResultSessionId(result); - const ownerSessionId = traceOwnerSessionId(traceId, authoritySessionId); - if (!ownerSessionId) return false; - const source = serverState.value.messagesBySessionId[ownerSessionId] ?? serverState.value.sessionsById[ownerSessionId]?.messages ?? []; - const existing = [...source].reverse().find((message) => messageMatchesTraceAuthority(message, traceId, authoritySessionId, ownerSessionId)) ?? null; - const next = existing ? projectTurnAuthorityMessage(existing, result) : terminalAuthorityMessageFromTurnResult(result); - if (!next) return false; - const message = next.sessionId ? next : { ...next, sessionId: ownerSessionId }; - reduceServerState({ type: "message.upsert", sessionId: ownerSessionId, message }); - const sealed = !messageHasSealedTerminalResult(existing) && messageHasSealedTerminalResult(message); - if (sealed) { - recordWorkbenchRuntimeDiagnostic({ module: "workbench-terminal-authority", sessionId: ownerSessionId, traceId, outcome: "ok", diagnostic: { code: "terminal_direct_seal", reason, source: "turn-authority", valuesRedacted: true } }); - } - return sealed; - } - - function projectTurnAuthorityMessage(message: ChatMessage, result: AgentChatResultResponse): ChatMessage { - const resultStatus = turnResultStatusForMerge(result); - const terminal = turnResultIsTerminalForMerge(result); - const resultError = normalizeAgentError(result.error ?? null); - const resultProjection = projectionFromResult(result); - const mergedRunnerTrace = mergeTerminalResultTrace(message.runnerTrace, result); - const clearCompletedDiagnostics = shouldClearCompletedTurnDiagnostics(resultStatus, resultError); - const runnerTrace = clearCompletedDiagnostics ? clearRunnerTraceTransientDiagnostics(mergedRunnerTrace) : mergedRunnerTrace; - rememberTraceAuthority(runnerTrace); - const error = resultError ?? (terminal || clearCompletedDiagnostics ? null : normalizeAgentError(runnerTrace?.error ?? message.error)); - const agentRun = agentRunFromResult(result, runnerTrace) ?? agentRunFromMessage(message); - const projection = clearCompletedDiagnostics ? nonBlockingProjection(resultProjection) : resultProjection ?? (terminal ? null : runnerTrace.projection ?? message.projection ?? null); - const terminalPatch = terminalMessagePatchFromTurnResult(message, result) ?? {}; - const updatedAt = firstNonEmptyString(result.updatedAt, runnerTrace?.updatedAt, message.updatedAt) ?? new Date().toISOString(); - return { ...message, ...messageTimingPatchForMerge(message, result), ...messageStatusPatchForTerminalMerge(message, resultStatus, terminal), runnerTrace, error, projection, projectionStatus: projection?.projectionStatus ?? null, projectionHealth: projection?.projectionHealth ?? null, blocker: projection?.blocker ?? null, agentRun: agentRun ?? undefined, ...terminalPatch, updatedAt }; - } - async function submitMessage(text: string): Promise { const value = text.trim(); if (!value) return false; @@ -681,12 +591,7 @@ export const useWorkbenchStore = defineStore("workbench", () => { threadId: firstNonEmptyString(response.data.threadId, threadId), status: response.data.status ?? "running" }; - applyTurnStatusSnapshot(canonicalTraceId, response.data); publishWorkbenchProjectionSignal(sessionId, canonicalTraceId, "submit-admitted"); - if (turnResultIsTerminalForMerge(response.data as AgentChatResultResponse)) { - completeTrace(canonicalTraceId, response.data as AgentChatResultResponse); - return true; - } restartRealtime("submit"); return true; } @@ -821,17 +726,6 @@ export const useWorkbenchStore = defineStore("workbench", () => { return source.filter(messageNeedsTraceDetailRead).slice(-runtimePolicy.traceDetailAutoQueueLimit).reverse(); } - function messagesWithTraceAuthority(source: ChatMessage[]): ChatMessage[] { - return source.map((message) => { - if (message.role !== "agent") return message; - const traceId = firstNonEmptyString(message.traceId, message.runnerTrace?.traceId); - const authority = traceId ? traceAuthorityById.value[traceId] ?? null : null; - if (!authority) return message; - const runnerTrace = mergeRunnerTrace(message.runnerTrace, authority); - return { ...message, ...messageTimingPatchForMerge(message, runnerTrace), runnerTrace }; - }); - } - function rememberTraceAuthority(trace: ChatMessage["runnerTrace"]): void { const traceId = firstNonEmptyString(trace?.traceId); if (!traceId || !trace) return; @@ -847,21 +741,23 @@ export const useWorkbenchStore = defineStore("workbench", () => { if (!ownerSessionId) return; const events = Array.isArray(result.events) ? result.events : Array.isArray(result.traceEvents) ? result.traceEvents : []; const activityLabel = firstNonEmptyString(result.lastEventLabel, result.status); - if (ownerSessionId === activeSessionId.value && (events.length > 0 || turnResultIsTerminalForMerge(result))) recordActivity(`trace:${activityLabel ?? "hydrated"}`); + if (ownerSessionId === activeSessionId.value && events.length > 0) recordActivity(`trace:${activityLabel ?? "detail"}`); markWorkbenchTraceEventsReceived({ traceId, events, transport: "detail_history" }); - if (turnResultIsTerminalForMerge(result)) { - const terminalSeal = terminalSealResultWithoutTraceEvents(result); - rememberTurnStatus(traceId, terminalSeal); - projectTurnAuthorityToMessages(traceId, terminalSeal, "trace-detail-read-terminal-seal"); - } updateSessionMessages(ownerSessionId, (source) => source.map((message) => { if (!messageMatchesTraceAuthority(message, traceId, authoritySessionId, ownerSessionId)) return message; - const traceDetailStatus = firstNonEmptyString(result.traceStatus, result.runnerTrace?.traceStatus, result.runnerTrace?.status); - const { traceId: _resultTraceId, sessionId: _resultSessionId, threadId: _resultThreadId, ...resultTraceRest } = result.runnerTrace ?? {}; + const { + traceId: _resultTraceId, + sessionId: _resultSessionId, + threadId: _resultThreadId, + status: _resultStatus, + traceStatus: _resultTraceStatus, + timing: _resultTiming, + terminalEvidence: _resultTerminalEvidence, + ...resultTraceRest + } = result.runnerTrace ?? {}; const nextTrace = { ...resultTraceRest, traceId: optionalString(result.traceId, traceId), - status: traceDetailStatus ?? undefined, sessionId: optionalString(authoritySessionId, message.runnerTrace?.sessionId), threadId: optionalString(result.threadId, message.runnerTrace?.threadId), events, @@ -873,9 +769,7 @@ export const useWorkbenchStore = defineStore("workbench", () => { truncated: result.truncated, nextProjectedSeq: typeof result.nextProjectedSeq === "number" ? result.nextProjectedSeq : null, range: result.range, - traceStatus: result.traceStatus, retention: result.retention, - terminalEvidence: result.terminalEvidence, traceSummary: result.traceSummary, agentRun: result.agentRun, projection: projectionFromResult(result), @@ -883,35 +777,14 @@ export const useWorkbenchStore = defineStore("workbench", () => { projectionHealth: result.projectionHealth ?? result.projection?.projectionHealth, staleMs: result.staleMs ?? result.projection?.staleMs, blocker: result.blocker ?? result.projection?.blocker, - ...messageTimingPatch(result), lastEventLabel: result.lastEventLabel ?? undefined, updatedAt: new Date().toISOString() } as NonNullable; - const resultStatus = turnResultStatusForMerge(result); - const terminal = turnResultIsTerminalForMerge(result); - const resultError = normalizeAgentError(result.error ?? null); - const resultProjection = projectionFromResult(result); - const mergedRunnerTrace = mergeRunnerTrace(message.runnerTrace, nextTrace); - const clearCompletedDiagnostics = shouldClearCompletedTurnDiagnostics(resultStatus, resultError); - const runnerTrace = clearCompletedDiagnostics ? clearRunnerTraceTransientDiagnostics(mergedRunnerTrace) : mergedRunnerTrace; - const error = resultError ?? (terminal || clearCompletedDiagnostics ? null : normalizeAgentError(runnerTrace?.error ?? message.error)); - const projection = clearCompletedDiagnostics ? nonBlockingProjection(resultProjection) : resultProjection ?? (terminal ? null : runnerTrace.projection ?? message.projection ?? null); - const agentRun = agentRunFromResult(result, runnerTrace) ?? agentRunFromMessage(message); + const runnerTrace = mergeRunnerTrace(message.runnerTrace, nextTrace); rememberTraceAuthority(runnerTrace); - return { ...message, ...messageTimingPatchForMerge(message, result), ...messageStatusPatchForTerminalMerge(message, resultStatus, terminal), runnerTrace, error, projection, projectionStatus: projection?.projectionStatus ?? null, projectionHealth: projection?.projectionHealth ?? null, blocker: projection?.blocker ?? null, agentRun: agentRun ?? undefined, updatedAt: new Date().toISOString() }; + return { ...message, runnerTrace }; })); markWorkbenchTraceProjected(traceId); - if (traceResultHasTerminalEvidence(result) && !turnResultIsTerminalForMerge(result)) { - } - if (turnResultIsTerminalForMerge(result)) { - rememberTurnStatus(traceId, result); - if (ownerSessionId === activeSessionId.value) { - chatPending.value = false; - currentRequest.value = null; - void clearActiveTrace(traceId, "trace-detail-read-terminal"); - restartRealtime("trace-detail-read-terminal"); - } - } } function setProviderProfile(value: ProviderProfile): void { @@ -1144,18 +1017,9 @@ export const useWorkbenchStore = defineStore("workbench", () => { function applyWorkbenchRealtimePlanStep(step: WorkbenchRealtimeApplyStep): void { switch (step.type) { - case "apply-trace-snapshot": - applyRealtimeTraceSnapshot(step.traceId, step.snapshot); - return; case "apply-trace-event": applyRealtimeTraceEvent(step.traceId, step.event, step.snapshot, step.realtimeEvent); return; - case "apply-message-snapshot": - applyRealtimeMessageSnapshot(step.realtimeEvent); - return; - case "apply-turn-snapshot": - applyRealtimeTurnSnapshot(step.turn); - return; case "apply-projection-error": applyRealtimeProjectionError(step.realtimeEvent); return; @@ -1165,27 +1029,6 @@ export const useWorkbenchStore = defineStore("workbench", () => { } } - function applyRealtimeMessageSnapshot(event: WorkbenchRealtimeEvent): void { - const message = event.message ? normalizeChatMessage(event.message) : null; - const sessionId = normalizeWorkbenchSessionId(firstNonEmptyString(event.sessionId, message?.sessionId)); - if (!message || !sessionId) return; - const activeId = activeSessionId.value; - if (activeId && sessionId !== activeId) return; - reduceServerState({ type: "message.snapshot", sessionId, message }); - } - - function applyRealtimeTraceSnapshot(traceId: string | null | undefined, snapshot: WorkbenchRealtimeEvent["snapshot"]): void { - const id = firstNonEmptyString(traceId, snapshot?.traceId); - if (!id || !snapshot) return; - const sessionId = traceResultSessionId(snapshot); - if (!shouldApplyActiveTraceAuthority(id, sessionId)) return; - if (traceTerminalBodyIsVisible(id, sessionId)) { - recordWorkbenchRuntimeDiagnostic({ module: "workbench-terminal-priority", sessionId, traceId: id, outcome: "ok", diagnostic: { code: "terminal_low_priority_sse_trace_skip", source: "realtime-trace-snapshot", valuesRedacted: true } }); - return; - } - applyTraceSnapshot(id, realtimeSnapshotToTraceSnapshot(id, snapshot)); - } - function applyRealtimeTraceEvent(traceId: string | null | undefined, event: WorkbenchRealtimeEvent["event"], snapshot: WorkbenchRealtimeEvent["snapshot"], realtimeEvent?: WorkbenchRealtimeEvent | null): void { const id = firstNonEmptyString(traceId, event?.traceId, snapshot?.traceId); if (!id) return; @@ -1199,29 +1042,30 @@ export const useWorkbenchStore = defineStore("workbench", () => { const events = event ? [event] : Array.isArray(snapshot?.events) ? snapshot.events : []; markWorkbenchTraceEventsReceived({ traceId: id, events, transport: "sse", serverSentAt: realtimeEvent?.serverSentAt, eventCreatedAt: realtimeEvent?.eventCreatedAt, traceSeq: realtimeEvent?.traceSeq ?? realtimeEvent?.cursor?.traceSeq }); if (liveKafkaEnvelope && event) { - applyLiveKafkaBusinessEvent(id, sessionId, event, new Date().toISOString()); + const context = recordValue(realtimeEvent?.context); + applyLiveKafkaBusinessEvent(id, sessionId, { ...event, threadId: firstNonEmptyString(event.threadId, context?.threadId) }, new Date().toISOString()); return; } - const eventSnapshot = snapshot ?? { traceId: id, status: event?.status, events }; - applyTraceSnapshot(id, realtimeSnapshotToTraceSnapshot(id, eventSnapshot, events)); + recordWorkbenchRuntimeDiagnostic({ module: "workbench-realtime-authority", sessionId, traceId: id, outcome: "warning", diagnostic: { code: "non_kafka_business_projection_ignored", source: "realtime-trace-event", valuesRedacted: true } }); } function applyLiveKafkaBusinessEvent(traceId: string, authoritySessionId: string | null, event: TraceEvent, receivedAt: string): void { const eventSessionId = firstNonEmptyString(event.sessionId); const ownerSessionId = traceOwnerSessionId(traceId, authoritySessionId ?? eventSessionId, true); if (!ownerSessionId) return; + const threadId = firstNonEmptyString(event.threadId); if (workbenchLiveKafkaProjectionTarget(event) === "user") { const messageId = firstNonEmptyString(event.userMessageId, event.messageId); const previousUserMessage = messageId ? (serverState.value.messagesBySessionId[ownerSessionId] ?? []).find((message) => (message.messageId ?? message.id) === messageId) ?? null : null; - const userMessage = projectWorkbenchLiveKafkaUserMessage({ previous: previousUserMessage, traceId, sessionId: ownerSessionId, event, receivedAt }); + const userMessage = projectWorkbenchLiveKafkaUserMessage({ previous: previousUserMessage, traceId, sessionId: ownerSessionId, event, receivedAt, threadId }); if (userMessage) reduceServerState({ type: "message.upsert", sessionId: ownerSessionId, message: userMessage }); else error.value = "workbench_live_user_message_invalid"; return; } const previousMessage = (serverState.value.messagesBySessionId[ownerSessionId] ?? []).find((message) => messageMatchesTraceAuthority(message, traceId, authoritySessionId ?? eventSessionId, ownerSessionId, true)) ?? null; - const message = projectWorkbenchLiveKafkaMessage({ previous: previousMessage, traceId, sessionId: ownerSessionId, event, receivedAt }); + const message = projectWorkbenchLiveKafkaMessage({ previous: previousMessage, traceId, sessionId: ownerSessionId, event, receivedAt, threadId }); reduceServerState({ type: "message.upsert", sessionId: ownerSessionId, message }); if (message.runnerTrace) rememberTraceAuthority(message.runnerTrace); markWorkbenchTraceProjected(traceId); @@ -1243,28 +1087,13 @@ export const useWorkbenchStore = defineStore("workbench", () => { updatedAt: receivedAt } as AgentChatResultResponse); const existing = sessions.value.find((session) => session.sessionId === ownerSessionId) ?? null; - if (existing) rememberSessionList(mergeSessionIntoList(sessions.value, { ...existing, status, lastTraceId: traceId, updatedAt: receivedAt })); + if (existing) rememberSessionList(mergeSessionIntoList(sessions.value, { ...existing, threadId: message.threadId ?? existing.threadId, status, lastTraceId: traceId, updatedAt: receivedAt })); if (terminal && currentRequest.value?.traceId === traceId) { chatPending.value = false; currentRequest.value = null; } } - function applyRealtimeTurnSnapshot(turn: Record): void { - const traceId = firstNonEmptyString(turn.traceId); - if (!traceId) return; - if (!shouldApplyActiveTraceAuthority(traceId, traceResultSessionId(turn))) return; - const activeId = activeSessionId.value; - const status = firstNonEmptyString(turn.status) ?? undefined; - const terminalTurn = turn.terminal === true || isTerminalMessageStatus(status); - if (activeId && !terminalTurn && !messages.value.some((message) => firstNonEmptyString(message.traceId, message.runnerTrace?.traceId) === traceId)) { - recordWorkbenchRuntimeDiagnostic({ module: "workbench-realtime-authority", sessionId: activeId, traceId, outcome: "ok", diagnostic: { code: "workbench_realtime_turn_gap_no_legacy_repair", reason: "realtime-turn-gap", source: "turn-snapshot", valuesRedacted: true } }); - } - const result = { ...turn, traceId, status, running: turn.running === true, terminal: turn.terminal === true, sessionId: firstNonEmptyString(turn.sessionId) ?? undefined, threadId: firstNonEmptyString(turn.threadId) ?? undefined, agentRun: turn.agentRun as AgentRunProvenance | undefined } as AgentChatResultResponse; - rememberTurnStatus(traceId, result); - scheduleRealtimeTurnProjection({ traceId, result, terminalTurn }); - } - function realtimeEventName(event: WorkbenchRealtimeEvent): string { switch (event.type) { case "trace.snapshot": @@ -1284,49 +1113,6 @@ export const useWorkbenchStore = defineStore("workbench", () => { } } - function scheduleRealtimeTurnProjection(item: RealtimeTurnProjectionItem): void { - realtimeTurnProjectionQueue.set(item.traceId, item); - scheduleRealtimeTurnProjectionFlush(); - } - - function scheduleRealtimeTurnProjectionFlush(): void { - if (realtimeTurnProjectionScheduled) return; - realtimeTurnProjectionScheduled = true; - const flush = () => flushRealtimeTurnProjection(); - if (typeof window !== "undefined" && typeof window.requestAnimationFrame === "function") { - window.requestAnimationFrame(flush); - return; - } - setTimeout(flush, Math.max(0, runtimePolicy.workbenchRealtimeFlushYieldMs)); - } - - function flushRealtimeTurnProjection(): void { - realtimeTurnProjectionScheduled = false; - const next = realtimeTurnProjectionQueue.values().next().value as RealtimeTurnProjectionItem | undefined; - if (!next) return; - realtimeTurnProjectionQueue.delete(next.traceId); - const startedAt = performanceNowMs(); - if (shouldApplyActiveTraceAuthority(next.traceId, traceResultSessionId(next.result))) { - syncTurnStatusToMessage(next.traceId, next.result); - } - recordWorkbenchRuntimeDiagnostic({ - module: "workbench-turn-status", - sessionId: traceResultSessionId(next.result), - traceId: next.traceId, - outcome: "ok", - diagnostic: { - code: "workbench_realtime_turn_projection_budget", - reason: "realtime-turn-snapshot", - source: "turn-snapshot-frame-queue", - remainingCount: realtimeTurnProjectionQueue.size, - terminal: next.terminalTurn, - flushDurationMs: Math.round(performanceNowMs() - startedAt), - valuesRedacted: true - } - }); - if (realtimeTurnProjectionQueue.size > 0) scheduleRealtimeTurnProjectionFlush(); - } - function installRealtimeVisibilityHandler(): void { if (typeof document === "undefined") return; document.addEventListener("visibilitychange", () => { @@ -1371,68 +1157,6 @@ export const useWorkbenchStore = defineStore("workbench", () => { restartRealtime("cross-tab-session-projection", true); } - function applyTraceSnapshot(traceId: string, snapshot: TraceSnapshot, canonicalSession = false): void { - const trace = snapshotToRunnerTrace({ - ...snapshot, - traceStatus: firstNonEmptyString(snapshot.traceStatus, snapshot.status) ?? undefined - }); - const authoritySessionId = canonicalSession - ? normalizeWorkbenchSessionId(firstNonEmptyString(trace.sessionId, snapshot.sessionId)) - : traceResultSessionId(trace) ?? traceResultSessionId(snapshot); - const ownerSessionId = traceOwnerSessionId(traceId, authoritySessionId, canonicalSession); - if (!ownerSessionId) return; - updateSessionMessages(ownerSessionId, (source) => source.map((message) => { - if (!messageMatchesTraceAuthority(message, traceId, authoritySessionId, ownerSessionId, canonicalSession)) return message; - const mergedRunnerTrace = mergeRunnerTrace(message.runnerTrace, trace); - const traceStatus = normalizedStatusText(trace.status ?? snapshot.status) ?? null; - const clearCompletedDiagnostics = shouldClearCompletedTurnDiagnostics(traceStatus, null); - const runnerTrace = clearCompletedDiagnostics ? clearRunnerTraceTransientDiagnostics(mergedRunnerTrace) : mergedRunnerTrace; - const terminal = turnResultIsTerminalForMerge({ ...(snapshot as AgentChatResultResponse), status: traceStatus } as AgentChatResultResponse); - rememberTraceAuthority(runnerTrace); - const error = clearCompletedDiagnostics ? null : message.role === "agent" ? normalizeAgentError(runnerTrace.error ?? message.error) : normalizeAgentError(message.error); - const projection = clearCompletedDiagnostics ? nonBlockingProjection(trace.projection ?? null) : trace.projection ?? runnerTrace.projection ?? message.projection ?? null; - return { ...message, ...messageTimingPatchForMerge(message, trace), ...messageStatusPatchForTerminalMerge(message, traceStatus, terminal), runnerTrace, error: clearCompletedDiagnostics ? null : error ?? message.error ?? null, projection, projectionStatus: projection?.projectionStatus ?? null, projectionHealth: projection?.projectionHealth ?? null, blocker: projection?.blocker ?? null, updatedAt: new Date().toISOString() }; - })); - markWorkbenchTraceProjected(traceId); - } - - function completeTrace(traceId: string, result: AgentChatResultResponse): void { - const authoritySessionId = traceResultSessionId(result); - const ownerSessionId = traceOwnerSessionId(traceId, authoritySessionId); - if (!ownerSessionId) return; - if (ownerSessionId === activeSessionId.value) recordActivity(`trace:terminal:${firstNonEmptyString(result.lastEventLabel, result.status, "completed") ?? "completed"}`); - projectTurnAuthorityToMessages(traceId, result, "complete-trace"); - rememberTurnStatus(traceId, result); - markWorkbenchTraceProjected(traceId); - if (ownerSessionId === activeSessionId.value) { - chatPending.value = false; - currentRequest.value = null; - void clearActiveTrace(traceId, "trace-terminal"); - restartRealtime("trace-terminal"); - } - } - - function applyTerminalResultDiagnostics(traceId: string, result: AgentChatResultResponse): void { - const authoritySessionId = traceResultSessionId(result); - const ownerSessionId = traceOwnerSessionId(traceId, authoritySessionId); - if (!ownerSessionId) return; - updateSessionMessages(ownerSessionId, (source) => source.map((message) => { - if (!messageMatchesTraceAuthority(message, traceId, authoritySessionId, ownerSessionId)) return message; - const resultStatus = turnResultStatusForMerge(result); - const terminal = turnResultIsTerminalForMerge(result); - const resultError = normalizeAgentError(result.error ?? null); - const resultProjection = projectionFromResult(result); - const mergedRunnerTrace = mergeTerminalResultTrace(message.runnerTrace, result); - const clearCompletedDiagnostics = shouldClearCompletedTurnDiagnostics(resultStatus, resultError); - const runnerTrace = clearCompletedDiagnostics ? clearRunnerTraceTransientDiagnostics(mergedRunnerTrace) : mergedRunnerTrace; - rememberTraceAuthority(runnerTrace); - const error = resultError ?? (terminal || clearCompletedDiagnostics ? null : normalizeAgentError(runnerTrace?.error ?? message.error)); - const agentRun = agentRunFromResult(result, runnerTrace) ?? agentRunFromMessage(message); - const projection = clearCompletedDiagnostics ? nonBlockingProjection(resultProjection) : resultProjection ?? (terminal ? null : runnerTrace.projection ?? message.projection ?? null); - return { ...message, ...messageTimingPatchForMerge(message, result), ...messageStatusPatchForTerminalMerge(message, resultStatus, terminal), title: normalizeWorkbenchMessageTitle(message.role, message.title), runnerTrace, error, projection, projectionStatus: projection?.projectionStatus ?? null, projectionHealth: projection?.projectionHealth ?? null, blocker: projection?.blocker ?? null, agentRun: agentRun ?? undefined, updatedAt: new Date().toISOString() }; - })); - } - function applyRealtimeProjectionError(event: WorkbenchRealtimeEvent): void { const traceId = firstNonEmptyString(event.traceId, realtimeTraceId()); const errorRecord = recordValue(event.error) ?? event; @@ -1608,12 +1332,10 @@ export const useWorkbenchStore = defineStore("workbench", () => { function applySelectedSessionDetail(session: WorkbenchSessionRecord, source: SessionSelectionSource = "system"): void { setActiveSessionSelection(session.sessionId, source); - const projectedMessages = messagesWithTraceAuthority(messagesFromSession(session)); + const projectedMessages = messagesFromSession(session); rememberSessionDetail({ ...session, messages: projectedMessages, messageCount: session.messageCount ?? projectedMessages.length }); sessionsReady.value = true; currentRequest.value = null; - void readTraceEventsForMessages(messages.value); - reattachRestoredActiveTrace(); restartRealtime("apply-selected-session"); } @@ -1633,26 +1355,15 @@ export const useWorkbenchStore = defineStore("workbench", () => { restartRealtime("session-read-unavailable"); } - async function loadWorkbenchSession(sessionId: string, seed: WorkbenchSessionRecord | null = null): Promise { + function loadWorkbenchSession(sessionId: string, seed: WorkbenchSessionRecord | null = null): WorkbenchSessionRecord | null { const requestId = normalizeWorkbenchSessionRouteId(sessionId); if (!requestId) return null; const normalizedRequestId = normalizeWorkbenchSessionId(requestId); - const messageLimit = sessionMessageProjectionWindowLimit(); - const [detail, messagePage] = await Promise.all([ - fetchSessionDetailPage(requestId, { reason: "load-session:detail", force: true }), - normalizedRequestId ? fetchSessionMessagesPage(normalizedRequestId, { limit: messageLimit, reason: "load-session:messages", force: true }) : Promise.resolve(null) - ]); - if (!detail.ok) return null; - const detailSession = sessionFromWorkbenchSession(detail.data?.session, { includeMessages: false }); - const id = detailSession?.sessionId ?? normalizedRequestId ?? seed?.sessionId; + const id = normalizedRequestId ?? seed?.sessionId; if (!id) return null; - const fallbackMessages = seed?.sessionId === id ? seed.messages ?? [] : []; - const base = detailSession ? { ...detailSession, messages: fallbackMessages } : seed; - if (!base) return null; - const hydratedMessages = messagePage?.ok && Array.isArray(messagePage.data?.messages) - ? messagePage.data.messages.map((message) => normalizeChatMessage(message as ChatMessage)) - : fallbackMessages; - return { ...base, sessionId: id, messages: hydratedMessages, messageCount: base.messageCount ?? hydratedMessages.length }; + const projected = serverState.value.sessionsById[id]; + if (projected) return projected; + return seed?.sessionId === id ? workbenchSessionNavigationSeed(id, seed) : workbenchSessionNavigationSeed(id); } function reattachRestoredActiveTrace(): void { @@ -1660,10 +1371,6 @@ export const useWorkbenchStore = defineStore("workbench", () => { if (traceId) reattachTrace(traceId); } - function performanceNowMs(): number { - return typeof performance !== "undefined" && typeof performance.now === "function" ? performance.now() : Date.now(); - } - function clearRawHwlabIngress(): void { rawHwlabIngress.value = createRawHwlabIngressState(rawHwlabIngress.value.scopeKey); } @@ -1677,7 +1384,19 @@ export const useWorkbenchStore = defineStore("workbench", () => { function workbenchSessionsFromPayload(payload: unknown): WorkbenchSessionRecord[] { const record = recordValue(payload); const sessions = Array.isArray(record?.sessions) ? record.sessions : []; - return sessions.map((item) => sessionFromWorkbenchSession(item)).filter((item): item is WorkbenchSessionRecord => Boolean(item)); + return sessions + .map((item) => sessionFromWorkbenchSession(item, { includeMessages: false })) + .filter((item): item is WorkbenchSessionRecord => Boolean(item)) + .map((session) => workbenchSessionNavigationSeed(session.sessionId, session)); +} + +function workbenchSessionNavigationSeed(sessionId: string, source: WorkbenchSessionRecord | null = null): WorkbenchSessionRecord { + return { + sessionId, + threadId: firstNonEmptyString(source?.threadId), + messages: [], + messageCount: 0 + }; } function sessionFromWorkbenchSession(value: unknown, options: { includeMessages?: boolean } = {}): WorkbenchSessionRecord | null {