Merge pull request #267 from pikasTech/fix/265-durable-dispatch-outbox
Pipelines as Code CI / agentrun-nc01-v02-ci-02adb14d56b8e0c84fe9d90c2716a12f1a674039 Success
Pipelines as Code CI / agentrun-nc01-v02-ci-02adb14d56b8e0c84fe9d90c2716a12f1a674039 Success
补全 Kafka 尾流的 durable 关联摘要
This commit is contained in:
@@ -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;
|
||||
|
||||
@@ -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> {
|
||||
|
||||
Reference in New Issue
Block a user