From 70c5c43012d1e98d6c8a0c20747b5366944490fd Mon Sep 17 00:00:00 2001 From: root Date: Fri, 10 Jul 2026 15:58:59 +0200 Subject: [PATCH] =?UTF-8?q?fix:=20=E5=BC=BA=E5=8C=96=20Kafka=20=E9=87=8D?= =?UTF-8?q?=E6=94=BE=20barrier=20=E5=AE=8C=E6=95=B4=E6=80=A7?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- internal/cloud/kafka-event-bridge.ts | 59 ++++++-- ...kbench-kafka-debug-replay-contract.test.ts | 130 +++++++++++++++++- .../workbench-kafka-debug-replay-contract.ts | 119 +++++++++++++++- .../cloud/workbench-kafka-sse-debug.test.ts | 42 +++++- internal/cloud/workbench-kafka-sse-debug.ts | 78 +++++++++-- tools/hwlab-cli/kafka-regenerate.test.ts | 8 ++ tools/src/hwlab-cli/kafka-regenerate.ts | 8 +- 7 files changed, 413 insertions(+), 31 deletions(-) diff --git a/internal/cloud/kafka-event-bridge.ts b/internal/cloud/kafka-event-bridge.ts index 5eaca755..58cb861a 100644 --- a/internal/cloud/kafka-event-bridge.ts +++ b/internal/cloud/kafka-event-bridge.ts @@ -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; diff --git a/internal/cloud/workbench-kafka-debug-replay-contract.test.ts b/internal/cloud/workbench-kafka-debug-replay-contract.test.ts index bc46a521..676d8331 100644 --- a/internal/cloud/workbench-kafka-debug-replay-contract.test.ts +++ b/internal/cloud/workbench-kafka-debug-replay-contract.test.ts @@ -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 = {}) { + 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 + }; +} diff --git a/internal/cloud/workbench-kafka-debug-replay-contract.ts b/internal/cloud/workbench-kafka-debug-replay-contract.ts index d01a4016..57c2ae6f 100644 --- a/internal/cloud/workbench-kafka-debug-replay-contract.ts +++ b/internal/cloud/workbench-kafka-debug-replay-contract.ts @@ -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; } diff --git a/internal/cloud/workbench-kafka-sse-debug.test.ts b/internal/cloud/workbench-kafka-sse-debug.test.ts index 2d718ad1..ac1318fc 100644 --- a/internal/cloud/workbench-kafka-sse-debug.test.ts +++ b/internal/cloud/workbench-kafka-sse-debug.test.ts @@ -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, transport: { offset?: string; partition?: number } = {}) => { assert.ok(eachMessage, "consumer.run must be called before emitting fake Kafka events"); diff --git a/internal/cloud/workbench-kafka-sse-debug.ts b/internal/cloud/workbench-kafka-sse-debug.ts index 1e9fbd20..883d5c01 100644 --- a/internal/cloud/workbench-kafka-sse-debug.ts +++ b/internal/cloud/workbench-kafka-sse-debug.ts @@ -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 }; diff --git a/tools/hwlab-cli/kafka-regenerate.test.ts b/tools/hwlab-cli/kafka-regenerate.test.ts index e48c619e..a07db059 100644 --- a/tools/hwlab-cli/kafka-regenerate.test.ts +++ b/tools/hwlab-cli/kafka-regenerate.test.ts @@ -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); }); diff --git a/tools/src/hwlab-cli/kafka-regenerate.ts b/tools/src/hwlab-cli/kafka-regenerate.ts index 42591860..5cca5b71 100644 --- a/tools/src/hwlab-cli/kafka-regenerate.ts +++ b/tools/src/hwlab-cli/kafka-regenerate.ts @@ -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: {