diff --git a/docs/reference/kafka-source-direct-debug.md b/docs/reference/kafka-source-direct-debug.md index 3ed316d8..c3f91bb7 100644 --- a/docs/reference/kafka-source-direct-debug.md +++ b/docs/reference/kafka-source-direct-debug.md @@ -8,6 +8,19 @@ - group prefix 必须标识独立 debug group,不复用产品 consumer group; - 命令只读取 Kafka,不发布事件,不启动 Cloud API,不访问数据库,也不依赖 projector。 +- `hwlab-cli kafka inspect order` 用于定位同一 session 在纯 Kafka 刷新链路中的首个顺序分歧: + - 已知 session 时使用 `--session-id `; + - 只知道用户输入时使用 `--contains-text ` 从 HWLAB retention 解析唯一 session; + - 文本发现只披露长度和 SHA-256,不回显用户正文; + - 命令分别扫描 `agentrun.event.v1` 与 `hwlab.event.v1` 到启动时捕获的 topic 上界; + - `hwlabSessionId` 及其 snake case 形式属于两个 topic 共用的正式 session identity 候选; + - HWLAB retained event 复用生产 refresh handoff、SSE frame 解码、queue/reducer 和 timeline row model; + - 输出逐层 transport、stable identity、`sourceSeq`、message index 和 DOM row index; + - `firstDivergence` 只报告首个改变顺序的层,不过滤、重排或补写任何事件; + - `--row-limit` 只限制返回的有界尾部明细,不改变全量计数和判定; + - 任一 topic 未到捕获上界时返回 `partial/source_scan_incomplete`,禁止根据部分扫描下结论; + - 命令只读取既有 Kafka retention,不发布事件,也不启动新的 AgentRun。 + - Trace 渲染与 Web 共享最后一个展示分叉之前的生产管线: - JSON frame 使用 `decodeWorkbenchRealtimeEventFrame` 解码; - 事件使用 `reduceWorkbenchRealtimeEvent` 与 `planWorkbenchRealtimeApply` 分类; diff --git a/internal/cloud/kafka-event-bridge.ts b/internal/cloud/kafka-event-bridge.ts index acfdcafe..c38e27d3 100644 --- a/internal/cloud/kafka-event-bridge.ts +++ b/internal/cloud/kafka-event-bridge.ts @@ -1443,7 +1443,7 @@ function eventFieldCandidates(value, name) { const nestedEvent = objectValue(value.event?.sourceEvent ?? value.agentRunEvent ?? payload.event); const nestedPayload = objectValue(nestedEvent.payload); if (name === "traceId") return textCandidates(value.traceId, value.trace_id, event.traceId, event.trace_id, context.traceId, context.trace_id, sourceEvent.traceId, sourceEvent.trace_id, nestedEvent.traceId, nestedEvent.trace_id, nestedPayload.traceId, nestedPayload.trace_id, payload.traceId, payload.trace_id, metadata.traceId, metadata.trace_id, trace.traceId, trace.trace_id, ids.traceId, ids.trace_id, stdio.traceId, stdio.trace_id); - if (name === "sessionId") return textCandidates(value.sessionId, value.session_id, event.sessionId, event.session_id, context.sessionId, context.session_id, run.sessionId, run.session_id, nestedPayload.sessionId, nestedPayload.session_id, payload.sessionId, payload.session_id, metadata.sessionId, metadata.session_id, ids.sessionId, ids.session_id, stdio.sessionId, stdio.session_id); + if (name === "sessionId") return textCandidates(value.sessionId, value.session_id, value.hwlabSessionId, value.hwlab_session_id, event.sessionId, event.session_id, event.hwlabSessionId, event.hwlab_session_id, context.sessionId, context.session_id, context.hwlabSessionId, context.hwlab_session_id, run.sessionId, run.session_id, run.hwlabSessionId, run.hwlab_session_id, nestedPayload.sessionId, nestedPayload.session_id, nestedPayload.hwlabSessionId, nestedPayload.hwlab_session_id, payload.sessionId, payload.session_id, payload.hwlabSessionId, payload.hwlab_session_id, metadata.sessionId, metadata.session_id, metadata.hwlabSessionId, metadata.hwlab_session_id, ids.sessionId, ids.session_id, ids.hwlabSessionId, ids.hwlab_session_id, stdio.sessionId, stdio.session_id, stdio.hwlabSessionId, stdio.hwlab_session_id); if (name === "runId") return textCandidates(value.runId, value.run_id, event.runId, event.run_id, context.runId, context.run_id, run.runId, run.run_id, nestedEvent.runId, nestedEvent.run_id, nestedPayload.runId, nestedPayload.run_id, payload.runId, payload.run_id, metadata.runId, metadata.run_id, ids.runId, ids.run_id, stdio.runId, stdio.run_id); if (name === "commandId") return textCandidates(value.commandId, value.command_id, event.commandId, event.command_id, context.commandId, context.command_id, command.commandId, command.command_id, nestedEvent.commandId, nestedEvent.command_id, nestedPayload.commandId, nestedPayload.command_id, payload.commandId, payload.command_id, metadata.commandId, metadata.command_id, ids.commandId, ids.command_id, stdio.commandId, stdio.command_id); if (name === "replayId") return textCandidates(value.replayId, value.replay_id, debugReplay.replayId, debugReplay.replay_id, debugLineage.replayId, debugLineage.replay_id); diff --git a/tools/hwlab-cli/kafka-regenerate.test.ts b/tools/hwlab-cli/kafka-regenerate.test.ts index 5b0b091c..494571b0 100644 --- a/tools/hwlab-cli/kafka-regenerate.test.ts +++ b/tools/hwlab-cli/kafka-regenerate.test.ts @@ -5,7 +5,7 @@ import path from "node:path"; import { test } from "bun:test"; -import { queryKafkaEventStream } from "../../internal/cloud/kafka-event-bridge.ts"; +import { projectAgentRunKafkaEventToHwlabEvent, queryKafkaEventStream } from "../../internal/cloud/kafka-event-bridge.ts"; import { mapAgentRunRecordsToHwlabDebugEvents, renderKafkaCliText, runKafkaCli } from "../src/hwlab-cli/kafka-regenerate.ts"; import { renderTraceRowsMarkdown, traceDisplayRows, traceToolSummary } from "../src/hwlab-cli/trace-renderer.ts"; @@ -20,6 +20,92 @@ test("Kafka help shows canonical default input and explicit reconstruction debug assert.ok(result.payload.commands.some((command: string) => /regenerate hwlab.*--input-topic agentrun\.event\.v1/u.test(command))); assert.ok(result.payload.commands.some((command: string) => /regenerate hwlab.*--input-topic agentrun\.event\.debug\.v1/u.test(command))); assert.ok(result.payload.commands.some((command: string) => /render trace.*--input-topic hwlab\.event\.v1/u.test(command))); + assert.ok(result.payload.commands.some((command: string) => /inspect order.*--session-id ses_/u.test(command))); +}); + +test("Kafka session order inspection discovers a retained session and proves source through DOM order", async () => { + const firstTrace = "trc_order_first"; + const secondTrace = "trc_order_second"; + const agentrun = [ + agentrunOrderRecord(1, firstTrace, "user_message", { userMessageId: "msg_order_first_user", text: "first input" }), + agentrunOrderRecord(2, firstTrace, "backend_status", { phase: "running", message: "first running" }), + agentrunOrderRecord(3, firstTrace, "assistant_message", { text: "first reply", replyAuthority: true, final: true }), + agentrunOrderRecord(4, firstTrace, "terminal_status", { terminalStatus: "completed", message: "completed" }), + agentrunOrderRecord(5, secondTrace, "user_message", { userMessageId: "msg_order_second_user", text: "second input" }), + agentrunOrderRecord(6, secondTrace, "backend_status", { phase: "running", message: "second running" }) + ]; + const hwlab = agentrun.map((record, index) => hwlabOrderRecord(record, index)); + const reads: any[] = []; + const result = await runKafkaCli([ + "inspect", "order", + "--contains-text", "second input", + "--group-prefix", "hwlab-v03-workbench-isolated-debug-order-test", + "--json" + ], { + env: {}, + now: () => "2026-07-11T12:00:00.000Z", + async readKafka(input) { + reads.push(input); + return completeKafkaRead(input.stream === "agentrun" ? agentrun : hwlab); + } + }); + + assert.equal(result.exitCode, 0); + assert.equal(result.payload.scope.resolvedSessionId, HWLAB_SESSION_ID); + assert.equal(result.payload.scope.discovery.containsTextChars, 12); + assert.equal(result.payload.scope.discovery.matchedEventCount, 1); + assert.equal(result.payload.scope.discovery.candidates[0].offset, hwlab[4].offset); + assert.equal(reads.some((input) => input.stream === "agentrun" && input.sessionId === HWLAB_SESSION_ID), true); + assert.equal(result.payload.input.refreshHandoff.counts.replayed, 6); + assert.equal(result.payload.order.eventCount, 6); + assert.equal(result.payload.order.firstDivergence, null); + assert.deepEqual(result.payload.order.turnPairs.map((pair: any) => [pair.traceId, pair.dom.user, pair.dom.agent]), [ + [firstTrace, 0, 1], + [secondTrace, 2, 3] + ]); + assert.deepEqual(result.payload.projection.finalMessageOrder.map((row: any) => [row.role, row.traceId]), [ + ["user", firstTrace], + ["agent", firstTrace], + ["user", secondTrace], + ["agent", secondTrace] + ]); + assert.equal(result.payload.validation.mapperOneToOne, true); + assert.equal(result.payload.validation.sourceOrderPreserved, true); + assert.equal(result.payload.validation.sseOrderPreserved, true); + assert.equal(result.payload.validation.reducerOrderPreserved, true); + assert.equal(result.payload.validation.domTurnOrderPreserved, true); + assert.equal(result.payload.validation.noFirstDivergence, true); +}); + +test("Kafka session order inspection identifies AgentRun as the first layer when user_message is committed late", async () => { + const traceId = "trc_order_late_user"; + const agentrun = [ + agentrunOrderRecord(1, traceId, "backend_status", { phase: "running", message: "running before input" }), + agentrunOrderRecord(2, traceId, "user_message", { userMessageId: "msg_order_late_user", text: "late input" }), + agentrunOrderRecord(3, traceId, "assistant_message", { text: "reply", replyAuthority: true, final: true }), + agentrunOrderRecord(4, traceId, "terminal_status", { terminalStatus: "completed", message: "completed" }) + ]; + const hwlab = agentrun.map((record, index) => hwlabOrderRecord(record, index)); + const result = await runKafkaCli([ + "inspect", "order", + "--session-id", HWLAB_SESSION_ID, + "--group-prefix", "hwlab-v03-workbench-isolated-debug-order-test", + "--json" + ], { + env: {}, + now: () => "2026-07-11T12:00:00.000Z", + async readKafka(input) { + return completeKafkaRead(input.stream === "agentrun" ? agentrun : hwlab); + } + }); + + assert.equal(result.exitCode, 0); + assert.equal(result.payload.order.firstDivergence.layer, "agentrun.event.v1"); + assert.equal(result.payload.order.firstDivergence.code, "user-after-agent-source-event"); + assert.deepEqual(result.payload.order.firstDivergence.evidence.agentrun, { user: 1, agent: 0 }); + assert.deepEqual(result.payload.order.firstDivergence.evidence.dom, { user: 1, agent: 0 }); + assert.equal(result.payload.validation.domTurnOrderPreserved, false); + assert.equal(result.payload.validation.noFirstDivergence, false); }); test("offline HWLAB JSONL uses the shared Web pipeline and renders final response only outside Trace", async () => { @@ -1033,6 +1119,37 @@ test("Kafka query prefers end-offset when the final record also reaches the even assert.equal(queried.reachedEndOffsets, true); }); +test("Kafka query session filter accepts canonical hwlabSessionId across AgentRun and HWLAB topics", async () => { + const matching = canonicalRecord(1).value; + const queried = await queryKafkaEventStream({ + env: { HWLAB_KAFKA_BOOTSTRAP_SERVERS: "kafka.test:9092", HWLAB_KAFKA_CLIENT_ID: "hwlab-test" }, + topic: "agentrun.event.v1", + sessionId: HWLAB_SESSION_ID, + limit: 10, + groupIdPrefix: "hwlab-v03-workbench-isolated-debug", + kafkaFactory() { + return { + admin() { return { async connect() {}, async fetchTopicOffsets() { return [{ partition: 0, low: "0", high: "1" }]; }, async disconnect() {} }; }, + consumer() { + return { + async connect() {}, + async subscribe() {}, + async run({ eachMessage }: any) { + await eachMessage({ topic: "agentrun.event.v1", partition: 0, message: { offset: "0", value: Buffer.from(JSON.stringify(matching)) } }); + }, + async stop() {}, + async disconnect() {} + }; + } + }; + } + }); + + assert.equal(queried.matchedCount, 1); + assert.equal(queried.filterRejectedCount, 0); + assert.equal(queried.completionReason, "end-offset"); +}); + test("Kafka query treats a retained empty partition as a completed bounded scan", async () => { let runCalled = false; let disconnected = false; @@ -1362,6 +1479,51 @@ function reconstructionRecord(seq: number): any { return { topic: "agentrun.event.debug.v1", partition: null, offset: null, key: SESSION_ID, value, headers: { "x-trace-id": TRACE_ID } }; } +function agentrunOrderRecord(seq: number, traceId: string, type: string, payload: Record): any { + const record = canonicalRecord(seq); + const eventId = `evt_order_${seq}`; + record.offset = String(1600 + seq); + record.value.eventId = eventId; + record.value.sourceSeq = seq; + record.value.traceId = traceId; + record.value.hwlabSessionId = HWLAB_SESSION_ID; + record.value.run.hwlabSessionId = HWLAB_SESSION_ID; + record.value.event = { + id: eventId, + runId: RUN_ID, + seq, + type, + payload: { + traceId, + sessionId: SESSION_ID, + hwlabSessionId: HWLAB_SESSION_ID, + ...payload + }, + createdAt: `2026-07-11T11:00:${String(seq).padStart(2, "0")}.000Z` + }; + return record; +} + +function hwlabOrderRecord(source: any, index: number): any { + const valueText = JSON.stringify(source.value); + const projected = projectAgentRunKafkaEventToHwlabEvent(source.value, { + producedAt: `2026-07-11T11:01:${String(index).padStart(2, "0")}.000Z`, + sourceTopic: source.topic, + sourcePartition: source.partition, + sourceOffset: source.offset, + sourceKey: source.key, + inputSha256: createSha256(valueText) + }); + return { + topic: "hwlab.event.v1", + partition: 0, + offset: String(1901 + index), + key: HWLAB_SESSION_ID, + valueSha256: createSha256(JSON.stringify(projected)), + value: projected + }; +} + function canonicalRecord(seq: number): any { const eventId = `evt_${seq}`; const value = { diff --git a/tools/src/hwlab-cli/kafka-regenerate.ts b/tools/src/hwlab-cli/kafka-regenerate.ts index d39bee37..fcc83339 100644 --- a/tools/src/hwlab-cli/kafka-regenerate.ts +++ b/tools/src/hwlab-cli/kafka-regenerate.ts @@ -12,6 +12,7 @@ import { DEFAULT_AGENTRUN_EVENT_TOPIC, DEFAULT_HWLAB_DEBUG_EVENT_TOPIC, DEFAULT_HWLAB_EVENT_TOPIC, + projectAgentRunKafkaEventToHwlabEvent, projectAgentRunKafkaMessageToHwlabDebugEvent, queryKafkaEventStream } from "../../../internal/cloud/kafka-event-bridge.ts"; @@ -20,10 +21,12 @@ import { decodeWorkbenchRealtimeEventFrame } from "../../../web/hwlab-cloud-web/ import { reduceWorkbenchRealtimeEvent } from "../../../web/hwlab-cloud-web/src/stores/workbench-event-reducer.ts"; import { projectWorkbenchLiveKafkaMessage, projectWorkbenchLiveKafkaUserMessage, workbenchLiveKafkaAssistantText, workbenchLiveKafkaProjectionTarget } from "../../../web/hwlab-cloud-web/src/stores/workbench-live-kafka-event.ts"; import { planWorkbenchRealtimeApply } from "../../../web/hwlab-cloud-web/src/stores/workbench-realtime-plan.ts"; +import { createWorkbenchServerState, reduceWorkbenchServerState, selectActiveMessages } from "../../../web/hwlab-cloud-web/src/stores/workbench-server-state.ts"; +import { buildWorkbenchTimelineRows, workbenchMessageIdentity } from "../../../web/hwlab-cloud-web/src/stores/workbench-timeline-model.ts"; import { renderTraceRowsMarkdown, traceDisplayRows } from "./trace-renderer.ts"; const CLI_NAME = "hwlab-cli"; -const VERSION = "0.3.5-kafka-trace-layer-duplicates"; +const VERSION = "0.3.6-kafka-session-order"; const DEFAULT_LIMIT = 500; const DEFAULT_TIMEOUT_MS = 5000; const DEBUG_OUTPUT_PARTITION = 0; @@ -69,6 +72,15 @@ export async function runKafkaCli(argv: string[], options: KafkaCliOptions = {}) }); return { ...result(rendered.payload.status === "partial" ? 2 : 0, rendered.payload, now), markdownOutput: rendered.markdown }; } + if (command === "inspect" && resource === "order") { + action = "kafka.inspect.order"; + const payload = await inspectKafkaSessionOrder(parsed, { + env: options.env ?? process.env, + now, + readKafka: options.readKafka ?? defaultKafkaReader + }); + return result(payload.status === "partial" ? 2 : 0, payload, now); + } if (command !== "regenerate" || resource !== "hwlab") throw cliError("unsupported_kafka_command", "supported commands: kafka regenerate hwlab | kafka render trace", { command, resource }); action = "kafka.regenerate.hwlab"; const payload = await regenerateHwlabDebugEvents(parsed, { @@ -673,11 +685,507 @@ export async function renderHwlabKafkaTrace(parsed: ParsedArgs, dependencies: Re }; } +export async function inspectKafkaSessionOrder(parsed: ParsedArgs, dependencies: { + env: EnvLike; + now: () => string; + readKafka: (input: Record) => Promise>; +}) { + const requestedSessionId = optionalSessionId(parsed.sessionId); + const containsText = text(parsed.containsText); + if (!requestedSessionId && !containsText) { + throw cliError("order_scope_required", "kafka inspect order requires --session-id or --contains-text", { fields: ["sessionId", "containsText"] }); + } + if (containsText.length > 500) throw cliError("order_contains_text_too_long", "containsText must not exceed 500 characters", { length: containsText.length }); + + const agentrunTopic = resolveConfig(parsed.agentrunTopic, dependencies.env.HWLAB_KAFKA_AGENTRUN_EVENT_TOPIC, DEFAULT_AGENTRUN_EVENT_TOPIC, "--agentrun-topic", "env:HWLAB_KAFKA_AGENTRUN_EVENT_TOPIC"); + const hwlabTopic = resolveConfig(parsed.hwlabTopic, dependencies.env.HWLAB_KAFKA_EVENT_TOPIC, DEFAULT_HWLAB_EVENT_TOPIC, "--hwlab-topic", "env:HWLAB_KAFKA_EVENT_TOPIC"); + const group = resolveKafkaGroupConfig(parsed.groupPrefix, dependencies.env.HWLAB_KAFKA_HWLAB_DEBUG_GROUP_PREFIX, true); + const groupPrefix = text(group.value); + assertDebugGroup(groupPrefix); + const limit = boundedInteger(parsed.limit, DEFAULT_LIMIT, 1, 5000, "limit"); + const timeoutMs = boundedInteger(parsed.timeoutMs, DEFAULT_TIMEOUT_MS, 250, 60000, "timeoutMs"); + const rowLimit = boundedInteger(parsed.rowLimit, 80, 1, 500, "rowLimit"); + + let sessionId = requestedSessionId; + let discovery: Record | null = null; + if (!sessionId) { + const discoveredRead = await dependencies.readKafka({ + env: dependencies.env, + stream: "hwlab", + topic: hwlabTopic.value, + limit, + timeoutMs, + fromBeginning: true, + groupIdPrefix: `${groupPrefix}-order-discovery` + }); + if (!kafkaReadComplete(discoveredRead)) return partialKafkaOrderResult({ + sessionId: null, + containsText, + agentrunTopic, + hwlabTopic, + group, + reads: { discovery: discoveredRead } + }); + const candidates = discoverKafkaOrderScopes(discoveredRead.events, containsText); + const sessions = uniqueText(candidates.flatMap((candidate) => candidate.sessionIds)); + discovery = { + containsTextChars: containsText.length, + containsTextSha256: sha256(containsText), + matchedEventCount: candidates.length, + candidateSessionIds: sessions, + candidates: candidates.slice(0, 20), + valuesPrinted: false + }; + if (sessions.length === 0) throw cliError("order_scope_not_found", "no retained HWLAB event matched --contains-text", { + containsTextChars: containsText.length, + containsTextSha256: sha256(containsText), + scannedCount: discoveredRead.scannedCount ?? null + }); + if (sessions.length !== 1) throw cliError("order_scope_ambiguous", "--contains-text matched more than one HWLAB session", { + candidateSessionIds: sessions, + matchedEventCount: candidates.length, + valuesPrinted: false + }); + sessionId = sessions[0]; + } + if (!sessionId) throw cliError("order_scope_not_found", "Kafka order inspection did not resolve a sessionId"); + + const [agentrunRead, hwlabRead] = await Promise.all([ + dependencies.readKafka({ + env: dependencies.env, + stream: "agentrun", + topic: agentrunTopic.value, + sessionId, + limit, + timeoutMs, + fromBeginning: true, + groupIdPrefix: `${groupPrefix}-order-agentrun` + }), + dependencies.readKafka({ + env: dependencies.env, + stream: "hwlab", + topic: hwlabTopic.value, + sessionId, + limit, + timeoutMs, + fromBeginning: true, + groupIdPrefix: `${groupPrefix}-order-hwlab` + }) + ]); + if (!kafkaReadComplete(agentrunRead) || !kafkaReadComplete(hwlabRead)) return partialKafkaOrderResult({ + sessionId, + containsText, + agentrunTopic, + hwlabTopic, + group, + reads: { agentrun: agentrunRead, hwlab: hwlabRead }, + discovery + }); + + const refreshHandoff = await replayHwlabKafkaScopeThroughProductionHandoff({ + traceId: null, + sessionId, + query: { stream: "hwlab", topic: hwlabTopic.value, sessionId }, + readKafka: async () => hwlabRead + }); + if (refreshHandoff.status !== "succeeded") { + throw cliError("order_refresh_handoff_failed", "production Kafka refresh handoff rejected the retained session order", { + sessionId, + refreshHandoff: refreshHandoff.evidence, + valuesPrinted: false + }); + } + + const agentrunRecords = Array.isArray(agentrunRead.events) ? agentrunRead.events : []; + const hwlabRecords = refreshHandoff.records; + const projected = projectHwlabSessionRecordsThroughWorkbench(hwlabRecords, { sessionId, observedAt: dependencies.now() }); + const agentrunByEventId = new Map(); + agentrunRecords.forEach((record, index) => { + const value = valueFromRecord(record); + const eventId = text(value.eventId ?? value.event?.id); + if (eventId) agentrunByEventId.set(eventId, { index, record, value }); + }); + const orderRows = projected.rows.map((row, index) => { + const hwlabRecord = hwlabRecords[index]; + const hwlabValue = valueFromRecord(hwlabRecord); + const sourceEventId = text(hwlabValue.sourceEventId ?? hwlabValue.event?.sourceEventId); + const source = agentrunByEventId.get(sourceEventId) ?? null; + const expected = source ? projectAgentRunKafkaEventToHwlabEvent(source.value, { + sourceTopic: source.record.topic, + sourcePartition: source.record.partition, + sourceOffset: source.record.offset, + sourceKey: source.record.key, + inputSha256: source.record.valueSha256 + }) : null; + return { + eventId: text(hwlabValue.eventId) || null, + sourceEventId: sourceEventId || null, + traceId: row.traceId, + targetTraceId: row.targetTraceId, + sourceSeq: row.sourceSeq, + kind: row.kind, + role: row.role, + agentrun: source ? { index: source.index, topic: source.record.topic, partition: source.record.partition, offset: source.record.offset } : null, + mapper: { + matched: Boolean(expected + && text(expected.eventId) === text(hwlabValue.eventId) + && text(expected.sourceEventId) === sourceEventId + && Number(expected.event?.sourceSeq) === Number(hwlabValue.event?.sourceSeq) + && text(expected.event?.type) === text(hwlabValue.event?.type)), + sourceTopic: text(hwlabValue.sourceEvent?.topic) || null, + sourcePartition: nullableIntegerOrNull(hwlabValue.sourceEvent?.partition), + sourceOffset: text(hwlabValue.sourceEvent?.offset) || null + }, + hwlab: { index, topic: hwlabRecord.topic, partition: hwlabRecord.partition, offset: hwlabRecord.offset }, + sse: { index: row.sseIndex, delivery: "replay" }, + reducer: { index: row.reducerIndex, applied: row.applied }, + web: { + messageId: row.messageId, + messageIndexAtApply: row.messageIndexAtApply, + finalMessageIndex: row.finalMessageIndex, + domRowIndexAtApply: row.domRowIndexAtApply, + finalDomRowIndex: row.finalDomRowIndex + }, + valuesPrinted: false + }; + }); + const turnPairs = kafkaOrderTurnPairs(orderRows); + const firstDivergence = firstKafkaOrderDivergence(turnPairs); + const mappedAgentRunIndices = orderRows.map((row) => row.agentrun?.index).filter((value): value is number => Number.isInteger(value)); + const validation = { + agentrunScanComplete: true, + hwlabScanComplete: true, + refreshHandoffComplete: true, + mapperOneToOne: orderRows.length === agentrunByEventId.size && orderRows.every((row) => row.agentrun && row.mapper.matched), + sourceOrderPreserved: strictlyIncreasing(mappedAgentRunIndices), + sseOrderPreserved: orderRows.every((row, index) => row.hwlab.index === index && row.sse.index === index), + allEventsApplied: projected.rejectedCount === 0 && projected.appliedCount === orderRows.length, + reducerOrderPreserved: orderRows.every((row, index) => row.reducer.applied && row.reducer.index === index), + domTurnOrderPreserved: turnPairs.every((pair) => pair.dom.user < pair.dom.agent), + noFirstDivergence: firstDivergence === null, + valuesPrinted: false + }; + const returnedRows = orderRows.slice(-rowLimit); + return { + ok: true, + action: "kafka.inspect.order", + status: "succeeded", + scope: { + requestedSessionId: requestedSessionId || null, + resolvedSessionId: sessionId, + discovery, + valuesPrinted: false + }, + input: { + agentrun: kafkaOrderReadEvidence(agentrunRead), + hwlab: kafkaOrderReadEvidence(hwlabRead), + refreshHandoff: refreshHandoff.evidence, + valuesPrinted: false + }, + order: { + pipeline: ["agentrun.event.v1", "HWLAB direct mapper", "hwlab.event.v1", "Kafka refresh handoff", "SSE frame", "Web decode/queue/reducer", "conversation DOM"], + eventCount: orderRows.length, + rowsReturned: returnedRows.length, + rowsOmitted: Math.max(0, orderRows.length - returnedRows.length), + rowWindow: orderRows.length > rowLimit ? "tail" : "full", + rows: returnedRows, + turnPairs, + firstDivergence, + valuesPrinted: false + }, + projection: { + decodedCount: projected.decodedCount, + plannedCount: projected.plannedCount, + appliedCount: projected.appliedCount, + rejectedCount: projected.rejectedCount, + finalMessageOrder: projected.messages.map((message: any, index: number) => ({ + index, + role: message.role, + traceId: text(message.traceId ?? message.runnerTrace?.traceId) || null, + messageId: text(message.messageId ?? message.id) || null + })), + finalDomOrder: projected.timelineRows.filter((row: any) => row.message).map((row: any, index: number) => ({ + index, + type: row.type, + role: row.role ?? null, + traceId: row.traceId, + messageId: row.identity + })), + valuesPrinted: false + }, + validation, + runtimeDependencies: { cloudApi: false, database: false, transactionalProjector: false, kafka: true }, + next: { + command: `hwlab-cli kafka inspect order --session-id ${sessionId} --group-prefix ${groupPrefix} --json`, + reason: firstDivergence + ? `首个顺序分歧位于 ${firstDivergence.layer};使用同一 session 和稳定身份继续定点修复。` + : "同一 session 在 source、mapper、SSE、reducer 与 DOM 的顺序一致。" + }, + valuesPrinted: false + }; +} + +export function projectHwlabSessionRecordsThroughWorkbench(records: JsonRecord[], options: { sessionId: string; observedAt: string }) { + let state = createWorkbenchServerState(); + let decodedCount = 0; + let plannedCount = 0; + let appliedCount = 0; + let rejectedCount = 0; + const rows: any[] = []; + const observedBase = Date.parse(options.observedAt); + const baseMs = Number.isFinite(observedBase) ? observedBase : 0; + for (const [index, record] of records.entries()) { + const envelope = valueFromRecord(record); + const decoded = decodeWorkbenchRealtimeEventFrame(JSON.stringify(envelope)); + if (!decoded.payload) { + rejectedCount += 1; + rows.push(rejectedKafkaOrderProjectionRow(index, envelope, `decode-${decoded.status}`)); + continue; + } + decodedCount += 1; + const reduced = reduceWorkbenchRealtimeEvent(decoded.payload as any, "hwlab.event.v1"); + const plan = planWorkbenchRealtimeApply(reduced.action); + const traceStep = plan.steps.find((step) => step.type === "apply-trace-event"); + if (!traceStep || traceStep.type !== "apply-trace-event" || !traceStep.event) { + rejectedCount += 1; + rows.push(rejectedKafkaOrderProjectionRow(index, envelope, reduced.action.type === "ignore" ? reduced.action.reason : "trace-step-missing")); + continue; + } + plannedCount += 1; + const traceId = text(envelope.traceId ?? traceStep.event.traceId); + const sessionId = text(envelope.hwlabSessionId ?? envelope.sessionId ?? traceStep.event.sessionId); + if (!traceId || sessionId !== options.sessionId) { + rejectedCount += 1; + rows.push(rejectedKafkaOrderProjectionRow(index, envelope, !traceId ? "trace-missing" : "session-mismatch")); + continue; + } + const event = traceStep.event as any; + const receivedAt = new Date(baseMs + index).toISOString(); + const role = workbenchLiveKafkaProjectionTarget(event) === "user" ? "user" : "agent"; + let message: any = null; + if (role === "user") { + const messageId = text(event.userMessageId ?? event.messageId); + const previous = selectActiveMessages(state, sessionId).find((item) => text(item.messageId ?? item.id) === messageId) ?? null; + message = projectWorkbenchLiveKafkaUserMessage({ previous, traceId, sessionId, event, receivedAt }); + } else { + const previous = selectActiveMessages(state, sessionId).find((item) => item.role === "agent" && text(item.traceId ?? item.runnerTrace?.traceId) === traceId) ?? null; + message = projectWorkbenchLiveKafkaMessage({ previous, traceId, sessionId, event, receivedAt }); + } + if (!message) { + rejectedCount += 1; + rows.push(rejectedKafkaOrderProjectionRow(index, envelope, "projection-invalid")); + continue; + } + state = reduceWorkbenchServerState(state, { type: "message.upsert", sessionId, message }); + appliedCount += 1; + const messages = selectActiveMessages(state, sessionId); + const messageId = text(message.messageId ?? message.id); + const timelineRows = buildWorkbenchTimelineRows(messages); + rows.push({ + sseIndex: index, + reducerIndex: index, + applied: true, + traceId, + targetTraceId: text(event.targetTraceId) || null, + sourceSeq: Number.isInteger(Number(event.sourceSeq)) ? Number(event.sourceSeq) : null, + kind: text(event.type ?? event.eventType) || "unknown", + role, + messageId, + messageIndexAtApply: messages.findIndex((item) => text(item.messageId ?? item.id) === messageId), + domRowIndexAtApply: timelineRows.findIndex((row) => row.message && workbenchMessageIdentity(row.message) === messageId), + finalMessageIndex: null, + finalDomRowIndex: null + }); + } + const messages = selectActiveMessages(state, options.sessionId); + const timelineRows = buildWorkbenchTimelineRows(messages); + const finalMessageIndex = new Map(messages.map((message, index) => [text(message.messageId ?? message.id), index])); + const finalDomIndex = new Map(timelineRows.filter((row) => row.message).map((row, index) => [workbenchMessageIdentity(row.message as any), index])); + for (const row of rows) { + row.finalMessageIndex = finalMessageIndex.get(row.messageId) ?? -1; + row.finalDomRowIndex = finalDomIndex.get(row.messageId) ?? -1; + } + return { rows, state, messages, timelineRows, decodedCount, plannedCount, appliedCount, rejectedCount }; +} + +function discoverKafkaOrderScopes(records: unknown, needle: string) { + if (!needle || !Array.isArray(records)) return []; + const normalizedNeedle = needle.toLocaleLowerCase(); + return records.flatMap((record: JsonRecord) => { + const value = valueFromRecord(record); + const event = recordObject(value.event) ?? {}; + const finalResponse = recordObject(event.finalResponse) ?? {}; + const visibleTexts = uniqueText([ + event.text, + event.message, + event.assistantText, + finalResponse.text + ]); + if (!visibleTexts.some((candidate) => candidate.toLocaleLowerCase().includes(normalizedNeedle))) return []; + return [{ + sessionIds: hwlabSessionCandidates(value), + traceIds: traceCandidates(value), + sourceEventId: text(value.sourceEventId ?? event.sourceEventId) || null, + topic: text(record.topic) || null, + partition: nullableIntegerOrNull(record.partition), + offset: text(record.offset) || null, + kind: text(event.type ?? event.eventType) || "unknown", + valuesPrinted: false + }]; + }); +} + +function kafkaOrderTurnPairs(rows: any[]) { + const byTrace = new Map(); + for (const row of rows) { + if (!row.traceId || !["user", "agent"].includes(row.role)) continue; + const pair = byTrace.get(row.traceId) ?? {}; + if (!pair[row.role as "user" | "agent"]) pair[row.role as "user" | "agent"] = row; + byTrace.set(row.traceId, pair); + } + return [...byTrace.entries()].flatMap(([traceId, pair]) => { + if (!pair.user || !pair.agent) return []; + return [{ + traceId, + userMessageId: pair.user.web.messageId, + agentMessageId: pair.agent.web.messageId, + agentrun: { user: pair.user.agentrun?.index ?? -1, agent: pair.agent.agentrun?.index ?? -1 }, + hwlab: { user: pair.user.hwlab.index, agent: pair.agent.hwlab.index }, + sse: { user: pair.user.sse.index, agent: pair.agent.sse.index }, + reducer: { user: pair.user.reducer.index, agent: pair.agent.reducer.index }, + message: { user: pair.user.web.finalMessageIndex, agent: pair.agent.web.finalMessageIndex }, + dom: { user: pair.user.web.finalDomRowIndex, agent: pair.agent.web.finalDomRowIndex }, + valuesPrinted: false + }]; + }).sort((left, right) => Math.min(left.hwlab.user, left.hwlab.agent) - Math.min(right.hwlab.user, right.hwlab.agent)); +} + +function firstKafkaOrderDivergence(turnPairs: any[]) { + for (const pair of turnPairs) { + if (pair.agentrun.user < 0 || pair.agentrun.agent < 0) return kafkaOrderDivergence(pair, "agentrun.event.v1 -> HWLAB direct mapper", "source-lineage-missing"); + if (pair.agentrun.user > pair.agentrun.agent) return kafkaOrderDivergence(pair, "agentrun.event.v1", "user-after-agent-source-event"); + if (pair.hwlab.user > pair.hwlab.agent) return kafkaOrderDivergence(pair, "HWLAB direct mapper -> hwlab.event.v1", "user-after-agent-hwlab-event"); + if (pair.sse.user > pair.sse.agent) return kafkaOrderDivergence(pair, "Kafka refresh handoff -> SSE frame", "user-after-agent-sse-frame"); + if (pair.reducer.user > pair.reducer.agent) return kafkaOrderDivergence(pair, "Web decode/queue/reducer", "user-after-agent-reducer-apply"); + if (pair.message.user > pair.message.agent) return kafkaOrderDivergence(pair, "Web message reducer", "user-after-agent-message-order"); + if (pair.dom.user > pair.dom.agent) return kafkaOrderDivergence(pair, "conversation DOM", "user-after-agent-dom-row"); + } + return null; +} + +function kafkaOrderDivergence(pair: any, layer: string, code: string) { + return { layer, code, traceId: pair.traceId, evidence: pair, valuesPrinted: false }; +} + +function rejectedKafkaOrderProjectionRow(index: number, envelope: JsonRecord, reason: string) { + const event = recordObject(envelope.event) ?? {}; + return { + sseIndex: index, + reducerIndex: index, + applied: false, + rejection: reason, + traceId: text(envelope.traceId ?? event.traceId) || null, + targetTraceId: text(event.targetTraceId) || null, + sourceSeq: Number.isInteger(Number(event.sourceSeq)) ? Number(event.sourceSeq) : null, + kind: text(event.type ?? event.eventType) || "unknown", + role: workbenchLiveKafkaProjectionTarget(event as any), + messageId: null, + messageIndexAtApply: -1, + domRowIndexAtApply: -1, + finalMessageIndex: -1, + finalDomRowIndex: -1 + }; +} + +function kafkaReadComplete(read: Record | null | undefined) { + return read?.completionReason === "end-offset" + && read?.reachedEndOffsets === true + && read?.completion?.complete === true; +} + +function kafkaOrderReadEvidence(read: Record | null | undefined) { + return { + topic: text(read?.topic) || null, + groupId: text(read?.groupId) || null, + scannedCount: read?.scannedCount ?? 0, + parsedCount: read?.parsedCount ?? 0, + matchedCount: read?.matchedCount ?? (Array.isArray(read?.events) ? read.events.length : 0), + invalidJsonCount: read?.invalidJsonCount ?? 0, + filterRejectedCount: read?.filterRejectedCount ?? 0, + completionReason: text(read?.completionReason) || null, + reachedEndOffsets: read?.reachedEndOffsets === true, + endOffsets: Array.isArray(read?.endOffsets) ? read.endOffsets : [], + lastScannedOffsets: Array.isArray(read?.lastScannedOffsets) ? read.lastScannedOffsets : [], + valuesPrinted: false + }; +} + +function partialKafkaOrderResult(input: { + sessionId: string | null; + containsText: string; + agentrunTopic: Record; + hwlabTopic: Record; + group: Record; + reads: Record>; + discovery?: Record | null; +}) { + const incomplete = Object.entries(input.reads).find(([, read]) => !kafkaReadComplete(read)); + const evidence = Object.fromEntries(Object.entries(input.reads).map(([name, read]) => [name, kafkaOrderReadEvidence(read)])); + return { + ok: false, + action: "kafka.inspect.order", + status: "partial", + error: { + code: "source_scan_incomplete", + message: `Kafka ${incomplete?.[0] ?? "order"} scan ended before its captured topic barrier`, + details: { stream: incomplete?.[0] ?? null, completionReason: incomplete?.[1]?.completionReason ?? null, valuesPrinted: false } + }, + scope: { + resolvedSessionId: input.sessionId, + discovery: input.discovery ?? (input.containsText ? { + containsTextChars: input.containsText.length, + containsTextSha256: sha256(input.containsText), + valuesPrinted: false + } : null), + valuesPrinted: false + }, + config: { agentrunTopic: input.agentrunTopic, hwlabTopic: input.hwlabTopic, groupPrefix: input.group }, + input: { ...evidence, valuesPrinted: false }, + validation: { agentrunScanComplete: kafkaReadComplete(input.reads.agentrun), hwlabScanComplete: kafkaReadComplete(input.reads.hwlab), valuesPrinted: false }, + next: { + command: input.sessionId + ? `hwlab-cli kafka inspect order --session-id ${input.sessionId} --group-prefix ${input.group.value} --json` + : `hwlab-cli kafka inspect order --contains-text --group-prefix ${input.group.value} --json`, + reason: "使用新的隔离 debug group 和足够的 YAML-owned limit/timeout 重试,只有完整 barrier 才能判定顺序。" + }, + valuesPrinted: false + }; +} + +function strictlyIncreasing(values: number[]) { + return values.every((value, index) => index === 0 || value > values[index - 1]!); +} + +function nullableIntegerOrNull(value: unknown) { + if (value === null || value === undefined || value === "") return null; + const parsed = Number(value); + return Number.isInteger(parsed) && parsed >= 0 ? parsed : null; +} + export async function replayHwlabKafkaTraceThroughProductionHandoff(input: { traceId: string; sessionId: string | null; query: Record; readKafka: (query: Record) => Promise>; +}) { + return replayHwlabKafkaScopeThroughProductionHandoff(input); +} + +export async function replayHwlabKafkaScopeThroughProductionHandoff(input: { + traceId: string | null; + sessionId: string | null; + query: Record; + readKafka: (query: Record) => Promise>; }) { let readResult: Record | null = null; let summary: Record | null = null; @@ -1116,6 +1624,8 @@ function kafkaHelp() { action: "kafka.help", status: "succeeded", commands: [ + "hwlab-cli kafka inspect order --session-id ses_... [--agentrun-topic agentrun.event.v1] [--hwlab-topic hwlab.event.v1] --group-prefix hwlab-...-debug [--row-limit 80] [--json]", + "hwlab-cli kafka inspect order --contains-text --group-prefix hwlab-...-debug [--json]", "hwlab-cli kafka render trace --from kafka --trace-id trc_... [--session-id ses_...] [--run-id run_...] [--command-id cmd_...] [--input-topic hwlab.event.v1] [--group-prefix hwlab-...-debug] [--format markdown] [--row-limit 20] [--json]", "hwlab-cli kafka render trace --from jsonl --trace-id trc_... [--run-id run_...] [--command-id cmd_...] --jsonl-file hwlab-events.jsonl [--format markdown] [--output-markdown trace.md] [--row-limit 20] [--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]", @@ -1439,6 +1949,7 @@ export function renderKafkaCliText(payload: Record) { return [ "ok", `action=${payload.action}`, + payload.scope?.resolvedSessionId ? `session=${payload.scope.resolvedSessionId}` : null, payload.input?.scannedCount !== undefined ? `scanned=${payload.input.scannedCount}` : null, payload.input?.parsedCount !== undefined ? `parsed=${payload.input.parsedCount}` : null, payload.input?.queryMatchedCount !== undefined ? `queryMatched=${payload.input.queryMatchedCount}` : null, @@ -1458,6 +1969,8 @@ export function renderKafkaCliText(payload: Record) { payload.output ? `output=${payload.output.count}` : null, payload.output ? `published=${payload.output.publishedCount}` : null, validation.orderPreserved !== undefined ? `order=${validation.orderPreserved ? "preserved" : "failed"}` : null, + validation.noFirstDivergence !== undefined ? `order=${validation.noFirstDivergence ? "preserved" : "diverged"}` : null, + payload.order?.firstDivergence?.layer ? `firstDivergence=${JSON.stringify(payload.order.firstDivergence.layer)}` : null, payload.output?.topic ? `topic=${payload.output.topic}` : null, payload.config?.outputTopic?.source ? `topicSource=${payload.config.outputTopic.source}` : null, payload.config?.groupPrefix?.source ? `groupSource=${payload.config.groupPrefix.source}` : null,