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()