Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
16 changes: 16 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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`,
Expand Down
51 changes: 51 additions & 0 deletions package-lock.json

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

1 change: 1 addition & 0 deletions package.json
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down
1 change: 1 addition & 0 deletions packages/harness/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -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"
},
Expand Down
27 changes: 27 additions & 0 deletions packages/harness/src/attempt.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,27 @@
// Internal Effect-backed fallback helpers. Deliberately duplicated in each
// 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. The fallback receives the raw thrown value. */
export function attempt<T>(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<T>(fn: () => Promise<T>, fallback: (error: unknown) => T): Promise<T> {
return Effect.runPromise(
Effect.tryPromise({ try: fn, catch: (error) => error }).pipe(
Effect.catchAll((error) => Effect.sync(() => fallback(error))),
),
);
}
Comment thread
coderabbitai[bot] marked this conversation as resolved.
60 changes: 32 additions & 28 deletions packages/harness/src/dispatch.ts
Original file line number Diff line number Diff line change
@@ -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";

Expand All @@ -14,6 +17,7 @@ export type ActionGate = (
) => JudgeDecision | Promise<JudgeDecision>;

export class AgentActionDeniedError extends Error {
readonly _tag = "AgentActionDenied" as const;
readonly action: AgentBusEntry;

constructor(action: AgentBusEntry) {
Expand Down Expand Up @@ -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<T>(
bus: WriteAheadAgentBus,
actor: string,
Expand All @@ -53,34 +58,33 @@ export async function dispatchAction<T>(
): Promise<T> {
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<Record<string, unknown>>) =>
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 = <A>(fn: () => Promise<A>) => 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<T>(effect: Effect.Effect<T, unknown>): Promise<T> {
const exit = await Effect.runPromiseExit(effect);
if (Exit.isSuccess(exit)) return exit.value;
throw Cause.squash(exit.cause);
}
57 changes: 22 additions & 35 deletions packages/harness/src/sandbox.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -90,34 +91,26 @@ export class HarnessSandbox extends BaseSandbox {

async uploadFiles(files: Array<[string, Uint8Array]>): Promise<FileUploadResponse[]> {
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<FileUploadResponse>(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<FileDownloadResponse[]> {
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<FileDownloadResponse>(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<T>(kind: string, payload: unknown, execute: () => Promise<T>): Promise<T> {
Expand All @@ -140,12 +133,8 @@ export class HarnessSandbox extends BaseSandbox {
}
}

async function isSymlink(path: string): Promise<boolean> {
try {
return (await lstat(path)).isSymbolicLink();
} catch {
return false;
}
function isSymlink(path: string): Promise<boolean> {
return attemptAsync(async () => (await lstat(path)).isSymbolicLink(), () => false);
}

async function existingAncestor(path: string, root: string): Promise<string> {
Expand All @@ -154,11 +143,9 @@ async function existingAncestor(path: string, root: string): Promise<string> {
return current;
}

async function exists(path: string): Promise<boolean> {
try {
function exists(path: string): Promise<boolean> {
return attemptAsync(async () => {
await lstat(path);
return true;
} catch {
return false;
}
}, () => false);
}
1 change: 1 addition & 0 deletions packages/training/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -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"
},
Expand Down
27 changes: 27 additions & 0 deletions packages/training/src/attempt.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,27 @@
// Internal Effect-backed fallback helpers. Deliberately duplicated in each
// 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. The fallback receives the raw thrown value. */
export function attempt<T>(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<T>(fn: () => Promise<T>, fallback: (error: unknown) => T): Promise<T> {
return Effect.runPromise(
Effect.tryPromise({ try: fn, catch: (error) => error }).pipe(
Effect.catchAll((error) => Effect.sync(() => fallback(error))),
),
);
}
Comment thread
coderabbitai[bot] marked this conversation as resolved.
4 changes: 4 additions & 0 deletions packages/training/src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@ export type {
Activation,
AppliedPromotion,
CaptureSettings,
ErrorPhase,
EvolutionSettings,
PromotionApplier,
TrainInput,
Expand All @@ -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,
Expand Down
Loading
Loading