fix: 删除 Workbench 旧实时架构分支
This commit is contained in:
@@ -6,7 +6,7 @@ import { test } from "bun:test";
|
||||
|
||||
import { createCloudApiBunServer } from "./bun-server.ts";
|
||||
import { buildCloudApiReadiness } from "./health-contract.ts";
|
||||
import { decodeCanonicalAgentRunKafkaMessage, kafkaEventBridgeConfig, liveKafkaOtelSpanAttributes, projectAgentRunKafkaEventToHwlabEvent, publishAgentRunKafkaMessageLive, relayHwlabKafkaOutboxOnce, shouldEmitLiveKafkaOtelSpan, startHwlabKafkaEventBridge } from "./kafka-event-bridge.ts";
|
||||
import { decodeCanonicalAgentRunKafkaMessage, kafkaEventBridgeConfig, projectAgentRunKafkaEventToHwlabEvent, relayHwlabKafkaOutboxOnce, startHwlabKafkaEventBridge } from "./kafka-event-bridge.ts";
|
||||
|
||||
const PROJECTOR_ENV = Object.freeze({
|
||||
HWLAB_KAFKA_DIRECT_PUBLISH_ENABLED: "false",
|
||||
@@ -31,23 +31,6 @@ const PROJECTOR_ENV = Object.freeze({
|
||||
HWLAB_KAFKA_OUTBOX_RELAY_RETRY_BACKOFF_MS: "1000"
|
||||
});
|
||||
|
||||
const LIVE_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",
|
||||
HWLAB_KAFKA_ENABLED: "true",
|
||||
HWLAB_KAFKA_AGENTRUN_EVENT_CONSUME_ENABLED: "true",
|
||||
HWLAB_KAFKA_BOOTSTRAP_SERVERS: "kafka.test:9092",
|
||||
HWLAB_KAFKA_AGENTRUN_EVENT_TOPIC: "agentrun.event.v1",
|
||||
HWLAB_KAFKA_EVENT_TOPIC: "hwlab.event.v1",
|
||||
HWLAB_KAFKA_CLIENT_ID: "hwlab-live-test",
|
||||
HWLAB_KAFKA_AGENTRUN_EVENT_GROUP_ID: "hwlab-live-bridge-fixed",
|
||||
HWLAB_KAFKA_HWLAB_EVENT_GROUP_ID: "hwlab-live-fanout-fixed"
|
||||
});
|
||||
|
||||
test("projects AgentRun assistant_message Kafka event into HWLAB trace event", () => {
|
||||
const projected = projectAgentRunKafkaEventToHwlabEvent({
|
||||
schema: "agentrun.event.v1",
|
||||
@@ -322,63 +305,6 @@ test("projects AgentRun terminal_status Kafka event into terminal HWLAB event",
|
||||
assert.equal(projected.event.message, "AgentRun completed");
|
||||
});
|
||||
|
||||
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"); } });
|
||||
void runtimeStore;
|
||||
const result = await publishAgentRunKafkaMessageLive({ topic: "agentrun.event.v1", partition: 3, message: canonicalMessage({ offset: "71" }) }, {
|
||||
producer: { async send(input) { sent.push(input); } },
|
||||
config: { clientId: "hwlab-live-test", hwlabTopic: "hwlab.event.v1" },
|
||||
logger: null
|
||||
});
|
||||
assert.equal(result.published, true);
|
||||
assert.equal(sent.length, 1);
|
||||
assert.equal(JSON.parse(sent[0].messages[0].value).sourceEvent.offset, "71");
|
||||
});
|
||||
|
||||
test("live Kafka OTel sampling is deterministic and bounded without payload attributes", () => {
|
||||
const envelope = (agentRunEventType, sourceSeq, event = {}) => ({
|
||||
schema: "hwlab.event.v1",
|
||||
eventId: `hwlab:evt_otel_${sourceSeq}`,
|
||||
sourceEventId: `evt_otel_${sourceSeq}`,
|
||||
traceId: "trc_live_otel",
|
||||
hwlabSessionId: "ses_live_otel",
|
||||
runId: "run_live_otel",
|
||||
commandId: "cmd_live_otel",
|
||||
context: { agentRunEventType, sourceSeq, valuesRedacted: true },
|
||||
event: { type: "backend", eventType: "backend", terminal: false, message: "must-not-leak", payload: { password: "must-not-leak" }, ...event }
|
||||
});
|
||||
|
||||
assert.equal(shouldEmitLiveKafkaOtelSpan(envelope("assistant_message", 1, { type: "assistant" })), true);
|
||||
assert.equal(shouldEmitLiveKafkaOtelSpan(envelope("tool_call", 2, { type: "tool" })), true);
|
||||
assert.equal(shouldEmitLiveKafkaOtelSpan(envelope("terminal_status", 3, { type: "result", terminal: true })), true);
|
||||
const decisions = Array.from({ length: 64 }, (_, index) => shouldEmitLiveKafkaOtelSpan(envelope("command_output", index + 1, { type: "output" })));
|
||||
assert.deepEqual(decisions, Array.from({ length: 64 }, (_, index) => (index + 1) % 32 === 0));
|
||||
assert.deepEqual(decisions, Array.from({ length: 64 }, (_, index) => shouldEmitLiveKafkaOtelSpan(envelope("command_output", index + 1, { type: "output" }))));
|
||||
|
||||
const attributes = liveKafkaOtelSpanAttributes(envelope("assistant_message", 1, { type: "assistant" }), {
|
||||
topic: "hwlab.event.v1",
|
||||
partition: 2,
|
||||
offset: "14"
|
||||
});
|
||||
assert.deepEqual(attributes, {
|
||||
businessTraceId: "trc_live_otel",
|
||||
hwlabSessionId: "ses_live_otel",
|
||||
runId: "run_live_otel",
|
||||
commandId: "cmd_live_otel",
|
||||
topic: "hwlab.event.v1",
|
||||
partition: 2,
|
||||
offset: "14",
|
||||
sourceTopic: null,
|
||||
sourcePartition: null,
|
||||
sourceOffset: null,
|
||||
eventType: "assistant_message",
|
||||
terminal: false,
|
||||
valuesRedacted: true
|
||||
});
|
||||
assert.equal(JSON.stringify(attributes).includes("must-not-leak"), false);
|
||||
});
|
||||
|
||||
test("Workbench bridge ignores obsolete architecture env and fixes transactional authority", () => {
|
||||
const config = kafkaEventBridgeConfig({
|
||||
...PROJECTOR_ENV,
|
||||
@@ -390,17 +316,15 @@ test("Workbench bridge ignores obsolete architecture env and fixes transactional
|
||||
HWLAB_WORKBENCH_PROJECTION_REALTIME_ENABLED: "false"
|
||||
});
|
||||
|
||||
assert.deepEqual(config?.capabilities, {
|
||||
directPublish: false,
|
||||
liveKafkaSse: false,
|
||||
kafkaRefreshReplay: false,
|
||||
transactionalProjector: true,
|
||||
projectionOutboxRelay: true,
|
||||
projectionRealtime: true
|
||||
assert.equal(config?.projectorGroupId, "hwlab-projector-test-v1");
|
||||
assert.deepEqual(config?.relay, {
|
||||
intervalMs: 60000,
|
||||
batchSize: 1,
|
||||
leaseMs: 30000,
|
||||
sendTimeoutMs: 10000,
|
||||
retryBackoffMs: 1000
|
||||
});
|
||||
assert.equal(config?.directPublishGroupId, null);
|
||||
assert.equal(config?.hwlabEventGroupId, null);
|
||||
assert.equal(config?.refreshReplay, null);
|
||||
assert.deepEqual(Object.keys(config ?? {}).sort(), ["agentRunTopic", "brokers", "clientId", "hwlabTopic", "projectorGroupId", "projectorHeartbeatIntervalMs", "relay"]);
|
||||
});
|
||||
|
||||
test("Workbench bridge does not require the six obsolete architecture env", () => {
|
||||
|
||||
@@ -4,9 +4,7 @@
|
||||
import { createHash, randomUUID } from "node:crypto";
|
||||
|
||||
import { Kafka, logLevel } from "kafkajs";
|
||||
import { emitCodeAgentOtelSpan } from "./otel-trace.ts";
|
||||
import { buildWorkbenchProjectionEventFacts } from "./workbench-projection-writer.ts";
|
||||
import { workbenchRealtimeCapabilities } from "./workbench-realtime-capabilities.ts";
|
||||
|
||||
const TRUE_VALUES = new Set(["1", "true", "yes", "on"]);
|
||||
export const DEFAULT_AGENTRUN_EVENT_TOPIC = "agentrun.event.v1";
|
||||
@@ -19,7 +17,6 @@ export const AGENTRUN_STDIO_RECONSTRUCTION_KIND = "stdio-derived partial reconst
|
||||
const DEFAULT_CLIENT_ID = "hwlab-v03-cloud-api";
|
||||
const DEFAULT_QUERY_TIMEOUT_MS = 5000;
|
||||
const DEFAULT_QUERY_LIMIT = 50;
|
||||
const LIVE_KAFKA_COMMAND_OUTPUT_OTEL_SAMPLE_MODULUS = 32;
|
||||
const REQUIRED_KAFKA_ENV = Object.freeze([
|
||||
"HWLAB_KAFKA_BOOTSTRAP_SERVERS",
|
||||
"HWLAB_KAFKA_AGENTRUN_EVENT_TOPIC",
|
||||
@@ -31,7 +28,6 @@ export function kafkaEventBridgeConfig(env = process.env) {
|
||||
const legacyKafkaEnabled = truthy(env.HWLAB_KAFKA_AGENTRUN_EVENT_CONSUME_ENABLED ?? env.HWLAB_KAFKA_AGENTRUN_CONSUME_ENABLED ?? env.HWLAB_KAFKA_ENABLED);
|
||||
const kafkaConfigured = REQUIRED_KAFKA_ENV.every((name) => stringValue(env[name]));
|
||||
if (!legacyKafkaEnabled && !kafkaConfigured) return null;
|
||||
const capabilities = workbenchRealtimeCapabilities();
|
||||
const required = [
|
||||
...REQUIRED_KAFKA_ENV,
|
||||
"HWLAB_KAFKA_PROJECTOR_GROUP_ID",
|
||||
@@ -66,21 +62,17 @@ export function kafkaEventBridgeConfig(env = process.env) {
|
||||
throw error;
|
||||
}
|
||||
return {
|
||||
capabilities,
|
||||
brokers,
|
||||
agentRunTopic,
|
||||
hwlabTopic,
|
||||
clientId,
|
||||
directPublishGroupId: null,
|
||||
projectorGroupId: stringValue(env.HWLAB_KAFKA_PROJECTOR_GROUP_ID),
|
||||
hwlabEventGroupId: null,
|
||||
refreshReplay: null,
|
||||
projectorHeartbeatIntervalMs: strictPositiveInteger(env.HWLAB_KAFKA_PROJECTOR_HEARTBEAT_INTERVAL_MS, "HWLAB_KAFKA_PROJECTOR_HEARTBEAT_INTERVAL_MS"),
|
||||
relay
|
||||
};
|
||||
}
|
||||
|
||||
export function startHwlabKafkaEventBridge({ env = process.env, logger = console, kafkaFactory = defaultKafkaFactory, runtimeStore = null, now = () => new Date().toISOString(), otelSpanEmitter = emitCodeAgentOtelSpan } = {}) {
|
||||
export function startHwlabKafkaEventBridge({ env = process.env, logger = console, kafkaFactory = defaultKafkaFactory, runtimeStore = null, now = () => new Date().toISOString() } = {}) {
|
||||
const config = kafkaEventBridgeConfig(env);
|
||||
if (!config) return { started: false, reason: "disabled_or_unconfigured", stop() {}, subscribeProjectionCommits() { return () => {}; }, valuesPrinted: false };
|
||||
return startTransactionalHwlabKafkaEventBridge({ config, logger, kafkaFactory, runtimeStore, now });
|
||||
@@ -157,7 +149,6 @@ function startTransactionalHwlabKafkaEventBridge({ config, logger = console, kaf
|
||||
hwlabTopic: config.hwlabTopic,
|
||||
clientId: config.clientId,
|
||||
projectorGroupId: config.projectorGroupId,
|
||||
capabilities: config.capabilities,
|
||||
valuesPrinted: false
|
||||
});
|
||||
})();
|
||||
@@ -191,7 +182,7 @@ function startTransactionalHwlabKafkaEventBridge({ config, logger = console, kaf
|
||||
? await runtimeStore.workbenchTransactionalRealtimeReadiness()
|
||||
: { ready: true, valuesRedacted: true };
|
||||
const projectorStatus = await runtimeStore.hwlabKafkaProjectorStatus();
|
||||
return { status: transactionalReadiness.ready === false ? "blocked" : "ready", capabilities: config.capabilities, transactionalReadiness, ...projectorStatus, valuesRedacted: true };
|
||||
return { status: transactionalReadiness.ready === false ? "blocked" : "ready", transactionalReadiness, ...projectorStatus, valuesRedacted: true };
|
||||
},
|
||||
startupStatus() {
|
||||
if (startupError) return { status: "blocked", errorCode: startupError.code ?? "hwlab_kafka_projector_start_failed", message: errorMessage(startupError), valuesRedacted: true };
|
||||
@@ -206,269 +197,6 @@ function startTransactionalHwlabKafkaEventBridge({ config, logger = console, kaf
|
||||
};
|
||||
}
|
||||
|
||||
function startLiveHwlabKafkaEventBridge({ config, env = process.env, logger = console, kafkaFactory = defaultKafkaFactory, otelSpanEmitter = emitCodeAgentOtelSpan } = {}) {
|
||||
let stopped = false;
|
||||
let bridgeConsumer = null;
|
||||
let fanoutConsumer = null;
|
||||
let producer = null;
|
||||
let startupError = null;
|
||||
let startupReady = false;
|
||||
const subscribers = new Set();
|
||||
const lagByPartition = new Map();
|
||||
const publishLiveEnvelope = (envelope, transport) => {
|
||||
for (const listener of [...subscribers]) {
|
||||
try {
|
||||
listener(envelope, transport);
|
||||
} catch (error) {
|
||||
logWarn(logger, "hwlab-live-kafka-subscriber-failed", { message: errorMessage(error), valuesPrinted: false });
|
||||
}
|
||||
}
|
||||
};
|
||||
const ready = (async () => {
|
||||
const kafka = kafkaFactory(config);
|
||||
if (config.capabilities.directPublish) {
|
||||
bridgeConsumer = kafka.consumer({ groupId: config.directPublishGroupId, allowAutoTopicCreation: false });
|
||||
producer = kafka.producer({ allowAutoTopicCreation: false });
|
||||
}
|
||||
if (config.capabilities.liveKafkaSse) fanoutConsumer = kafka.consumer({ groupId: config.hwlabEventGroupId, allowAutoTopicCreation: false });
|
||||
if (producer) await producer.connect();
|
||||
if (bridgeConsumer) await bridgeConsumer.connect();
|
||||
if (fanoutConsumer) await fanoutConsumer.connect();
|
||||
if (bridgeConsumer) await bridgeConsumer.subscribe({ topic: config.agentRunTopic, fromBeginning: false });
|
||||
if (fanoutConsumer) await fanoutConsumer.subscribe({ topic: config.hwlabTopic, fromBeginning: false });
|
||||
if (fanoutConsumer) await fanoutConsumer.run({
|
||||
eachBatchAutoResolve: false,
|
||||
eachBatch: async ({ batch, resolveOffset, heartbeat, isRunning, isStale }) => {
|
||||
let lastOffset = null;
|
||||
for (const message of batch.messages) {
|
||||
if (stopped || !isRunning() || isStale()) break;
|
||||
const envelope = parseJson(message?.value ? Buffer.from(message.value).toString("utf8") : "");
|
||||
if (!envelope || envelope.schema !== "hwlab.event.v1") {
|
||||
logWarn(logger, "hwlab-live-kafka-envelope-ignored", { reason: "schema-invalid", valuesPrinted: false });
|
||||
} else {
|
||||
const transport = { topic: batch.topic, partition: batch.partition, offset: message.offset };
|
||||
emitLiveKafkaOtelSpan("hwlab.kafka.live.fanout_receive", envelope, transport, { env, otelSpanEmitter });
|
||||
publishLiveEnvelope(envelope, transport);
|
||||
}
|
||||
resolveOffset(message.offset);
|
||||
lastOffset = message.offset;
|
||||
await heartbeat();
|
||||
}
|
||||
if (lastOffset !== null) lagByPartition.set(batch.partition, kafkaPartitionLag(batch.highWatermark, lastOffset));
|
||||
}
|
||||
});
|
||||
if (bridgeConsumer) await bridgeConsumer.run({
|
||||
eachMessage: async ({ topic, partition, message }) => {
|
||||
if (stopped) return;
|
||||
await publishAgentRunKafkaMessageLive({ topic, partition, message }, { producer, config, env, logger, otelSpanEmitter });
|
||||
}
|
||||
});
|
||||
startupReady = true;
|
||||
logInfo(logger, "hwlab-live-kafka-event-bridge-started", {
|
||||
capabilities: config.capabilities,
|
||||
agentRunTopic: config.agentRunTopic,
|
||||
hwlabTopic: config.hwlabTopic,
|
||||
clientId: config.clientId,
|
||||
directPublishGroupId: config.directPublishGroupId,
|
||||
hwlabEventGroupId: config.hwlabEventGroupId,
|
||||
valuesPrinted: false
|
||||
});
|
||||
})();
|
||||
void ready.catch((error) => {
|
||||
startupError = error;
|
||||
logWarn(logger, "hwlab-live-kafka-event-bridge-start-failed", {
|
||||
message: errorMessage(error),
|
||||
agentRunTopic: config.agentRunTopic,
|
||||
hwlabTopic: config.hwlabTopic,
|
||||
valuesPrinted: false
|
||||
});
|
||||
});
|
||||
return {
|
||||
started: true,
|
||||
...config,
|
||||
ready,
|
||||
liveReady: ready,
|
||||
async stop() {
|
||||
stopped = true;
|
||||
subscribers.clear();
|
||||
await Promise.allSettled([bridgeConsumer?.stop?.(), fanoutConsumer?.stop?.()]);
|
||||
await Promise.allSettled([bridgeConsumer?.disconnect?.(), fanoutConsumer?.disconnect?.(), producer?.disconnect?.()]);
|
||||
},
|
||||
async status() {
|
||||
if (startupError) return { status: "blocked", capabilities: config.capabilities, errorCode: startupError.code ?? "hwlab_live_kafka_start_failed", message: errorMessage(startupError), valuesRedacted: true };
|
||||
return { status: startupReady ? "ready" : "initializing", capabilities: config.capabilities, subscriberCount: subscribers.size, consumerLag: liveKafkaLagSummary(lagByPartition), lossPossible: true, valuesRedacted: true };
|
||||
},
|
||||
startupStatus() {
|
||||
if (startupError) return { status: "blocked", capabilities: config.capabilities, errorCode: startupError.code ?? "hwlab_live_kafka_start_failed", message: errorMessage(startupError), valuesRedacted: true };
|
||||
return { status: startupReady ? "ready" : "initializing", capabilities: config.capabilities, valuesRedacted: true };
|
||||
},
|
||||
subscribeProjectionCommits() { return () => {}; },
|
||||
subscribeLiveHwlabEvents(listener) {
|
||||
if (typeof listener !== "function") return () => {};
|
||||
subscribers.add(listener);
|
||||
return () => subscribers.delete(listener);
|
||||
},
|
||||
queryHwlabEventRetention(params = {}) {
|
||||
return queryKafkaEventStream({ ...params, env, stream: "hwlab", topic: config.hwlabTopic, kafkaFactory });
|
||||
},
|
||||
liveSubscriberCount() { return subscribers.size; },
|
||||
valuesPrinted: false
|
||||
};
|
||||
}
|
||||
|
||||
function kafkaPartitionLag(highWatermark, lastOffset) {
|
||||
try {
|
||||
const lag = BigInt(String(highWatermark ?? "0")) - (BigInt(String(lastOffset ?? "0")) + 1n);
|
||||
return Number(lag > 0n ? lag : 0n);
|
||||
} catch {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
function liveKafkaLagSummary(lagByPartition) {
|
||||
const partitions = [...lagByPartition.entries()].map(([partition, lag]) => ({ partition, lag })).sort((left, right) => left.partition - right.partition);
|
||||
const known = partitions.map((item) => item.lag).filter((lag) => Number.isFinite(lag));
|
||||
return {
|
||||
known: known.length > 0,
|
||||
total: known.length > 0 ? known.reduce((sum, lag) => sum + lag, 0) : null,
|
||||
partitions,
|
||||
valuesRedacted: true
|
||||
};
|
||||
}
|
||||
|
||||
function combineKafkaEventBridgeComponents(config, components) {
|
||||
const ready = Promise.all(components.map((component) => component.ready));
|
||||
const liveOwner = components.find((component) => typeof component.subscribeLiveHwlabEvents === "function");
|
||||
const componentStartup = () => components.map((component) => component.startupStatus?.() ?? { status: "initializing" });
|
||||
return {
|
||||
started: true,
|
||||
...config,
|
||||
ready,
|
||||
liveReady: liveOwner?.liveReady ?? liveOwner?.ready ?? ready,
|
||||
async stop() { await Promise.allSettled(components.map((component) => component.stop())); },
|
||||
async status() {
|
||||
const statuses = await Promise.all(components.map((component) => component.status()));
|
||||
const blocked = statuses.find((status) => status.status === "blocked");
|
||||
const initializing = statuses.some((status) => status.status !== "ready");
|
||||
return {
|
||||
status: blocked ? "blocked" : initializing ? "initializing" : "ready",
|
||||
capabilities: config.capabilities,
|
||||
components: statuses,
|
||||
errorCode: blocked?.errorCode ?? null,
|
||||
message: blocked?.message ?? null,
|
||||
valuesRedacted: true
|
||||
};
|
||||
},
|
||||
startupStatus() {
|
||||
const statuses = componentStartup();
|
||||
const blocked = statuses.find((status) => status.status === "blocked");
|
||||
return {
|
||||
status: blocked ? "blocked" : statuses.every((status) => status.status === "ready") ? "ready" : "initializing",
|
||||
capabilities: config.capabilities,
|
||||
components: statuses,
|
||||
errorCode: blocked?.errorCode ?? null,
|
||||
message: blocked?.message ?? null,
|
||||
valuesRedacted: true
|
||||
};
|
||||
},
|
||||
subscribeProjectionCommits(listener) {
|
||||
const stops = components.map((component) => component.subscribeProjectionCommits?.(listener)).filter((stop) => typeof stop === "function");
|
||||
return () => stops.forEach((stop) => stop());
|
||||
},
|
||||
subscribeLiveHwlabEvents(listener) {
|
||||
return liveOwner?.subscribeLiveHwlabEvents(listener) ?? (() => {});
|
||||
},
|
||||
queryHwlabEventRetention(params = {}) {
|
||||
if (typeof liveOwner?.queryHwlabEventRetention !== "function") throw contractError("hwlab_kafka_refresh_query_unconfigured", "Kafka refresh retention query is not configured.");
|
||||
return liveOwner.queryHwlabEventRetention(params);
|
||||
},
|
||||
liveSubscriberCount() {
|
||||
const owner = components.find((component) => typeof component.liveSubscriberCount === "function");
|
||||
return owner?.liveSubscriberCount() ?? 0;
|
||||
},
|
||||
valuesPrinted: false
|
||||
};
|
||||
}
|
||||
|
||||
export async function publishAgentRunKafkaMessageLive(kafkaMessage = {}, { producer, config, env = process.env, logger = console, otelSpanEmitter = emitCodeAgentOtelSpan } = {}) {
|
||||
const transport = kafkaTransport(kafkaMessage);
|
||||
const decoded = decodeCanonicalAgentRunKafkaMessage(kafkaMessage);
|
||||
if (!decoded.ok) {
|
||||
logWarn(logger, "hwlab-live-kafka-message-invalid", { topic: transport.sourceTopic, partition: transport.sourcePartition, offset: transport.sourceOffset, errorCode: decoded.error.code, valuesPrinted: false });
|
||||
return { handled: true, invalid: true, published: false, valuesPrinted: false };
|
||||
}
|
||||
if (!stringValue(decoded.event.traceId) || !stringValue(decoded.event.hwlabSessionId)) {
|
||||
return { handled: true, ignored: true, published: false, valuesPrinted: false };
|
||||
}
|
||||
const projected = projectAgentRunKafkaEventToHwlabEvent(decoded.event, {
|
||||
source: config.clientId,
|
||||
sourceTopic: transport.sourceTopic,
|
||||
sourcePartition: transport.sourcePartition,
|
||||
sourceOffset: transport.sourceOffset,
|
||||
sourceKey: transport.sourceKey,
|
||||
inputSha256: transport.inputSha256
|
||||
});
|
||||
const publishMetadata = await producer.send({ topic: config.hwlabTopic, messages: [buildHwlabKafkaProducerMessage(projected)] });
|
||||
const publishedRecord = Array.isArray(publishMetadata) ? publishMetadata[0] : null;
|
||||
emitLiveKafkaOtelSpan("hwlab.kafka.live.direct_publish", projected, {
|
||||
topic: config.hwlabTopic,
|
||||
partition: publishedRecord?.partition,
|
||||
offset: publishedRecord?.baseOffset ?? publishedRecord?.offset,
|
||||
sourceTopic: transport.sourceTopic,
|
||||
sourcePartition: transport.sourcePartition,
|
||||
sourceOffset: transport.sourceOffset
|
||||
}, { env, otelSpanEmitter });
|
||||
return { handled: true, published: true, envelope: projected, sessionId: projected.sessionId, traceId: projected.traceId, valuesPrinted: false };
|
||||
}
|
||||
|
||||
export function shouldEmitLiveKafkaOtelSpan(envelope = {}) {
|
||||
const event = objectValue(envelope.event);
|
||||
const context = objectValue(envelope.context);
|
||||
const eventType = firstText(context.agentRunEventType, event.eventType, event.type, envelope.eventType)?.toLowerCase();
|
||||
if (event.terminal === true || eventType === "terminal" || eventType === "terminal_status") return true;
|
||||
if (["assistant", "assistant_progress", "assistant_message", "tool", "tool_call"].includes(eventType)) return true;
|
||||
if (eventType !== "command_output" && event.type !== "output") return true;
|
||||
const sourceSeq = integerValue(context.sourceSeq ?? event.sourceSeq);
|
||||
if (Number.isInteger(sourceSeq) && sourceSeq > 0) return sourceSeq % LIVE_KAFKA_COMMAND_OUTPUT_OTEL_SAMPLE_MODULUS === 0;
|
||||
const identity = firstText(envelope.sourceEventId, envelope.eventId, envelope.traceId, envelope.hwlabSessionId, envelope.sessionId);
|
||||
if (!identity) return false;
|
||||
return createHash("sha256").update(identity).digest().readUInt32BE(0) % LIVE_KAFKA_COMMAND_OUTPUT_OTEL_SAMPLE_MODULUS === 0;
|
||||
}
|
||||
|
||||
export function liveKafkaOtelSpanAttributes(envelope = {}, transport = {}) {
|
||||
const event = objectValue(envelope.event);
|
||||
const context = objectValue(envelope.context);
|
||||
const sourceEvent = objectValue(envelope.sourceEvent);
|
||||
return {
|
||||
businessTraceId: firstText(envelope.traceId, event.traceId),
|
||||
hwlabSessionId: firstText(envelope.hwlabSessionId, envelope.sessionId, event.sessionId),
|
||||
runId: firstText(envelope.runId, context.runId, event.runId),
|
||||
commandId: firstText(envelope.commandId, context.commandId, event.commandId),
|
||||
topic: firstText(transport.topic, sourceEvent.topic),
|
||||
partition: integerValue(transport.partition ?? sourceEvent.partition),
|
||||
offset: firstText(transport.offset, sourceEvent.offset),
|
||||
sourceTopic: firstText(transport.sourceTopic),
|
||||
sourcePartition: integerValue(transport.sourcePartition),
|
||||
sourceOffset: firstText(transport.sourceOffset),
|
||||
eventType: firstText(context.agentRunEventType, event.eventType, event.type, envelope.eventType),
|
||||
terminal: event.terminal === true,
|
||||
valuesRedacted: true
|
||||
};
|
||||
}
|
||||
|
||||
export function emitLiveKafkaOtelSpan(name, envelope = {}, transport = {}, { env = process.env, otelSpanEmitter = emitCodeAgentOtelSpan } = {}) {
|
||||
if (typeof otelSpanEmitter !== "function" || !shouldEmitLiveKafkaOtelSpan(envelope)) return { emitted: false, valuesRedacted: true };
|
||||
const attributes = liveKafkaOtelSpanAttributes(envelope, transport);
|
||||
if (!attributes.businessTraceId) return { emitted: false, valuesRedacted: true };
|
||||
try {
|
||||
const pending = otelSpanEmitter(name, attributes.businessTraceId, env, { attributes });
|
||||
void Promise.resolve(pending).catch(() => undefined);
|
||||
} catch {
|
||||
// OTel is best effort and must never change the live Kafka delivery path.
|
||||
}
|
||||
return { emitted: true, valuesRedacted: true };
|
||||
}
|
||||
|
||||
export async function projectAgentRunKafkaMessage(kafkaMessage = {}, { runtimeStore, config, logger = console } = {}) {
|
||||
const transport = kafkaTransport(kafkaMessage);
|
||||
const decoded = decodeCanonicalAgentRunKafkaMessage(kafkaMessage);
|
||||
|
||||
@@ -20,9 +20,6 @@ import { codeAgentSessionLifecycleSummary } from "./code-agent-session-lifecycle
|
||||
import { messageAuthorityTextValue } from "./code-agent-agentrun-prompt.ts";
|
||||
import { promoteWorkbenchTurnAdmission as persistWorkbenchTurnAdmissionPromotion, writeWorkbenchSessionAdmissionFact } from "./workbench-projection-writer.ts";
|
||||
import { createWorkbenchReadModel } from "./workbench-read-model.ts";
|
||||
import {
|
||||
workbenchRealtimeCapabilities
|
||||
} from "./workbench-realtime-capabilities.ts";
|
||||
import { codeAgentOtelTraceFields, emitCodeAgentOtelSpan } from "./otel-trace.ts";
|
||||
import {
|
||||
firstHeaderValue,
|
||||
@@ -245,7 +242,6 @@ export async function handleCodeAgentChatHttp(request, response, options) {
|
||||
lastEventAt: submitted?.lastEventAt ?? admittedTiming.lastEventAt,
|
||||
finishedAt: submitted?.finishedAt ?? admittedTiming.finishedAt,
|
||||
durationMs: submitted?.durationMs ?? admittedTiming.durationMs,
|
||||
realtimeCapabilities: workbenchRealtimeCapabilities(options.env ?? process.env),
|
||||
admission: submitted?.agentRun ? {
|
||||
state: submitted.admissionState ?? "promoted",
|
||||
runId: submitted.agentRun.runId ?? null,
|
||||
@@ -1323,7 +1319,6 @@ async function promoteCodeAgentTurnAdmission({ payload = {}, params = {}, option
|
||||
throw error;
|
||||
}
|
||||
const lifecycle = codeAgentTurnLifecycleFields(traceId, params);
|
||||
const realtimeCapabilities = workbenchRealtimeCapabilities();
|
||||
const inputFact = buildCodeAgentSessionInputFact({
|
||||
params,
|
||||
options,
|
||||
|
||||
@@ -46,29 +46,10 @@ const TRANSACTIONAL_REALTIME_ENV = Object.freeze({
|
||||
HWLAB_WORKBENCH_PROJECTION_REALTIME_ENABLED: "true"
|
||||
});
|
||||
|
||||
const REFRESH_REALTIME_ENV = Object.freeze({
|
||||
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"
|
||||
});
|
||||
|
||||
function projectionRealtimeBridge(capabilities = {}) {
|
||||
function projectionRealtimeBridge() {
|
||||
return {
|
||||
started: true,
|
||||
capabilities: {
|
||||
directPublish: false,
|
||||
liveKafkaSse: false,
|
||||
kafkaRefreshReplay: false,
|
||||
transactionalProjector: false,
|
||||
projectionOutboxRelay: false,
|
||||
projectionRealtime: true,
|
||||
...capabilities
|
||||
},
|
||||
ready: Promise.resolve(),
|
||||
subscribeLiveHwlabEvents() { return () => {}; },
|
||||
subscribeProjectionCommits() { return () => {}; },
|
||||
async stop() {}
|
||||
};
|
||||
@@ -146,7 +127,7 @@ test("workbench realtime initial connection emits current snapshot without repla
|
||||
}
|
||||
});
|
||||
|
||||
test("product realtime uses projection outbox even when obsolete live capability values remain present", async () => {
|
||||
test("product realtime uses the projection outbox authority", async () => {
|
||||
const sessionId = "ses_composable_realtime";
|
||||
const traceId = "trc_composable_realtime";
|
||||
let projectionReads = 0;
|
||||
@@ -169,7 +150,7 @@ test("product realtime uses projection outbox even when obsolete live capability
|
||||
const server = createCloudApiServer({
|
||||
accessController: realtimeAccessController(),
|
||||
workbenchRuntime: runtime,
|
||||
kafkaEventBridge: projectionRealtimeBridge({ directPublish: true, liveKafkaSse: true }),
|
||||
kafkaEventBridge: projectionRealtimeBridge(),
|
||||
env: {
|
||||
HWLAB_KAFKA_DIRECT_PUBLISH_ENABLED: "true",
|
||||
HWLAB_WORKBENCH_LIVE_KAFKA_SSE_ENABLED: "true",
|
||||
@@ -827,57 +808,6 @@ test("workbench trace event page keeps per-trace terminal status after later ses
|
||||
}
|
||||
});
|
||||
|
||||
function refreshReplayBridge({ records = [], endOffset = "0", onQuery = null, beforeQueryResult = null, ready = Promise.resolve(), liveReady = Promise.resolve() } = {}) {
|
||||
let liveListener = null;
|
||||
return {
|
||||
capabilities: { directPublish: true, liveKafkaSse: true, kafkaRefreshReplay: true, transactionalProjector: false, projectionOutboxRelay: false, projectionRealtime: false },
|
||||
refreshReplay: {
|
||||
groupIdPrefix: "hwlab-v03-refresh-test",
|
||||
timeoutMs: 2500,
|
||||
scanLimit: 500,
|
||||
matchedEventLimit: 50,
|
||||
liveBufferLimit: 20
|
||||
},
|
||||
ready,
|
||||
liveReady,
|
||||
subscribeLiveHwlabEvents(listener) { liveListener = listener; return () => { if (liveListener === listener) liveListener = null; }; },
|
||||
async queryHwlabEventRetention(params) {
|
||||
onQuery?.(params);
|
||||
beforeQueryResult?.(liveListener);
|
||||
return {
|
||||
topic: "hwlab.event.v1",
|
||||
events: records,
|
||||
completionReason: "end-offset",
|
||||
reachedEndOffsets: true,
|
||||
endOffsetsAvailable: true,
|
||||
endOffsets: [{ partition: 0, startOffset: "0", endOffset }],
|
||||
completion: { reason: "end-offset", complete: true, barrierReached: true, retentionStartVerified: true }
|
||||
};
|
||||
},
|
||||
async stop() {}
|
||||
};
|
||||
}
|
||||
|
||||
function refreshRecord(offset, sessionId, traceId, type, extra = {}) {
|
||||
const sourceEventId = `evt_refresh_${offset}`;
|
||||
return {
|
||||
topic: "hwlab.event.v1",
|
||||
partition: 0,
|
||||
offset: String(offset),
|
||||
value: {
|
||||
schema: "hwlab.event.v1",
|
||||
eventType: "hwlab.trace.event.projected",
|
||||
eventId: `hwlab:${sourceEventId}`,
|
||||
sourceEventId,
|
||||
traceId,
|
||||
hwlabSessionId: sessionId,
|
||||
sessionId,
|
||||
event: { type, eventType: type === "result" ? "terminal" : type, traceId, sessionId, sourceEventId, status: "running", ...extra },
|
||||
valuesPrinted: false
|
||||
}
|
||||
};
|
||||
}
|
||||
|
||||
function realtimeAccessController({ actor = ACTOR, sessions = [] } = {}) {
|
||||
const bySessionId = new Map(sessions.map((session) => [session.id, session]));
|
||||
const byTraceId = new Map(sessions.filter((session) => session.lastTraceId).map((session) => [session.lastTraceId, session]));
|
||||
|
||||
@@ -14,16 +14,11 @@ import {
|
||||
} from "./server-http-utils.ts";
|
||||
import { createWorkbenchReadModel } from "./workbench-read-model.ts";
|
||||
import { createWorkbenchRuntimeClient } from "./workbench-runtime-client.ts";
|
||||
import { emitLiveKafkaOtelSpan } from "./kafka-event-bridge.ts";
|
||||
import { createWorkbenchKafkaRefreshHandoff, workbenchKafkaRefreshErrorPayload } from "./workbench-kafka-refresh-handoff.ts";
|
||||
import { buildWorkbenchSessionDetail, compactLaunchContext, includeMessagesForSessionDetail } from "./workbench-session-detail-response.ts";
|
||||
import { handleWorkbenchSyncHttp } from "./workbench-realtime-authority.ts";
|
||||
import { projectionOutboxRealtimeEvents } from "./workbench-projection-outbox-events.ts";
|
||||
import { durableTraceStatus, RUNNING_STATUSES, TERMINAL_STATUSES } from "./workbench-turn-projection.ts";
|
||||
import { emitCodeAgentOtelSpan, emitHttpServerRequestSpan } from "./otel-trace.ts";
|
||||
import {
|
||||
workbenchRealtimeCapabilities
|
||||
} from "./workbench-realtime-capabilities.ts";
|
||||
import * as workbenchFacts from "./server-workbench-facts.ts";
|
||||
|
||||
const DEFAULT_PAGE_LIMIT = 50;
|
||||
@@ -183,7 +178,6 @@ export async function drainWorkbenchRealtimeConnections(options = {}) {
|
||||
export async function handleWorkbenchRealtimeHttp(request, response, url, options = {}) {
|
||||
try {
|
||||
if (request.method !== "GET") return methodNotAllowed(response, "GET");
|
||||
const realtimeCapabilities = workbenchRealtimeCapabilities(options.env ?? process.env);
|
||||
const projectionOnly = url.pathname === "/v1/workbench/projection-events";
|
||||
if (projectionOnly) {
|
||||
sendJson(response, 410, workbenchError("workbench_projection_events_route_removed", "Use /v1/workbench/events for the single transactional projection authority."));
|
||||
@@ -449,405 +443,6 @@ export async function handleWorkbenchRealtimeHttp(request, response, url, option
|
||||
}
|
||||
}
|
||||
|
||||
async function handleLiveKafkaWorkbenchRealtimeHttp(request, response, url, options = {}) {
|
||||
const requestedSessionId = safeSessionId(url.searchParams.get("sessionId") ?? url.searchParams.get("includeSessionId"));
|
||||
const requestedTraceId = safeTraceId(url.searchParams.get("traceId"));
|
||||
const realtimeCapabilities = workbenchRealtimeCapabilities(options.env ?? process.env);
|
||||
const heartbeatMs = parsePositiveInteger(options.env?.HWLAB_WORKBENCH_SSE_HEARTBEAT_MS, DEFAULT_WORKBENCH_SSE_HEARTBEAT_MS);
|
||||
attachWorkbenchRealtimeOtelContext(request, {
|
||||
sessionId: requestedSessionId,
|
||||
traceId: requestedTraceId,
|
||||
heartbeatMs,
|
||||
realtimeSource: "live-kafka-sse",
|
||||
realtimeCapabilities
|
||||
});
|
||||
const perf = options.backendPerformance;
|
||||
const auth = perf ? await perf.measure("workbench_auth", () => authenticateWorkbenchRead(request, response, options)) : await authenticateWorkbenchRead(request, response, options);
|
||||
if (!auth) return;
|
||||
if (url.searchParams.has("projectId") || url.searchParams.has("workspaceId")) {
|
||||
sendJson(response, 400, workbenchError("workbench_authority_removed", "Workbench realtime is keyed by sessionId/traceId only."));
|
||||
return;
|
||||
}
|
||||
if (!requestedSessionId && !requestedTraceId) {
|
||||
sendJson(response, 400, workbenchError("workbench_realtime_scope_required", "Workbench realtime requires sessionId or traceId."));
|
||||
return;
|
||||
}
|
||||
const authorization = await authorizeLiveKafkaWorkbenchRealtimeScope(options, auth.actor, {
|
||||
sessionId: requestedSessionId,
|
||||
traceId: requestedTraceId
|
||||
});
|
||||
if (!authorization.ok) {
|
||||
sendJson(response, authorization.status, authorization.body);
|
||||
return;
|
||||
}
|
||||
const authorizedSessionId = authorization.sessionId;
|
||||
const bridge = options.kafkaEventBridge;
|
||||
if (bridge?.capabilities?.liveKafkaSse !== true || typeof bridge?.subscribeLiveHwlabEvents !== "function") {
|
||||
sendJson(response, 503, workbenchError("workbench_live_kafka_unconfigured", "Workbench live Kafka fanout is not configured."));
|
||||
return;
|
||||
}
|
||||
if (realtimeCapabilities.kafkaRefreshReplay) {
|
||||
await handleKafkaRefreshReplayWorkbenchRealtimeHttp(request, response, options, {
|
||||
bridge,
|
||||
heartbeatMs,
|
||||
realtimeCapabilities,
|
||||
requestedSessionId: authorizedSessionId,
|
||||
requestedTraceId
|
||||
});
|
||||
return;
|
||||
}
|
||||
|
||||
let closed = false;
|
||||
let realtimeCloseReason = "client_close";
|
||||
let realtimeCloseSignal = null;
|
||||
const realtimeStartedAtMs = Date.now();
|
||||
const cleanup = [];
|
||||
const isActive = () => !closed && !response.destroyed && !response.writableEnded;
|
||||
const closeConnection = () => {
|
||||
if (closed) return;
|
||||
closed = true;
|
||||
emitWorkbenchRealtimeClosedOtelSpan(request, options.env, {
|
||||
reason: realtimeCloseReason,
|
||||
signal: realtimeCloseSignal,
|
||||
startedAtMs: realtimeStartedAtMs,
|
||||
sessionId: authorizedSessionId,
|
||||
threadId: null,
|
||||
traceId: requestedTraceId,
|
||||
activeConnectionCount: activeWorkbenchRealtimeConnections.size
|
||||
});
|
||||
for (const item of cleanup.splice(0)) item();
|
||||
};
|
||||
|
||||
response.once("close", closeConnection);
|
||||
request.once?.("aborted", closeConnection);
|
||||
request.socket?.once?.("close", closeConnection);
|
||||
cleanup.push(() => request.off?.("aborted", closeConnection));
|
||||
cleanup.push(() => request.socket?.off?.("close", closeConnection));
|
||||
|
||||
const bufferedEnvelopes = [];
|
||||
let liveDeliveryStarted = false;
|
||||
let enqueueEvent = null;
|
||||
const unsubscribe = bridge.subscribeLiveHwlabEvents((envelope, transport) => {
|
||||
if (!liveKafkaEnvelopeMatches(envelope, authorizedSessionId, requestedTraceId)) return;
|
||||
if (!liveDeliveryStarted || typeof enqueueEvent !== "function") {
|
||||
bufferedEnvelopes.push({ envelope, transport });
|
||||
return;
|
||||
}
|
||||
void enqueueEvent("hwlab.event.v1", envelope, transport);
|
||||
});
|
||||
if (typeof unsubscribe === "function") cleanup.push(unsubscribe);
|
||||
|
||||
try {
|
||||
await (bridge.liveReady ?? bridge.ready);
|
||||
} catch (error) {
|
||||
closeConnection();
|
||||
throw error;
|
||||
}
|
||||
if (!isActive()) {
|
||||
closeConnection();
|
||||
return;
|
||||
}
|
||||
|
||||
response.writeHead(200, {
|
||||
"content-type": "text/event-stream; charset=utf-8",
|
||||
"cache-control": "no-store, no-transform",
|
||||
connection: "keep-alive",
|
||||
"x-accel-buffering": "no",
|
||||
"x-content-type-options": "nosniff"
|
||||
});
|
||||
if (typeof response.flushHeaders === "function") response.flushHeaders();
|
||||
|
||||
const writeEvent = async (name, payload, transport = null) => {
|
||||
if (!isActive()) return false;
|
||||
const block = `event: ${name}\ndata: ${JSON.stringify(payload)}\n\n`;
|
||||
const writable = response.write(block);
|
||||
if (!writable) await waitForResponseDrain(response, isActive);
|
||||
const active = isActive();
|
||||
if (active && name === "hwlab.event.v1") {
|
||||
emitLiveKafkaOtelSpan("hwlab.workbench.live_sse.business_event_write", payload, transport, {
|
||||
env: options.env,
|
||||
otelSpanEmitter: options.otelSpanEmitter ?? emitCodeAgentOtelSpan
|
||||
});
|
||||
}
|
||||
return active;
|
||||
};
|
||||
let writeChain = Promise.resolve(true);
|
||||
enqueueEvent = (name, payload, transport = null) => {
|
||||
writeChain = writeChain.then(() => writeEvent(name, payload, transport));
|
||||
return writeChain;
|
||||
};
|
||||
|
||||
const realtimeConnection = {
|
||||
close(fields = {}) {
|
||||
if (!isActive()) return false;
|
||||
realtimeCloseReason = textValue(fields.reason) || "server_shutdown";
|
||||
realtimeCloseSignal = textValue(fields.signal) || null;
|
||||
void enqueueEvent("workbench.server_draining", {
|
||||
type: "server.draining",
|
||||
status: "closing",
|
||||
reason: realtimeCloseReason,
|
||||
signal: realtimeCloseSignal,
|
||||
capabilities: realtimeCapabilities,
|
||||
lossPossible: true
|
||||
}).finally(() => {
|
||||
if (!response.writableEnded) response.end();
|
||||
});
|
||||
return true;
|
||||
},
|
||||
isActive
|
||||
};
|
||||
activeWorkbenchRealtimeConnections.add(realtimeConnection);
|
||||
cleanup.push(() => activeWorkbenchRealtimeConnections.delete(realtimeConnection));
|
||||
|
||||
await enqueueEvent("workbench.connected", {
|
||||
type: "connected",
|
||||
status: "connected",
|
||||
capabilities: realtimeCapabilities,
|
||||
realtimeSource: "hwlab.event.v1",
|
||||
deliverySemantics: "live-only",
|
||||
liveOnly: true,
|
||||
replay: false,
|
||||
replaySupported: false,
|
||||
lossPossible: true,
|
||||
filters: { sessionId: authorizedSessionId, traceId: requestedTraceId }
|
||||
});
|
||||
liveDeliveryStarted = true;
|
||||
const readyBufferedEnvelopes = bufferedEnvelopes.splice(0);
|
||||
for (const { envelope, transport } of readyBufferedEnvelopes) void enqueueEvent("hwlab.event.v1", envelope, transport);
|
||||
await writeChain;
|
||||
if (!isActive()) {
|
||||
closeConnection();
|
||||
return;
|
||||
}
|
||||
const heartbeatTimer = setInterval(() => {
|
||||
void enqueueEvent("workbench.heartbeat", {
|
||||
type: "heartbeat",
|
||||
status: "connected",
|
||||
realtimeSource: "hwlab.event.v1",
|
||||
liveOnly: true,
|
||||
replay: false,
|
||||
lossPossible: true,
|
||||
serverSentAt: new Date().toISOString(),
|
||||
valuesPrinted: false
|
||||
});
|
||||
}, heartbeatMs);
|
||||
heartbeatTimer.unref?.();
|
||||
cleanup.push(() => clearInterval(heartbeatTimer));
|
||||
emitWorkbenchRealtimeAcceptedOtelSpan(request, options.env);
|
||||
}
|
||||
|
||||
async function handleKafkaRefreshReplayWorkbenchRealtimeHttp(request, response, options, input) {
|
||||
const {
|
||||
bridge,
|
||||
heartbeatMs,
|
||||
realtimeCapabilities,
|
||||
requestedSessionId,
|
||||
requestedTraceId
|
||||
} = input;
|
||||
const refreshReplay = bridge?.refreshReplay;
|
||||
if (!refreshReplay || typeof bridge?.queryHwlabEventRetention !== "function") {
|
||||
sendJson(response, 503, workbenchError("workbench_kafka_refresh_unconfigured", "Workbench Kafka refresh replay is enabled without its retention query runtime."));
|
||||
return;
|
||||
}
|
||||
let closed = false;
|
||||
let realtimeCloseReason = "client_close";
|
||||
let realtimeCloseSignal = null;
|
||||
let failureClosing = false;
|
||||
const realtimeStartedAtMs = Date.now();
|
||||
const cleanup = [];
|
||||
const isActive = () => !closed && !response.destroyed && !response.writableEnded;
|
||||
const closeConnection = () => {
|
||||
if (closed) return;
|
||||
closed = true;
|
||||
emitWorkbenchRealtimeClosedOtelSpan(request, options.env, {
|
||||
reason: realtimeCloseReason,
|
||||
signal: realtimeCloseSignal,
|
||||
startedAtMs: realtimeStartedAtMs,
|
||||
sessionId: requestedSessionId,
|
||||
threadId: null,
|
||||
traceId: requestedTraceId,
|
||||
activeConnectionCount: activeWorkbenchRealtimeConnections.size
|
||||
});
|
||||
for (const item of cleanup.splice(0)) item();
|
||||
};
|
||||
|
||||
response.once("close", closeConnection);
|
||||
request.once?.("aborted", closeConnection);
|
||||
request.socket?.once?.("close", closeConnection);
|
||||
cleanup.push(() => request.off?.("aborted", closeConnection));
|
||||
cleanup.push(() => request.socket?.off?.("close", closeConnection));
|
||||
|
||||
try {
|
||||
await (bridge.liveReady ?? bridge.ready);
|
||||
} catch (error) {
|
||||
for (const item of cleanup.splice(0)) item();
|
||||
throw error;
|
||||
}
|
||||
if (!isActive()) {
|
||||
closeConnection();
|
||||
return;
|
||||
}
|
||||
|
||||
response.writeHead(200, {
|
||||
"content-type": "text/event-stream; charset=utf-8",
|
||||
"cache-control": "no-store, no-transform",
|
||||
connection: "keep-alive",
|
||||
"x-accel-buffering": "no",
|
||||
"x-content-type-options": "nosniff"
|
||||
});
|
||||
if (typeof response.flushHeaders === "function") response.flushHeaders();
|
||||
|
||||
const writeEvent = async (name, payload, transport = null) => {
|
||||
if (!isActive()) return false;
|
||||
const block = `event: ${name}\ndata: ${JSON.stringify(payload)}\n\n`;
|
||||
const writable = response.write(block);
|
||||
if (!writable) await waitForResponseDrain(response, isActive);
|
||||
const active = isActive();
|
||||
if (active && name === "hwlab.event.v1") {
|
||||
emitLiveKafkaOtelSpan("hwlab.workbench.live_sse.business_event_write", payload, transport, {
|
||||
env: options.env,
|
||||
otelSpanEmitter: options.otelSpanEmitter ?? emitCodeAgentOtelSpan
|
||||
});
|
||||
}
|
||||
return active;
|
||||
};
|
||||
let writeChain = Promise.resolve(true);
|
||||
const enqueueEvent = (name, payload, transport = null) => {
|
||||
writeChain = writeChain.then(() => writeEvent(name, payload, transport));
|
||||
return writeChain;
|
||||
};
|
||||
const closeWithRefreshFailure = async (error) => {
|
||||
if (failureClosing || !isActive()) return;
|
||||
failureClosing = true;
|
||||
await enqueueEvent("workbench.error", workbenchKafkaRefreshErrorPayload(error, {
|
||||
sessionId: requestedSessionId,
|
||||
traceId: requestedTraceId
|
||||
}));
|
||||
await writeChain;
|
||||
if (!response.writableEnded) response.end();
|
||||
};
|
||||
|
||||
const realtimeConnection = {
|
||||
close(fields = {}) {
|
||||
if (!isActive()) return false;
|
||||
realtimeCloseReason = textValue(fields.reason) || "server_shutdown";
|
||||
realtimeCloseSignal = textValue(fields.signal) || null;
|
||||
void enqueueEvent("workbench.server_draining", {
|
||||
type: "server.draining",
|
||||
status: "closing",
|
||||
reason: realtimeCloseReason,
|
||||
signal: realtimeCloseSignal,
|
||||
capabilities: realtimeCapabilities,
|
||||
lossPossible: false
|
||||
}).finally(() => {
|
||||
if (!response.writableEnded) response.end();
|
||||
});
|
||||
return true;
|
||||
},
|
||||
isActive
|
||||
};
|
||||
activeWorkbenchRealtimeConnections.add(realtimeConnection);
|
||||
cleanup.push(() => activeWorkbenchRealtimeConnections.delete(realtimeConnection));
|
||||
|
||||
const handoff = createWorkbenchKafkaRefreshHandoff({
|
||||
sessionId: requestedSessionId,
|
||||
traceId: requestedTraceId,
|
||||
liveBufferLimit: refreshReplay.liveBufferLimit,
|
||||
subscribeLive: (listener) => bridge.subscribeLiveHwlabEvents(listener),
|
||||
queryRetention: ({ signal }) => bridge.queryHwlabEventRetention({
|
||||
sessionId: requestedSessionId,
|
||||
traceId: requestedTraceId,
|
||||
limit: refreshReplay.matchedEventLimit,
|
||||
scanLimit: refreshReplay.scanLimit,
|
||||
timeoutMs: refreshReplay.timeoutMs,
|
||||
groupIdPrefix: refreshReplay.groupIdPrefix,
|
||||
fromBeginning: true,
|
||||
signal
|
||||
}),
|
||||
deliverEvent: (envelope, transport) => enqueueEvent("hwlab.event.v1", envelope, transport),
|
||||
deliverConnected: (summary) => enqueueEvent("workbench.connected", {
|
||||
type: "connected",
|
||||
status: "connected",
|
||||
capabilities: realtimeCapabilities,
|
||||
realtimeSource: "hwlab.event.v1",
|
||||
deliverySemantics: "kafka-retention-then-live",
|
||||
liveOnly: false,
|
||||
replay: true,
|
||||
replaySupported: true,
|
||||
lossPossible: false,
|
||||
filters: { sessionId: requestedSessionId, traceId: requestedTraceId },
|
||||
refreshReplay: summary
|
||||
}),
|
||||
onFailure: closeWithRefreshFailure
|
||||
});
|
||||
cleanup.push(() => handoff.stop("connection-closed"));
|
||||
|
||||
try {
|
||||
await handoff.start();
|
||||
} catch (error) {
|
||||
if (isActive()) await closeWithRefreshFailure(error);
|
||||
return;
|
||||
}
|
||||
if (!isActive()) return;
|
||||
|
||||
const heartbeatTimer = setInterval(() => {
|
||||
void enqueueEvent("workbench.heartbeat", {
|
||||
type: "heartbeat",
|
||||
status: "connected",
|
||||
realtimeSource: "hwlab.event.v1",
|
||||
deliverySemantics: "kafka-retention-then-live",
|
||||
liveOnly: false,
|
||||
replay: true,
|
||||
lossPossible: false,
|
||||
serverSentAt: new Date().toISOString(),
|
||||
valuesPrinted: false
|
||||
});
|
||||
}, heartbeatMs);
|
||||
heartbeatTimer.unref?.();
|
||||
cleanup.push(() => clearInterval(heartbeatTimer));
|
||||
emitWorkbenchRealtimeAcceptedOtelSpan(request, options.env);
|
||||
}
|
||||
|
||||
async function authorizeLiveKafkaWorkbenchRealtimeScope(options, actor, { sessionId, traceId }) {
|
||||
const access = options.accessController;
|
||||
if ((sessionId && typeof access?.getAgentSession !== "function") || (traceId && typeof access?.getAgentSessionByTraceId !== "function")) {
|
||||
return {
|
||||
ok: false,
|
||||
status: 503,
|
||||
body: workbenchError("workbench_realtime_authorization_unconfigured", "Workbench live realtime ownership authorization is not configured.")
|
||||
};
|
||||
}
|
||||
|
||||
const [sessionById, sessionByTrace] = await Promise.all([
|
||||
sessionId ? access.getAgentSession(sessionId) : null,
|
||||
traceId ? access.getAgentSessionByTraceId(traceId) : null
|
||||
]);
|
||||
const sessionByIdId = sessionById ? safeSessionId(sessionById.id ?? sessionById.sessionId) : null;
|
||||
const sessionByTraceId = sessionByTrace ? safeSessionId(sessionByTrace.id ?? sessionByTrace.sessionId) : null;
|
||||
const scopeMissing = (sessionId && sessionByIdId !== sessionId)
|
||||
|| (traceId && !sessionByTraceId)
|
||||
|| (sessionId && traceId && sessionByIdId !== sessionByTraceId);
|
||||
const resolvedSessions = [sessionById, sessionByTrace].filter(Boolean);
|
||||
const ownedByActor = actor?.role === "admin" || (Boolean(actor?.id)
|
||||
&& resolvedSessions.length > 0
|
||||
&& resolvedSessions.every((session) => textValue(session.ownerUserId) === actor.id));
|
||||
if (scopeMissing || !ownedByActor) {
|
||||
return {
|
||||
ok: false,
|
||||
status: 404,
|
||||
body: workbenchError("workbench_realtime_scope_not_found", "Workbench realtime scope is not visible to the current actor.")
|
||||
};
|
||||
}
|
||||
return { ok: true, sessionId: sessionByIdId ?? sessionByTraceId };
|
||||
}
|
||||
|
||||
function liveKafkaEnvelopeMatches(envelope, sessionId, traceId) {
|
||||
if (!envelope || envelope.schema !== "hwlab.event.v1") return false;
|
||||
if (sessionId && safeSessionId(envelope.hwlabSessionId ?? envelope.sessionId) !== sessionId) return false;
|
||||
if (traceId && safeTraceId(envelope.traceId) !== traceId) return false;
|
||||
return true;
|
||||
}
|
||||
|
||||
export function attachWorkbenchRealtimeOtelContext(request, fields = {}) {
|
||||
const context = request?.hwlabHttpRequestContext;
|
||||
if (!context) return;
|
||||
|
||||
@@ -372,10 +372,7 @@ export async function buildHealthPayload(options = {}) {
|
||||
async function kafkaProjectorHealth(projector, { liveProbe = false } = {}) {
|
||||
if (!projector?.started) return { started: false, reason: projector?.reason ?? "disabled", valuesRedacted: true };
|
||||
const identity = {
|
||||
capabilities: projector.capabilities ?? null,
|
||||
directPublishGroupId: projector.directPublishGroupId ?? null,
|
||||
projectorGroupId: projector.projectorGroupId ?? null,
|
||||
hwlabEventGroupId: projector.hwlabEventGroupId ?? null,
|
||||
agentRunTopic: projector.agentRunTopic,
|
||||
hwlabTopic: projector.hwlabTopic
|
||||
};
|
||||
|
||||
@@ -1,15 +0,0 @@
|
||||
// SPEC: PJ2026-0104010803 Workbench fixed transactional realtime authority.
|
||||
// Responsibility: expose the code-owned Workbench projector/outbox/SSE invariant.
|
||||
|
||||
export const WORKBENCH_REALTIME_CAPABILITIES = Object.freeze({
|
||||
directPublish: false,
|
||||
liveKafkaSse: false,
|
||||
kafkaRefreshReplay: false,
|
||||
transactionalProjector: true,
|
||||
projectionOutboxRelay: true,
|
||||
projectionRealtime: true
|
||||
});
|
||||
|
||||
export function workbenchRealtimeCapabilities() {
|
||||
return WORKBENCH_REALTIME_CAPABILITIES;
|
||||
}
|
||||
@@ -0,0 +1,57 @@
|
||||
import assert from "node:assert/strict";
|
||||
import { readFile } from "node:fs/promises";
|
||||
import test from "node:test";
|
||||
import { fileURLToPath } from "node:url";
|
||||
|
||||
const repositoryRoot = fileURLToPath(new URL("../../..", import.meta.url));
|
||||
const productionFiles = [
|
||||
"internal/cloud/kafka-event-bridge.ts",
|
||||
"internal/cloud/server-workbench-realtime-http.ts",
|
||||
"internal/cloud/server-code-agent-admission-http.ts",
|
||||
"internal/cloud/server.ts",
|
||||
"web/hwlab-cloud-web/src/config/runtime.ts",
|
||||
"web/hwlab-cloud-web/src/api/workbench-events.ts",
|
||||
"web/hwlab-cloud-web/src/utils/workbench-stream-transport.ts",
|
||||
"web/hwlab-cloud-web/src/utils/workbench-realtime-runtime.ts",
|
||||
"web/hwlab-cloud-web/src/stores/workbench.ts",
|
||||
"web/hwlab-cloud-web/src/stores/workbench-event-reducer.ts",
|
||||
"web/hwlab-cloud-web/src/stores/workbench-realtime-plan.ts"
|
||||
];
|
||||
const forbiddenArchitectureIdentifiers = [
|
||||
"directPublish",
|
||||
"liveKafkaSse",
|
||||
"kafkaRefreshReplay",
|
||||
"transactionalProjector",
|
||||
"projectionOutboxRelay",
|
||||
"projectionRealtime",
|
||||
"startLiveHwlabKafkaEventBridge",
|
||||
"handleLiveKafkaWorkbenchRealtimeHttp",
|
||||
"handleKafkaRefreshReplayWorkbenchRealtimeHttp",
|
||||
"subscribeLiveHwlabEvents",
|
||||
"queryHwlabEventRetention",
|
||||
"WorkbenchRealtimeCapabilities",
|
||||
"workbenchRealtimeCapabilities",
|
||||
"workbenchHistoryAuthorityPolicy",
|
||||
"workbenchRealtimeTransportEnabled",
|
||||
"workbenchRealtimeTraceIdForCapabilities",
|
||||
"workbenchProjectionEventStreamPath"
|
||||
];
|
||||
|
||||
test("Workbench production source has one fixed transactional projection architecture", async () => {
|
||||
const violations: string[] = [];
|
||||
for (const relativePath of productionFiles) {
|
||||
const source = await readFile(`${repositoryRoot}/${relativePath}`, "utf8");
|
||||
for (const identifier of forbiddenArchitectureIdentifiers) {
|
||||
if (source.includes(identifier)) violations.push(`${relativePath}: ${identifier}`);
|
||||
}
|
||||
}
|
||||
assert.deepEqual(violations, []);
|
||||
});
|
||||
|
||||
test("raw hwlab Kafka fixtures remain isolated to the explicit debug path", async () => {
|
||||
const fixtureSource = await readFile(`${repositoryRoot}/web/hwlab-cloud-web/src/stores/workbench-live-kafka-event.ts`, "utf8");
|
||||
const debugSource = await readFile(`${repositoryRoot}/web/hwlab-cloud-web/src/stores/workbench-isolated-kafka-debug.ts`, "utf8");
|
||||
assert.match(fixtureSource, /historical hwlab\.event\.debug\.v1 fixtures/u);
|
||||
assert.match(fixtureSource, /never imported by the product Workbench SSE path/u);
|
||||
assert.match(debugSource, /workbench-live-kafka-event/u);
|
||||
});
|
||||
@@ -6,9 +6,9 @@ import test from "node:test";
|
||||
|
||||
import type { ChatMessage } from "../src/types/index.ts";
|
||||
import { workbenchSessionDetailPathForTest, workbenchSessionMessagesPathForTest } from "../src/api/workbench.ts";
|
||||
import { workbenchEventStreamPath, workbenchProjectionEventStreamPath } from "../src/api/workbench-events.ts";
|
||||
import { workbenchEventStreamPath } from "../src/api/workbench-events.ts";
|
||||
import type { WorkbenchStreamTransportRecovery } from "../src/utils/workbench-realtime-runtime.ts";
|
||||
import { workbenchRealtimeTraceIdForCapabilities, workbenchRealtimeTransportEnabled } from "../src/utils/workbench-stream-transport.ts";
|
||||
import { workbenchRealtimeTraceId } from "../src/utils/workbench-stream-transport.ts";
|
||||
import { workbenchRuntimePolicy } from "../src/config/workbench-runtime-policy.ts";
|
||||
import { AsyncQueue, work } from "../src/utils/scheduler/async-queue.ts";
|
||||
import { createCoalescedEventQueue } from "../src/utils/scheduler/coalesced-event-queue.ts";
|
||||
@@ -20,8 +20,7 @@ import { checkWorkbenchHealth, createWorkbenchHealthProbeCache } from "../src/ut
|
||||
import { messageDiagnosticView } from "../src/utils/workbench-error-runtime.ts";
|
||||
import { projectRejectedWorkbenchAdmission } from "../src/stores/workbench-admission-failure.ts";
|
||||
import { WORKBENCH_TIMELINE_OPENCODE_PARITY, buildWorkbenchTimelineRows, normalizeWorkbenchTimelineMessages, workbenchTimelineSignature } from "../src/stores/workbench-timeline-model.ts";
|
||||
import { reduceWorkbenchRealtimeEvent, workbenchRealtimeEventIsBusinessActivity } from "../src/stores/workbench-event-reducer.ts";
|
||||
import { reduceWorkbenchLiveKafkaMessageState, workbenchAgentMessageIdForTrace } from "../src/stores/workbench-live-kafka-event.ts";
|
||||
import { reduceWorkbenchRealtimeEvent } from "../src/stores/workbench-event-reducer.ts";
|
||||
import { planWorkbenchRealtimeApply, planWorkbenchRealtimeRecovery } from "../src/stores/workbench-realtime-plan.ts";
|
||||
import { WORKBENCH_REALTIME_AUTHORITY_VERSION, workbenchRealtimePrimaryAuthorityDecision } from "../src/stores/workbench-realtime-authority.ts";
|
||||
import { cleanupWorkbenchServerStateDroppedSessions, cleanupWorkbenchServerStateSessions, createWorkbenchServerState, reduceWorkbenchServerState } from "../src/stores/workbench-server-state.ts";
|
||||
@@ -83,24 +82,16 @@ test("Workbench API uses metadata-only session detail and bounded messages paths
|
||||
assert.equal(workbenchSessionMessagesPathForTest("ses_metadata", { limit: 9 }), "/v1/workbench/sessions/ses_metadata/messages?limit=9");
|
||||
});
|
||||
|
||||
test("only the projection authority sends the durable afterSeq cursor", () => {
|
||||
assert.equal(workbenchEventStreamPath({ realtimeCapabilities: { liveKafkaSse: true, kafkaRefreshReplay: false, projectionRealtime: false }, sessionId: "ses_live", traceId: "trc_live", afterSeq: 42 }), "/v1/workbench/events?sessionId=ses_live&traceId=trc_live");
|
||||
assert.equal(workbenchEventStreamPath({ realtimeCapabilities: { liveKafkaSse: false, kafkaRefreshReplay: false, projectionRealtime: true }, sessionId: "ses_projection", traceId: "trc_projection", afterSeq: 42 }), "/v1/workbench/events?sessionId=ses_projection&traceId=trc_projection&afterSeq=42");
|
||||
assert.equal(workbenchEventStreamPath({ realtimeCapabilities: { liveKafkaSse: true, kafkaRefreshReplay: true, projectionRealtime: true }, sessionId: "ses_both", traceId: "trc_both", afterSeq: 42 }), "/v1/workbench/events?sessionId=ses_both&traceId=trc_both&afterSeq=42");
|
||||
assert.equal(workbenchProjectionEventStreamPath({ sessionId: "ses_both", traceId: "trc_both", afterSeq: 42 }), "/v1/workbench/projection-events?sessionId=ses_both&traceId=trc_both&afterSeq=42");
|
||||
assert.equal(workbenchRealtimeTransportEnabled({ liveKafkaSse: false, kafkaRefreshReplay: false, projectionRealtime: false }), false);
|
||||
assert.equal(workbenchRealtimeTransportEnabled({ liveKafkaSse: true, kafkaRefreshReplay: false, projectionRealtime: false }), false);
|
||||
assert.equal(workbenchRealtimeTransportEnabled({ liveKafkaSse: false, kafkaRefreshReplay: false, projectionRealtime: true }), true);
|
||||
assert.equal(workbenchRealtimeTransportEnabled({ liveKafkaSse: true, kafkaRefreshReplay: true, projectionRealtime: true }), true);
|
||||
test("the product EventSource always sends the durable outbox cursor", () => {
|
||||
assert.equal(workbenchEventStreamPath({ sessionId: "ses_projection", traceId: "trc_projection", afterSeq: 42 }), "/v1/workbench/events?sessionId=ses_projection&traceId=trc_projection&afterSeq=42");
|
||||
});
|
||||
|
||||
test("obsolete live flags cannot erase the active projection trace scope", () => {
|
||||
const capabilities = { liveKafkaSse: true, kafkaRefreshReplay: true, projectionRealtime: false };
|
||||
const beforeSubmit = workbenchRealtimeScopeKey("ses_live", workbenchRealtimeTraceIdForCapabilities(capabilities, null, null));
|
||||
const afterSubmit = workbenchRealtimeScopeKey("ses_live", workbenchRealtimeTraceIdForCapabilities(capabilities, "trc_current_request", "trc_message"));
|
||||
test("the active projection trace scope follows the current request", () => {
|
||||
const beforeSubmit = workbenchRealtimeScopeKey("ses_projection", workbenchRealtimeTraceId(null, null));
|
||||
const afterSubmit = workbenchRealtimeScopeKey("ses_projection", workbenchRealtimeTraceId("trc_current_request", "trc_message"));
|
||||
assert.notEqual(afterSubmit, beforeSubmit);
|
||||
assert.equal(afterSubmit, workbenchRealtimeScopeKey("ses_live", "trc_current_request"));
|
||||
assert.equal(workbenchRealtimeTraceIdForCapabilities({ liveKafkaSse: false, kafkaRefreshReplay: false, projectionRealtime: true }, "trc_current_request", "trc_message"), "trc_current_request");
|
||||
assert.equal(afterSubmit, workbenchRealtimeScopeKey("ses_projection", "trc_current_request"));
|
||||
assert.equal(workbenchRealtimeTraceId(null, "trc_message"), "trc_message");
|
||||
});
|
||||
|
||||
test("Error runtime owns Workbench message diagnostic view model", () => {
|
||||
@@ -533,53 +524,6 @@ test("realtime event reducer classifies SSE payloads before store side effects",
|
||||
assert.equal(error.diagnostic.traceId, "trc_2");
|
||||
});
|
||||
|
||||
test("live transport connected and heartbeat frames do not extend business activity timeouts", () => {
|
||||
assert.equal(workbenchRealtimeEventIsBusinessActivity({ type: "connected" }, "workbench.connected", true), false);
|
||||
assert.equal(workbenchRealtimeEventIsBusinessActivity({ type: "heartbeat" }, "workbench.heartbeat", true), false);
|
||||
assert.equal(workbenchRealtimeEventIsBusinessActivity({ schema: "hwlab.event.v1", event: { type: "assistant" } }, "hwlab.event.v1", true), true);
|
||||
assert.equal(workbenchRealtimeEventIsBusinessActivity({ type: "trace.event" }, "workbench.trace.event", false), true);
|
||||
});
|
||||
|
||||
test("live hwlab envelope projects assistant, tool/output, and terminal without replay or finalizer", () => {
|
||||
const envelope = (event: Record<string, unknown>) => ({
|
||||
schema: "hwlab.event.v1",
|
||||
eventType: "hwlab.trace.event.projected",
|
||||
eventId: `hwlab:${String(event.sourceEventId ?? event.type)}`,
|
||||
hwlabSessionId: "ses_live_web",
|
||||
sessionId: "ses_live_web",
|
||||
traceId: "trc_live_web",
|
||||
runId: "run_live_web",
|
||||
commandId: "cmd_live_web",
|
||||
context: { runId: "run_live_web", commandId: "cmd_live_web", valuesRedacted: true },
|
||||
event
|
||||
});
|
||||
const assistantEnvelope = envelope({ type: "assistant", eventType: "assistant", sourceEventId: "evt_assistant", traceId: "trc_live_web", sessionId: "ses_live_web", status: "running", assistantText: "running increment", terminal: false });
|
||||
const assistant = reduceWorkbenchRealtimeEvent(assistantEnvelope, "hwlab.event.v1");
|
||||
assert.equal(assistant.action.type, "trace.event");
|
||||
let visible = reduceWorkbenchLiveKafkaMessageState({ text: "", status: "running", terminal: false }, assistantEnvelope.event);
|
||||
assert.deepEqual(visible, { text: "running increment", status: "running", terminal: false });
|
||||
|
||||
const toolEnvelope = envelope({ type: "tool", eventType: "tool", sourceEventId: "evt_tool", traceId: "trc_live_web", sessionId: "ses_live_web", status: "completed", toolName: "commandExecution", outputSummary: "ok", terminal: false });
|
||||
assert.equal(reduceWorkbenchRealtimeEvent(toolEnvelope, "hwlab.event.v1").action.type, "trace.event");
|
||||
visible = reduceWorkbenchLiveKafkaMessageState(visible, toolEnvelope.event);
|
||||
assert.deepEqual(visible, { text: "running increment", status: "running", terminal: false });
|
||||
|
||||
const outputEnvelope = envelope({ type: "output", eventType: "status", sourceEventId: "evt_output", traceId: "trc_live_web", sessionId: "ses_live_web", status: "running", message: "stdout increment", terminal: false });
|
||||
visible = reduceWorkbenchLiveKafkaMessageState(visible, outputEnvelope.event);
|
||||
assert.deepEqual(visible, { text: "running increment", status: "running", terminal: false });
|
||||
|
||||
const terminalEnvelope = envelope({ type: "result", eventType: "terminal", sourceEventId: "evt_terminal", traceId: "trc_live_web", sessionId: "ses_live_web", status: "completed", terminal: true });
|
||||
visible = reduceWorkbenchLiveKafkaMessageState(visible, terminalEnvelope.event);
|
||||
assert.deepEqual(visible, { text: "running increment", status: "completed", terminal: true });
|
||||
assert.equal(workbenchAgentMessageIdForTrace("trc_live_web"), "msg_live_web_agent");
|
||||
|
||||
const failed = reduceWorkbenchLiveKafkaMessageState(
|
||||
{ text: "partial response", status: "running", terminal: false },
|
||||
{ type: "result", eventType: "terminal", status: "failed", message: "provider stream disconnected", terminal: true }
|
||||
);
|
||||
assert.deepEqual(failed, { text: "provider stream disconnected", status: "failed", terminal: true });
|
||||
});
|
||||
|
||||
test("realtime apply planner turns reducer actions into store steps", () => {
|
||||
const reduced = reduceWorkbenchRealtimeEvent(realtimeEvent({ type: "trace.event", traceId: "trc_1", event: { traceId: "trc_1", label: "delta" }, snapshot: { traceId: "trc_1", status: "running" }, entity: { family: "traceEvents", id: "trc_1:1", version: 1, projectionRevision: "prj_1" } }), "workbench.trace.event");
|
||||
const tracePlan = planWorkbenchRealtimeApply(reduced.action);
|
||||
|
||||
@@ -2,7 +2,6 @@ import assert from "node:assert/strict";
|
||||
import { test } from "bun:test";
|
||||
|
||||
import { connectWorkbenchEvents, realtimeCoalesceKey, workbenchEventStreamPath, type WorkbenchRealtimeEvent, type WorkbenchSseIngressFrame } from "./workbench-events";
|
||||
import { createCoalescedEventQueue } from "../utils/scheduler/coalesced-event-queue";
|
||||
|
||||
test("turn snapshot coalescing keys by trace instead of per sequence", () => {
|
||||
const first: WorkbenchRealtimeEvent = { type: "turn.snapshot", turn: { sessionId: "ses_queue", traceId: "trc_queue", status: "running" }, cursor: { traceSeq: 10 } };
|
||||
@@ -21,7 +20,6 @@ test("trace events keep sequence-specific coalescing keys", () => {
|
||||
test("product EventSource resumes the transactional projection authority with one durable outbox cursor", () => {
|
||||
assert.equal(
|
||||
workbenchEventStreamPath({
|
||||
realtimeCapabilities: { liveKafkaSse: true, kafkaRefreshReplay: false, projectionRealtime: true },
|
||||
sessionId: "ses_projection_authority",
|
||||
traceId: "trc_projection_authority",
|
||||
afterSeq: 42
|
||||
@@ -30,29 +28,7 @@ test("product EventSource resumes the transactional projection authority with on
|
||||
);
|
||||
});
|
||||
|
||||
test("live Kafka envelopes from one trace retain every assistant, tool, output, and terminal event", () => {
|
||||
const delivered: WorkbenchRealtimeEvent[] = [];
|
||||
const queue = createCoalescedEventQueue<WorkbenchRealtimeEvent>({
|
||||
keyOf: (event) => realtimeCoalesceKey(event, "hwlab.event.v1"),
|
||||
onFlush: (events) => delivered.push(...events)
|
||||
});
|
||||
for (const sourceEventId of ["evt_assistant", "evt_tool", "evt_output", "evt_terminal"]) {
|
||||
queue.push({
|
||||
schema: "hwlab.event.v1",
|
||||
eventId: `hwlab:${sourceEventId}`,
|
||||
sourceEventId,
|
||||
sessionId: "ses_live_queue",
|
||||
traceId: "trc_live_queue",
|
||||
event: { sourceEventId, traceId: "trc_live_queue", sessionId: "ses_live_queue" }
|
||||
});
|
||||
}
|
||||
|
||||
queue.flush();
|
||||
|
||||
assert.deepEqual(delivered.map((event) => event.sourceEventId), ["evt_assistant", "evt_tool", "evt_output", "evt_terminal"]);
|
||||
});
|
||||
|
||||
test("product EventSource tees every raw HWLAB frame once before decode and reducer delivery", async () => {
|
||||
test("product EventSource tees projection frames once before reducer delivery", async () => {
|
||||
const originalEventSource = globalThis.EventSource;
|
||||
const ingress: WorkbenchSseIngressFrame[] = [];
|
||||
const delivered: WorkbenchRealtimeEvent[] = [];
|
||||
@@ -60,7 +36,6 @@ test("product EventSource tees every raw HWLAB frame once before decode and redu
|
||||
globalThis.EventSource = FakeProductEventSource as unknown as typeof EventSource;
|
||||
try {
|
||||
const stream = connectWorkbenchEvents({
|
||||
realtimeCapabilities: { liveKafkaSse: true, kafkaRefreshReplay: true, projectionRealtime: false },
|
||||
sessionId: "ses_raw_tee",
|
||||
flushYieldMs: 1,
|
||||
onIngress: (frame) => ingress.push(frame),
|
||||
@@ -69,17 +44,17 @@ test("product EventSource tees every raw HWLAB frame once before decode and redu
|
||||
assert.ok(stream);
|
||||
const source = FakeProductEventSource.instances[0];
|
||||
assert.ok(source);
|
||||
source.emitRaw("hwlab.event.v1", '{"schema":"hwlab.event.v1","eventId":"evt_raw_tee"}');
|
||||
source.emitRaw("hwlab.event.v1", "{invalid-json");
|
||||
source.emitRaw("workbench.trace.event", '{"type":"trace.event","traceId":"trc_raw_tee","event":{"traceId":"trc_raw_tee","projectedSeq":1}}');
|
||||
source.emitRaw("workbench.trace.event", "{invalid-json");
|
||||
await new Promise((resolve) => setTimeout(resolve, 10));
|
||||
|
||||
assert.equal(FakeProductEventSource.instances.length, 1);
|
||||
assert.deepEqual(ingress.map((frame) => [frame.eventName, frame.decodeStatus]), [
|
||||
["hwlab.event.v1", "accepted"],
|
||||
["hwlab.event.v1", "invalid-json"]
|
||||
["workbench.trace.event", "accepted"],
|
||||
["workbench.trace.event", "invalid-json"]
|
||||
]);
|
||||
assert.equal(delivered.length, 1);
|
||||
assert.equal(delivered[0]?.eventId, "evt_raw_tee");
|
||||
assert.equal(delivered[0]?.traceId, "trc_raw_tee");
|
||||
stream.close();
|
||||
} finally {
|
||||
globalThis.EventSource = originalEventSource;
|
||||
|
||||
@@ -6,7 +6,6 @@ import type { ApiResult, ChatMessage, ProjectionDiagnostic, TraceEvent } from "@
|
||||
import { createCoalescedEventQueue } from "@/utils/scheduler/coalesced-event-queue";
|
||||
import { composeWorkbenchScopedKey, firstScopePart } from "@/utils/workbench-key";
|
||||
import { recordWorkbenchRuntimeDiagnostic, recordWorkbenchSseLifecycle } from "@/utils/workbench-performance";
|
||||
import type { WorkbenchRealtimeCapabilities } from "@/config/runtime";
|
||||
import { decodeWorkbenchRealtimeEventFrame } from "./workbench-realtime-codec";
|
||||
|
||||
export interface WorkbenchRealtimeTraceSnapshot {
|
||||
@@ -29,11 +28,6 @@ export interface WorkbenchRealtimeEvent {
|
||||
eventType?: string | null;
|
||||
eventId?: string | null;
|
||||
sourceEventId?: string | null;
|
||||
capabilities?: WorkbenchRealtimeCapabilities | null;
|
||||
liveOnly?: boolean | null;
|
||||
replay?: boolean | null;
|
||||
deliverySemantics?: string | null;
|
||||
lossPossible?: boolean | null;
|
||||
type?: string;
|
||||
contractVersion?: string | null;
|
||||
realtimeAuthority?: string | null;
|
||||
@@ -76,7 +70,6 @@ export interface WorkbenchRealtimeEvent {
|
||||
}
|
||||
|
||||
export interface WorkbenchEventStreamOptions {
|
||||
realtimeCapabilities: WorkbenchRealtimeCapabilities;
|
||||
sessionId?: string | null;
|
||||
traceId?: string | null;
|
||||
afterSeq?: number | null;
|
||||
@@ -105,7 +98,6 @@ interface QueuedRealtimeEvent {
|
||||
}
|
||||
|
||||
const WORKBENCH_EVENT_NAMES = [
|
||||
"hwlab.event.v1",
|
||||
"workbench.connected",
|
||||
"workbench.trace.snapshot",
|
||||
"workbench.trace.event",
|
||||
@@ -183,20 +175,12 @@ export function connectWorkbenchEvents(options: WorkbenchEventStreamOptions): Wo
|
||||
};
|
||||
}
|
||||
|
||||
export function workbenchEventStreamPath(options: Pick<WorkbenchEventStreamOptions, "realtimeCapabilities" | "sessionId" | "traceId" | "afterSeq">): string {
|
||||
const params = new URLSearchParams();
|
||||
appendParam(params, "sessionId", options.sessionId);
|
||||
appendParam(params, "traceId", options.traceId);
|
||||
if (options.realtimeCapabilities.projectionRealtime) appendNumberParam(params, "afterSeq", options.afterSeq);
|
||||
return `/v1/workbench/events?${params.toString()}`;
|
||||
}
|
||||
|
||||
export function workbenchProjectionEventStreamPath(options: Pick<WorkbenchEventStreamOptions, "sessionId" | "traceId" | "afterSeq">): string {
|
||||
export function workbenchEventStreamPath(options: Pick<WorkbenchEventStreamOptions, "sessionId" | "traceId" | "afterSeq">): string {
|
||||
const params = new URLSearchParams();
|
||||
appendParam(params, "sessionId", options.sessionId);
|
||||
appendParam(params, "traceId", options.traceId);
|
||||
appendNumberParam(params, "afterSeq", options.afterSeq);
|
||||
return `/v1/workbench/projection-events?${params.toString()}`;
|
||||
return `/v1/workbench/events?${params.toString()}`;
|
||||
}
|
||||
|
||||
function appendParam(params: URLSearchParams, key: string, value: string | null | undefined): void {
|
||||
@@ -228,10 +212,6 @@ function scheduleRealtimeFlushYield(flush: () => void, yieldMs: number | null |
|
||||
}
|
||||
|
||||
export function realtimeCoalesceKey(event: WorkbenchRealtimeEvent, eventName: string): string | null {
|
||||
if (eventName === "hwlab.event.v1" || event.schema === "hwlab.event.v1") {
|
||||
const eventId = firstScopePart(event.eventId, event.sourceEventId, event.event?.sourceEventId);
|
||||
return eventId ? composeWorkbenchScopedKey("workbench.sse.live-kafka-event", eventId) : null;
|
||||
}
|
||||
const entity = event.entity;
|
||||
const entityFamily = firstScopePart(entity?.family);
|
||||
const entityId = firstScopePart(entity?.id);
|
||||
|
||||
@@ -10,12 +10,6 @@ export interface OpenCodeFrameConfig {
|
||||
url: string;
|
||||
}
|
||||
|
||||
export interface WorkbenchRealtimeCapabilities {
|
||||
liveKafkaSse: boolean;
|
||||
kafkaRefreshReplay: boolean;
|
||||
projectionRealtime: boolean;
|
||||
}
|
||||
|
||||
export interface WorkbenchDebugCapabilities {
|
||||
isolatedKafka: boolean;
|
||||
rawHwlabEventWindow: WorkbenchRawHwlabEventWindowCapability;
|
||||
@@ -97,10 +91,6 @@ export function workbenchTraceTimelinePolicy(): WorkbenchTraceTimelinePolicy {
|
||||
};
|
||||
}
|
||||
|
||||
export function workbenchRealtimeCapabilities(): WorkbenchRealtimeCapabilities {
|
||||
return { liveKafkaSse: false, kafkaRefreshReplay: false, projectionRealtime: true };
|
||||
}
|
||||
|
||||
export function workbenchDebugCapabilities(): WorkbenchDebugCapabilities {
|
||||
const features = window.HWLAB_CLOUD_WEB_CONFIG?.workbench?.debugCapabilities;
|
||||
const rawWindow = features?.rawHwlabEventWindow;
|
||||
|
||||
@@ -48,15 +48,7 @@ export function reduceWorkbenchRealtimeEvent(event: WorkbenchRealtimeEvent, even
|
||||
};
|
||||
}
|
||||
|
||||
export function workbenchRealtimeEventIsBusinessActivity(event: WorkbenchRealtimeEvent, eventName: string, liveKafkaSse: boolean): boolean {
|
||||
if (!liveKafkaSse) return true;
|
||||
return eventName === "hwlab.event.v1" || event.schema === "hwlab.event.v1";
|
||||
}
|
||||
|
||||
function reduceRealtimeAction(event: WorkbenchRealtimeEvent, eventName: string): WorkbenchRealtimeAction {
|
||||
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) {
|
||||
|
||||
@@ -5,9 +5,7 @@ import type { WorkbenchKafkaSseDebugEvent } from "@/api/workbench-debug";
|
||||
import type { WorkbenchRealtimeEvent } from "@/api/workbench-events";
|
||||
import type { ChatMessage } from "@/types";
|
||||
import { firstNonEmptyString } from "@/utils";
|
||||
import { reduceWorkbenchRealtimeEvent } from "./workbench-event-reducer";
|
||||
import { projectWorkbenchLiveKafkaMessage } from "./workbench-live-kafka-event";
|
||||
import { planWorkbenchRealtimeApply } from "./workbench-realtime-plan";
|
||||
|
||||
export interface WorkbenchIsolatedKafkaDebugLog {
|
||||
id: string;
|
||||
@@ -48,22 +46,17 @@ export function applyWorkbenchIsolatedKafkaDebugEvent(
|
||||
if (!traceId) return rejected(base, debugEvent, null, "debug-envelope-trace-missing");
|
||||
if (expected && traceId !== expected) return rejected(base, debugEvent, traceId, "debug-envelope-trace-mismatch");
|
||||
|
||||
const event = { ...debugEvent, schema: "hwlab.event.v1" } as WorkbenchRealtimeEvent;
|
||||
const reduced = reduceWorkbenchRealtimeEvent(event, "hwlab.event.v1");
|
||||
const plan = planWorkbenchRealtimeApply(reduced.action);
|
||||
const traceStep = plan.steps.find((step) => step.type === "apply-trace-event");
|
||||
if (!traceStep || traceStep.type !== "apply-trace-event" || !traceStep.event) {
|
||||
return withLog(base, event, traceId, reduced.action.type, plan.steps.map((step) => step.type), false, reduced.action.type === "ignore" ? reduced.action.reason : "debug-envelope-not-a-trace-event");
|
||||
}
|
||||
const traceEvent = debugEvent.event;
|
||||
if (!traceEvent) return withLog(base, debugEvent, traceId, "ignore", [], false, "debug-envelope-not-a-trace-event");
|
||||
const message = projectWorkbenchLiveKafkaMessage({
|
||||
previous: state.message,
|
||||
traceId,
|
||||
sessionId: firstNonEmptyString(event.hwlabSessionId, event.sessionId, traceStep.event.sessionId, state.message?.sessionId) ?? "ses_workbench_isolated_debug",
|
||||
event: traceStep.event,
|
||||
sessionId: firstNonEmptyString(debugEvent.hwlabSessionId, debugEvent.sessionId, traceEvent.sessionId, state.message?.sessionId) ?? "ses_workbench_isolated_debug",
|
||||
event: traceEvent,
|
||||
receivedAt: new Date().toISOString(),
|
||||
title: "Code Agent · 隔离调试"
|
||||
});
|
||||
return withLog({ ...base, message, appliedCount: state.appliedCount + 1, error: null }, event, traceId, reduced.action.type, plan.steps.map((step) => step.type), true, null);
|
||||
return withLog({ ...base, message, appliedCount: state.appliedCount + 1, error: null }, debugEvent, traceId, "debug.trace.event", ["debug-project-message"], true, null);
|
||||
}
|
||||
|
||||
export function workbenchCurrentDebugTraceId(messages: ChatMessage[], sessionLastTraceId?: string | null): string | null {
|
||||
|
||||
@@ -1,28 +0,0 @@
|
||||
import assert from "node:assert/strict";
|
||||
import { test } from "bun:test";
|
||||
|
||||
import { workbenchHistoryAuthorityPolicy } from "./workbench-kafka-refresh-policy";
|
||||
|
||||
test("every capability combination keeps automatic HTTP terminal and history readers disabled", () => {
|
||||
for (const liveKafkaSse of [false, true]) {
|
||||
for (const kafkaRefreshReplay of [false, true]) {
|
||||
for (const projectionRealtime of [false, true]) {
|
||||
const policy = workbenchHistoryAuthorityPolicy({ liveKafkaSse, kafkaRefreshReplay, projectionRealtime });
|
||||
const calls = { fetchSessionMessagesPage: 0, fetchTurn: 0, fetchTraceEvents: 0 };
|
||||
if (policy.sessionMessagesHydrate) calls.fetchSessionMessagesPage += 1;
|
||||
if (policy.turnStatusHydrate) calls.fetchTurn += 1;
|
||||
if (policy.traceEventsHydrate) calls.fetchTraceEvents += 1;
|
||||
|
||||
assert.deepEqual(policy, {
|
||||
kafkaRetention: false,
|
||||
sessionMetadataRead: true,
|
||||
sessionMessagesHydrate: false,
|
||||
turnStatusHydrate: false,
|
||||
traceEventsHydrate: false,
|
||||
syncReplay: projectionRealtime
|
||||
});
|
||||
assert.deepEqual(calls, { fetchSessionMessagesPage: 0, fetchTurn: 0, fetchTraceEvents: 0 });
|
||||
}
|
||||
}
|
||||
}
|
||||
});
|
||||
@@ -1,21 +0,0 @@
|
||||
import type { WorkbenchRealtimeCapabilities } from "@/config/runtime";
|
||||
|
||||
export interface WorkbenchHistoryAuthorityPolicy {
|
||||
kafkaRetention: boolean;
|
||||
sessionMetadataRead: true;
|
||||
sessionMessagesHydrate: false;
|
||||
turnStatusHydrate: false;
|
||||
traceEventsHydrate: false;
|
||||
syncReplay: boolean;
|
||||
}
|
||||
|
||||
export function workbenchHistoryAuthorityPolicy(capabilities: WorkbenchRealtimeCapabilities): WorkbenchHistoryAuthorityPolicy {
|
||||
return {
|
||||
kafkaRetention: false,
|
||||
sessionMetadataRead: true,
|
||||
sessionMessagesHydrate: false,
|
||||
turnStatusHydrate: false,
|
||||
traceEventsHydrate: false,
|
||||
syncReplay: capabilities.projectionRealtime
|
||||
};
|
||||
}
|
||||
@@ -1,5 +1,5 @@
|
||||
// SPEC: PJ2026-0104010803 Workbench live Kafka SSE.
|
||||
// Responsibility: project one transparent hwlab.event.v1 business event into visible turn state without replay/finalizer/polling.
|
||||
// SPEC: PJ2026-0104010803 isolated Kafka debug fixture.
|
||||
// Responsibility: project historical hwlab.event.debug.v1 fixtures for the admin-only isolated debugger; never imported by the product Workbench SSE path.
|
||||
|
||||
import { mergeRunnerTrace } from "../composables/workbench-trace-snapshot";
|
||||
import type { ChatMessage, TraceEvent, WorkbenchTurnTimingProjection } from "../types";
|
||||
|
||||
@@ -4,13 +4,13 @@
|
||||
import { computed, nextTick, ref } from "vue";
|
||||
import { defineStore } from "pinia";
|
||||
import { api } from "@/api";
|
||||
import { workbenchDebugCapabilities, workbenchRealtimeCapabilities } from "@/config/runtime";
|
||||
import { workbenchDebugCapabilities } from "@/config/runtime";
|
||||
import { workbenchRuntimePolicy } from "@/config/workbench-runtime-policy";
|
||||
import { createKeyedSingleflight } from "@/utils/scheduler/keyed-singleflight";
|
||||
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 { createWorkbenchStreamTransportRuntime, workbenchRealtimeTraceId, 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";
|
||||
@@ -21,10 +21,8 @@ import { RECENT_DRAFTS_STORAGE_KEY, appendSessionPage, defaultProviderProfileOpt
|
||||
import { initialWorkbenchSessionIdFromLocation } from "./workbench-projection";
|
||||
import { cleanupWorkbenchServerStateSessions, selectActiveMessages, selectActiveSession, selectSessionList, selectSessionStatusAuthority, selectTraceAuthorityById, selectTurnStatusAuthority, type WorkbenchServerAction } from "./workbench-server-state";
|
||||
import { cleanupDroppedWorkbenchSessionCaches, trimWorkbenchSessionCache } from "./workbench-session-cache";
|
||||
import { reduceWorkbenchRealtimeEvent, workbenchRealtimeEventIsBusinessActivity, type WorkbenchRealtimeAction } from "./workbench-event-reducer";
|
||||
import { projectWorkbenchLiveKafkaMessage, projectWorkbenchLiveKafkaUserMessage, workbenchAgentMessageIdForTrace, workbenchLiveKafkaEnvelope, workbenchLiveKafkaProjectionTarget, workbenchUserMessageIdForTrace } from "./workbench-live-kafka-event";
|
||||
import { reduceWorkbenchRealtimeEvent, type WorkbenchRealtimeAction } from "./workbench-event-reducer";
|
||||
import { projectRejectedWorkbenchAdmission } from "./workbench-admission-failure";
|
||||
import { workbenchHistoryAuthorityPolicy } from "./workbench-kafka-refresh-policy";
|
||||
import { messageHasSealedTerminalResult, messageIsSealedTerminal, traceAuthorityIsSealed } from "./workbench-terminal-authority";
|
||||
import { boundedProjectionMessageLimit, mergeBoundedProjectionMessages, selectProjectionMessageWindow, traceProjectionIsTerminalSealed } from "./workbench-message-projection-budget";
|
||||
import {
|
||||
@@ -108,10 +106,17 @@ interface RealtimeTurnProjectionItem {
|
||||
terminalTurn: boolean;
|
||||
}
|
||||
|
||||
function workbenchMessageIdForTrace(traceId: string, role: "user" | "agent"): string {
|
||||
const suffix = firstNonEmptyString(traceId)
|
||||
?.replace(/^trc_/u, "")
|
||||
.replace(/[^A-Za-z0-9_.:-]/gu, "_")
|
||||
.slice(0, 48);
|
||||
if (!suffix) throw new Error("traceId must produce a stable Workbench message identity");
|
||||
return `msg_${suffix}_${role}`;
|
||||
}
|
||||
|
||||
export const useWorkbenchStore = defineStore("workbench", () => {
|
||||
const runtimePolicy = workbenchRuntimePolicy();
|
||||
const realtimeCapabilities = workbenchRealtimeCapabilities();
|
||||
const historyAuthorityPolicy = workbenchHistoryAuthorityPolicy(realtimeCapabilities);
|
||||
const debugCapabilities = workbenchDebugCapabilities();
|
||||
const workbenchColadaReducer = useWorkbenchColadaReducer();
|
||||
const workbenchColadaQueries = useWorkbenchColadaQueries();
|
||||
@@ -138,7 +143,6 @@ export const useWorkbenchStore = defineStore("workbench", () => {
|
||||
const sessionListLoadingMore = ref(false);
|
||||
const sessionListLoadMoreError = ref<string | null>(null);
|
||||
const chatPending = ref(false);
|
||||
const liveRealtimeReadySessionId = ref<string | null>(null);
|
||||
const rawHwlabIngress = ref(createRawHwlabIngressState());
|
||||
const error = ref<string | null>(null);
|
||||
const activityRef = ref({ lastActivityAt: Date.now(), lastActivityIso: new Date().toISOString(), waitingFor: "idle", lastEventLabel: null as string | null });
|
||||
@@ -169,7 +173,7 @@ export const useWorkbenchStore = defineStore("workbench", () => {
|
||||
const sessionListLoadedCount = computed(() => sessionTabs.value.length);
|
||||
const sessionListLoading = computed(() => shouldShowSessionListLoading({ loading: loading.value, sessionsReady: sessionsReady.value }));
|
||||
const sessionDetailLoading = computed(() => Boolean(sessionDetailLoadingId.value || switchingSessionId.value || (loading.value && messages.value.length === 0)));
|
||||
const composer = computed(() => resolveComposerState({ messages: messages.value, sessions: sessions.value, activeSessionId: activeSessionId.value, chatPending: chatPending.value, realtimeReady: !realtimeCapabilities.liveKafkaSse || liveRealtimeReadySessionId.value === activeSessionId.value, currentRequest: currentRequest.value, turnStatusAuthority: turnStatusAuthority.value }));
|
||||
const composer = computed(() => resolveComposerState({ messages: messages.value, sessions: sessions.value, activeSessionId: activeSessionId.value, chatPending: chatPending.value, realtimeReady: true, currentRequest: currentRequest.value, turnStatusAuthority: turnStatusAuthority.value }));
|
||||
|
||||
function recordActivity(label = "user-activity"): void {
|
||||
const now = Date.now();
|
||||
@@ -376,7 +380,6 @@ export const useWorkbenchStore = defineStore("workbench", () => {
|
||||
rememberSessionList(next);
|
||||
applySessionPagination(response.data);
|
||||
sessionsReady.value = true;
|
||||
if (historyAuthorityPolicy.turnStatusHydrate) await refreshSessionStatusAuthority(next);
|
||||
return;
|
||||
}
|
||||
sessionsReady.value = sessions.value.length > 0;
|
||||
@@ -384,7 +387,6 @@ export const useWorkbenchStore = defineStore("workbench", () => {
|
||||
}
|
||||
|
||||
function scheduleSessionListRefresh(includeSessionId: string | null | undefined = activeSessionId.value, delayMs = runtimePolicy.sessionListRealtimeRefreshDelayMs): void {
|
||||
if (realtimeCapabilities.liveKafkaSse) return;
|
||||
void includeSessionId;
|
||||
if (typeof window === "undefined") {
|
||||
void workbenchColadaQueries.invalidateSessionList();
|
||||
@@ -412,7 +414,6 @@ export const useWorkbenchStore = defineStore("workbench", () => {
|
||||
rememberSessionList(next);
|
||||
applySessionPagination(response.data);
|
||||
sessionsReady.value = true;
|
||||
if (historyAuthorityPolicy.turnStatusHydrate) await refreshSessionStatusAuthority(next);
|
||||
}
|
||||
|
||||
function currentSessionListLimit(): number {
|
||||
@@ -763,14 +764,9 @@ export const useWorkbenchStore = defineStore("workbench", () => {
|
||||
const submitEntry = steerMode ? "steer" : "existing";
|
||||
const traceId = steerMode && composer.value.targetTraceId ? composer.value.targetTraceId : nextProtocolId("trc");
|
||||
const steerTraceId = steerMode ? nextProtocolId("trc_steer") : null;
|
||||
const userMessageId = workbenchUserMessageIdForTrace(steerTraceId ?? traceId);
|
||||
const agentMessageId = workbenchAgentMessageIdForTrace(traceId);
|
||||
const userMessageId = workbenchMessageIdForTrace(steerTraceId ?? traceId, "user");
|
||||
const agentMessageId = workbenchMessageIdForTrace(traceId, "agent");
|
||||
const sessionId = composer.value.sessionId;
|
||||
if (realtimeCapabilities.liveKafkaSse && liveRealtimeReadySessionId.value !== sessionId) {
|
||||
error.value = "realtime_connecting";
|
||||
restartRealtime("submit-readiness");
|
||||
return false;
|
||||
}
|
||||
const threadId = composer.value.threadId;
|
||||
const providerThreadId = providerThreadIdForRequest(threadId);
|
||||
const submittedAt = new Date().toISOString();
|
||||
@@ -1128,10 +1124,8 @@ export const useWorkbenchStore = defineStore("workbench", () => {
|
||||
void reason;
|
||||
const realtimeScopeKey = workbenchRealtimeScopeKey(sessionId, traceId);
|
||||
const scopeChanged = realtimeTransport.currentKey() !== realtimeScopeKey;
|
||||
if (realtimeCapabilities.liveKafkaSse && scopeChanged) liveRealtimeReadySessionId.value = null;
|
||||
if (debugCapabilities.rawHwlabEventWindow.enabled && scopeChanged) rawHwlabIngress.value = createRawHwlabIngressState(realtimeScopeKey);
|
||||
realtimeTransport.restart({
|
||||
realtimeCapabilities,
|
||||
sessionId,
|
||||
traceId,
|
||||
forceReconnect,
|
||||
@@ -1142,31 +1136,8 @@ export const useWorkbenchStore = defineStore("workbench", () => {
|
||||
onIngress: debugCapabilities.rawHwlabEventWindow.enabled
|
||||
? (frame) => { rawHwlabIngress.value = appendRawHwlabIngressFrame(rawHwlabIngress.value, frame, debugCapabilities.rawHwlabEventWindow); }
|
||||
: undefined,
|
||||
onOpen: () => undefined,
|
||||
onError: () => {
|
||||
if (realtimeCapabilities.liveKafkaSse) liveRealtimeReadySessionId.value = null;
|
||||
},
|
||||
onState: (state) => {
|
||||
if (realtimeCapabilities.liveKafkaSse && ["connecting", "error", "closed", "blocked"].includes(state.phase)) liveRealtimeReadySessionId.value = null;
|
||||
},
|
||||
onRecovery: (recovery) => handleRealtimeRecovery(recovery),
|
||||
onEvent: (event, eventName) => {
|
||||
if (realtimeCapabilities.liveKafkaSse && eventName === "workbench.connected") {
|
||||
const filters = recordValue(event.filters);
|
||||
const connectedSessionId = normalizeWorkbenchSessionId(filters?.sessionId);
|
||||
const refreshReplay = realtimeCapabilities.kafkaRefreshReplay;
|
||||
const deliveryValid = refreshReplay
|
||||
? event.deliverySemantics === "kafka-retention-then-live" && event.liveOnly === false && event.replay === true && event.lossPossible === false
|
||||
: event.deliverySemantics === "live-only" && event.liveOnly === true && event.replay === false;
|
||||
const valid = deliveryValid
|
||||
&& event.capabilities?.liveKafkaSse === true
|
||||
&& event.capabilities?.kafkaRefreshReplay === refreshReplay
|
||||
&& connectedSessionId === sessionId;
|
||||
liveRealtimeReadySessionId.value = valid ? connectedSessionId : null;
|
||||
if (!valid) error.value = "workbench_live_realtime_contract_invalid";
|
||||
}
|
||||
applyRealtimeEvent(event, eventName);
|
||||
}
|
||||
onEvent: (event, eventName) => applyRealtimeEvent(event, eventName)
|
||||
});
|
||||
}
|
||||
|
||||
@@ -1193,13 +1164,11 @@ export const useWorkbenchStore = defineStore("workbench", () => {
|
||||
}
|
||||
|
||||
function stopRealtime(): void {
|
||||
liveRealtimeReadySessionId.value = null;
|
||||
realtimeTransport.stop();
|
||||
}
|
||||
|
||||
function realtimeTraceId(): string | null {
|
||||
return workbenchRealtimeTraceIdForCapabilities(
|
||||
realtimeCapabilities,
|
||||
return workbenchRealtimeTraceId(
|
||||
currentRequest.value?.traceId,
|
||||
activeTraceIdFromMessages(messages.value, turnStatusAuthority.value)
|
||||
);
|
||||
@@ -1217,7 +1186,7 @@ export const useWorkbenchStore = defineStore("workbench", () => {
|
||||
const snapshot = recordValue(event.snapshot);
|
||||
const traceEvent = recordValue(event.event);
|
||||
const sessionId = firstNonEmptyString(event.hwlabSessionId, event.sessionId, turn?.sessionId, snapshot?.sessionId, traceEvent?.sessionId);
|
||||
return event.schema === "hwlab.event.v1" ? normalizeWorkbenchSessionId(sessionId) : normalizeTraceAuthoritySessionId(sessionId);
|
||||
return normalizeTraceAuthoritySessionId(sessionId);
|
||||
}
|
||||
|
||||
function traceResultSessionId(result: AgentChatResultResponse | TraceSnapshot | Record<string, unknown> | null | undefined): string | null {
|
||||
@@ -1284,7 +1253,7 @@ export const useWorkbenchStore = defineStore("workbench", () => {
|
||||
|
||||
function applyRealtimeEvent(event: WorkbenchRealtimeEvent, eventName: string): void {
|
||||
const reduced = reduceWorkbenchRealtimeEvent(event, eventName);
|
||||
if (workbenchRealtimeEventIsBusinessActivity(event, eventName, realtimeCapabilities.liveKafkaSse)) recordActivity(reduced.activityLabel);
|
||||
recordActivity(reduced.activityLabel);
|
||||
applyWorkbenchRealtimeAction(reduced.action);
|
||||
}
|
||||
|
||||
@@ -1340,77 +1309,18 @@ export const useWorkbenchStore = defineStore("workbench", () => {
|
||||
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;
|
||||
const liveKafkaEnvelope = workbenchLiveKafkaEnvelope(realtimeEvent?.schema);
|
||||
const sessionId = realtimeEvent ? realtimeEventSessionId(realtimeEvent) : traceResultSessionId(snapshot ?? event ?? null);
|
||||
if (!shouldApplyActiveTraceAuthority(id, sessionId, liveKafkaEnvelope)) return;
|
||||
if (traceTerminalBodyIsVisible(id, sessionId) && !liveKafkaEnvelope) {
|
||||
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-event", valuesRedacted: true } });
|
||||
return;
|
||||
}
|
||||
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());
|
||||
return;
|
||||
}
|
||||
const eventSnapshot = snapshot ?? { traceId: id, status: event?.status, events };
|
||||
applyTraceSnapshot(id, realtimeSnapshotToTraceSnapshot(id, eventSnapshot, events));
|
||||
}
|
||||
|
||||
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;
|
||||
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 });
|
||||
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
|
||||
});
|
||||
reduceServerState({
|
||||
type: "message.upsert",
|
||||
sessionId: ownerSessionId,
|
||||
message
|
||||
});
|
||||
if (message.runnerTrace) rememberTraceAuthority(message.runnerTrace);
|
||||
markWorkbenchTraceProjected(traceId);
|
||||
const terminal = message.traceAutoLifecycle === "terminal";
|
||||
const status = firstNonEmptyString(message.status) ?? "running";
|
||||
const finalResponse = terminal && message.text ? { text: message.text, status, traceId } : undefined;
|
||||
rememberTurnStatus(traceId, {
|
||||
traceId,
|
||||
sessionId: ownerSessionId,
|
||||
status,
|
||||
running: !terminal,
|
||||
terminal,
|
||||
finalResponse,
|
||||
timing: message.timing,
|
||||
startedAt: message.timing?.startedAt,
|
||||
lastEventAt: message.timing?.lastEventAt,
|
||||
finishedAt: message.timing?.finishedAt,
|
||||
durationMs: message.timing?.durationMs,
|
||||
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 (terminal && currentRequest.value?.traceId === traceId) {
|
||||
chatPending.value = false;
|
||||
currentRequest.value = null;
|
||||
}
|
||||
}
|
||||
|
||||
function applyRealtimeTurnSnapshot(turn: Record<string, unknown>): void {
|
||||
const traceId = firstNonEmptyString(turn.traceId);
|
||||
if (!traceId) return;
|
||||
@@ -1513,7 +1423,6 @@ export const useWorkbenchStore = defineStore("workbench", () => {
|
||||
}
|
||||
|
||||
function publishWorkbenchProjectionSignal(sessionId: string | null | undefined, traceId: string | null | undefined, reason: string): void {
|
||||
if (!realtimeCapabilities.projectionRealtime) return;
|
||||
if (typeof window === "undefined") return;
|
||||
const id = normalizeWorkbenchSessionId(sessionId);
|
||||
if (!id) return;
|
||||
@@ -1523,7 +1432,6 @@ export const useWorkbenchStore = defineStore("workbench", () => {
|
||||
}
|
||||
|
||||
function handleWorkbenchProjectionSignal(value: unknown): void {
|
||||
if (!realtimeCapabilities.projectionRealtime) return;
|
||||
const record = recordValue(value);
|
||||
if (!record || record.sourceId === workbenchProjectionSignalSourceId) return;
|
||||
if (firstNonEmptyString(record.type) !== "session-projection") return;
|
||||
@@ -1619,8 +1527,7 @@ export const useWorkbenchStore = defineStore("workbench", () => {
|
||||
const code = firstNonEmptyString(errorRecord.code, "workbench_realtime_error") ?? "workbench_realtime_error";
|
||||
const phase = firstNonEmptyString(event.phase, errorRecord.phase, "realtime") ?? "realtime";
|
||||
const message = firstNonEmptyString(errorRecord.message, errorRecord.summary, "Workbench 实时状态更新异常,最新运行记录暂不可见。") ?? "Workbench 实时状态更新异常,最新运行记录暂不可见。";
|
||||
liveRealtimeReadySessionId.value = null;
|
||||
error.value = `${phase} · ${code} · ${message}`;
|
||||
error.value = `${phase} · ${code} · ${message}`;
|
||||
return;
|
||||
}
|
||||
const projection = normalizeProjectionDiagnostic(event) ?? projectionDiagnosticFromFailure({ code: firstNonEmptyString(errorRecord.code, "workbench_realtime_error") ?? "workbench_realtime_error", message: firstNonEmptyString(errorRecord.message, errorRecord.summary, "Workbench 实时状态更新异常,最新运行记录暂不可见。") ?? "Workbench 实时状态更新异常,最新运行记录暂不可见。", health: "degraded", diagnostic: normalizeErrorDiagnostic(errorRecord.diagnostic, recordValue(event.projection)?.diagnostic, recordValue(event)?.diagnostic), apiError: normalizeApiErrorRecord(errorRecord, "Workbench 实时状态更新异常,最新运行记录暂不可见。") });
|
||||
@@ -1793,10 +1700,6 @@ export const useWorkbenchStore = defineStore("workbench", () => {
|
||||
rememberSessionDetail({ ...session, messages: projectedMessages, messageCount: session.messageCount ?? projectedMessages.length });
|
||||
sessionsReady.value = true;
|
||||
currentRequest.value = null;
|
||||
if (historyAuthorityPolicy.kafkaRetention) {
|
||||
restartRealtime("apply-selected-session:kafka-retention");
|
||||
return;
|
||||
}
|
||||
void hydrateTurnStatusAuthority(messages.value);
|
||||
void hydrateTerminalMessageDiagnostics();
|
||||
void readTraceEventsForMessages(messages.value);
|
||||
@@ -1824,8 +1727,6 @@ export const useWorkbenchStore = defineStore("workbench", () => {
|
||||
const requestId = normalizeWorkbenchSessionRouteId(sessionId);
|
||||
if (!requestId) return null;
|
||||
const normalizedRequestId = normalizeWorkbenchSessionId(requestId);
|
||||
const messageLimit = sessionMessageProjectionWindowLimit();
|
||||
const eagerMessages = historyAuthorityPolicy.sessionMessagesHydrate && normalizedRequestId ? fetchSessionMessagesPage(normalizedRequestId, { limit: messageLimit, reason: "load-session:eager", force: true }) : null;
|
||||
const detail = await fetchSessionDetailPage(requestId, { reason: "load-session:detail", force: true });
|
||||
if (!detail.ok) return null;
|
||||
const detailSession = sessionFromWorkbenchSession(detail.data?.session, { includeMessages: false });
|
||||
@@ -1834,11 +1735,7 @@ export const useWorkbenchStore = defineStore("workbench", () => {
|
||||
const fallbackMessages = seed?.sessionId === id ? seed.messages ?? [] : [];
|
||||
const base = detailSession ? { ...detailSession, messages: fallbackMessages } : seed;
|
||||
if (!base) return null;
|
||||
if (!historyAuthorityPolicy.sessionMessagesHydrate) return { ...base, sessionId: id, messages: [], messageCount: base.messageCount ?? 0 };
|
||||
const messages = eagerMessages && id === normalizedRequestId ? await eagerMessages : await fetchSessionMessagesPage(id, { limit: messageLimit, reason: "load-session", force: true });
|
||||
const page = messages.ok ? messages.data : null;
|
||||
const pageMessages = Array.isArray(page?.messages) ? await sealRestoredActiveTurnMessages(page.messages.map((message) => normalizeChatMessage(message as ChatMessage))) : fallbackMessages;
|
||||
return { ...base, sessionId: id, messages: pageMessages, messageCount: page?.total ?? pageMessages?.length ?? base.messageCount };
|
||||
return { ...base, sessionId: id, messages: [], messageCount: base.messageCount ?? 0 };
|
||||
}
|
||||
|
||||
async function sealRestoredActiveTurnMessages(source: ChatMessage[]): Promise<ChatMessage[]> {
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
// SPEC: PJ2026-0106050514 Workbench实时运行面 draft-2026-06-30-p0-1297-spec-first.
|
||||
// Responsibility: Production Workbench realtime runtime entry point for SSE transport, event coalescing, and scoped cursor ownership.
|
||||
|
||||
export { createWorkbenchStreamTransportRuntime, WorkbenchStreamTransportRuntime, workbenchRealtimeTraceIdForCapabilities, type WorkbenchStreamTransportRecovery, type WorkbenchStreamTransportRestartInput, type WorkbenchStreamTransportRestartResult, type WorkbenchStreamTransportState, type WorkbenchRealtimeEvent } from "@/utils/workbench-stream-transport";
|
||||
export { createWorkbenchStreamTransportRuntime, WorkbenchStreamTransportRuntime, workbenchRealtimeTraceId, type WorkbenchStreamTransportRecovery, type WorkbenchStreamTransportRestartInput, type WorkbenchStreamTransportRestartResult, type WorkbenchStreamTransportState, type WorkbenchRealtimeEvent } from "@/utils/workbench-stream-transport";
|
||||
|
||||
@@ -8,13 +8,11 @@
|
||||
// - OpenCode stream.transport.ts:504-530 reduce output/trace.
|
||||
|
||||
import { connectWorkbenchEvents, type WorkbenchEventStream, type WorkbenchRealtimeEvent, type WorkbenchSseIngressFrame } from "@/api/workbench-events";
|
||||
import type { WorkbenchRealtimeCapabilities } from "@/config/runtime";
|
||||
import { workbenchRealtimeScopeKey } from "@/utils/workbench-key";
|
||||
|
||||
export type { WorkbenchRealtimeEvent };
|
||||
|
||||
export interface WorkbenchStreamTransportRestartInput {
|
||||
realtimeCapabilities: WorkbenchRealtimeCapabilities;
|
||||
sessionId?: string | null;
|
||||
traceId?: string | null;
|
||||
afterSeq?: number | null;
|
||||
@@ -78,8 +76,7 @@ interface WorkbenchTransportCursor {
|
||||
traceSeq: number | null;
|
||||
}
|
||||
|
||||
export function workbenchRealtimeTraceIdForCapabilities(
|
||||
_capabilities: WorkbenchRealtimeCapabilities,
|
||||
export function workbenchRealtimeTraceId(
|
||||
...candidates: Array<string | null | undefined>
|
||||
): string | null {
|
||||
for (const candidate of candidates) {
|
||||
@@ -89,10 +86,6 @@ export function workbenchRealtimeTraceIdForCapabilities(
|
||||
return null;
|
||||
}
|
||||
|
||||
export function workbenchRealtimeTransportEnabled(capabilities: WorkbenchRealtimeCapabilities): boolean {
|
||||
return capabilities.projectionRealtime;
|
||||
}
|
||||
|
||||
export class WorkbenchStreamTransportRuntime {
|
||||
private stream: WorkbenchEventStream | null = null;
|
||||
private key = "";
|
||||
@@ -107,11 +100,9 @@ export class WorkbenchStreamTransportRuntime {
|
||||
this.stop();
|
||||
this.key = key;
|
||||
this.armWait(key);
|
||||
if (!workbenchRealtimeTransportEnabled(input.realtimeCapabilities)) return { key, changed: true, streamStarted: false };
|
||||
if (!input.sessionId && !input.traceId) return { key, changed: true, streamStarted: false };
|
||||
this.emitState(input, "connecting", null);
|
||||
this.stream = connectWorkbenchEvents({
|
||||
realtimeCapabilities: input.realtimeCapabilities,
|
||||
sessionId: input.sessionId ?? null,
|
||||
traceId: input.traceId ?? null,
|
||||
afterSeq: input.afterSeq ?? this.cursorByKey.get(key)?.outboxSeq ?? null,
|
||||
@@ -202,7 +193,7 @@ export class WorkbenchStreamTransportRuntime {
|
||||
const sessionId = input.sessionId ?? null;
|
||||
const traceId = input.traceId ?? null;
|
||||
const cursor = this.currentCursor(key);
|
||||
const actions: WorkbenchStreamTransportRecoveryAction[] = input.realtimeCapabilities.projectionRealtime && (sessionId || traceId) ? ["events-reconnect"] : [];
|
||||
const actions: WorkbenchStreamTransportRecoveryAction[] = 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 });
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user