From 347f82ab08de5f66c6731225d0525c198f3589bb Mon Sep 17 00:00:00 2001 From: root Date: Fri, 10 Jul 2026 10:19:36 +0200 Subject: [PATCH] =?UTF-8?q?fix:=20=E4=BF=AE=E6=AD=A3=20warm-run=20command?= =?UTF-8?q?=20trace=20=E5=BD=92=E5=B1=9E?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- src/backend/codex-stdio.ts | 2 +- src/common/otel-trace.ts | 2 +- src/selftest/cases/35-kafka-durable-outbox.ts | 59 ++++++++++++++++++- 3 files changed, 59 insertions(+), 4 deletions(-) diff --git a/src/backend/codex-stdio.ts b/src/backend/codex-stdio.ts index 28b5684..a40762f 100644 --- a/src/backend/codex-stdio.ts +++ b/src/backend/codex-stdio.ts @@ -455,7 +455,7 @@ function truthy(value: string | undefined): boolean { return new Set(["1", "true", "yes", "on"]).has(String(value ?? "").trim().toLowerCase()); } -function codexStdioKafkaContext(options: CodexStdioTurnOptions): CodexStdioKafkaContext { +export function codexStdioKafkaContext(options: CodexStdioTurnOptions): CodexStdioKafkaContext { const run = options.otelContext?.run; const command = options.otelContext?.command; return { diff --git a/src/common/otel-trace.ts b/src/common/otel-trace.ts index f0fe9e6..9a74734 100644 --- a/src/common/otel-trace.ts +++ b/src/common/otel-trace.ts @@ -10,11 +10,11 @@ export function agentRunBusinessTraceId(run: RunRecord | null | undefined, comma const sessionMetadata = asRecord(run?.sessionRef?.metadata); const commandPayload = asRecord(command?.payload); for (const value of [ + commandPayload?.traceId, traceSink?.traceId, traceSink?.businessTraceId, sessionMetadata?.hwlabTraceId, sessionMetadata?.traceId, - commandPayload?.traceId, ]) { const text = typeof value === "string" ? value.trim() : ""; if (/^trc_[A-Za-z0-9_.:-]+$/u.test(text)) return text; diff --git a/src/selftest/cases/35-kafka-durable-outbox.ts b/src/selftest/cases/35-kafka-durable-outbox.ts index bd90a45..cc85449 100644 --- a/src/selftest/cases/35-kafka-durable-outbox.ts +++ b/src/selftest/cases/35-kafka-durable-outbox.ts @@ -1,8 +1,10 @@ import assert from "node:assert/strict"; 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 { agentRunOtelTraceContext } from "../../common/otel-trace.js"; +import type { CreateRunInput, JsonRecord, KafkaEventOutboxRecord } from "../../common/types.js"; import { validateCreateCommand } from "../../common/validation.js"; +import { codexStdioKafkaContext } from "../../backend/codex-stdio.js"; import { relayKafkaEventOutboxOnce, type KafkaOutboxRelayOptions } from "../../mgr/kafka-outbox-relay.js"; import { startManagerServer } from "../../mgr/server.js"; import { MemoryAgentRunStore, openAgentRunStoreFromEnv } from "../../mgr/store.js"; @@ -124,6 +126,7 @@ const selfTest: SelfTestCase = async () => { assert.equal(stdioMetadata.frameSeq, 7); assert.equal(stdioMetadata.runId, run.id); assert.equal(stdioMetadata.commandId, command.id); + assertWarmRunCommandTraceAuthority(); const claimedIntent = store.claimRunnerDispatchIntents({ owner: "terminalizer", leaseMs: 60_000, limit: 1 })[0]; assert.ok(claimedIntent); @@ -137,9 +140,61 @@ const selfTest: SelfTestCase = async () => { await assertStaleClaimsCannotOverwrite(); await assertProducerRecreatedAfterDisconnected(); - 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"] }; + return { name: "kafka-durable-outbox", tests: ["explicit-enable", "canonical-event", "warm-run-command-trace-authority", "run-level-trace-fallback", "stdio-command-trace-context", "kafka-tail-correlation-summary", "durable-head-of-line", "atomic-dispatch-terminal", "stale-claim-fencing", "producer-reconnect"] }; }; +function assertWarmRunCommandTraceAuthority(): void { + const store = new MemoryAgentRunStore({ eventOutbox: { enabled: true, topic: config.agentrunEventTopic, source: config.clientId } }); + const run = store.createRun({ + ...runInput("ses_warm_trace"), + traceSink: { kind: "hwlab", traceId: "trc_run_original" }, + }); + const firstCommand = store.createCommand(run.id, validateCreateCommand({ + type: "turn", + idempotencyKey: "warm-trace-first", + payload: { prompt: "first warm turn", traceId: "trc_warm_first" }, + })); + const secondCommand = store.createCommand(run.id, validateCreateCommand({ + type: "turn", + idempotencyKey: "warm-trace-second", + payload: { prompt: "second warm turn", traceId: "trc_warm_second" }, + })); + + const outbox: KafkaEventOutboxRecord[] = []; + while (true) { + const item = store.claimKafkaEventOutbox({ owner: "warm-trace-test", leaseMs: 60_000, limit: 1 })[0]; + if (!item) break; + outbox.push(item); + store.completeKafkaEventOutbox(item); + } + const runOnlyEnvelope = outbox.find((item) => item.value.commandId === null)?.value; + const firstEnvelope = outbox.find((item) => item.value.commandId === firstCommand.id)?.value; + const secondEnvelope = outbox.find((item) => item.value.commandId === secondCommand.id)?.value; + assert.equal(runOnlyEnvelope?.traceId, "trc_run_original"); + assert.equal(firstEnvelope?.traceId, "trc_warm_first"); + assert.equal(secondEnvelope?.traceId, "trc_warm_second"); + + const runOnlyOtel = agentRunOtelTraceContext(run, null); + const firstOtel = agentRunOtelTraceContext(run, firstCommand); + const secondOtel = agentRunOtelTraceContext(run, secondCommand); + assert.equal(runOnlyOtel.businessTraceId, "trc_run_original"); + assert.equal(firstOtel.businessTraceId, "trc_warm_first"); + assert.equal(secondOtel.businessTraceId, "trc_warm_second"); + assert.notEqual(firstOtel.traceId, secondOtel.traceId); + assert.notEqual(runOnlyOtel.traceId, secondOtel.traceId); + + const stdioContext = codexStdioKafkaContext({ + prompt: "second warm turn", + cwd: ".", + approvalPolicy: "never", + sandbox: "workspace-write", + timeoutMs: 30_000, + otelContext: { run, command: secondCommand }, + }); + assert.equal(stdioContext.traceId, "trc_warm_second"); + assert.equal(stdioContext.commandId, secondCommand.id); +} + async function assertStaleClaimsCannotOverwrite(): Promise { const store = new MemoryAgentRunStore({ eventOutbox: { enabled: true, topic: config.agentrunEventTopic, source: config.clientId } }); const run = store.createRun(runInput("ses_stale_claim"));