Skip to content
Open
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
5 changes: 5 additions & 0 deletions .changeset/tired-sites-shave.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
"@toapi/cache": minor
---

Add S3LfsCache to efficiently cache large files
6 changes: 6 additions & 0 deletions compose.yml
Original file line number Diff line number Diff line change
Expand Up @@ -9,3 +9,9 @@ services:
POSTGRES_PASSWORD: postgres
ports:
- "5432:5432"
s3:
image: ghcr.io/shyim/local-s3:latest
environment:
S3_ACCOUNT_ADMIN: toapi:toapi-secret
ports:
- "9000:9000"
11 changes: 8 additions & 3 deletions packages/toapi-cache/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -41,18 +41,20 @@
"test": "vitest run"
},
"devDependencies": {
"@aws-sdk/client-s3": "^3.1141.0",
"@redis/client": "^6.0.0",
"@types/node": "^25.0.3",
"@types/pg": "^8.0.0",
"vitest": "^5.0.0",
"pg": "^8.0.0",
"typescript": "^7.0.0",
"@redis/client": "^6.0.0",
"pg": "^8.0.0"
"vitest": "^5.0.0"
},
"repository": {
"type": "git",
"url": "git+https://github.com/toapi-js/toapi.git"
},
"peerDependencies": {
"@aws-sdk/client-s3": "^3.1141.0",
"@redis/client": "^5.10.0 || ^6.0.0",
"pg": "^8.0.0"
},
Expand All @@ -62,6 +64,9 @@
},
"pg": {
"optional": true
},
"@aws-sdk/client-s3": {
"optional": true
}
}
}
43 changes: 43 additions & 0 deletions packages/toapi-cache/src/no-lfs-cache.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,43 @@
import {
GetObjectCommand,
PutObjectCommand,
type S3Client,
} from "@aws-sdk/client-s3";
import type { Cache, CacheEntry, Json, Subscription } from "./index.js";

interface Options {
cutoffBytes?: number;
}

export class NoLfsCache implements Cache {
cutoffBytes: number;

constructor(
private base: Cache,
options?: Options,
) {
this.cutoffBytes = options?.cutoffBytes ?? 1_000_000;
}

get(key: string) {
return this.base.get(key);
}

set(input: CacheEntry & { key: string; ttl: number; tags: string[] }) {
if (input.attachment && input.attachment.byteLength > this.cutoffBytes) {
return Promise.resolve();
}

return this.base.set(input);
}

delete(tags: string[], meta?: { clientId?: string }) {
return this.invalidate(tags, meta);
}
invalidate(tags: string[], meta?: { clientId?: string }) {
return this.base.invalidate(tags);
}
subscribe(callback: Subscription) {
return this.base.subscribe(callback);
}
}
267 changes: 267 additions & 0 deletions packages/toapi-cache/src/s3-lfs-cache.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,267 @@
import { randomUUID } from "node:crypto";
import {
CreateBucketCommand,
GetObjectCommand,
HeadObjectCommand,
S3Client,
} from "@aws-sdk/client-s3";
import { beforeAll, beforeEach, describe, expect, test } from "vitest";
import { InMemoryCache } from "./in-memory-cache.js";
import { S3LfsCache } from "./s3-lfs-cache.js";

const S3_URL = process.env.S3_URL ?? "http://localhost:9000";
const S3_ACCESS_KEY_ID = process.env.S3_ACCESS_KEY_ID ?? "toapi";
const S3_SECRET_ACCESS_KEY = process.env.S3_SECRET_ACCESS_KEY ?? "toapi-secret";
const BUCKET = "toapi-cache-test";
const CUTOFF = 16;

// Skip the suite when no S3 server is reachable (e.g. local runs without
// `docker compose up`). CI provides an S3 service so it always runs there.
async function isS3Available(url: string): Promise<boolean> {
try {
await fetch(url, { signal: AbortSignal.timeout(1000) });
return true;
} catch {
return false;
}
}

const s3Available = await isS3Available(S3_URL);

function bytes(length: number, seed = 0) {
return Uint8Array.from({ length }, (_, i) => (i + seed) % 256);
}

describe.skipIf(!s3Available)("S3LfsCache", () => {
const client = new S3Client({
endpoint: S3_URL,
region: "us-east-1",
forcePathStyle: true,
credentials: {
accessKeyId: S3_ACCESS_KEY_ID,
secretAccessKey: S3_SECRET_ACCESS_KEY,
},
});

let base: InMemoryCache;
let sut: S3LfsCache;
// The bucket outlives a test run, so every test gets its own key namespace.
let prefix: string;

beforeAll(async () => {
try {
await client.send(new CreateBucketCommand({ Bucket: BUCKET }));
} catch (error) {
if (
!(error instanceof Error) ||
(error.name !== "BucketAlreadyOwnedByYou" &&
error.name !== "BucketAlreadyExists")
) {
throw error;
}
}
});

beforeEach(() => {
base = new InMemoryCache();
sut = new S3LfsCache(base, client, { bucket: BUCKET, cutoffBytes: CUTOFF });
prefix = `${randomUUID()}/`;
});

async function getObject(key: string) {
const response = await client.send(
new GetObjectCommand({ Bucket: BUCKET, Key: key }),
);
return response.Body?.transformToByteArray();
}

async function objectExists(key: string) {
try {
await client.send(new HeadObjectCommand({ Bucket: BUCKET, Key: key }));
return true;
} catch (error) {
if (error instanceof Error && error.name === "NotFound") return false;
throw error;
}
}

test("returns null for unknown keys", async () => {
expect(await sut.get(`${prefix}missing`)).toEqual(null);
});

test("basic store and retrieve", async () => {
const key = `${prefix}test`;

await sut.set({
key,
data: { foo: 1, bar: "baz" },
ttl: 1000,
tags: [],
});

expect(await sut.get(key)).toEqual({
data: { foo: 1, bar: "baz" },
attachment: null,
});
expect(await objectExists(key)).toBe(false);
});

test("keeps small attachments in the base cache", async () => {
const key = `${prefix}small`;
const attachment = bytes(CUTOFF - 1);

await sut.set({ key, attachment, ttl: 1000, tags: [] });

expect(await sut.get(key)).toEqual({ data: null, attachment });
expect((await base.get(key))?.attachment).toEqual(attachment);
expect(await objectExists(key)).toBe(false);
});

test("keeps attachments of exactly cutoffBytes in the base cache", async () => {
const key = `${prefix}boundary`;
const attachment = bytes(CUTOFF);

await sut.set({ key, attachment, ttl: 1000, tags: [] });

expect(await sut.get(key)).toEqual({ data: null, attachment });
expect(await objectExists(key)).toBe(false);
});

test("offloads large attachments to s3", async () => {
const key = `${prefix}large`;
const attachment = bytes(CUTOFF * 64);

await sut.set({ key, attachment, ttl: 1000, tags: [] });

expect(await sut.get(key)).toEqual({ data: null, attachment });
expect(await getObject(key)).toEqual(attachment);
expect((await base.get(key))?.attachment).toBeNull();
});

test("store both data and large attachment", async () => {
const key = `${prefix}both`;
const attachment = bytes(CUTOFF * 64);

await sut.set({
key,
data: { message: "hello" },
attachment,
ttl: 1000,
tags: ["tag1"],
});

expect(await sut.get(key)).toEqual({
data: { message: "hello" },
attachment,
});
});

test("does not mutate the input entry", async () => {
const attachment = bytes(CUTOFF * 64);
const input = {
key: `${prefix}input`,
data: { message: "hello" },
attachment,
ttl: 1000,
tags: [],
};

await sut.set(input);

expect(input.data).toEqual({ message: "hello" });
expect(input.attachment).toBe(attachment);
});

test("overwriting a key returns the latest attachment", async () => {
const key = `${prefix}overwrite`;

await sut.set({ key, attachment: bytes(CUTOFF * 64), ttl: 1000, tags: [] });
await sut.set({
key,
attachment: bytes(CUTOFF * 32, 7),
ttl: 1000,
tags: [],
});

expect((await sut.get(key))?.attachment).toEqual(bytes(CUTOFF * 32, 7));
});

test("overwriting a large attachment with a small one", async () => {
const key = `${prefix}shrink`;

await sut.set({ key, attachment: bytes(CUTOFF * 64), ttl: 1000, tags: [] });
await sut.set({ key, attachment: bytes(4), ttl: 1000, tags: [] });

expect(await sut.get(key)).toEqual({ data: null, attachment: bytes(4) });
});

test("expire by ttl", async () => {
const key = `${prefix}ttl`;

await sut.set({
key,
data: { foo: 1 },
attachment: bytes(CUTOFF * 64),
ttl: 1,
tags: [],
});

// margin over the 1s TTL to avoid a boundary race on loaded CI runners
await new Promise((resolve) => setTimeout(resolve, 1500));

expect(await sut.get(key)).toEqual(null);
});

test("expire by tags", async () => {
const first = `${prefix}first`;
const second = `${prefix}second`;

await sut.set({
key: first,
data: { foo: 1 },
attachment: bytes(CUTOFF * 64),
ttl: 1000,
tags: ["tag1", "tag2"],
});
await sut.set({
key: second,
data: { foo: 2 },
attachment: bytes(CUTOFF * 64, 1),
ttl: 1000,
tags: ["tag2", "tag3"],
});

await sut.invalidate(["tag1"]);

expect(await sut.get(first)).toEqual(null);
expect(await sut.get(second)).toEqual({
data: { foo: 2 },
attachment: bytes(CUTOFF * 64, 1),
});

await sut.delete(["tag2"]);

expect(await sut.get(second)).toEqual(null);
});

test("forwards invalidations to subscribers", async () => {
const calls: string[][] = [];
const unsubscribe = sut.subscribe((tags) => {
calls.push(tags);
});

await sut.invalidate(["tag1"]);
unsubscribe();
await sut.invalidate(["tag2"]);

expect(calls).toEqual([["tag1"]]);
});

test("ignores base entries not written by S3LfsCache", async () => {
const key = `${prefix}foreign`;

await base.set({ key, data: { foo: 1 }, ttl: 1000, tags: [] });

expect(await sut.get(key)).toEqual(null);
});
});
Loading
Loading