import assert from "node:assert/strict"; import { test } from "bun:test"; import { createWorkbenchKafkaRefreshHandoff, workbenchKafkaRefreshErrorPayload } from "./workbench-kafka-refresh-handoff.ts"; const SESSION_ID = "ses_refresh_handoff"; const TRACE_ID = "trc_refresh_handoff"; test("retention handoff replays once, deduplicates overlap and late pre-barrier live, and drains growth before connected", async () => { let publishLive: any = null; const output: string[] = []; let growDuringFlush = true; const replay = [record(0), record(1), record(2)]; const handoff = createWorkbenchKafkaRefreshHandoff({ sessionId: SESSION_ID, liveBufferLimit: 8, subscribeLive(listener: any) { publishLive = listener; return () => { publishLive = null; }; }, async queryRetention() { publishLive(record(1).value, transport(1)); publishLive(record(3).value, transport(3)); return queryResult(replay, [{ partition: 0, startOffset: "0", endOffset: "3" }]); }, async deliverEvent(envelope: any, kafkaTransport: any, delivery: string) { output.push(`${delivery}:${kafkaTransport.offset}:${envelope.eventId}`); if (delivery === "buffered-live" && kafkaTransport.offset === "3" && growDuringFlush) { growDuringFlush = false; publishLive(record(4).value, transport(4)); } return true; }, async deliverConnected(summary: any) { output.push(`connected:${summary.counts.deduplicated}`); return true; } }); const summary = await handoff.start(); assert.equal(handoff.phase(), "live"); assert.deepEqual(output, [ "replay:0:hwlab:evt_0_0", "replay:1:hwlab:evt_0_1", "replay:2:hwlab:evt_0_2", "buffered-live:3:hwlab:evt_0_3", "buffered-live:4:hwlab:evt_0_4", "connected:1" ]); assert.deepEqual(summary.counts, { matched: 3, filteredNonHwlab: 0, filteredMissingIdentity: 0, replayed: 3, buffered: 3, bufferedDelivered: 2, liveDelivered: 0, deduplicated: 1, resumeSkipped: 0 }); publishLive(record(2).value, transport(2)); publishLive(record(5).value, transport(5)); await waitFor(() => output.includes("live:5:hwlab:evt_0_5")); assert.equal(output.filter((item) => item.includes("hwlab:evt_0_2")).length, 1); assert.equal(handoff.summary().counts.deduplicated, 2); handoff.stop(); }); test("retention resume delivers only events after a found identity and marks a missing identity as window advanced", async () => { const foundOutput: string[] = []; const found = createWorkbenchKafkaRefreshHandoff({ sessionId: SESSION_ID, resumeAfterId: "evt_0_1", liveBufferLimit: 4, subscribeLive() { return () => {}; }, async queryRetention() { return queryResult([record(0), record(1), record(2)], [{ partition: 0, startOffset: "0", endOffset: "3" }]); }, async deliverEvent(envelope: any) { foundOutput.push(envelope.sourceEventId); return true; }, async deliverConnected() { return true; } }); const foundSummary = await found.start(); assert.deepEqual(foundOutput, ["evt_0_2"]); assert.deepEqual(foundSummary.resume, { requested: true, found: true, windowAdvanced: false }); assert.equal(foundSummary.counts.resumeSkipped, 2); const missingOutput: string[] = []; const missing = createWorkbenchKafkaRefreshHandoff({ sessionId: SESSION_ID, resumeAfterId: "evt_expired", liveBufferLimit: 4, subscribeLive() { return () => {}; }, async queryRetention() { return queryResult([record(4), record(5)], [{ partition: 0, startOffset: "4", endOffset: "6" }]); }, async deliverEvent(envelope: any) { missingOutput.push(envelope.sourceEventId); return true; }, async deliverConnected() { return true; } }); const missingSummary = await missing.start(); assert.deepEqual(missingOutput, ["evt_0_4", "evt_0_5"]); assert.deepEqual(missingSummary.resume, { requested: true, found: false, windowAdvanced: true }); }); test("retention emits handoff before buffered live flush and connected live", async () => { let publishLive: any = null; const phases: string[] = []; const handoff = createWorkbenchKafkaRefreshHandoff({ sessionId: SESSION_ID, liveBufferLimit: 4, subscribeLive(listener: any) { publishLive = listener; return () => {}; }, async queryRetention() { publishLive(record(2).value, transport(2)); return queryResult([record(0), record(1)], [{ partition: 0, startOffset: "0", endOffset: "2" }]); }, async deliverEvent(_envelope: any, transportValue: any, delivery: string) { phases.push(`${delivery}:${transportValue.offset}`); return true; }, async deliverHandoff(summary: any) { phases.push(summary.phase); return true; }, async deliverConnected(summary: any) { phases.push(summary.phase); return true; } }); await handoff.start(); assert.deepEqual(phases, ["replay:0", "replay:1", "handoff", "buffered-live:2", "live"]); }); test("unscoped observer filters non-deliverable history without weakening scoped handoff", async () => { let publishLive: any = null; const delivered: string[] = []; const legacy = { ...record(0), value: { schema: "platform-infra.kafka.shadow-produce.v1", eventId: "legacy-smoke" } }; const handoff = createWorkbenchKafkaRefreshHandoff({ allowUnscoped: true, liveBufferLimit: 4, subscribeLive(listener: any) { publishLive = listener; return () => {}; }, async queryRetention() { publishLive({ schema: "platform-infra.kafka.shadow-produce.v1", eventId: "legacy-live" }, transport(2)); return queryResult([legacy, { ...record(1), value: { schema: "hwlab.event.v1", eventType: "legacy-without-identity" } }, record(2)], [{ partition: 0, startOffset: "0", endOffset: "3" }]); }, async deliverEvent(envelope: any) { delivered.push(envelope.eventId); return true; }, async deliverConnected() { return true; } }); const summary = await handoff.start(); assert.deepEqual(delivered, ["hwlab:evt_0_2"]); assert.equal(summary.counts.matched, 3); assert.equal(summary.counts.filteredNonHwlab, 2); assert.equal(summary.counts.filteredMissingIdentity, 1); assert.equal(summary.counts.replayed, 1); }); test("retention handoff fails closed on a pre-barrier live gap", async () => { let publishLive: any = null; const handoff = createWorkbenchKafkaRefreshHandoff({ sessionId: SESSION_ID, liveBufferLimit: 4, subscribeLive(listener: any) { publishLive = listener; return () => {}; }, async queryRetention() { publishLive(record(1).value, transport(1)); return queryResult([record(0)], [{ partition: 0, startOffset: "0", endOffset: "3" }]); }, async deliverEvent() { return true; }, async deliverConnected() { return true; } }); await assert.rejects(handoff.start(), (error: any) => error?.code === "workbench_kafka_refresh_barrier_gap" && error?.phase === "flushing"); assert.equal(handoff.phase(), "failed"); }); test("connected write failure blocks a concurrently queued live business event", async () => { let publishLive: any = null; let connectedStarted = false; let releaseConnected: (() => void) | null = null; const connectedGate = new Promise((resolve) => { releaseConnected = resolve; }); const delivered: string[] = []; const handoff = createWorkbenchKafkaRefreshHandoff({ sessionId: SESSION_ID, liveBufferLimit: 4, subscribeLive(listener: any) { publishLive = listener; return () => { publishLive = null; }; }, async queryRetention() { return queryResult([], [{ partition: 0, startOffset: "0", endOffset: "1" }]); }, async deliverEvent(_envelope: any, kafkaTransport: any, delivery: string) { delivered.push(`${delivery}:${kafkaTransport.offset}`); return true; }, async deliverConnected() { connectedStarted = true; await connectedGate; return false; } }); const started = handoff.start(); await waitFor(() => connectedStarted); publishLive(record(1).value, transport(1)); releaseConnected?.(); await assert.rejects(started, (error: any) => error?.code === "workbench_kafka_refresh_sse_write_failed"); await Promise.resolve(); assert.equal(handoff.phase(), "failed"); assert.deepEqual(delivered, []); }); test("retention handoff fails closed on overflow, cross-partition scope, and missing stable identity", async () => { let overflowLive: any = null; const overflow = createWorkbenchKafkaRefreshHandoff({ sessionId: SESSION_ID, liveBufferLimit: 1, subscribeLive(listener: any) { overflowLive = listener; return () => {}; }, async queryRetention() { overflowLive(record(3).value, transport(3)); overflowLive(record(4).value, transport(4)); return queryResult([], [{ partition: 0, startOffset: "0", endOffset: "3" }]); }, async deliverEvent() { return true; }, async deliverConnected() { return true; } }); await assert.rejects(overflow.start(), (error: any) => error?.code === "workbench_kafka_refresh_live_buffer_overflow"); const multiPartition = createWorkbenchKafkaRefreshHandoff({ sessionId: SESSION_ID, liveBufferLimit: 4, subscribeLive() { return () => {}; }, async queryRetention() { return queryResult([record(0, 0), record(0, 1)], [{ partition: 0, startOffset: "0", endOffset: "1" }, { partition: 1, startOffset: "0", endOffset: "1" }]); }, async deliverEvent() { return true; }, async deliverConnected() { return true; } }); await assert.rejects(multiPartition.start(), (error: any) => error?.code === "workbench_kafka_refresh_multi_partition"); const missingIdentityRecord = record(0); delete missingIdentityRecord.value.eventId; delete missingIdentityRecord.value.sourceEventId; const missingIdentity = createWorkbenchKafkaRefreshHandoff({ sessionId: SESSION_ID, liveBufferLimit: 4, subscribeLive() { return () => {}; }, async queryRetention() { return queryResult([missingIdentityRecord], [{ partition: 0, startOffset: "0", endOffset: "1" }]); }, async deliverEvent() { return true; }, async deliverConnected() { return true; } }); await assert.rejects(missingIdentity.start(), (error: any) => error?.code === "workbench_kafka_refresh_stable_identity_missing"); }); test("multi-partition topic barrier succeeds when the authorized scope remains on one partition", async () => { const delivered: string[] = []; const handoff = createWorkbenchKafkaRefreshHandoff({ sessionId: SESSION_ID, liveBufferLimit: 4, subscribeLive() { return () => {}; }, async queryRetention() { return queryResult([record(4, 1)], [ { partition: 0, startOffset: "3", endOffset: "3" }, { partition: 1, startOffset: "4", endOffset: "5" } ]); }, async deliverEvent(_envelope: any, kafkaTransport: any) { delivered.push(`${kafkaTransport.partition}:${kafkaTransport.offset}`); return true; }, async deliverConnected() { return true; } }); const summary = await handoff.start(); assert.deepEqual(delivered, ["1:4"]); assert.deepEqual(summary.topicPartitions, [1]); assert.deepEqual(summary.barrier, [{ partition: 0, endOffset: "3" }, { partition: 1, endOffset: "5" }]); }); test("stable eventId and sourceEventId are independently conflict checked", async () => { const first = record(0); const conflicting = record(1); conflicting.value.eventId = first.value.eventId; const handoff = createWorkbenchKafkaRefreshHandoff({ sessionId: SESSION_ID, liveBufferLimit: 4, subscribeLive() { return () => {}; }, async queryRetention() { return queryResult([first, conflicting], [{ partition: 0, startOffset: "0", endOffset: "2" }]); }, async deliverEvent() { return true; }, async deliverConnected() { return true; } }); await assert.rejects(handoff.start(), (error: any) => error?.code === "workbench_kafka_refresh_stable_identity_conflict"); }); test("session-empty refresh is legal while trace-empty refresh is typed fail-closed", async () => { let sessionLive: any = null; const sessionOutput: string[] = []; const sessionHandoff = createWorkbenchKafkaRefreshHandoff({ sessionId: SESSION_ID, liveBufferLimit: 4, subscribeLive(listener: any) { sessionLive = listener; return () => {}; }, async queryRetention() { return queryResult([], [{ partition: 0, startOffset: "7", endOffset: "7" }, { partition: 1, startOffset: "10", endOffset: "10" }]); }, async deliverEvent(_envelope: any, kafkaTransport: any, delivery: string) { sessionOutput.push(`${delivery}:${kafkaTransport.offset}`); return true; }, async deliverConnected() { sessionOutput.push("connected"); return true; } }); const summary = await sessionHandoff.start(); assert.equal(summary.counts.matched, 0); assert.deepEqual(sessionOutput, ["connected"]); sessionLive(record(10, 1).value, transport(10, 1)); await waitFor(() => sessionOutput.includes("live:10")); const traceHandoff = createWorkbenchKafkaRefreshHandoff({ traceId: TRACE_ID, liveBufferLimit: 4, subscribeLive() { return () => {}; }, async queryRetention() { return queryResult([], [{ partition: 0, startOffset: "7", endOffset: "7" }]); }, async deliverEvent() { return true; }, async deliverConnected() { return true; } }); await assert.rejects(traceHandoff.start(), (error: any) => error?.code === "workbench_kafka_refresh_trace_empty"); }); test("post-handoff partition failure is reported through a typed session-scoped workbench error", async () => { let publishLive: any = null; let failure: any = null; const handoff = createWorkbenchKafkaRefreshHandoff({ sessionId: SESSION_ID, liveBufferLimit: 4, subscribeLive(listener: any) { publishLive = listener; return () => {}; }, async queryRetention() { return queryResult([record(0)], [{ partition: 0, startOffset: "0", endOffset: "1" }]); }, async deliverEvent() { return true; }, async deliverConnected() { return true; }, async onFailure(error: any) { failure = error; } }); await handoff.start(); publishLive(record(1, 1).value, transport(1, 1)); await waitFor(() => Boolean(failure)); const payload = workbenchKafkaRefreshErrorPayload(failure, { sessionId: SESSION_ID }); assert.equal(payload.sessionId, SESSION_ID); assert.equal(payload.traceId, null); assert.equal(payload.phase, "live"); assert.equal(payload.error.code, "workbench_kafka_refresh_multi_partition"); assert.deepEqual(payload.error.details, { expectedPartition: 0, actualPartition: 1 }); assert.equal(payload.fallback, false); }); test("typed refresh failure exposes only bounded safe scan diagnostics", async () => { const handoff = createWorkbenchKafkaRefreshHandoff({ traceId: TRACE_ID, liveBufferLimit: 4, subscribeLive() { return () => {}; }, async queryRetention() { return { topic: "hwlab.event.v1", events: [], completionReason: "timeout", completion: { reason: "timeout", complete: false, barrierReached: false, retentionStartVerified: true }, reachedEndOffsets: false, endOffsetsAvailable: true, endOffsets: [{ partition: 3, startOffset: "4", endOffset: "9" }], limit: 50, scanLimit: 500, timeoutMs: 2500, scannedCount: 17, matchedCount: 2 }; }, async deliverEvent() { return true; }, async deliverConnected() { return true; } }); let failure: any = null; await assert.rejects(handoff.start(), (error: any) => { failure = error; return error?.code === "workbench_kafka_refresh_timeout"; }); failure.details.untrustedPayload = "must-not-leak"; const payload = workbenchKafkaRefreshErrorPayload(failure, { traceId: TRACE_ID }); assert.deepEqual(payload.error.details, { completionReason: "timeout", matchedEventLimit: 50, scanLimit: 500, timeoutMs: 2500, scannedCount: 17, matchedCount: 2, partitions: [3] }); assert.equal(JSON.stringify(payload).includes("must-not-leak"), false); }); test("post-handoff topic drift is rejected for the lifetime of the retention proof", async () => { let publishLive: any = null; let failure: any = null; const handoff = createWorkbenchKafkaRefreshHandoff({ sessionId: SESSION_ID, liveBufferLimit: 4, subscribeLive(listener: any) { publishLive = listener; return () => {}; }, async queryRetention() { return queryResult([], [{ partition: 0, startOffset: "0", endOffset: "1" }]); }, async deliverEvent() { return true; }, async deliverConnected() { return true; }, async onFailure(error: any) { failure = error; } }); await handoff.start(); publishLive(record(1).value, { ...transport(1), topic: "hwlab.event.other.v1" }); await waitFor(() => Boolean(failure)); assert.equal(failure.code, "workbench_kafka_refresh_transport_topic_mismatch"); assert.equal(failure.phase, "live"); assert.equal(handoff.phase(), "failed"); }); function queryResult(events: any[], endOffsets: any[]) { return { topic: "hwlab.event.v1", events, completionReason: "end-offset", reachedEndOffsets: true, endOffsetsAvailable: true, endOffsets, completion: { reason: "end-offset", complete: true, barrierReached: true, retentionStartVerified: true } }; } function record(offset: number, partition = 0) { return { ...transport(offset, partition), value: envelope(offset, partition) }; } function transport(offset: number, partition = 0) { return { topic: "hwlab.event.v1", partition, offset: String(offset) }; } function envelope(offset: number, partition = 0) { const eventId = `evt_${partition}_${offset}`; return { schema: "hwlab.event.v1", eventType: "hwlab.trace.event.projected", eventId: `hwlab:${eventId}`, sourceEventId: eventId, traceId: TRACE_ID, hwlabSessionId: SESSION_ID, sessionId: SESSION_ID, event: { type: "tool", eventType: "tool", traceId: TRACE_ID, sessionId: SESSION_ID, sourceEventId: eventId } }; } async function waitFor(predicate: () => boolean, timeoutMs = 250) { const deadline = Date.now() + timeoutMs; while (!predicate()) { if (Date.now() >= deadline) throw new Error("condition timeout"); await new Promise((resolve) => setTimeout(resolve, 1)); } }