Merge pull request #2505 from pikasTech/fix/2498-kafka-order-debug
Pipelines as Code CI / hwlab-nc01-v03-ci-poll- Success
Pipelines as Code CI / hwlab-nc01-v03-ci-poll- Success
feat: 增加 Kafka 会话顺序单步诊断
This commit is contained in:
@@ -8,6 +8,19 @@
|
||||
- group prefix 必须标识独立 debug group,不复用产品 consumer group;
|
||||
- 命令只读取 Kafka,不发布事件,不启动 Cloud API,不访问数据库,也不依赖 projector。
|
||||
|
||||
- `hwlab-cli kafka inspect order` 用于定位同一 session 在纯 Kafka 刷新链路中的首个顺序分歧:
|
||||
- 已知 session 时使用 `--session-id <sessionId>`;
|
||||
- 只知道用户输入时使用 `--contains-text <bounded-user-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` 分类;
|
||||
|
||||
@@ -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);
|
||||
|
||||
@@ -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<string, unknown>): 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 = {
|
||||
|
||||
@@ -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<string, any>) => Promise<Record<string, any>>;
|
||||
}) {
|
||||
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<string, any> | 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<string, { index: number; record: JsonRecord; value: JsonRecord }>();
|
||||
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<string, { user?: any; agent?: any }>();
|
||||
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<string, any> | null | undefined) {
|
||||
return read?.completionReason === "end-offset"
|
||||
&& read?.reachedEndOffsets === true
|
||||
&& read?.completion?.complete === true;
|
||||
}
|
||||
|
||||
function kafkaOrderReadEvidence(read: Record<string, any> | 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<string, any>;
|
||||
hwlabTopic: Record<string, any>;
|
||||
group: Record<string, any>;
|
||||
reads: Record<string, Record<string, any>>;
|
||||
discovery?: Record<string, any> | 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 <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<string, any>;
|
||||
readKafka: (query: Record<string, any>) => Promise<Record<string, any>>;
|
||||
}) {
|
||||
return replayHwlabKafkaScopeThroughProductionHandoff(input);
|
||||
}
|
||||
|
||||
export async function replayHwlabKafkaScopeThroughProductionHandoff(input: {
|
||||
traceId: string | null;
|
||||
sessionId: string | null;
|
||||
query: Record<string, any>;
|
||||
readKafka: (query: Record<string, any>) => Promise<Record<string, any>>;
|
||||
}) {
|
||||
let readResult: Record<string, any> | null = null;
|
||||
let summary: Record<string, any> | 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 <bounded-user-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<string, any>) {
|
||||
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<string, any>) {
|
||||
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,
|
||||
|
||||
Reference in New Issue
Block a user