fix: 补全 Kafka 尾流关联摘要

This commit is contained in:
root
2026-07-10 06:47:31 +02:00
parent d25733a8e0
commit 5159e11b0b
2 changed files with 48 additions and 11 deletions
+29 -9
View File
@@ -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<string, unknown>)[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;
+19 -2
View File
@@ -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<void> {