Merge pull request #275 from pikasTech/fix/2467-command-trace-authority
Pipelines as Code CI / agentrun-nc01-v02-ci-0ca23b128bc8b577ae40aa284d90f0f9f610be6e Success

[v0.2][P0] 修正 warm-run command trace 归属
This commit is contained in:
Lyon
2026-07-10 16:26:34 +08:00
committed by GitHub
3 changed files with 59 additions and 4 deletions
+1 -1
View File
@@ -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 {
+1 -1
View File
@@ -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;
+57 -2
View File
@@ -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<void> {
const store = new MemoryAgentRunStore({ eventOutbox: { enabled: true, topic: config.agentrunEventTopic, source: config.clientId } });
const run = store.createRun(runInput("ses_stale_claim"));