From acfc6792b7a437860724464913d3e8de944a0d7b Mon Sep 17 00:00:00 2001 From: James Sterling Date: Fri, 21 Aug 2026 14:00:32 +0000 Subject: [PATCH 01/11] feat: persist completed sample outcomes so interrupted runs resume instead of restarting Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- package.json | 1 + src/harness/run.test.ts | 38 +++++++++++----- src/harness/run.ts | 41 +++++++++++++++-- src/harness/sample-result-store.ts | 69 +++++++++++++++++++++++++++++ src/runner/run-by-id.ts | 11 +++++ test/helpers/noop-progress-layer.ts | 9 ++++ 6 files changed, 155 insertions(+), 14 deletions(-) create mode 100644 src/harness/sample-result-store.ts diff --git a/package.json b/package.json index 2e74d4b..9557fa3 100644 --- a/package.json +++ b/package.json @@ -32,6 +32,7 @@ "./parquet": "./src/results/parquet.ts", "./parquet-schema": "./src/results/parquet-schema.ts", "./progress": "./src/harness/progress.ts", + "./sample-result-store": "./src/harness/sample-result-store.ts", "./result-store": "./src/results/result-store.ts", "./internal/log": "./src/internal/log.ts", "./internal/effect-logger": "./src/internal/effect-logger.ts" diff --git a/src/harness/run.test.ts b/src/harness/run.test.ts index 803a334..efcce8a 100644 --- a/src/harness/run.test.ts +++ b/src/harness/run.test.ts @@ -16,6 +16,7 @@ import { fromChunk } from "effect/Stream"; import { noopProgressLayer, noopCheckpointLayer, + noopSampleResultLayer, } from "../../test/helpers/noop-progress-layer"; import { mcqScorer } from "../benchmarks/scorers/mcq/scorer"; import { runHarnessPromise } from "../internal/effect-logger"; @@ -27,6 +28,7 @@ import type { ModelService } from "./model"; import { Model } from "./model"; import type { CheckpointStore, ProgressReporter } from "./progress"; import { runBenchmark } from "./run"; +import type { SampleResultStore } from "./sample-result-store"; import { Scorer } from "./scorer"; import type { SolverService } from "./solver"; import { systemMessage, chain, generate, Solver } from "./solver"; @@ -94,14 +96,22 @@ function makeLayers( layer: Layer; }, solverService: ReturnType -): Layer { +): Layer< + | Dataset + | Solver + | Scorer + | ProgressReporter + | CheckpointStore + | SampleResultStore +> { return mergeAll( fakeDatasetLayer(SAMPLES), layerSucceed(Solver, Solver.of(solverService)), layerSucceed(Scorer, Scorer.of(mcqScorer)), model.layer, noopProgressLayer, - noopCheckpointLayer + noopCheckpointLayer, + noopSampleResultLayer ); } describe("runBenchmark", () => { @@ -130,7 +140,8 @@ describe("runBenchmark", () => { layerSucceed(Solver, Solver.of(solver)), layerSucceed(Scorer, Scorer.of(mcqScorer)), noopProgressLayer, - noopCheckpointLayer + noopCheckpointLayer, + noopSampleResultLayer ); await runHarnessPromise( runBenchmark({ @@ -184,7 +195,8 @@ describe("runBenchmark", () => { layerSucceed(Scorer, Scorer.of(mcqScorer)), model.layer, noopProgressLayer, - noopCheckpointLayer + noopCheckpointLayer, + noopSampleResultLayer ); const result = await runPromise( runBenchmark({ epochs: 1, maxConcurrency: 1 }).pipe(provide(layers)) @@ -225,7 +237,8 @@ describe("runBenchmark", () => { layerSucceed(Solver, Solver.of(solver)), layerSucceed(Scorer, Scorer.of(mcqScorer)), noopProgressLayer, - noopCheckpointLayer + noopCheckpointLayer, + noopSampleResultLayer ); const result = await runPromise( runBenchmark({ epochs: 1, maxConcurrency: 1 }).pipe(provide(layers)) @@ -252,7 +265,8 @@ describe("runBenchmark", () => { layerSucceed(Solver, Solver.of(solver)), layerSucceed(Scorer, Scorer.of(mcqScorer)), noopProgressLayer, - noopCheckpointLayer + noopCheckpointLayer, + noopSampleResultLayer ); const result = await runPromise( runBenchmark({ epochs: 1, maxConcurrency: 1 }).pipe(provide(layers)) @@ -268,7 +282,8 @@ describe("runBenchmark", () => { layerSucceed(Scorer, Scorer.of(mcqScorer)), model.layer, noopProgressLayer, - noopCheckpointLayer + noopCheckpointLayer, + noopSampleResultLayer ); const result = await runPromise( runBenchmark({ @@ -307,7 +322,8 @@ describe("runBenchmark", () => { layerSucceed(Scorer, Scorer.of(mcqScorer)), model.layer, noopProgressLayer, - noopCheckpointLayer + noopCheckpointLayer, + noopSampleResultLayer ); const result = await runPromise( runBenchmark({ epochs: 1, maxConcurrency: 2 }).pipe(provide(layers)) @@ -346,7 +362,8 @@ describe("runBenchmark", () => { layerSucceed(Scorer, Scorer.of(mcqScorer)), model.layer, noopProgressLayer, - noopCheckpointLayer + noopCheckpointLayer, + noopSampleResultLayer ); const result = await runPromise( runBenchmark({ epochs: 1, maxConcurrency: 2 }).pipe(provide(layers)) @@ -389,7 +406,8 @@ describe("runBenchmark", () => { layerSucceed(Scorer, Scorer.of(mcqScorer)), model.layer, noopProgressLayer, - noopCheckpointLayer + noopCheckpointLayer, + noopSampleResultLayer ); const result = await runPromise( runBenchmark({ epochs: 1, maxConcurrency: 2 }).pipe(provide(layers)) diff --git a/src/harness/run.ts b/src/harness/run.ts index 9135416..b946397 100644 --- a/src/harness/run.ts +++ b/src/harness/run.ts @@ -6,6 +6,8 @@ import { gen as effectGen, map as effectMap, annotateLogs, + orElseSucceed, + tryPromise, withLogSpan, succeed as effectSucceed, } from "effect/Effect"; @@ -42,6 +44,7 @@ import { Dataset } from "./dataset"; import type { AggregateMetrics, SampleScore } from "./metric"; import { aggregateScores } from "./metric"; import { CheckpointStore, ProgressReporter } from "./progress"; +import { SampleResultStore, sampleResultKey } from "./sample-result-store"; import { Scorer } from "./scorer"; import { Solver } from "./solver"; @@ -109,12 +112,12 @@ function evalWithProgress( evaluate: Effect< EvalOutcome, ModelError | SolverError, - Solver | Scorer | ProgressReporter | CheckpointStore + Solver | Scorer | ProgressReporter | CheckpointStore | SampleResultStore > ): Effect< EvalOutcome, ModelError | SolverError, - Solver | Scorer | ProgressReporter | CheckpointStore + Solver | Scorer | ProgressReporter | CheckpointStore | SampleResultStore > { const { sample, epoch, sampleIndex } = sampleEpoch; return effectGen(function* () { @@ -190,6 +193,31 @@ interface EvaluateOneOpts { readonly degradeSolverErrors: boolean; } +function evaluateOneResumable( + opts: EvaluateOneOpts +): Effect< + EvalOutcome, + ModelError | SolverError, + Solver | Scorer | ProgressReporter | CheckpointStore | SampleResultStore +> { + const { sample, epoch } = opts.sampleEpoch; + const key = sampleResultKey(sample.id, epoch); + return effectGen(function* () { + const store = yield* SampleResultStore; + const persisted = yield* tryPromise(() => store.read(key)).pipe( + orElseSucceed(() => null) + ); + if (persisted !== null) { + return persisted; + } + const outcome = yield* evaluateOne(opts); + yield* tryPromise(() => store.write(key, outcome)).pipe( + orElseSucceed(() => undefined) + ); + return outcome; + }); +} + function evaluateOne( opts: EvaluateOneOpts ): Effect< @@ -330,7 +358,12 @@ export function runBenchmark( ): Effect< RunResult, ModelError | SolverError | DatasetError, - Dataset | Solver | Scorer | ProgressReporter | CheckpointStore + | Dataset + | Solver + | Scorer + | ProgressReporter + | CheckpointStore + | SampleResultStore > { return Dataset.pipe( effectFlatMap((dataset) => { @@ -348,7 +381,7 @@ export function runBenchmark( (se) => evalWithProgress( se, - evaluateOne({ + evaluateOneResumable({ sampleEpoch: se, degradeSolverErrors: config.degradeSolverErrors ?? false, }).pipe( diff --git a/src/harness/sample-result-store.ts b/src/harness/sample-result-store.ts new file mode 100644 index 0000000..f04908b --- /dev/null +++ b/src/harness/sample-result-store.ts @@ -0,0 +1,69 @@ +import { Tag } from "effect/Context"; + +import { z } from "../internal/zod"; +import { ChatMessageSchema, ScoreValue } from "./core"; + +export const PersistedSampleOutcomeSchema = z.object({ + sampleScore: z.object({ + sampleId: z.string(), + epoch: z.number(), + score: z.object({ + value: z.enum([ + ScoreValue.Correct, + ScoreValue.Incorrect, + ScoreValue.Skipped, + ]), + answer: z.string().nullable(), + explanation: z.string(), + }), + messages: z.array(ChatMessageSchema).readonly().optional(), + responseItems: z + .array(z.record(z.string(), z.unknown()).readonly()) + .readonly() + .optional(), + requestBody: z.record(z.string(), z.unknown()).readonly().optional(), + generationIds: z.array(z.string()).readonly().optional(), + metadata: z.record(z.string(), z.unknown()).readonly().optional(), + input: z.string().optional(), + target: z.string().optional(), + }), + usage: z + .object({ + inputTokens: z.number().optional(), + outputTokens: z.number().optional(), + totalTokens: z.number().optional(), + reasoningTokens: z.number().optional(), + totalCost: z.number().optional(), + serverToolUse: z + .object({ + webSearchRequests: z.number().optional(), + toolCallsRequested: z.number().optional(), + toolCallsExecuted: z.number().optional(), + }) + .optional(), + }) + .optional(), + generationTimeMs: z.number().optional(), +}); + +export type PersistedSampleOutcome = z.infer< + typeof PersistedSampleOutcomeSchema +>; + +export function sampleResultKey(sampleId: string, epoch: number): string { + return `${sampleId}/${epoch}`; +} + +export interface SampleResultStoreService { + readonly read: (key: string) => Promise; + readonly write: (key: string, data: PersistedSampleOutcome) => Promise; +} + +export class SampleResultStore extends Tag( + "@openrouter/bench-harness/sample-result-store" +)() {} + +export const NOOP_SAMPLE_RESULT_STORE: SampleResultStoreService = { + read: async () => null, + write: async () => {}, +}; diff --git a/src/runner/run-by-id.ts b/src/runner/run-by-id.ts index a035212..99a5cdd 100644 --- a/src/runner/run-by-id.ts +++ b/src/runner/run-by-id.ts @@ -22,6 +22,11 @@ import { } from "../harness/progress"; import type { RunResult, RunConfig } from "../harness/run"; import { runBenchmark } from "../harness/run"; +import type { SampleResultStoreService } from "../harness/sample-result-store"; +import { + NOOP_SAMPLE_RESULT_STORE, + SampleResultStore, +} from "../harness/sample-result-store"; import { runHarnessPromise } from "../internal/effect-logger"; import type { AsyncEither } from "../internal/either"; import { Either } from "../internal/either"; @@ -50,6 +55,7 @@ export interface RunBenchmarkInput { readonly datasetRetry?: RetryConfig; readonly progressReporter?: ProgressReporterService; readonly checkpointStore?: CheckpointStoreService; + readonly sampleResultStore?: SampleResultStoreService; readonly abortSignal?: AbortSignal; readonly resultStore?: ResultStoreService; readonly maxOutputTokensCeiling?: number; @@ -91,6 +97,10 @@ export function runBenchmarkById( CheckpointStore, input.checkpointStore ?? NOOP_CHECKPOINT_STORE ); + const sampleResultLayer = layerSucceed( + SampleResultStore, + input.sampleResultStore ?? NOOP_SAMPLE_RESULT_STORE + ); const model = modelFromConfig(input.benchmarkConfig); const runConfig: RunConfig = { epochs: input.epochs, @@ -122,6 +132,7 @@ export function runBenchmarkById( fullBenchmarkLayer, progressLayer, checkpointLayer, + sampleResultLayer, resolverLayer ); const runOpts = diff --git a/test/helpers/noop-progress-layer.ts b/test/helpers/noop-progress-layer.ts index 04b819b..f98fde9 100644 --- a/test/helpers/noop-progress-layer.ts +++ b/test/helpers/noop-progress-layer.ts @@ -6,6 +6,10 @@ import { NOOP_PROGRESS_REPORTER, ProgressReporter, } from "../../src/harness/progress"; +import { + NOOP_SAMPLE_RESULT_STORE, + SampleResultStore, +} from "../../src/harness/sample-result-store"; export const noopProgressLayer = layerSucceed( ProgressReporter, @@ -16,3 +20,8 @@ export const noopCheckpointLayer = layerSucceed( CheckpointStore, NOOP_CHECKPOINT_STORE ); + +export const noopSampleResultLayer = layerSucceed( + SampleResultStore, + NOOP_SAMPLE_RESULT_STORE +); From 0b60aa8690a34a86ac07c19f9151d697b4941e36 Mon Sep 17 00:00:00 2001 From: James Sterling Date: Fri, 21 Aug 2026 14:51:01 +0000 Subject: [PATCH 02/11] test: cover sample-result resume behavior in the shared harness Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- src/harness/sample-result-resume.test.ts | 225 +++++++++++++++++++++++ 1 file changed, 225 insertions(+) create mode 100644 src/harness/sample-result-resume.test.ts diff --git a/src/harness/sample-result-resume.test.ts b/src/harness/sample-result-resume.test.ts new file mode 100644 index 0000000..29aa253 --- /dev/null +++ b/src/harness/sample-result-resume.test.ts @@ -0,0 +1,225 @@ +import { describe, expect, it } from "bun:test"; + +import { fromIterable } from "effect/Chunk"; +import { + flatMap as effectFlatMap, + runPromise, + succeed as effectSucceed, + provide, +} from "effect/Effect"; +import type { Layer } from "effect/Layer"; +import { mergeAll, succeed as layerSucceed } from "effect/Layer"; +import { fromChunk } from "effect/Stream"; + +import { + noopProgressLayer, + noopCheckpointLayer, +} from "../../test/helpers/noop-progress-layer"; +import { mcqScorer } from "../benchmarks/scorers/mcq/scorer"; +import { recordGenerationId } from "../runtime/generation-ids"; +import type { Sample } from "./core"; +import { MessageRole } from "./core"; +import { Dataset } from "./dataset"; +import type { ModelService } from "./model"; +import type { RunResult } from "./run"; +import { runBenchmark } from "./run"; +import type { SampleResultStoreService } from "./sample-result-store"; +import { + PersistedSampleOutcomeSchema, + SampleResultStore, +} from "./sample-result-store"; +import { Scorer } from "./scorer"; +import type { ScorerService } from "./scorer"; +import { systemMessage, chain, generate, Solver } from "./solver"; + +const SAMPLES: readonly Sample[] = [ + { id: "s1", input: "Q1 target B", target: { text: "B" } }, + { id: "s2", input: "Q2 target B", target: { text: "B" } }, + { id: "s3", input: "Q3 target B", target: { text: "B" } }, +]; +const EPOCHS = 2; +const ALL_KEYS = ["s1/0", "s1/1", "s2/0", "s2/1", "s3/0", "s3/1"]; + +function fakeDatasetLayer(samples: readonly Sample[]): Layer { + return layerSucceed( + Dataset, + Dataset.of({ + stream: (opts) => { + const start = opts?.start ?? 0; + const end = opts?.end ?? samples.length; + return fromChunk(fromIterable(samples.slice(start, end))); + }, + size: effectSucceed(samples.length), + }) + ); +} + +function fakeModel(): ModelService { + return { + generate: (messages) => { + const userMsg = + messages.find((m) => m.role === MessageRole.User)?.content ?? ""; + const completion = userMsg.includes("Q2") ? "Answer: A" : "Answer: B"; + return recordGenerationId(`fake-${userMsg}`).pipe( + effectFlatMap(() => + effectSucceed({ + completion, + message: { role: MessageRole.Assistant, content: completion }, + usage: { + inputTokens: 10, + outputTokens: 5, + totalTokens: 15, + totalCost: 0.001, + }, + generationTimeMs: 100, + }) + ) + ); + }, + }; +} + +interface Counters { + solver: number; + scorer: number; +} + +interface InMemoryStore extends SampleResultStoreService { + readonly map: Map; + readonly writtenKeys: string[]; +} + +function makeInMemoryStore(seed?: Map): InMemoryStore { + const map = new Map(seed ?? []); + const writtenKeys: string[] = []; + return { + map, + writtenKeys, + read: async (key) => { + const raw = map.get(key); + if (raw === undefined) { + return null; + } + const parsed = PersistedSampleOutcomeSchema.safeParse(JSON.parse(raw)); + return parsed.success ? parsed.data : null; + }, + write: async (key, data) => { + writtenKeys.push(key); + map.set(key, JSON.stringify(data)); + }, + }; +} + +function throwingStore(): SampleResultStoreService { + return { + read: async () => { + throw new Error("read exploded"); + }, + write: async () => { + throw new Error("write exploded"); + }, + }; +} + +async function runWithStore( + store: SampleResultStoreService +): Promise<{ result: RunResult; counters: Counters }> { + const counters: Counters = { solver: 0, scorer: 0 }; + const model = fakeModel(); + const innerSolver = chain( + systemMessage("You are a helpful assistant."), + generate(model, { temperature: 0.5 }) + ); + const countingSolver = Solver.of((state) => { + counters.solver += 1; + return innerSolver(state); + }); + const countingScorer: ScorerService = (state, target) => { + counters.scorer += 1; + return mcqScorer(state, target); + }; + const layers = mergeAll( + fakeDatasetLayer(SAMPLES), + layerSucceed(Solver, countingSolver), + layerSucceed(Scorer, Scorer.of(countingScorer)), + noopProgressLayer, + noopCheckpointLayer, + layerSucceed(SampleResultStore, store) + ); + const result = await runPromise( + runBenchmark({ epochs: EPOCHS, maxConcurrency: 2 }).pipe(provide(layers)) + ); + return { result, counters }; +} + +function summarize(result: RunResult): { key: string; value: string }[] { + return result.sampleScores + .map((s) => ({ key: `${s.sampleId}/${s.epoch}`, value: s.score.value })) + .sort((a, b) => a.key.localeCompare(b.key)); +} + +let run1Store: InMemoryStore; +let run1Summary: { key: string; value: string }[]; +let run1Metrics: { accuracy: number; total: number; correct: number }; + +describe("sample-result resume", () => { + it("T1: fresh run evaluates all 6 and persists all 6 keys", async () => { + run1Store = makeInMemoryStore(); + const { result, counters } = await runWithStore(run1Store); + expect(result.sampleScores.length).toBe(6); + expect(counters.solver).toBe(6); + expect(counters.scorer).toBe(6); + expect([...run1Store.map.keys()].sort()).toEqual(ALL_KEYS); + for (const raw of run1Store.map.values()) { + const parsed = PersistedSampleOutcomeSchema.safeParse(JSON.parse(raw)); + expect(parsed.success).toBe(true); + } + run1Summary = summarize(result); + run1Metrics = { + accuracy: result.metrics.accuracy, + total: result.metrics.totalQuestions, + correct: result.metrics.correctAnswers, + }; + expect(run1Metrics.total).toBe(3); + expect(run1Metrics.accuracy).toBeCloseTo(2 / 3, 5); + }); + + it("T2: full resume — solver/scorer never invoked, identical outcome", async () => { + const store = makeInMemoryStore(run1Store.map); + const { result, counters } = await runWithStore(store); + expect(counters.solver).toBe(0); + expect(counters.scorer).toBe(0); + expect(summarize(result)).toEqual(run1Summary); + expect(result.metrics.accuracy).toBeCloseTo(run1Metrics.accuracy, 5); + expect(result.metrics.totalQuestions).toBe(run1Metrics.total); + expect(result.metrics.correctAnswers).toBe(run1Metrics.correct); + expect(store.writtenKeys).toEqual([]); + }); + + it("T3: partial resume — only the 4 missing keys are evaluated", async () => { + const seeded = new Map( + [...run1Store.map.entries()].filter(([k]) => ["s1/0", "s2/1"].includes(k)) + ); + const store = makeInMemoryStore(seeded); + const { result, counters } = await runWithStore(store); + expect(counters.solver).toBe(4); + expect(counters.scorer).toBe(4); + expect([...store.writtenKeys].sort()).toEqual([ + "s1/1", + "s2/0", + "s3/0", + "s3/1", + ]); + expect(summarize(result)).toEqual(run1Summary); + expect(result.metrics.accuracy).toBeCloseTo(run1Metrics.accuracy, 5); + }); + + it("T4: store whose read/write throw — run completes, all evaluated", async () => { + const { result, counters } = await runWithStore(throwingStore()); + expect(counters.solver).toBe(6); + expect(counters.scorer).toBe(6); + expect(result.sampleScores.length).toBe(6); + expect(summarize(result)).toEqual(run1Summary); + expect(result.metrics.accuracy).toBeCloseTo(run1Metrics.accuracy, 5); + }); +}); From c55e2bc799a9077ba63611231aea8f9920df0681 Mon Sep 17 00:00:00 2001 From: James Sterling Date: Mon, 24 Aug 2026 22:18:13 +0000 Subject: [PATCH 03/11] fix(harness): log sample-result store read/write failures Merges main and replaces silent orElseSucceed fallbacks in evaluateOneResumable with wLog'd catchAll handlers so dropped persistence (and therefore lost resume coverage) is observable. Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- src/harness/run.ts | 20 +++++++++++++--- src/harness/sample-result-resume.test.ts | 30 +++++++++++++++++------- 2 files changed, 39 insertions(+), 11 deletions(-) diff --git a/src/harness/run.ts b/src/harness/run.ts index b946397..917346f 100644 --- a/src/harness/run.ts +++ b/src/harness/run.ts @@ -1,12 +1,12 @@ import type { Effect } from "effect/Effect"; import { + catchAll, catchTags, fail as effectFail, flatMap as effectFlatMap, gen as effectGen, map as effectMap, annotateLogs, - orElseSucceed, tryPromise, withLogSpan, succeed as effectSucceed, @@ -20,6 +20,8 @@ import { zipWithIndex as streamZipWithIndex, } from "effect/Stream"; +import { unknownErrorToString } from "../internal/errors"; +import { wLog } from "../internal/log"; import { resetGenerationIds } from "../runtime/generation-ids"; import type { ReplayedUsage } from "../runtime/generation-resolver"; import { resolveCollectedGenerations } from "../runtime/generation-resolver"; @@ -205,14 +207,26 @@ function evaluateOneResumable( return effectGen(function* () { const store = yield* SampleResultStore; const persisted = yield* tryPromise(() => store.read(key)).pipe( - orElseSucceed(() => null) + catchAll((error) => { + wLog("Failed to read persisted sample outcome; re-evaluating", { + sample_result_key: key, + error: unknownErrorToString(error), + }); + return effectSucceed(null); + }) ); if (persisted !== null) { return persisted; } const outcome = yield* evaluateOne(opts); yield* tryPromise(() => store.write(key, outcome)).pipe( - orElseSucceed(() => undefined) + catchAll((error) => { + wLog("Failed to persist sample outcome; a resumed run will re-run it", { + sample_result_key: key, + error: unknownErrorToString(error), + }); + return effectSucceed(undefined); + }) ); return outcome; }); diff --git a/src/harness/sample-result-resume.test.ts b/src/harness/sample-result-resume.test.ts index 29aa253..2ad99c8 100644 --- a/src/harness/sample-result-resume.test.ts +++ b/src/harness/sample-result-resume.test.ts @@ -1,4 +1,4 @@ -import { describe, expect, it } from "bun:test"; +import { describe, expect, it, spyOn } from "bun:test"; import { fromIterable } from "effect/Chunk"; import { @@ -214,12 +214,26 @@ describe("sample-result resume", () => { expect(result.metrics.accuracy).toBeCloseTo(run1Metrics.accuracy, 5); }); - it("T4: store whose read/write throw — run completes, all evaluated", async () => { - const { result, counters } = await runWithStore(throwingStore()); - expect(counters.solver).toBe(6); - expect(counters.scorer).toBe(6); - expect(result.sampleScores.length).toBe(6); - expect(summarize(result)).toEqual(run1Summary); - expect(result.metrics.accuracy).toBeCloseTo(run1Metrics.accuracy, 5); + it("T4: store whose read/write throw — run completes, all evaluated, failures logged", async () => { + const warn = spyOn(console, "warn").mockImplementation(() => {}); + try { + const { result, counters } = await runWithStore(throwingStore()); + expect(counters.solver).toBe(6); + expect(counters.scorer).toBe(6); + expect(result.sampleScores.length).toBe(6); + expect(summarize(result)).toEqual(run1Summary); + expect(result.metrics.accuracy).toBeCloseTo(run1Metrics.accuracy, 5); + const warned = warn.mock.calls.map((call) => String(call[0])); + expect( + warned.some((line) => + line.includes("Failed to read persisted sample outcome") + ) + ).toBe(true); + expect( + warned.some((line) => line.includes("Failed to persist sample outcome")) + ).toBe(true); + } finally { + warn.mockRestore(); + } }); }); From da40a0696a25a63a33007fe065ec80df4d0a64c0 Mon Sep 17 00:00:00 2001 From: James Sterling Date: Mon, 24 Aug 2026 22:25:57 +0000 Subject: [PATCH 04/11] fix(harness): persist scorer trajectory, skip degraded outcomes, surface store error causes Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- src/harness/run.ts | 19 +++- src/harness/sample-result-resume.test.ts | 89 +++++++++++++++- src/harness/sample-result-store.ts | 125 +++++++++++++++-------- 3 files changed, 180 insertions(+), 53 deletions(-) diff --git a/src/harness/run.ts b/src/harness/run.ts index 917346f..09f465d 100644 --- a/src/harness/run.ts +++ b/src/harness/run.ts @@ -82,6 +82,7 @@ type EvalOutcome = { sampleScore: SampleScore; usage?: ModelUsage; generationTimeMs?: number; + isDegraded?: boolean; }; function sampleEpochStream( @@ -206,11 +207,14 @@ function evaluateOneResumable( const key = sampleResultKey(sample.id, epoch); return effectGen(function* () { const store = yield* SampleResultStore; - const persisted = yield* tryPromise(() => store.read(key)).pipe( + const persisted = yield* tryPromise({ + try: () => store.read(key), + catch: unknownErrorToString, + }).pipe( catchAll((error) => { wLog("Failed to read persisted sample outcome; re-evaluating", { sample_result_key: key, - error: unknownErrorToString(error), + error, }); return effectSucceed(null); }) @@ -219,11 +223,17 @@ function evaluateOneResumable( return persisted; } const outcome = yield* evaluateOne(opts); - yield* tryPromise(() => store.write(key, outcome)).pipe( + if (outcome.isDegraded === true) { + return outcome; + } + yield* tryPromise({ + try: () => store.write(key, outcome), + catch: unknownErrorToString, + }).pipe( catchAll((error) => { wLog("Failed to persist sample outcome; a resumed run will re-run it", { sample_result_key: key, - error: unknownErrorToString(error), + error, }); return effectSucceed(undefined); }) @@ -355,6 +365,7 @@ function errorOutcome(opts: ErrorOutcomeOpts): EvalOutcome { input: sample.input, target: sample.target.text, }, + isDegraded: true, }; } diff --git a/src/harness/sample-result-resume.test.ts b/src/harness/sample-result-resume.test.ts index 2ad99c8..97a35ed 100644 --- a/src/harness/sample-result-resume.test.ts +++ b/src/harness/sample-result-resume.test.ts @@ -2,6 +2,7 @@ import { describe, expect, it, spyOn } from "bun:test"; import { fromIterable } from "effect/Chunk"; import { + fail as effectFail, flatMap as effectFlatMap, runPromise, succeed as effectSucceed, @@ -18,7 +19,7 @@ import { import { mcqScorer } from "../benchmarks/scorers/mcq/scorer"; import { recordGenerationId } from "../runtime/generation-ids"; import type { Sample } from "./core"; -import { MessageRole } from "./core"; +import { MessageRole, ModelError, ScoreValue } from "./core"; import { Dataset } from "./dataset"; import type { ModelService } from "./model"; import type { RunResult } from "./run"; @@ -122,10 +123,14 @@ function throwingStore(): SampleResultStoreService { } async function runWithStore( - store: SampleResultStoreService + store: SampleResultStoreService, + overrides?: { + readonly model?: ModelService; + readonly scorer?: ScorerService; + } ): Promise<{ result: RunResult; counters: Counters }> { const counters: Counters = { solver: 0, scorer: 0 }; - const model = fakeModel(); + const model = overrides?.model ?? fakeModel(); const innerSolver = chain( systemMessage("You are a helpful assistant."), generate(model, { temperature: 0.5 }) @@ -134,9 +139,10 @@ async function runWithStore( counters.solver += 1; return innerSolver(state); }); + const baseScorer = overrides?.scorer ?? mcqScorer; const countingScorer: ScorerService = (state, target) => { counters.scorer += 1; - return mcqScorer(state, target); + return baseScorer(state, target); }; const layers = mergeAll( fakeDatasetLayer(SAMPLES), @@ -223,7 +229,9 @@ describe("sample-result resume", () => { expect(result.sampleScores.length).toBe(6); expect(summarize(result)).toEqual(run1Summary); expect(result.metrics.accuracy).toBeCloseTo(run1Metrics.accuracy, 5); - const warned = warn.mock.calls.map((call) => String(call[0])); + const warned = warn.mock.calls.map((call) => + call.map((arg) => JSON.stringify(arg)).join(" ") + ); expect( warned.some((line) => line.includes("Failed to read persisted sample outcome") @@ -232,8 +240,79 @@ describe("sample-result resume", () => { expect( warned.some((line) => line.includes("Failed to persist sample outcome")) ).toBe(true); + expect(warned.some((line) => line.includes("read exploded"))).toBe(true); + expect(warned.some((line) => line.includes("write exploded"))).toBe(true); } finally { warn.mockRestore(); } }); + + it("T5: degraded error outcomes are not persisted and are re-evaluated on resume", async () => { + const failingModel: ModelService = { + generate: (messages) => { + const userMsg = + messages.find((m) => m.role === MessageRole.User)?.content ?? ""; + if (userMsg.includes("Q1")) { + return effectFail( + new ModelError({ message: "OpenRouter HTTP 429", status: 429 }) + ); + } + return effectSucceed({ + completion: "Answer: B", + message: { role: MessageRole.Assistant, content: "Answer: B" }, + generationTimeMs: 100, + }); + }, + }; + const store = makeInMemoryStore(); + const run1 = await runWithStore(store, { model: failingModel }); + const skippedRun1 = run1.result.sampleScores.filter( + (s) => s.sampleId === "s1" + ); + expect(skippedRun1.every((s) => s.score.value === ScoreValue.Skipped)).toBe( + true + ); + expect([...store.map.keys()].sort()).toEqual([ + "s2/0", + "s2/1", + "s3/0", + "s3/1", + ]); + + const run2 = await runWithStore(store); + expect(run2.counters.solver).toBe(2); + expect(run2.counters.scorer).toBe(2); + const s1Run2 = run2.result.sampleScores.filter((s) => s.sampleId === "s1"); + expect(s1Run2.every((s) => s.score.value === ScoreValue.Correct)).toBe( + true + ); + expect([...store.map.keys()].sort()).toEqual(ALL_KEYS); + }); + + it("T6: scorer trajectory survives persistence and resume", async () => { + const trajectoryScorer: ScorerService = (state, target) => + mcqScorer(state, target).pipe( + effectFlatMap((score) => + effectSucceed({ + ...score, + trajectory: { + kind: "verifier_log" as const, + log: "pytest: 12 passed", + }, + }) + ) + ); + const store = makeInMemoryStore(); + await runWithStore(store, { scorer: trajectoryScorer }); + + const resumed = await runWithStore(store); + expect(resumed.counters.solver).toBe(0); + expect( + resumed.result.sampleScores.every( + (s) => + s.score.trajectory?.kind === "verifier_log" && + s.score.trajectory.log === "pytest: 12 passed" + ) + ).toBe(true); + }); }); diff --git a/src/harness/sample-result-store.ts b/src/harness/sample-result-store.ts index f04908b..0c4d505 100644 --- a/src/harness/sample-result-store.ts +++ b/src/harness/sample-result-store.ts @@ -1,54 +1,91 @@ import { Tag } from "effect/Context"; +import type { ZodShape } from "../internal/zod"; import { z } from "../internal/zod"; +import type { + ModelUsage, + ResponseItem, + Score, + ScorerTrajectory, + ServerToolUseCounts, +} from "./core"; import { ChatMessageSchema, ScoreValue } from "./core"; +import type { SampleScore } from "./metric"; -export const PersistedSampleOutcomeSchema = z.object({ - sampleScore: z.object({ - sampleId: z.string(), - epoch: z.number(), - score: z.object({ - value: z.enum([ - ScoreValue.Correct, - ScoreValue.Incorrect, - ScoreValue.Skipped, - ]), - answer: z.string().nullable(), - explanation: z.string(), - }), - messages: z.array(ChatMessageSchema).readonly().optional(), - responseItems: z - .array(z.record(z.string(), z.unknown()).readonly()) - .readonly() - .optional(), - requestBody: z.record(z.string(), z.unknown()).readonly().optional(), - generationIds: z.array(z.string()).readonly().optional(), - metadata: z.record(z.string(), z.unknown()).readonly().optional(), - input: z.string().optional(), - target: z.string().optional(), - }), - usage: z - .object({ - inputTokens: z.number().optional(), - outputTokens: z.number().optional(), - totalTokens: z.number().optional(), - reasoningTokens: z.number().optional(), - totalCost: z.number().optional(), - serverToolUse: z - .object({ - webSearchRequests: z.number().optional(), - toolCallsRequested: z.number().optional(), - toolCallsExecuted: z.number().optional(), - }) - .optional(), - }) - .optional(), - generationTimeMs: z.number().optional(), -}); +export interface PersistedSampleOutcome { + readonly sampleScore: SampleScore; + readonly usage?: ModelUsage; + readonly generationTimeMs?: number; +} + +const ScoreValueSchema = z.enum([ + ScoreValue.Correct, + ScoreValue.Incorrect, + ScoreValue.Skipped, +]); -export type PersistedSampleOutcome = z.infer< - typeof PersistedSampleOutcomeSchema +type VerifierLogTrajectory = Extract< + ScorerTrajectory, + { kind: "verifier_log" } >; +type JudgeRunsTrajectory = Extract; + +const ScorerTrajectorySchema: z.ZodType = + z.discriminatedUnion("kind", [ + z.object({ + kind: z.literal("verifier_log"), + log: z.string(), + } satisfies ZodShape), + z.object({ + kind: z.literal("judge_runs"), + runs: z.array(z.unknown()).readonly(), + } satisfies ZodShape), + ]); + +const ScoreSchema = z.object({ + value: ScoreValueSchema, + answer: z.string().nullable(), + explanation: z.string(), + trajectory: ScorerTrajectorySchema.optional(), +} satisfies ZodShape); + +const ResponseItemSchema: z.ZodType = z + .record(z.string(), z.unknown()) + .readonly(); + +const SampleScoreSchema = z.object({ + sampleId: z.string(), + epoch: z.number(), + score: ScoreSchema, + messages: z.array(ChatMessageSchema).readonly().optional(), + responseItems: z.array(ResponseItemSchema).readonly().optional(), + requestBody: z.record(z.string(), z.unknown()).readonly().optional(), + generationIds: z.array(z.string()).readonly().optional(), + metadata: z.record(z.string(), z.unknown()).readonly().optional(), + input: z.string().optional(), + target: z.string().optional(), +} satisfies ZodShape); + +const ServerToolUseCountsSchema = z.object({ + webSearchRequests: z.number().optional(), + toolCallsRequested: z.number().optional(), + toolCallsExecuted: z.number().optional(), +} satisfies ZodShape); + +const ModelUsageSchema = z.object({ + inputTokens: z.number().optional(), + outputTokens: z.number().optional(), + totalTokens: z.number().optional(), + reasoningTokens: z.number().optional(), + totalCost: z.number().optional(), + serverToolUse: ServerToolUseCountsSchema.optional(), +} satisfies ZodShape); + +export const PersistedSampleOutcomeSchema = z.object({ + sampleScore: SampleScoreSchema, + usage: ModelUsageSchema.optional(), + generationTimeMs: z.number().optional(), +} satisfies ZodShape); export function sampleResultKey(sampleId: string, epoch: number): string { return `${sampleId}/${epoch}`; From f03ffb01b5497e858ee2a0a19a77c4821e8ad940 Mon Sep 17 00:00:00 2001 From: James Sterling Date: Mon, 24 Aug 2026 22:50:57 +0000 Subject: [PATCH 05/11] fix(internal): make ZodShape require optional keys so schemas cannot silently omit them Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- src/internal/zod.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/internal/zod.ts b/src/internal/zod.ts index 772cc0f..b53c8d2 100644 --- a/src/internal/zod.ts +++ b/src/internal/zod.ts @@ -5,7 +5,7 @@ import { Either } from "./either"; export { z }; export type ZodShape = { - [key in keyof T]: ZodType; + [key in keyof Required]: ZodType; }; export function parseSchema( From c462e9355a8f6e85ab0a8dd416774d4e7a44bad8 Mon Sep 17 00:00:00 2001 From: James Sterling Date: Mon, 24 Aug 2026 22:58:12 +0000 Subject: [PATCH 06/11] refactor(harness): use zInt() for persisted epoch field Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- src/harness/sample-result-store.ts | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/src/harness/sample-result-store.ts b/src/harness/sample-result-store.ts index 0c4d505..a8d2f59 100644 --- a/src/harness/sample-result-store.ts +++ b/src/harness/sample-result-store.ts @@ -1,7 +1,7 @@ import { Tag } from "effect/Context"; import type { ZodShape } from "../internal/zod"; -import { z } from "../internal/zod"; +import { z, zInt } from "../internal/zod"; import type { ModelUsage, ResponseItem, @@ -55,7 +55,7 @@ const ResponseItemSchema: z.ZodType = z const SampleScoreSchema = z.object({ sampleId: z.string(), - epoch: z.number(), + epoch: zInt(), score: ScoreSchema, messages: z.array(ChatMessageSchema).readonly().optional(), responseItems: z.array(ResponseItemSchema).readonly().optional(), From e4b574ff4862cc1261bb1efe455d8dd75158ceec Mon Sep 17 00:00:00 2001 From: James Sterling Date: Tue, 25 Aug 2026 03:41:49 +0000 Subject: [PATCH 07/11] refactor(harness): use zInt() for persisted token and tool-use count fields Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- src/harness/sample-result-store.ts | 14 +++++++------- 1 file changed, 7 insertions(+), 7 deletions(-) diff --git a/src/harness/sample-result-store.ts b/src/harness/sample-result-store.ts index a8d2f59..6cb9152 100644 --- a/src/harness/sample-result-store.ts +++ b/src/harness/sample-result-store.ts @@ -67,16 +67,16 @@ const SampleScoreSchema = z.object({ } satisfies ZodShape); const ServerToolUseCountsSchema = z.object({ - webSearchRequests: z.number().optional(), - toolCallsRequested: z.number().optional(), - toolCallsExecuted: z.number().optional(), + webSearchRequests: zInt().optional(), + toolCallsRequested: zInt().optional(), + toolCallsExecuted: zInt().optional(), } satisfies ZodShape); const ModelUsageSchema = z.object({ - inputTokens: z.number().optional(), - outputTokens: z.number().optional(), - totalTokens: z.number().optional(), - reasoningTokens: z.number().optional(), + inputTokens: zInt().optional(), + outputTokens: zInt().optional(), + totalTokens: zInt().optional(), + reasoningTokens: zInt().optional(), totalCost: z.number().optional(), serverToolUse: ServerToolUseCountsSchema.optional(), } satisfies ZodShape); From e26ebb9a04395c663ee29a5e4852bfc878d39d36 Mon Sep 17 00:00:00 2001 From: James Sterling Date: Tue, 25 Aug 2026 03:49:47 +0000 Subject: [PATCH 08/11] fix(harness): validate persisted sample outcomes in the harness and log schema rejections Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- src/harness/run.ts | 21 ++++++++--- src/harness/sample-result-resume.test.ts | 44 +++++++++++++++++++++--- src/harness/sample-result-store.ts | 2 +- 3 files changed, 57 insertions(+), 10 deletions(-) diff --git a/src/harness/run.ts b/src/harness/run.ts index 09f465d..fa86c1c 100644 --- a/src/harness/run.ts +++ b/src/harness/run.ts @@ -20,8 +20,10 @@ import { zipWithIndex as streamZipWithIndex, } from "effect/Stream"; +import { Either } from "../internal/either"; import { unknownErrorToString } from "../internal/errors"; import { wLog } from "../internal/log"; +import { firstZodIssueMessage, parseSchema } from "../internal/zod"; import { resetGenerationIds } from "../runtime/generation-ids"; import type { ReplayedUsage } from "../runtime/generation-resolver"; import { resolveCollectedGenerations } from "../runtime/generation-resolver"; @@ -46,7 +48,11 @@ import { Dataset } from "./dataset"; import type { AggregateMetrics, SampleScore } from "./metric"; import { aggregateScores } from "./metric"; import { CheckpointStore, ProgressReporter } from "./progress"; -import { SampleResultStore, sampleResultKey } from "./sample-result-store"; +import { + PersistedSampleOutcomeSchema, + SampleResultStore, + sampleResultKey, +} from "./sample-result-store"; import { Scorer } from "./scorer"; import { Solver } from "./solver"; @@ -207,7 +213,7 @@ function evaluateOneResumable( const key = sampleResultKey(sample.id, epoch); return effectGen(function* () { const store = yield* SampleResultStore; - const persisted = yield* tryPromise({ + const raw = yield* tryPromise({ try: () => store.read(key), catch: unknownErrorToString, }).pipe( @@ -219,8 +225,15 @@ function evaluateOneResumable( return effectSucceed(null); }) ); - if (persisted !== null) { - return persisted; + if (raw !== null && raw !== undefined) { + const persisted = parseSchema(PersistedSampleOutcomeSchema, raw); + if (Either.isRight(persisted)) { + return persisted.right; + } + wLog("Persisted sample outcome failed validation; re-evaluating", { + sample_result_key: key, + error: firstZodIssueMessage(persisted.left), + }); } const outcome = yield* evaluateOne(opts); if (outcome.isDegraded === true) { diff --git a/src/harness/sample-result-resume.test.ts b/src/harness/sample-result-resume.test.ts index 97a35ed..0cb3cde 100644 --- a/src/harness/sample-result-resume.test.ts +++ b/src/harness/sample-result-resume.test.ts @@ -98,11 +98,7 @@ function makeInMemoryStore(seed?: Map): InMemoryStore { writtenKeys, read: async (key) => { const raw = map.get(key); - if (raw === undefined) { - return null; - } - const parsed = PersistedSampleOutcomeSchema.safeParse(JSON.parse(raw)); - return parsed.success ? parsed.data : null; + return raw === undefined ? null : JSON.parse(raw); }, write: async (key, data) => { writtenKeys.push(key); @@ -289,6 +285,44 @@ describe("sample-result resume", () => { expect([...store.map.keys()].sort()).toEqual(ALL_KEYS); }); + it("T7: unparseable persisted record is logged and re-evaluated", async () => { + const seeded = new Map(run1Store.map); + const corruptRaw = seeded.get("s1/0"); + if (corruptRaw === undefined) { + throw new Error("expected s1/0 in run1 store"); + } + const corrupt: unknown = JSON.parse(corruptRaw); + if (typeof corrupt !== "object" || corrupt === null) { + throw new Error("expected persisted record to be an object"); + } + seeded.set( + "s1/0", + JSON.stringify({ ...corrupt, usage: { inputTokens: 10.5 } }) + ); + const store = makeInMemoryStore(seeded); + const warn = spyOn(console, "warn").mockImplementation(() => {}); + try { + const { result, counters } = await runWithStore(store); + expect(counters.solver).toBe(1); + expect(counters.scorer).toBe(1); + expect(summarize(result)).toEqual(run1Summary); + expect(store.writtenKeys).toEqual(["s1/0"]); + const warned = warn.mock.calls.map((call) => + call.map((arg) => JSON.stringify(arg)).join(" ") + ); + expect( + warned.some( + (line) => + line.includes("Persisted sample outcome failed validation") && + line.includes("s1/0") && + line.includes("inputTokens") + ) + ).toBe(true); + } finally { + warn.mockRestore(); + } + }); + it("T6: scorer trajectory survives persistence and resume", async () => { const trajectoryScorer: ScorerService = (state, target) => mcqScorer(state, target).pipe( diff --git a/src/harness/sample-result-store.ts b/src/harness/sample-result-store.ts index 6cb9152..f0f24de 100644 --- a/src/harness/sample-result-store.ts +++ b/src/harness/sample-result-store.ts @@ -92,7 +92,7 @@ export function sampleResultKey(sampleId: string, epoch: number): string { } export interface SampleResultStoreService { - readonly read: (key: string) => Promise; + readonly read: (key: string) => Promise; readonly write: (key: string, data: PersistedSampleOutcome) => Promise; } From 3294d70dfb1da8d5e2eab4cd98f6632d24e2b41f Mon Sep 17 00:00:00 2001 From: James Sterling Date: Tue, 25 Aug 2026 17:26:32 +0000 Subject: [PATCH 09/11] fix(runner): namespace sample result keys by sessionId to avoid cross-run collisions Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- src/harness/sample-result-store.test.ts | 40 +++++++++++++++++++++++++ src/harness/sample-result-store.ts | 10 +++++++ src/runner/run-by-id.ts | 5 +++- 3 files changed, 54 insertions(+), 1 deletion(-) create mode 100644 src/harness/sample-result-store.test.ts diff --git a/src/harness/sample-result-store.test.ts b/src/harness/sample-result-store.test.ts new file mode 100644 index 0000000..b9f1ea2 --- /dev/null +++ b/src/harness/sample-result-store.test.ts @@ -0,0 +1,40 @@ +import { describe, expect, it } from "bun:test"; + +import { ScoreValue } from "./core"; +import type { + PersistedSampleOutcome, + SampleResultStoreService, +} from "./sample-result-store"; +import { + namespacedSampleResultStore, + sampleResultKey, +} from "./sample-result-store"; + +const OUTCOME: PersistedSampleOutcome = { + sampleScore: { + sampleId: "s1", + epoch: 0, + score: { value: ScoreValue.Correct, answer: "B", explanation: "" }, + }, +}; + +describe("namespacedSampleResultStore", () => { + it("prefixes reads and writes so different sessions do not collide", async () => { + const map = new Map(); + const backing: SampleResultStoreService = { + read: async (key) => map.get(key) ?? null, + write: async (key, data) => { + map.set(key, data); + }, + }; + const sessionA = namespacedSampleResultStore("session-a", backing); + const sessionB = namespacedSampleResultStore("session-b", backing); + const key = sampleResultKey("s1", 0); + + await sessionA.write(key, OUTCOME); + + expect([...map.keys()]).toEqual(["session-a/s1/0"]); + expect(await sessionA.read(key)).toEqual(OUTCOME); + expect(await sessionB.read(key)).toBeNull(); + }); +}); diff --git a/src/harness/sample-result-store.ts b/src/harness/sample-result-store.ts index f0f24de..677d917 100644 --- a/src/harness/sample-result-store.ts +++ b/src/harness/sample-result-store.ts @@ -104,3 +104,13 @@ export const NOOP_SAMPLE_RESULT_STORE: SampleResultStoreService = { read: async () => null, write: async () => {}, }; + +export function namespacedSampleResultStore( + prefix: string, + store: SampleResultStoreService +): SampleResultStoreService { + return { + read: (key) => store.read(`${prefix}/${key}`), + write: (key, data) => store.write(`${prefix}/${key}`, data), + }; +} diff --git a/src/runner/run-by-id.ts b/src/runner/run-by-id.ts index 6a37bbf..1863cc6 100644 --- a/src/runner/run-by-id.ts +++ b/src/runner/run-by-id.ts @@ -35,6 +35,7 @@ import type { RunResult, RunConfig } from "../harness/run"; import { runBenchmark } from "../harness/run"; import type { SampleResultStoreService } from "../harness/sample-result-store"; import { + namespacedSampleResultStore, NOOP_SAMPLE_RESULT_STORE, SampleResultStore, } from "../harness/sample-result-store"; @@ -96,7 +97,9 @@ export function runBenchmarkById( ); const sampleResultLayer = layerSucceed( SampleResultStore, - input.sampleResultStore ?? NOOP_SAMPLE_RESULT_STORE + input.sampleResultStore === undefined + ? NOOP_SAMPLE_RESULT_STORE + : namespacedSampleResultStore(input.sessionId, input.sampleResultStore) ); const model = modelFromConfig(input.benchmarkConfig); const runConfig: RunConfig = { From fb7f0337ada5df5c8f9f775d9a46883bcc21f96f Mon Sep 17 00:00:00 2001 From: James Sterling Date: Tue, 25 Aug 2026 17:49:08 +0000 Subject: [PATCH 10/11] refactor(harness): encode run-relative key invariant in SampleResultStore parameter names Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- src/harness/sample-result-store.ts | 14 +++++++++----- 1 file changed, 9 insertions(+), 5 deletions(-) diff --git a/src/harness/sample-result-store.ts b/src/harness/sample-result-store.ts index 677d917..39f6c89 100644 --- a/src/harness/sample-result-store.ts +++ b/src/harness/sample-result-store.ts @@ -92,8 +92,11 @@ export function sampleResultKey(sampleId: string, epoch: number): string { } export interface SampleResultStoreService { - readonly read: (key: string) => Promise; - readonly write: (key: string, data: PersistedSampleOutcome) => Promise; + readonly read: (runRelativeKey: string) => Promise; + readonly write: ( + runRelativeKey: string, + data: PersistedSampleOutcome + ) => Promise; } export class SampleResultStore extends Tag( @@ -106,11 +109,12 @@ export const NOOP_SAMPLE_RESULT_STORE: SampleResultStoreService = { }; export function namespacedSampleResultStore( - prefix: string, + sessionId: string, store: SampleResultStoreService ): SampleResultStoreService { return { - read: (key) => store.read(`${prefix}/${key}`), - write: (key, data) => store.write(`${prefix}/${key}`, data), + read: (runRelativeKey) => store.read(`${sessionId}/${runRelativeKey}`), + write: (runRelativeKey, data) => + store.write(`${sessionId}/${runRelativeKey}`, data), }; } From c8821d090f2ce61fa1ab539979de77fa7da42f9e Mon Sep 17 00:00:00 2001 From: James Sterling Date: Wed, 26 Aug 2026 02:22:24 +0000 Subject: [PATCH 11/11] fix(harness): re-evaluate persisted sample outcomes whose identity mismatches the requested key Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- src/harness/run.ts | 19 +++++++++++---- src/harness/sample-result-resume.test.ts | 30 ++++++++++++++++++++++++ 2 files changed, 44 insertions(+), 5 deletions(-) diff --git a/src/harness/run.ts b/src/harness/run.ts index fa86c1c..ab822d7 100644 --- a/src/harness/run.ts +++ b/src/harness/run.ts @@ -228,12 +228,21 @@ function evaluateOneResumable( if (raw !== null && raw !== undefined) { const persisted = parseSchema(PersistedSampleOutcomeSchema, raw); if (Either.isRight(persisted)) { - return persisted.right; + const { sampleId, epoch: persistedEpoch } = persisted.right.sampleScore; + if (sampleId === sample.id && persistedEpoch === epoch) { + return persisted.right; + } + wLog("Persisted sample outcome identity mismatch; re-evaluating", { + sample_result_key: key, + persisted_sample_id: sampleId, + persisted_epoch: persistedEpoch, + }); + } else { + wLog("Persisted sample outcome failed validation; re-evaluating", { + sample_result_key: key, + error: firstZodIssueMessage(persisted.left), + }); } - wLog("Persisted sample outcome failed validation; re-evaluating", { - sample_result_key: key, - error: firstZodIssueMessage(persisted.left), - }); } const outcome = yield* evaluateOne(opts); if (outcome.isDegraded === true) { diff --git a/src/harness/sample-result-resume.test.ts b/src/harness/sample-result-resume.test.ts index 0cb3cde..ea8edaf 100644 --- a/src/harness/sample-result-resume.test.ts +++ b/src/harness/sample-result-resume.test.ts @@ -323,6 +323,36 @@ describe("sample-result resume", () => { } }); + it("T8: persisted record with mismatched identity is logged and re-evaluated", async () => { + const backing = makeInMemoryStore(run1Store.map); + const misMappingStore: SampleResultStoreService = { + read: (key) => backing.read(`s1/${key.split("/")[1]}`), + write: backing.write, + }; + const warn = spyOn(console, "warn").mockImplementation(() => {}); + try { + const { result, counters } = await runWithStore(misMappingStore); + expect(counters.solver).toBe(4); + expect(counters.scorer).toBe(4); + expect(summarize(result)).toEqual(run1Summary); + expect(result.metrics.totalQuestions).toBe(run1Metrics.total); + expect(result.metrics.accuracy).toBeCloseTo(run1Metrics.accuracy, 5); + const warned = warn.mock.calls.map((call) => + call.map((arg) => JSON.stringify(arg)).join(" ") + ); + expect( + warned.some( + (line) => + line.includes("Persisted sample outcome identity mismatch") && + line.includes("s2/0") && + line.includes("s1") + ) + ).toBe(true); + } finally { + warn.mockRestore(); + } + }); + it("T6: scorer trajectory survives persistence and resume", async () => { const trajectoryScorer: ScorerService = (state, target) => mcqScorer(state, target).pipe(