fix: 强化 Kafka 重放 barrier 完整性

This commit is contained in:
root
2026-07-10 15:58:59 +02:00
parent 92e9345309
commit 70c5c43012
7 changed files with 413 additions and 31 deletions
+48 -11
View File
@@ -990,7 +990,7 @@ export async function queryKafkaEventStream({ env = process.env, stream = "hwlab
};
}
export async function openKafkaEventStream({ env = process.env, stream = "hwlab", topic = null, groupIdPrefix = null, traceId = null, sessionId = null, runId = null, commandId = null, replayId = null, offsetRange = null, fromBeginning = false, onEvent, onRecord = null, onError = null, kafkaFactory = defaultKafkaFactory } = {}) {
export async function openKafkaEventStream({ env = process.env, stream = "hwlab", topic = null, groupIdPrefix = null, traceId = null, sessionId = null, runId = null, commandId = null, replayId = null, offsetRange = null, fromBeginning = false, onEvent, onRecord = null, onBarrierReady = null, onBarrierComplete = null, onError = null, kafkaFactory = defaultKafkaFactory } = {}) {
if (typeof onEvent !== "function") throw new Error("onEvent callback is required for Kafka event streaming.");
const brokers = csv(env.HWLAB_KAFKA_BOOTSTRAP_SERVERS);
if (brokers.length === 0) throw new Error("HWLAB_KAFKA_BOOTSTRAP_SERVERS is required for Kafka event streams.");
@@ -1004,6 +1004,9 @@ export async function openKafkaEventStream({ env = process.env, stream = "hwlab"
const consumer = kafka.consumer({ groupId, allowAutoTopicCreation: false });
let running = false;
let stopped = false;
let seekApplied = false;
let barrierCompleted = false;
const seenRecordOffsets = new Set();
await consumer.connect();
await consumer.subscribe({ topic: resolvedTopic, fromBeginning: fromBeginning === true });
running = true;
@@ -1011,35 +1014,61 @@ export async function openKafkaEventStream({ env = process.env, stream = "hwlab"
await consumer.run({
eachMessage: async ({ topic: messageTopic, partition, message }) => {
if (stopped) return;
if (resolvedOffsetRange && !seekApplied) return;
const valueText = message.value ? Buffer.from(message.value).toString("utf8") : "";
const value = parseJson(valueText);
const parsed = Boolean(value && typeof value === "object" && !Array.isArray(value));
const inOffsetRange = kafkaRecordInOffsetRange(partition, message.offset, resolvedOffsetRange);
const filterMatches = parsed ? eventFilterResults(value, filters) : {};
const matched = parsed && inOffsetRange && Object.values(filterMatches).every(Boolean);
const offsetKey = `${messageTopic}:${partition}:${message.offset}`;
const duplicate = Boolean(resolvedOffsetRange && seenRecordOffsets.has(offsetKey));
if (resolvedOffsetRange) seenRecordOffsets.add(offsetKey);
const matched = !duplicate && parsed && inOffsetRange && Object.values(filterMatches).every(Boolean);
if (typeof onRecord === "function") {
await onRecord({
topic: messageTopic,
partition,
offset: message.offset,
parsed,
duplicate,
inOffsetRange,
filterMatches,
matched,
valuesPrinted: false
});
}
if (!matched) return;
await onEvent({
topic: messageTopic,
partition,
offset: message.offset,
key: message.key ? Buffer.from(message.key).toString("utf8") : null,
timestamp: message.timestamp ?? null,
value
});
if (matched) {
await onEvent({
topic: messageTopic,
partition,
offset: message.offset,
key: message.key ? Buffer.from(message.key).toString("utf8") : null,
timestamp: message.timestamp ?? null,
value
});
}
if (!barrierCompleted && kafkaRecordCompletesOffsetRange(partition, message.offset, resolvedOffsetRange)) {
barrierCompleted = true;
if (typeof onBarrierComplete === "function") {
await onBarrierComplete({
topic: messageTopic,
partition,
offset: message.offset,
offsetRange: resolvedOffsetRange,
valuesPrinted: false
});
}
}
}
});
if (resolvedOffsetRange) {
if (typeof consumer.seek !== "function") throw new Error("Kafka consumer seek is required for correlated replay offset barriers.");
await consumer.seek({ topic: resolvedTopic, partition: resolvedOffsetRange.partition, offset: resolvedOffsetRange.firstOffset });
seekApplied = true;
if (typeof onBarrierReady === "function") {
await onBarrierReady({ topic: resolvedTopic, offsetRange: resolvedOffsetRange, seekApplied: true, valuesPrinted: false });
}
}
} catch (error) {
if (typeof onError === "function") onError(error);
stopped = true;
@@ -1054,6 +1083,7 @@ export async function openKafkaEventStream({ env = process.env, stream = "hwlab"
groupId,
filters,
offsetRange: resolvedOffsetRange,
seekApplied,
async stop() {
stopped = true;
if (running) await consumer.stop().catch(() => undefined);
@@ -1188,6 +1218,13 @@ function kafkaRecordInOffsetRange(partitionValue, offsetValue, range) {
return partition === range.partition && offset !== null && BigInt(offset) >= BigInt(range.firstOffset) && BigInt(offset) <= BigInt(range.lastOffset);
}
function kafkaRecordCompletesOffsetRange(partitionValue, offsetValue, range) {
if (!range) return false;
const partition = integerValue(partitionValue);
const offset = kafkaOffsetValue(offsetValue);
return partition === range.partition && offset !== null && BigInt(offset) === BigInt(range.lastOffset);
}
function kafkaOffsetValue(value) {
const text = stringValue(value);
if (!text || !/^\d+$/u.test(text)) return null;
@@ -1,7 +1,16 @@
import assert from "node:assert/strict";
import { test } from "bun:test";
import { classifyWorkbenchKafkaDebugReplayResult } from "./workbench-kafka-debug-replay-contract.ts";
import {
classifyWorkbenchKafkaDebugReplayClientOverlay,
classifyWorkbenchKafkaDebugReplayResult,
completeWorkbenchKafkaDebugReplay,
createWorkbenchKafkaDebugReplayTracker,
markWorkbenchKafkaDebugReplayBarrierComplete,
markWorkbenchKafkaDebugReplayConsumerReady,
observeWorkbenchKafkaDebugReplayRecord,
recordWorkbenchKafkaDebugReplayDelivery
} from "./workbench-kafka-debug-replay-contract.ts";
const readyConsumer = { ready: true, groupId: "hwlab-debug-test" };
const deliveredCounts = {
@@ -106,3 +115,122 @@ test("isolated replay keeps source, filter, and decoder failures distinct", () =
assert.equal(filterMismatch.code, "records_scanned_but_filter_mismatch");
assert.equal(decoderRejected.code, "decoder_rejected");
});
test("correlated replay requires seek, lastOffset, and complete unique batch cardinality", () => {
const complete = correlatedReplay();
assert.equal(classifyWorkbenchKafkaDebugReplayResult(complete).code, "terminal_complete");
const seekMissing = correlatedReplay({ barrier: { seekApplied: false, completed: true, lastOffsetObserved: true } });
assert.equal(classifyWorkbenchKafkaDebugReplayResult(seekMissing).code, "consumer_seek_not_applied");
const lastOffsetMissing = correlatedReplay({ barrier: { seekApplied: true, completed: false, lastOffsetObserved: false } });
assert.equal(classifyWorkbenchKafkaDebugReplayResult(lastOffsetMissing).code, "barrier_not_completed");
const middleOffsetMissing = correlatedReplay({
requestedOffsetRange: { partition: 0, firstOffset: "10", lastOffset: "12" },
producer: { invoked: true, sourceMatched: true, publishedCount: 3 }
});
assert.equal(classifyWorkbenchKafkaDebugReplayResult(middleOffsetMissing).code, "barrier_cardinality_mismatch");
const duplicateOffset = correlatedReplay({ rejectedByReason: { duplicate_offset: 1 } });
assert.equal(classifyWorkbenchKafkaDebugReplayResult(duplicateOffset).code, "barrier_duplicate_offset");
});
test("correlated replay fails typed parse, filter, delivery, and producer evidence loss", () => {
const parseLost = correlatedReplay({ counts: { ...correlatedCounts(), parsed: 1, traceMatched: 1, replayMatched: 1, matched: 1, delivered: 1 } });
const filterLost = correlatedReplay({ counts: { ...correlatedCounts(), traceMatched: 1, replayMatched: 1, matched: 1, delivered: 1 } });
const deliveryLost = correlatedReplay({ counts: { ...correlatedCounts(), delivered: 1 } });
const producerFalse = correlatedReplay({ producer: { invoked: false, sourceMatched: true, publishedCount: 2 } });
assert.equal(classifyWorkbenchKafkaDebugReplayResult(parseLost).code, "record_parse_rejected");
assert.equal(classifyWorkbenchKafkaDebugReplayResult(filterLost).code, "records_scanned_but_filter_mismatch");
assert.equal(classifyWorkbenchKafkaDebugReplayResult(deliveryLost).code, "sse_write_failed");
assert.equal(classifyWorkbenchKafkaDebugReplayResult(producerFalse).code, "producer_not_invoked");
});
test("server result declares client counts unavailable and client overlay reuses the pure classifier", () => {
const serverResult = {
...correlatedReplay(),
offsetRange: { requested: { partition: 0, firstOffset: "10", lastOffset: "11" } },
completionReason: "barrier",
classificationScope: "server",
requiresClientOverlay: true,
clientCounts: { available: false, received: null, decoded: null, applied: null }
};
const overlay = classifyWorkbenchKafkaDebugReplayClientOverlay(serverResult, { received: 2, decoded: 2, applied: 1 });
assert.equal(overlay.code, "reducer_rejected");
assert.equal(overlay.classificationScope, "client");
assert.equal(overlay.requiresClientOverlay, false);
assert.deepEqual(overlay.clientCounts, { available: true, received: 2, decoded: 2, applied: 1, source: "client-local-overlay" });
});
test("correlated tracker counts unique offsets and rejects a duplicate even after lastOffset", () => {
const tracker = createWorkbenchKafkaDebugReplayTracker({
mode: "correlated-v2",
replayId: "rpl_unique_offsets",
traceId: "trc_unique_offsets",
offsetRange: { partition: 0, firstOffset: "10", lastOffset: "11" },
producer: { invoked: true, sourceMatched: true, publishedCount: 2 }
});
markWorkbenchKafkaDebugReplayConsumerReady(tracker, { groupId: "hwlab-debug-unique", seekApplied: true });
const observation = (offset: string, duplicate = false) => ({
partition: 0,
offset,
duplicate,
inOffsetRange: true,
parsed: true,
filterMatches: { traceId: true, replayId: true },
matched: !duplicate
});
observeWorkbenchKafkaDebugReplayRecord(tracker, observation("10"));
observeWorkbenchKafkaDebugReplayRecord(tracker, observation("10", true));
observeWorkbenchKafkaDebugReplayRecord(tracker, observation("11"));
recordWorkbenchKafkaDebugReplayDelivery(tracker, { partition: 0, offset: "10", value: {} });
recordWorkbenchKafkaDebugReplayDelivery(tracker, { partition: 0, offset: "11", value: {} });
markWorkbenchKafkaDebugReplayBarrierComplete(tracker);
const result = completeWorkbenchKafkaDebugReplay(tracker, { reason: "barrier", terminalObserved: true });
assert.equal(result.counts.scanned, 3);
assert.equal(result.counts.barrierScanned, 2);
assert.equal(result.barrier.uniqueOffsetCount, 2);
assert.equal(result.rejectedByReason.duplicate_offset, 1);
assert.equal(result.code, "barrier_duplicate_offset");
assert.equal(result.classificationScope, "server");
assert.equal(result.requiresClientOverlay, true);
assert.deepEqual(result.clientCounts, {
available: false,
received: null,
decoded: null,
applied: null,
source: "client-local-overlay-required"
});
});
function correlatedReplay(overrides: Record<string, any> = {}) {
return {
mode: "correlated-v2",
replayId: "rpl_correlated",
traceId: "trc_correlated",
producer: { invoked: true, sourceMatched: true, publishedCount: 2 },
consumer: readyConsumer,
requestedOffsetRange: { partition: 0, firstOffset: "10", lastOffset: "11" },
barrier: { seekApplied: true, completed: true, lastOffsetObserved: true },
rejectedByReason: {},
counts: correlatedCounts(),
terminalObserved: true,
...overrides
};
}
function correlatedCounts() {
return {
scanned: 2,
barrierScanned: 2,
parsed: 2,
traceMatched: 2,
replayMatched: 2,
matched: 2,
delivered: 2,
clientReceived: null,
decoded: null,
applied: null
};
}
@@ -4,6 +4,7 @@ export const WORKBENCH_KAFKA_DEBUG_REPLAY_CONTRACT_VERSION = "workbench-kafka-de
export function createWorkbenchKafkaDebugReplayTracker(input = {}) {
return {
mode: replayMode(input.mode),
replayId: textValue(input.replayId),
traceId: textValue(input.traceId),
topic: textValue(input.topic),
@@ -13,6 +14,14 @@ export function createWorkbenchKafkaDebugReplayTracker(input = {}) {
publishedCount: integerOrNull(input.producer?.publishedCount)
},
consumer: { ready: false, groupId: null },
barrier: {
seekApplied: false,
completed: false,
lastOffsetObserved: false,
firstOffsetObserved: null,
finalOffsetObserved: null,
uniqueOffsetCount: 0
},
requestedOffsetRange: normalizeOffsetRange(input.offsetRange),
counts: {
scanned: 0,
@@ -29,6 +38,7 @@ export function createWorkbenchKafkaDebugReplayTracker(input = {}) {
rejectedByReason: {},
offsets: {
scanned: new Map(),
barrier: new Map(),
matched: new Map(),
delivered: new Map()
},
@@ -39,17 +49,33 @@ export function createWorkbenchKafkaDebugReplayTracker(input = {}) {
export function markWorkbenchKafkaDebugReplayConsumerReady(tracker, input = {}) {
tracker.consumer.ready = true;
tracker.consumer.groupId = textValue(input.groupId);
tracker.barrier.seekApplied = input.seekApplied === true;
return tracker;
}
export function markWorkbenchKafkaDebugReplayBarrierComplete(tracker) {
tracker.barrier.completed = true;
tracker.barrier.lastOffsetObserved = true;
return tracker;
}
export function observeWorkbenchKafkaDebugReplayRecord(tracker, observation = {}) {
tracker.counts.scanned += 1;
recordOffset(tracker.offsets.scanned, observation.partition, observation.offset);
if (observation.duplicate === true) {
reject(tracker, "duplicate_offset");
return tracker;
}
if (observation.inOffsetRange === false) {
reject(tracker, "offset_range_mismatch");
return tracker;
}
tracker.counts.barrierScanned += 1;
recordOffset(tracker.offsets.barrier, observation.partition, observation.offset);
const observed = summarizeOffsets(tracker.offsets.barrier);
tracker.barrier.firstOffsetObserved = observed[0]?.firstOffset ?? null;
tracker.barrier.finalOffsetObserved = observed.at(-1)?.lastOffset ?? null;
tracker.barrier.uniqueOffsetCount = tracker.counts.barrierScanned;
if (observation.parsed !== true) {
reject(tracker, "record_parse_rejected");
return tracker;
@@ -86,6 +112,7 @@ export function completeWorkbenchKafkaDebugReplay(tracker, input = {}) {
};
const terminalObserved = input.terminalObserved === true;
const result = classifyWorkbenchKafkaDebugReplayResult({
mode: tracker.mode,
replayId: tracker.replayId,
traceId: tracker.traceId,
producer: tracker.producer,
@@ -93,15 +120,19 @@ export function completeWorkbenchKafkaDebugReplay(tracker, input = {}) {
requestedOffsetRange: tracker.requestedOffsetRange,
counts,
completionReason: textValue(input.reason),
terminalObserved
terminalObserved,
barrier: tracker.barrier,
rejectedByReason: tracker.rejectedByReason
});
return {
contractVersion: WORKBENCH_KAFKA_DEBUG_REPLAY_CONTRACT_VERSION,
mode: tracker.mode,
replayId: tracker.replayId,
traceId: tracker.traceId,
topic: tracker.topic,
producer: tracker.producer,
consumer: tracker.consumer,
barrier: tracker.barrier,
phase: result.phase,
code: result.code,
ok: result.ok,
@@ -113,15 +144,26 @@ export function completeWorkbenchKafkaDebugReplay(tracker, input = {}) {
offsetRange: {
requested: tracker.requestedOffsetRange,
scanned: summarizeOffsets(tracker.offsets.scanned),
observed: summarizeOffsets(tracker.offsets.barrier),
matched: summarizeOffsets(tracker.offsets.matched),
delivered: summarizeOffsets(tracker.offsets.delivered)
},
sourceLineage: tracker.sourceLineage,
classificationScope: "server",
requiresClientOverlay: true,
clientCounts: {
available: false,
received: null,
decoded: null,
applied: null,
source: "client-local-overlay-required"
},
valuesPrinted: false
};
}
export function classifyWorkbenchKafkaDebugReplayResult(input = {}) {
const mode = replayMode(input.mode);
const counts = normalizeCounts(input.counts);
const producer = input.producer ?? {};
const consumer = input.consumer ?? {};
@@ -129,9 +171,13 @@ export function classifyWorkbenchKafkaDebugReplayResult(input = {}) {
const hasReplayFilter = Boolean(textValue(input.replayId));
const hasOffsetBarrier = Boolean(normalizeOffsetRange(input.requestedOffsetRange));
if (producer.invoked === false) return failure("producer", "producer_not_invoked", "No isolated debug producer was invoked for this replay.");
if (consumer.ready === false) return failure("consumer", "consumer_not_ready", "The isolated Kafka consumer did not reach its ready barrier.");
if (producer.sourceMatched === false) return failure("producer", "source_trace_missing", "The source AgentRun topic did not contain the requested trace lineage.");
if (producer.invoked === false && counts.matched === 0) return failure("producer", "producer_not_invoked", "No isolated debug producer was invoked for this replay.");
if (mode === "correlated-v2") {
const correlated = classifyCorrelatedReplay(input, counts);
if (correlated) return correlated;
}
if (hasOffsetBarrier && counts.barrierScanned === 0) return failure("kafka", "topic_no_append_for_replay", "The requested replay offset range was not appended or observed.");
if (counts.scanned === 0) return failure("producer", "producer_not_invoked", "The isolated debug topic contained no records for this replay.");
if (counts.barrierScanned > 0 && counts.parsed === 0) return failure("server", "record_parse_rejected", "Kafka records were scanned but none decoded as JSON objects.");
@@ -150,14 +196,83 @@ export function classifyWorkbenchKafkaDebugReplayResult(input = {}) {
return { ok: true, phase: "terminal", code: "terminal_complete", message: "The isolated replay reached a terminal event." };
}
export function classifyWorkbenchKafkaDebugReplayClientOverlay(serverResult = {}, clientCounts = {}) {
const received = nullableCount(clientCounts.received ?? clientCounts.clientReceived, null);
const decoded = nullableCount(clientCounts.decoded, null);
const applied = nullableCount(clientCounts.applied, null);
const counts = { ...normalizeCounts(serverResult.counts), clientReceived: received, decoded, applied };
const complete = received !== null && decoded !== null && applied !== null;
const result = complete
? classifyWorkbenchKafkaDebugReplayResult({
mode: serverResult.mode,
replayId: serverResult.replayId,
traceId: serverResult.traceId,
producer: serverResult.producer,
consumer: serverResult.consumer,
requestedOffsetRange: serverResult.offsetRange?.requested,
counts,
completionReason: serverResult.completionReason,
terminalObserved: serverResult.terminalObserved,
barrier: serverResult.barrier,
rejectedByReason: serverResult.rejectedByReason
})
: failure("client", "client_counts_incomplete", "Client replay classification requires received, decoded, and applied counts.");
return {
...serverResult,
...result,
counts,
classificationScope: "client",
requiresClientOverlay: !complete,
clientCounts: {
available: complete,
received,
decoded,
applied,
source: "client-local-overlay"
},
valuesPrinted: false
};
}
export function normalizeWorkbenchKafkaDebugOffsetRange(value) {
return normalizeOffsetRange(value);
}
export function workbenchKafkaDebugOffsetRangeWidth(value) {
const range = normalizeOffsetRange(value);
if (!range) return null;
const width = BigInt(range.lastOffset) - BigInt(range.firstOffset) + 1n;
return width <= BigInt(Number.MAX_SAFE_INTEGER) ? Number(width) : null;
}
function classifyCorrelatedReplay(input, counts) {
const range = normalizeOffsetRange(input.requestedOffsetRange);
const expected = workbenchKafkaDebugOffsetRangeWidth(range);
const publishedCount = integerOrNull(input.producer?.publishedCount);
if (!textValue(input.replayId) || !range) {
return failure("contract", "v2_correlation_incomplete", "Correlated replay requires replayId and a complete offset range.");
}
if (expected === null || (publishedCount !== null && publishedCount !== expected)) return failure("producer", "producer_cardinality_mismatch", "Producer publishedCount does not equal the requested replay offset range width.");
if (input.barrier?.seekApplied !== true) return failure("consumer", "consumer_seek_not_applied", "The correlated replay consumer did not seek to the requested firstOffset.");
if (input.barrier?.completed !== true || input.barrier?.lastOffsetObserved !== true) return failure("kafka", "barrier_not_completed", "The correlated replay did not observe its requested lastOffset completion boundary.");
if (Number(input.rejectedByReason?.duplicate_offset ?? 0) > 0) return failure("kafka", "barrier_duplicate_offset", "The correlated replay observed a duplicate Kafka offset.");
if (counts.barrierScanned !== expected) return failure("kafka", "barrier_cardinality_mismatch", "The consumer did not observe every offset in the requested replay barrier.");
if (counts.parsed !== expected) return failure("server", "record_parse_rejected", "One or more records in the replay barrier failed JSON object decoding.");
if (counts.traceMatched !== expected || counts.replayMatched !== expected || counts.matched !== expected) {
return failure("server", "records_scanned_but_filter_mismatch", "One or more records in the replay barrier failed traceId or replayId correlation.");
}
if (counts.delivered !== expected) return failure("server", "sse_write_failed", "One or more correlated replay records were not delivered over SSE.");
return null;
}
function failure(phase, code, message) {
return { ok: false, phase, code, message };
}
function replayMode(value) {
return value === "correlated-v2" ? "correlated-v2" : "trace-only-v1";
}
function reject(tracker, reason) {
tracker.rejectedByReason[reason] = (tracker.rejectedByReason[reason] ?? 0) + 1;
}
@@ -105,27 +105,33 @@ test("isolated Workbench Kafka replay applies replayId and offset barriers with
});
await listen(server);
try {
const query = "stream=hwlab-debug&traceId=trc_barrier&replayId=rpl_barrier&fromBeginning=true&partition=0&firstOffset=10&lastOffset=11&producerInvoked=true&sourceMatched=true&publishedCount=2";
const query = "stream=hwlab-debug&traceId=trc_barrier&replayId=rpl_barrier&fromBeginning=true&partition=0&firstOffset=10&lastOffset=11";
const response = await fetch(`${serverUrl(server)}/v1/workbench/debug/kafka-sse/events?${query}`);
const reader = response.body?.getReader();
assert.ok(reader);
const connected = await readUntil(reader, "hwlab.kafka.connected");
assert.match(connected, /"replayId":"rpl_barrier"/u);
assert.match(connected, /"firstOffset":"10"/u);
assert.match(connected, /"replayMode":"correlated-v2"/u);
assert.match(connected, /"seekApplied":true/u);
assert.equal(fakeKafka.fromBeginning, false);
assert.deepEqual(fakeKafka.seekCalls, [{ topic: "hwlab.event.debug.v1", partition: 0, offset: "10" }]);
await fakeKafka.emit({ schema: "hwlab.event.debug.v1", traceId: "trc_barrier", replayId: "rpl_barrier", event: { type: "assistant" } }, { offset: "9" });
await fakeKafka.emit({ schema: "hwlab.event.debug.v1", traceId: "trc_barrier", replayId: "rpl_barrier", debugLineage: { inputTraceId: "trc_barrier", inputSessionId: "ses_agentrun_barrier", sourceSeq: 1 }, event: { type: "assistant" } }, { offset: "10" });
await fakeKafka.emit({ schema: "hwlab.event.debug.v1", traceId: "trc_barrier", replayId: "rpl_barrier", debugLineage: { inputTraceId: "trc_barrier", inputSessionId: "ses_agentrun_barrier", sourceSeq: 2 }, eventType: "terminal", event: { type: "result", terminal: true } }, { offset: "11" });
const replay = await readUntil(reader, "hwlab.kafka.replay-complete");
assert.match(replay, /"code":"terminal_complete"/u);
assert.match(replay, /"scanned":3/u);
assert.match(replay, /"reason":"barrier"/u);
assert.match(replay, /"scanned":2/u);
assert.match(replay, /"barrierScanned":2/u);
assert.match(replay, /"parsed":2/u);
assert.match(replay, /"matched":2/u);
assert.match(replay, /"delivered":2/u);
assert.match(replay, /"offset_range_mismatch":1/u);
assert.match(replay, /"firstOffset":"10","lastOffset":"11","count":2/u);
assert.match(replay, /"inputSessionId":"ses_agentrun_barrier"/u);
assert.match(replay, /"barrierComplete":true/u);
assert.match(replay, /"clientCountsAvailable":false/u);
await waitUntil(() => fakeKafka.stopCalls > 0);
} finally {
await close(server);
}
@@ -190,6 +196,26 @@ test("isolated Workbench Kafka debug fails closed when disabled and rejects repl
assert.equal(invalidReplayId.status, 400);
assert.match(await invalidReplayId.text(), /workbench_kafka_debug_replay_id_invalid/u);
assert.equal(fakeKafka.consumerCreated, false);
const replayOnly = await fetch(`${serverUrl(server)}/v1/workbench/debug/kafka-sse/events?stream=hwlab-debug&traceId=trc_replay&fromBeginning=true&replayId=rpl_incomplete`);
assert.equal(replayOnly.status, 400);
assert.match(await replayOnly.text(), /workbench_kafka_debug_v2_correlation_incomplete/u);
assert.equal(fakeKafka.consumerCreated, false);
const rangeOnly = await fetch(`${serverUrl(server)}/v1/workbench/debug/kafka-sse/events?stream=hwlab-debug&traceId=trc_replay&fromBeginning=true&partition=0&firstOffset=10&lastOffset=11`);
assert.equal(rangeOnly.status, 400);
assert.match(await rangeOnly.text(), /workbench_kafka_debug_v2_correlation_incomplete/u);
assert.equal(fakeKafka.consumerCreated, false);
const producerNotInvoked = await fetch(`${serverUrl(server)}/v1/workbench/debug/kafka-sse/events?stream=hwlab-debug&traceId=trc_replay&fromBeginning=true&replayId=rpl_not_invoked&partition=0&firstOffset=10&lastOffset=11&producerInvoked=false`);
assert.equal(producerNotInvoked.status, 409);
assert.match(await producerNotInvoked.text(), /workbench_kafka_debug_producer_not_invoked/u);
assert.equal(fakeKafka.consumerCreated, false);
const cardinalityMismatch = await fetch(`${serverUrl(server)}/v1/workbench/debug/kafka-sse/events?stream=hwlab-debug&traceId=trc_replay&fromBeginning=true&replayId=rpl_bad_count&partition=0&firstOffset=10&lastOffset=11&publishedCount=3`);
assert.equal(cardinalityMismatch.status, 400);
assert.match(await cardinalityMismatch.text(), /workbench_kafka_debug_producer_cardinality_mismatch/u);
assert.equal(fakeKafka.consumerCreated, false);
} finally {
await close(server);
}
@@ -408,6 +434,8 @@ function createFakeKafkaFactory(options: { deferRun?: boolean } = {}) {
let stopCalls = 0;
let disconnectCalls = 0;
let consumerCreated = false;
let runStarted = false;
const seekCalls: Array<{ topic: string; partition: number; offset: string }> = [];
const consumer = {
connect: async () => undefined,
subscribe: async (input: { topic?: string; fromBeginning?: boolean }) => {
@@ -416,8 +444,13 @@ function createFakeKafkaFactory(options: { deferRun?: boolean } = {}) {
},
run: async (input: any) => {
if (runGate) await runGate;
runStarted = true;
eachMessage = input.eachMessage;
},
seek: async (input: { topic: string; partition: number; offset: string }) => {
assert.equal(runStarted, true, "seek must happen after consumer.run assignment");
seekCalls.push(input);
},
stop: async () => { stopCalls += 1; },
disconnect: async () => { disconnectCalls += 1; }
};
@@ -435,6 +468,7 @@ function createFakeKafkaFactory(options: { deferRun?: boolean } = {}) {
get stopCalls() { return stopCalls; },
get disconnectCalls() { return disconnectCalls; },
get consumerCreated() { return consumerCreated; },
get seekCalls() { return seekCalls; },
releaseRun: () => releaseRun?.(),
emit: async (value: Record<string, unknown>, transport: { offset?: string; partition?: number } = {}) => {
assert.ok(eachMessage, "consumer.run must be called before emitting fake Kafka events");
+68 -10
View File
@@ -10,10 +10,12 @@ import { workbenchKafkaDebugCapability } from "./workbench-kafka-debug-capabilit
import {
completeWorkbenchKafkaDebugReplay,
createWorkbenchKafkaDebugReplayTracker,
markWorkbenchKafkaDebugReplayBarrierComplete,
markWorkbenchKafkaDebugReplayConsumerReady,
normalizeWorkbenchKafkaDebugOffsetRange,
observeWorkbenchKafkaDebugReplayRecord,
recordWorkbenchKafkaDebugReplayDelivery
recordWorkbenchKafkaDebugReplayDelivery,
workbenchKafkaDebugOffsetRangeWidth
} from "./workbench-kafka-debug-replay-contract.ts";
const CONTRACT_VERSION = "workbench-debug-kafka-sse-v1";
@@ -83,7 +85,9 @@ async function openKafkaDebugSse(request, response, url, options, isolatedDebug
const isolatedReplay = stream === "hwlab-debug";
const resolvedFilters = isolatedReplay ? compactObject({ traceId: filters.traceId, replayId: filters.replayId }) : resolveKafkaDebugFilters(filters, options);
const replayRequest = isolatedDebug?.replayRequest ?? replayRequestFromUrl(url);
const correlatedReplay = isolatedReplay && replayRequest.mode === "correlated-v2";
const replayTracker = isolatedReplay ? createWorkbenchKafkaDebugReplayTracker({
mode: replayRequest.mode,
replayId: resolvedFilters.replayId,
traceId: resolvedFilters.traceId,
topic: isolatedDebug.topic,
@@ -105,6 +109,7 @@ async function openKafkaDebugSse(request, response, url, options, isolatedDebug
let replayTimer = null;
let replayCount = 0;
let replayTerminalObserved = false;
let barrierCompletionPending = false;
const bufferedRecords = [];
const close = () => {
if (closed) return;
@@ -130,6 +135,7 @@ async function openKafkaDebugSse(request, response, url, options, isolatedDebug
topic: isolatedDebug.topic,
traceId: resolvedFilters.traceId,
replayComplete: true,
replayMode: replayResult.mode,
replayId: replayResult.replayId,
phase: replayResult.phase,
code: replayResult.code,
@@ -143,6 +149,9 @@ async function openKafkaDebugSse(request, response, url, options, isolatedDebug
limit: isolatedDebug.replayLimit,
timeoutMs: isolatedDebug.replayTimeoutMs,
terminalObserved: replayTerminalObserved,
barrierComplete: replayResult.barrier.completed,
stopBoundary: "sse-end",
clientCountsAvailable: false,
valuesPrinted: false
});
if (!response.writableEnded) response.end();
@@ -166,8 +175,8 @@ async function openKafkaDebugSse(request, response, url, options, isolatedDebug
recordWorkbenchKafkaDebugReplayDelivery(replayTracker, record, delivered);
replayCount = replayTracker.counts.delivered;
replayTerminalObserved ||= debugRecordIsTerminal(record);
if (replayTerminalObserved) finishReplay("terminal");
else if (replayCount >= isolatedDebug.replayLimit) finishReplay("limit");
if (!correlatedReplay && replayTerminalObserved) finishReplay("terminal");
else if (!correlatedReplay && replayCount >= isolatedDebug.replayLimit) finishReplay("limit");
};
if (isolatedReplay) {
replayTimer = setTimeout(() => finishReplay("timeout"), isolatedDebug.replayTimeoutMs);
@@ -179,13 +188,18 @@ async function openKafkaDebugSse(request, response, url, options, isolatedDebug
stream,
topic: isolatedDebug?.topic ?? null,
groupIdPrefix: isolatedDebug?.groupIdPrefix ?? null,
fromBeginning: isolatedReplay || url.searchParams.get("fromBeginning") === "1" || url.searchParams.get("fromBeginning") === "true",
fromBeginning: isolatedReplay ? !correlatedReplay : url.searchParams.get("fromBeginning") === "1" || url.searchParams.get("fromBeginning") === "true",
...resolvedFilters,
offsetRange: replayRequest.offsetRange,
kafkaFactory: options.kafkaFactory,
onRecord: isolatedReplay ? async (observation) => {
observeWorkbenchKafkaDebugReplayRecord(replayTracker, observation);
} : null,
onBarrierComplete: correlatedReplay ? async () => {
markWorkbenchKafkaDebugReplayBarrierComplete(replayTracker);
if (consumerReady) finishReplay("barrier");
else barrierCompletionPending = true;
} : null,
onEvent: async (record) => {
if (isolatedReplay && !consumerReady) {
if (bufferedRecords.length < isolatedDebug.replayLimit) bufferedRecords.push(record);
@@ -198,7 +212,7 @@ async function openKafkaDebugSse(request, response, url, options, isolatedDebug
await kafkaStream.stop?.();
return;
}
if (isolatedReplay) markWorkbenchKafkaDebugReplayConsumerReady(replayTracker, { groupId: kafkaStream.groupId });
if (isolatedReplay) markWorkbenchKafkaDebugReplayConsumerReady(replayTracker, { groupId: kafkaStream.groupId, seekApplied: kafkaStream.seekApplied });
writeSse(response, "hwlab.kafka.connected", {
ok: true,
contractVersion: CONTRACT_VERSION,
@@ -213,8 +227,10 @@ async function openKafkaDebugSse(request, response, url, options, isolatedDebug
replayLimit: isolatedDebug?.replayLimit ?? null,
replayTimeoutMs: isolatedDebug?.replayTimeoutMs ?? null,
replayId: replayTracker?.replayId ?? null,
replayMode: replayTracker?.mode ?? null,
offsetRange: replayTracker?.requestedOffsetRange ?? null,
replayContractVersion: replayTracker ? "workbench-kafka-debug-replay-v2" : null,
replayContractVersion: replayTracker?.mode === "correlated-v2" ? "workbench-kafka-debug-replay-v2" : CONTRACT_VERSION,
seekApplied: kafkaStream.seekApplied === true,
filters,
resolvedFilters,
serverSentAt: new Date().toISOString(),
@@ -225,6 +241,7 @@ async function openKafkaDebugSse(request, response, url, options, isolatedDebug
if (replayFinished || closed) break;
writeRecord(record);
}
if (barrierCompletionPending && !replayFinished && !closed) finishReplay("barrier");
} catch (error) {
if (replayTimer) clearTimeout(replayTimer);
writeSse(response, "hwlab.kafka.error", { ok: false, error: errorMessagePayload(error), valuesPrinted: false });
@@ -334,10 +351,41 @@ function isolatedDebugConfig(response, url, env, options = {}) {
sendJson(response, 400, debugError("workbench_kafka_debug_offset_range_invalid", "Replay offset barrier requires partition, firstOffset, and lastOffset with firstOffset <= lastOffset."));
return false;
}
if (options.requireReplay && replayRequest.mode === "correlated-v2") {
if (!replayRequest.replayId || !replayRequest.offsetRange) {
sendJson(response, 400, debugError("workbench_kafka_debug_v2_correlation_incomplete", "Correlated V2 replay requires replayId and partition/firstOffset/lastOffset."));
return false;
}
if (replayRequest.producerMetadataInvalid) {
sendJson(response, 400, debugError("workbench_kafka_debug_v2_producer_metadata_invalid", "Optional producerInvoked, sourceMatched, and publishedCount must use valid boolean/integer values."));
return false;
}
if (replayRequest.producer.invoked === false) {
sendJson(response, 409, debugError("workbench_kafka_debug_producer_not_invoked", "Correlated V2 replay cannot use producerInvoked=false."));
return false;
}
if (replayRequest.producer.sourceMatched === false) {
sendJson(response, 409, debugError("workbench_kafka_debug_source_trace_missing", "Correlated V2 replay cannot start because the producer did not match the requested source trace."));
return false;
}
const rangeWidth = workbenchKafkaDebugOffsetRangeWidth(replayRequest.offsetRange);
if (!rangeWidth || (replayRequest.producer.publishedCount !== null && replayRequest.producer.publishedCount !== rangeWidth)) {
sendJson(response, 400, debugError("workbench_kafka_debug_producer_cardinality_mismatch", "publishedCount must equal the requested offset range width."));
return false;
}
if (rangeWidth > config.replayLimit) {
sendJson(response, 400, debugError("workbench_kafka_debug_replay_limit_exceeded", "The correlated replay batch exceeds the YAML-owned replay limit."));
return false;
}
}
return { ...config, replayRequest };
}
function replayRequestFromUrl(url) {
const rawReplayId = url.searchParams.get("replayId") || url.searchParams.get("replay-id");
const rawProducerInvoked = url.searchParams.get("producerInvoked") || url.searchParams.get("producer-invoked");
const rawSourceMatched = url.searchParams.get("sourceMatched") || url.searchParams.get("source-matched");
const rawPublishedCount = url.searchParams.get("publishedCount") || url.searchParams.get("published-count");
const rawRange = {
partition: url.searchParams.get("partition"),
firstOffset: url.searchParams.get("firstOffset") || url.searchParams.get("first-offset") || url.searchParams.get("startOffset") || url.searchParams.get("start-offset"),
@@ -345,14 +393,24 @@ function replayRequestFromUrl(url) {
};
const offsetRangeProvided = Object.values(rawRange).some((entry) => textValue(entry));
const offsetRange = normalizeWorkbenchKafkaDebugOffsetRange(rawRange);
const correlationProvided = Boolean(textValue(rawReplayId) || offsetRangeProvided || textValue(rawProducerInvoked) || textValue(rawSourceMatched) || textValue(rawPublishedCount));
const producerInvoked = optionalBoolean(rawProducerInvoked);
const sourceMatched = optionalBoolean(rawSourceMatched);
const publishedCount = optionalNonNegativeInteger(rawPublishedCount);
return {
replayId: safeReplayId(url.searchParams.get("replayId") || url.searchParams.get("replay-id")),
mode: correlationProvided ? "correlated-v2" : "trace-only-v1",
replayId: safeReplayId(rawReplayId),
offsetRange,
offsetRangeInvalid: offsetRangeProvided && !offsetRange,
producerMetadataInvalid: Boolean(
(textValue(rawProducerInvoked) && producerInvoked === null)
|| (textValue(rawSourceMatched) && sourceMatched === null)
|| (textValue(rawPublishedCount) && publishedCount === null)
),
producer: {
invoked: optionalBoolean(url.searchParams.get("producerInvoked") || url.searchParams.get("producer-invoked")),
sourceMatched: optionalBoolean(url.searchParams.get("sourceMatched") || url.searchParams.get("source-matched")),
publishedCount: optionalNonNegativeInteger(url.searchParams.get("publishedCount") || url.searchParams.get("published-count"))
invoked: producerInvoked,
sourceMatched,
publishedCount
},
valuesPrinted: false
};
+8
View File
@@ -13,6 +13,13 @@ const HWLAB_SESSION_ID = "ses_5ec4e141-6abc-466d-9afe-049f7c0ac105";
const TRACE_ID = "trc_mretx18t3jl4tg";
const RUN_ID = "run_a78cbcb05f2c4afaa748a3158db44998";
test("Kafka help shows canonical default input and explicit reconstruction debug input", async () => {
const result = await runKafkaCli(["help", "--json"], { env: {}, now: () => "2026-07-10T12:00:00.000Z" });
assert.equal(result.exitCode, 0);
assert.match(result.payload.commands[0], /--input-topic agentrun\.event\.v1/u);
assert.match(result.payload.commands[1], /--input-topic agentrun\.event\.debug\.v1/u);
});
test("offline JSONL maps 35 stdio reconstructions to 35 ordered HWLAB debug events without runtime dependencies", async () => {
const cwd = await mkdtemp(path.join(os.tmpdir(), "hwlab-kafka-debug-jsonl-"));
const inputFile = path.join(cwd, "agentrun-events.jsonl");
@@ -204,6 +211,7 @@ test("canonical Kafka regeneration preflights without publishing and preserves r
assert.equal(result.payload.replay.producerInvoked, false);
assert.equal(result.payload.replay.outputBarrier, null);
assert.match(result.payload.next.command, /--replay-id rpl_canonical_preflight/u);
assert.match(result.payload.next.command, /--group-prefix hwlab-v03-workbench-isolated-debug/u);
assert.match(result.payload.next.command, /--publish/u);
assert.equal(producerCreated, false);
});
+5 -3
View File
@@ -170,9 +170,10 @@ export async function regenerateHwlabDebugEvents(parsed: ParsedArgs, dependencie
}
const sourceFile = text(readResult.sourceFile);
const traceArgument = traceId ? ` --trace-id ${traceId}` : "";
const groupArgument = groupPrefix ? ` --group-prefix ${groupPrefix}` : "";
const baseCommand = sourceMode === "jsonl"
? `hwlab-cli kafka regenerate hwlab --from jsonl --session-id ${sessionId}${traceArgument} --replay-id ${replayId} --jsonl-file ${shellArg(sourceFile || String(parsed.jsonlFile))}`
: `hwlab-cli kafka regenerate hwlab --from kafka --session-id ${sessionId}${traceArgument} --replay-id ${replayId} --input-topic ${inputTopic}`;
? `hwlab-cli kafka regenerate hwlab --from jsonl --session-id ${sessionId}${traceArgument} --replay-id ${replayId}${groupArgument} --jsonl-file ${shellArg(sourceFile || String(parsed.jsonlFile))}`
: `hwlab-cli kafka regenerate hwlab --from kafka --session-id ${sessionId}${traceArgument} --replay-id ${replayId}${groupArgument} --input-topic ${inputTopic}`;
const nextCommand = shouldPublish
? `${baseCommand} --no-publish --output-topic ${outputTopic} --expect-count ${mapped.events.length} --json`
: `${baseCommand} --publish --output-topic ${outputTopic} --expect-count ${mapped.events.length} --json`;
@@ -317,7 +318,8 @@ function kafkaHelp() {
action: "kafka.help",
status: "succeeded",
commands: [
"hwlab-cli kafka regenerate hwlab --from kafka --session-id ses_... [--trace-id trc_...] [--replay-id rpl_...] [--input-topic agentrun.event.debug.v1] [--expect-count 35] [--json]",
"hwlab-cli kafka regenerate hwlab --from kafka --session-id ses_... [--trace-id trc_...] [--replay-id rpl_...] [--input-topic agentrun.event.v1] [--group-prefix hwlab-...-debug] [--expect-count 35] [--json]",
"hwlab-cli kafka regenerate hwlab --from kafka --session-id ses_... --input-topic agentrun.event.debug.v1 --group-prefix hwlab-...-debug [--replay-id rpl_...] [--json]",
"hwlab-cli kafka regenerate hwlab --from jsonl --session-id ses_... [--replay-id rpl_...] --jsonl-file agentrun-events.jsonl --no-publish [--output-jsonl hwlab-events.jsonl] [--json]"
],
defaults: {