Skip to content
Draft
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
1 change: 1 addition & 0 deletions package.json
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down
38 changes: 28 additions & 10 deletions src/harness/run.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand All @@ -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";
Expand Down Expand Up @@ -94,14 +96,22 @@ function makeLayers(
layer: Layer<Model>;
},
solverService: ReturnType<typeof chain>
): Layer<Dataset | Solver | Scorer | ProgressReporter | CheckpointStore> {
): 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", () => {
Expand Down Expand Up @@ -130,7 +140,8 @@ describe("runBenchmark", () => {
layerSucceed(Solver, Solver.of(solver)),
layerSucceed(Scorer, Scorer.of(mcqScorer)),
noopProgressLayer,
noopCheckpointLayer
noopCheckpointLayer,
noopSampleResultLayer
);
await runHarnessPromise(
runBenchmark({
Expand Down Expand Up @@ -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))
Expand Down Expand Up @@ -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))
Expand All @@ -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))
Expand All @@ -268,7 +282,8 @@ describe("runBenchmark", () => {
layerSucceed(Scorer, Scorer.of(mcqScorer)),
model.layer,
noopProgressLayer,
noopCheckpointLayer
noopCheckpointLayer,
noopSampleResultLayer
);
const result = await runPromise(
runBenchmark({
Expand Down Expand Up @@ -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))
Expand Down Expand Up @@ -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))
Expand Down Expand Up @@ -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))
Expand Down
88 changes: 84 additions & 4 deletions src/harness/run.ts
Original file line number Diff line number Diff line change
@@ -1,11 +1,13 @@
import type { Effect } from "effect/Effect";
import {
catchAll,
catchTags,
fail as effectFail,
flatMap as effectFlatMap,
gen as effectGen,
map as effectMap,
annotateLogs,
tryPromise,
withLogSpan,
succeed as effectSucceed,
} from "effect/Effect";
Expand All @@ -18,6 +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";
Expand All @@ -42,6 +48,11 @@ import { Dataset } from "./dataset";
import type { AggregateMetrics, SampleScore } from "./metric";
import { aggregateScores } from "./metric";
import { CheckpointStore, ProgressReporter } from "./progress";
import {
PersistedSampleOutcomeSchema,
SampleResultStore,
sampleResultKey,
} from "./sample-result-store";
import { Scorer } from "./scorer";
import { Solver } from "./solver";

Expand Down Expand Up @@ -77,6 +88,7 @@ type EvalOutcome = {
sampleScore: SampleScore;
usage?: ModelUsage;
generationTimeMs?: number;
isDegraded?: boolean;
};

function sampleEpochStream(
Expand Down Expand Up @@ -109,12 +121,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* () {
Expand Down Expand Up @@ -190,6 +202,68 @@ 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 raw = 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,
});
return effectSucceed(null);
})
);
if (raw !== null && raw !== undefined) {
const persisted = parseSchema(PersistedSampleOutcomeSchema, raw);
if (Either.isRight(persisted)) {
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),
});
}
}
const outcome = yield* evaluateOne(opts);
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,
});
return effectSucceed(undefined);
})
);
return outcome;
});
}

function evaluateOne(
opts: EvaluateOneOpts
): Effect<
Expand Down Expand Up @@ -313,6 +387,7 @@ function errorOutcome(opts: ErrorOutcomeOpts): EvalOutcome {
input: sample.input,
target: sample.target.text,
},
isDegraded: true,
};
}

Expand All @@ -330,7 +405,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) => {
Expand All @@ -348,7 +428,7 @@ export function runBenchmark(
(se) =>
evalWithProgress(
se,
evaluateOne({
evaluateOneResumable({
sampleEpoch: se,
degradeSolverErrors: config.degradeSolverErrors ?? false,
}).pipe(
Expand Down
Loading
Loading