From 3ac98464847e783367bdb0c223c12eb123c438e2 Mon Sep 17 00:00:00 2001 From: Claude Date: Sun, 12 Jul 2026 22:29:28 +0000 Subject: [PATCH 1/2] Adopt Effect internally for typed errors and configurable resilience Replace the scattered try/catch adapters with Effect-based boundaries while keeping the public API Promise-based: - New ts-autocode-training resilience module: named per-operation policies (propose/evaluate/store) composing a per-attempt timeout (typed OperationTimeoutError) with jittered exponential-backoff retries via Effect Schedule, wired through the additive TrainingSettings.resilience field. Without a policy every operation behaves exactly as before. - Training runtime routes all background failures (capture, store, evolve) through one #report boundary into onError; JSON fallbacks become attempt() one-liners. - Harness dispatchAction becomes a linear Effect pipeline whose failure records are best-effort tapError taps; original errors are rethrown unwrapped. AgentActionDeniedError gains a _tag. - Sandbox file operations map errors to their typed fallback values declaratively; ax provider scoring/parsing fallbacks likewise. - rewrite package, loop.ts, and the execution.ts try/finally are left as-is: each is already a minimal, intentional boundary. Co-Authored-By: Claude Fable 5 Claude-Session: https://claude.ai/code/session_01W4PmJx4HSwBjFnCuGXdfwK --- README.md | 16 +++ package-lock.json | 51 ++++++++++ package.json | 1 + packages/harness/package.json | 1 + packages/harness/src/attempt.ts | 17 ++++ packages/harness/src/dispatch.ts | 60 +++++------ packages/harness/src/sandbox.ts | 57 +++++------ packages/training/package.json | 1 + packages/training/src/attempt.ts | 22 +++++ packages/training/src/index.ts | 4 + packages/training/src/resilience.ts | 115 ++++++++++++++++++++++ packages/training/src/training.ts | 101 +++++++++++-------- packages/training/test/resilience.test.ts | 114 +++++++++++++++++++++ packages/training/test/training.test.ts | 84 ++++++++++++++++ src/attempt.ts | 18 ++++ src/index.ts | 7 ++ src/providers/ax.ts | 14 ++- 17 files changed, 569 insertions(+), 114 deletions(-) create mode 100644 packages/harness/src/attempt.ts create mode 100644 packages/training/src/attempt.ts create mode 100644 packages/training/src/resilience.ts create mode 100644 packages/training/test/resilience.test.ts create mode 100644 src/attempt.ts diff --git a/README.md b/README.md index f45a84e..1742273 100644 --- a/README.md +++ b/README.md @@ -232,6 +232,22 @@ Runtime dependencies enter through `TrainingSettings`: - `loop` replaces the default harness orchestration with any `TrainingLoop`. - `secrets` and `variables` are passed to engine factories without entering traces. - `store`, `capture`, and `tracing` configure recording globally. +- `resilience` attaches named timeout/retry policies to runtime operations — + `propose` (the engine/LLM call), `evaluate` (each candidate execution inside + an eval run), and `store` (capture writes). A policy composes a per-attempt + `timeoutMs` (surfacing as a typed `OperationTimeoutError`) with jittered + exponential-backoff retries; operations without a policy behave exactly as + before. For example, to retry rate-limited proposals up to three times with + a 30-second cap per attempt: + + ```ts + configureTraining({ + resilience: { + propose: { timeoutMs: 30_000, retry: { attempts: 3 } }, + }, + }); + ``` + - `source` overrides TypeScript project discovery when the default `tsconfig.json` is not the desired project. - `outputDir` relocates run artifacts and eval output (default `.agentv`, diff --git a/package-lock.json b/package-lock.json index b0f4822..0b67d1b 100644 --- a/package-lock.json +++ b/package-lock.json @@ -14,6 +14,7 @@ ], "dependencies": { "@ax-llm/ax": "^23.0.0", + "effect": "^3.21.4", "ts-autocode-harness": "0.1.0", "ts-autocode-rewrite": "0.1.0", "ts-autocode-training": "0.1.0", @@ -2159,6 +2160,16 @@ "safe-buffer": "^5.0.1" } }, + "node_modules/effect": { + "version": "3.21.4", + "resolved": "https://registry.npmjs.org/effect/-/effect-3.21.4.tgz", + "integrity": "sha512-B89v/xSgPbl1J2Ai2u18jxq3odpFauU1rC6/eSs4FeNHi72kwKdJp12VGigvRV2lK+kRnx+OOz41XV8guZd4gQ==", + "license": "MIT", + "dependencies": { + "@standard-schema/spec": "^1.0.0", + "fast-check": "^3.23.1" + } + }, "node_modules/es-module-lexer": { "version": "2.3.0", "resolved": "https://registry.npmjs.org/es-module-lexer/-/es-module-lexer-2.3.0.tgz", @@ -2199,6 +2210,28 @@ "integrity": "sha512-fjquC59cD7CyW6urNXK0FBufkZcoiGG80wTuPujX590cB5Ttln20E2UB4S/WARVqhXffZl2LNgS+gQdPIIim/g==", "license": "MIT" }, + "node_modules/fast-check": { + "version": "3.23.2", + "resolved": "https://registry.npmjs.org/fast-check/-/fast-check-3.23.2.tgz", + "integrity": "sha512-h5+1OzzfCC3Ef7VbtKdcv7zsstUQwUDlYpUTvjeUsJAssPgLn7QzbboPtL5ro04Mq0rPOsMzl7q5hIbRs2wD1A==", + "funding": [ + { + "type": "individual", + "url": "https://github.com/sponsors/dubzzz" + }, + { + "type": "opencollective", + "url": "https://opencollective.com/fast-check" + } + ], + "license": "MIT", + "dependencies": { + "pure-rand": "^6.1.0" + }, + "engines": { + "node": ">=8.0.0" + } + }, "node_modules/fast-glob": { "version": "3.3.3", "resolved": "https://registry.npmjs.org/fast-glob/-/fast-glob-3.3.3.tgz", @@ -3197,6 +3230,22 @@ "node": ">=12.0.0" } }, + "node_modules/pure-rand": { + "version": "6.1.0", + "resolved": "https://registry.npmjs.org/pure-rand/-/pure-rand-6.1.0.tgz", + "integrity": "sha512-bVWawvoZoBYpp6yIoQtQXHZjmz35RSVHnUOTefl8Vcjr8snTPY1wnpSPMWekcFwbxI6gtmT7rSYPFvz71ldiOA==", + "funding": [ + { + "type": "individual", + "url": "https://github.com/sponsors/dubzzz" + }, + { + "type": "opencollective", + "url": "https://opencollective.com/fast-check" + } + ], + "license": "MIT" + }, "node_modules/queue-microtask": { "version": "1.2.3", "resolved": "https://registry.npmjs.org/queue-microtask/-/queue-microtask-1.2.3.tgz", @@ -3859,6 +3908,7 @@ "dependencies": { "@microsoft/mxc-sdk": "^0.7.0", "deepagents": "^1.10.7", + "effect": "^3.21.4", "unstorage": "^1.17.5", "zod": "^4.4.3" }, @@ -3887,6 +3937,7 @@ "@agentv/core": "^4.42.4", "@arizeai/openinference-semantic-conventions": "^2.5.0", "@opentelemetry/api": "^1.9.1", + "effect": "^3.21.4", "typescript": "^5.9.3", "zod": "^4.4.3" }, diff --git a/package.json b/package.json index 7733010..e48f87f 100644 --- a/package.json +++ b/package.json @@ -62,6 +62,7 @@ }, "dependencies": { "@ax-llm/ax": "^23.0.0", + "effect": "^3.21.4", "ts-autocode-harness": "0.1.0", "ts-autocode-rewrite": "0.1.0", "ts-autocode-training": "0.1.0", diff --git a/packages/harness/package.json b/packages/harness/package.json index 0167a51..7664587 100644 --- a/packages/harness/package.json +++ b/packages/harness/package.json @@ -35,6 +35,7 @@ "dependencies": { "@microsoft/mxc-sdk": "^0.7.0", "deepagents": "^1.10.7", + "effect": "^3.21.4", "unstorage": "^1.17.5", "zod": "^4.4.3" }, diff --git a/packages/harness/src/attempt.ts b/packages/harness/src/attempt.ts new file mode 100644 index 0000000..9a97d1a --- /dev/null +++ b/packages/harness/src/attempt.ts @@ -0,0 +1,17 @@ +// Internal Effect-backed fallback helpers. Deliberately duplicated in each +// workspace package (training, root) instead of adding a shared package for a +// handful of lines; keep the copies in sync. +import { Effect } from "effect"; + +export function errorMessage(error: unknown): string { + return error instanceof Error ? error.message : String(error); +} + +/** Async error-to-value boundary: resolves `fallback(error)` when `fn` throws or rejects. */ +export function attemptAsync(fn: () => Promise, fallback: (error: unknown) => T): Promise { + return Effect.runPromise( + Effect.tryPromise({ try: fn, catch: (error) => error }).pipe( + Effect.catchAll((error) => Effect.sync(() => fallback(error))), + ), + ); +} diff --git a/packages/harness/src/dispatch.ts b/packages/harness/src/dispatch.ts index 6a7a3f6..f1fb5cf 100644 --- a/packages/harness/src/dispatch.ts +++ b/packages/harness/src/dispatch.ts @@ -1,3 +1,6 @@ +import { Cause, Effect, Exit } from "effect"; + +import { errorMessage } from "./attempt.js"; import type { WriteAheadAgentBus } from "./bus.js"; import type { AgentBusEntry } from "./schema.js"; @@ -14,6 +17,7 @@ export type ActionGate = ( ) => JudgeDecision | Promise; export class AgentActionDeniedError extends Error { + readonly _tag = "AgentActionDenied" as const; readonly action: AgentBusEntry; constructor(action: AgentBusEntry) { @@ -42,7 +46,8 @@ export function recordDecision( /** The write-ahead convention, layered on top of the plain message bus: * record the intent, ask the gate, record the verdict as the judge's own * message, execute only after a pass, and record the outcome. Without a gate - * the action is logged and executed. */ + * the action is logged and executed. Failure records are best-effort taps — + * the gate or execution error stays the outcome, rethrown as itself. */ export async function dispatchAction( bus: WriteAheadAgentBus, actor: string, @@ -53,34 +58,33 @@ export async function dispatchAction( ): Promise { const agent = bus.agent(actor); const action = await agent(kind, payload); - if (gate) { - let decision: JudgeDecision; - try { - decision = await gate(action, await bus.read()); - } catch (error) { - // The gate error is the outcome; a failing failure record must not replace it. - await agent(failureOf(kind), { actionId: action.id, stage: "gate", message: errorMessage(error) }) - .catch(() => undefined); - throw error; + const recordFailure = (error: unknown, detail: Readonly>) => + Effect.ignore(promised(() => agent(failureOf(kind), { actionId: action.id, ...detail, message: errorMessage(error) }))); + const program = Effect.gen(function* () { + if (gate) { + const decision = yield* promised(async () => gate(action, await bus.read())).pipe( + Effect.tapError((error) => recordFailure(error, { stage: "gate" })), + ); + yield* promised(() => recordDecision(bus, { subject: "action", actionId: action.id, decision })); + if (decision === "fail") return yield* Effect.fail(new AgentActionDeniedError(action)); } - await recordDecision(bus, { subject: "action", actionId: action.id, decision }); - if (decision === "fail") throw new AgentActionDeniedError(action); - } - let result: T; - try { - result = await execute(); - } catch (error) { - // Likewise: the execution error is the outcome, recorded best-effort. - await agent(failureOf(kind), { actionId: action.id, message: errorMessage(error) }) - .catch(() => undefined); - throw error; - } - // Recorded after the fact, outside the catch: a failing completion append - // surfaces as a bus error, never as a failed action. - await agent(completionOf(kind), { actionId: action.id, result }); - return result; + const result = yield* promised(async () => execute()).pipe( + Effect.tapError((error) => recordFailure(error, {})), + ); + // Recorded outside the failure tap: a failing completion append surfaces + // as a bus error, never as a failed action. + yield* promised(() => agent(completionOf(kind), { actionId: action.id, result })); + return result; + }); + return runUnwrapped(program); } -function errorMessage(error: unknown): string { - return error instanceof Error ? error.message : String(error); +const promised = (fn: () => Promise) => Effect.tryPromise({ try: fn, catch: (error) => error }); + +/** Settles the pipeline into promise semantics, rejecting with the original + * error rather than Effect's fiber wrapper. */ +async function runUnwrapped(effect: Effect.Effect): Promise { + const exit = await Effect.runPromiseExit(effect); + if (Exit.isSuccess(exit)) return exit.value; + throw Cause.squash(exit.cause); } diff --git a/packages/harness/src/sandbox.ts b/packages/harness/src/sandbox.ts index 2731693..6c0ccc4 100644 --- a/packages/harness/src/sandbox.ts +++ b/packages/harness/src/sandbox.ts @@ -9,6 +9,7 @@ import { type FileUploadResponse, } from "deepagents"; +import { attemptAsync } from "./attempt.js"; import { WriteAheadAgentBus } from "./bus.js"; import { dispatchAction, type ActionGate } from "./dispatch.js"; import { createSandboxPolicy } from "./policy.js"; @@ -77,34 +78,26 @@ export class HarnessSandbox extends BaseSandbox { async uploadFiles(files: Array<[string, Uint8Array]>): Promise { return this.#perform("sandbox.upload", { sandbox: this.id, paths: files.map(([path]) => path) }, () => - Promise.all(files.map(async ([path, content]) => { - try { - const target = this.#path(path); - await mkdir(this.#workspace, { recursive: true }); - // Preflight before mkdir: a symlinked ancestor would otherwise create directories outside. - await this.#assertContained(await existingAncestor(dirname(target), this.#workspace)); - await mkdir(dirname(target), { recursive: true }); - await this.#assertContained(dirname(target)); - if (await isSymlink(target)) throw new Error("path escapes sandbox workspace"); - await writeFile(target, content); - return { path, error: null }; - } catch { - return { path, error: "permission_denied" as const }; - } - }))); + Promise.all(files.map(([path, content]) => attemptAsync(async () => { + const target = this.#path(path); + await mkdir(this.#workspace, { recursive: true }); + // Preflight before mkdir: a symlinked ancestor would otherwise create directories outside. + await this.#assertContained(await existingAncestor(dirname(target), this.#workspace)); + await mkdir(dirname(target), { recursive: true }); + await this.#assertContained(dirname(target)); + if (await isSymlink(target)) throw new Error("path escapes sandbox workspace"); + await writeFile(target, content); + return { path, error: null }; + }, () => ({ path, error: "permission_denied" as const }))))); } async downloadFiles(paths: string[]): Promise { return this.#perform("sandbox.download", { sandbox: this.id, paths }, () => - Promise.all(paths.map(async (path) => { - try { - const target = this.#path(path); - await this.#assertContained(target); - return { path, content: await readFile(target), error: null }; - } catch { - return { path, content: null, error: "file_not_found" as const }; - } - }))); + Promise.all(paths.map((path) => attemptAsync(async () => { + const target = this.#path(path); + await this.#assertContained(target); + return { path, content: await readFile(target), error: null }; + }, () => ({ path, content: null, error: "file_not_found" as const }))))); } #perform(kind: string, payload: unknown, execute: () => Promise): Promise { @@ -127,12 +120,8 @@ export class HarnessSandbox extends BaseSandbox { } } -async function isSymlink(path: string): Promise { - try { - return (await lstat(path)).isSymbolicLink(); - } catch { - return false; - } +function isSymlink(path: string): Promise { + return attemptAsync(async () => (await lstat(path)).isSymbolicLink(), () => false); } async function existingAncestor(path: string, root: string): Promise { @@ -141,11 +130,9 @@ async function existingAncestor(path: string, root: string): Promise { return current; } -async function exists(path: string): Promise { - try { +function exists(path: string): Promise { + return attemptAsync(async () => { await lstat(path); return true; - } catch { - return false; - } + }, () => false); } diff --git a/packages/training/package.json b/packages/training/package.json index 7f5d77a..ee37d70 100644 --- a/packages/training/package.json +++ b/packages/training/package.json @@ -36,6 +36,7 @@ "@agentv/core": "^4.42.4", "@arizeai/openinference-semantic-conventions": "^2.5.0", "@opentelemetry/api": "^1.9.1", + "effect": "^3.21.4", "typescript": "^5.9.3", "zod": "^4.4.3" }, diff --git a/packages/training/src/attempt.ts b/packages/training/src/attempt.ts new file mode 100644 index 0000000..42098ec --- /dev/null +++ b/packages/training/src/attempt.ts @@ -0,0 +1,22 @@ +// Internal Effect-backed fallback helpers. Deliberately duplicated in each +// workspace package (harness, root) instead of adding a shared package for a +// handful of lines; keep the copies in sync. +import { Effect } from "effect"; + +export function errorMessage(error: unknown): string { + return error instanceof Error ? error.message : String(error); +} + +/** Runs `fn`, mapping a throw to `fallback(error)` — a sync error-to-value boundary. */ +export function attempt(fn: () => T, fallback: (error: unknown) => T): T { + return Effect.runSync(Effect.try(fn).pipe(Effect.catchAll((error) => Effect.sync(() => fallback(error))))); +} + +/** Async variant: resolves `fallback(error)` when `fn` throws or rejects. */ +export function attemptAsync(fn: () => Promise, fallback: (error: unknown) => T): Promise { + return Effect.runPromise( + Effect.tryPromise({ try: fn, catch: (error) => error }).pipe( + Effect.catchAll((error) => Effect.sync(() => fallback(error))), + ), + ); +} diff --git a/packages/training/src/index.ts b/packages/training/src/index.ts index b01500d..668814e 100644 --- a/packages/training/src/index.ts +++ b/packages/training/src/index.ts @@ -11,6 +11,7 @@ export type { Activation, AppliedPromotion, CaptureSettings, + ErrorPhase, EvolutionSettings, PromotionApplier, TrainInput, @@ -21,6 +22,9 @@ export type { TracingSettings, } from "./training.js"; +export { defaultRetry, OperationTimeoutError, withPolicy } from "./resilience.js"; +export type { ResiliencePolicy, ResilienceSettings, RetryOptions } from "./resilience.js"; + export { defaultFanOut, defaultMaxRounds, sequentialLoop, trainingRounds } from "./loop.js"; export type { CandidateReview, diff --git a/packages/training/src/resilience.ts b/packages/training/src/resilience.ts new file mode 100644 index 0000000..1105fe0 --- /dev/null +++ b/packages/training/src/resilience.ts @@ -0,0 +1,115 @@ +import { Cause, Data, Duration, Effect, Exit, Option, Schedule } from "effect"; + +/** A per-attempt timeout imposed by a {@link ResiliencePolicy}. Retryable by + * default, so a policy with both `timeoutMs` and `retry` re-attempts timed-out + * operations. */ +export class OperationTimeoutError extends Data.TaggedError("OperationTimeout")<{ + readonly operation: string; + readonly timeoutMs: number; +}> { + override get message(): string { + return `${this.operation} timed out after ${this.timeoutMs}ms`; + } +} + +export interface RetryOptions { + /** Total attempts including the first; 1 (the default) means no retries. */ + readonly attempts?: number; + /** Base delay of the exponential backoff between attempts. */ + readonly delayMs?: number; + /** Ceiling for the backoff delay. */ + readonly maxDelayMs?: number; + /** Randomize each delay to avoid thundering herds; on by default. */ + readonly jitter?: boolean; + /** Which failures are worth another attempt. Defaults to all of them; + * narrow this to skip deterministic failures like validation errors. */ + readonly retryable?: (error: unknown) => boolean; +} + +/** Timeout and retry behavior for one named operation. An empty policy is the + * identity: the operation runs exactly as it would without one. */ +export interface ResiliencePolicy { + /** Per-attempt time limit; expiry fails the attempt with {@link OperationTimeoutError}. */ + readonly timeoutMs?: number; + readonly retry?: RetryOptions; +} + +/** Named policies for the operations the training runtime performs, in the + * style of a resilience-pipeline registry. Unnamed operations run bare. */ +export interface ResilienceSettings { + /** Candidate proposal — the engine/LLM call. */ + readonly propose?: ResiliencePolicy; + /** Each candidate execution inside an evaluation run. Retries apply per + * eval case, which suits flaky sandboxes. */ + readonly evaluate?: ResiliencePolicy; + /** Capture writes via `TrainingStore.append`. Retrying a store whose + * failures are ambiguous can append duplicates; keep retries off unless + * the store is idempotent or failures are known-clean. */ + readonly store?: ResiliencePolicy; +} + +export const defaultRetry: Required> = Object.freeze({ + attempts: 1, + delayMs: 250, + maxDelayMs: 10_000, + jitter: true, +}); + +/** Runs `fn` under `policy`: per-attempt timeout, then jittered exponential + * retry, composed with Effect and settled back into an ordinary promise. + * Failures reject with the original error (or {@link OperationTimeoutError}), + * never a wrapped one. Without a policy this is exactly `fn(signal)`. Each + * attempt receives a signal that fires on the caller's `signal`, its own + * timeout, or teardown; an aborted caller signal is never retried. */ +export function withPolicy( + policy: ResiliencePolicy | undefined, + operation: string, + fn: (signal?: AbortSignal) => Promise, + signal?: AbortSignal, +): Promise { + if (policy?.timeoutMs === undefined && policy?.retry === undefined) return fn(signal); + const { timeoutMs, retry } = policy; + const bare: Effect.Effect = Effect.tryPromise({ + try: (attemptSignal) => fn(signal === undefined ? attemptSignal : AbortSignal.any([signal, attemptSignal])), + catch: (error) => error, + }); + const attempt = timeoutMs === undefined ? bare : bare.pipe(Effect.timeoutFail({ + duration: Duration.millis(timeoutMs), + onTimeout: () => new OperationTimeoutError({ operation, timeoutMs }), + })); + return unwrapExit( + Effect.runPromiseExit( + Effect.retry(attempt, retrySchedule(retry ?? {}, signal)), + signal === undefined ? undefined : { signal }, + ), + operation, + signal, + ); +} + +function retrySchedule(retry: RetryOptions, signal: AbortSignal | undefined): Schedule.Schedule { + const attempts = Math.max(1, Math.floor(retry.attempts ?? defaultRetry.attempts)); + const retryable = retry.retryable ?? (() => true); + const capped = Schedule.exponential(Duration.millis(retry.delayMs ?? defaultRetry.delayMs)).pipe( + // Union takes the shorter interval, capping the exponential growth. + Schedule.union(Schedule.spaced(Duration.millis(retry.maxDelayMs ?? defaultRetry.maxDelayMs))), + ); + const backoff = (retry.jitter ?? defaultRetry.jitter) ? Schedule.jittered(capped) : capped; + return backoff.pipe( + Schedule.intersect(Schedule.recurs(attempts - 1)), + Schedule.whileInput((error: unknown) => signal?.aborted !== true && retryable(error)), + ); +} + +/** Settles an Effect exit into promise semantics: the typed failure or defect + * rejects as itself, interruption surfaces the caller's abort reason. */ +async function unwrapExit(exit: Promise>, operation: string, signal?: AbortSignal): Promise { + const settled = await exit; + if (Exit.isSuccess(settled)) return settled.value; + const failure = Cause.failureOption(settled.cause); + if (Option.isSome(failure)) throw failure.value; + const defect = Cause.dieOption(settled.cause); + if (Option.isSome(defect)) throw defect.value; + signal?.throwIfAborted(); + throw new Error(`${operation} was interrupted`); +} diff --git a/packages/training/src/training.ts b/packages/training/src/training.ts index 3d9df47..b00660d 100644 --- a/packages/training/src/training.ts +++ b/packages/training/src/training.ts @@ -5,6 +5,7 @@ import { z } from "zod"; import { OpenInferenceSpanKind, SemanticConventions } from "@arizeai/openinference-semantic-conventions"; import { SpanStatusCode, trace, type Attributes, type Span, type Tracer } from "@opentelemetry/api"; +import { attempt, errorMessage } from "./attempt.js"; import { CandidateEngine, type BoundEvaluation, @@ -13,6 +14,7 @@ import { type SecretProvider, type TrainingEngine, } from "./engine.js"; +import { withPolicy, type ResilienceSettings } from "./resilience.js"; import { evaluateTrainable, type TrainableEvalRun } from "./evaluation.js"; import { sequentialLoop, type TrainingLoop, type TrainingRound } from "./loop.js"; import { evaluatePromotionGate, type PromotionDecision, type PromotionGate } from "./promotion.js"; @@ -82,9 +84,16 @@ export interface TrainingSettings { readonly variables?: Readonly>; readonly capture?: CaptureSettings; readonly tracing?: TracingSettings; - readonly onError?: (error: unknown, phase: "capture" | "store" | "evolve") => void; + /** Timeout/retry policies for named runtime operations; operations without + * a policy behave exactly as before. */ + readonly resilience?: ResilienceSettings; + readonly onError?: (error: unknown, phase: ErrorPhase) => void; } +/** Where a background failure was routed from: trace capture, store writes, + * or background evolution. These never fail the traced call itself. */ +export type ErrorPhase = "capture" | "store" | "evolve"; + export interface TrainInput { readonly trainable: TrainableIdentity; /** Optimization goal; defaults to preserving the evaluated behavior. */ @@ -214,7 +223,7 @@ class TrainingRuntime implements Training { } evolution.onEvolved?.(await run.activate()); })() - .catch((error) => this.#settings.onError?.(error, "evolve")) + .catch(this.#report("evolve")) .finally(() => { state.running = false; if (state.queued) { @@ -244,11 +253,16 @@ class TrainingRuntime implements Training { const evaluated = await evaluateTrainable(token, { ...evaluation, task: async (input) => { - const output = await execute( - candidate.target, - candidate.implementation, - evaluationArgs(input), - signal === undefined ? {} : { signal }, + const output = await withPolicy( + this.#settings.resilience?.evaluate, + "candidate.execute", + (attemptSignal) => execute( + candidate.target, + candidate.implementation, + evaluationArgs(input), + attemptSignal === undefined ? {} : { signal: attemptSignal }, + ), + signal, ); return typeof output === "string" ? output : JSON.stringify(output) ?? String(output); }, @@ -352,20 +366,25 @@ class TrainingRuntime implements Training { }): Promise { const target = findTrainable(token.id, this.#settings.source); const records = await this.records(token); - return this.#engineFor(input.engine).propose( - { - trainableId: token.id, - objective: input.objective, - target, - records, - evaluations: this.#evaluations.get(token.id) ?? [], - ...(input.constraints.length === 0 ? {} : { constraints: input.constraints }), - }, - { - variables: this.#variables, - ...(this.#settings.secrets === undefined ? {} : { secrets: this.#settings.secrets }), - ...(input.signal === undefined ? {} : { signal: input.signal }), - }, + return withPolicy( + this.#settings.resilience?.propose, + "engine.propose", + (signal) => this.#engineFor(input.engine).propose( + { + trainableId: token.id, + objective: input.objective, + target, + records, + evaluations: this.#evaluations.get(token.id) ?? [], + ...(input.constraints.length === 0 ? {} : { constraints: input.constraints }), + }, + { + variables: this.#variables, + ...(this.#settings.secrets === undefined ? {} : { secrets: this.#settings.secrets }), + ...(signal === undefined ? {} : { signal }), + }, + ), + input.signal, ); } @@ -462,7 +481,7 @@ class TrainingRuntime implements Training { } const capture = this.#settings.capture ?? {}; if (capture.enabled === false) return; - try { + attempt(() => { const spanContext = span?.spanContext(); const input = capture.mapInput ? capture.mapInput(args, token) : args; const output = error === undefined @@ -491,21 +510,27 @@ class TrainingRuntime implements Training { ...(error === undefined ? {} : { error: errorMessage(error) }), }), }; - this.#enqueue(this.#store.append(record)); + this.#enqueue(() => this.#store.append(record)); if (error === undefined) this.#maybeEvolve(token); - } catch (captureError) { - this.#settings.onError?.(captureError, "capture"); - } + }, this.#report("capture")); + } + + /** The boundary sink for background failures: every capture, store, and + * evolution error funnels through here into `TrainingSettings.onError`. */ + #report(phase: ErrorPhase): (error: unknown) => void { + return (error) => this.#settings.onError?.(error, phase); } #serialize(value: unknown): string { return (this.#settings.capture?.serialize ?? defaultSerialize)(value); } - #enqueue(write: Promise): void { - const pending = write.catch((error) => this.#settings.onError?.(error, "store")).finally(() => { - this.#pending.delete(pending); - }); + #enqueue(write: () => Promise): void { + const pending = withPolicy(this.#settings.resilience?.store, "store.append", write) + .catch(this.#report("store")) + .finally(() => { + this.#pending.delete(pending); + }); this.#pending.add(pending); } } @@ -570,24 +595,14 @@ function isPromise(value: T): value is T & Promise> { function defaultSerialize(value: unknown): string { if (typeof value === "string") return value; - try { - return JSON.stringify(value) ?? String(value); - } catch { - return String(value); - } -} - -function errorMessage(error: unknown): string { - return error instanceof Error ? error.message : String(error); + return attempt(() => JSON.stringify(value) ?? String(value), () => String(value)); } function evaluationArgs(input: string): readonly unknown[] { - try { + return attempt(() => { const parsed = JSON.parse(input) as unknown; return Array.isArray(parsed) ? parsed : [parsed]; - } catch { - return [input]; - } + }, () => [input]); } function liveEvalCases(records: readonly TrainingRecord[]): readonly EvalTestInput[] { diff --git a/packages/training/test/resilience.test.ts b/packages/training/test/resilience.test.ts new file mode 100644 index 0000000..1562c85 --- /dev/null +++ b/packages/training/test/resilience.test.ts @@ -0,0 +1,114 @@ +import { describe, expect, it, vi } from "vitest"; + +import { OperationTimeoutError, withPolicy } from "../src/index.js"; + +describe("withPolicy", () => { + it("is a passthrough without a policy", async () => { + const signal = new AbortController().signal; + const fn = vi.fn(async (received?: AbortSignal) => { + expect(received).toBe(signal); + return 42; + }); + + await expect(withPolicy(undefined, "op", fn, signal)).resolves.toBe(42); + await expect(withPolicy({}, "op", fn, signal)).resolves.toBe(42); + expect(fn).toHaveBeenCalledTimes(2); + }); + + it("rejects with the original error without a policy", async () => { + const failure = new Error("boom"); + + await expect(withPolicy(undefined, "op", async () => { throw failure; })).rejects.toBe(failure); + }); + + it("retries until an attempt succeeds", async () => { + let attempts = 0; + + const result = await withPolicy({ retry: { attempts: 3, delayMs: 1, jitter: false } }, "op", async () => { + attempts += 1; + if (attempts < 3) throw new Error(`attempt ${attempts} failed`); + return "ok"; + }); + + expect(result).toBe("ok"); + expect(attempts).toBe(3); + }); + + it("rethrows the last original error unwrapped when retries are exhausted", async () => { + const failures = [new Error("first"), new Error("second")]; + let attempts = 0; + + await expect(withPolicy( + { retry: { attempts: 2, delayMs: 1, jitter: false } }, + "op", + async () => { throw failures[attempts++]; }, + )).rejects.toBe(failures[1]); + expect(attempts).toBe(2); + }); + + it("does not retry failures the policy marks non-retryable", async () => { + let attempts = 0; + + await expect(withPolicy( + { retry: { attempts: 3, delayMs: 1, retryable: (error) => !(error instanceof TypeError) } }, + "op", + async () => { + attempts += 1; + throw new TypeError("bad input"); + }, + )).rejects.toThrow("bad input"); + expect(attempts).toBe(1); + }); + + it("fails a hung attempt with a typed timeout and aborts the attempt signal", async () => { + let observed: AbortSignal | undefined; + + const rejection = withPolicy({ timeoutMs: 20 }, "slow.op", (signal) => { + observed = signal; + return new Promise((_, reject) => signal?.addEventListener("abort", () => reject(signal.reason))); + }); + + const error: unknown = await rejection.then(() => undefined, (thrown: unknown) => thrown); + expect(error).toBeInstanceOf(OperationTimeoutError); + expect(error).toMatchObject({ _tag: "OperationTimeout", operation: "slow.op", timeoutMs: 20 }); + expect((error as Error).message).toBe("slow.op timed out after 20ms"); + expect(observed?.aborted).toBe(true); + }); + + it("retries timed-out attempts by default", async () => { + let attempts = 0; + + const result = await withPolicy( + { timeoutMs: 100, retry: { attempts: 2, delayMs: 1, jitter: false } }, + "op", + async () => { + attempts += 1; + if (attempts === 1) await new Promise((resolve) => setTimeout(resolve, 1_000)); + return "ok"; + }, + ); + + expect(result).toBe("ok"); + expect(attempts).toBe(2); + }); + + it("stops retrying when the caller aborts during backoff", async () => { + const controller = new AbortController(); + const reason = new Error("stop now"); + let attempts = 0; + + const pending = withPolicy( + { retry: { attempts: 5, delayMs: 5_000, jitter: false } }, + "op", + async () => { + attempts += 1; + throw new Error("flaky"); + }, + controller.signal, + ); + setTimeout(() => controller.abort(reason), 20); + + await expect(pending).rejects.toBe(reason); + expect(attempts).toBe(1); + }); +}); diff --git a/packages/training/test/training.test.ts b/packages/training/test/training.test.ts index 8728eef..6a1958a 100644 --- a/packages/training/test/training.test.ts +++ b/packages/training/test/training.test.ts @@ -11,10 +11,12 @@ import { captureTrainable, configureTraining, defineTrainable, + MemoryTrainingStore, training as defaultTraining, type Activation, type ImplementationExecutor, type TrainingEngine, + type TrainingStore, } from "../src/index.js"; describe("trainable identity", () => { @@ -235,5 +237,87 @@ describe("training execution", () => { }); }); +describe("training resilience policies", () => { + it("retries store appends under a store policy without surfacing an error", async () => { + const inner = new MemoryTrainingStore(); + let failures = 1; + const store: TrainingStore = { + append: async (record) => { + if (failures > 0) { + failures -= 1; + throw new Error("flaky store"); + } + await inner.append(record); + }, + list: (trainableId) => inner.list(trainableId), + }; + const onError = vi.fn(); + const training = configureTraining({ + store, + onError, + tracing: { enabled: false }, + resilience: { store: { retry: { attempts: 2, delayMs: 1, jitter: false } } }, + }); + + captureTrainable("Router.storeRetry", "route", undefined, (input: string) => input, ["billing"]); + await training.flush(); + + expect(onError).not.toHaveBeenCalled(); + expect(await training.records(defineTrainable("Router.storeRetry"))).toHaveLength(1); + }); + + it("routes store failures to onError when no store policy is configured", async () => { + const failure = new Error("store down"); + const onError = vi.fn(); + const training = configureTraining({ + store: { append: async () => { throw failure; }, list: async () => [] }, + onError, + tracing: { enabled: false }, + }); + + captureTrainable("Router.storeDown", "route", undefined, (input: string) => input, ["billing"]); + await training.flush(); + + expect(onError).toHaveBeenCalledWith(failure, "store"); + }); + + it("retries engine proposals under a propose policy", async () => { + const directory = await mkdtemp(join(tmpdir(), "ts-autocode-propose-retry-")); + const artifact = join(directory, "echo.ts"); + await writeFile(artifact, `export function echoRetry(input: string): string { + "use training"; + return input; +}\n`); + let failures = 1; + const optimize = vi.fn(async () => { + if (failures > 0) { + failures -= 1; + throw new Error("rate limited"); + } + return { implementation: "return input.toUpperCase();" }; + }); + const training = configureTraining({ + engine: { id: "propose-retry-test", optimize }, + executor: functionExecutor, + source: { files: [artifact] }, + tracing: { enabled: false }, + resilience: { propose: { retry: { attempts: 2, delayMs: 1, jitter: false } } }, + }); + + const run = await training.train({ + trainable: defineTrainable("echoRetry").symbol, + objective: "Uppercase the input", + evaluation: { + tests: [{ id: "upper", input: "abc", assert: [{ type: "equals", value: "ABC" }] }], + task: (input) => input.toUpperCase(), + outputDir: join(directory, "agentv"), + }, + }); + + expect(run.outcome).toBe("ready"); + expect(optimize).toHaveBeenCalledTimes(2); + }); +}); + const functionExecutor: ImplementationExecutor = async (target, implementation, args) => new Function(...target.parameters.map((parameter) => parameter.name), implementation)(...args); diff --git a/src/attempt.ts b/src/attempt.ts new file mode 100644 index 0000000..6637a16 --- /dev/null +++ b/src/attempt.ts @@ -0,0 +1,18 @@ +// Internal Effect-backed fallback helpers. Deliberately duplicated in each +// workspace package (training, harness) instead of adding a shared package for +// a handful of lines; keep the copies in sync. +import { Effect } from "effect"; + +/** Runs `fn`, mapping a throw to `fallback(error)` — a sync error-to-value boundary. */ +export function attempt(fn: () => T, fallback: (error: unknown) => T): T { + return Effect.runSync(Effect.try(fn).pipe(Effect.catchAll((error) => Effect.sync(() => fallback(error))))); +} + +/** Async variant: resolves `fallback(error)` when `fn` throws or rejects. */ +export function attemptAsync(fn: () => Promise, fallback: (error: unknown) => T): Promise { + return Effect.runPromise( + Effect.tryPromise({ try: fn, catch: (error) => error }).pipe( + Effect.catchAll((error) => Effect.sync(() => fallback(error))), + ), + ); +} diff --git a/src/index.ts b/src/index.ts index 8bc356e..b16992b 100644 --- a/src/index.ts +++ b/src/index.ts @@ -29,9 +29,11 @@ export { captureTrainable, configureTraining, MemoryTrainingStore, + OperationTimeoutError, defaultEvolution, defaultObjective, defaultOutputDir, + defaultRetry, defaultTsconfig, defineTrainable, discoverTrainables, @@ -39,6 +41,7 @@ export { provideTrainingDefaults, training, trainingMarker, + withPolicy, } from "ts-autocode-training"; export type { Activation, @@ -49,6 +52,7 @@ export type { CaptureSettings, EngineCandidate, EngineContext, + ErrorPhase, EvolutionSettings, ImplementationExecutor, OptimizeRequest, @@ -57,6 +61,9 @@ export type { PromotionGate, PromotionGateContext, PromotionGateInput, + ResiliencePolicy, + ResilienceSettings, + RetryOptions, RoundObserver, RoundSequence, SecretProvider, diff --git a/src/providers/ax.ts b/src/providers/ax.ts index 440e13f..66de0a2 100644 --- a/src/providers/ax.ts +++ b/src/providers/ax.ts @@ -10,6 +10,7 @@ import { import type { EngineContext, OptimizeRequest, TrainingEngine } from "ts-autocode-training"; +import { attempt, attemptAsync } from "../attempt.js"; import { defaultExecutionTimeoutMs, executeImplementation } from "../execution.js"; export { defaultExecutionTimeoutMs } from "../execution.js"; @@ -170,15 +171,14 @@ async function scoreImplementation( ): Promise { if (!prediction?.[field.output]?.trim()) return 0; const args = JSON.parse(String(exampleValue[field.args] ?? "[]")) as unknown[]; - try { + // A candidate that fails to run simply scores zero. + return attemptAsync(async () => { const actual = await executeImplementation(request.target, prediction[field.output], args, { timeoutMs: timeout, ...(signal === undefined ? {} : { signal }), }); return outputText(actual) === String(exampleValue[field.expected] ?? "") ? 1 : 0; - } catch { - return 0; - } + }, () => 0); } async function service(value: Service | undefined, context: EngineContext): Promise { @@ -211,12 +211,10 @@ function argsFromMessages(messages: OptimizeRequest["evaluations"][number]["resu function argsFromContent(content: unknown): unknown[] | undefined { const text = contentText(content); if (text === undefined) return undefined; - try { + return attempt(() => { const parsed = JSON.parse(text) as unknown; return Array.isArray(parsed) ? parsed : [parsed]; - } catch { - return [text]; - } + }, () => [text]); } function contentText(content: unknown): string | undefined { From 544814a2ed87300c2b0e12ff4be607dfa8126929 Mon Sep 17 00:00:00 2001 From: Claude Date: Sun, 12 Jul 2026 23:02:13 +0000 Subject: [PATCH 2/2] Preserve raw thrown values in attempt() and align the duplicated helpers Address review findings: the Effect.try shorthand wrapped sync throws in UnknownException, so attempt()'s fallback (and therefore onError in the capture phase) received the wrapper instead of the original error. Use the object form with an identity catch, add a regression test locking the raw-value contract, and make the three duplicated attempt.ts copies byte-identical (errorMessage + attempt + attemptAsync) as their headers promise. Co-Authored-By: Claude Fable 5 Claude-Session: https://claude.ai/code/session_01W4PmJx4HSwBjFnCuGXdfwK --- packages/harness/src/attempt.ts | 16 ++++++++-- packages/training/src/attempt.ts | 13 +++++--- packages/training/test/attempt.test.ts | 41 ++++++++++++++++++++++++++ src/attempt.ts | 17 ++++++++--- 4 files changed, 76 insertions(+), 11 deletions(-) create mode 100644 packages/training/test/attempt.test.ts diff --git a/packages/harness/src/attempt.ts b/packages/harness/src/attempt.ts index 9a97d1a..e7d0a83 100644 --- a/packages/harness/src/attempt.ts +++ b/packages/harness/src/attempt.ts @@ -1,13 +1,23 @@ // Internal Effect-backed fallback helpers. Deliberately duplicated in each -// workspace package (training, root) instead of adding a shared package for a -// handful of lines; keep the copies in sync. +// workspace package (root src/, packages/training, packages/harness) instead +// of adding a shared package for a handful of lines; keep the copies identical. import { Effect } from "effect"; export function errorMessage(error: unknown): string { return error instanceof Error ? error.message : String(error); } -/** Async error-to-value boundary: resolves `fallback(error)` when `fn` throws or rejects. */ +/** Runs `fn`, mapping a throw to `fallback(error)` — a sync error-to-value + * boundary. The fallback receives the raw thrown value. */ +export function attempt(fn: () => T, fallback: (error: unknown) => T): T { + return Effect.runSync( + Effect.try({ try: fn, catch: (error) => error }).pipe( + Effect.catchAll((error) => Effect.sync(() => fallback(error))), + ), + ); +} + +/** Async variant: resolves `fallback(error)` when `fn` throws or rejects. */ export function attemptAsync(fn: () => Promise, fallback: (error: unknown) => T): Promise { return Effect.runPromise( Effect.tryPromise({ try: fn, catch: (error) => error }).pipe( diff --git a/packages/training/src/attempt.ts b/packages/training/src/attempt.ts index 42098ec..e7d0a83 100644 --- a/packages/training/src/attempt.ts +++ b/packages/training/src/attempt.ts @@ -1,15 +1,20 @@ // Internal Effect-backed fallback helpers. Deliberately duplicated in each -// workspace package (harness, root) instead of adding a shared package for a -// handful of lines; keep the copies in sync. +// workspace package (root src/, packages/training, packages/harness) instead +// of adding a shared package for a handful of lines; keep the copies identical. import { Effect } from "effect"; export function errorMessage(error: unknown): string { return error instanceof Error ? error.message : String(error); } -/** Runs `fn`, mapping a throw to `fallback(error)` — a sync error-to-value boundary. */ +/** Runs `fn`, mapping a throw to `fallback(error)` — a sync error-to-value + * boundary. The fallback receives the raw thrown value. */ export function attempt(fn: () => T, fallback: (error: unknown) => T): T { - return Effect.runSync(Effect.try(fn).pipe(Effect.catchAll((error) => Effect.sync(() => fallback(error))))); + return Effect.runSync( + Effect.try({ try: fn, catch: (error) => error }).pipe( + Effect.catchAll((error) => Effect.sync(() => fallback(error))), + ), + ); } /** Async variant: resolves `fallback(error)` when `fn` throws or rejects. */ diff --git a/packages/training/test/attempt.test.ts b/packages/training/test/attempt.test.ts new file mode 100644 index 0000000..227f934 --- /dev/null +++ b/packages/training/test/attempt.test.ts @@ -0,0 +1,41 @@ +import { describe, expect, it } from "vitest"; + +import { attempt, attemptAsync, errorMessage } from "../src/attempt.js"; + +describe("attempt fallbacks", () => { + it("passes the raw thrown value to the sync fallback", () => { + const failure = new Error("sync boom"); + let received: unknown; + + const result = attempt(() => { throw failure; }, (error) => { + received = error; + return "fell back"; + }); + + expect(result).toBe("fell back"); + expect(received).toBe(failure); + }); + + it("passes the raw rejection to the async fallback", async () => { + const failure = new Error("async boom"); + let received: unknown; + + const result = await attemptAsync(async () => { throw failure; }, (error) => { + received = error; + return "fell back"; + }); + + expect(result).toBe("fell back"); + expect(received).toBe(failure); + }); + + it("returns the function result untouched on success", async () => { + expect(attempt(() => 7, () => 0)).toBe(7); + await expect(attemptAsync(async () => 7, () => 0)).resolves.toBe(7); + }); + + it("stringifies non-Error values in errorMessage", () => { + expect(errorMessage(new Error("named"))).toBe("named"); + expect(errorMessage("plain")).toBe("plain"); + }); +}); diff --git a/src/attempt.ts b/src/attempt.ts index 6637a16..e7d0a83 100644 --- a/src/attempt.ts +++ b/src/attempt.ts @@ -1,11 +1,20 @@ // Internal Effect-backed fallback helpers. Deliberately duplicated in each -// workspace package (training, harness) instead of adding a shared package for -// a handful of lines; keep the copies in sync. +// workspace package (root src/, packages/training, packages/harness) instead +// of adding a shared package for a handful of lines; keep the copies identical. import { Effect } from "effect"; -/** Runs `fn`, mapping a throw to `fallback(error)` — a sync error-to-value boundary. */ +export function errorMessage(error: unknown): string { + return error instanceof Error ? error.message : String(error); +} + +/** Runs `fn`, mapping a throw to `fallback(error)` — a sync error-to-value + * boundary. The fallback receives the raw thrown value. */ export function attempt(fn: () => T, fallback: (error: unknown) => T): T { - return Effect.runSync(Effect.try(fn).pipe(Effect.catchAll((error) => Effect.sync(() => fallback(error))))); + return Effect.runSync( + Effect.try({ try: fn, catch: (error) => error }).pipe( + Effect.catchAll((error) => Effect.sync(() => fallback(error))), + ), + ); } /** Async variant: resolves `fallback(error)` when `fn` throws or rejects. */