From cd62fa61fbc970b6f35db65af3404330b87dd43f Mon Sep 17 00:00:00 2001 From: Bulat Yapparov Date: Mon, 14 Sep 2026 12:20:07 +0100 Subject: [PATCH 1/5] feat(cli): expose execution retry measurements --- CHANGELOG.md | 1 + EVENTS.md | 65 +++++- packages/cli/src/cli/cmd/run.invocation.ts | 2 + packages/cli/src/cli/cmd/run.ts | 90 ++++++- packages/cli/src/session/processor.ts | 6 + packages/cli/src/session/retry.ts | 30 +++ packages/cli/src/session/status.ts | 6 + .../cli/test/cli/invocation-lifecycle.test.ts | 24 +- .../cli/test/cli/run-retry-telemetry.test.ts | 220 ++++++++++++++++++ packages/cli/test/session/retry.test.ts | 24 +- 10 files changed, 446 insertions(+), 22 deletions(-) create mode 100644 packages/cli/test/cli/run-retry-telemetry.test.ts diff --git a/CHANGELOG.md b/CHANGELOG.md index 8cc2381..1d8899e 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -4,6 +4,7 @@ ### Features +- **Headless execution measurement inputs** — NDJSON events now carry CLI release identity and correlated retry lifecycle events with bounded reason and outcome categories. Existing retry behavior is unchanged. (#110) - **GPT-5.6 Codex models** — Added OpenAI's Sol, Terra, and Luna models with API and subscription-backed reasoning effort variants, including the Codex-only `ultra` alias for Sol and Terra. - **GLM-5.3 model support** — Added the latest Z.AI Coding Plan model with its 1M-token context window and native `low`, `high`, and `max` reasoning efforts. diff --git a/EVENTS.md b/EVENTS.md index 911798f..dd97dc7 100644 --- a/EVENTS.md +++ b/EVENTS.md @@ -6,12 +6,14 @@ When the parsed `aictrl run --format json` handler starts, the CLI emits newline { "type": "", "timestamp": 1741500000000, + "schemaVersion": "1", + "cliVersion": "0.4.3", "invocationID": "7d142250-8bdc-43df-99af-efa252db62a7", "sessionID": "session_01abc..." } ``` -`invocationID` is present on every event from `run --format json`. `sessionID` is present only after a real session has been created; invocation events never fabricate one. +`schemaVersion`, `cliVersion`, and `invocationID` are present on every event from `run --format json`. `cliVersion` identifies the emitting CLI release and is independent of the event schema version. `sessionID` is present only after a real session has been created; invocation events never fabricate one. The schema is versioned via `invocation_start.schemaVersion` and `session_start.schemaVersion`. This document describes **schema version `"1"`**. Consumers should pin to this version and treat unknown fields as forward-compatible additions. @@ -28,6 +30,7 @@ Emitted once, before piped stdin is read and before run validation or bootstrap "type": "invocation_start", "timestamp": 1741500000000, "schemaVersion": "1", + "cliVersion": "0.4.3", "invocationID": "7d142250-8bdc-43df-99af-efa252db62a7" } ``` @@ -174,6 +177,66 @@ Emitted immediately before `session_complete` when the session terminates abnorm - `code` (string, optional) — provider HTTP status code, error code, or conventional signal-derived exit code (`130` for `SIGINT`, `143` for `SIGTERM`) when available. - `message` (string, **required**) — human-readable error message. +### `retry_scheduled` + +Emitted when the primary session schedules an existing automatic retry. This event observes the retry policy; it does not add or change retry behavior. + +```json +{ + "type": "retry_scheduled", + "retryID": "96766fef-4f5c-47bb-a292-8fd43761ac3a", + "messageID": "msg_01abc...", + "providerID": "google", + "modelID": "gemini-2.5-flash", + "attempt": 1, + "reason": "rate_limit", + "delayMs": 2000 +} +``` + +- `retryID` (string, **required**) — opaque identity shared with the matching `retry_complete` event. +- `messageID` (string or null, **required**) — assistant message that owns the retry. It is null only when attached to an older server that did not supply the additive correlation fields. +- `providerID` / `modelID` (string or null, **required**) — resolved provider and model when the retry source supplied them. +- `attempt` (number, **required**) — one-based retry ordinal for the assistant message. +- `reason` (string, **required**) — bounded category: `rate_limit`, `timeout`, `network`, `provider`, or `unknown`. Treat this as an open set. Free-text provider errors are not emitted as metric dimensions. +- `delayMs` (number, **required**) — scheduled backoff. For an older attached server this is the non-negative time remaining when the event is observed. + +### `retry_complete` + +Emitted when the scheduled retry is resolved by a terminal assistant message, another scheduled retry, a session error, cancellation, or an idle session boundary. + +```json +{ + "type": "retry_complete", + "retryID": "96766fef-4f5c-47bb-a292-8fd43761ac3a", + "messageID": "msg_01abc...", + "providerID": "google", + "modelID": "gemini-2.5-flash", + "attempt": 1, + "reason": "rate_limit", + "delayMs": 2000, + "outcome": "recovered" +} +``` + +Identity and dimension fields match `retry_scheduled`. `outcome` is `recovered`, `failed`, `aborted`, or `unknown`. A retry superseded by another retry is `failed`; a successful terminal assistant message is `recovered`; cancellation is `aborted`; an idle boundary without an observable terminal message is `unknown`. If the event stream ends before `retry_complete`, the attempt is censored and must remain unknown. + +## Measurement contract + +The NDJSON stream supplies measurement inputs; it does not calculate product-specific results. Aggregate top-level invocations by `invocationID`, and use `sessionID` and `messageID` only for correlation. Do not count child sessions as new invocations. + +| Metric | Numerator | Denominator | Unknown / censored handling | +|---|---|---|---| +| Invocation outcome rate | `invocation_complete` grouped by `status` | distinct `invocation_start.invocationID` | A start without a complete event is unknown, never success | +| Provider-turn error rate | primary `message_complete.status = error` | all primary `message_complete` events | A started invocation without a terminal primary message is unknown | +| Contradiction count | completed invocation whose primary message or session is terminally erroneous | completed invocations | Keep as a data-quality count; the terminal error fix is tracked separately | +| Diagnostic coverage | error turns with the separately defined provider diagnostic | provider-error turns | Raw provider reason availability and redaction are defined separately; do not infer them from `reason` | +| Retry recovery rate | `retry_complete.outcome = recovered` | `retry_complete` with outcome `recovered` or `failed` | Exclude `aborted`, `unknown`, and scheduled retries without a completion event | + +Sum `retry_scheduled.delayMs` once per distinct `retryID` for additional scheduled backoff. Join usage and cost from the owning `message_complete.messageID`; `retry_complete` is a lifecycle correlation event and carries no duplicate usage or cost. Missing usage or cost remains unknown rather than zero. + +Safe metric dimensions are `cliVersion`, `providerID`, `modelID`, and bounded `reason`. `invocationID`, `sessionID`, `messageID`, and `retryID` are high-cardinality correlation attributes. Never use free-text errors as metric labels. + ## Message Events ### `message_complete` diff --git a/packages/cli/src/cli/cmd/run.invocation.ts b/packages/cli/src/cli/cmd/run.invocation.ts index b2df3cb..9a23033 100644 --- a/packages/cli/src/cli/cmd/run.invocation.ts +++ b/packages/cli/src/cli/cmd/run.invocation.ts @@ -1,5 +1,6 @@ import { Stdout } from "../stdout" import { SCHEMA_VERSION } from "./run.errors" +import { Installation } from "../../installation" export type RunInvocationPhase = "validation" | "stdin" | "bootstrap" | "session" @@ -19,6 +20,7 @@ export function createRunInvocation(enabled: boolean) { type, timestamp: Date.now(), schemaVersion: SCHEMA_VERSION, + cliVersion: Installation.VERSION, invocationID: id, ...data, }), diff --git a/packages/cli/src/cli/cmd/run.ts b/packages/cli/src/cli/cmd/run.ts index ca65b00..ed9915c 100644 --- a/packages/cli/src/cli/cmd/run.ts +++ b/packages/cli/src/cli/cmd/run.ts @@ -43,6 +43,7 @@ import { Shutdown } from "../shutdown" import { Stdout } from "../stdout" import { attempt, signals, type Signals } from "../signals" import { createRunInvocation } from "./run.invocation" +import { Installation } from "../../installation" type ToolProps = { input: Tool.InferParameters @@ -527,6 +528,7 @@ export const RunCommand = cmd({ type, timestamp: Date.now(), schemaVersion: SCHEMA_VERSION, + cliVersion: Installation.VERSION, invocationID: invocation.id, sessionID, ...data, @@ -542,6 +544,53 @@ export const RunCommand = cmd({ const childSessions = new Set() const emitted = new Set() const seqBySession = new Map() + type Retry = { + retryID: string + messageID: string | null + providerID: string | null + modelID: string | null + attempt: number + reason: string + delayMs: number + } + const retries = new Map() + + function resolveRetry(sid: string, outcome: "recovered" | "failed" | "aborted" | "unknown") { + const retry = retries.get(sid) + if (!retry) return + retries.delete(sid) + emit("retry_complete", { ...retry, outcome }) + } + + function scheduleRetry( + sid: string, + status: { + attempt: number + next: number + retryID?: string + messageID?: string + providerID?: string + modelID?: string + reason?: string + delayMs?: number + }, + ) { + const current = retries.get(sid) + if (current && current.retryID === status.retryID) return + resolveRetry(sid, "failed") + const retry = { + retryID: status.retryID ?? crypto.randomUUID(), + messageID: status.messageID ?? null, + providerID: status.providerID ?? null, + modelID: status.modelID ?? null, + attempt: status.attempt, + reason: status.reason ?? "unknown", + delayMs: status.delayMs ?? Math.max(0, status.next - Date.now()), + } + retries.set(sid, retry) + emit("retry_scheduled", retry) + } + function nextSeq(sid: string): number { const n = (seqBySession.get(sid) ?? 0) + 1 seqBySession.set(sid, n) @@ -572,6 +621,16 @@ export const RunCommand = cmd({ if (info.sessionID === sessionID && info.time.completed !== undefined && !emitted.has(info.id)) { emitted.add(info.id) const usage = terminalUsage(info) + const status = + info.error?.name === "MessageAbortedError" ? "aborted" : info.error ? "error" : "completed" + resolveRetry( + info.sessionID, + status === "completed" && info.finish !== "error" && info.finish !== "content-filter" + ? "recovered" + : status === "aborted" + ? "aborted" + : "failed", + ) // Context-window utilization: used = input + cache.read + cache.write // (all prompt tokens that occupy the model's context window this turn). @@ -606,7 +665,7 @@ export const RunCommand = cmd({ cost: info.cost, tokens: usage.tokens, usageStatus: usage.usageStatus, - status: info.error?.name === "MessageAbortedError" ? "aborted" : info.error ? "error" : "completed", + status, finish: info.finish, }) } @@ -711,6 +770,10 @@ export const RunCommand = cmd({ if (!control.current) process.exitCode = 1 invocation.error(props.error) const classified = classifySessionError(props.error) + resolveRetry( + props.sessionID, + classified.reason === "interrupted" || classified.reason === "terminated" ? "aborted" : "failed", + ) // Structured session_error is the telemetry/CI channel for the // primary session. The legacy "error" event below is the raw // pass-through for both primary and child-session failures. @@ -760,15 +823,22 @@ export const RunCommand = cmd({ } } - if (event.type === "session.status" && event.properties.status.type === "idle") { - if (event.properties.sessionID === sessionID) { - break + if (event.type === "session.status") { + const status = event.properties.status + if (status.type === "retry" && event.properties.sessionID === sessionID) { + scheduleRetry(event.properties.sessionID, status) } - if (childSessions.has(event.properties.sessionID)) { - emit("subagent_complete", { - subagentSessionID: event.properties.sessionID, - parentSessionID: sessionID, - }) + if (status.type === "idle") { + if (event.properties.sessionID === sessionID) { + resolveRetry(event.properties.sessionID, "unknown") + break + } + if (childSessions.has(event.properties.sessionID)) { + emit("subagent_complete", { + subagentSessionID: event.properties.sessionID, + parentSessionID: sessionID, + }) + } } } @@ -862,6 +932,7 @@ export const RunCommand = cmd({ function interrupt(signal: Signals.Info) { error ??= signal.message invocation.error(signal.message) + resolveRetry(sessionID, "aborted") report(signal.reason, String(signal.code), signal.message) abort() } @@ -878,6 +949,7 @@ export const RunCommand = cmd({ const classified = classifySessionError(cause) error ??= classified.message invocation.error(cause) + resolveRetry(sessionID, "failed") report(classified.reason, classified.code, classified.message) complete(error) if (control.current) { diff --git a/packages/cli/src/session/processor.ts b/packages/cli/src/session/processor.ts index 2577169..9fba1c5 100644 --- a/packages/cli/src/session/processor.ts +++ b/packages/cli/src/session/processor.ts @@ -380,6 +380,12 @@ export namespace SessionProcessor { attempt, message: retry, next: Date.now() + delay, + retryID: crypto.randomUUID(), + messageID: input.assistantMessage.id, + providerID: input.model.providerID, + modelID: input.model.id, + reason: SessionRetry.reason(error), + delayMs: delay, }) await SessionRetry.sleep(delay, input.abort).catch(() => {}) continue diff --git a/packages/cli/src/session/retry.ts b/packages/cli/src/session/retry.ts index e3b20f3..96a96fa 100644 --- a/packages/cli/src/session/retry.ts +++ b/packages/cli/src/session/retry.ts @@ -3,6 +3,8 @@ import { MessageV2 } from "./message-v2" import { iife } from "@/util/iife" export namespace SessionRetry { + export type Reason = "rate_limit" | "timeout" | "network" | "provider" | "unknown" + export const RETRY_INITIAL_DELAY = 2000 export const RETRY_BACKOFF_FACTOR = 2 export const RETRY_MAX_DELAY_NO_HEADERS = 30_000 // 30 seconds @@ -99,4 +101,32 @@ export namespace SessionRetry { return undefined } } + + /** A bounded category for telemetry dimensions. This does not affect retry policy. */ + export function reason(error: ReturnType): Reason { + const data = error.data && typeof error.data === "object" ? error.data : undefined + const status = + data && "statusCode" in data && typeof data.statusCode === "number" + ? data.statusCode + : data && "status" in data && typeof data.status === "number" + ? data.status + : undefined + const detail = (() => { + if (!data) return "" + const values = ["message", "responseBody"].flatMap((key) => { + if (!(key in data)) return [] + const value = data[key as keyof typeof data] + return typeof value === "string" ? [value] : [] + }) + return values.join(" ") + })() + + if (status === 429 || /rate[-_\s]?limit|too[-_\s]many[-_\s]requests/i.test(detail)) return "rate_limit" + if (status === 408 || /\btimeout\b|timed out|time out/i.test(detail)) return "timeout" + if (/network[-_\s]?error|fetch failed|connection|socket|ECONN|ENOTFOUND|EAI_AGAIN/i.test(detail)) return "network" + if ((status !== undefined && status >= 500) || /overloaded|unavailable|resource[-_\s]exhausted/i.test(detail)) { + return "provider" + } + return "unknown" + } } diff --git a/packages/cli/src/session/status.ts b/packages/cli/src/session/status.ts index e524c40..800f7cd 100644 --- a/packages/cli/src/session/status.ts +++ b/packages/cli/src/session/status.ts @@ -14,6 +14,12 @@ export namespace SessionStatus { attempt: z.number(), message: z.string(), next: z.number(), + retryID: z.string().optional(), + messageID: z.string().optional(), + providerID: z.string().optional(), + modelID: z.string().optional(), + reason: z.enum(["rate_limit", "timeout", "network", "provider", "unknown"]).optional(), + delayMs: z.number().optional(), }), z.object({ type: z.literal("busy"), diff --git a/packages/cli/test/cli/invocation-lifecycle.test.ts b/packages/cli/test/cli/invocation-lifecycle.test.ts index 3ae752c..90e30ca 100644 --- a/packages/cli/test/cli/invocation-lifecycle.test.ts +++ b/packages/cli/test/cli/invocation-lifecycle.test.ts @@ -44,6 +44,7 @@ function expectFailure(result: Awaited>, phase: string const [start, error, complete] = result.events expect(start.schemaVersion).toBe("1") + expect(typeof start.cliVersion).toBe("string") expect(typeof start.invocationID).toBe("string") expect(start).not.toHaveProperty("sessionID") expect(error.invocationID).toBe(start.invocationID) @@ -77,6 +78,7 @@ describe("run --format json invocation lifecycle (#90)", () => { type: "invocation_start", schemaVersion: "1", }) + expect(typeof events(new TextDecoder().decode(first.value))[0].cliVersion).toBe("string") }, 15_000) test("reports an invalid directory as a validation error", async () => { @@ -88,18 +90,21 @@ describe("run --format json invocation lifecycle (#90)", () => { test("completes when guarded pre-run I/O rejects", async () => { const source = path.resolve(import.meta.dir, "../../src/cli/cmd/run.invocation.ts") - const proc = Bun.spawn([ - process.execPath, - "--conditions=browser", - "-e", - `import { createRunInvocation } from ${JSON.stringify(source)} + const proc = Bun.spawn( + [ + process.execPath, + "--conditions=browser", + "-e", + `import { createRunInvocation } from ${JSON.stringify(source)} const invocation = createRunInvocation(true) invocation.phase("stdin") await invocation.guard(() => Promise.reject(new Error("stdin failed")))`, - ], { - stdout: "pipe", - stderr: "pipe", - }) + ], + { + stdout: "pipe", + stderr: "pipe", + }, + ) expectFailure(await output(proc), "stdin") }) @@ -196,6 +201,7 @@ await invocation.guard(() => Promise.reject(new Error("stdin failed")))`, .every((event) => event.invocationID === start?.invocationID), ).toBe(true) expect(result.events.every((event) => event.schemaVersion === "1")).toBe(true) + expect(result.events.every((event) => typeof event.cliVersion === "string")).toBe(true) } finally { server.stop(true) } diff --git a/packages/cli/test/cli/run-retry-telemetry.test.ts b/packages/cli/test/cli/run-retry-telemetry.test.ts new file mode 100644 index 0000000..224a729 --- /dev/null +++ b/packages/cli/test/cli/run-retry-telemetry.test.ts @@ -0,0 +1,220 @@ +import path from "path" +import { afterEach, describe, expect, test } from "bun:test" + +const cli = path.resolve(import.meta.dir, "../../src/index.ts") +const models = path.resolve(import.meta.dir, "../tool/fixtures/models-api.json") +const sessionID = "ses_retry_measurement" +const messageID = "msg_retry_measurement" +const servers: Bun.Server[] = [] + +function server(events: unknown[]) { + const server = Bun.serve({ + port: 0, + fetch(req) { + const url = new URL(req.url) + if (req.method === "POST" && url.pathname === "/session") return Response.json({ id: sessionID }) + if (req.method === "GET" && url.pathname === "/config") return Response.json({}) + if (req.method === "POST" && url.pathname.endsWith("/message")) return Response.json({}) + if (req.method === "GET" && url.pathname === "/event") { + return new Response(events.map((event) => `data: ${JSON.stringify(event)}\n\n`).join(""), { + headers: { "content-type": "text/event-stream" }, + }) + } + return Response.json({ error: "not found" }, { status: 404 }) + }, + }) + servers.push(server) + return `http://localhost:${server.port}` +} + +function status(value: Record) { + return { + type: "session.status", + properties: { sessionID, status: value }, + } +} + +function completed(finish = "stop") { + return { + type: "message.updated", + properties: { + info: { + id: messageID, + sessionID, + role: "assistant", + time: { created: 1, completed: 2 }, + parentID: "msg_user", + modelID: "glm-4.7", + providerID: "zai", + agent: "build", + path: { cwd: "/", root: "/" }, + cost: 0, + tokens: { input: 0, output: 0, reasoning: 0, cache: { read: 0, write: 0 } }, + usageStatus: "reported", + finish, + }, + }, + } +} + +afterEach(() => servers.splice(0).map((item) => item.stop(true))) + +describe("run --format json retry telemetry (#110)", () => { + test("correlates retry ordinals and resolves eventual recovery", async () => { + const first = "c76dbc08-64ce-48cc-b79f-ecc1ba16be2c" + const second = "52d9f7cf-d667-4a82-a6c8-d65f75fb965b" + const retry = (retryID: string, attempt: number, reason: string, delayMs: number) => + status({ + type: "retry", + retryID, + messageID, + providerID: "zai", + modelID: "glm-4.7", + attempt, + reason, + delayMs, + message: "redacted from NDJSON", + next: Date.now() + delayMs, + }) + + const proc = Bun.spawn( + [ + "bun", + "run", + cli, + "run", + "--format", + "json", + "--attach", + server([ + retry(first, 1, "rate_limit", 2_000), + retry(second, 2, "provider", 4_000), + completed(), + status({ type: "idle" }), + ]), + "prompt", + ], + { + cwd: process.cwd(), + env: { ...process.env, AICTRL_MODELS_PATH: models }, + stdout: "pipe", + stderr: "pipe", + }, + ) + const [stdout, stderr, code] = await Promise.all([ + new Response(proc.stdout).text(), + new Response(proc.stderr).text(), + proc.exited, + ]) + expect(code, stderr).toBe(0) + const output = stdout + .split("\n") + .filter(Boolean) + .map((line) => JSON.parse(line)) + .filter((event) => event.type === "retry_scheduled" || event.type === "retry_complete") + + expect(output.map((event) => [event.type, event.retryID, event.outcome])).toEqual([ + ["retry_scheduled", first, undefined], + ["retry_complete", first, "failed"], + ["retry_scheduled", second, undefined], + ["retry_complete", second, "recovered"], + ]) + expect(output[0]).toMatchObject({ + messageID, + providerID: "zai", + modelID: "glm-4.7", + attempt: 1, + reason: "rate_limit", + delayMs: 2_000, + }) + expect(JSON.stringify(output)).not.toContain("redacted from NDJSON") + expect(output.every((event) => typeof event.cliVersion === "string")).toBe(true) + }, 20_000) + + test("marks an uncorrelated older-server retry unknown at the idle boundary", async () => { + const proc = Bun.spawn( + [ + "bun", + "run", + cli, + "run", + "--format", + "json", + "--attach", + server([ + status({ type: "retry", attempt: 1, message: "legacy", next: Date.now() + 1_000 }), + status({ type: "idle" }), + ]), + "prompt", + ], + { + cwd: process.cwd(), + env: { ...process.env, AICTRL_MODELS_PATH: models }, + stdout: "pipe", + stderr: "pipe", + }, + ) + const [stdout, stderr, code] = await Promise.all([ + new Response(proc.stdout).text(), + new Response(proc.stderr).text(), + proc.exited, + ]) + expect(code, stderr).toBe(0) + const output = stdout + .split("\n") + .filter(Boolean) + .map((line) => JSON.parse(line)) + .filter((event) => event.type === "retry_scheduled" || event.type === "retry_complete") + + expect(output).toHaveLength(2) + expect(output[0]).toMatchObject({ messageID: null, providerID: null, modelID: null, reason: "unknown" }) + expect(output[1]).toMatchObject({ retryID: output[0].retryID, outcome: "unknown" }) + }, 20_000) + + test("does not count a normalized error finish as recovered", async () => { + const retryID = "3c45ab98-65fe-4efe-9309-4d130341e31c" + const proc = Bun.spawn( + [ + "bun", + "run", + cli, + "run", + "--format", + "json", + "--attach", + server([ + status({ + type: "retry", + retryID, + messageID, + providerID: "zai", + modelID: "glm-4.7", + attempt: 1, + reason: "provider", + delayMs: 2_000, + message: "provider error", + next: Date.now() + 2_000, + }), + completed("error"), + status({ type: "idle" }), + ]), + "prompt", + ], + { + cwd: process.cwd(), + env: { ...process.env, AICTRL_MODELS_PATH: models }, + stdout: "pipe", + stderr: "pipe", + }, + ) + const stdout = await new Response(proc.stdout).text() + await proc.exited + const result = stdout + .split("\n") + .filter(Boolean) + .map((line) => JSON.parse(line)) + .find((event) => event.type === "retry_complete") + + expect(result).toMatchObject({ retryID, outcome: "failed" }) + }, 20_000) +}) diff --git a/packages/cli/test/session/retry.test.ts b/packages/cli/test/session/retry.test.ts index f457bfa..cc3b61c 100644 --- a/packages/cli/test/session/retry.test.ts +++ b/packages/cli/test/session/retry.test.ts @@ -124,6 +124,26 @@ describe("session.retry.retryable", () => { }) }) +describe("session.retry.reason", () => { + test.each([ + [new MessageV2.APIError({ message: "rate limited", statusCode: 429, isRetryable: true }).toObject(), "rate_limit"], + [new MessageV2.APIError({ message: "request timed out", isRetryable: true }).toObject(), "timeout"], + [new MessageV2.APIError({ message: "ECONNRESET socket closed", isRetryable: true }).toObject(), "network"], + [ + new MessageV2.APIError({ message: "provider overloaded", statusCode: 503, isRetryable: true }).toObject(), + "provider", + ], + [new MessageV2.APIError({ message: "transient failure", isRetryable: true }).toObject(), "unknown"], + ] as const)("uses the bounded %s category", (error, expected) => { + expect(SessionRetry.reason(error)).toBe(expected) + }) + + test("classifies existing JSON retry messages without exposing them as dimensions", () => { + expect(SessionRetry.reason(wrap('{"error":{"type":"too_many_requests"}}'))).toBe("rate_limit") + expect(SessionRetry.reason(wrap('{"code":"resource_exhausted"}'))).toBe("provider") + }) +}) + describe("session.message-v2.fromError", () => { test.concurrent( "converts ECONNRESET socket errors to retryable APIError", @@ -195,9 +215,7 @@ describe("session.retry.max_attempts", () => { }) test("processor checks max retry attempts", async () => { - const source = await Bun.file( - path.join(import.meta.dir, "../../src/session/processor.ts"), - ).text() + const source = await Bun.file(path.join(import.meta.dir, "../../src/session/processor.ts")).text() expect(source).toMatch(/attempt\s*>\s*SessionRetry\.MAX_RETRY_ATTEMPTS/) }) }) From fce09ed11f99e422511ba1ce6610f803dd99bc5c Mon Sep 17 00:00:00 2001 From: Bulat Yapparov Date: Mon, 14 Sep 2026 12:24:39 +0100 Subject: [PATCH 2/5] fix(cli): finalize retries from terminal state --- packages/cli/src/cli/cmd/run.errors.ts | 6 +- packages/cli/src/cli/cmd/run.ts | 29 +++-- .../test/cli/classify-session-error.test.ts | 8 ++ .../cli/test/cli/run-retry-telemetry.test.ts | 115 ++++++++++++++++++ 4 files changed, 145 insertions(+), 13 deletions(-) diff --git a/packages/cli/src/cli/cmd/run.errors.ts b/packages/cli/src/cli/cmd/run.errors.ts index 1b86c54..92afbf1 100644 --- a/packages/cli/src/cli/cmd/run.errors.ts +++ b/packages/cli/src/cli/cmd/run.errors.ts @@ -24,6 +24,9 @@ export function classifySessionError(err: unknown): ClassifiedSessionError { if (status === 429) return { reason: "rate_limit", code: "429", message } if (status === 401 || status === 403) return { reason: "auth", code: String(status), message } if (name === "ProviderAuthError") return { reason: "auth", code: status ? String(status) : undefined, message } + if (name === "MessageAbortedError") { + return { reason: "interrupted", code: status ? String(status) : undefined, message } + } if (name === "AbortError" || /timeout/i.test(message)) { return { reason: "timeout", code: status ? String(status) : undefined, message } } @@ -50,7 +53,8 @@ function extractMessage(err: unknown): string { function extractStatus(err: unknown): number | undefined { if (err && typeof err === "object") { const e = err as { status?: unknown; statusCode?: unknown; response?: { status?: unknown }; data?: unknown } - const data = e.data && typeof e.data === "object" ? (e.data as { status?: unknown; statusCode?: unknown }) : undefined + const data = + e.data && typeof e.data === "object" ? (e.data as { status?: unknown; statusCode?: unknown }) : undefined const raw = e.status ?? e.statusCode ?? e.response?.status ?? data?.status ?? data?.statusCode if (typeof raw === "number") return raw if (typeof raw === "string" && /^\d+$/.test(raw)) return Number(raw) diff --git a/packages/cli/src/cli/cmd/run.ts b/packages/cli/src/cli/cmd/run.ts index ed9915c..5c59c1c 100644 --- a/packages/cli/src/cli/cmd/run.ts +++ b/packages/cli/src/cli/cmd/run.ts @@ -554,6 +554,17 @@ export const RunCommand = cmd({ delayMs: number } const retries = new Map() + const outcomes = new Map() + + function retryOutcome(sid: string) { + const outcome = outcomes.get(sid) + if (!outcome) return "unknown" as const + if (outcome.status === "aborted") return "aborted" as const + if (outcome.status === "error" || outcome.finish === "error" || outcome.finish === "content-filter") { + return "failed" as const + } + return "recovered" as const + } function resolveRetry(sid: string, outcome: "recovered" | "failed" | "aborted" | "unknown") { const retry = retries.get(sid) @@ -618,19 +629,13 @@ export const RunCommand = cmd({ if (event.type === "message.updated" && event.properties.info.role === "assistant") { const info = event.properties.info if (args.format === "json") { - if (info.sessionID === sessionID && info.time.completed !== undefined && !emitted.has(info.id)) { - emitted.add(info.id) - const usage = terminalUsage(info) + if (info.sessionID === sessionID && info.time.completed !== undefined) { const status = info.error?.name === "MessageAbortedError" ? "aborted" : info.error ? "error" : "completed" - resolveRetry( - info.sessionID, - status === "completed" && info.finish !== "error" && info.finish !== "content-filter" - ? "recovered" - : status === "aborted" - ? "aborted" - : "failed", - ) + outcomes.set(info.sessionID, { status, finish: info.finish }) + if (emitted.has(info.id)) continue + emitted.add(info.id) + const usage = terminalUsage(info) // Context-window utilization: used = input + cache.read + cache.write // (all prompt tokens that occupy the model's context window this turn). @@ -830,7 +835,7 @@ export const RunCommand = cmd({ } if (status.type === "idle") { if (event.properties.sessionID === sessionID) { - resolveRetry(event.properties.sessionID, "unknown") + resolveRetry(event.properties.sessionID, retryOutcome(event.properties.sessionID)) break } if (childSessions.has(event.properties.sessionID)) { diff --git a/packages/cli/test/cli/classify-session-error.test.ts b/packages/cli/test/cli/classify-session-error.test.ts index f88c532..8d1dd3c 100644 --- a/packages/cli/test/cli/classify-session-error.test.ts +++ b/packages/cli/test/cli/classify-session-error.test.ts @@ -33,6 +33,14 @@ describe("classifySessionError (#63)", () => { expect(classifySessionError(err).reason).toBe("timeout") }) + test("stored MessageAbortedError → interrupted", () => { + const res = classifySessionError({ + name: "MessageAbortedError", + data: { message: "Session cancelled" }, + }) + expect(res.reason).toBe("interrupted") + }) + test("heap OOM → oom", () => { const err = new Error("JavaScript heap out of memory") expect(classifySessionError(err).reason).toBe("oom") diff --git a/packages/cli/test/cli/run-retry-telemetry.test.ts b/packages/cli/test/cli/run-retry-telemetry.test.ts index 224a729..b6d29e6 100644 --- a/packages/cli/test/cli/run-retry-telemetry.test.ts +++ b/packages/cli/test/cli/run-retry-telemetry.test.ts @@ -217,4 +217,119 @@ describe("run --format json retry telemetry (#110)", () => { expect(result).toMatchObject({ retryID, outcome: "failed" }) }, 20_000) + + test("waits for structured-output validation before resolving recovery", async () => { + const retryID = "e62be190-51b4-4564-8bd6-7401ccf60476" + const initial = completed() + const corrected = { + ...initial, + properties: { + ...initial.properties, + info: { + ...initial.properties.info, + error: { + name: "StructuredOutputError", + data: { message: "Model did not produce structured output" }, + }, + }, + }, + } + const proc = Bun.spawn( + [ + "bun", + "run", + cli, + "run", + "--format", + "json", + "--attach", + server([ + status({ + type: "retry", + retryID, + messageID, + providerID: "zai", + modelID: "glm-4.7", + attempt: 1, + reason: "provider", + delayMs: 2_000, + message: "provider error", + next: Date.now() + 2_000, + }), + completed(), + corrected, + status({ type: "idle" }), + ]), + "prompt", + ], + { + cwd: process.cwd(), + env: { ...process.env, AICTRL_MODELS_PATH: models }, + stdout: "pipe", + stderr: "pipe", + }, + ) + const stdout = await new Response(proc.stdout).text() + await proc.exited + const result = stdout + .split("\n") + .filter(Boolean) + .map((line) => JSON.parse(line)) + .find((event) => event.type === "retry_complete") + + expect(result).toMatchObject({ retryID, outcome: "failed" }) + }, 20_000) + + test("treats server-side message cancellation as aborted", async () => { + const retryID = "01fffcfe-4381-47dd-83f9-257175014779" + const proc = Bun.spawn( + [ + "bun", + "run", + cli, + "run", + "--format", + "json", + "--attach", + server([ + status({ + type: "retry", + retryID, + messageID, + providerID: "zai", + modelID: "glm-4.7", + attempt: 1, + reason: "network", + delayMs: 2_000, + message: "network error", + next: Date.now() + 2_000, + }), + { + type: "session.error", + properties: { + sessionID, + error: { name: "MessageAbortedError", data: { message: "Session cancelled" } }, + }, + }, + status({ type: "idle" }), + ]), + "prompt", + ], + { + cwd: process.cwd(), + env: { ...process.env, AICTRL_MODELS_PATH: models }, + stdout: "pipe", + stderr: "pipe", + }, + ) + const stdout = await new Response(proc.stdout).text() + await proc.exited + const output = stdout + .split("\n") + .filter(Boolean) + .map((line) => JSON.parse(line)) + + expect(output.find((event) => event.type === "retry_complete")).toMatchObject({ retryID, outcome: "aborted" }) + expect(output.find((event) => event.type === "session_error")).toMatchObject({ reason: "interrupted" }) + }, 20_000) }) From 0cbc591ffb8094a8fa210453b518d002117cd08e Mon Sep 17 00:00:00 2001 From: Bulat Yapparov Date: Mon, 14 Sep 2026 12:26:43 +0100 Subject: [PATCH 3/5] fix(cli): scope retry outcomes by message --- packages/cli/src/cli/cmd/run.ts | 17 ++++-- .../cli/test/cli/run-retry-telemetry.test.ts | 60 ++++++++++++++++++- 2 files changed, 70 insertions(+), 7 deletions(-) diff --git a/packages/cli/src/cli/cmd/run.ts b/packages/cli/src/cli/cmd/run.ts index 5c59c1c..ddd24e8 100644 --- a/packages/cli/src/cli/cmd/run.ts +++ b/packages/cli/src/cli/cmd/run.ts @@ -556,8 +556,9 @@ export const RunCommand = cmd({ const retries = new Map() const outcomes = new Map() - function retryOutcome(sid: string) { - const outcome = outcomes.get(sid) + function retryOutcome(retry: Retry) { + if (!retry.messageID) return "unknown" as const + const outcome = outcomes.get(retry.messageID) if (!outcome) return "unknown" as const if (outcome.status === "aborted") return "aborted" as const if (outcome.status === "error" || outcome.finish === "error" || outcome.finish === "content-filter") { @@ -588,7 +589,12 @@ export const RunCommand = cmd({ ) { const current = retries.get(sid) if (current && current.retryID === status.retryID) return - resolveRetry(sid, "failed") + if (current) { + resolveRetry( + sid, + current.messageID && status.messageID === current.messageID ? "failed" : retryOutcome(current), + ) + } const retry = { retryID: status.retryID ?? crypto.randomUUID(), messageID: status.messageID ?? null, @@ -632,7 +638,7 @@ export const RunCommand = cmd({ if (info.sessionID === sessionID && info.time.completed !== undefined) { const status = info.error?.name === "MessageAbortedError" ? "aborted" : info.error ? "error" : "completed" - outcomes.set(info.sessionID, { status, finish: info.finish }) + outcomes.set(info.id, { status, finish: info.finish }) if (emitted.has(info.id)) continue emitted.add(info.id) const usage = terminalUsage(info) @@ -835,7 +841,8 @@ export const RunCommand = cmd({ } if (status.type === "idle") { if (event.properties.sessionID === sessionID) { - resolveRetry(event.properties.sessionID, retryOutcome(event.properties.sessionID)) + const retry = retries.get(event.properties.sessionID) + if (retry) resolveRetry(event.properties.sessionID, retryOutcome(retry)) break } if (childSessions.has(event.properties.sessionID)) { diff --git a/packages/cli/test/cli/run-retry-telemetry.test.ts b/packages/cli/test/cli/run-retry-telemetry.test.ts index b6d29e6..c33c35b 100644 --- a/packages/cli/test/cli/run-retry-telemetry.test.ts +++ b/packages/cli/test/cli/run-retry-telemetry.test.ts @@ -34,12 +34,12 @@ function status(value: Record) { } } -function completed(finish = "stop") { +function completed(finish = "stop", id = messageID) { return { type: "message.updated", properties: { info: { - id: messageID, + id, sessionID, role: "assistant", time: { created: 1, completed: 2 }, @@ -332,4 +332,60 @@ describe("run --format json retry telemetry (#110)", () => { expect(output.find((event) => event.type === "retry_complete")).toMatchObject({ retryID, outcome: "aborted" }) expect(output.find((event) => event.type === "session_error")).toMatchObject({ reason: "interrupted" }) }, 20_000) + + test("keeps retry outcomes scoped to their owning message across tool turns", async () => { + const first = "6c0f1b08-87a0-48ac-8c33-cccf68d591f0" + const second = "6907e224-9112-44f5-810d-2cf4f491646c" + const firstMessage = "msg_tool_turn" + const secondMessage = "msg_followup_turn" + const retry = (retryID: string, owner: string) => + status({ + type: "retry", + retryID, + messageID: owner, + providerID: "zai", + modelID: "glm-4.7", + attempt: 1, + reason: "network", + delayMs: 2_000, + message: "network error", + next: Date.now() + 2_000, + }) + const proc = Bun.spawn( + [ + "bun", + "run", + cli, + "run", + "--format", + "json", + "--attach", + server([ + retry(first, firstMessage), + completed("tool-calls", firstMessage), + retry(second, secondMessage), + status({ type: "idle" }), + ]), + "prompt", + ], + { + cwd: process.cwd(), + env: { ...process.env, AICTRL_MODELS_PATH: models }, + stdout: "pipe", + stderr: "pipe", + }, + ) + const stdout = await new Response(proc.stdout).text() + await proc.exited + const output = stdout + .split("\n") + .filter(Boolean) + .map((line) => JSON.parse(line)) + .filter((event) => event.type === "retry_complete") + + expect(output.map((event) => [event.retryID, event.messageID, event.outcome])).toEqual([ + [first, firstMessage, "recovered"], + [second, secondMessage, "unknown"], + ]) + }, 20_000) }) From 4611c378c3599fd58854fa5e7f881128b9e0bf06 Mon Sep 17 00:00:00 2001 From: Bulat Yapparov Date: Mon, 14 Sep 2026 12:28:53 +0100 Subject: [PATCH 4/5] fix(cli): close retries before later messages --- packages/cli/src/cli/cmd/run.ts | 6 ++ .../cli/test/cli/run-retry-telemetry.test.ts | 78 +++++++++++++++++++ 2 files changed, 84 insertions(+) diff --git a/packages/cli/src/cli/cmd/run.ts b/packages/cli/src/cli/cmd/run.ts index ddd24e8..ab6cff1 100644 --- a/packages/cli/src/cli/cmd/run.ts +++ b/packages/cli/src/cli/cmd/run.ts @@ -635,6 +635,12 @@ export const RunCommand = cmd({ if (event.type === "message.updated" && event.properties.info.role === "assistant") { const info = event.properties.info if (args.format === "json") { + if (info.sessionID === sessionID) { + const retry = retries.get(info.sessionID) + if (retry?.messageID && retry.messageID !== info.id) { + resolveRetry(info.sessionID, retryOutcome(retry)) + } + } if (info.sessionID === sessionID && info.time.completed !== undefined) { const status = info.error?.name === "MessageAbortedError" ? "aborted" : info.error ? "error" : "completed" diff --git a/packages/cli/test/cli/run-retry-telemetry.test.ts b/packages/cli/test/cli/run-retry-telemetry.test.ts index c33c35b..9bfd9e7 100644 --- a/packages/cli/test/cli/run-retry-telemetry.test.ts +++ b/packages/cli/test/cli/run-retry-telemetry.test.ts @@ -57,6 +57,20 @@ function completed(finish = "stop", id = messageID) { } } +function started(id: string) { + const event = completed("stop", id) + return { + ...event, + properties: { + ...event.properties, + info: { + ...event.properties.info, + time: { created: 3 }, + }, + }, + } +} + afterEach(() => servers.splice(0).map((item) => item.stop(true))) describe("run --format json retry telemetry (#110)", () => { @@ -388,4 +402,68 @@ describe("run --format json retry telemetry (#110)", () => { [second, secondMessage, "unknown"], ]) }, 20_000) + + test.each([ + ["ProviderAuthError", "nonretryable failure"], + ["MessageAbortedError", "cancellation"], + ])( + "does not attribute a later message %s to the prior recovered retry", + async (name, label) => { + const retryID = "a03f99e1-a49f-4dd9-a0fd-777eec13d8f7" + const firstMessage = "msg_retried_turn" + const secondMessage = "msg_later_failure" + const proc = Bun.spawn( + [ + "bun", + "run", + cli, + "run", + "--format", + "json", + "--attach", + server([ + status({ + type: "retry", + retryID, + messageID: firstMessage, + providerID: "zai", + modelID: "glm-4.7", + attempt: 1, + reason: "network", + delayMs: 2_000, + message: "network error", + next: Date.now() + 2_000, + }), + completed("tool-calls", firstMessage), + started(secondMessage), + { + type: "session.error", + properties: { + sessionID, + error: { name, data: { message: label } }, + }, + }, + status({ type: "idle" }), + ]), + "prompt", + ], + { + cwd: process.cwd(), + env: { ...process.env, AICTRL_MODELS_PATH: models }, + stdout: "pipe", + stderr: "pipe", + }, + ) + const stdout = await new Response(proc.stdout).text() + await proc.exited + const output = stdout + .split("\n") + .filter(Boolean) + .map((line) => JSON.parse(line)) + .filter((event) => event.type === "retry_complete") + + expect(output).toEqual([expect.objectContaining({ retryID, messageID: firstMessage, outcome: "recovered" })]) + }, + 20_000, + ) }) From a250a0ceced9a89c9edf6b09e256761e0b55e981 Mon Sep 17 00:00:00 2001 From: Bulat Yapparov Date: Mon, 14 Sep 2026 12:37:53 +0100 Subject: [PATCH 5/5] fix(sdk): mirror retry status telemetry --- packages/cli/test/cli/usage-token-breakdown.test.ts | 4 ++-- packages/sdk/src/gen/types.gen.ts | 6 ++++++ 2 files changed, 8 insertions(+), 2 deletions(-) diff --git a/packages/cli/test/cli/usage-token-breakdown.test.ts b/packages/cli/test/cli/usage-token-breakdown.test.ts index af48818..6a4e7f5 100644 --- a/packages/cli/test/cli/usage-token-breakdown.test.ts +++ b/packages/cli/test/cli/usage-token-breakdown.test.ts @@ -163,7 +163,7 @@ describe("message_complete emit block shape (source-verified, #86)", () => { describe("EVENTS.md documents token breakdown and context (#86)", () => { test("EVENTS.md message_complete section includes reasoning token field", async () => { const doc = await Bun.file(EVENTS_MD).text() - const idx = doc.indexOf("message_complete") + const idx = doc.indexOf("### `message_complete`") expect(idx).toBeGreaterThan(-1) const section = doc.slice(idx, idx + 1500) expect(section).toContain("reasoning") @@ -171,7 +171,7 @@ describe("EVENTS.md documents token breakdown and context (#86)", () => { test("EVENTS.md message_complete section documents context field", async () => { const doc = await Bun.file(EVENTS_MD).text() - const idx = doc.indexOf("message_complete") + const idx = doc.indexOf("### `message_complete`") expect(idx).toBeGreaterThan(-1) const section = doc.slice(idx, idx + 1500) expect(section).toContain("context") diff --git a/packages/sdk/src/gen/types.gen.ts b/packages/sdk/src/gen/types.gen.ts index 5b92b1b..2cead17 100644 --- a/packages/sdk/src/gen/types.gen.ts +++ b/packages/sdk/src/gen/types.gen.ts @@ -459,6 +459,12 @@ export type SessionStatus = attempt: number message: string next: number + retryID?: string + messageID?: string + providerID?: string + modelID?: string + reason?: "rate_limit" | "timeout" | "network" | "provider" | "unknown" + delayMs?: number } | { type: "busy"