Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
59 changes: 48 additions & 11 deletions packages/cloudflare/src/publisher-object.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand All @@ -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)
}
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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')

Expand Down Expand Up @@ -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)
Expand All @@ -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)
Expand All @@ -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)
})
})
13 changes: 13 additions & 0 deletions packages/cloudflare/src/publisher-object.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
}
}
Expand Down
4 changes: 2 additions & 2 deletions packages/cloudflare/src/publisher.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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()
Expand Down
Loading