411 lines
18 KiB
TypeScript
411 lines
18 KiB
TypeScript
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<void>((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));
|
|
}
|
|
}
|