From 8cbf30a35003c85afbae7e92cb7bb4276aec1446 Mon Sep 17 00:00:00 2001 From: James Date: Thu, 10 Sep 2026 11:30:59 +0100 Subject: [PATCH 01/17] perf(cache): reduce response store Durable Object load --- examples/workers-cache/README.md | 2 +- .../cache/response-store-adapter.worker.ts | 7 +- .../src/cache/response-store-data.runtime.ts | 1 + .../tests/response-store-adapter.e2e.test.ts | 2 +- .../tests/response-store-data.test.ts | 3 +- packages/workers-response-store/README.md | 4 +- .../example/service-binding/user-worker.ts | 5 + .../workers-response-store/example/worker.ts | 13 + .../workers-response-store/src/binding.ts | 279 ++++---- .../workers-response-store/src/metadata-do.ts | 621 +++++++++++++----- .../workers-response-store/tests/e2e.test.mjs | 117 +++- .../tests/service-binding-e2e.test.mjs | 10 + 12 files changed, 725 insertions(+), 339 deletions(-) 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..9d43c7add 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; @@ -99,6 +106,7 @@ 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", }); @@ -188,6 +196,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..403f3fdc8 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; }; @@ -102,6 +104,9 @@ type CacheKey = { }; type WriteReservation = CacheKey & { + claimId?: string; + coalesced?: boolean; + objectKey: string; revision: number; }; @@ -111,8 +116,8 @@ type StoreResult = { }; type PublicationResult = { + entry: StoredEntry | null; published: boolean; - previousObjectKey?: string; }; export type WorkersResponseStoreEnv = { @@ -207,9 +212,6 @@ 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]); @@ -339,6 +341,30 @@ 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, + coalesce = false, + ): Promise { + const reservation = await metadata.reserveWrite( + keyHash, + cacheKey, + this.objectKeyPrefix(keyHash), + Date.now(), + coalesce, + ); + return { cacheKey, keyHash, ...reservation }; + } + private logCleanupFailure(objectKey: string, error: unknown): void { console.error( JSON.stringify({ @@ -349,56 +375,13 @@ export class ResponseStoreBinding extends WorkerEntrypoint< ); } - private async deletePendingObjects( - metadata: CacheMetadataStub, - objectKeys: string[], - ): 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( + private async releaseFailedWrite( metadata: CacheMetadataStub, - objectKeys: string[], - createdAt: number, + write: Pick, ): 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 { @@ -456,46 +439,45 @@ export class ResponseStoreBinding extends WorkerEntrypoint< revalidator: ResponseStorePutOptions["revalidator"], reservation?: WriteReservation, ): 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, false)); + 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 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, + }; + // 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 +489,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 +513,30 @@ 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)); - 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 +559,7 @@ export class ResponseStoreBinding extends WorkerEntrypoint< entry.keyHash, entry.activeRevision, entry.cacheKey, + this.objectKeyPrefix(entry.keyHash), Date.now(), BACKGROUND_REVALIDATION_LEASE_MS, ); @@ -604,7 +570,9 @@ export class ResponseStoreBinding extends WorkerEntrypoint< try { await this.regenerateEntry(metadata, entry, "swr", { cacheKey: entry.cacheKey, + claimId: claim.claimId, keyHash: entry.keyHash, + objectKey: claim.objectKey, revision: claim.revision, }); } catch (error) { @@ -615,8 +583,6 @@ export class ResponseStoreBinding extends WorkerEntrypoint< error: error instanceof Error ? error.message : String(error), }), ); - } finally { - await metadata.finishRevalidation(entry.keyHash, claim.claimId); } } @@ -668,9 +634,21 @@ 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); + const { cacheKey, keyHash } = await this.deriveCacheKey(request); + const reservation = await this.reserveWrite(metadata, keyHash, cacheKey, options.coalesce); + if (reservation.coalesced) { + await response.body?.cancel().catch(() => {}); + return { backingStoreUpdated: false, edgePurgeAccepted: false }; + } + + const result = await this.storeResponse( + metadata, + request, + response, + options.revalidator, + reservation, + ); if (!result.published || !result.entry) { return { backingStoreUpdated: false, edgePurgeAccepted: false }; } @@ -688,14 +666,25 @@ 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, + keyHash: entry.keyHash, + ...reservation, + } + : undefined, + ); return result.published ? result.entry : null; }), ); @@ -720,7 +709,7 @@ export class ResponseStoreBinding extends WorkerEntrypoint< } return { - backingStoreUpdated: refreshed.length === activeEntries.length, + backingStoreUpdated: refreshed.length === candidates.length, edgePurgeAccepted, }; } @@ -742,12 +731,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..cdb217d8b 100644 --- a/packages/workers-response-store/src/metadata-do.ts +++ b/packages/workers-response-store/src/metadata-do.ts @@ -11,44 +11,64 @@ 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 = { + coalesced: boolean; + 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, + coalesce?: boolean, + ): 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 & { @@ -70,18 +90,28 @@ 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; +const WRITE_COALESCE_LEASE_MS = 30_000; + +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 || @@ -138,8 +168,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(` @@ -170,12 +202,6 @@ 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 @@ -190,40 +216,33 @@ export class CacheMetadata extends DurableObject> { CREATE INDEX IF NOT EXISTS pending_objects_created_at ON pending_objects(created_at); `); - const hasBackfilledTagIndex = + const legacyTagIndexIsActive = 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`, - ) - .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, - ); - } - } - - ctx.storage.sql.exec("INSERT INTO metadata_schema_migrations (version) VALUES (1)"); - }); + if (legacyTagIndexIsActive) { + // A rollback rebuilds this index and restores the marker. Clear it once + // when moving forward so ordinary writes do not retain duplicate rows. + ctx.storage.sql.exec(` + DELETE FROM entry_tags; + DELETE FROM metadata_schema_migrations WHERE version = 1; + `); } }); } + private async ensureCleanupAlarm(createdAt: number): Promise { + if (this.cleanupAlarmKnown) return; + + const current = await this.ctx.storage.getAlarm(); + if (current === null) { + await this.ctx.storage.setAlarm(createdAt + ORPHAN_RETENTION_MS); + } + this.cleanupAlarmKnown = true; + } + private findMatchingEntryRows(options: ResponseStorePurgeOptions): EntryRow[] { if (options.purgeEverything) { return this.ctx.storage.sql @@ -233,63 +252,58 @@ export class CacheMetadata extends DurableObject> { .toArray(); } - const matches = new Map(); - const tags = normalizeTags(options.tags ?? []); - - 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 tags = new Set(normalizeTags(options.tags ?? [])); 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(); - - for (const row of rows) { - if (prefixes.some((prefix) => row.cache_key.startsWith(prefix))) { - matches.set(row.key_hash, row); - } - } - } - - return [...matches.values()]; + return this.ctx.storage.sql + .exec("SELECT * FROM entries WHERE tombstoned = 0 AND active_revision IS NOT NULL") + .toArray() + .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) VALUES ${batch + .map(() => "(?, ?)") + .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 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 +312,32 @@ 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) { + console.error( + JSON.stringify({ + message: "Workers Response Store R2 cleanup failed", + objectKeys: batch, + error: error instanceof Error ? error.message : String(error), + }), + ); + } + } + } + listExpiredPendingObjects(cutoff: number, limit: number): string[] { return this.ctx.storage.sql .exec<{ object_key: string }>( @@ -320,18 +354,82 @@ 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 }>( + `SELECT pending_objects.object_key, pending_objects.created_at, + 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 expiredObjectKeys = batch + .filter(({ active, created_at }) => !active && created_at <= cutoff) + .map(({ object_key }) => object_key); + this.finishPendingObjects(activeObjectKeys); + if (expiredObjectKeys.length) { + await this.env.CACHE_BODIES.delete(expiredObjectKeys); + this.finishPendingObjects(expiredObjectKeys); + } + + 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) { + await this.ensureCleanupAlarm(next.created_at); + } + + 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 +438,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,17 +463,42 @@ export class CacheMetadata extends DurableObject> { now, now + leaseMs, ); + this.ctx.storage.sql.exec( + "INSERT OR REPLACE INTO pending_objects (object_key, created_at) VALUES (?, ?)", + 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 { @@ -391,38 +509,93 @@ export class CacheMetadata extends DurableObject> { keyHash, ) .toArray()[0]; - const revision = (current?.latest_revision ?? 0) + 1; + return this.reserveRevision(keyHash, cacheKey, current); + }); + } - 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, - ); + async reserveWrite( + keyHash: string, + cacheKey: string, + objectKeyPrefix: string, + createdAt: number, + coalesce = false, + ): Promise { + const reservation = this.ctx.storage.transactionSync(() => { + const current: + | { + active_revision: number | null; + latest_revision: number; + write_pending?: number; + } + | undefined = coalesce + ? this.ctx.storage.sql + .exec<{ + active_revision: number | null; + latest_revision: number; + write_pending: number; + }>( + `SELECT entries.active_revision, entries.latest_revision, + pending_objects.object_key IS NOT NULL AS write_pending + FROM entries + LEFT JOIN pending_objects ON pending_objects.object_key = ? || '/' || entries.latest_revision + AND pending_objects.created_at > ? + WHERE entries.key_hash = ?`, + objectKeyPrefix, + createdAt - WRITE_COALESCE_LEASE_MS, + keyHash, + ) + .toArray()[0] + : this.ctx.storage.sql + .exec<{ active_revision: number | null; latest_revision: number }>( + "SELECT active_revision, latest_revision FROM entries WHERE key_hash = ?", + keyHash, + ) + .toArray()[0]; + if ( + coalesce && + current?.active_revision !== current?.latest_revision && + current?.write_pending + ) { + const objectKey = `${objectKeyPrefix}/${current.latest_revision}`; + return { coalesced: true, objectKey, revision: current.latest_revision }; } - 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) VALUES (?, ?)", + objectKey, + createdAt, + ); + return { coalesced: false, 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 = ?", - keyHash, - ) + .exec("SELECT * FROM entries WHERE key_hash = ?", keyHash) .toArray()[0]; if (!current || current.latest_revision !== revision) { - return { published: false }; + 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; @@ -450,22 +623,67 @@ export class CacheMetadata extends DurableObject> { 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) VALUES (?, ?)", + 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: revision, + objectKey: metadata.objectKey, + statusText: metadata.statusText, + responseHeaders: metadata.responseHeaders, + freshUntil: metadata.freshUntil, + swrUntil: metadata.swrUntil, + revalidator: metadata.revalidator, + cacheTags: metadata.cacheTags, + ...(responseMetadataIsInR2 + ? {} + : { + legacyResponseMetadata: { + status: metadata.status, + createdAt: metadata.createdAt, + initialAge: metadata.initialAge, + }, + }), + }; 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.ensureCleanupAlarm(Date.now()); + await this.deleteTrackedObjects([cleanupObjectKey]); + } + return result; } getEntry(keyHash: string): StoredEntry | null { @@ -479,8 +697,7 @@ export class CacheMetadata extends DurableObject> { 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 }>( @@ -495,26 +712,97 @@ export class CacheMetadata extends DurableObject> { return expiration; } - getEntriesMatching(options: ResponseStoreRefreshOptions): StoredEntry[] { - return storedEntriesFromRows(this.findMatchingEntryRows(options)); + async reserveRefresh( + options: ResponseStoreRefreshOptions, + objectKeyRoot: string, + createdAt: number, + ): Promise { + const candidates = this.ctx.storage.transactionSync(() => { + const matches = this.findMatchingEntryRows(options); + 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) VALUES ${batch + .map(() => "(?, ?)") + .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; } - purgeMatching(options: ResponseStorePurgeOptions): PurgedEntry[] { - return this.ctx.storage.transactionSync(() => { + async purgeMatching(options: ResponseStorePurgeOptions): Promise { + const invalidatedAt = Date.now(); + const matches = this.ctx.storage.transactionSync(() => { const matches = this.findMatchingEntryRows(options); - const invalidatedAt = Date.now(); - for (const tag of normalizeTags(options.tags ?? [])) { + const tags = normalizeTags(options.tags ?? []); + 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) VALUES ${batch + .map(() => "(?, ?)") + .join(", ")} ON CONFLICT(tag) DO UPDATE SET invalidated_at = MAX(tag_invalidations.invalidated_at, excluded.invalidated_at)`, - tag, - invalidatedAt, + ...batch.flatMap((tag) => [tag, invalidatedAt]), ); } - for (const row of matches) { + for (const batch of batches(matches, MAX_SQL_PARAMETERS)) { + 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) + SELECT object_key, ? 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, @@ -531,14 +819,13 @@ export class CacheMetadata extends DurableObject> { 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) => ({ @@ -547,6 +834,11 @@ export class CacheMetadata extends DurableObject> { objectKey: row.object_key!, })); }); + if (matches.length) { + await this.ensureCleanupAlarm(invalidatedAt); + await this.deleteTrackedObjects(matches.map((entry) => entry.objectKey)); + } + return matches; } inspect(): StoredEntry[] { @@ -555,11 +847,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..f03b0927e 100644 --- a/packages/workers-response-store/tests/e2e.test.mjs +++ b/packages/workers-response-store/tests/e2e.test.mjs @@ -15,6 +15,7 @@ beforeEach(async () => { compatibilityDate: "2026-04-08", compatibilityFlags: ["nodejs_compat"], unsafeEphemeralDurableObjects: true, + unsafeInspectDurableObjects: true, workers: [ { name: "user-worker", @@ -64,6 +65,8 @@ 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)); const response = await worker.fetch(`https://user.test/admin/put${path}`, { method: "PUT", @@ -104,6 +107,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 +127,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() { @@ -155,6 +173,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 () => { @@ -304,6 +323,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 +443,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 +472,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 +497,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"); @@ -520,7 +562,34 @@ 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: false, edgePurgeAccepted: false }); + assert.equal(await (await read("/coalesced")).text(), "first"); + assert.equal((await metadata())[0].activeRevision, 1); + assert.equal((await r2Objects()).objects.length, 1); +}); + +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("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 +599,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 +611,21 @@ test("retention cleanup removes orphaned candidates without deleting active R2 o assert.deepEqual(await stub.listExpiredPendingObjects(1, finishedKeys.length), []); }); +test("replacement and purge clean their durable object markers", async () => { + await put("/legacy-cleanup", "first"); + const stub = await metadataStub(); + assert.equal(await metadataRowCount("pending_objects"), 0); + + await put("/legacy-cleanup", "second"); + assert.equal(await metadataRowCount("pending_objects"), 0); + assert.equal((await r2Objects()).objects.length, 1); + + await purge({ pathPrefixes: ["/legacy-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" }); From d04d5f8242d632646383b6e72b43073ae3a75cf0 Mon Sep 17 00:00:00 2001 From: James Date: Sun, 13 Sep 2026 22:53:37 +0100 Subject: [PATCH 02/17] refactor(cache): remove response store compatibility shims --- .../workers-response-store/src/binding.ts | 29 ++-------- .../workers-response-store/src/metadata-do.ts | 56 +------------------ .../workers-response-store/tests/e2e.test.mjs | 37 +----------- 3 files changed, 10 insertions(+), 112 deletions(-) diff --git a/packages/workers-response-store/src/binding.ts b/packages/workers-response-store/src/binding.ts index 403f3fdc8..6335057f0 100644 --- a/packages/workers-response-store/src/binding.ts +++ b/packages/workers-response-store/src/binding.ts @@ -69,23 +69,13 @@ type EntryMetadata = { cacheTags: string[]; }; -export type CandidateMetadata = EntryMetadata & { - status: number; - createdAt: number; - initialAge: number; - responseMetadataInR2?: true; -}; +export type CandidateMetadata = EntryMetadata; export type StoredEntry = EntryMetadata & { keyHash: string; cacheKey: string; activeRevision: number; latestRevision: number; - legacyResponseMetadata?: { - status: number; - createdAt: number; - initialAge: number; - }; }; export type PurgedEntry = { @@ -390,13 +380,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 || @@ -463,19 +449,12 @@ export class ResponseStoreBinding extends WorkerEntrypoint< ]; 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, }; // RPC-transferred Response streams do not retain the fixed-length marker diff --git a/packages/workers-response-store/src/metadata-do.ts b/packages/workers-response-store/src/metadata-do.ts index cdb217d8b..8ee3c04b5 100644 --- a/packages/workers-response-store/src/metadata-do.ts +++ b/packages/workers-response-store/src/metadata-do.ts @@ -77,11 +77,8 @@ type EntryRow = Record & { active_revision: number | null; latest_revision: number; object_key: string | null; - status: 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; @@ -124,7 +121,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, @@ -143,16 +140,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[] { @@ -181,11 +168,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, @@ -206,30 +190,12 @@ export class CacheMetadata extends DurableObject { tag TEXT PRIMARY KEY, invalidated_at INTEGER NOT NULL ) WITHOUT ROWID; - CREATE TABLE IF NOT EXISTS metadata_schema_migrations ( - version INTEGER PRIMARY KEY - ); CREATE TABLE IF NOT EXISTS pending_objects ( object_key TEXT PRIMARY KEY, created_at INTEGER NOT NULL ); CREATE INDEX IF NOT EXISTS pending_objects_created_at ON pending_objects(created_at); `); - - const legacyTagIndexIsActive = - ctx.storage.sql - .exec<{ version: number }>( - "SELECT version FROM metadata_schema_migrations WHERE version = 1", - ) - .toArray().length > 0; - if (legacyTagIndexIsActive) { - // A rollback rebuilds this index and restores the marker. Clear it once - // when moving forward so ordinary writes do not retain duplicate rows. - ctx.storage.sql.exec(` - DELETE FROM entry_tags; - DELETE FROM metadata_schema_migrations WHERE version = 1; - `); - } }); } @@ -598,20 +564,16 @@ export class CacheMetadata extends DurableObject { }; } - 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 = ?`, 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, @@ -655,15 +617,6 @@ export class CacheMetadata extends DurableObject { swrUntil: metadata.swrUntil, revalidator: metadata.revalidator, cacheTags: metadata.cacheTags, - ...(responseMetadataIsInR2 - ? {} - : { - legacyResponseMetadata: { - status: metadata.status, - createdAt: metadata.createdAt, - initialAge: metadata.initialAge, - }, - }), }; return { @@ -808,11 +761,8 @@ export class CacheMetadata extends DurableObject { latest_revision = latest_revision + 1, active_revision = NULL, 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, diff --git a/packages/workers-response-store/tests/e2e.test.mjs b/packages/workers-response-store/tests/e2e.test.mjs index f03b0927e..617a0ebee 100644 --- a/packages/workers-response-store/tests/e2e.test.mjs +++ b/packages/workers-response-store/tests/e2e.test.mjs @@ -156,7 +156,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"); @@ -186,36 +185,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", @@ -612,15 +581,15 @@ test("retention sweep removes orphaned candidates without deleting active R2 obj }); test("replacement and purge clean their durable object markers", async () => { - await put("/legacy-cleanup", "first"); + await put("/replacement-cleanup", "first"); const stub = await metadataStub(); assert.equal(await metadataRowCount("pending_objects"), 0); - await put("/legacy-cleanup", "second"); + await put("/replacement-cleanup", "second"); assert.equal(await metadataRowCount("pending_objects"), 0); assert.equal((await r2Objects()).objects.length, 1); - await purge({ pathPrefixes: ["/legacy-cleanup"] }); + 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); From 8b006188068f3889354603c390b2accfd88747da Mon Sep 17 00:00:00 2001 From: James Date: Sun, 13 Sep 2026 23:22:08 +0100 Subject: [PATCH 03/17] fix(cache): batch response store purge parameters --- packages/workers-response-store/src/metadata-do.ts | 2 +- packages/workers-response-store/tests/e2e.test.mjs | 13 +++++++++++++ 2 files changed, 14 insertions(+), 1 deletion(-) diff --git a/packages/workers-response-store/src/metadata-do.ts b/packages/workers-response-store/src/metadata-do.ts index 8ee3c04b5..dca4945a5 100644 --- a/packages/workers-response-store/src/metadata-do.ts +++ b/packages/workers-response-store/src/metadata-do.ts @@ -746,7 +746,7 @@ export class CacheMetadata extends DurableObject { ); } - for (const batch of batches(matches, MAX_SQL_PARAMETERS)) { + 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( diff --git a/packages/workers-response-store/tests/e2e.test.mjs b/packages/workers-response-store/tests/e2e.test.mjs index 617a0ebee..c8b7260a9 100644 --- a/packages/workers-response-store/tests/e2e.test.mjs +++ b/packages/workers-response-store/tests/e2e.test.mjs @@ -519,6 +519,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)); From b1481f91b3ee9b6b40b367f8bcea62e9ae6dc524 Mon Sep 17 00:00:00 2001 From: James Date: Sun, 13 Sep 2026 23:26:40 +0100 Subject: [PATCH 04/17] fix(cache): preserve coalesced response store writes --- .../workers-response-store/src/binding.ts | 71 +++++++++++++------ .../workers-response-store/src/metadata-do.ts | 50 ++----------- .../workers-response-store/tests/e2e.test.mjs | 19 ++++- 3 files changed, 73 insertions(+), 67 deletions(-) diff --git a/packages/workers-response-store/src/binding.ts b/packages/workers-response-store/src/binding.ts index 6335057f0..bae800c68 100644 --- a/packages/workers-response-store/src/binding.ts +++ b/packages/workers-response-store/src/binding.ts @@ -95,7 +95,6 @@ type CacheKey = { type WriteReservation = CacheKey & { claimId?: string; - coalesced?: boolean; objectKey: string; revision: number; }; @@ -206,6 +205,7 @@ 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) { @@ -343,14 +343,12 @@ export class ResponseStoreBinding extends WorkerEntrypoint< metadata: CacheMetadataStub, keyHash: string, cacheKey: string, - coalesce = false, ): Promise { const reservation = await metadata.reserveWrite( keyHash, cacheKey, this.objectKeyPrefix(keyHash), Date.now(), - coalesce, ); return { cacheKey, keyHash, ...reservation }; } @@ -427,8 +425,7 @@ export class ResponseStoreBinding extends WorkerEntrypoint< ): Promise { const cacheKey = reservation ?? (await this.deriveCacheKey(request)); const write = - reservation ?? - (await this.reserveWrite(metadata, cacheKey.keyHash, cacheKey.cacheKey, false)); + reservation ?? (await this.reserveWrite(metadata, cacheKey.keyHash, cacheKey.cacheKey)); const { keyHash, objectKey, revision } = write; let publication: PublicationResult; @@ -615,28 +612,56 @@ export class ResponseStoreBinding extends WorkerEntrypoint< const metadata = this.getMetadata(); const { cacheKey, keyHash } = await this.deriveCacheKey(request); - const reservation = await this.reserveWrite(metadata, keyHash, cacheKey, options.coalesce); - if (reservation.coalesced) { - await response.body?.cancel().catch(() => {}); - return { backingStoreUpdated: false, edgePurgeAccepted: false }; + const pendingPutKey = `${this.getVersionId()}:${keyHash}`; + if (options.coalesce) { + for (;;) { + const pending = pendingPuts.get(pendingPutKey); + if (!pending) break; + + try { + const result = await pending; + if (result.backingStoreUpdated) { + await response.body?.cancel().catch(() => {}); + return result; + } + } catch { + // Preserve this response as the fallback when the leading write fails. + } + if (pendingPuts.get(pendingPutKey) === pending) { + pendingPuts.delete(pendingPutKey); + } + } } - const result = await this.storeResponse( - metadata, - request, - response, - options.revalidator, - reservation, - ); - if (!result.published || !result.entry) { - return { backingStoreUpdated: false, edgePurgeAccepted: false }; - } + const write = (async (): Promise => { + const reservation = await this.reserveWrite(metadata, keyHash, cacheKey); + const result = await this.storeResponse( + metadata, + request, + response, + options.revalidator, + reservation, + ); + if (!result.published || !result.entry) { + return { backingStoreUpdated: false, edgePurgeAccepted: false }; + } - const edgePurgeAccepted = options.purgeExisting - ? await this.purgeEdgeCacheByTags([purgeTagForEntry(result.entry)]) - : true; + const edgePurgeAccepted = options.purgeExisting + ? await this.purgeEdgeCacheByTags([purgeTagForEntry(result.entry)]) + : true; - return { backingStoreUpdated: true, edgePurgeAccepted }; + return { backingStoreUpdated: true, edgePurgeAccepted }; + })(); + if (!options.coalesce) return write; + + pendingPuts.set(pendingPutKey, write); + try { + return await write; + } finally { + if (pendingPuts.get(pendingPutKey) === write) { + pendingPuts.delete(pendingPutKey); + } + } } async refresh(options: ResponseStoreRefreshOptions): Promise { diff --git a/packages/workers-response-store/src/metadata-do.ts b/packages/workers-response-store/src/metadata-do.ts index dca4945a5..7e27d4b14 100644 --- a/packages/workers-response-store/src/metadata-do.ts +++ b/packages/workers-response-store/src/metadata-do.ts @@ -21,7 +21,6 @@ type PublicationResult = { }; type WriteReservation = { - coalesced: boolean; objectKey: string; revision: number; }; @@ -38,7 +37,6 @@ export type CacheMetadataStub = DurableObjectStub & { cacheKey: string, objectKeyPrefix: string, createdAt: number, - coalesce?: boolean, ): Promise; claimRevalidation( keyHash: string, @@ -92,7 +90,6 @@ 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; -const WRITE_COALESCE_LEASE_MS = 30_000; type CacheMetadataEnv = { CACHE_BODIES: R2Bucket; @@ -484,47 +481,14 @@ export class CacheMetadata extends DurableObject { cacheKey: string, objectKeyPrefix: string, createdAt: number, - coalesce = false, ): Promise { const reservation = this.ctx.storage.transactionSync(() => { - const current: - | { - active_revision: number | null; - latest_revision: number; - write_pending?: number; - } - | undefined = coalesce - ? this.ctx.storage.sql - .exec<{ - active_revision: number | null; - latest_revision: number; - write_pending: number; - }>( - `SELECT entries.active_revision, entries.latest_revision, - pending_objects.object_key IS NOT NULL AS write_pending - FROM entries - LEFT JOIN pending_objects ON pending_objects.object_key = ? || '/' || entries.latest_revision - AND pending_objects.created_at > ? - WHERE entries.key_hash = ?`, - objectKeyPrefix, - createdAt - WRITE_COALESCE_LEASE_MS, - keyHash, - ) - .toArray()[0] - : this.ctx.storage.sql - .exec<{ active_revision: number | null; latest_revision: number }>( - "SELECT active_revision, latest_revision FROM entries WHERE key_hash = ?", - keyHash, - ) - .toArray()[0]; - if ( - coalesce && - current?.active_revision !== current?.latest_revision && - current?.write_pending - ) { - const objectKey = `${objectKeyPrefix}/${current.latest_revision}`; - return { coalesced: true, objectKey, revision: current.latest_revision }; - } + const current = this.ctx.storage.sql + .exec<{ active_revision: number | null; latest_revision: number }>( + "SELECT active_revision, latest_revision FROM entries WHERE key_hash = ?", + keyHash, + ) + .toArray()[0]; const revision = this.reserveRevision(keyHash, cacheKey, current); const objectKey = `${objectKeyPrefix}/${revision}`; @@ -533,7 +497,7 @@ export class CacheMetadata extends DurableObject { objectKey, createdAt, ); - return { coalesced: false, objectKey, revision }; + return { objectKey, revision }; }); await this.ensureCleanupAlarm(createdAt); return reservation; diff --git a/packages/workers-response-store/tests/e2e.test.mjs b/packages/workers-response-store/tests/e2e.test.mjs index c8b7260a9..30d514c86 100644 --- a/packages/workers-response-store/tests/e2e.test.mjs +++ b/packages/workers-response-store/tests/e2e.test.mjs @@ -551,7 +551,7 @@ test("overlapping framework writes can be coalesced", async () => { const firstResult = await first; assert.deepEqual(firstResult.json, { backingStoreUpdated: true, edgePurgeAccepted: true }); - assert.deepEqual(second.json, { backingStoreUpdated: false, edgePurgeAccepted: false }); + 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); @@ -571,6 +571,23 @@ test("a failed coalesced write does not suppress an immediate retry", async () = 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("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"); From 47686c3cce47dab8f58026fae48c8b0069fd29b8 Mon Sep 17 00:00:00 2001 From: James Date: Sun, 13 Sep 2026 23:35:48 +0100 Subject: [PATCH 05/17] fix(cache): preserve concurrent response store writes --- .../workers-response-store/src/binding.ts | 5 +- .../workers-response-store/src/metadata-do.ts | 60 +++++++++++-------- .../workers-response-store/tests/e2e.test.mjs | 49 +++++++++++++++ 3 files changed, 87 insertions(+), 27 deletions(-) diff --git a/packages/workers-response-store/src/binding.ts b/packages/workers-response-store/src/binding.ts index bae800c68..adf8a486e 100644 --- a/packages/workers-response-store/src/binding.ts +++ b/packages/workers-response-store/src/binding.ts @@ -613,14 +613,17 @@ export class ResponseStoreBinding extends WorkerEntrypoint< const { cacheKey, keyHash } = await this.deriveCacheKey(request); const pendingPutKey = `${this.getVersionId()}:${keyHash}`; + let reservation: WriteReservation | undefined; if (options.coalesce) { for (;;) { const pending = pendingPuts.get(pendingPutKey); if (!pending) break; + reservation ??= await this.reserveWrite(metadata, keyHash, cacheKey); try { const result = await pending; if (result.backingStoreUpdated) { + await metadata.finishPendingObjects([reservation.objectKey]); await response.body?.cancel().catch(() => {}); return result; } @@ -634,7 +637,7 @@ export class ResponseStoreBinding extends WorkerEntrypoint< } const write = (async (): Promise => { - const reservation = await this.reserveWrite(metadata, keyHash, cacheKey); + reservation ??= await this.reserveWrite(metadata, keyHash, cacheKey); const result = await this.storeResponse( metadata, request, diff --git a/packages/workers-response-store/src/metadata-do.ts b/packages/workers-response-store/src/metadata-do.ts index 7e27d4b14..55fca8436 100644 --- a/packages/workers-response-store/src/metadata-do.ts +++ b/packages/workers-response-store/src/metadata-do.ts @@ -206,27 +206,29 @@ export class CacheMetadata extends DurableObject { this.cleanupAlarmKnown = true; } - private findMatchingEntryRows(options: ResponseStorePurgeOptions): EntryRow[] { + 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` + : "SELECT * FROM entries WHERE tombstoned = 0 AND active_revision IS NOT NULL", + ) + .toArray(); if (options.purgeEverything) { - return this.ctx.storage.sql - .exec( - "SELECT * FROM entries WHERE tombstoned = 0 AND active_revision IS NOT NULL", - ) - .toArray(); + return rows; } const tags = new Set(normalizeTags(options.tags ?? [])); const prefixes = options.pathPrefixes ?? []; - return this.ctx.storage.sql - .exec("SELECT * FROM entries WHERE tombstoned = 0 AND active_revision IS NOT NULL") - .toArray() - .filter( - (row) => - prefixes.some((prefix) => row.cache_key.startsWith(prefix)) || - normalizeTags(JSON.parse(row.cache_tags ?? "[]") as string[]).some((tag) => - tags.has(tag), - ), - ); + 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)), + ); } async trackPendingObjects(objectKeys: string[], createdAt: number): Promise { @@ -513,7 +515,11 @@ export class CacheMetadata extends DurableObject { const current = this.ctx.storage.sql .exec("SELECT * FROM entries WHERE key_hash = ?", keyHash) .toArray()[0]; - if (!current || current.latest_revision !== revision) { + if ( + !current || + revision > current.latest_revision || + (current.active_revision !== null && revision <= current.active_revision) + ) { if (claimId) { this.ctx.storage.sql.exec( "DELETE FROM revalidation_claims WHERE key_hash = ? AND claim_id = ?", @@ -533,7 +539,8 @@ export class CacheMetadata extends DurableObject { 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, metadata.statusText, @@ -545,6 +552,7 @@ export class CacheMetadata extends DurableObject { JSON.stringify(metadata.cacheTags), keyHash, revision, + revision, ); const published = update.rowsWritten === 1; @@ -573,7 +581,7 @@ export class CacheMetadata extends DurableObject { keyHash, cacheKey: current.cache_key, activeRevision: revision, - latestRevision: revision, + latestRevision: current.latest_revision, objectKey: metadata.objectKey, statusText: metadata.statusText, responseHeaders: metadata.responseHeaders, @@ -696,7 +704,7 @@ export class CacheMetadata extends DurableObject { async purgeMatching(options: ResponseStorePurgeOptions): Promise { const invalidatedAt = Date.now(); const matches = this.ctx.storage.transactionSync(() => { - const matches = this.findMatchingEntryRows(options); + const matches = this.findMatchingEntryRows(options, true); const tags = normalizeTags(options.tags ?? []); for (const batch of batches(tags, MAX_SQL_PARAMETERS / 2)) { @@ -723,7 +731,7 @@ export class CacheMetadata extends DurableObject { 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_text = NULL, response_headers = NULL, @@ -742,11 +750,11 @@ export class CacheMetadata extends DurableObject { ); } - 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.ensureCleanupAlarm(invalidatedAt); diff --git a/packages/workers-response-store/tests/e2e.test.mjs b/packages/workers-response-store/tests/e2e.test.mjs index 30d514c86..624d45b01 100644 --- a/packages/workers-response-store/tests/e2e.test.mjs +++ b/packages/workers-response-store/tests/e2e.test.mjs @@ -555,6 +555,7 @@ test("overlapping framework writes can be coalesced", async () => { 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 failed coalesced write does not suppress an immediate retry", async () => { @@ -588,6 +589,54 @@ test("a failed coalesced write preserves an overlapping successful write", async 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 }); + 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-cold-write")).status, 404); + assert.equal((await r2Objects()).objects.length, 0); +}); + 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"); From 29f6277a4a5446220fe26b2b04e703893be5910b Mon Sep 17 00:00:00 2001 From: James Date: Sun, 13 Sep 2026 23:42:35 +0100 Subject: [PATCH 06/17] fix(cache): fence response store tag purges --- .../workers-response-store/src/binding.ts | 45 ++++++++++++------- .../workers-response-store/src/metadata-do.ts | 32 ++++++------- .../workers-response-store/tests/e2e.test.mjs | 29 ++++++++++++ 3 files changed, 75 insertions(+), 31 deletions(-) diff --git a/packages/workers-response-store/src/binding.ts b/packages/workers-response-store/src/binding.ts index adf8a486e..d3c540cd6 100644 --- a/packages/workers-response-store/src/binding.ts +++ b/packages/workers-response-store/src/binding.ts @@ -69,7 +69,9 @@ type EntryMetadata = { cacheTags: string[]; }; -export type CandidateMetadata = EntryMetadata; +export type CandidateMetadata = EntryMetadata & { + fenceTags: string[]; +}; export type StoredEntry = EntryMetadata & { keyHash: string; @@ -95,6 +97,7 @@ type CacheKey = { type WriteReservation = CacheKey & { claimId?: string; + fenceTags: string[]; objectKey: string; revision: number; }; @@ -251,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 @@ -343,6 +357,7 @@ export class ResponseStoreBinding extends WorkerEntrypoint< metadata: CacheMetadataStub, keyHash: string, cacheKey: string, + cacheTags: string[], ): Promise { const reservation = await metadata.reserveWrite( keyHash, @@ -350,7 +365,7 @@ export class ResponseStoreBinding extends WorkerEntrypoint< this.objectKeyPrefix(keyHash), Date.now(), ); - return { cacheKey, keyHash, ...reservation }; + return { cacheKey, fenceTags: cacheTags, keyHash, ...reservation }; } private logCleanupFailure(objectKey: string, error: unknown): void { @@ -422,10 +437,12 @@ export class ResponseStoreBinding extends WorkerEntrypoint< response: Response, revalidator: ResponseStorePutOptions["revalidator"], reservation?: WriteReservation, + cacheTags = cacheTagsFromResponse(response), ): Promise { const cacheKey = reservation ?? (await this.deriveCacheKey(request)); const write = - reservation ?? (await this.reserveWrite(metadata, cacheKey.keyHash, cacheKey.cacheKey)); + reservation ?? + (await this.reserveWrite(metadata, cacheKey.keyHash, cacheKey.cacheKey, cacheTags)); const { keyHash, objectKey, revision } = write; let publication: PublicationResult; @@ -436,22 +453,15 @@ export class ResponseStoreBinding extends WorkerEntrypoint< 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 = { + fenceTags: [...new Set([...write.fenceTags, ...cacheTags])], objectKey, statusText: response.statusText, responseHeaders, freshUntil: policy.freshUntil, swrUntil: policy.swrUntil, revalidator: revalidator ?? null, - cacheTags: responseCacheTags, + cacheTags, }; // RPC-transferred Response streams do not retain the fixed-length marker @@ -491,7 +501,8 @@ export class ResponseStoreBinding extends WorkerEntrypoint< const cacheRequest = new Request(`https://runtime-cache.invalid${entry.cacheKey}`); const writeReservation = - reservation ?? (await this.reserveWrite(metadata, entry.keyHash, entry.cacheKey)); + reservation ?? + (await this.reserveWrite(metadata, entry.keyHash, entry.cacheKey, entry.cacheTags)); let response: Response; try { @@ -547,6 +558,7 @@ export class ResponseStoreBinding extends WorkerEntrypoint< await this.regenerateEntry(metadata, entry, "swr", { cacheKey: entry.cacheKey, claimId: claim.claimId, + fenceTags: entry.cacheTags, keyHash: entry.keyHash, objectKey: claim.objectKey, revision: claim.revision, @@ -612,6 +624,7 @@ export class ResponseStoreBinding extends WorkerEntrypoint< const metadata = this.getMetadata(); const { cacheKey, keyHash } = await this.deriveCacheKey(request); + const cacheTags = cacheTagsFromResponse(response); const pendingPutKey = `${this.getVersionId()}:${keyHash}`; let reservation: WriteReservation | undefined; if (options.coalesce) { @@ -619,7 +632,7 @@ export class ResponseStoreBinding extends WorkerEntrypoint< const pending = pendingPuts.get(pendingPutKey); if (!pending) break; - reservation ??= await this.reserveWrite(metadata, keyHash, cacheKey); + reservation ??= await this.reserveWrite(metadata, keyHash, cacheKey, cacheTags); try { const result = await pending; if (result.backingStoreUpdated) { @@ -637,13 +650,14 @@ export class ResponseStoreBinding extends WorkerEntrypoint< } const write = (async (): Promise => { - reservation ??= await this.reserveWrite(metadata, keyHash, cacheKey); + reservation ??= await this.reserveWrite(metadata, keyHash, cacheKey, cacheTags); const result = await this.storeResponse( metadata, request, response, options.revalidator, reservation, + cacheTags, ); if (!result.published || !result.entry) { return { backingStoreUpdated: false, edgePurgeAccepted: false }; @@ -687,6 +701,7 @@ export class ResponseStoreBinding extends WorkerEntrypoint< reservation ? { cacheKey: entry.cacheKey, + fenceTags: entry.cacheTags, keyHash: entry.keyHash, ...reservation, } diff --git a/packages/workers-response-store/src/metadata-do.ts b/packages/workers-response-store/src/metadata-do.ts index 55fca8436..697acd0ec 100644 --- a/packages/workers-response-store/src/metadata-do.ts +++ b/packages/workers-response-store/src/metadata-do.ts @@ -31,7 +31,6 @@ type RefreshCandidate = { }; export type CacheMetadataStub = DurableObjectStub & { - beginWrite(keyHash: string, cacheKey: string): Promise; reserveWrite( keyHash: string, cacheKey: string, @@ -75,6 +74,7 @@ type EntryRow = Record & { active_revision: number | null; latest_revision: number; object_key: string | null; + pending_created_at?: number | null; status_text: string | null; response_headers: string | null; fresh_until: number | null; @@ -214,7 +214,8 @@ export class CacheMetadata extends DurableObject { .exec( includePending ? `SELECT * FROM entries - WHERE (tombstoned = 0 AND active_revision IS NOT NULL) OR active_revision IS NULL` + 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(); @@ -466,18 +467,6 @@ export class CacheMetadata extends DurableObject { return revision; } - beginWrite(keyHash: string, cacheKey: string): number { - return this.ctx.storage.transactionSync(() => { - const current = this.ctx.storage.sql - .exec<{ latest_revision: number }>( - "SELECT latest_revision FROM entries WHERE key_hash = ?", - keyHash, - ) - .toArray()[0]; - return this.reserveRevision(keyHash, cacheKey, current); - }); - } - async reserveWrite( keyHash: string, cacheKey: string, @@ -513,12 +502,23 @@ export class CacheMetadata extends DurableObject { ): Promise { const { cleanupObjectKey, result } = this.ctx.storage.transactionSync(() => { const current = this.ctx.storage.sql - .exec("SELECT * FROM entries WHERE key_hash = ?", keyHash) + .exec( + `SELECT entries.*, pending_objects.created_at AS pending_created_at + FROM entries + LEFT JOIN pending_objects ON pending_objects.object_key = ? + WHERE entries.key_hash = ?`, + metadata.objectKey, + keyHash, + ) .toArray()[0]; if ( !current || + current.pending_created_at === null || + current.pending_created_at === undefined || revision > current.latest_revision || - (current.active_revision !== null && revision <= current.active_revision) + (current.active_revision !== null && revision <= current.active_revision) || + (metadata.fenceTags.length > 0 && + this.getTagExpiration(metadata.fenceTags) >= current.pending_created_at) ) { if (claimId) { this.ctx.storage.sql.exec( diff --git a/packages/workers-response-store/tests/e2e.test.mjs b/packages/workers-response-store/tests/e2e.test.mjs index 624d45b01..bea07c343 100644 --- a/packages/workers-response-store/tests/e2e.test.mjs +++ b/packages/workers-response-store/tests/e2e.test.mjs @@ -637,6 +637,35 @@ test("purge prevents an initial slow write from creating an entry", async () => 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("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"); From 95ffadee85d2b1bf68792fbfacf736d2f9b131fc Mon Sep 17 00:00:00 2001 From: James Date: Sun, 13 Sep 2026 23:49:30 +0100 Subject: [PATCH 07/17] fix(cache): order tag purge fences durably --- .../workers-response-store/src/metadata-do.ts | 76 ++++++++++++++----- .../workers-response-store/tests/e2e.test.mjs | 25 ++++++ 2 files changed, 84 insertions(+), 17 deletions(-) diff --git a/packages/workers-response-store/src/metadata-do.ts b/packages/workers-response-store/src/metadata-do.ts index 697acd0ec..737583d6c 100644 --- a/packages/workers-response-store/src/metadata-do.ts +++ b/packages/workers-response-store/src/metadata-do.ts @@ -74,7 +74,8 @@ type EntryRow = Record & { active_revision: number | null; latest_revision: number; object_key: string | null; - pending_created_at?: number | null; + current_invalidation_sequence?: number; + pending_invalidation_sequence?: number | null; status_text: string | null; response_headers: string | null; fresh_until: number | null; @@ -185,10 +186,17 @@ export class CacheMetadata extends DurableObject { ); 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 ) WITHOUT ROWID; + 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, created_at INTEGER NOT NULL ); CREATE INDEX IF NOT EXISTS pending_objects_created_at ON pending_objects(created_at); @@ -430,7 +438,9 @@ export class CacheMetadata extends DurableObject { now + leaseMs, ); this.ctx.storage.sql.exec( - "INSERT OR REPLACE INTO pending_objects (object_key, created_at) VALUES (?, ?)", + `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, ); @@ -484,7 +494,9 @@ export class CacheMetadata extends DurableObject { 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) VALUES (?, ?)", + `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, ); @@ -503,8 +515,11 @@ export class CacheMetadata extends DurableObject { const { cleanupObjectKey, result } = this.ctx.storage.transactionSync(() => { const current = this.ctx.storage.sql .exec( - `SELECT entries.*, pending_objects.created_at AS pending_created_at + `SELECT entries.*, + pending_objects.invalidation_sequence AS pending_invalidation_sequence, + metadata_state.tag_invalidation_sequence AS current_invalidation_sequence FROM entries + CROSS JOIN metadata_state LEFT JOIN pending_objects ON pending_objects.object_key = ? WHERE entries.key_hash = ?`, metadata.objectKey, @@ -513,12 +528,14 @@ export class CacheMetadata extends DurableObject { .toArray()[0]; if ( !current || - current.pending_created_at === null || - current.pending_created_at === undefined || + current.pending_invalidation_sequence === null || + current.pending_invalidation_sequence === undefined || revision > current.latest_revision || (current.active_revision !== null && revision <= current.active_revision) || (metadata.fenceTags.length > 0 && - this.getTagExpiration(metadata.fenceTags) >= current.pending_created_at) + 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( @@ -618,7 +635,10 @@ 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; @@ -626,7 +646,7 @@ export class CacheMetadata extends DurableObject { 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, ) @@ -637,6 +657,10 @@ export class CacheMetadata extends DurableObject { return expiration; } + getTagExpiration(tags: string[]): number { + return this.getTagInvalidationMaximum(tags, "invalidated_at"); + } + async reserveRefresh( options: ResponseStoreRefreshOptions, objectKeyRoot: string, @@ -667,9 +691,13 @@ export class CacheMetadata extends DurableObject { } 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) VALUES ${batch - .map(() => "(?, ?)") - .join(", ")}`, + `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]), ); } @@ -707,13 +735,27 @@ export class CacheMetadata extends DurableObject { const matches = this.findMatchingEntryRows(options, true); 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 ${batch - .map(() => "(?, ?)") - .join(", ")} + `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)`, + MAX(tag_invalidations.invalidated_at, excluded.invalidated_at), + invalidation_sequence = MAX( + tag_invalidations.invalidation_sequence, + excluded.invalidation_sequence + )`, ...batch.flatMap((tag) => [tag, invalidatedAt]), ); } diff --git a/packages/workers-response-store/tests/e2e.test.mjs b/packages/workers-response-store/tests/e2e.test.mjs index bea07c343..d7659225a 100644 --- a/packages/workers-response-store/tests/e2e.test.mjs +++ b/packages/workers-response-store/tests/e2e.test.mjs @@ -666,6 +666,31 @@ test("tag purge prevents a pending tagged write from publishing", async () => { assert.equal((await read("/purge-pending-tag")).status, 404); }); +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("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"); From c09351bbed0a0c17ebd31d1ae57f4f6fb9c3582a Mon Sep 17 00:00:00 2001 From: James Date: Sun, 13 Sep 2026 23:58:55 +0100 Subject: [PATCH 08/17] fix(cache): migrate response store metadata schema --- .../workers-response-store/src/metadata-do.ts | 40 ++++++- .../workers-response-store/tests/e2e.test.mjs | 106 ++++++++++++++++++ 2 files changed, 145 insertions(+), 1 deletion(-) diff --git a/packages/workers-response-store/src/metadata-do.ts b/packages/workers-response-store/src/metadata-do.ts index 737583d6c..8c0fba906 100644 --- a/packages/workers-response-store/src/metadata-do.ts +++ b/packages/workers-response-store/src/metadata-do.ts @@ -187,8 +187,11 @@ export class CacheMetadata extends DurableObject { CREATE TABLE IF NOT EXISTS tag_invalidations ( tag TEXT PRIMARY KEY, invalidated_at INTEGER NOT NULL, - invalidation_sequence 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 @@ -201,6 +204,41 @@ export class CacheMetadata extends DurableObject { ); CREATE INDEX IF NOT EXISTS pending_objects_created_at ON pending_objects(created_at); `); + + ctx.storage.transactionSync(() => { + const migrated = ctx.storage.sql + .exec<{ version: number }>( + "SELECT version FROM metadata_schema_migrations WHERE version = 2", + ) + .toArray().length; + if (migrated) return; + + 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 ( + !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 ( + !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", + ); + } + ctx.storage.sql.exec("INSERT INTO metadata_schema_migrations (version) VALUES (2)"); + }); }); } diff --git a/packages/workers-response-store/tests/e2e.test.mjs b/packages/workers-response-store/tests/e2e.test.mjs index d7659225a..33f6e1ed5 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"; @@ -691,6 +694,109 @@ test("a write reserved after a tag purge is not rejected by its timestamp", asyn 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"); From d8c6587c2449572c1c1888fe549bf182a8c6b3b8 Mon Sep 17 00:00:00 2001 From: James Date: Mon, 14 Sep 2026 00:06:11 +0100 Subject: [PATCH 09/17] fix(cache): preserve coalesced purge semantics --- .../workers-response-store/src/binding.ts | 2 +- .../workers-response-store/tests/e2e.test.mjs | 20 +++++++++++++++++++ 2 files changed, 21 insertions(+), 1 deletion(-) diff --git a/packages/workers-response-store/src/binding.ts b/packages/workers-response-store/src/binding.ts index d3c540cd6..c7aa7c6bd 100644 --- a/packages/workers-response-store/src/binding.ts +++ b/packages/workers-response-store/src/binding.ts @@ -625,7 +625,7 @@ export class ResponseStoreBinding extends WorkerEntrypoint< const { cacheKey, keyHash } = await this.deriveCacheKey(request); const cacheTags = cacheTagsFromResponse(response); - const pendingPutKey = `${this.getVersionId()}:${keyHash}`; + const pendingPutKey = `${this.getVersionId()}:${keyHash}:${Boolean(options.purgeExisting)}`; let reservation: WriteReservation | undefined; if (options.coalesce) { for (;;) { diff --git a/packages/workers-response-store/tests/e2e.test.mjs b/packages/workers-response-store/tests/e2e.test.mjs index 33f6e1ed5..010d1b359 100644 --- a/packages/workers-response-store/tests/e2e.test.mjs +++ b/packages/workers-response-store/tests/e2e.test.mjs @@ -561,6 +561,26 @@ test("overlapping framework writes can be coalesced", async () => { assert.equal(await metadataRowCount("pending_objects"), 0); }); +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", { From 3958d51aeb2c347d419190f777d90789c44c8023 Mon Sep 17 00:00:00 2001 From: James Date: Mon, 14 Sep 2026 00:13:39 +0100 Subject: [PATCH 10/17] fix(cache): isolate coalesced cleanup failures --- packages/workers-response-store/src/binding.ts | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/packages/workers-response-store/src/binding.ts b/packages/workers-response-store/src/binding.ts index c7aa7c6bd..1d59a5fc9 100644 --- a/packages/workers-response-store/src/binding.ts +++ b/packages/workers-response-store/src/binding.ts @@ -636,7 +636,10 @@ export class ResponseStoreBinding extends WorkerEntrypoint< try { const result = await pending; if (result.backingStoreUpdated) { - await metadata.finishPendingObjects([reservation.objectKey]); + const objectKey = reservation.objectKey; + await metadata + .finishPendingObjects([objectKey]) + .catch((error) => this.logCleanupFailure(objectKey, error)); await response.body?.cancel().catch(() => {}); return result; } From e528112c782ab1a3848edff915f335c6cbabbb55 Mon Sep 17 00:00:00 2001 From: James Date: Mon, 14 Sep 2026 00:22:48 +0100 Subject: [PATCH 11/17] fix(cache): harden coalesced response store writes --- .../workers-response-store/example/worker.ts | 5 ++ .../workers-response-store/src/binding.ts | 41 +++++++------- .../workers-response-store/src/metadata-do.ts | 7 +++ .../workers-response-store/tests/e2e.test.mjs | 54 +++++++++++++++++++ 4 files changed, 88 insertions(+), 19 deletions(-) diff --git a/packages/workers-response-store/example/worker.ts b/packages/workers-response-store/example/worker.ts index 9d43c7add..ba93929be 100644 --- a/packages/workers-response-store/example/worker.ts +++ b/packages/workers-response-store/example/worker.ts @@ -89,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, @@ -110,6 +114,7 @@ async function handlePut(request: Request, store: WorkersResponseStore): Promise revalidator, purgeExisting: request.headers.get("X-Purge-Existing") === "1", }); + await new Response(teeSibling).arrayBuffer(); return json(result); } diff --git a/packages/workers-response-store/src/binding.ts b/packages/workers-response-store/src/binding.ts index 1d59a5fc9..3aefdb136 100644 --- a/packages/workers-response-store/src/binding.ts +++ b/packages/workers-response-store/src/binding.ts @@ -208,7 +208,7 @@ 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>(); +const pendingPuts = new Map>(); function* batches(values: readonly T[], size: number): Generator { for (let offset = 0; offset < values.length; offset += size) { @@ -635,13 +635,18 @@ export class ResponseStoreBinding extends WorkerEntrypoint< reservation ??= await this.reserveWrite(metadata, keyHash, cacheKey, cacheTags); try { const result = await pending; - if (result.backingStoreUpdated) { + if (result.published && result.entry) { const objectKey = reservation.objectKey; await metadata .finishPendingObjects([objectKey]) .catch((error) => this.logCleanupFailure(objectKey, error)); - await response.body?.cancel().catch(() => {}); - return result; + void response.body?.cancel().catch(() => {}); + return { + backingStoreUpdated: true, + edgePurgeAccepted: options.purgeExisting + ? await this.purgeEdgeCacheByTags([purgeTagForEntry(result.entry)]) + : true, + }; } } catch { // Preserve this response as the fallback when the leading write fails. @@ -652,9 +657,9 @@ export class ResponseStoreBinding extends WorkerEntrypoint< } } - const write = (async (): Promise => { + const write = (async (): Promise => { reservation ??= await this.reserveWrite(metadata, keyHash, cacheKey, cacheTags); - const result = await this.storeResponse( + return this.storeResponse( metadata, request, response, @@ -662,23 +667,21 @@ export class ResponseStoreBinding extends WorkerEntrypoint< reservation, cacheTags, ); + })(); + if (options.coalesce) pendingPuts.set(pendingPutKey, write); + try { + const result = await write; if (!result.published || !result.entry) { return { backingStoreUpdated: false, edgePurgeAccepted: false }; } - - const edgePurgeAccepted = options.purgeExisting - ? await this.purgeEdgeCacheByTags([purgeTagForEntry(result.entry)]) - : true; - - return { backingStoreUpdated: true, edgePurgeAccepted }; - })(); - if (!options.coalesce) return write; - - pendingPuts.set(pendingPutKey, write); - try { - return await write; + return { + backingStoreUpdated: true, + edgePurgeAccepted: options.purgeExisting + ? await this.purgeEdgeCacheByTags([purgeTagForEntry(result.entry)]) + : true, + }; } finally { - if (pendingPuts.get(pendingPutKey) === write) { + if (options.coalesce && pendingPuts.get(pendingPutKey) === write) { pendingPuts.delete(pendingPutKey); } } diff --git a/packages/workers-response-store/src/metadata-do.ts b/packages/workers-response-store/src/metadata-do.ts index 8c0fba906..413b533e1 100644 --- a/packages/workers-response-store/src/metadata-do.ts +++ b/packages/workers-response-store/src/metadata-do.ts @@ -74,6 +74,8 @@ type EntryRow = Record & { active_revision: number | null; latest_revision: number; object_key: string | null; + claim_id?: string | null; + claim_revision?: number | null; current_invalidation_sequence?: number; pending_invalidation_sequence?: number | null; status_text: string | null; @@ -554,11 +556,14 @@ export class CacheMetadata extends DurableObject { const current = this.ctx.storage.sql .exec( `SELECT entries.*, + revalidation_claims.claim_id AS claim_id, + revalidation_claims.revision AS claim_revision, pending_objects.invalidation_sequence AS pending_invalidation_sequence, 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, @@ -570,6 +575,8 @@ export class CacheMetadata extends DurableObject { current.pending_invalidation_sequence === undefined || revision > current.latest_revision || (current.active_revision !== null && revision <= current.active_revision) || + (claimId !== undefined && + (current.claim_id !== claimId || current.claim_revision !== revision)) || (metadata.fenceTags.length > 0 && current.current_invalidation_sequence! > current.pending_invalidation_sequence && this.getTagInvalidationMaximum(metadata.fenceTags, "invalidation_sequence") > diff --git a/packages/workers-response-store/tests/e2e.test.mjs b/packages/workers-response-store/tests/e2e.test.mjs index 010d1b359..8020a4d0c 100644 --- a/packages/workers-response-store/tests/e2e.test.mjs +++ b/packages/workers-response-store/tests/e2e.test.mjs @@ -71,6 +71,7 @@ async function put(path, body, options = {}) { 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, @@ -561,6 +562,18 @@ test("overlapping framework writes can be coalesced", async () => { 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)); @@ -689,6 +702,47 @@ test("tag purge prevents a pending tagged write from publishing", async () => { 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 write reserved after a tag purge is not rejected by its timestamp", async () => { const createdAt = Date.now(); await purge({ tags: ["already-purged"] }); From 9a5fbea828fd6098454d9f53773dd81dfa76da68 Mon Sep 17 00:00:00 2001 From: James Date: Mon, 14 Sep 2026 00:26:41 +0100 Subject: [PATCH 12/17] fix(cache): isolate coalesced purge failures --- .../workers-response-store/src/binding.ts | 33 +++++++++++-------- .../workers-response-store/tests/e2e.test.mjs | 6 +++- 2 files changed, 24 insertions(+), 15 deletions(-) diff --git a/packages/workers-response-store/src/binding.ts b/packages/workers-response-store/src/binding.ts index 3aefdb136..22059f498 100644 --- a/packages/workers-response-store/src/binding.ts +++ b/packages/workers-response-store/src/binding.ts @@ -633,23 +633,28 @@ export class ResponseStoreBinding extends WorkerEntrypoint< if (!pending) break; reservation ??= await this.reserveWrite(metadata, keyHash, cacheKey, cacheTags); + let result: StoreResult; try { - const result = await pending; - 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, - }; - } + 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); diff --git a/packages/workers-response-store/tests/e2e.test.mjs b/packages/workers-response-store/tests/e2e.test.mjs index 8020a4d0c..eceddb5c8 100644 --- a/packages/workers-response-store/tests/e2e.test.mjs +++ b/packages/workers-response-store/tests/e2e.test.mjs @@ -662,7 +662,11 @@ test("purge prevents coalesced writes from resurrecting an entry", async () => { test("purge prevents an initial slow write from creating an entry", async () => { const write = put("/purge-cold-write", "too-late", { bodyDelayMs: 300 }); - await new Promise((resolve) => setTimeout(resolve, 50)); + 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, { From d8f4afa27ca4a22f831f62474f6ea1a4e5c450e1 Mon Sep 17 00:00:00 2001 From: James Date: Mon, 14 Sep 2026 00:34:03 +0100 Subject: [PATCH 13/17] fix(cache): isolate post-commit alarm failures --- packages/workers-response-store/src/metadata-do.ts | 9 ++++++++- 1 file changed, 8 insertions(+), 1 deletion(-) diff --git a/packages/workers-response-store/src/metadata-do.ts b/packages/workers-response-store/src/metadata-do.ts index 413b533e1..3d3e8f786 100644 --- a/packages/workers-response-store/src/metadata-do.ts +++ b/packages/workers-response-store/src/metadata-do.ts @@ -667,7 +667,14 @@ export class CacheMetadata extends DurableObject { }); if (cleanupObjectKey) { - await this.ensureCleanupAlarm(Date.now()); + await this.ensureCleanupAlarm(Date.now()).catch((error) => { + console.error( + JSON.stringify({ + message: "Workers Response Store cleanup alarm update failed", + error: error instanceof Error ? error.message : String(error), + }), + ); + }); await this.deleteTrackedObjects([cleanupObjectKey]); } return result; From ff288723fc02ba657f07ea9afd31a041cc7f4acc Mon Sep 17 00:00:00 2001 From: James Date: Mon, 14 Sep 2026 00:41:18 +0100 Subject: [PATCH 14/17] fix(cache): continue cleanup after alarm errors --- .../workers-response-store/src/metadata-do.ts | 22 +++++++++++-------- 1 file changed, 13 insertions(+), 9 deletions(-) diff --git a/packages/workers-response-store/src/metadata-do.ts b/packages/workers-response-store/src/metadata-do.ts index 3d3e8f786..118b7b4b1 100644 --- a/packages/workers-response-store/src/metadata-do.ts +++ b/packages/workers-response-store/src/metadata-do.ts @@ -254,6 +254,17 @@ export class CacheMetadata extends DurableObject { this.cleanupAlarmKnown = true; } + 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), + }), + ); + }); + } + private findMatchingEntryRows( options: ResponseStorePurgeOptions, includePending = false, @@ -667,14 +678,7 @@ export class CacheMetadata extends DurableObject { }); if (cleanupObjectKey) { - await this.ensureCleanupAlarm(Date.now()).catch((error) => { - console.error( - JSON.stringify({ - message: "Workers Response Store cleanup alarm update failed", - error: error instanceof Error ? error.message : String(error), - }), - ); - }); + await this.maintainCleanupAlarm(Date.now()); await this.deleteTrackedObjects([cleanupObjectKey]); } return result; @@ -851,7 +855,7 @@ export class CacheMetadata extends DurableObject { ); }); if (matches.length) { - await this.ensureCleanupAlarm(invalidatedAt); + await this.maintainCleanupAlarm(invalidatedAt); await this.deleteTrackedObjects(matches.map((entry) => entry.objectKey)); } return matches; From 58975ddcd0783bb1262871f5d423fa2234ace6a4 Mon Sep 17 00:00:00 2001 From: James Date: Mon, 14 Sep 2026 00:47:15 +0100 Subject: [PATCH 15/17] fix(cache): fence expired response store cleanup --- .../workers-response-store/src/metadata-do.ts | 56 ++++++++++++++++--- .../workers-response-store/tests/e2e.test.mjs | 51 +++++++++++++++++ 2 files changed, 98 insertions(+), 9 deletions(-) diff --git a/packages/workers-response-store/src/metadata-do.ts b/packages/workers-response-store/src/metadata-do.ts index 118b7b4b1..89796d552 100644 --- a/packages/workers-response-store/src/metadata-do.ts +++ b/packages/workers-response-store/src/metadata-do.ts @@ -74,6 +74,7 @@ type EntryRow = Record & { active_revision: number | null; latest_revision: number; object_key: string | null; + claim_active_revision?: number | null; claim_id?: string | null; claim_revision?: number | null; current_invalidation_sequence?: number; @@ -359,6 +360,17 @@ export class CacheMetadata extends DurableObject { 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), + }), + ); + } } } } @@ -382,8 +394,14 @@ export class CacheMetadata extends DurableObject { 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 }>( + .exec<{ + active: number; + created_at: number; + invalidation_sequence: number; + object_key: string; + }>( `SELECT pending_objects.object_key, pending_objects.created_at, + pending_objects.invalidation_sequence, entries.object_key IS NOT NULL AS active FROM pending_objects LEFT JOIN entries @@ -397,13 +415,30 @@ export class CacheMetadata extends DurableObject { const activeObjectKeys = batch .filter(({ active }) => active) .map(({ object_key }) => object_key); - const expiredObjectKeys = batch - .filter(({ active, created_at }) => !active && created_at <= cutoff) - .map(({ object_key }) => object_key); - this.finishPendingObjects(activeObjectKeys); - if (expiredObjectKeys.length) { - await this.env.CACHE_BODIES.delete(expiredObjectKeys); - this.finishPendingObjects(expiredObjectKeys); + const expired = batch.filter(({ active, created_at }) => !active && created_at <= cutoff); + const expiredObjectKeys = expired.map(({ object_key }) => object_key); + this.finishPendingObjects([...activeObjectKeys, ...expiredObjectKeys]); + try { + if (expiredObjectKeys.length) { + await this.env.CACHE_BODIES.delete(expiredObjectKeys); + } + } catch (error) { + this.ctx.storage.transactionSync(() => { + for (const batch of batches(expired, Math.floor(MAX_SQL_PARAMETERS / 3))) { + this.ctx.storage.sql.exec( + `INSERT OR REPLACE INTO pending_objects + (object_key, created_at, invalidation_sequence) VALUES ${batch + .map(() => "(?, ?, ?)") + .join(", ")}`, + ...batch.flatMap(({ object_key, invalidation_sequence }) => [ + object_key, + 0, + invalidation_sequence, + ]), + ); + } + }); + throw error; } const next = batch.find(({ active, created_at }) => !active && created_at > cutoff); @@ -567,6 +602,7 @@ export class CacheMetadata extends DurableObject { const current = this.ctx.storage.sql .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, @@ -587,7 +623,9 @@ export class CacheMetadata extends DurableObject { revision > current.latest_revision || (current.active_revision !== null && revision <= current.active_revision) || (claimId !== undefined && - (current.claim_id !== claimId || current.claim_revision !== revision)) || + (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") > diff --git a/packages/workers-response-store/tests/e2e.test.mjs b/packages/workers-response-store/tests/e2e.test.mjs index eceddb5c8..9dcdf78b2 100644 --- a/packages/workers-response-store/tests/e2e.test.mjs +++ b/packages/workers-response-store/tests/e2e.test.mjs @@ -747,6 +747,57 @@ test("an expired revalidation claim cannot publish after its replacement", async 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"] }); From b9b4ca6296fdf950b6d091fc59c753981455e8ad Mon Sep 17 00:00:00 2001 From: James Date: Mon, 14 Sep 2026 00:54:27 +0100 Subject: [PATCH 16/17] fix(cache): retain durable cleanup ownership --- .../workers-response-store/src/metadata-do.ts | 95 +++++++++++-------- 1 file changed, 57 insertions(+), 38 deletions(-) diff --git a/packages/workers-response-store/src/metadata-do.ts b/packages/workers-response-store/src/metadata-do.ts index 89796d552..911f2ed3c 100644 --- a/packages/workers-response-store/src/metadata-do.ts +++ b/packages/workers-response-store/src/metadata-do.ts @@ -79,6 +79,7 @@ type EntryRow = Record & { 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; fresh_until: number | null; @@ -203,18 +204,22 @@ export class CacheMetadata extends DurableObject { 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); `); ctx.storage.transactionSync(() => { - const migrated = ctx.storage.sql - .exec<{ version: number }>( - "SELECT version FROM metadata_schema_migrations WHERE version = 2", - ) - .toArray().length; - if (migrated) return; + const migrations = new Set( + ctx.storage.sql + .exec<{ version: number }>( + "SELECT version FROM metadata_schema_migrations WHERE version IN (2, 3)", + ) + .toArray() + .map(({ version }) => version), + ); + if (migrations.size === 2) return; const schemas = ctx.storage.sql .exec<{ name: string; sql: string }>( @@ -223,6 +228,7 @@ export class CacheMetadata extends DurableObject { ) .toArray(); if ( + !migrations.has(2) && !schemas .find(({ name }) => name === "tag_invalidations") ?.sql.includes("invalidation_sequence") @@ -232,6 +238,7 @@ export class CacheMetadata extends DurableObject { ); } if ( + !migrations.has(2) && !schemas .find(({ name }) => name === "pending_objects") ?.sql.includes("invalidation_sequence") @@ -240,7 +247,19 @@ export class CacheMetadata extends DurableObject { "ALTER TABLE pending_objects ADD COLUMN invalidation_sequence INTEGER NOT NULL DEFAULT 0", ); } - ctx.storage.sql.exec("INSERT INTO metadata_schema_migrations (version) VALUES (2)"); + 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)"); + } }); }); } @@ -301,9 +320,10 @@ export class CacheMetadata extends DurableObject { 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 ${batch - .map(() => "(?, ?)") - .join(", ")}`, + `INSERT OR REPLACE INTO pending_objects + (object_key, created_at, publishable) VALUES ${batch + .map(() => "(?, ?, 0)") + .join(", ")}`, ...values, ); } @@ -319,7 +339,7 @@ export class CacheMetadata extends DurableObject { // 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 WHERE object_key = ?", + "UPDATE pending_objects SET created_at = 0, publishable = 0 WHERE object_key = ?", objectKey, ); if (claimId) { @@ -353,6 +373,15 @@ export class CacheMetadata extends DurableObject { 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", @@ -397,11 +426,9 @@ export class CacheMetadata extends DurableObject { .exec<{ active: number; created_at: number; - invalidation_sequence: number; object_key: string; }>( `SELECT pending_objects.object_key, pending_objects.created_at, - pending_objects.invalidation_sequence, entries.object_key IS NOT NULL AS active FROM pending_objects LEFT JOIN entries @@ -417,28 +444,17 @@ export class CacheMetadata extends DurableObject { .map(({ object_key }) => object_key); const expired = batch.filter(({ active, created_at }) => !active && created_at <= cutoff); const expiredObjectKeys = expired.map(({ object_key }) => object_key); - this.finishPendingObjects([...activeObjectKeys, ...expiredObjectKeys]); - try { - if (expiredObjectKeys.length) { - await this.env.CACHE_BODIES.delete(expiredObjectKeys); - } - } catch (error) { - this.ctx.storage.transactionSync(() => { - for (const batch of batches(expired, Math.floor(MAX_SQL_PARAMETERS / 3))) { - this.ctx.storage.sql.exec( - `INSERT OR REPLACE INTO pending_objects - (object_key, created_at, invalidation_sequence) VALUES ${batch - .map(() => "(?, ?, ?)") - .join(", ")}`, - ...batch.flatMap(({ object_key, invalidation_sequence }) => [ - object_key, - 0, - invalidation_sequence, - ]), - ); - } - }); - throw error; + this.finishPendingObjects(activeObjectKeys); + for (const expiredBatch of batches(expiredObjectKeys, MAX_SQL_PARAMETERS)) { + this.ctx.storage.sql.exec( + `UPDATE pending_objects SET created_at = 0, publishable = 0 + WHERE object_key IN (${expiredBatch.map(() => "?").join(", ")})`, + ...expiredBatch, + ); + } + if (expiredObjectKeys.length) { + await this.env.CACHE_BODIES.delete(expiredObjectKeys); + this.finishPendingObjects(expiredObjectKeys); } const next = batch.find(({ active, created_at }) => !active && created_at > cutoff); @@ -606,6 +622,7 @@ export class CacheMetadata extends DurableObject { 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 @@ -620,6 +637,7 @@ export class CacheMetadata extends DurableObject { !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 && @@ -675,7 +693,8 @@ export class CacheMetadata extends DurableObject { } 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) VALUES (?, ?)", + `INSERT OR IGNORE INTO pending_objects + (object_key, created_at, publishable) VALUES (?, ?, 0)`, current.object_key, Date.now(), ); @@ -858,8 +877,8 @@ export class CacheMetadata extends DurableObject { 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) - SELECT object_key, ? FROM entries + `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, From 70dae5cd342a81df8671d16aa2522242b8eb4df6 Mon Sep 17 00:00:00 2001 From: James Date: Mon, 14 Sep 2026 00:59:18 +0100 Subject: [PATCH 17/17] fix(cache): retain cleanup fences across deletion --- .../workers-response-store/src/metadata-do.ts | 24 ++++++++++----- .../workers-response-store/tests/e2e.test.mjs | 30 +++++++++++++++++++ 2 files changed, 47 insertions(+), 7 deletions(-) diff --git a/packages/workers-response-store/src/metadata-do.ts b/packages/workers-response-store/src/metadata-do.ts index 911f2ed3c..e3ea77641 100644 --- a/packages/workers-response-store/src/metadata-do.ts +++ b/packages/workers-response-store/src/metadata-do.ts @@ -427,8 +427,10 @@ export class CacheMetadata extends DurableObject { 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 @@ -443,26 +445,34 @@ export class CacheMetadata extends DurableObject { .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 expiredBatch of batches(expiredObjectKeys, MAX_SQL_PARAMETERS)) { + for (const reservationBatch of batches(reservations, MAX_SQL_PARAMETERS - 1)) { this.ctx.storage.sql.exec( - `UPDATE pending_objects SET created_at = 0, publishable = 0 - WHERE object_key IN (${expiredBatch.map(() => "?").join(", ")})`, - ...expiredBatch, + `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(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) { - await this.ensureCleanupAlarm(next.created_at); + } else if (next || reservations.length) { + await this.ensureCleanupAlarm( + Math.min(next?.created_at ?? Number.POSITIVE_INFINITY, fencedAt), + ); } return expiredObjectKeys.length; diff --git a/packages/workers-response-store/tests/e2e.test.mjs b/packages/workers-response-store/tests/e2e.test.mjs index 9dcdf78b2..1d87bd065 100644 --- a/packages/workers-response-store/tests/e2e.test.mjs +++ b/packages/workers-response-store/tests/e2e.test.mjs @@ -948,6 +948,36 @@ test("retention sweep removes orphaned candidates without deleting active R2 obj 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();