From 5159e11b0be9acae7d6d938f1e4f9f33d09f81c8 Mon Sep 17 00:00:00 2001 From: root Date: Fri, 10 Jul 2026 06:47:31 +0200 Subject: [PATCH] =?UTF-8?q?fix:=20=E8=A1=A5=E5=85=A8=20Kafka=20=E5=B0=BE?= =?UTF-8?q?=E6=B5=81=E5=85=B3=E8=81=94=E6=91=98=E8=A6=81?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- src/common/kafka-events.ts | 38 ++++++++++++++----- src/selftest/cases/35-kafka-durable-outbox.ts | 21 +++++++++- 2 files changed, 48 insertions(+), 11 deletions(-) diff --git a/src/common/kafka-events.ts b/src/common/kafka-events.ts index f018960..1df394e 100644 --- a/src/common/kafka-events.ts +++ b/src/common/kafka-events.ts @@ -217,15 +217,7 @@ function kafkaMessageSummary(topic: string, payload: EachMessagePayload, include keySha256: key ? sha256Hex(key) : null, valueSha256: sha256Hex(raw), valueBytes: Buffer.byteLength(raw, "utf8"), - schema: typeof parsed?.schema === "string" ? parsed.schema : null, - eventType: typeof parsed?.eventType === "string" ? parsed.eventType : null, - producedAt: typeof parsed?.producedAt === "string" ? parsed.producedAt : null, - runId: nestedString(parsed, ["run", "runId"]) ?? nestedString(parsed, ["context", "runId"]), - commandId: nestedString(parsed, ["command", "commandId"]) ?? nestedString(parsed, ["context", "commandId"]) ?? nestedString(parsed, ["event", "commandId"]), - sessionId: nestedString(parsed, ["run", "sessionId"]) ?? nestedString(parsed, ["context", "sessionId"]), - traceId: nestedString(parsed, ["context", "traceId"]) ?? nestedString(parsed, ["traceId"]), - method: nestedString(parsed, ["stdio", "method"]), - direction: nestedString(parsed, ["stdio", "direction"]), + ...kafkaMessageMetadata(parsed), valuesPrinted: includeValue, }; if (includeValue) { @@ -235,6 +227,25 @@ function kafkaMessageSummary(topic: string, payload: EachMessagePayload, include return base; } +export function kafkaMessageMetadata(parsed: JsonRecord | null): JsonRecord { + return { + schema: nestedString(parsed, ["schema"]), + eventType: nestedString(parsed, ["eventType"]), + producedAt: nestedString(parsed, ["producedAt"]), + committedAt: nestedString(parsed, ["committedAt"]), + eventId: nestedString(parsed, ["eventId"]) ?? nestedString(parsed, ["event", "id"]), + outboxSeq: nestedNumber(parsed, ["outboxSeq"]), + sourceSeq: nestedNumber(parsed, ["sourceSeq"]) ?? nestedNumber(parsed, ["event", "seq"]), + frameSeq: nestedNumber(parsed, ["frameSeq"]) ?? nestedNumber(parsed, ["stdio", "frameSeq"]), + runId: nestedString(parsed, ["run", "runId"]) ?? nestedString(parsed, ["context", "runId"]), + commandId: nestedString(parsed, ["command", "commandId"]) ?? nestedString(parsed, ["context", "commandId"]) ?? nestedString(parsed, ["event", "commandId"]), + sessionId: nestedString(parsed, ["run", "sessionId"]) ?? nestedString(parsed, ["context", "sessionId"]), + traceId: nestedString(parsed, ["context", "traceId"]) ?? nestedString(parsed, ["traceId"]), + method: nestedString(parsed, ["stdio", "method"]), + direction: nestedString(parsed, ["stdio", "direction"]), + }; +} + function matchesKafkaFilter(message: JsonRecord, filter: KafkaTailOptions["filter"]): boolean { if (!filter) return true; if (filter.traceId && message.traceId !== filter.traceId) return false; @@ -305,6 +316,15 @@ function nestedString(record: JsonRecord | null, path: string[]): string | null return typeof current === "string" && current.length > 0 ? current : null; } +function nestedNumber(record: JsonRecord | null, path: string[]): number | null { + let current: unknown = record; + for (const key of path) { + if (!current || typeof current !== "object" || Array.isArray(current)) return null; + current = (current as Record)[key]; + } + return typeof current === "number" && Number.isFinite(current) ? current : null; +} + function sortJson(value: JsonValue): JsonValue { if (Array.isArray(value)) return value.map(sortJson); if (!value || typeof value !== "object") return value; diff --git a/src/selftest/cases/35-kafka-durable-outbox.ts b/src/selftest/cases/35-kafka-durable-outbox.ts index 5809fd7..bd90a45 100644 --- a/src/selftest/cases/35-kafka-durable-outbox.ts +++ b/src/selftest/cases/35-kafka-durable-outbox.ts @@ -1,5 +1,5 @@ import assert from "node:assert/strict"; -import { closeAgentRunKafkaProducer, publishAgentRunKafkaMessage, setAgentRunKafkaProducerFactoryForSelfTest, type AgentRunKafkaConfig } from "../../common/kafka-events.js"; +import { closeAgentRunKafkaProducer, kafkaMessageMetadata, publishAgentRunKafkaMessage, rawFrameValue, setAgentRunKafkaProducerFactoryForSelfTest, type AgentRunKafkaConfig } from "../../common/kafka-events.js"; import { AgentRunError } from "../../common/errors.js"; import type { CreateRunInput, JsonRecord } from "../../common/types.js"; import { validateCreateCommand } from "../../common/validation.js"; @@ -107,6 +107,23 @@ const selfTest: SelfTestCase = async () => { assert.equal(published[1]?.schema, "agentrun.event.v1"); assert.equal(published[1]?.eventType, "agentrun.event.committed"); assert.equal(((published[1]?.event as JsonRecord).payload as JsonRecord).phase, "command-created"); + const canonicalMetadata = kafkaMessageMetadata(published[1] ?? null); + assert.equal(canonicalMetadata.eventId, published[1]?.eventId); + assert.equal(canonicalMetadata.outboxSeq, published[1]?.outboxSeq); + assert.equal(canonicalMetadata.sourceSeq, published[1]?.sourceSeq); + assert.equal(canonicalMetadata.frameSeq, null); + const stdioMetadata = kafkaMessageMetadata(rawFrameValue({ + schema: "codex-stdio.raw.v1", + eventType: "codex.stdio.stdout", + source: "selftest", + frameSeq: 7, + context: { runId: run.id, commandId: command.id, sessionId: "ses_kafka_durable", traceId: "trc_durable" }, + stdio: { frameSeq: 7, direction: "stdout", method: "response" }, + rawJson: { ok: true }, + })); + assert.equal(stdioMetadata.frameSeq, 7); + assert.equal(stdioMetadata.runId, run.id); + assert.equal(stdioMetadata.commandId, command.id); const claimedIntent = store.claimRunnerDispatchIntents({ owner: "terminalizer", leaseMs: 60_000, limit: 1 })[0]; assert.ok(claimedIntent); @@ -120,7 +137,7 @@ const selfTest: SelfTestCase = async () => { await assertStaleClaimsCannotOverwrite(); await assertProducerRecreatedAfterDisconnected(); - return { name: "kafka-durable-outbox", tests: ["explicit-enable", "canonical-event", "durable-head-of-line", "atomic-dispatch-terminal", "stale-claim-fencing", "producer-reconnect"] }; + return { name: "kafka-durable-outbox", tests: ["explicit-enable", "canonical-event", "kafka-tail-correlation-summary", "durable-head-of-line", "atomic-dispatch-terminal", "stale-claim-fencing", "producer-reconnect"] }; }; async function assertStaleClaimsCannotOverwrite(): Promise {