mirror of
https://github.com/stablyai/orca.git
synced 2026-10-06 00:02:43 +00:00
* fix(push): bound the delivery claim and stop keeping finished batches * fix(push): bound claim scans to the notification TTL and document the queue * test(push): boot waits out a table writer for the queue indexes; queue stays correct without them * fix(push): share the claim lock so a previous-revision claim cannot re-lease a delivery During a deploy overlap the previous revision's claim scans under an exclusive push-worker-claim lock and re-reads the row without checking lease_until, so it could overwrite a lease this revision had just committed and send twice. The new claim now takes the same key shared: new claimers never block each other, and the previous claim waits until their leases commit before it scans. Droppable one release after every worker runs this revision. * fix(push): install one queue index at boot, not three push_batches_leased_device indexed lease_until, so every lease, renew and finish UPDATE lost heap-only eligibility and rewrote every index. push_batches_pending_due was unused: on a synthetic 902k-row table the candidate scan plans onto the existing (state, due_at) and expiry indexes with or without it. The per-device pending index stays; the head check and the busy anti-join use it. Fewer boot-time builds also shorten the SHARE lock the first boot takes on the table. * test(push): pin the claim's TTL scan bound and the server's worker connection cap Removing either guard left the suite green. The claim test captures every row the candidate scan returns and plants one row that only the TTL term excludes; the server test drives the real worker through createPushServer and fails when the request-connection reservation is unwired (peak 4 instead of 2). * fix(push): renew delivery leases outside the background connection cap Renew shared the single background slot with claim retries and prune batches, so a heartbeat could wait long enough for a lease to lapse and the delivery to be re-leased mid-send. It is a keyed one-row UPDATE, so request traffic cannot starve it on the ungated pool. * docs(push): describe the shared claim lock for mixed-revision deploys
235 lines
9.6 KiB
TypeScript
235 lines
9.6 KiB
TypeScript
import { afterEach, describe, expect, it } from 'vitest'
|
|
import { openPushDatabase } from './push-database.js'
|
|
import { DurablePushStore, DELIVERY_LEASE_MS } from './durable-push-store.js'
|
|
import {
|
|
cleanupDurablePushFixtures,
|
|
durablePushTestDatabaseUrl,
|
|
fixture,
|
|
notification
|
|
} from './durable-push-store.test-fixture.js'
|
|
|
|
afterEach(cleanupDurablePushFixtures)
|
|
|
|
describe('durable push acceptance', () => {
|
|
it('counts a logical event once across phones and separates the 300/15min dismissal budget', async () => {
|
|
const { store, advance } = await fixture()
|
|
for (let i = 0; i < 300; i++) {
|
|
expect(await store.accept('host', 'phone1', notification(i))).toBe('queued')
|
|
expect(await store.accept('host', 'phone2', notification(i))).toBe('queued')
|
|
expect(await store.accept('host', 'phone1', notification(i, 'dismiss'))).toBe('queued')
|
|
}
|
|
expect(await store.accept('host', 'phone1', notification(300))).toBe('rate_limited')
|
|
expect(await store.accept('host', 'phone1', notification(300, 'dismiss'))).toBe('rate_limited')
|
|
expect(await store.accept('another-host', 'phone3', notification(300))).toBe('queued')
|
|
advance(15 * 60_000)
|
|
expect(await store.accept('host', 'phone1', notification(301))).toBe('queued')
|
|
})
|
|
|
|
it('queues one delivery per event and recovers work across service instances', async () => {
|
|
const { db, store, clock, advance } = await fixture()
|
|
await store.accept('host', 'phone', notification(1))
|
|
const restarted = new DurablePushStore(db, clock)
|
|
await restarted.accept('host', 'phone', notification(1))
|
|
advance(1)
|
|
await restarted.accept('host', 'phone', notification(2))
|
|
const rows = await db.query(
|
|
"SELECT payload_json, due_at, created_at FROM push_delivery_batches WHERE registration_id = ? AND state = 'pending' ORDER BY created_at, batch_id",
|
|
['phone']
|
|
)
|
|
expect(rows).toHaveLength(2)
|
|
expect(rows.map((row) => JSON.parse(String(row.payload_json)))).toEqual([
|
|
notification(1),
|
|
notification(2)
|
|
])
|
|
expect(rows.every((row) => Number(row.due_at) >= Number(row.created_at))).toBe(true)
|
|
const delivery = await restarted.claim()
|
|
expect(delivery?.notification.notificationSeq).toBe(1)
|
|
expect(await store.claim()).toBeNull()
|
|
advance(DELIVERY_LEASE_MS)
|
|
const reclaimed = await store.claim()
|
|
expect(reclaimed?.id).toBe(delivery?.id)
|
|
expect(reclaimed?.lease).not.toBe(delivery?.lease)
|
|
await restarted.finish(delivery!)
|
|
expect(await store.claim()).toBeNull()
|
|
await store.finish(reclaimed!)
|
|
const second = await restarted.claim()
|
|
expect(second?.notification.notificationSeq).toBe(2)
|
|
await restarted.finish(second!)
|
|
expect(await restarted.claim()).toBeNull()
|
|
})
|
|
|
|
it('never extends expiry and refuses conflicting duplicate content', async () => {
|
|
const { store, advance } = await fixture()
|
|
await store.accept('host', 'phone', notification(1))
|
|
expect(await store.accept('host', 'phone', { ...notification(1), body: 'changed' })).toBe(
|
|
'error'
|
|
)
|
|
const delivery = (await store.claim())!
|
|
await store.finish(delivery, 10 * 60_000)
|
|
advance(60_000)
|
|
expect(await store.claim()).toBeNull()
|
|
advance(5 * 60_000)
|
|
expect(await store.accept('host', 'phone', notification(1))).toBe('error')
|
|
})
|
|
|
|
it('orders a due retry before a fresh first attempt without delaying the retry', async () => {
|
|
const { store, advance } = await fixture()
|
|
await store.accept('host', 'phone', notification(1))
|
|
const first = (await store.claim())!
|
|
await store.finish(first, 1000)
|
|
expect(await store.claim()).toBeNull()
|
|
|
|
advance(1000)
|
|
await store.accept('host', 'phone', notification(2))
|
|
const retry = (await store.claim())!
|
|
expect(retry.notification.notificationSeq).toBe(1)
|
|
await store.finish(retry)
|
|
const fresh = (await store.claim())!
|
|
expect(fresh?.notification.notificationSeq).toBe(2)
|
|
await store.finish(fresh!)
|
|
})
|
|
|
|
it('orders an expired first-attempt lease by creation time after a retry becomes due', async () => {
|
|
const { db, store, clock, advance } = await fixture()
|
|
await store.accept('host', 'phone', notification(1))
|
|
const retry = (await store.claim())!
|
|
await store.finish(retry, 1000)
|
|
|
|
advance(2000)
|
|
await db.query(
|
|
`INSERT INTO push_delivery_batches(batch_id, host_fingerprint, registration_id, kind, payload_json, state, due_at, expires_at, lease_until, attempts, created_at)
|
|
VALUES ('crashed-singleton', 'host', 'phone', 'alert', ?, 'pending', ?, ?, 0, 1, ?)`,
|
|
[JSON.stringify(notification(2)), clock() - 1, clock() + 300_000, clock()]
|
|
)
|
|
const reclaimedRetry = (await store.claim())!
|
|
expect(reclaimedRetry.notification.notificationSeq).toBe(1)
|
|
await store.finish(reclaimedRetry)
|
|
const reclaimedCrash = (await store.claim())!
|
|
expect(reclaimedCrash.notification.notificationSeq).toBe(2)
|
|
await store.finish(reclaimedCrash)
|
|
})
|
|
|
|
it('rolls quota and payload back together if persistence fails', async () => {
|
|
const { db, store } = await fixture()
|
|
await db.query('ALTER TABLE push_delivery_batches RENAME TO push_delivery_batches_unavailable')
|
|
try {
|
|
const databaseUrl = durablePushTestDatabaseUrl
|
|
if (databaseUrl) {
|
|
const concurrent = await openPushDatabase({ databaseUrl, dataDir: '' })
|
|
try {
|
|
await expect(
|
|
concurrent.query('SELECT COUNT(*) FROM push_delivery_batches')
|
|
).resolves.toHaveLength(1)
|
|
} finally {
|
|
await concurrent.close()
|
|
}
|
|
}
|
|
await expect(store.accept('host', 'phone', notification(1))).rejects.toThrow()
|
|
expect(await db.query('SELECT * FROM push_events')).toEqual([])
|
|
expect(await db.query('SELECT * FROM push_event_recipients')).toEqual([])
|
|
} finally {
|
|
await db.query(
|
|
'ALTER TABLE push_delivery_batches_unavailable RENAME TO push_delivery_batches'
|
|
)
|
|
}
|
|
})
|
|
})
|
|
|
|
it('serializes concurrent instances at the quota boundary', async () => {
|
|
const { db, store, clock } = await fixture()
|
|
for (let seq = 0; seq < 299; seq++) await store.accept('host', 'phone', notification(seq))
|
|
const second = new DurablePushStore(db, clock)
|
|
const results = await Promise.all(
|
|
Array.from({ length: 6 }, (_, index) =>
|
|
(index % 2 ? store : second).accept('host', 'phone', notification(400 + index))
|
|
)
|
|
)
|
|
expect(results.filter((result) => result === 'queued')).toHaveLength(1)
|
|
expect(results.filter((result) => result === 'rate_limited')).toHaveLength(5)
|
|
})
|
|
|
|
it('cancels unsent alerts and prevents an older replay after dismissal', async () => {
|
|
const { store } = await fixture()
|
|
const alert = notification(1)
|
|
await store.accept('host', 'phone', alert)
|
|
await store.accept('host', 'phone', {
|
|
...notification(2, 'dismiss'),
|
|
notificationId: alert.notificationId
|
|
})
|
|
const delivery = (await store.claim())!
|
|
expect(delivery.notification.kind).toBe('dismiss')
|
|
await store.finish(delivery)
|
|
expect(await store.claim()).toBeNull()
|
|
await store.accept('host', 'another-phone', alert)
|
|
expect(await store.claim()).toBeNull()
|
|
})
|
|
|
|
it('does not resurrect an in-flight alert after a dismissal and transient provider failure', async () => {
|
|
const { store, advance } = await fixture()
|
|
await store.accept('host', 'phone', notification(1))
|
|
const inFlight = (await store.claim())!
|
|
await store.accept('host', 'phone', {
|
|
...notification(2, 'dismiss'),
|
|
notificationId: notification(1).notificationId
|
|
})
|
|
await store.finish(inFlight, 1000)
|
|
const dismissal = (await store.claim())!
|
|
expect(dismissal.notification.kind).toBe('dismiss')
|
|
await store.finish(dismissal)
|
|
advance(1000)
|
|
expect(await store.claim()).toBeNull()
|
|
expect(await store.pendingCount('phone')).toBe(0)
|
|
})
|
|
|
|
it.each([false, true])(
|
|
'normalizes default alert kind (explicit first: %s)',
|
|
async (explicitFirst) => {
|
|
const { db, store } = await fixture()
|
|
const { kind: _kind, ...implicit } = notification(1)
|
|
const explicit = { kind: 'alert' as const, ...implicit }
|
|
for (const event of explicitFirst ? [explicit, implicit] : [implicit, explicit]) {
|
|
expect(await store.accept('host', 'phone', event)).toBe('queued')
|
|
}
|
|
expect(await store.pendingCount('phone')).toBe(1)
|
|
expect(await db.query('SELECT event_id FROM push_events')).toHaveLength(1)
|
|
expect(await store.accept('host', 'phone', { ...explicit, body: 'changed' })).toBe('error')
|
|
expect(await store.accept('host', 'phone', { ...implicit, kind: 'dismiss' })).toBe('queued')
|
|
expect(await db.query('SELECT event_id FROM push_events')).toHaveLength(2)
|
|
}
|
|
)
|
|
|
|
it('fences late renew and finish after an expired claim is dismissed', async () => {
|
|
const { db, store, advance } = await fixture()
|
|
const alert = notification(1)
|
|
await store.accept('host', 'phone', alert)
|
|
const stale = (await store.claim())!
|
|
advance(DELIVERY_LEASE_MS)
|
|
await store.accept('host', 'phone', {
|
|
...notification(2, 'dismiss'),
|
|
notificationId: alert.notificationId
|
|
})
|
|
const read = async () =>
|
|
(
|
|
await db.query(
|
|
'SELECT state, payload_json, lease_until FROM push_delivery_batches WHERE batch_id = ?',
|
|
[stale.id]
|
|
)
|
|
)[0]
|
|
expect(await read()).toBeUndefined()
|
|
await store.renew(stale)
|
|
await store.finish(stale, 1000)
|
|
await store.finish(stale)
|
|
expect(await read()).toBeUndefined()
|
|
const dismissal = (await store.claim())!
|
|
expect(dismissal.notification.kind).toBe('dismiss')
|
|
await store.finish(dismissal)
|
|
await store.accept('host', 'phone', notification(3))
|
|
const fresh = (await store.claim())!
|
|
await store.finish(fresh, 1000)
|
|
advance(1000)
|
|
const retry = (await store.claim())!
|
|
expect(retry.id).toBe(fresh.id)
|
|
await store.finish(retry)
|
|
expect(await store.claim()).toBeNull()
|
|
})
|