diff --git a/examples/workers-cache/README.md b/examples/workers-cache/README.md index 06e5ef3a2..840e1b922 100644 --- a/examples/workers-cache/README.md +++ b/examples/workers-cache/README.md @@ -1,6 +1,6 @@ # Workers Response Store adapter POC -This example uses one `responseStoreAdapter()` from `@vinext/cloudflare` in place of both `cdnAdapter()` and `kvDataAdapter()`. The application Worker keeps Workers Cache disabled. Its `RESPONSE_STORE` service binding calls a separately deployed cache Worker that owns Workers Cache, R2 response bodies, SQLite Durable Object metadata, tag indexes, and SWR regeneration. +This example uses one `responseStoreAdapter()` from `@vinext/cloudflare` in place of both `cdnAdapter()` and `kvDataAdapter()`. The application Worker keeps Workers Cache disabled. Its `RESPONSE_STORE` service binding calls a separately deployed cache Worker that owns Workers Cache, R2 response bodies, SQLite Durable Object metadata and tag invalidation timestamps, and SWR regeneration. The cache Worker is shared infrastructure, not a second deployment of the application. Each application version has one ordinary build and deploy. Cached entries retain a loopback to that application version's vinext response-stage entrypoint. Route and fetch-cache entries can replay that stage, while a transformed public `"use cache"` entry records its encrypted arguments and server-reference identity so regeneration invokes only that function. If its arguments cannot be safely recorded, the adapter falls back to replaying the cacheable route. diff --git a/packages/cloudflare/src/cache/response-store-adapter.worker.ts b/packages/cloudflare/src/cache/response-store-adapter.worker.ts index d8eb031d1..d54e3c71d 100644 --- a/packages/cloudflare/src/cache/response-store-adapter.worker.ts +++ b/packages/cloudflare/src/cache/response-store-adapter.worker.ts @@ -400,6 +400,7 @@ export default { return; } await responseStore.put(key, admitted, { + coalesce: true, revalidator: { id: ROUTE_REVALIDATOR_ID, args: [invocation] }, }); }) @@ -426,6 +427,7 @@ export default { const [foreground, cacheBody] = rendered.body ? rendered.body.tee() : [null, null]; const cacheResponse = new Response(cacheBody, rendered); await responseStore.put(key, cacheResponse, { + coalesce: true, revalidator: { id: ROUTE_REVALIDATOR_ID, args: [invocation] }, }); @@ -446,7 +448,10 @@ export default { headers: rscHeaders, status: 200, }), - { revalidator: { id: ROUTE_REVALIDATOR_ID, args: [rscInvocation] } }, + { + coalesce: true, + revalidator: { id: ROUTE_REVALIDATOR_ID, args: [rscInvocation] }, + }, ); } return publicResponse(new Response(foreground, rendered), "MISS"); diff --git a/packages/cloudflare/src/cache/response-store-data.runtime.ts b/packages/cloudflare/src/cache/response-store-data.runtime.ts index 5d70e4604..9b808c7b9 100644 --- a/packages/cloudflare/src/cache/response-store-data.runtime.ts +++ b/packages/cloudflare/src/cache/response-store-data.runtime.ts @@ -442,6 +442,7 @@ export class WorkersResponseStoreCacheHandler implements CacheHandler { await this.store.put(await cacheRequest(key), response, { ...(revalidator ? { revalidator } : {}), + coalesce: true, purgeExisting: true, }); } diff --git a/packages/cloudflare/tests/response-store-adapter.e2e.test.ts b/packages/cloudflare/tests/response-store-adapter.e2e.test.ts index 9245185e2..987185ece 100644 --- a/packages/cloudflare/tests/response-store-adapter.e2e.test.ts +++ b/packages/cloudflare/tests/response-store-adapter.e2e.test.ts @@ -285,7 +285,7 @@ describe("Cloudflare Workers Response Store adapter", () => { assert.doesNotMatch(serialized, /first-secret|second-secret/); }); - test("keeps concurrent cold renders successful when a cache write loses CAS", async () => { + test("keeps concurrent cold renders successful", async () => { const responses = await Promise.all( Array.from({ length: 8 }, () => request("/cached/concurrent")), ); diff --git a/packages/cloudflare/tests/response-store-data.test.ts b/packages/cloudflare/tests/response-store-data.test.ts index b8240eca3..0208f4894 100644 --- a/packages/cloudflare/tests/response-store-data.test.ts +++ b/packages/cloudflare/tests/response-store-data.test.ts @@ -70,6 +70,7 @@ test("only attaches loopback regeneration to replayable requests", async () => { handler.set("get", null, { cacheControl: { revalidate: 1, expire: 2 } }), ); expect(store.options).toMatchObject({ + coalesce: true, purgeExisting: true, revalidator: { id: "vinext:data", args: ["get", "safe-get"] }, }); @@ -78,7 +79,7 @@ test("only attaches loopback regeneration to replayable requests", async () => { await runWithResponseStoreInvocation("unsafe-post", false, () => handler.set("post", null, { cacheControl: { revalidate: 1, expire: 2 } }), ); - expect(store.options).toEqual({ purgeExisting: true }); + expect(store.options).toEqual({ coalesce: true, purgeExisting: true }); expect(store.response?.headers.get("X-Vinext-Response-Store-Replayable")).toBeNull(); expect(store.response?.headers.get("Cache-Control")).toBe("public, max-age=315360000"); }); diff --git a/packages/workers-response-store/README.md b/packages/workers-response-store/README.md index ab1b79e47..d075ce106 100644 --- a/packages/workers-response-store/README.md +++ b/packages/workers-response-store/README.md @@ -74,7 +74,9 @@ The complete pair lives in `example/service-binding`: `cache-worker.ts` owns the ## API and storage model -The library exposes `fetch`, `put`, `refresh({ tags, pathPrefixes })`, and `purge({ tags, pathPrefixes, purgeEverything })`. Cache identity is the request pathname plus query string. Response bodies live only in revision-specific R2 objects; SQLite stores metadata, freshness, revalidator descriptors, tag invalidation timestamps, and the reverse tag-to-entry index needed by refresh and purge. +The library exposes `fetch`, `put`, `refresh({ tags, pathPrefixes })`, and `purge({ tags, pathPrefixes, purgeEverything })`. Cache identity is the request pathname plus query string. Response bodies live only in revision-specific R2 objects; SQLite stores metadata, freshness, revalidator descriptors, and tag invalidation timestamps. Explicit refresh and purge operations scan stored entry metadata instead of maintaining a write-heavy tag index. + +Framework soft-tag checks query their expiration only after a candidate cache hit and memoize each distinct tag set within the current request. Misses and repeated lookups therefore add no metadata DO request, while every check still reads the authoritative invalidation state directly. SQLite assigns monotonically increasing revisions and conditionally publishes metadata, so slow writes cannot replace newer writes or resurrect purged entries. User RPC, R2, and purge I/O happen outside SQLite transactions. Stale R2 responses within their SWR window return immediately while `ctx.waitUntil()` runs one claimed regeneration; hard-expired responses are never returned. diff --git a/packages/workers-response-store/example/service-binding/user-worker.ts b/packages/workers-response-store/example/service-binding/user-worker.ts index 0d49f25f1..f84c1e634 100644 --- a/packages/workers-response-store/example/service-binding/user-worker.ts +++ b/packages/workers-response-store/example/service-binding/user-worker.ts @@ -104,6 +104,11 @@ export default { return json(await responseStore.purge((await request.json()) as ResponseStorePurgeOptions)); } + if (request.method === "POST" && url.pathname === "/admin/tag-expiration") { + const { tags } = (await request.json()) as { tags: string[] }; + return json({ expiration: await responseStore.getTagExpiration(tags) }); + } + return new Response("Not found", { status: 404 }); } catch (error) { return json({ error: error instanceof Error ? error.message : String(error) }, 500); diff --git a/packages/workers-response-store/example/worker.ts b/packages/workers-response-store/example/worker.ts index ce6ea2bdd..ba93929be 100644 --- a/packages/workers-response-store/example/worker.ts +++ b/packages/workers-response-store/example/worker.ts @@ -62,6 +62,13 @@ async function handlePut(request: Request, store: WorkersResponseStore): Promise if (cdnCacheControl) headers.set("CDN-Cache-Control", cdnCacheControl); let body = request.body; + if (request.headers.get("X-Body-Failure") === "1") { + body = new ReadableStream({ + start(controller) { + controller.error(new Error("Fixture body failure")); + }, + }); + } if (body && bodyDelayMs > 0) { const reader = body.getReader(); let delayed = false; @@ -82,6 +89,10 @@ async function handlePut(request: Request, store: WorkersResponseStore): Promise }, }); } + let teeSibling: ReadableStream | undefined; + if (body && request.headers.get("X-Tee-Body") === "1") { + [body, teeSibling] = body.tee(); + } const response = new Response(NULL_BODY_STATUSES.has(status) ? null : body, { status, @@ -99,9 +110,11 @@ async function handlePut(request: Request, store: WorkersResponseStore): Promise }; const result = await store.put(target, response, { + coalesce: request.headers.get("X-Coalesce") === "1", revalidator, purgeExisting: request.headers.get("X-Purge-Existing") === "1", }); + await new Response(teeSibling).arrayBuffer(); return json(result); } @@ -188,6 +201,11 @@ export default { return json(await responseStore.purge(options)); } + if (request.method === "POST" && url.pathname === "/admin/tag-expiration") { + const { tags } = (await request.json()) as { tags: string[] }; + return json({ expiration: await responseStore.getTagExpiration(tags) }); + } + if (request.method === "GET" && url.pathname === "/admin/stats") { return json({ regenerationCount }); } diff --git a/packages/workers-response-store/src/binding.ts b/packages/workers-response-store/src/binding.ts index 7f706a977..22059f498 100644 --- a/packages/workers-response-store/src/binding.ts +++ b/packages/workers-response-store/src/binding.ts @@ -17,6 +17,8 @@ export type SerializableValue = | { [key: string]: SerializableValue }; export type ResponseStorePutOptions = { + /** @internal Collapse overlapping framework writes for the same cache key. */ + coalesce?: boolean; revalidator?: RevalidatorDescriptor; purgeExisting?: boolean; }; @@ -68,10 +70,7 @@ type EntryMetadata = { }; export type CandidateMetadata = EntryMetadata & { - status: number; - createdAt: number; - initialAge: number; - responseMetadataInR2?: true; + fenceTags: string[]; }; export type StoredEntry = EntryMetadata & { @@ -79,11 +78,6 @@ export type StoredEntry = EntryMetadata & { cacheKey: string; activeRevision: number; latestRevision: number; - legacyResponseMetadata?: { - status: number; - createdAt: number; - initialAge: number; - }; }; export type PurgedEntry = { @@ -102,6 +96,9 @@ type CacheKey = { }; type WriteReservation = CacheKey & { + claimId?: string; + fenceTags: string[]; + objectKey: string; revision: number; }; @@ -111,8 +108,8 @@ type StoreResult = { }; type PublicationResult = { + entry: StoredEntry | null; published: boolean; - previousObjectKey?: string; }; export type WorkersResponseStoreEnv = { @@ -207,13 +204,11 @@ const MISS_HEADERS = { }; const BACKGROUND_REVALIDATION_LEASE_MS = 30_000; -const ORPHAN_RETENTION_MS = 60 * 60 * 1000; -const ORPHAN_CLEANUP_LIMIT = 100; -const R2_DELETE_BATCH_SIZE = 1_000; const CACHE_PURGE_BATCH_SIZE = 100; const MAX_CACHE_TAG_HEADER_BYTES = 16 * 1024; const NULL_BODY_STATUSES = new Set([204, 205, 304]); const AGE_BASIS_HEADER = "X-Workers-Response-Store-Age-Basis"; +const pendingPuts = new Map>(); function* batches(values: readonly T[], size: number): Generator { for (let offset = 0; offset < values.length; offset += size) { @@ -259,6 +254,17 @@ function cacheTagHeader(entry: Pick): stri return tags.join(","); } +function cacheTagsFromResponse(response: Response): string[] { + return [ + ...new Set( + (response.headers.get("Cache-Tag") ?? "") + .split(",") + .map((tag) => tag.trim()) + .filter(Boolean), + ), + ]; +} + export class ResponseStoreBinding extends WorkerEntrypoint< WorkersResponseStoreEnv, WorkersResponseStoreProps @@ -339,6 +345,29 @@ export class ResponseStoreBinding extends WorkerEntrypoint< return accepted; } + private objectKeyRoot(): string { + return ["runtime-cache", this.getVersionId()].join("/"); + } + + private objectKeyPrefix(keyHash: string): string { + return `${this.objectKeyRoot()}/${keyHash}`; + } + + private async reserveWrite( + metadata: CacheMetadataStub, + keyHash: string, + cacheKey: string, + cacheTags: string[], + ): Promise { + const reservation = await metadata.reserveWrite( + keyHash, + cacheKey, + this.objectKeyPrefix(keyHash), + Date.now(), + ); + return { cacheKey, fenceTags: cacheTags, keyHash, ...reservation }; + } + private logCleanupFailure(objectKey: string, error: unknown): void { console.error( JSON.stringify({ @@ -349,56 +378,13 @@ export class ResponseStoreBinding extends WorkerEntrypoint< ); } - private async deletePendingObjects( + private async releaseFailedWrite( metadata: CacheMetadataStub, - objectKeys: string[], + write: Pick, ): Promise { - if (!objectKeys.length) { - return; - } - - for (const batch of batches(objectKeys, R2_DELETE_BATCH_SIZE)) { - try { - await this.env.CACHE_BODIES.delete(batch); - await metadata.finishPendingObjects(batch); - } catch (error) { - this.logCleanupFailure(batch.join(","), error); - } - } - } - - private async trackPendingObjects( - metadata: CacheMetadataStub, - objectKeys: string[], - createdAt: number, - ): Promise { - try { - await metadata.trackPendingObjects(objectKeys, createdAt); - } catch (bulkError) { - try { - for (const objectKey of objectKeys) { - await metadata.trackPendingObject(objectKey, createdAt); - } - } catch (fallbackError) { - throw new AggregateError([bulkError, fallbackError], "Failed to track pending R2 objects"); - } - } - } - - private async cleanupExpiredPendingObjects(metadata: CacheMetadataStub): Promise { - try { - const objectKeys = await metadata.listExpiredPendingObjects( - Date.now() - ORPHAN_RETENTION_MS, - ORPHAN_CLEANUP_LIMIT, - ); - if (!objectKeys.length) { - return; - } - - await this.deletePendingObjects(metadata, objectKeys); - } catch (error) { - this.logCleanupFailure("expired-pending-objects", error); - } + await metadata + .releaseWrite(write.keyHash, write.objectKey, write.claimId) + .catch((error) => this.logCleanupFailure(write.objectKey, error)); } private async readStoredResponse(entry: StoredEntry, now = Date.now()): Promise { @@ -407,13 +393,9 @@ export class ResponseStoreBinding extends WorkerEntrypoint< return null; } - const status = - metadataInteger(object.customMetadata?.status) ?? entry.legacyResponseMetadata?.status; - const createdAt = - metadataInteger(object.customMetadata?.createdAt) ?? entry.legacyResponseMetadata?.createdAt; - const initialAge = - metadataInteger(object.customMetadata?.initialAge) ?? - entry.legacyResponseMetadata?.initialAge; + const status = metadataInteger(object.customMetadata?.status); + const createdAt = metadataInteger(object.customMetadata?.createdAt); + const initialAge = metadataInteger(object.customMetadata?.initialAge); if ( status === undefined || status < 200 || @@ -455,47 +437,33 @@ export class ResponseStoreBinding extends WorkerEntrypoint< response: Response, revalidator: ResponseStorePutOptions["revalidator"], reservation?: WriteReservation, + cacheTags = cacheTagsFromResponse(response), ): Promise { - const { cacheKey, keyHash } = reservation ?? (await this.deriveCacheKey(request)); - const revision = reservation?.revision ?? (await metadata.beginWrite(keyHash, cacheKey)); - - const objectKey = ["runtime-cache", this.getVersionId(), keyHash, String(revision)].join("/"); - - const now = Date.now(); - const policy = deriveCachePolicy(response.headers, now); - const responseHeaders = [...response.headers].filter(([name]) => { - const lower = name.toLowerCase(); - return lower !== "age" && lower !== "cf-cache-status" && lower !== "content-length"; - }); - const responseCacheTags = [ - ...new Set( - (response.headers.get("Cache-Tag") ?? "") - .split(",") - .map((tag) => tag.trim()) - .filter(Boolean), - ), - ]; - - const candidate: CandidateMetadata = { - objectKey, - status: response.status, - statusText: response.statusText, - responseHeaders, - createdAt: policy.createdAt, - initialAge: policy.initialAge, - freshUntil: policy.freshUntil, - swrUntil: policy.swrUntil, - revalidator: revalidator ?? null, - cacheTags: responseCacheTags, - // Keep the scalar values in the RPC payload so an older DO can still - // publish during a rolling deployment. The marker tells the current DO - // not to duplicate them in SQLite. - responseMetadataInR2: true, - }; - - await this.trackPendingObjects(metadata, [objectKey], now); + const cacheKey = reservation ?? (await this.deriveCacheKey(request)); + const write = + reservation ?? + (await this.reserveWrite(metadata, cacheKey.keyHash, cacheKey.cacheKey, cacheTags)); + const { keyHash, objectKey, revision } = write; + let publication: PublicationResult; try { + const now = Date.now(); + const policy = deriveCachePolicy(response.headers, now); + const responseHeaders = [...response.headers].filter(([name]) => { + const lower = name.toLowerCase(); + return lower !== "age" && lower !== "cf-cache-status" && lower !== "content-length"; + }); + const candidate: CandidateMetadata = { + fenceTags: [...new Set([...write.fenceTags, ...cacheTags])], + objectKey, + statusText: response.statusText, + responseHeaders, + freshUntil: policy.freshUntil, + swrUntil: policy.swrUntil, + revalidator: revalidator ?? null, + cacheTags, + }; + // RPC-transferred Response streams do not retain the fixed-length marker // required by R2's single-part put API. Materialise only in the cache // Worker; bodies are never stored in the metadata Durable Object. @@ -507,37 +475,18 @@ export class ResponseStoreBinding extends WorkerEntrypoint< initialAge: String(policy.initialAge), }, }); - } catch (error) { - await this.deletePendingObjects(metadata, [objectKey]); - throw error; - } - let publication: PublicationResult; - try { - publication = await metadata.publish(keyHash, revision, candidate); + publication = await metadata.publish(keyHash, revision, candidate, write.claimId); } catch (error) { - await this.deletePendingObjects(metadata, [objectKey]); + await this.releaseFailedWrite(metadata, write); throw error; } if (!publication.published) { - await this.deletePendingObjects(metadata, [objectKey]); - return { published: false, entry: await metadata.getEntry(keyHash) }; + return { published: false, entry: publication.entry }; } - await metadata - .finishPendingObjects([objectKey]) - .catch((error) => this.logCleanupFailure(objectKey, error)); - - const previousObjectKey = publication.previousObjectKey; - if (previousObjectKey && previousObjectKey !== objectKey) { - await this.trackPendingObjects(metadata, [previousObjectKey], Date.now()).catch((error) => - this.logCleanupFailure(previousObjectKey, error), - ); - await this.deletePendingObjects(metadata, [previousObjectKey]); - } - - return { published: true, entry: await metadata.getEntry(keyHash) }; + return { published: true, entry: publication.entry }; } private async regenerateEntry( @@ -550,28 +499,31 @@ export class ResponseStoreBinding extends WorkerEntrypoint< throw new Error("Cache entry has no configured revalidator"); } - const origin = - this.ctx.props?.revalidator ?? - (Reflect.get(this.ctx.exports, "ResponseStoreRevalidator") as - | RevalidationService - | undefined); - if (typeof origin?.regenerate !== "function") { - throw new Error("The ResponseStoreRevalidator entrypoint is unavailable"); - } - const cacheRequest = new Request(`https://runtime-cache.invalid${entry.cacheKey}`); - const writeReservation = reservation ?? { - cacheKey: entry.cacheKey, - keyHash: entry.keyHash, - revision: await metadata.beginWrite(entry.keyHash, entry.cacheKey), - }; + const writeReservation = + reservation ?? + (await this.reserveWrite(metadata, entry.keyHash, entry.cacheKey, entry.cacheTags)); - const response = await origin.regenerate({ - request: cacheRequest, - id: entry.revalidator.id, - args: entry.revalidator.args, - reason, - }); + let response: Response; + try { + const origin = + this.ctx.props?.revalidator ?? + (Reflect.get(this.ctx.exports, "ResponseStoreRevalidator") as + | RevalidationService + | undefined); + if (typeof origin?.regenerate !== "function") { + throw new Error("The ResponseStoreRevalidator entrypoint is unavailable"); + } + response = await origin.regenerate({ + request: cacheRequest, + id: entry.revalidator.id, + args: entry.revalidator.args, + reason, + }); + } catch (error) { + await this.releaseFailedWrite(metadata, writeReservation); + throw error; + } return this.storeResponse( metadata, @@ -594,6 +546,7 @@ export class ResponseStoreBinding extends WorkerEntrypoint< entry.keyHash, entry.activeRevision, entry.cacheKey, + this.objectKeyPrefix(entry.keyHash), Date.now(), BACKGROUND_REVALIDATION_LEASE_MS, ); @@ -604,7 +557,10 @@ export class ResponseStoreBinding extends WorkerEntrypoint< try { await this.regenerateEntry(metadata, entry, "swr", { cacheKey: entry.cacheKey, + claimId: claim.claimId, + fenceTags: entry.cacheTags, keyHash: entry.keyHash, + objectKey: claim.objectKey, revision: claim.revision, }); } catch (error) { @@ -615,8 +571,6 @@ export class ResponseStoreBinding extends WorkerEntrypoint< error: error instanceof Error ? error.message : String(error), }), ); - } finally { - await metadata.finishRevalidation(entry.keyHash, claim.claimId); } } @@ -668,18 +622,74 @@ export class ResponseStoreBinding extends WorkerEntrypoint< options: ResponseStorePutOptions = {}, ): Promise { const metadata = this.getMetadata(); - this.ctx.waitUntil(this.cleanupExpiredPendingObjects(metadata)); - const result = await this.storeResponse(metadata, request, response, options.revalidator); - if (!result.published || !result.entry) { - return { backingStoreUpdated: false, edgePurgeAccepted: false }; + const { cacheKey, keyHash } = await this.deriveCacheKey(request); + const cacheTags = cacheTagsFromResponse(response); + const pendingPutKey = `${this.getVersionId()}:${keyHash}:${Boolean(options.purgeExisting)}`; + let reservation: WriteReservation | undefined; + if (options.coalesce) { + for (;;) { + const pending = pendingPuts.get(pendingPutKey); + if (!pending) break; + + reservation ??= await this.reserveWrite(metadata, keyHash, cacheKey, cacheTags); + let result: StoreResult; + try { + result = await pending; + } catch { + // Preserve this response as the fallback when the leading write fails. + if (pendingPuts.get(pendingPutKey) === pending) { + pendingPuts.delete(pendingPutKey); + } + continue; + } + if (result.published && result.entry) { + const objectKey = reservation.objectKey; + await metadata + .finishPendingObjects([objectKey]) + .catch((error) => this.logCleanupFailure(objectKey, error)); + void response.body?.cancel().catch(() => {}); + return { + backingStoreUpdated: true, + edgePurgeAccepted: options.purgeExisting + ? await this.purgeEdgeCacheByTags([purgeTagForEntry(result.entry)]) + : true, + }; + } + if (pendingPuts.get(pendingPutKey) === pending) { + pendingPuts.delete(pendingPutKey); + } + } } - const edgePurgeAccepted = options.purgeExisting - ? await this.purgeEdgeCacheByTags([purgeTagForEntry(result.entry)]) - : true; - - return { backingStoreUpdated: true, edgePurgeAccepted }; + const write = (async (): Promise => { + reservation ??= await this.reserveWrite(metadata, keyHash, cacheKey, cacheTags); + return this.storeResponse( + metadata, + request, + response, + options.revalidator, + reservation, + cacheTags, + ); + })(); + if (options.coalesce) pendingPuts.set(pendingPutKey, write); + try { + const result = await write; + if (!result.published || !result.entry) { + return { backingStoreUpdated: false, edgePurgeAccepted: false }; + } + return { + backingStoreUpdated: true, + edgePurgeAccepted: options.purgeExisting + ? await this.purgeEdgeCacheByTags([purgeTagForEntry(result.entry)]) + : true, + }; + } finally { + if (options.coalesce && pendingPuts.get(pendingPutKey) === write) { + pendingPuts.delete(pendingPutKey); + } + } } async refresh(options: ResponseStoreRefreshOptions): Promise { @@ -688,14 +698,26 @@ export class ResponseStoreBinding extends WorkerEntrypoint< } const metadata = this.getMetadata(); - const activeEntries = await metadata.getEntriesMatching(options); - if (activeEntries.length === 0) { + const candidates = await metadata.reserveRefresh(options, this.objectKeyRoot(), Date.now()); + if (candidates.length === 0) { return { backingStoreUpdated: false, edgePurgeAccepted: false }; } const settled = await Promise.allSettled( - activeEntries.map(async (entry) => { - const result = await this.regenerateEntry(metadata, entry, "manual"); + candidates.map(async ({ entry, reservation }) => { + const result = await this.regenerateEntry( + metadata, + entry, + "manual", + reservation + ? { + cacheKey: entry.cacheKey, + fenceTags: entry.cacheTags, + keyHash: entry.keyHash, + ...reservation, + } + : undefined, + ); return result.published ? result.entry : null; }), ); @@ -720,7 +742,7 @@ export class ResponseStoreBinding extends WorkerEntrypoint< } return { - backingStoreUpdated: refreshed.length === activeEntries.length, + backingStoreUpdated: refreshed.length === candidates.length, edgePurgeAccepted, }; } @@ -742,12 +764,6 @@ export class ResponseStoreBinding extends WorkerEntrypoint< ); } - if (purged.length > 0) { - const objectKeys = purged.map((entry) => entry.objectKey); - await this.trackPendingObjects(metadata, objectKeys, Date.now()); - await this.deletePendingObjects(metadata, objectKeys); - } - return { backingStoreUpdated: true, edgePurgeAccepted }; } } diff --git a/packages/workers-response-store/src/metadata-do.ts b/packages/workers-response-store/src/metadata-do.ts index eaf083f68..e3ea77641 100644 --- a/packages/workers-response-store/src/metadata-do.ts +++ b/packages/workers-response-store/src/metadata-do.ts @@ -11,44 +11,61 @@ import type { type RevalidationClaim = { claimId: string; + objectKey: string; revision: number; }; type PublicationResult = { + entry: StoredEntry | null; published: boolean; - previousObjectKey?: string; }; -type TagIndexEntry = { - keyHash: string; - tag: string; +type WriteReservation = { + objectKey: string; + revision: number; +}; + +type RefreshCandidate = { + entry: StoredEntry; + reservation?: Pick; }; export type CacheMetadataStub = DurableObjectStub & { - beginWrite(keyHash: string, cacheKey: string): Promise; + reserveWrite( + keyHash: string, + cacheKey: string, + objectKeyPrefix: string, + createdAt: number, + ): Promise; claimRevalidation( keyHash: string, activeRevision: number, cacheKey: string, + objectKeyPrefix: string, now: number, leaseMs: number, ): Promise; - finishRevalidation(keyHash: string, claimId: string): Promise; trackPendingObject(objectKey: string, createdAt: number): Promise; trackPendingObjects(objectKeys: string[], createdAt: number): Promise; + releaseWrite(keyHash: string, objectKey: string, claimId?: string): Promise; finishPendingObjects(objectKeys: string[]): Promise; listExpiredPendingObjects(cutoff: number, limit: number): Promise; + sweepExpiredPendingObjects(cutoff?: number): Promise; publish( keyHash: string, revision: number, metadata: CandidateMetadata, + claimId?: string, ): Promise; getEntry(keyHash: string): Promise; getTagExpiration(tags: string[]): Promise; - getEntriesMatching(options: ResponseStoreRefreshOptions): Promise; + reserveRefresh( + options: ResponseStoreRefreshOptions, + objectKeyRoot: string, + createdAt: number, + ): Promise; purgeMatching(options: ResponseStorePurgeOptions): Promise; inspect(): Promise; - inspectTagIndex(): Promise; }; type EntryRow = Record & { @@ -57,11 +74,14 @@ type EntryRow = Record & { active_revision: number | null; latest_revision: number; object_key: string | null; - status: number | null; + claim_active_revision?: number | null; + claim_id?: string | null; + claim_revision?: number | null; + current_invalidation_sequence?: number; + pending_invalidation_sequence?: number | null; + pending_publishable?: number | null; status_text: string | null; response_headers: string | null; - created_at: number | null; - initial_age: number | null; fresh_until: number | null; swr_until: number | null; revalidator_id: string | null; @@ -70,18 +90,27 @@ type EntryRow = Record & { tombstoned: number; }; -type TagRow = Record & { - key_hash: string; - tag: string; -}; - const MAX_SQL_PARAMETERS = 100; +const R2_DELETE_BATCH_SIZE = 1_000; +const ORPHAN_RETENTION_MS = 60 * 60 * 1000; +const ORPHAN_CLEANUP_LIMIT = 100; +const ORPHAN_CLEANUP_RETRY_MS = 60 * 1000; + +type CacheMetadataEnv = { + CACHE_BODIES: R2Bucket; +}; function normalizeTags(tags: string[]): string[] { const normalized = tags.map((tag) => tag.trim().toLowerCase()).filter(Boolean); return [...new Set(normalized)]; } +function* batches(values: readonly T[], size: number): Generator { + for (let offset = 0; offset < values.length; offset += size) { + yield values.slice(offset, offset + size); + } +} + function storedEntryFromRow(row: EntryRow): StoredEntry | null { if ( row.tombstoned || @@ -94,7 +123,7 @@ function storedEntryFromRow(row: EntryRow): StoredEntry | null { return null; } - const entry: StoredEntry = { + return { keyHash: row.key_hash, cacheKey: row.cache_key, activeRevision: row.active_revision, @@ -113,16 +142,6 @@ function storedEntryFromRow(row: EntryRow): StoredEntry | null { }, cacheTags: JSON.parse(row.cache_tags ?? "[]") as string[], }; - - if (row.status !== null && row.created_at !== null && row.initial_age !== null) { - entry.legacyResponseMetadata = { - status: row.status, - createdAt: row.created_at, - initialAge: row.initial_age, - }; - } - - return entry; } function storedEntriesFromRows(rows: EntryRow[]): StoredEntry[] { @@ -138,8 +157,10 @@ function storedEntriesFromRows(rows: EntryRow[]): StoredEntry[] { return entries; } -export class CacheMetadata extends DurableObject> { - constructor(ctx: DurableObjectState, env: Record) { +export class CacheMetadata extends DurableObject { + private cleanupAlarmKnown = false; + + constructor(ctx: DurableObjectState, env: CacheMetadataEnv) { super(ctx, env); void ctx.blockConcurrencyWhile(async () => { ctx.storage.sql.exec(` @@ -149,11 +170,8 @@ export class CacheMetadata extends DurableObject> { active_revision INTEGER, latest_revision INTEGER NOT NULL, object_key TEXT, - status INTEGER, status_text TEXT, response_headers TEXT, - created_at INTEGER, - initial_age INTEGER, fresh_until INTEGER, swr_until INTEGER, revalidator_id TEXT, @@ -170,126 +188,168 @@ export class CacheMetadata extends DurableObject> { claimed_at INTEGER NOT NULL, expires_at INTEGER NOT NULL ); - CREATE TABLE IF NOT EXISTS entry_tags ( - tag TEXT NOT NULL, - key_hash TEXT NOT NULL, - PRIMARY KEY (tag, key_hash) - ) WITHOUT ROWID; - CREATE INDEX IF NOT EXISTS entry_tags_key_hash ON entry_tags(key_hash); CREATE TABLE IF NOT EXISTS tag_invalidations ( tag TEXT PRIMARY KEY, - invalidated_at INTEGER NOT NULL + invalidated_at INTEGER NOT NULL, + invalidation_sequence INTEGER NOT NULL DEFAULT 0 ) WITHOUT ROWID; CREATE TABLE IF NOT EXISTS metadata_schema_migrations ( version INTEGER PRIMARY KEY ); + CREATE TABLE IF NOT EXISTS metadata_state ( + singleton INTEGER PRIMARY KEY CHECK (singleton = 1), + tag_invalidation_sequence INTEGER NOT NULL + ); + INSERT OR IGNORE INTO metadata_state (singleton, tag_invalidation_sequence) VALUES (1, 0); CREATE TABLE IF NOT EXISTS pending_objects ( object_key TEXT PRIMARY KEY, + invalidation_sequence INTEGER NOT NULL DEFAULT 0, + publishable INTEGER NOT NULL DEFAULT 1, created_at INTEGER NOT NULL ); CREATE INDEX IF NOT EXISTS pending_objects_created_at ON pending_objects(created_at); `); - const hasBackfilledTagIndex = - ctx.storage.sql - .exec<{ version: number }>( - "SELECT version FROM metadata_schema_migrations WHERE version = 1", - ) - .toArray().length > 0; - - if (!hasBackfilledTagIndex) { - ctx.storage.transactionSync(() => { - const rows = ctx.storage.sql - .exec<{ key_hash: string; cache_tags: string | null }>( - `SELECT key_hash, cache_tags FROM entries - WHERE tombstoned = 0 AND active_revision IS NOT NULL`, + ctx.storage.transactionSync(() => { + const migrations = new Set( + ctx.storage.sql + .exec<{ version: number }>( + "SELECT version FROM metadata_schema_migrations WHERE version IN (2, 3)", ) - .toArray(); - - for (const row of rows) { - const tags = JSON.parse(row.cache_tags ?? "[]") as string[]; - - for (const tag of normalizeTags(tags)) { - ctx.storage.sql.exec( - "INSERT OR IGNORE INTO entry_tags (tag, key_hash) VALUES (?, ?)", - tag, - row.key_hash, - ); - } - } + .toArray() + .map(({ version }) => version), + ); + if (migrations.size === 2) return; - ctx.storage.sql.exec("INSERT INTO metadata_schema_migrations (version) VALUES (1)"); - }); - } + const schemas = ctx.storage.sql + .exec<{ name: string; sql: string }>( + `SELECT name, sql FROM sqlite_schema + WHERE type = 'table' AND name IN ('tag_invalidations', 'pending_objects')`, + ) + .toArray(); + if ( + !migrations.has(2) && + !schemas + .find(({ name }) => name === "tag_invalidations") + ?.sql.includes("invalidation_sequence") + ) { + ctx.storage.sql.exec( + "ALTER TABLE tag_invalidations ADD COLUMN invalidation_sequence INTEGER NOT NULL DEFAULT 0", + ); + } + if ( + !migrations.has(2) && + !schemas + .find(({ name }) => name === "pending_objects") + ?.sql.includes("invalidation_sequence") + ) { + ctx.storage.sql.exec( + "ALTER TABLE pending_objects ADD COLUMN invalidation_sequence INTEGER NOT NULL DEFAULT 0", + ); + } + if (!migrations.has(2)) { + ctx.storage.sql.exec("INSERT INTO metadata_schema_migrations (version) VALUES (2)"); + } + if (!migrations.has(3)) { + if ( + !schemas.find(({ name }) => name === "pending_objects")?.sql.includes("publishable") + ) { + ctx.storage.sql.exec( + "ALTER TABLE pending_objects ADD COLUMN publishable INTEGER NOT NULL DEFAULT 1", + ); + } + ctx.storage.sql.exec("INSERT INTO metadata_schema_migrations (version) VALUES (3)"); + } + }); }); } - private findMatchingEntryRows(options: ResponseStorePurgeOptions): EntryRow[] { - if (options.purgeEverything) { - return this.ctx.storage.sql - .exec( - "SELECT * FROM entries WHERE tombstoned = 0 AND active_revision IS NOT NULL", - ) - .toArray(); - } - - const matches = new Map(); - const tags = normalizeTags(options.tags ?? []); + private async ensureCleanupAlarm(createdAt: number): Promise { + if (this.cleanupAlarmKnown) return; - for (let offset = 0; offset < tags.length; offset += MAX_SQL_PARAMETERS) { - const batch = tags.slice(offset, offset + MAX_SQL_PARAMETERS); - const placeholders = batch.map(() => "?").join(", "); - const rows = this.ctx.storage.sql - .exec( - `SELECT DISTINCT entries.* FROM entries - INNER JOIN entry_tags ON entry_tags.key_hash = entries.key_hash - WHERE entries.tombstoned = 0 AND entries.active_revision IS NOT NULL - AND entry_tags.tag IN (${placeholders})`, - ...batch, - ) - .toArray(); - - for (const row of rows) { - matches.set(row.key_hash, row); - } + const current = await this.ctx.storage.getAlarm(); + if (current === null) { + await this.ctx.storage.setAlarm(createdAt + ORPHAN_RETENTION_MS); } + this.cleanupAlarmKnown = true; + } - const prefixes = options.pathPrefixes ?? []; - if (prefixes.length) { - const rows = this.ctx.storage.sql - .exec( - "SELECT * FROM entries WHERE tombstoned = 0 AND active_revision IS NOT NULL", - ) - .toArray(); + private async maintainCleanupAlarm(createdAt: number): Promise { + await this.ensureCleanupAlarm(createdAt).catch((error) => { + console.error( + JSON.stringify({ + message: "Workers Response Store cleanup alarm update failed", + error: error instanceof Error ? error.message : String(error), + }), + ); + }); + } - for (const row of rows) { - if (prefixes.some((prefix) => row.cache_key.startsWith(prefix))) { - matches.set(row.key_hash, row); - } - } + private findMatchingEntryRows( + options: ResponseStorePurgeOptions, + includePending = false, + ): EntryRow[] { + const rows = this.ctx.storage.sql + .exec( + includePending + ? `SELECT * FROM entries + WHERE (tombstoned = 0 AND active_revision IS NOT NULL) + OR active_revision IS NULL OR latest_revision > active_revision` + : "SELECT * FROM entries WHERE tombstoned = 0 AND active_revision IS NOT NULL", + ) + .toArray(); + if (options.purgeEverything) { + return rows; } - return [...matches.values()]; + const tags = new Set(normalizeTags(options.tags ?? [])); + const prefixes = options.pathPrefixes ?? []; + return rows.filter( + (row) => + prefixes.some((prefix) => row.cache_key.startsWith(prefix)) || + normalizeTags(JSON.parse(row.cache_tags ?? "[]") as string[]).some((tag) => tags.has(tag)), + ); } - trackPendingObjects(objectKeys: string[], createdAt: number): void { + async trackPendingObjects(objectKeys: string[], createdAt: number): Promise { if (!objectKeys.length) { return; } this.ctx.storage.transactionSync(() => { - for (const objectKey of objectKeys) { + for (const batch of batches(objectKeys, MAX_SQL_PARAMETERS / 2)) { + const values = batch.flatMap((objectKey) => [objectKey, createdAt]); this.ctx.storage.sql.exec( - "INSERT OR REPLACE INTO pending_objects (object_key, created_at) VALUES (?, ?)", - objectKey, - createdAt, + `INSERT OR REPLACE INTO pending_objects + (object_key, created_at, publishable) VALUES ${batch + .map(() => "(?, ?, 0)") + .join(", ")}`, + ...values, ); } }); + await this.ensureCleanupAlarm(createdAt); } - trackPendingObject(objectKey: string, createdAt: number): void { - this.trackPendingObjects([objectKey], createdAt); + trackPendingObject(objectKey: string, createdAt: number): Promise { + return this.trackPendingObjects([objectKey], createdAt); + } + + releaseWrite(keyHash: string, objectKey: string, claimId?: string): void { + // Keep the object registered for cleanup while expiring its overlap lease. + this.ctx.storage.transactionSync(() => { + this.ctx.storage.sql.exec( + "UPDATE pending_objects SET created_at = 0, publishable = 0 WHERE object_key = ?", + objectKey, + ); + if (claimId) { + this.ctx.storage.sql.exec( + "DELETE FROM revalidation_claims WHERE key_hash = ? AND claim_id = ?", + keyHash, + claimId, + ); + } + }); } finishPendingObjects(objectKeys: string[]): void { @@ -298,12 +358,52 @@ export class CacheMetadata extends DurableObject> { } this.ctx.storage.transactionSync(() => { - for (const objectKey of objectKeys) { - this.ctx.storage.sql.exec("DELETE FROM pending_objects WHERE object_key = ?", objectKey); + for (const batch of batches(objectKeys, MAX_SQL_PARAMETERS)) { + this.ctx.storage.sql.exec( + `DELETE FROM pending_objects WHERE object_key IN (${batch.map(() => "?").join(", ")})`, + ...batch, + ); } }); } + private async deleteTrackedObjects(objectKeys: string[]): Promise { + for (const batch of batches(objectKeys, R2_DELETE_BATCH_SIZE)) { + try { + await this.env.CACHE_BODIES.delete(batch); + this.finishPendingObjects(batch); + } catch (error) { + for (const retryBatch of batches(batch, MAX_SQL_PARAMETERS)) { + this.ctx.storage.sql.exec( + `INSERT OR REPLACE INTO pending_objects + (object_key, created_at, publishable) VALUES ${retryBatch + .map(() => "(?, 0, 0)") + .join(", ")}`, + ...retryBatch, + ); + } + console.error( + JSON.stringify({ + message: "Workers Response Store R2 cleanup failed", + objectKeys: batch, + error: error instanceof Error ? error.message : String(error), + }), + ); + try { + await this.ctx.storage.setAlarm(Date.now() + ORPHAN_CLEANUP_RETRY_MS); + this.cleanupAlarmKnown = true; + } catch (alarmError) { + console.error( + JSON.stringify({ + message: "Workers Response Store cleanup retry scheduling failed", + error: alarmError instanceof Error ? alarmError.message : String(alarmError), + }), + ); + } + } + } + } + listExpiredPendingObjects(cutoff: number, limit: number): string[] { return this.ctx.storage.sql .exec<{ object_key: string }>( @@ -320,18 +420,102 @@ export class CacheMetadata extends DurableObject> { .map((row) => row.object_key); } - claimRevalidation( + async sweepExpiredPendingObjects(cutoff = Date.now() - ORPHAN_RETENTION_MS): Promise { + this.cleanupAlarmKnown = false; + const rows = this.ctx.storage.sql + .exec<{ + active: number; + created_at: number; + object_key: string; + publishable: number; + }>( + `SELECT pending_objects.object_key, pending_objects.created_at, + pending_objects.publishable, + entries.object_key IS NOT NULL AS active + FROM pending_objects + LEFT JOIN entries + ON entries.object_key = pending_objects.object_key AND entries.tombstoned = 0 + ORDER BY pending_objects.created_at + LIMIT ?`, + ORPHAN_CLEANUP_LIMIT + 1, + ) + .toArray(); + const batch = rows.slice(0, ORPHAN_CLEANUP_LIMIT); + const activeObjectKeys = batch + .filter(({ active }) => active) + .map(({ object_key }) => object_key); + const expired = batch.filter(({ active, created_at }) => !active && created_at <= cutoff); + const reservations = expired.filter(({ publishable }) => publishable === 1); + const cleanupObjectKeys = expired + .filter(({ publishable }) => publishable === 0) + .map(({ object_key }) => object_key); + const fencedAt = Date.now(); + const expiredObjectKeys = expired.map(({ object_key }) => object_key); + this.finishPendingObjects(activeObjectKeys); + for (const reservationBatch of batches(reservations, MAX_SQL_PARAMETERS - 1)) { + this.ctx.storage.sql.exec( + `UPDATE pending_objects SET created_at = ?, publishable = 0 + WHERE object_key IN (${reservationBatch.map(() => "?").join(", ")})`, + fencedAt, + ...reservationBatch.map(({ object_key }) => object_key), + ); + } + if (expiredObjectKeys.length) { + await this.env.CACHE_BODIES.delete(expiredObjectKeys); + this.finishPendingObjects(cleanupObjectKeys); + } + + const next = batch.find(({ active, created_at }) => !active && created_at > cutoff); + if (!next && rows.length > ORPHAN_CLEANUP_LIMIT) { + await this.ctx.storage.setAlarm(Date.now()); + this.cleanupAlarmKnown = true; + } else if (next || reservations.length) { + await this.ensureCleanupAlarm( + Math.min(next?.created_at ?? Number.POSITIVE_INFINITY, fencedAt), + ); + } + + return expiredObjectKeys.length; + } + + async alarm(): Promise { + try { + await this.sweepExpiredPendingObjects(); + } catch (error) { + console.error( + JSON.stringify({ + message: "Workers Response Store orphan cleanup failed", + error: error instanceof Error ? error.message : String(error), + }), + ); + await this.ctx.storage.setAlarm(Date.now() + ORPHAN_CLEANUP_RETRY_MS); + this.cleanupAlarmKnown = true; + } + } + + async claimRevalidation( keyHash: string, activeRevision: number, cacheKey: string, + objectKeyPrefix: string, now: number, leaseMs: number, - ): RevalidationClaim | null { - return this.ctx.storage.transactionSync(() => { + ): Promise { + const claim = this.ctx.storage.transactionSync(() => { const entry = this.ctx.storage.sql - .exec<{ active_revision: number | null; latest_revision: number; tombstoned: number }>( - `SELECT active_revision, latest_revision, tombstoned - FROM entries WHERE key_hash = ? AND cache_key = ?`, + .exec<{ + active_revision: number | null; + latest_revision: number; + tombstoned: number; + claim_active_revision: number | null; + claim_expires_at: number | null; + }>( + `SELECT entries.active_revision, entries.latest_revision, entries.tombstoned, + revalidation_claims.active_revision AS claim_active_revision, + revalidation_claims.expires_at AS claim_expires_at + FROM entries + LEFT JOIN revalidation_claims ON revalidation_claims.key_hash = entries.key_hash + WHERE entries.key_hash = ? AND entries.cache_key = ?`, keyHash, cacheKey, ) @@ -340,18 +524,13 @@ export class CacheMetadata extends DurableObject> { return null; } - const existing = this.ctx.storage.sql - .exec<{ active_revision: number; expires_at: number }>( - "SELECT active_revision, expires_at FROM revalidation_claims WHERE key_hash = ?", - keyHash, - ) - .toArray()[0]; - if (existing?.active_revision === activeRevision && existing.expires_at > now) { + if (entry.claim_active_revision === activeRevision && (entry.claim_expires_at ?? 0) > now) { return null; } const revision = entry.latest_revision + 1; const claimId = crypto.randomUUID(); + const objectKey = `${objectKeyPrefix}/${revision}`; this.ctx.storage.sql.exec( "UPDATE entries SET latest_revision = ? WHERE key_hash = ? AND active_revision = ?", @@ -370,75 +549,141 @@ export class CacheMetadata extends DurableObject> { now, now + leaseMs, ); + this.ctx.storage.sql.exec( + `INSERT OR REPLACE INTO pending_objects + (object_key, created_at, invalidation_sequence) + VALUES (?, ?, (SELECT tag_invalidation_sequence FROM metadata_state WHERE singleton = 1))`, + objectKey, + now, + ); - return { claimId, revision }; + return { claimId, objectKey, revision }; }); + if (claim) await this.ensureCleanupAlarm(now); + return claim; } - finishRevalidation(keyHash: string, claimId: string): void { - this.ctx.storage.sql.exec( - "DELETE FROM revalidation_claims WHERE key_hash = ? AND claim_id = ?", - keyHash, - claimId, - ); + private reserveRevision( + keyHash: string, + cacheKey: string, + current: { latest_revision: number } | undefined, + ): number { + const revision = (current?.latest_revision ?? 0) + 1; + + if (!current) { + this.ctx.storage.sql.exec( + "INSERT INTO entries (key_hash, cache_key, latest_revision, tombstoned) VALUES (?, ?, ?, 1)", + keyHash, + cacheKey, + revision, + ); + } else { + this.ctx.storage.sql.exec( + "UPDATE entries SET cache_key = ?, latest_revision = ? WHERE key_hash = ?", + cacheKey, + revision, + keyHash, + ); + } + + return revision; } - beginWrite(keyHash: string, cacheKey: string): number { - return this.ctx.storage.transactionSync(() => { + async reserveWrite( + keyHash: string, + cacheKey: string, + objectKeyPrefix: string, + createdAt: number, + ): Promise { + const reservation = this.ctx.storage.transactionSync(() => { const current = this.ctx.storage.sql - .exec<{ latest_revision: number }>( - "SELECT latest_revision FROM entries WHERE key_hash = ?", + .exec<{ active_revision: number | null; latest_revision: number }>( + "SELECT active_revision, latest_revision FROM entries WHERE key_hash = ?", keyHash, ) .toArray()[0]; - const revision = (current?.latest_revision ?? 0) + 1; - - if (!current) { - this.ctx.storage.sql.exec( - "INSERT INTO entries (key_hash, cache_key, latest_revision, tombstoned) VALUES (?, ?, ?, 1)", - keyHash, - cacheKey, - revision, - ); - } else { - this.ctx.storage.sql.exec( - "UPDATE entries SET cache_key = ?, latest_revision = ? WHERE key_hash = ?", - cacheKey, - revision, - keyHash, - ); - } - return revision; + const revision = this.reserveRevision(keyHash, cacheKey, current); + const objectKey = `${objectKeyPrefix}/${revision}`; + this.ctx.storage.sql.exec( + `INSERT OR REPLACE INTO pending_objects + (object_key, created_at, invalidation_sequence) + VALUES (?, ?, (SELECT tag_invalidation_sequence FROM metadata_state WHERE singleton = 1))`, + objectKey, + createdAt, + ); + return { objectKey, revision }; }); + await this.ensureCleanupAlarm(createdAt); + return reservation; } - publish(keyHash: string, revision: number, metadata: CandidateMetadata): PublicationResult { - return this.ctx.storage.transactionSync(() => { + async publish( + keyHash: string, + revision: number, + metadata: CandidateMetadata, + claimId?: string, + ): Promise { + const { cleanupObjectKey, result } = this.ctx.storage.transactionSync(() => { const current = this.ctx.storage.sql - .exec<{ latest_revision: number; object_key: string | null }>( - "SELECT latest_revision, object_key FROM entries WHERE key_hash = ?", + .exec( + `SELECT entries.*, + revalidation_claims.active_revision AS claim_active_revision, + revalidation_claims.claim_id AS claim_id, + revalidation_claims.revision AS claim_revision, + pending_objects.invalidation_sequence AS pending_invalidation_sequence, + pending_objects.publishable AS pending_publishable, + metadata_state.tag_invalidation_sequence AS current_invalidation_sequence + FROM entries + CROSS JOIN metadata_state + LEFT JOIN pending_objects ON pending_objects.object_key = ? + LEFT JOIN revalidation_claims ON revalidation_claims.key_hash = entries.key_hash + WHERE entries.key_hash = ?`, + metadata.objectKey, keyHash, ) .toArray()[0]; - if (!current || current.latest_revision !== revision) { - return { published: false }; + if ( + !current || + current.pending_invalidation_sequence === null || + current.pending_invalidation_sequence === undefined || + current.pending_publishable !== 1 || + revision > current.latest_revision || + (current.active_revision !== null && revision <= current.active_revision) || + (claimId !== undefined && + (current.claim_id !== claimId || + current.claim_revision !== revision || + current.claim_active_revision !== current.active_revision)) || + (metadata.fenceTags.length > 0 && + current.current_invalidation_sequence! > current.pending_invalidation_sequence && + this.getTagInvalidationMaximum(metadata.fenceTags, "invalidation_sequence") > + current.pending_invalidation_sequence) + ) { + if (claimId) { + this.ctx.storage.sql.exec( + "DELETE FROM revalidation_claims WHERE key_hash = ? AND claim_id = ?", + keyHash, + claimId, + ); + } + return { + cleanupObjectKey: + current?.object_key === metadata.objectKey ? undefined : metadata.objectKey, + result: { entry: current ? storedEntryFromRow(current) : null, published: false }, + }; } - const responseMetadataIsInR2 = metadata.responseMetadataInR2 === true; const update = this.ctx.storage.sql.exec( `UPDATE entries SET - active_revision = ?, object_key = ?, status = ?, status_text = ?, - response_headers = ?, created_at = ?, initial_age = ?, fresh_until = ?, swr_until = ?, + active_revision = ?, object_key = ?, status_text = ?, response_headers = ?, + fresh_until = ?, swr_until = ?, revalidator_id = ?, revalidator_args = ?, cache_tags = ?, tombstoned = 0 - WHERE key_hash = ? AND latest_revision = ?`, + WHERE key_hash = ? AND latest_revision >= ? + AND (active_revision IS NULL OR active_revision < ?)`, revision, metadata.objectKey, - responseMetadataIsInR2 ? null : metadata.status, metadata.statusText, JSON.stringify(metadata.responseHeaders), - responseMetadataIsInR2 ? null : metadata.createdAt, - responseMetadataIsInR2 ? null : metadata.initialAge, metadata.freshUntil, metadata.swrUntil, metadata.revalidator?.id ?? null, @@ -446,26 +691,64 @@ export class CacheMetadata extends DurableObject> { JSON.stringify(metadata.cacheTags), keyHash, revision, + revision, ); const published = update.rowsWritten === 1; if (published) { - this.ctx.storage.sql.exec("DELETE FROM entry_tags WHERE key_hash = ?", keyHash); - - for (const tag of normalizeTags(metadata.cacheTags)) { - this.ctx.storage.sql.exec( - "INSERT INTO entry_tags (tag, key_hash) VALUES (?, ?)", - tag, - keyHash, - ); - } + this.ctx.storage.sql.exec( + "DELETE FROM pending_objects WHERE object_key = ?", + metadata.objectKey, + ); + } + if (published && current.object_key && current.object_key !== metadata.objectKey) { + this.ctx.storage.sql.exec( + `INSERT OR IGNORE INTO pending_objects + (object_key, created_at, publishable) VALUES (?, ?, 0)`, + current.object_key, + Date.now(), + ); + } + if (claimId) { + this.ctx.storage.sql.exec( + "DELETE FROM revalidation_claims WHERE key_hash = ? AND claim_id = ?", + keyHash, + claimId, + ); } + const entry: StoredEntry = { + keyHash, + cacheKey: current.cache_key, + activeRevision: revision, + latestRevision: current.latest_revision, + objectKey: metadata.objectKey, + statusText: metadata.statusText, + responseHeaders: metadata.responseHeaders, + freshUntil: metadata.freshUntil, + swrUntil: metadata.swrUntil, + revalidator: metadata.revalidator, + cacheTags: metadata.cacheTags, + }; + return { - published, - ...(current.object_key ? { previousObjectKey: current.object_key } : {}), + cleanupObjectKey: published + ? current.object_key && current.object_key !== metadata.objectKey + ? current.object_key + : undefined + : metadata.objectKey, + result: { + entry: published ? entry : null, + published, + }, }; }); + + if (cleanupObjectKey) { + await this.maintainCleanupAlarm(Date.now()); + await this.deleteTrackedObjects([cleanupObjectKey]); + } + return result; } getEntry(keyHash: string): StoredEntry | null { @@ -475,16 +758,18 @@ export class CacheMetadata extends DurableObject> { return row ? storedEntryFromRow(row) : null; } - getTagExpiration(tags: string[]): number { + private getTagInvalidationMaximum( + tags: string[], + column: "invalidated_at" | "invalidation_sequence", + ): number { const normalized = normalizeTags(tags); let expiration = 0; - for (let offset = 0; offset < normalized.length; offset += MAX_SQL_PARAMETERS) { - const batch = normalized.slice(offset, offset + MAX_SQL_PARAMETERS); + for (const batch of batches(normalized, MAX_SQL_PARAMETERS)) { const placeholders = batch.map(() => "?").join(", "); const row = this.ctx.storage.sql .exec<{ invalidated_at: number | null }>( - `SELECT MAX(invalidated_at) AS invalidated_at + `SELECT MAX(${column}) AS invalidated_at FROM tag_invalidations WHERE tag IN (${placeholders})`, ...batch, ) @@ -495,58 +780,152 @@ export class CacheMetadata extends DurableObject> { return expiration; } - getEntriesMatching(options: ResponseStoreRefreshOptions): StoredEntry[] { - return storedEntriesFromRows(this.findMatchingEntryRows(options)); + getTagExpiration(tags: string[]): number { + return this.getTagInvalidationMaximum(tags, "invalidated_at"); } - purgeMatching(options: ResponseStorePurgeOptions): PurgedEntry[] { - return this.ctx.storage.transactionSync(() => { + async reserveRefresh( + options: ResponseStoreRefreshOptions, + objectKeyRoot: string, + createdAt: number, + ): Promise { + const candidates = this.ctx.storage.transactionSync(() => { const matches = this.findMatchingEntryRows(options); - const invalidatedAt = Date.now(); + const reservations = matches.flatMap((row) => { + const entry = storedEntryFromRow(row); + return entry?.revalidator + ? [ + { + keyHash: row.key_hash, + objectKey: `${objectKeyRoot}/${row.key_hash}/${row.latest_revision + 1}`, + revision: row.latest_revision + 1, + }, + ] + : []; + }); + + for (const batch of batches(reservations, MAX_SQL_PARAMETERS)) { + const keyHashes = batch.map(({ keyHash }) => keyHash); + this.ctx.storage.sql.exec( + `UPDATE entries SET latest_revision = latest_revision + 1 + WHERE key_hash IN (${keyHashes.map(() => "?").join(", ")})`, + ...keyHashes, + ); + } + for (const batch of batches(reservations, MAX_SQL_PARAMETERS / 2)) { + this.ctx.storage.sql.exec( + `INSERT OR REPLACE INTO pending_objects + (object_key, created_at, invalidation_sequence) VALUES ${batch + .map( + () => + "(?, ?, (SELECT tag_invalidation_sequence FROM metadata_state WHERE singleton = 1))", + ) + .join(", ")}`, + ...batch.flatMap(({ objectKey }) => [objectKey, createdAt]), + ); + } + + const byKey = new Map(reservations.map((reservation) => [reservation.keyHash, reservation])); + return matches.flatMap((row) => { + const entry = storedEntryFromRow(row); + if (!entry) return []; + const reservation = byKey.get(row.key_hash); + return [ + { + entry, + ...(reservation + ? { + reservation: { + objectKey: reservation.objectKey, + revision: reservation.revision, + }, + } + : {}), + }, + ]; + }); + }); + + if (candidates.some(({ reservation }) => reservation)) { + await this.ensureCleanupAlarm(createdAt); + } + return candidates; + } + + async purgeMatching(options: ResponseStorePurgeOptions): Promise { + const invalidatedAt = Date.now(); + const matches = this.ctx.storage.transactionSync(() => { + const matches = this.findMatchingEntryRows(options, true); - for (const tag of normalizeTags(options.tags ?? [])) { + const tags = normalizeTags(options.tags ?? []); + if (tags.length) { + this.ctx.storage.sql.exec( + `UPDATE metadata_state SET tag_invalidation_sequence = tag_invalidation_sequence + 1 + WHERE singleton = 1`, + ); + } + for (const batch of batches(tags, MAX_SQL_PARAMETERS / 2)) { this.ctx.storage.sql.exec( - `INSERT INTO tag_invalidations (tag, invalidated_at) VALUES (?, ?) + `INSERT INTO tag_invalidations + (tag, invalidated_at, invalidation_sequence) VALUES ${batch + .map( + () => + "(?, ?, (SELECT tag_invalidation_sequence FROM metadata_state WHERE singleton = 1))", + ) + .join(", ")} ON CONFLICT(tag) DO UPDATE SET invalidated_at = - MAX(tag_invalidations.invalidated_at, excluded.invalidated_at)`, - tag, - invalidatedAt, + MAX(tag_invalidations.invalidated_at, excluded.invalidated_at), + invalidation_sequence = MAX( + tag_invalidations.invalidation_sequence, + excluded.invalidation_sequence + )`, + ...batch.flatMap((tag) => [tag, invalidatedAt]), ); } - for (const row of matches) { + for (const batch of batches(matches, MAX_SQL_PARAMETERS - 1)) { + const keyHashes = batch.map((row) => row.key_hash); + const placeholders = keyHashes.map(() => "?").join(", "); + this.ctx.storage.sql.exec( + `INSERT OR IGNORE INTO pending_objects (object_key, created_at, publishable) + SELECT object_key, ?, 0 FROM entries + WHERE key_hash IN (${placeholders}) AND object_key IS NOT NULL`, + invalidatedAt, + ...keyHashes, + ); this.ctx.storage.sql.exec( `UPDATE entries SET latest_revision = latest_revision + 1, - active_revision = NULL, + active_revision = latest_revision + 1, object_key = NULL, - status = NULL, status_text = NULL, response_headers = NULL, - created_at = NULL, - initial_age = NULL, fresh_until = NULL, swr_until = NULL, revalidator_id = NULL, revalidator_args = NULL, cache_tags = NULL, tombstoned = 1 - WHERE key_hash = ?`, - row.key_hash, + WHERE key_hash IN (${placeholders})`, + ...keyHashes, ); this.ctx.storage.sql.exec( - "DELETE FROM revalidation_claims WHERE key_hash = ?", - row.key_hash, + `DELETE FROM revalidation_claims WHERE key_hash IN (${placeholders})`, + ...keyHashes, ); - this.ctx.storage.sql.exec("DELETE FROM entry_tags WHERE key_hash = ?", row.key_hash); } - return matches.map((row) => ({ - keyHash: row.key_hash, - cacheKey: row.cache_key, - objectKey: row.object_key!, - })); + return matches.flatMap((row) => + row.object_key === null + ? [] + : [{ keyHash: row.key_hash, cacheKey: row.cache_key, objectKey: row.object_key }], + ); }); + if (matches.length) { + await this.maintainCleanupAlarm(invalidatedAt); + await this.deleteTrackedObjects(matches.map((entry) => entry.objectKey)); + } + return matches; } inspect(): StoredEntry[] { @@ -555,11 +934,4 @@ export class CacheMetadata extends DurableObject> { .toArray(); return storedEntriesFromRows(rows); } - - inspectTagIndex(): TagIndexEntry[] { - return this.ctx.storage.sql - .exec("SELECT key_hash, tag FROM entry_tags ORDER BY tag, key_hash") - .toArray() - .map((row) => ({ keyHash: row.key_hash, tag: row.tag })); - } } diff --git a/packages/workers-response-store/tests/e2e.test.mjs b/packages/workers-response-store/tests/e2e.test.mjs index cdf0c0300..1d87bd065 100644 --- a/packages/workers-response-store/tests/e2e.test.mjs +++ b/packages/workers-response-store/tests/e2e.test.mjs @@ -1,4 +1,7 @@ import assert from "node:assert/strict"; +import { mkdtemp, rm } from "node:fs/promises"; +import { tmpdir } from "node:os"; +import path from "node:path"; import { fileURLToPath } from "node:url"; import { Miniflare } from "miniflare"; @@ -15,6 +18,7 @@ beforeEach(async () => { compatibilityDate: "2026-04-08", compatibilityFlags: ["nodejs_compat"], unsafeEphemeralDurableObjects: true, + unsafeInspectDurableObjects: true, workers: [ { name: "user-worker", @@ -64,7 +68,10 @@ async function put(path, body, options = {}) { } if (options.noRevalidator) headers.set("X-No-Revalidator", "1"); if (options.purgeExisting) headers.set("X-Purge-Existing", "1"); + if (options.coalesce) headers.set("X-Coalesce", "1"); + if (options.bodyFailure) headers.set("X-Body-Failure", "1"); if (options.bodyDelayMs) headers.set("X-Body-Delay-Ms", String(options.bodyDelayMs)); + if (options.teeBody) headers.set("X-Tee-Body", "1"); const response = await worker.fetch(`https://user.test/admin/put${path}`, { method: "PUT", headers, @@ -104,6 +111,17 @@ async function purge(options) { return { response, json: await response.json() }; } +async function tagExpiration(tags) { + const response = await worker.fetch("https://user.test/admin/tag-expiration", { + method: "POST", + headers: { "Content-Type": "application/json" }, + body: JSON.stringify({ tags }), + }); + const body = await response.json(); + assert.equal(response.status, 200, JSON.stringify(body)); + return body.expiration; +} + async function metadataStub() { const namespace = await mf.getDurableObjectNamespace("CACHE_METADATA", "user-worker"); return namespace.getByName(metadataName); @@ -113,8 +131,12 @@ async function metadata() { return (await metadataStub()).inspect(); } -async function tagIndex() { - return (await metadataStub()).inspectTagIndex(); +async function metadataRowCount(table) { + const storage = await mf.unsafeGetDurableObjectStorage("user-worker", "CacheMetadata", { + name: metadataName, + }); + const [row] = await storage.exec(`SELECT COUNT(*) AS count FROM ${table}`); + return row.count; } async function r2Objects() { @@ -138,7 +160,6 @@ test("put and fetch use pathname plus query, excluding host", async () => { assert.equal(entries.length, 1); assert.equal(entries[0].cacheKey, "/identity?a=1"); assert.equal("body" in entries[0], false); - assert.equal("legacyResponseMetadata" in entries[0], false); const objects = await r2Objects(); assert.equal(objects.objects.length, 1); const bucket = await mf.getR2Bucket("CACHE_BODIES", "user-worker"); @@ -155,6 +176,7 @@ test("put and fetch use pathname plus query, excluding host", async () => { new TextEncoder().encode(JSON.stringify(object.customMetadata)).byteLength < 128, "R2 custom metadata should remain tiny relative to the 8 KiB object metadata limit", ); + assert.equal(await metadataRowCount("pending_objects"), 0); }); test("null-body response statuses refill without an R2 body stream", async () => { @@ -167,36 +189,6 @@ test("null-body response statuses refill without an R2 body stream", async () => assert.equal((await r2Objects()).objects.length, 1); }); -test("metadata remains readable when published by the pre-R2-metadata binding", async () => { - const cacheKey = "/legacy-object-metadata"; - const keyHash = [ - ...new Uint8Array(await crypto.subtle.digest("SHA-256", new TextEncoder().encode(cacheKey))), - ] - .map((byte) => byte.toString(16).padStart(2, "0")) - .join(""); - const objectKey = "runtime-cache/legacy-object-metadata"; - const bucket = await mf.getR2Bucket("CACHE_BODIES", "user-worker"); - await bucket.put(objectKey, "legacy-body"); - const stub = await metadataStub(); - const revision = await stub.beginWrite(keyHash, cacheKey); - await stub.publish(keyHash, revision, { - objectKey, - status: 200, - statusText: "", - responseHeaders: [["content-type", "text/plain"]], - createdAt: Date.now(), - initialAge: 0, - freshUntil: Date.now() + 60_000, - swrUntil: Date.now() + 60_000, - revalidator: null, - cacheTags: [], - }); - - const response = await read(cacheKey); - assert.equal(response.status, 200); - assert.equal(await response.text(), "legacy-body"); -}); - test("a cold refill preserves downstream headers, representation age, and remaining freshness", async () => { await put("/freshness", "aged", { cacheControl: "public, max-age=100, stale-while-revalidate=100", @@ -304,6 +296,8 @@ test("stale R2 content returns immediately and deduplicates background regenerat assert.equal(fresh.headers.get("X-Workers-Response-Store-Revision"), "2"); const stats = await (await worker.fetch("https://user.test/admin/stats")).json(); assert.equal(stats.regenerationCount, 1); + assert.equal(await metadataRowCount("revalidation_claims"), 0); + assert.equal(await metadataRowCount("pending_objects"), 0); }); test("a failed background regeneration releases its claim for a later retry", async () => { @@ -422,7 +416,26 @@ test("refresh accepts more tag selectors than one SQLite parameter batch", async assert.equal(await (await read("/refresh-many-tags")).text(), "refreshed"); }); -test("refresh and purge use a reverse tag-to-entry index that follows publication", async () => { +test("refresh reserves more than one SQLite batch in one metadata call", async () => { + await Promise.all( + Array.from({ length: 101 }, (_, index) => + put(`/refresh-batch/${index}`, "seed", { + tags: ["refresh-batch"], + revalidator: { body: "refreshed", cacheControl: "public, max-age=60" }, + }), + ), + ); + + assert.deepEqual((await refreshSelectors({ tags: ["refresh-batch"] })).json, { + backingStoreUpdated: true, + edgePurgeAccepted: false, + }); + assert.equal((await metadata()).filter((entry) => entry.activeRevision === 2).length, 101); + assert.equal(await metadataRowCount("pending_objects"), 0); + assert.equal((await r2Objects()).objects.length, 101); +}); + +test("refresh and purge select entries from their stored tags", async () => { await put("/tag-index", "seed", { tags: ["Original", "Shared"], revalidator: { @@ -432,25 +445,18 @@ test("refresh and purge use a reverse tag-to-entry index that follows publicatio }, }); - assert.deepEqual( - (await tagIndex()).map(({ tag }) => tag), - ["original", "shared"], - ); + assert.deepEqual((await metadata())[0].cacheTags, ["Original", "Shared"]); assert.deepEqual((await refreshSelectors({ tags: ["ORIGINAL"] })).json, { backingStoreUpdated: true, edgePurgeAccepted: false, }); - assert.deepEqual( - (await tagIndex()).map(({ tag }) => tag), - ["replacement"], - ); + assert.deepEqual((await metadata())[0].cacheTags, ["Replacement"]); assert.deepEqual((await refreshSelectors({ tags: ["original"] })).json, { backingStoreUpdated: false, edgePurgeAccepted: false, }); await purge({ tags: ["REPLACEMENT"] }); - assert.deepEqual(await tagIndex(), []); assert.equal((await read("/tag-index")).status, 404); assert.ok((await (await metadataStub()).getTagExpiration(["replacement"])) > 0); }); @@ -464,11 +470,20 @@ test("tag expiration is recorded without creating an R2 marker object", async () assert.equal(await stub.getTagExpiration(["other-tag"]), 0); const batchedTags = Array.from({ length: 101 }, (_, index) => `tag-${index}`); - await purge({ tags: [batchedTags.at(-1)] }); + await purge({ tags: batchedTags }); assert.ok((await stub.getTagExpiration(batchedTags)) >= before); assert.equal((await r2Objects()).objects.length, 0); }); +test("tag expiration lookup reads authoritative invalidation state", async () => { + const before = Date.now(); + assert.equal(await tagExpiration(["unchanged"]), 0); + + await purge({ tags: ["changed"] }); + assert.ok((await tagExpiration(["changed"])) >= before); + assert.equal(await tagExpiration(["unchanged"]), 0); +}); + test("the internal purge tag is first and large tag sets remain selectable", async () => { await put("/cache-tag-order", "tagged", { tags: ["user-tag"] }); const taggedResponse = await read("/cache-tag-order"); @@ -508,6 +523,19 @@ test("purge supports tags, path prefixes, and purgeEverything", async () => { assert.equal((await r2Objects()).objects.length, 0); }); +test("purge batches more entries than the SQL parameter limit", async () => { + await Promise.all( + Array.from({ length: 101 }, (_, index) => put(`/large-purge/${index}`, `${index}`)), + ); + + assert.deepEqual((await purge({ pathPrefixes: ["/large-purge/"] })).json, { + backingStoreUpdated: true, + edgePurgeAccepted: false, + }); + assert.equal((await metadata()).length, 0); + assert.equal((await r2Objects()).objects.length, 0); +}); + test("a newer put wins and the superseded candidate is cleaned up", async () => { const slow = put("/race", "slow", { bodyDelayMs: 300 }); await new Promise((resolve) => setTimeout(resolve, 50)); @@ -520,7 +548,385 @@ test("a newer put wins and the superseded candidate is cleaned up", async () => assert.equal((await r2Objects()).objects.length, 1); }); -test("retention cleanup removes orphaned candidates without deleting active R2 objects", async () => { +test("overlapping framework writes can be coalesced", async () => { + const first = put("/coalesced", "first", { bodyDelayMs: 300, coalesce: true }); + await new Promise((resolve) => setTimeout(resolve, 50)); + const second = await put("/coalesced", "second", { coalesce: true }); + const firstResult = await first; + + assert.deepEqual(firstResult.json, { backingStoreUpdated: true, edgePurgeAccepted: true }); + assert.deepEqual(second.json, { backingStoreUpdated: true, edgePurgeAccepted: true }); + assert.equal(await (await read("/coalesced")).text(), "first"); + assert.equal((await metadata())[0].activeRevision, 1); + assert.equal((await r2Objects()).objects.length, 1); + assert.equal(await metadataRowCount("pending_objects"), 0); +}); + +test("a coalesced tee body does not block on its unread sibling", async () => { + const first = put("/coalesced-tee", "first", { bodyDelayMs: 300, coalesce: true }); + await new Promise((resolve) => setTimeout(resolve, 50)); + const second = put("/coalesced-tee", "second", { coalesce: true, teeBody: true }); + + assert.deepEqual((await second).json, { + backingStoreUpdated: true, + edgePurgeAccepted: true, + }); + await first; +}); + +test("writes with different purge requirements are not coalesced", async () => { + const first = put("/coalesced-purge", "first", { bodyDelayMs: 300, coalesce: true }); + await new Promise((resolve) => setTimeout(resolve, 50)); + const second = await put("/coalesced-purge", "second", { + coalesce: true, + purgeExisting: true, + }); + const firstResult = await first; + + assert.deepEqual(firstResult.json, { + backingStoreUpdated: false, + edgePurgeAccepted: false, + }); + assert.deepEqual(second.json, { + backingStoreUpdated: true, + edgePurgeAccepted: false, + }); + assert.equal(await (await read("/coalesced-purge")).text(), "second"); +}); + +test("a failed coalesced write does not suppress an immediate retry", async () => { + await assert.rejects( + put("/coalesced-retry", "fails", { + bodyFailure: true, + coalesce: true, + }), + /put fixture returned 500/, + ); + + const retry = await put("/coalesced-retry", "succeeds", { coalesce: true }); + assert.deepEqual(retry.json, { backingStoreUpdated: true, edgePurgeAccepted: true }); + assert.equal(await (await read("/coalesced-retry")).text(), "succeeds"); +}); + +test("a failed coalesced write preserves an overlapping successful write", async () => { + const failing = put("/coalesced-fallback", "fails", { + bodyDelayMs: 300, + bodyFailure: true, + coalesce: true, + }); + await new Promise((resolve) => setTimeout(resolve, 50)); + const fallback = put("/coalesced-fallback", "succeeds", { coalesce: true }); + + await assert.rejects(failing, /put fixture returned 500/); + assert.deepEqual((await fallback).json, { + backingStoreUpdated: true, + edgePurgeAccepted: true, + }); + assert.equal(await (await read("/coalesced-fallback")).text(), "succeeds"); +}); + +test("a failed newer write does not discard an overlapping successful write", async () => { + const successful = put("/write-fallback", "succeeds", { bodyDelayMs: 300 }); + await new Promise((resolve) => setTimeout(resolve, 50)); + + await assert.rejects( + put("/write-fallback", "fails", { bodyFailure: true }), + /put fixture returned 500/, + ); + assert.deepEqual((await successful).json, { + backingStoreUpdated: true, + edgePurgeAccepted: true, + }); + assert.equal(await (await read("/write-fallback")).text(), "succeeds"); +}); + +test("purge prevents coalesced writes from resurrecting an entry", async () => { + await put("/purge-coalesced", "seed"); + const first = put("/purge-coalesced", "first", { bodyDelayMs: 300, coalesce: true }); + await new Promise((resolve) => setTimeout(resolve, 50)); + const second = put("/purge-coalesced", "second", { coalesce: true }); + await new Promise((resolve) => setTimeout(resolve, 50)); + + await purge({ purgeEverything: true }); + assert.deepEqual((await first).json, { + backingStoreUpdated: false, + edgePurgeAccepted: false, + }); + assert.deepEqual((await second).json, { + backingStoreUpdated: false, + edgePurgeAccepted: false, + }); + assert.equal((await read("/purge-coalesced")).status, 404); + assert.equal((await r2Objects()).objects.length, 0); +}); + +test("purge prevents an initial slow write from creating an entry", async () => { + const write = put("/purge-cold-write", "too-late", { bodyDelayMs: 300 }); + for (let attempt = 0; attempt < 50; attempt++) { + if ((await metadataRowCount("pending_objects")) === 1) break; + await new Promise((resolve) => setTimeout(resolve, 10)); + } + assert.equal(await metadataRowCount("pending_objects"), 1); + + await purge({ purgeEverything: true }); + assert.deepEqual((await write).json, { + backingStoreUpdated: false, + edgePurgeAccepted: false, + }); + assert.equal((await read("/purge-cold-write")).status, 404); + assert.equal((await r2Objects()).objects.length, 0); +}); + +test("repeated purge prevents a post-tombstone write from resurrecting an entry", async () => { + await put("/purge-twice", "seed"); + await purge({ purgeEverything: true }); + const write = put("/purge-twice", "too-late", { bodyDelayMs: 300 }); + await new Promise((resolve) => setTimeout(resolve, 50)); + + await purge({ purgeEverything: true }); + assert.deepEqual((await write).json, { + backingStoreUpdated: false, + edgePurgeAccepted: false, + }); + assert.equal((await read("/purge-twice")).status, 404); +}); + +test("tag purge prevents a pending tagged write from publishing", async () => { + const write = put("/purge-pending-tag", "too-late", { + bodyDelayMs: 300, + tags: ["pending-tag"], + }); + await new Promise((resolve) => setTimeout(resolve, 50)); + + await purge({ tags: ["pending-tag"] }); + assert.deepEqual((await write).json, { + backingStoreUpdated: false, + edgePurgeAccepted: false, + }); + assert.equal((await read("/purge-pending-tag")).status, 404); +}); + +test("an expired revalidation claim cannot publish after its replacement", async () => { + await put("/claim-replacement", "seed"); + const [entry] = await metadata(); + const stub = await metadataStub(); + const first = await stub.claimRevalidation( + entry.keyHash, + entry.activeRevision, + entry.cacheKey, + "runtime-cache/poc-v2/claim-replacement", + 100, + 1, + ); + const second = await stub.claimRevalidation( + entry.keyHash, + entry.activeRevision, + entry.cacheKey, + "runtime-cache/poc-v2/claim-replacement", + 102, + 100, + ); + assert.ok(first); + assert.ok(second); + + const result = await stub.publish( + entry.keyHash, + first.revision, + { + objectKey: first.objectKey, + statusText: "", + responseHeaders: [], + freshUntil: 1_000, + swrUntil: 1_000, + revalidator: null, + cacheTags: [], + fenceTags: [], + }, + first.claimId, + ); + assert.equal(result.published, false); +}); + +test("a revalidation claim cannot replace a newer active revision", async () => { + await put("/claim-active-revision", "seed"); + const [entry] = await metadata(); + const stub = await metadataStub(); + const write = await stub.reserveWrite( + entry.keyHash, + entry.cacheKey, + "runtime-cache/poc-v2/claim-active-revision", + Date.now(), + ); + const claim = await stub.claimRevalidation( + entry.keyHash, + entry.activeRevision, + entry.cacheKey, + "runtime-cache/poc-v2/claim-active-revision", + 100, + 100, + ); + assert.ok(claim); + + const candidate = { + statusText: "", + responseHeaders: [], + freshUntil: 1_000, + swrUntil: 1_000, + revalidator: null, + cacheTags: [], + fenceTags: [], + }; + assert.equal( + ( + await stub.publish(entry.keyHash, write.revision, { + ...candidate, + objectKey: write.objectKey, + }) + ).published, + true, + ); + assert.equal( + ( + await stub.publish( + entry.keyHash, + claim.revision, + { ...candidate, objectKey: claim.objectKey }, + claim.claimId, + ) + ).published, + false, + ); +}); + +test("a write reserved after a tag purge is not rejected by its timestamp", async () => { + const createdAt = Date.now(); + await purge({ tags: ["already-purged"] }); + + const stub = await metadataStub(); + const reservation = await stub.reserveWrite( + "post-purge-write", + "/post-purge-write", + "runtime-cache/poc-v2/post-purge-write", + createdAt, + ); + const result = await stub.publish("post-purge-write", reservation.revision, { + objectKey: reservation.objectKey, + statusText: "", + responseHeaders: [], + freshUntil: createdAt + 60_000, + swrUntil: createdAt + 60_000, + revalidator: null, + cacheTags: ["already-purged"], + fenceTags: ["already-purged"], + }); + + assert.equal(result.published, true); +}); + +test("the previous metadata schema is upgraded in place", async () => { + const persistencePath = await mkdtemp(path.join(tmpdir(), "response-store-migration-")); + let legacy; + let upgraded; + + try { + legacy = new Miniflare({ + compatibilityDate: "2026-04-08", + resourcePersistencePath: persistencePath, + unsafeEphemeralDurableObjects: true, + workers: [ + { + name: "migration-worker", + modules: true, + script: ` + import { DurableObject } from "cloudflare:workers"; + export class CacheMetadata extends DurableObject { + constructor(ctx, env) { + super(ctx, env); + ctx.blockConcurrencyWhile(async () => ctx.storage.sql.exec(\` + CREATE TABLE tag_invalidations ( + tag TEXT PRIMARY KEY, + invalidated_at INTEGER NOT NULL + ) WITHOUT ROWID; + CREATE TABLE metadata_schema_migrations (version INTEGER PRIMARY KEY); + INSERT INTO metadata_schema_migrations (version) VALUES (1); + CREATE TABLE pending_objects ( + object_key TEXT PRIMARY KEY, + created_at INTEGER NOT NULL + ); + \`)); + } + seed() { + this.ctx.storage.sql.exec( + "INSERT INTO tag_invalidations (tag, invalidated_at) VALUES ('old-tag', 123)" + ); + } + } + export default { fetch() { return new Response("ok"); } }; + `, + durableObjects: { + CACHE_METADATA: { className: "CacheMetadata", useSQLite: true }, + }, + }, + ], + }); + const legacyNamespace = await legacy.getDurableObjectNamespace( + "CACHE_METADATA", + "migration-worker", + ); + await legacyNamespace.getByName(metadataName).seed(); + await legacy.dispose(); + legacy = undefined; + + upgraded = new Miniflare({ + compatibilityDate: "2026-04-08", + compatibilityFlags: ["nodejs_compat"], + resourcePersistencePath: persistencePath, + unsafeEphemeralDurableObjects: true, + workers: [ + { + name: "migration-worker", + compatibilityDate: "2026-04-08", + compatibilityFlags: ["nodejs_compat"], + modules: true, + scriptPath: workerScript, + durableObjects: { + CACHE_METADATA: { className: "CacheMetadata", useSQLite: true }, + }, + r2Buckets: { CACHE_BODIES: "migration-test" }, + bindings: { + CF_VERSION_METADATA: { + id: metadataName, + tag: "test", + timestamp: "2026-09-04T00:00:00Z", + }, + }, + }, + ], + }); + const upgradedNamespace = await upgraded.getDurableObjectNamespace( + "CACHE_METADATA", + "migration-worker", + ); + const stub = upgradedNamespace.getByName(metadataName); + const reservation = await stub.reserveWrite( + "migrated-write", + "/migrated-write", + "runtime-cache/poc-v2/migrated-write", + Date.now(), + ); + await stub.purgeMatching({ tags: ["new-tag"] }); + + assert.ok(reservation.objectKey); + assert.equal(await stub.getTagExpiration(["old-tag"]), 123); + assert.ok((await stub.getTagExpiration(["new-tag"])) > 123); + } finally { + await legacy?.dispose(); + await upgraded?.dispose(); + await rm(persistencePath, { force: true, recursive: true }); + } +}); + +test("retention sweep removes orphaned candidates without deleting active R2 objects", async () => { await put("/active-cleanup", "active"); const bucket = await mf.getR2Bucket("CACHE_BODIES", "user-worker"); const stub = await metadataStub(); @@ -530,10 +936,7 @@ test("retention cleanup removes orphaned candidates without deleting active R2 o await stub.trackPendingObject(activeObjectKey, 0); await stub.trackPendingObjects([orphanObjectKey], 0); - await put("/cleanup-trigger", "trigger"); - for (let attempt = 0; attempt < 20 && (await bucket.head(orphanObjectKey)); attempt++) { - await new Promise((resolve) => setTimeout(resolve, 25)); - } + assert.equal(await stub.sweepExpiredPendingObjects(1), 1); assert.equal(await bucket.head(orphanObjectKey), null); assert.notEqual(await bucket.head(activeObjectKey), null); @@ -545,6 +948,51 @@ test("retention cleanup removes orphaned candidates without deleting active R2 o assert.deepEqual(await stub.listExpiredPendingObjects(1, finishedKeys.length), []); }); +test("retention cleanup fences a body recreated after its first delete", async () => { + const stub = await metadataStub(); + const bucket = await mf.getR2Bucket("CACHE_BODIES", "user-worker"); + const createdAt = Date.now(); + const reservation = await stub.reserveWrite( + "expired-reservation", + "/expired-reservation", + "runtime-cache/poc-v2/expired-reservation", + createdAt, + ); + + assert.equal(await stub.sweepExpiredPendingObjects(createdAt + 1), 1); + assert.equal(await metadataRowCount("pending_objects"), 1); + await bucket.put(reservation.objectKey, "recreated-after-delete"); + + const result = await stub.publish("expired-reservation", reservation.revision, { + objectKey: reservation.objectKey, + statusText: "", + responseHeaders: [], + freshUntil: 1_000, + swrUntil: 1_000, + revalidator: null, + cacheTags: [], + fenceTags: [], + }); + assert.equal(result.published, false); + assert.equal(await bucket.head(reservation.objectKey), null); + assert.equal(await metadataRowCount("pending_objects"), 0); +}); + +test("replacement and purge clean their durable object markers", async () => { + await put("/replacement-cleanup", "first"); + const stub = await metadataStub(); + assert.equal(await metadataRowCount("pending_objects"), 0); + + await put("/replacement-cleanup", "second"); + assert.equal(await metadataRowCount("pending_objects"), 0); + assert.equal((await r2Objects()).objects.length, 1); + + await purge({ pathPrefixes: ["/replacement-cleanup"] }); + assert.equal(await metadataRowCount("pending_objects"), 0); + assert.deepEqual(await stub.listExpiredPendingObjects(Date.now() + 1, 10), []); + assert.equal((await r2Objects()).objects.length, 0); +}); + test("purge tombstones an entry before a slow regeneration can publish", async () => { await put("/purge-race", "seed", { tags: ["purge-race"], diff --git a/packages/workers-response-store/tests/service-binding-e2e.test.mjs b/packages/workers-response-store/tests/service-binding-e2e.test.mjs index aa893d6b4..00afcb42a 100644 --- a/packages/workers-response-store/tests/service-binding-e2e.test.mjs +++ b/packages/workers-response-store/tests/service-binding-e2e.test.mjs @@ -92,6 +92,16 @@ test("a service-bound cache Worker stores and returns responses", async () => { assert.equal(response.headers.get("X-Workers-Response-Store"), "R2-FRESH"); }); +test("a service-bound cache Worker resolves authoritative tag expirations", async () => { + const initial = await worker.fetch("https://user.test/admin/tag-expiration", { + method: "POST", + body: JSON.stringify({ tags: ["unchanged"] }), + }); + const body = await initial.json(); + assert.equal(initial.status, 200, JSON.stringify(body)); + assert.deepEqual(body, { expiration: 0 }); +}); + test("manual refresh calls back into the user Worker version", async () => { await put("/manual", "seed", { regeneratedBody: "manually-regenerated" });