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
10 changes: 8 additions & 2 deletions src/benchmarks/tau-bench-airline/environment.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,10 @@ import type { Effect, Semaphore } from "effect/Effect";
import { gen } from "effect/Effect";

import type { CachedFileError } from "../../datasets/cached-file";
import { fetchCachedTextFile } from "../../datasets/cached-file";
import {
fetchCachedTextFile,
jsonTextValidator,
} from "../../datasets/cached-file";
import { Either } from "../../internal/either";
import { isRecord } from "../../internal/guards";
import type { AirlineData } from "./types";
Expand All @@ -32,7 +35,10 @@ export function ensureAirlineData(
if (airlineDbCache) {
return;
}
airlineDbCache = yield* fetchCachedTextFile({ url: AIRLINE_DB_URL });
airlineDbCache = yield* fetchCachedTextFile({
url: AIRLINE_DB_URL,
validate: jsonTextValidator("object"),
});
})
);
}
Expand Down
18 changes: 18 additions & 0 deletions src/benchmarks/tau3-bench-banking/dataset.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -103,5 +103,23 @@ describe("tau3-bench-banking dataset", () => {
expect(size).toBe(fixtureTasks.length);
expect(fetchCalls).toBe(2);
});
it("does not retry when maxRetries is 0", async () => {
seedBankingTasksRawCache("");
let fetchCalls = 0;
global.fetch = async () => {
fetchCalls++;
return new Response("upstream hiccup", { status: 503 });
};
const layer = makeBankingDatasetLayer({ baseDelayMs: 0, maxRetries: 0 });
await expect(
runPromise(
Dataset.pipe(
flatMap((dataset) => dataset.size),
provide(layer)
)
)
).rejects.toThrow();
expect(fetchCalls).toBe(1);
});
});
});
5 changes: 1 addition & 4 deletions src/benchmarks/tau3-bench-banking/dataset.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,6 @@ import {
map,
mapError,
provideService,
retry,
} from "effect/Effect";
import type { Layer as LayerType } from "effect/Layer";
import { effect as layerEffect, provide as layerProvide } from "effect/Layer";
Expand All @@ -17,7 +16,6 @@ import {
fromIterable,
} from "effect/Stream";

import { hfFetchRetrySchedule } from "../../datasets/huggingface";
import type { Sample } from "../../harness/core";
import { DatasetError } from "../../harness/core";
import type { DatasetStreamOptions } from "../../harness/dataset";
Expand All @@ -41,9 +39,8 @@ export function makeBankingDatasetLayer(
const makeService = gen(function* () {
const client = yield* HttpClient.HttpClient;
const fetchLock = yield* makeSemaphore(1);
const loadTasks = ensureBankingTasks(fetchLock).pipe(
const loadTasks = ensureBankingTasks(fetchLock, retryConfig).pipe(
provideService(HttpClient.HttpClient, client),
retry(hfFetchRetrySchedule(retryConfig)),
mapError(
(cause) =>
new DatasetError({
Expand Down
36 changes: 28 additions & 8 deletions src/benchmarks/tau3-bench-banking/environment.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5,10 +5,18 @@ import type { Effect, Semaphore } from "effect/Effect";
import { gen } from "effect/Effect";

import type { CachedFileError } from "../../datasets/cached-file";
import { fetchCachedTextFile } from "../../datasets/cached-file";
import {
fetchCachedTextFile,
jsonTextValidator,
} from "../../datasets/cached-file";
import { Either } from "../../internal/either";
import { isDefinedAndNotNull, isRecord } from "../../internal/guards";
import {
definedValues,
isDefinedAndNotNull,
isRecord,
} from "../../internal/guards";
import { parseSchema } from "../../internal/zod";
import type { RetryConfig } from "../../runtime/retry";
import type { BankingData, BankingTable, Tau3Task } from "./types";
import { BANKING_TABLES, isBankingTableName, Tau3TaskSchema } from "./types";

Expand All @@ -22,17 +30,24 @@ let bankingDbCache: string | undefined;
let bankingTasksCache: string | undefined;

function fetchGithubFile(
filename: string
filename: string,
expected: "object" | "array",
retryConfig?: RetryConfig
): Effect<
string,
CachedFileError | HttpClientError.HttpClientError,
HttpClient.HttpClient
> {
return fetchCachedTextFile({ url: `${BANKING_SOURCE_BASE_URL}/${filename}` });
return fetchCachedTextFile({
url: `${BANKING_SOURCE_BASE_URL}/${filename}`,
validate: jsonTextValidator(expected),
...definedValues({ retry: retryConfig }),
});
}

export function ensureBankingData(
fetchLock: Semaphore
fetchLock: Semaphore,
retryConfig?: RetryConfig
): Effect<
void,
CachedFileError | HttpClientError.HttpClientError,
Expand All @@ -43,13 +58,14 @@ export function ensureBankingData(
if (bankingDbCache) {
return;
}
bankingDbCache = yield* fetchGithubFile("db.json");
bankingDbCache = yield* fetchGithubFile("db.json", "object", retryConfig);
})
);
}

export function ensureBankingTasks(
fetchLock: Semaphore
fetchLock: Semaphore,
retryConfig?: RetryConfig
): Effect<
void,
CachedFileError | HttpClientError.HttpClientError,
Expand All @@ -60,7 +76,11 @@ export function ensureBankingTasks(
if (bankingTasksCache) {
return;
}
bankingTasksCache = yield* fetchGithubFile("tasks.json");
bankingTasksCache = yield* fetchGithubFile(
"tasks.json",
"array",
retryConfig
);
})
);
}
Expand Down
77 changes: 77 additions & 0 deletions src/datasets/cached-file.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ import type { CacheStore } from "./cache-store";
import {
fetchCachedTextFile,
isHuggingFaceUrl,
jsonTextValidator,
parseRetryAfterMs,
} from "./cached-file";

Expand Down Expand Up @@ -245,6 +246,72 @@ describe("fetchCachedTextFile", () => {
expect(entries.size).toBe(0);
});

it("fails without caching when the validator rejects a 200 body", async () => {
stubFetch([
new Response("<html>rate limited</html>", { status: 200 }),
new Response('{"users":{}}', { status: 200 }),
]);
const { store, entries } = makeMemoryStore();

const result = await runPromise(
fetchCachedTextFile({
...REQUEST,
cacheStore: store,
validate: jsonTextValidator("object"),
retry: { maxRetries: 3, baseDelayMs: 1 },
}).pipe(either, provide(FetchHttpClient.layer))
);

assert(Either.isLeft(result));
assert(result.left._tag === "CachedFileError");
expect(result.left.status).toBeUndefined();
expect(requestCount).toBe(1);
expect(entries.size).toBe(0);
});

it("skips the store entirely when caching is disabled", async () => {
stubFetch([new Response("body", { status: 200 })]);
let reads = 0;
let writes = 0;
const { store } = makeMemoryStore({
enabled: false,
async readJson() {
reads += 1;
return undefined;
},
async writeJson() {
writes += 1;
},
});

await expect(run({ ...REQUEST, cacheStore: store })).resolves.toBe("body");
expect(reads).toBe(0);
expect(writes).toBe(0);
expect(requestCount).toBe(1);
});

it("drains the body of a failed response", async () => {
const failed = new Response("slow down", { status: 429 });
const recovered = new Response("recovered", { status: 200 });
const served: Response[] = [];
global.fetch = () => {
const response = served.length === 0 ? failed : recovered;
served.push(response);
requestCount += 1;
return Promise.resolve(response);
};
const { store } = makeMemoryStore();

await expect(
run({
...REQUEST,
cacheStore: store,
retry: { maxRetries: 1, baseDelayMs: 1 },
})
).resolves.toBe("recovered");
expect(failed.bodyUsed).toBe(true);
});

it("returns the body even when the cache write fails", async () => {
stubFetch([new Response("body", { status: 200 })]);
const { store } = makeMemoryStore({
Expand All @@ -267,6 +334,16 @@ describe("isHuggingFaceUrl", () => {
});
});

describe("jsonTextValidator", () => {
it("accepts only the expected JSON shape", () => {
expect(jsonTextValidator("object")("{}")).toBeUndefined();
expect(jsonTextValidator("array")("[]")).toBeUndefined();
expect(jsonTextValidator("object")("[]")).toBeDefined();
expect(jsonTextValidator("array")("{}")).toBeDefined();
expect(jsonTextValidator("object")("<html>")).toBeDefined();
});
});

describe("parseRetryAfterMs", () => {
it("parses delay seconds and http dates", () => {
const now = Date.UTC(2026, 0, 1, 0, 0, 0);
Expand Down
48 changes: 40 additions & 8 deletions src/datasets/cached-file.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,7 @@ import type { Effect } from "effect/Effect";
import { fail, gen, ignore, promise, retry, tryPromise } from "effect/Effect";

import { Either } from "../internal/either";
import { definedValues } from "../internal/guards";
import { definedValues, isRecord } from "../internal/guards";
import { parseSchema, z } from "../internal/zod";
import type { RetryConfig } from "../runtime/retry";
import type { CacheStore } from "./cache-store";
Expand All @@ -19,11 +19,14 @@ export class CachedFileError extends TaggedError("CachedFileError")<{
readonly retryAfterMs?: number;
}> {}

export type CachedTextValidator = (text: string) => string | undefined;

export interface CachedTextFileRequest {
readonly url: string;
readonly retry?: RetryConfig;
readonly cacheStore?: CacheStore;
readonly hfToken?: string;
readonly validate?: CachedTextValidator;
}

type CachedFileFailure = CachedFileError | HttpClientError.HttpClientError;
Expand Down Expand Up @@ -66,6 +69,22 @@ export function isRetryableCachedFileFailure(
return error.status === 429 || (error.status ?? 0) >= 500;
}

export function jsonTextValidator(
expected: "object" | "array"
): CachedTextValidator {
return (text) => {
const parsed = Either.try((): unknown => JSON.parse(text));
if (Either.isLeft(parsed)) {
return "body is not valid JSON";
}
const value = parsed.right;
if (expected === "array") {
return Array.isArray(value) ? undefined : "body is not a JSON array";
}
return isRecord(value) ? undefined : "body is not a JSON object";
};
}

function cachedFileRetryAfterMs(error: CachedFileFailure): number | undefined {
return error._tag === "CachedFileError" ? error.retryAfterMs : undefined;
}
Expand All @@ -92,6 +111,7 @@ function download(
headers !== undefined ? { headers } : undefined
);
if (response.status < 200 || response.status >= 300) {
yield* ignore(response.text);
return yield* fail(
new CachedFileError(
definedValues({
Expand All @@ -113,12 +133,14 @@ export function fetchCachedTextFile(
const client = yield* HttpClient.HttpClient;
const store = request.cacheStore ?? resolveCacheStore();
const key = `files/${encodeCacheKeySegment(request.url)}.json`;
const cached = parseSchema(
CachedTextSchema,
yield* promise(() => store.readJson(key))
);
if (Either.isRight(cached)) {
return cached.right.text;
if (store.enabled) {
const cached = parseSchema(
CachedTextSchema,
yield* promise(() => store.readJson(key))
);
if (Either.isRight(cached)) {
return cached.right.text;
}
}
const hfToken = request.hfToken ?? (yield* resolveHfToken());
const text = yield* download(request.url, hfToken, client).pipe(
Expand All @@ -130,7 +152,17 @@ export function fetchCachedTextFile(
)
)
);
yield* tryPromise(() => store.writeJson(key, { text })).pipe(ignore);
const invalid = request.validate?.(text);
if (invalid !== undefined) {
return yield* fail(
new CachedFileError({
message: `Invalid response for ${request.url}: ${invalid}`,
})
);
}
if (store.enabled) {
yield* tryPromise(() => store.writeJson(key, { text })).pipe(ignore);
}
return text;
});
}
Loading