From 4bd6e5b1ef69cde80fe8cfab5b7079182d891ebd Mon Sep 17 00:00:00 2001 From: Claude Date: Tue, 6 Oct 2026 02:09:56 +0000 Subject: [PATCH 1/4] fix(cloudflare): keep resume ids increasing after the events table is recreated The idle alarm's deleteAll() and resetSchema's DROP TABLE reset the AUTOINCREMENT sequence, so new events got ids starting at 1 again. Replay uses `WHERE id > ?`, so a subscriber resuming with an id issued before the reset silently skipped every newer event stored since. Seed the sequence with the current time in microseconds whenever the table is created, so ids issued after a reset stay above earlier ones. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_017T4x44m5ATzLkRn6JNgLr2 --- .../cloudflare/src/publisher-object.test.ts | 59 +++++++++++++++---- packages/cloudflare/src/publisher-object.ts | 13 ++++ packages/cloudflare/src/publisher.test.ts | 4 +- 3 files changed, 63 insertions(+), 13 deletions(-) diff --git a/packages/cloudflare/src/publisher-object.test.ts b/packages/cloudflare/src/publisher-object.test.ts index bebb6ba86..1dc094138 100644 --- a/packages/cloudflare/src/publisher-object.test.ts +++ b/packages/cloudflare/src/publisher-object.test.ts @@ -144,22 +144,23 @@ describe('durable publisher object', () => { })).status).toBe(204) expect((await publish(stub, { data: { text: 'third' } })).status).toBe(204) - const liveMessages = await readMessages(liveSubscriber, 3) + const liveMessages = await readMessages<{ meta: { id: string } }>(liveSubscriber, 3) + const firstId = BigInt(liveMessages[0]!.meta.id) expect(liveMessages).toEqual([ - { data: { text: 'first' }, meta: { id: '1' } }, - { data: { text: 'second' }, meta: { id: '2', comments: ['keep me'] } }, - { data: { text: 'third' }, meta: { id: '3' } }, + { data: { text: 'first' }, meta: { id: String(firstId) } }, + { data: { text: 'second' }, meta: { id: String(firstId + 1n), comments: ['keep me'] } }, + { data: { text: 'third' }, meta: { id: String(firstId + 2n) } }, ]) expect(liveSubscriber.replayedEvents).toBe('0') - const resumeSubscriber = await openSocket(stub, '2') + const resumeSubscriber = await openSocket(stub, liveMessages[1]!.meta.id) const resumedMessages = await readMessages(resumeSubscriber, 1) expect(resumedMessages).toEqual([liveMessages[2]]) expect(resumeSubscriber.replayedEvents).toBe('1') - const tailSubscriber = await openSocket(stub, '3') + const tailSubscriber = await openSocket(stub, liveMessages[2]!.meta.id) await sleep(2) expect(tailSubscriber.messages).toHaveLength(0) @@ -178,6 +179,12 @@ describe('durable publisher object', () => { it('resumes messages in numeric id order', async () => { const stub = env.PUBLISHER_RESUME3S_DON.getByName(crypto.randomUUID()) + // create the events table, then restart its ids so they cross a digit boundary + await closeSocket(await openSocket(stub, '0')) + await runInDurableObject(stub, async (_, state) => { + state.storage.sql.exec(`UPDATE sqlite_sequence SET seq = 0 WHERE name = 'prefix:events'`) + }) + for (let order = 1; order <= 11; order++) { expect((await publish(stub, { data: { order } })).status).toBe(204) } @@ -326,7 +333,7 @@ describe('durable publisher object', () => { expect((await publish(stub, { data: { text: 'after-error' } })).status).toBe(204) expect((await readMessages(subscriber, 1))[0]).toEqual({ data: { text: 'after-error' }, - meta: { id: '1' }, + meta: { id: expect.any(String) }, }) await closeSocket(subscriber) @@ -401,7 +408,7 @@ describe('durable publisher object', () => { const resumeSubscriber = await openSocket(stub, '0') expect(await readMessages(resumeSubscriber, 1)).toEqual([{ data: { text: 'recovered' }, - meta: { id: '1' }, + meta: { id: expect.any(String) }, }]) expect(resumeSubscriber.replayedEvents).toBe('1') @@ -464,7 +471,7 @@ describe('durable publisher object', () => { expect((await publish(stub, { data: { text: 'after-alarm' } })).status).toBe(204) expect((await readMessages(subscriber, 1))[0]).toEqual({ data: { text: 'after-alarm' }, - meta: { id: '1' }, + meta: { id: expect.any(String) }, }) await runDurableObjectAlarm(stub) @@ -485,7 +492,7 @@ describe('durable publisher object', () => { expect((await readMessages(beforeExpirySubscriber, 1))[0]).toEqual({ data: { text: 'fresh resume event' }, - meta: { id: '1' }, + meta: { id: expect.any(String) }, }) await closeSocket(beforeExpirySubscriber) @@ -509,9 +516,39 @@ describe('durable publisher object', () => { expect((await publish(stub, { data: { text: 'after cleanup' } })).status).toBe(204) expect((await readMessages(newLiveSubscriber, 1))[0]).toEqual({ data: { text: 'after cleanup' }, - meta: { id: '1' }, + meta: { id: expect.any(String) }, }) await closeSocket(newLiveSubscriber) }) + + it('resumes events stored after idle cleanup for an id issued before it', async () => { + const stub = env.PUBLISHER_RESUME3S_DON.getByName(crypto.randomUUID()) + + const subscriber = await openSocket(stub) + for (const text of ['first', 'second', 'third']) { + expect((await publish(stub, { data: { text } })).status).toBe(204) + } + const seen = await readMessages<{ meta: { id: string } }>(subscriber, 3) + await closeSocket(subscriber) + + await runInDurableObject(stub, async (_, state) => { + state.storage.sql.exec('UPDATE "prefix:events" SET stored_at = unixepoch() - 10') + }) + await evictDurableObject(stub) + await runDurableObjectAlarm(stub) + expect(await getAlarm(stub)).toBeNull() + + for (const text of ['after cleanup 1', 'after cleanup 2']) { + expect((await publish(stub, { data: { text } })).status).toBe(204) + } + + const resumeSubscriber = await openSocket(stub, seen[2]!.meta.id) + const resumed = await readMessages<{ data: { text: string }, meta: { id: string } }>(resumeSubscriber, 2) + + expect(resumeSubscriber.replayedEvents).toBe('2') + expect(resumed.map(message => message.data.text)).toEqual(['after cleanup 1', 'after cleanup 2']) + + await closeSocket(resumeSubscriber) + }) }) diff --git a/packages/cloudflare/src/publisher-object.ts b/packages/cloudflare/src/publisher-object.ts index 764f81a21..869803301 100644 --- a/packages/cloudflare/src/publisher-object.ts +++ b/packages/cloudflare/src/publisher-object.ts @@ -296,6 +296,19 @@ class ResumeStorage { this.isInitedSchema = true if (initTableResult.rowsWritten > 0) { + /** + * Recreating the table (after the idle cleanup or a schema reset) restarts the + * AUTOINCREMENT sequence, while subscribers may still hold ids issued before it, + * and resuming with one of them would skip newer events (`WHERE id > ?`). + * Start the sequence at the current time in microseconds so new ids stay above + * every id issued before. + */ + this.ctx.storage.sql.exec( + `INSERT INTO sqlite_sequence (name, seq) VALUES (?, CAST(? AS INTEGER))`, + `${this.schemaPrefix}events`, + Date.now() * 1000, + ) + this.lastCleanupTime = Date.now() // schema just created, nothing to cleanup } } diff --git a/packages/cloudflare/src/publisher.test.ts b/packages/cloudflare/src/publisher.test.ts index 8d44405ba..d93149345 100644 --- a/packages/cloudflare/src/publisher.test.ts +++ b/packages/cloudflare/src/publisher.test.ts @@ -145,10 +145,10 @@ describe('durable publisher', () => { const second = live.mock.calls[1]![0] expect(first).toEqual({ text: 'first' }) - expect(getEventMeta(first)?.id).toBe('1') + expect(getEventMeta(first)?.id).toMatch(/^\d+$/) expect(second).toEqual({ text: 'second' }) - expect(getEventMeta(second)?.id).toBe('2') + expect(getEventMeta(second)?.id).toBe(String(BigInt(getEventMeta(first)!.id!) + 1n)) expect(getEventMeta(second)?.comments).toEqual(['keep me']) await stopLive() From e18a7f2a4849a7c07b5df605ff61bc5b35426c19 Mon Sep 17 00:00:00 2001 From: Claude Date: Tue, 6 Oct 2026 06:46:46 +0000 Subject: [PATCH 2/4] fix(cloudflare): prefix resume ids with the events table generation The previous fix seeded the id sequence with the current time, which only kept new ids above old ones while the old table averaged fewer than 1,000 events per millisecond. Give each events table a generation instead, stored in a meta table and changed whenever the table is recreated, and issue `-` ids. On resume, an id from another generation replays every stored event, however the sequences compare, so no assumption about the event rate is left. Tables created before this change keep issuing plain ids, so ids already held by subscribers keep working. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_017T4x44m5ATzLkRn6JNgLr2 --- .../cloudflare/src/publisher-object.test.ts | 102 +++++++++++++----- packages/cloudflare/src/publisher-object.ts | 71 +++++++++--- packages/cloudflare/src/publisher.test.ts | 4 +- 3 files changed, 136 insertions(+), 41 deletions(-) diff --git a/packages/cloudflare/src/publisher-object.test.ts b/packages/cloudflare/src/publisher-object.test.ts index 1dc094138..d3bd11845 100644 --- a/packages/cloudflare/src/publisher-object.test.ts +++ b/packages/cloudflare/src/publisher-object.test.ts @@ -67,6 +67,12 @@ async function getAlarm(stub: DurableObjectStub): Promise { return runInDurableObject(stub, async (_, state) => state.storage.getAlarm()) } +async function getGeneration(stub: DurableObjectStub): Promise { + return runInDurableObject(stub, async (_, state) => { + return state.storage.sql.exec(`SELECT value FROM "prefix:meta" WHERE key = 'generation'`).one().value as string + }) +} + async function closeSocket(socket: OpenSocket): Promise { socket.socket.close(1000, 'done') await sleep(0) @@ -144,23 +150,24 @@ describe('durable publisher object', () => { })).status).toBe(204) expect((await publish(stub, { data: { text: 'third' } })).status).toBe(204) - const liveMessages = await readMessages<{ meta: { id: string } }>(liveSubscriber, 3) - const firstId = BigInt(liveMessages[0]!.meta.id) + const liveMessages = await readMessages(liveSubscriber, 3) + const generation = await getGeneration(stub) + expect(generation).toMatch(/^\d+$/) expect(liveMessages).toEqual([ - { data: { text: 'first' }, meta: { id: String(firstId) } }, - { data: { text: 'second' }, meta: { id: String(firstId + 1n), comments: ['keep me'] } }, - { data: { text: 'third' }, meta: { id: String(firstId + 2n) } }, + { data: { text: 'first' }, meta: { id: `${generation}-1` } }, + { data: { text: 'second' }, meta: { id: `${generation}-2`, comments: ['keep me'] } }, + { data: { text: 'third' }, meta: { id: `${generation}-3` } }, ]) expect(liveSubscriber.replayedEvents).toBe('0') - const resumeSubscriber = await openSocket(stub, liveMessages[1]!.meta.id) + const resumeSubscriber = await openSocket(stub, `${generation}-2`) const resumedMessages = await readMessages(resumeSubscriber, 1) expect(resumedMessages).toEqual([liveMessages[2]]) expect(resumeSubscriber.replayedEvents).toBe('1') - const tailSubscriber = await openSocket(stub, liveMessages[2]!.meta.id) + const tailSubscriber = await openSocket(stub, `${generation}-3`) await sleep(2) expect(tailSubscriber.messages).toHaveLength(0) @@ -179,25 +186,20 @@ describe('durable publisher object', () => { it('resumes messages in numeric id order', async () => { const stub = env.PUBLISHER_RESUME3S_DON.getByName(crypto.randomUUID()) - // create the events table, then restart its ids so they cross a digit boundary - await closeSocket(await openSocket(stub, '0')) - await runInDurableObject(stub, async (_, state) => { - state.storage.sql.exec(`UPDATE sqlite_sequence SET seq = 0 WHERE name = 'prefix:events'`) - }) - for (let order = 1; order <= 11; order++) { expect((await publish(stub, { data: { order } })).status).toBe(204) } // sorting ids as text would replay '10' and '11' before '8' and '9' - const subscriber = await openSocket(stub, '7') + const generation = await getGeneration(stub) + const subscriber = await openSocket(stub, `${generation}-7`) const messages = await readMessages(subscriber, 4) expect(messages).toEqual([ - { data: { order: 8 }, meta: { id: '8' } }, - { data: { order: 9 }, meta: { id: '9' } }, - { data: { order: 10 }, meta: { id: '10' } }, - { data: { order: 11 }, meta: { id: '11' } }, + { data: { order: 8 }, meta: { id: `${generation}-8` } }, + { data: { order: 9 }, meta: { id: `${generation}-9` } }, + { data: { order: 10 }, meta: { id: `${generation}-10` } }, + { data: { order: 11 }, meta: { id: `${generation}-11` } }, ]) await closeSocket(subscriber) @@ -400,12 +402,17 @@ describe('durable publisher object', () => { const stub = env.PUBLISHER_RESUME3S_DON.getByName(crypto.randomUUID()) vi.spyOn(console, 'error').mockImplementation(() => {}) + const subscriber = await openSocket(stub) expect((await publish(stub, { data: { text: 'initial' } })).status).toBe(204) + const [initial] = await readMessages<{ meta: { id: string } }>(subscriber, 1) + await closeSocket(subscriber) + await runInDurableObject(stub, async (_, state) => breakTable(state.storage.sql)) expect((await publish(stub, { data: { text: 'recovered' } })).status).toBe(204) - const resumeSubscriber = await openSocket(stub, '0') + // the recreated table restarts its sequence, so the id issued before must still replay + const resumeSubscriber = await openSocket(stub, initial!.meta.id) expect(await readMessages(resumeSubscriber, 1)).toEqual([{ data: { text: 'recovered' }, meta: { id: expect.any(String) }, @@ -526,10 +533,16 @@ describe('durable publisher object', () => { const stub = env.PUBLISHER_RESUME3S_DON.getByName(crypto.randomUUID()) const subscriber = await openSocket(stub) - for (const text of ['first', 'second', 'third']) { - expect((await publish(stub, { data: { text } })).status).toBe(204) - } - const seen = await readMessages<{ meta: { id: string } }>(subscriber, 3) + expect((await publish(stub, { data: { text: 'first' } })).status).toBe(204) + + // however far the old sequence got, it must not hide events stored after the cleanup + await runInDurableObject(stub, async (_, state) => { + state.storage.sql.exec(`UPDATE sqlite_sequence SET seq = 9000000000000000000 WHERE name = 'prefix:events'`) + }) + expect((await publish(stub, { data: { text: 'second' } })).status).toBe(204) + + const seen = await readMessages<{ meta: { id: string } }>(subscriber, 2) + expect(seen[1]!.meta.id).toMatch(/^\d+-9000000000000000001$/) await closeSocket(subscriber) await runInDurableObject(stub, async (_, state) => { @@ -543,12 +556,49 @@ describe('durable publisher object', () => { expect((await publish(stub, { data: { text } })).status).toBe(204) } - const resumeSubscriber = await openSocket(stub, seen[2]!.meta.id) - const resumed = await readMessages<{ data: { text: string }, meta: { id: string } }>(resumeSubscriber, 2) + const generation = await getGeneration(stub) + expect(seen[1]!.meta.id.startsWith(`${generation}-`)).toBe(false) + + const resumeSubscriber = await openSocket(stub, seen[1]!.meta.id) expect(resumeSubscriber.replayedEvents).toBe('2') - expect(resumed.map(message => message.data.text)).toEqual(['after cleanup 1', 'after cleanup 2']) + expect(await readMessages(resumeSubscriber, 2)).toEqual([ + { data: { text: 'after cleanup 1' }, meta: { id: `${generation}-1` } }, + { data: { text: 'after cleanup 2' }, meta: { id: `${generation}-2` } }, + ]) await closeSocket(resumeSubscriber) }) + + it('keeps plain ids for a table created before generations existed', async () => { + const stub = env.PUBLISHER_RESUME3S_DON.getByName(crypto.randomUUID()) + + expect((await publish(stub, { data: { text: 'first' } })).status).toBe(204) + await runInDurableObject(stub, async (_, state) => { + state.storage.sql.exec('DROP TABLE "prefix:meta"') + }) + await evictDurableObject(stub) + + const subscriber = await openSocket(stub) + expect((await publish(stub, { data: { text: 'second' } })).status).toBe(204) + expect(await readMessages(subscriber, 1)).toEqual([{ data: { text: 'second' }, meta: { id: '2' } }]) + await closeSocket(subscriber) + + // a plain id still resumes within the same table + const resumeSubscriber = await openSocket(stub, '1') + expect(await readMessages(resumeSubscriber, 1)).toEqual([{ data: { text: 'second' }, meta: { id: '2' } }]) + await closeSocket(resumeSubscriber) + + // and replays everything once the table is recreated with a generation + vi.spyOn(console, 'error').mockImplementation(() => {}) + await runInDurableObject(stub, async (_, state) => { + state.storage.sql.exec('DROP TABLE "prefix:events"') + }) + expect((await publish(stub, { data: { text: 'recovered' } })).status).toBe(204) + + const generation = await getGeneration(stub) + const afterResetSubscriber = await openSocket(stub, '2') + expect(await readMessages(afterResetSubscriber, 1)).toEqual([{ data: { text: 'recovered' }, meta: { id: `${generation}-1` } }]) + await closeSocket(afterResetSubscriber) + }) }) diff --git a/packages/cloudflare/src/publisher-object.ts b/packages/cloudflare/src/publisher-object.ts index 869803301..d89aff16a 100644 --- a/packages/cloudflare/src/publisher-object.ts +++ b/packages/cloudflare/src/publisher-object.ts @@ -131,6 +131,17 @@ interface SerializedPayload { meta?: EventMeta } +/** + * Generation of a table created before generations existed, which keeps + * issuing plain sequence ids so ids already held by subscribers stay valid. + */ +const LEGACY_GENERATION = '0' + +/** + * Matches `-`, or a plain `` from the legacy generation. + */ +const EVENT_ID_REGEX = /^(?:(\d+)-)?(\d+)$/ + class ResumeStorage { private readonly enabled: boolean private readonly seconds: number @@ -141,6 +152,13 @@ class ResumeStorage { private isInitedAlarm = false private lastCleanupTime: number | undefined + /** + * Identifies the current events table. Its sequence restarts at 1 whenever the table is + * recreated (by the idle cleanup or a schema reset), so event ids are `-` + * to tell an id issued by an earlier table apart from one issued by the current table. + */ + private generation = LEGACY_GENERATION + constructor( private readonly ctx: DurableObjectState, options: DurablePublisherObjectResumeOptions = { enabled: false }, @@ -215,6 +233,19 @@ class ResumeStorage { this.ensureSchemaAndCleanup() + const match = EVENT_ID_REGEX.exec(lastEventId) + if (!match) { + return [] // not an id this object issues, so there is no position to resume from + } + + const [, generation = LEGACY_GENERATION, sequence] = match + + /** + * An id from another generation was issued by an earlier table, so every + * stored event is newer than it, however their sequences compare. + */ + const afterSequence = generation === this.generation ? sequence : '0' + /** * SQLite INTEGER can exceed JavaScript's safe integer range, * so we cast to TEXT for safe resume ID comparison. @@ -227,7 +258,7 @@ class ResumeStorage { FROM "${this.schemaPrefix}events" WHERE id > ? ORDER BY id ASC - `, lastEventId) + `, afterSequence) const events: string[] = [] for (const record of result.toArray()) { @@ -293,24 +324,36 @@ class ResumeStorage { CREATE INDEX IF NOT EXISTS "${this.schemaPrefix}idx_events_stored_at" ON "${this.schemaPrefix}events" (stored_at) `) - this.isInitedSchema = true + this.ctx.storage.sql.exec(` + CREATE TABLE IF NOT EXISTS "${this.schemaPrefix}meta" ( + key TEXT PRIMARY KEY, + value TEXT NOT NULL + ) + `) if (initTableResult.rowsWritten > 0) { /** - * Recreating the table (after the idle cleanup or a schema reset) restarts the - * AUTOINCREMENT sequence, while subscribers may still hold ids issued before it, - * and resuming with one of them would skip newer events (`WHERE id > ?`). - * Start the sequence at the current time in microseconds so new ids stay above - * every id issued before. + * A new generation only has to differ from the earlier ones: the time keeps it + * distinct across Durable Object restarts, and the previous generation this + * instance knows keeps it distinct even if the clock has not moved. */ - this.ctx.storage.sql.exec( - `INSERT INTO sqlite_sequence (name, seq) VALUES (?, CAST(? AS INTEGER))`, - `${this.schemaPrefix}events`, - Date.now() * 1000, - ) + this.generation = String(Math.max(Date.now(), Number(this.generation) + 1)) + + this.ctx.storage.sql.exec(` + INSERT OR REPLACE INTO "${this.schemaPrefix}meta" (key, value) VALUES ('generation', ?) + `, this.generation) this.lastCleanupTime = Date.now() // schema just created, nothing to cleanup } + else { + const generationRow = this.ctx.storage.sql.exec(` + SELECT value FROM "${this.schemaPrefix}meta" WHERE key = 'generation' + `).toArray()[0] + + this.generation = generationRow ? generationRow.value as string : LEGACY_GENERATION + } + + this.isInitedSchema = true } const now = Date.now() @@ -354,7 +397,9 @@ class ResumeStorage { return this.ctx.storage.setAlarm(Date.now() + this.cleanupIntervalSeconds * 1000) } - private attachEventId(message: SerializedPayload, id: string): SerializedPayload { + private attachEventId(message: SerializedPayload, sequence: string): SerializedPayload { + const id = this.generation === LEGACY_GENERATION ? sequence : `${this.generation}-${sequence}` + return { ...message, meta: { ...message.meta, id }, diff --git a/packages/cloudflare/src/publisher.test.ts b/packages/cloudflare/src/publisher.test.ts index d93149345..70afe7a72 100644 --- a/packages/cloudflare/src/publisher.test.ts +++ b/packages/cloudflare/src/publisher.test.ts @@ -145,10 +145,10 @@ describe('durable publisher', () => { const second = live.mock.calls[1]![0] expect(first).toEqual({ text: 'first' }) - expect(getEventMeta(first)?.id).toMatch(/^\d+$/) + expect(getEventMeta(first)?.id).toMatch(/^\d+-1$/) expect(second).toEqual({ text: 'second' }) - expect(getEventMeta(second)?.id).toBe(String(BigInt(getEventMeta(first)!.id!) + 1n)) + expect(getEventMeta(second)?.id).toMatch(/^\d+-2$/) expect(getEventMeta(second)?.comments).toEqual(['keep me']) await stopLive() From b7a396c3fc00af73f09d64d6ed5bd17b5c7555b4 Mon Sep 17 00:00:00 2001 From: Claude Date: Tue, 6 Oct 2026 06:56:45 +0000 Subject: [PATCH 3/4] refactor(cloudflare): simplify generation-prefixed resume ids - Issue `-` ids from every table, including tables created before generations existed (`0-`). Plain ids already held by subscribers still parse as generation `0`, so one id format remains. - Declare the meta table `WITHOUT ROWID`, so its primary key needs no hidden index and each write touches one B-tree instead of two. - Tighten test assertions loosened by the earlier time-seed approach and drop redundant test steps. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_017T4x44m5ATzLkRn6JNgLr2 --- .../cloudflare/src/publisher-object.test.ts | 27 ++++++++----------- packages/cloudflare/src/publisher-object.ts | 18 ++++++------- 2 files changed, 19 insertions(+), 26 deletions(-) diff --git a/packages/cloudflare/src/publisher-object.test.ts b/packages/cloudflare/src/publisher-object.test.ts index d3bd11845..d8fa45f1c 100644 --- a/packages/cloudflare/src/publisher-object.test.ts +++ b/packages/cloudflare/src/publisher-object.test.ts @@ -153,7 +153,6 @@ describe('durable publisher object', () => { const liveMessages = await readMessages(liveSubscriber, 3) const generation = await getGeneration(stub) - expect(generation).toMatch(/^\d+$/) expect(liveMessages).toEqual([ { data: { text: 'first' }, meta: { id: `${generation}-1` } }, { data: { text: 'second' }, meta: { id: `${generation}-2`, comments: ['keep me'] } }, @@ -335,7 +334,7 @@ describe('durable publisher object', () => { expect((await publish(stub, { data: { text: 'after-error' } })).status).toBe(204) expect((await readMessages(subscriber, 1))[0]).toEqual({ data: { text: 'after-error' }, - meta: { id: expect.any(String) }, + meta: { id: expect.stringMatching(/^\d+-1$/) }, }) await closeSocket(subscriber) @@ -402,20 +401,18 @@ describe('durable publisher object', () => { const stub = env.PUBLISHER_RESUME3S_DON.getByName(crypto.randomUUID()) vi.spyOn(console, 'error').mockImplementation(() => {}) - const subscriber = await openSocket(stub) expect((await publish(stub, { data: { text: 'initial' } })).status).toBe(204) - const [initial] = await readMessages<{ meta: { id: string } }>(subscriber, 1) - await closeSocket(subscriber) + const initialId = `${await getGeneration(stub)}-1` await runInDurableObject(stub, async (_, state) => breakTable(state.storage.sql)) expect((await publish(stub, { data: { text: 'recovered' } })).status).toBe(204) // the recreated table restarts its sequence, so the id issued before must still replay - const resumeSubscriber = await openSocket(stub, initial!.meta.id) + const resumeSubscriber = await openSocket(stub, initialId) expect(await readMessages(resumeSubscriber, 1)).toEqual([{ data: { text: 'recovered' }, - meta: { id: expect.any(String) }, + meta: { id: expect.stringMatching(/^\d+-1$/) }, }]) expect(resumeSubscriber.replayedEvents).toBe('1') @@ -478,7 +475,7 @@ describe('durable publisher object', () => { expect((await publish(stub, { data: { text: 'after-alarm' } })).status).toBe(204) expect((await readMessages(subscriber, 1))[0]).toEqual({ data: { text: 'after-alarm' }, - meta: { id: expect.any(String) }, + meta: { id: expect.stringMatching(/^\d+-1$/) }, }) await runDurableObjectAlarm(stub) @@ -499,7 +496,7 @@ describe('durable publisher object', () => { expect((await readMessages(beforeExpirySubscriber, 1))[0]).toEqual({ data: { text: 'fresh resume event' }, - meta: { id: expect.any(String) }, + meta: { id: expect.stringMatching(/^\d+-1$/) }, }) await closeSocket(beforeExpirySubscriber) @@ -523,7 +520,7 @@ describe('durable publisher object', () => { expect((await publish(stub, { data: { text: 'after cleanup' } })).status).toBe(204) expect((await readMessages(newLiveSubscriber, 1))[0]).toEqual({ data: { text: 'after cleanup' }, - meta: { id: expect.any(String) }, + meta: { id: expect.stringMatching(/^\d+-1$/) }, }) await closeSocket(newLiveSubscriber) @@ -557,8 +554,6 @@ describe('durable publisher object', () => { } const generation = await getGeneration(stub) - expect(seen[1]!.meta.id.startsWith(`${generation}-`)).toBe(false) - const resumeSubscriber = await openSocket(stub, seen[1]!.meta.id) expect(resumeSubscriber.replayedEvents).toBe('2') @@ -570,7 +565,7 @@ describe('durable publisher object', () => { await closeSocket(resumeSubscriber) }) - it('keeps plain ids for a table created before generations existed', async () => { + it('resumes plain ids issued by a table created before generations existed', async () => { const stub = env.PUBLISHER_RESUME3S_DON.getByName(crypto.randomUUID()) expect((await publish(stub, { data: { text: 'first' } })).status).toBe(204) @@ -581,12 +576,12 @@ describe('durable publisher object', () => { const subscriber = await openSocket(stub) expect((await publish(stub, { data: { text: 'second' } })).status).toBe(204) - expect(await readMessages(subscriber, 1)).toEqual([{ data: { text: 'second' }, meta: { id: '2' } }]) + expect(await readMessages(subscriber, 1)).toEqual([{ data: { text: 'second' }, meta: { id: '0-2' } }]) await closeSocket(subscriber) - // a plain id still resumes within the same table + // a plain id issued before generations existed still resumes within the same table const resumeSubscriber = await openSocket(stub, '1') - expect(await readMessages(resumeSubscriber, 1)).toEqual([{ data: { text: 'second' }, meta: { id: '2' } }]) + expect(await readMessages(resumeSubscriber, 1)).toEqual([{ data: { text: 'second' }, meta: { id: '0-2' } }]) await closeSocket(resumeSubscriber) // and replays everything once the table is recreated with a generation diff --git a/packages/cloudflare/src/publisher-object.ts b/packages/cloudflare/src/publisher-object.ts index d89aff16a..3348ff084 100644 --- a/packages/cloudflare/src/publisher-object.ts +++ b/packages/cloudflare/src/publisher-object.ts @@ -132,13 +132,12 @@ interface SerializedPayload { } /** - * Generation of a table created before generations existed, which keeps - * issuing plain sequence ids so ids already held by subscribers stay valid. + * Generation of a table created before generations existed, which issued plain sequence ids. */ const LEGACY_GENERATION = '0' /** - * Matches `-`, or a plain `` from the legacy generation. + * Matches `-`, or a plain `` issued by the legacy generation. */ const EVENT_ID_REGEX = /^(?:(\d+)-)?(\d+)$/ @@ -156,6 +155,7 @@ class ResumeStorage { * Identifies the current events table. Its sequence restarts at 1 whenever the table is * recreated (by the idle cleanup or a schema reset), so event ids are `-` * to tell an id issued by an earlier table apart from one issued by the current table. + * Kept across resets, since the next generation is derived from it. */ private generation = LEGACY_GENERATION @@ -328,7 +328,7 @@ class ResumeStorage { CREATE TABLE IF NOT EXISTS "${this.schemaPrefix}meta" ( key TEXT PRIMARY KEY, value TEXT NOT NULL - ) + ) WITHOUT ROWID `) if (initTableResult.rowsWritten > 0) { @@ -346,11 +346,11 @@ class ResumeStorage { this.lastCleanupTime = Date.now() // schema just created, nothing to cleanup } else { - const generationRow = this.ctx.storage.sql.exec(` + const generation = this.ctx.storage.sql.exec(` SELECT value FROM "${this.schemaPrefix}meta" WHERE key = 'generation' - `).toArray()[0] + `).toArray()[0]?.value as string | undefined - this.generation = generationRow ? generationRow.value as string : LEGACY_GENERATION + this.generation = generation ?? LEGACY_GENERATION } this.isInitedSchema = true @@ -398,11 +398,9 @@ class ResumeStorage { } private attachEventId(message: SerializedPayload, sequence: string): SerializedPayload { - const id = this.generation === LEGACY_GENERATION ? sequence : `${this.generation}-${sequence}` - return { ...message, - meta: { ...message.meta, id }, + meta: { ...message.meta, id: `${this.generation}-${sequence}` }, } } } From 94c99db3dfed573ed3a77aa7fa8cde3e75979958 Mon Sep 17 00:00:00 2001 From: Claude Date: Tue, 6 Oct 2026 08:16:57 +0000 Subject: [PATCH 4/4] revert(cloudflare): go back to time-seeded resume ids Revert the generation-prefixed ids (e18a7f2, b7a396c) and keep the time-seeded AUTOINCREMENT sequence from 4bd6e5b, so resume ids stay plain integers compared with `WHERE id > ?` and no extra storage is needed. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_017T4x44m5ATzLkRn6JNgLr2 --- .../cloudflare/src/publisher-object.test.ts | 107 +++++------------- packages/cloudflare/src/publisher-object.ts | 71 +++--------- packages/cloudflare/src/publisher.test.ts | 4 +- 3 files changed, 47 insertions(+), 135 deletions(-) diff --git a/packages/cloudflare/src/publisher-object.test.ts b/packages/cloudflare/src/publisher-object.test.ts index d8fa45f1c..1dc094138 100644 --- a/packages/cloudflare/src/publisher-object.test.ts +++ b/packages/cloudflare/src/publisher-object.test.ts @@ -67,12 +67,6 @@ async function getAlarm(stub: DurableObjectStub): Promise { return runInDurableObject(stub, async (_, state) => state.storage.getAlarm()) } -async function getGeneration(stub: DurableObjectStub): Promise { - return runInDurableObject(stub, async (_, state) => { - return state.storage.sql.exec(`SELECT value FROM "prefix:meta" WHERE key = 'generation'`).one().value as string - }) -} - async function closeSocket(socket: OpenSocket): Promise { socket.socket.close(1000, 'done') await sleep(0) @@ -150,23 +144,23 @@ describe('durable publisher object', () => { })).status).toBe(204) expect((await publish(stub, { data: { text: 'third' } })).status).toBe(204) - const liveMessages = await readMessages(liveSubscriber, 3) - const generation = await getGeneration(stub) + const liveMessages = await readMessages<{ meta: { id: string } }>(liveSubscriber, 3) + const firstId = BigInt(liveMessages[0]!.meta.id) expect(liveMessages).toEqual([ - { data: { text: 'first' }, meta: { id: `${generation}-1` } }, - { data: { text: 'second' }, meta: { id: `${generation}-2`, comments: ['keep me'] } }, - { data: { text: 'third' }, meta: { id: `${generation}-3` } }, + { data: { text: 'first' }, meta: { id: String(firstId) } }, + { data: { text: 'second' }, meta: { id: String(firstId + 1n), comments: ['keep me'] } }, + { data: { text: 'third' }, meta: { id: String(firstId + 2n) } }, ]) expect(liveSubscriber.replayedEvents).toBe('0') - const resumeSubscriber = await openSocket(stub, `${generation}-2`) + const resumeSubscriber = await openSocket(stub, liveMessages[1]!.meta.id) const resumedMessages = await readMessages(resumeSubscriber, 1) expect(resumedMessages).toEqual([liveMessages[2]]) expect(resumeSubscriber.replayedEvents).toBe('1') - const tailSubscriber = await openSocket(stub, `${generation}-3`) + const tailSubscriber = await openSocket(stub, liveMessages[2]!.meta.id) await sleep(2) expect(tailSubscriber.messages).toHaveLength(0) @@ -185,20 +179,25 @@ describe('durable publisher object', () => { it('resumes messages in numeric id order', async () => { const stub = env.PUBLISHER_RESUME3S_DON.getByName(crypto.randomUUID()) + // create the events table, then restart its ids so they cross a digit boundary + await closeSocket(await openSocket(stub, '0')) + await runInDurableObject(stub, async (_, state) => { + state.storage.sql.exec(`UPDATE sqlite_sequence SET seq = 0 WHERE name = 'prefix:events'`) + }) + for (let order = 1; order <= 11; order++) { expect((await publish(stub, { data: { order } })).status).toBe(204) } // sorting ids as text would replay '10' and '11' before '8' and '9' - const generation = await getGeneration(stub) - const subscriber = await openSocket(stub, `${generation}-7`) + const subscriber = await openSocket(stub, '7') const messages = await readMessages(subscriber, 4) expect(messages).toEqual([ - { data: { order: 8 }, meta: { id: `${generation}-8` } }, - { data: { order: 9 }, meta: { id: `${generation}-9` } }, - { data: { order: 10 }, meta: { id: `${generation}-10` } }, - { data: { order: 11 }, meta: { id: `${generation}-11` } }, + { data: { order: 8 }, meta: { id: '8' } }, + { data: { order: 9 }, meta: { id: '9' } }, + { data: { order: 10 }, meta: { id: '10' } }, + { data: { order: 11 }, meta: { id: '11' } }, ]) await closeSocket(subscriber) @@ -334,7 +333,7 @@ describe('durable publisher object', () => { expect((await publish(stub, { data: { text: 'after-error' } })).status).toBe(204) expect((await readMessages(subscriber, 1))[0]).toEqual({ data: { text: 'after-error' }, - meta: { id: expect.stringMatching(/^\d+-1$/) }, + meta: { id: expect.any(String) }, }) await closeSocket(subscriber) @@ -402,17 +401,14 @@ describe('durable publisher object', () => { vi.spyOn(console, 'error').mockImplementation(() => {}) expect((await publish(stub, { data: { text: 'initial' } })).status).toBe(204) - const initialId = `${await getGeneration(stub)}-1` - await runInDurableObject(stub, async (_, state) => breakTable(state.storage.sql)) expect((await publish(stub, { data: { text: 'recovered' } })).status).toBe(204) - // the recreated table restarts its sequence, so the id issued before must still replay - const resumeSubscriber = await openSocket(stub, initialId) + const resumeSubscriber = await openSocket(stub, '0') expect(await readMessages(resumeSubscriber, 1)).toEqual([{ data: { text: 'recovered' }, - meta: { id: expect.stringMatching(/^\d+-1$/) }, + meta: { id: expect.any(String) }, }]) expect(resumeSubscriber.replayedEvents).toBe('1') @@ -475,7 +471,7 @@ describe('durable publisher object', () => { expect((await publish(stub, { data: { text: 'after-alarm' } })).status).toBe(204) expect((await readMessages(subscriber, 1))[0]).toEqual({ data: { text: 'after-alarm' }, - meta: { id: expect.stringMatching(/^\d+-1$/) }, + meta: { id: expect.any(String) }, }) await runDurableObjectAlarm(stub) @@ -496,7 +492,7 @@ describe('durable publisher object', () => { expect((await readMessages(beforeExpirySubscriber, 1))[0]).toEqual({ data: { text: 'fresh resume event' }, - meta: { id: expect.stringMatching(/^\d+-1$/) }, + meta: { id: expect.any(String) }, }) await closeSocket(beforeExpirySubscriber) @@ -520,7 +516,7 @@ describe('durable publisher object', () => { expect((await publish(stub, { data: { text: 'after cleanup' } })).status).toBe(204) expect((await readMessages(newLiveSubscriber, 1))[0]).toEqual({ data: { text: 'after cleanup' }, - meta: { id: expect.stringMatching(/^\d+-1$/) }, + meta: { id: expect.any(String) }, }) await closeSocket(newLiveSubscriber) @@ -530,16 +526,10 @@ describe('durable publisher object', () => { const stub = env.PUBLISHER_RESUME3S_DON.getByName(crypto.randomUUID()) const subscriber = await openSocket(stub) - expect((await publish(stub, { data: { text: 'first' } })).status).toBe(204) - - // however far the old sequence got, it must not hide events stored after the cleanup - await runInDurableObject(stub, async (_, state) => { - state.storage.sql.exec(`UPDATE sqlite_sequence SET seq = 9000000000000000000 WHERE name = 'prefix:events'`) - }) - expect((await publish(stub, { data: { text: 'second' } })).status).toBe(204) - - const seen = await readMessages<{ meta: { id: string } }>(subscriber, 2) - expect(seen[1]!.meta.id).toMatch(/^\d+-9000000000000000001$/) + for (const text of ['first', 'second', 'third']) { + expect((await publish(stub, { data: { text } })).status).toBe(204) + } + const seen = await readMessages<{ meta: { id: string } }>(subscriber, 3) await closeSocket(subscriber) await runInDurableObject(stub, async (_, state) => { @@ -553,47 +543,12 @@ describe('durable publisher object', () => { expect((await publish(stub, { data: { text } })).status).toBe(204) } - const generation = await getGeneration(stub) - const resumeSubscriber = await openSocket(stub, seen[1]!.meta.id) + const resumeSubscriber = await openSocket(stub, seen[2]!.meta.id) + const resumed = await readMessages<{ data: { text: string }, meta: { id: string } }>(resumeSubscriber, 2) expect(resumeSubscriber.replayedEvents).toBe('2') - expect(await readMessages(resumeSubscriber, 2)).toEqual([ - { data: { text: 'after cleanup 1' }, meta: { id: `${generation}-1` } }, - { data: { text: 'after cleanup 2' }, meta: { id: `${generation}-2` } }, - ]) + expect(resumed.map(message => message.data.text)).toEqual(['after cleanup 1', 'after cleanup 2']) await closeSocket(resumeSubscriber) }) - - it('resumes plain ids issued by a table created before generations existed', async () => { - const stub = env.PUBLISHER_RESUME3S_DON.getByName(crypto.randomUUID()) - - expect((await publish(stub, { data: { text: 'first' } })).status).toBe(204) - await runInDurableObject(stub, async (_, state) => { - state.storage.sql.exec('DROP TABLE "prefix:meta"') - }) - await evictDurableObject(stub) - - const subscriber = await openSocket(stub) - expect((await publish(stub, { data: { text: 'second' } })).status).toBe(204) - expect(await readMessages(subscriber, 1)).toEqual([{ data: { text: 'second' }, meta: { id: '0-2' } }]) - await closeSocket(subscriber) - - // a plain id issued before generations existed still resumes within the same table - const resumeSubscriber = await openSocket(stub, '1') - expect(await readMessages(resumeSubscriber, 1)).toEqual([{ data: { text: 'second' }, meta: { id: '0-2' } }]) - await closeSocket(resumeSubscriber) - - // and replays everything once the table is recreated with a generation - vi.spyOn(console, 'error').mockImplementation(() => {}) - await runInDurableObject(stub, async (_, state) => { - state.storage.sql.exec('DROP TABLE "prefix:events"') - }) - expect((await publish(stub, { data: { text: 'recovered' } })).status).toBe(204) - - const generation = await getGeneration(stub) - const afterResetSubscriber = await openSocket(stub, '2') - expect(await readMessages(afterResetSubscriber, 1)).toEqual([{ data: { text: 'recovered' }, meta: { id: `${generation}-1` } }]) - await closeSocket(afterResetSubscriber) - }) }) diff --git a/packages/cloudflare/src/publisher-object.ts b/packages/cloudflare/src/publisher-object.ts index 3348ff084..869803301 100644 --- a/packages/cloudflare/src/publisher-object.ts +++ b/packages/cloudflare/src/publisher-object.ts @@ -131,16 +131,6 @@ interface SerializedPayload { meta?: EventMeta } -/** - * Generation of a table created before generations existed, which issued plain sequence ids. - */ -const LEGACY_GENERATION = '0' - -/** - * Matches `-`, or a plain `` issued by the legacy generation. - */ -const EVENT_ID_REGEX = /^(?:(\d+)-)?(\d+)$/ - class ResumeStorage { private readonly enabled: boolean private readonly seconds: number @@ -151,14 +141,6 @@ class ResumeStorage { private isInitedAlarm = false private lastCleanupTime: number | undefined - /** - * Identifies the current events table. Its sequence restarts at 1 whenever the table is - * recreated (by the idle cleanup or a schema reset), so event ids are `-` - * to tell an id issued by an earlier table apart from one issued by the current table. - * Kept across resets, since the next generation is derived from it. - */ - private generation = LEGACY_GENERATION - constructor( private readonly ctx: DurableObjectState, options: DurablePublisherObjectResumeOptions = { enabled: false }, @@ -233,19 +215,6 @@ class ResumeStorage { this.ensureSchemaAndCleanup() - const match = EVENT_ID_REGEX.exec(lastEventId) - if (!match) { - return [] // not an id this object issues, so there is no position to resume from - } - - const [, generation = LEGACY_GENERATION, sequence] = match - - /** - * An id from another generation was issued by an earlier table, so every - * stored event is newer than it, however their sequences compare. - */ - const afterSequence = generation === this.generation ? sequence : '0' - /** * SQLite INTEGER can exceed JavaScript's safe integer range, * so we cast to TEXT for safe resume ID comparison. @@ -258,7 +227,7 @@ class ResumeStorage { FROM "${this.schemaPrefix}events" WHERE id > ? ORDER BY id ASC - `, afterSequence) + `, lastEventId) const events: string[] = [] for (const record of result.toArray()) { @@ -324,36 +293,24 @@ class ResumeStorage { CREATE INDEX IF NOT EXISTS "${this.schemaPrefix}idx_events_stored_at" ON "${this.schemaPrefix}events" (stored_at) `) - this.ctx.storage.sql.exec(` - CREATE TABLE IF NOT EXISTS "${this.schemaPrefix}meta" ( - key TEXT PRIMARY KEY, - value TEXT NOT NULL - ) WITHOUT ROWID - `) + this.isInitedSchema = true if (initTableResult.rowsWritten > 0) { /** - * A new generation only has to differ from the earlier ones: the time keeps it - * distinct across Durable Object restarts, and the previous generation this - * instance knows keeps it distinct even if the clock has not moved. + * Recreating the table (after the idle cleanup or a schema reset) restarts the + * AUTOINCREMENT sequence, while subscribers may still hold ids issued before it, + * and resuming with one of them would skip newer events (`WHERE id > ?`). + * Start the sequence at the current time in microseconds so new ids stay above + * every id issued before. */ - this.generation = String(Math.max(Date.now(), Number(this.generation) + 1)) - - this.ctx.storage.sql.exec(` - INSERT OR REPLACE INTO "${this.schemaPrefix}meta" (key, value) VALUES ('generation', ?) - `, this.generation) + this.ctx.storage.sql.exec( + `INSERT INTO sqlite_sequence (name, seq) VALUES (?, CAST(? AS INTEGER))`, + `${this.schemaPrefix}events`, + Date.now() * 1000, + ) this.lastCleanupTime = Date.now() // schema just created, nothing to cleanup } - else { - const generation = this.ctx.storage.sql.exec(` - SELECT value FROM "${this.schemaPrefix}meta" WHERE key = 'generation' - `).toArray()[0]?.value as string | undefined - - this.generation = generation ?? LEGACY_GENERATION - } - - this.isInitedSchema = true } const now = Date.now() @@ -397,10 +354,10 @@ class ResumeStorage { return this.ctx.storage.setAlarm(Date.now() + this.cleanupIntervalSeconds * 1000) } - private attachEventId(message: SerializedPayload, sequence: string): SerializedPayload { + private attachEventId(message: SerializedPayload, id: string): SerializedPayload { return { ...message, - meta: { ...message.meta, id: `${this.generation}-${sequence}` }, + meta: { ...message.meta, id }, } } } diff --git a/packages/cloudflare/src/publisher.test.ts b/packages/cloudflare/src/publisher.test.ts index 70afe7a72..d93149345 100644 --- a/packages/cloudflare/src/publisher.test.ts +++ b/packages/cloudflare/src/publisher.test.ts @@ -145,10 +145,10 @@ describe('durable publisher', () => { const second = live.mock.calls[1]![0] expect(first).toEqual({ text: 'first' }) - expect(getEventMeta(first)?.id).toMatch(/^\d+-1$/) + expect(getEventMeta(first)?.id).toMatch(/^\d+$/) expect(second).toEqual({ text: 'second' }) - expect(getEventMeta(second)?.id).toMatch(/^\d+-2$/) + expect(getEventMeta(second)?.id).toBe(String(BigInt(getEventMeta(first)!.id!) + 1n)) expect(getEventMeta(second)?.comments).toEqual(['keep me']) await stopLive()