From da95225ba78dd17b409c017776d8bf5b2c09ba01 Mon Sep 17 00:00:00 2001 From: Neil <4138956+nwparker@users.noreply.github.com> Date: Sat, 26 Sep 2026 20:39:57 -0700 Subject: [PATCH] perf(push): keep retention sweeps from overlapping without slowing the drain (#23001) * perf(push): keep retention sweeps from overlapping * perf(push): drain a saturated retention sweep instead of idling out the tick The overlap guard on the shared prune timer removed a side effect the sweeper had been relying on: overlap was the only thing that let a backlog exceed the 50-batch-per-call cap inside one 60-second tick. With the guard, a sweep that spent its whole budget went idle for the rest of the interval, so a large backlog drained far slower exactly when retention matters most. The timer is now a chained setTimeout rather than an interval. `deleteInBatches` reports whether it exhausted its batch budget, `prune()` returns `{ deleted, saturated }`, and a saturated sweep is rescheduled immediately. The connection gate hands a freed slot to the longest waiter, so one serial sweeper looping back to back still parks a single statement ahead of a worker claim: claim latency keeps the value the guard bought while the maximum drain rate returns to what it was before. The loop is self-limiting and stops once the backlog clears. A sweep that has not settled a full interval after it started now logs `orca_push_prune_overdue` with its target. Admission waits have no timeout, so a lost slot release could previously wedge retention permanently and silently. Chaining makes the overlap guard structural, so there is no flag to scope. The two single-DELETE sweeps state through `unbatchedSweep` that they have no batch budget to exhaust, which keeps the immediate-resume path readable as delivery-only. Co-Authored-By: Claude --------- Co-authored-by: Claude --- .../apps/push/src/durable-push-claim.test.ts | 10 +- cloud/apps/push/src/durable-push-store.ts | 19 +- cloud/apps/push/src/push-background.test.ts | 174 ++++++++++++++++++ cloud/apps/push/src/push-background.ts | 73 ++++++-- 4 files changed, 246 insertions(+), 30 deletions(-) create mode 100644 cloud/apps/push/src/push-background.test.ts diff --git a/cloud/apps/push/src/durable-push-claim.test.ts b/cloud/apps/push/src/durable-push-claim.test.ts index f43588bd2c3..c7a228a0559 100644 --- a/cloud/apps/push/src/durable-push-claim.test.ts +++ b/cloud/apps/push/src/durable-push-claim.test.ts @@ -157,7 +157,7 @@ it('keeps claim, exclusion and prune correct without the queue index', async () for (const claim of leased) await store.finish(claim) expect((await store.claim())?.notification.notificationSeq).toBe(2) advance(10 * 60_000) - expect(await store.prune()).toBe(1) + expect(await store.prune()).toEqual({ deleted: 1, saturated: false }) expect(await batchCount(db)).toBe(0) }) @@ -211,11 +211,11 @@ it('prunes a large backlog in bounded calls without touching live or leased work [JSON.stringify(notification(9)), now - 1, now - 1, now + 1000, now - 1] ) await store.accept('host', 'phone-live', notification(1)) - expect(await store.prune()).toBe(perCall) - expect(await store.prune()).toBe(500) - expect(await store.prune()).toBe(0) + expect(await store.prune()).toEqual({ deleted: perCall, saturated: true }) + expect(await store.prune()).toEqual({ deleted: 500, saturated: false }) + expect(await store.prune()).toEqual({ deleted: 0, saturated: false }) advance(1000) - expect(await store.prune()).toBe(1) + expect(await store.prune()).toEqual({ deleted: 1, saturated: false }) expect(await batchCount(db)).toBe(1) expect(await store.pendingCount('phone-live')).toBe(1) }) diff --git a/cloud/apps/push/src/durable-push-store.ts b/cloud/apps/push/src/durable-push-store.ts index 61fed235018..02d26f20b3e 100644 --- a/cloud/apps/push/src/durable-push-store.ts +++ b/cloud/apps/push/src/durable-push-store.ts @@ -12,6 +12,8 @@ const CLAIM_CANDIDATE_ATTEMPTS = 4 export const PRUNE_BATCH_ROWS = 2_000 export const PRUNE_MAX_BATCHES = 50 export const DELIVERY_LEASE_MS = 30_000 +// `saturated` means the batch budget ran out with rows still matching, so a backlog remains. +export type PushPruneSweep = { deleted: number; saturated: boolean } export type QueuedPushDelivery = { id: string registrationId: string @@ -200,10 +202,10 @@ export class DurablePushStore { } // Bounded per call and per statement, so it drains any backlog on its own without holding locks. - async prune(): Promise { + async prune(): Promise { const now = this.now() // Also clears terminal rows older revisions kept, since each carries a past expires_at. - let deleted = await this.deleteInBatches( + let { deleted, saturated } = await this.deleteInBatches( 'push_delivery_batches', 'batch_id', 'expires_at <= ? AND lease_until <= ?', @@ -215,9 +217,11 @@ export class DurablePushStore { ['push_event_recipients', 'event_id, registration_id'], ['push_events', 'event_id'] ] as const) { - deleted += await this.deleteInBatches(table, key, 'created_at < ?', [now - RETENTION_MS]) + const sweep = await this.deleteInBatches(table, key, 'created_at < ?', [now - RETENTION_MS]) + deleted += sweep.deleted + saturated ||= sweep.saturated } - return deleted + return { deleted, saturated } } private async deleteInBatches( @@ -225,7 +229,7 @@ export class DurablePushStore { key: string, where: string, params: unknown[] - ): Promise { + ): Promise { const lockRows = this.background.dialect === 'postgres' ? ' FOR UPDATE SKIP LOCKED' : '' let total = 0 for (let batch = 0; batch < PRUNE_MAX_BATCHES; batch++) { @@ -235,8 +239,9 @@ export class DurablePushStore { ) const changes = Number(result?.changes ?? 0) total += changes - if (changes < PRUNE_BATCH_ROWS) break + // A short batch drained the predicate; only a full last batch leaves rows behind. + if (changes < PRUNE_BATCH_ROWS) return { deleted: total, saturated: false } } - return total + return { deleted: total, saturated: true } } } diff --git a/cloud/apps/push/src/push-background.test.ts b/cloud/apps/push/src/push-background.test.ts new file mode 100644 index 00000000000..21598ef67a2 --- /dev/null +++ b/cloud/apps/push/src/push-background.test.ts @@ -0,0 +1,174 @@ +import { afterEach, expect, it, vi } from 'vitest' +import { startPushBackground } from './push-background.js' +import { createPushServerHarness } from './push-server-harness.test-fixture.js' +import { reserveRequestConnection } from './push-background-database.js' +import { DurablePushStore, PRUNE_BATCH_ROWS, type PushPruneSweep } from './durable-push-store.js' +import type { PushDatabase } from './push-database.js' + +const cleanups: (() => Promise)[] = [] +afterEach(async () => { + for (const cleanup of cleanups.splice(0)) await cleanup() + vi.useRealTimers() + vi.restoreAllMocks() +}) + +async function fixture(mode: 'active' | 'validation' = 'active') { + const harness = await createPushServerHarness() + const runtime = harness.server + const challenges = vi.spyOn(runtime.challenges, 'pruneExpired').mockResolvedValue(0) + const sessions = vi.spyOn(runtime.sessions, 'pruneExpired').mockResolvedValue(0) + const deliveries = vi + .spyOn(runtime.deliveryStore, 'prune') + .mockResolvedValue({ deleted: 0, saturated: false }) + vi.spyOn(runtime.worker, 'start').mockImplementation(() => {}) + const warn = vi.spyOn(console, 'warn').mockImplementation(() => {}) + vi.useFakeTimers() + const stop = startPushBackground({ mode }, runtime) + cleanups.push(async () => { + await stop() + await harness.close() + }) + return { stop, challenges, sessions, deliveries, warn } +} + +it('keeps one slow sweep per store while other stores keep their cadence', async () => { + const h = await fixture() + let finish!: (sweep: PushPruneSweep) => void + h.deliveries.mockImplementationOnce( + () => new Promise((resolve) => (finish = resolve)) + ) + try { + await vi.advanceTimersByTimeAsync(10 * 60_000) + expect(h.deliveries).toHaveBeenCalledTimes(1) + expect(h.challenges).toHaveBeenCalledTimes(10) + expect(h.sessions).toHaveBeenCalledTimes(1) + } finally { + finish({ deleted: 100_000, saturated: false }) + } + await vi.advanceTimersByTimeAsync(60_000) + expect(h.deliveries).toHaveBeenCalledTimes(2) + await h.stop() + await vi.advanceTimersByTimeAsync(10 * 60_000) + expect(h.deliveries).toHaveBeenCalledTimes(2) + expect(h.challenges).toHaveBeenCalledTimes(11) + expect(h.sessions).toHaveBeenCalledTimes(1) +}) + +it('reports a sweep that is still running a full interval after it started', async () => { + const h = await fixture() + let finish!: (sweep: PushPruneSweep) => void + h.deliveries.mockImplementationOnce( + () => new Promise((resolve) => (finish = resolve)) + ) + try { + await vi.advanceTimersByTimeAsync(119_999) + expect(h.warn).not.toHaveBeenCalled() + await vi.advanceTimersByTimeAsync(1) + expect(h.warn).toHaveBeenCalledWith( + JSON.stringify({ event: 'orca_push_prune_overdue', target: 'deliveries' }) + ) + } finally { + finish({ deleted: 0, saturated: false }) + } + // One report per sweep: the settled sweep clears its own watchdog. + await vi.advanceTimersByTimeAsync(10 * 60_000) + expect(h.warn).toHaveBeenCalledTimes(1) +}) + +it('releases a failed sweep so the next scheduled sweep can recover', async () => { + const h = await fixture() + h.deliveries.mockRejectedValueOnce(new Error('database unavailable')) + await vi.advanceTimersByTimeAsync(60_000) + expect(h.deliveries).toHaveBeenCalledTimes(1) + expect(h.warn).toHaveBeenCalledWith( + JSON.stringify({ event: 'orca_push_prune_failed', target: 'deliveries', error: 'Error' }) + ) + await vi.advanceTimersByTimeAsync(60_000) + expect(h.deliveries).toHaveBeenCalledTimes(2) + expect(h.warn).toHaveBeenCalledTimes(1) +}) + +it('resumes a saturated sweep at once rather than waiting out the interval', async () => { + const h = await fixture() + let backlogSweeps = 3 + h.deliveries.mockImplementation(async () => { + const saturated = backlogSweeps-- > 0 + return { deleted: saturated ? PRUNE_BATCH_ROWS : 0, saturated } + }) + await vi.advanceTimersByTimeAsync(60_000) + expect(h.deliveries).toHaveBeenCalledTimes(1) + // Three saturated sweeps resume within milliseconds instead of costing an interval each. + await vi.advanceTimersByTimeAsync(10) + expect(h.deliveries).toHaveBeenCalledTimes(4) + await vi.advanceTimersByTimeAsync(59_000) + expect(h.deliveries).toHaveBeenCalledTimes(4) + await vi.advanceTimersByTimeAsync(1_000) + expect(h.deliveries).toHaveBeenCalledTimes(5) +}) + +it('keeps a delivery claim behind one statement while a backlog drains back to back', async () => { + const h = await fixture() + let backlog = true + let finishedDeletes = 0 + const database: PushDatabase = { + dialect: 'postgres', + query: async (sql) => { + if (!sql.startsWith('DELETE')) return [] + await new Promise((resolve) => setTimeout(resolve, 4_000)) + finishedDeletes++ + // Only deliveries hold a backlog, so each sweep spends its whole budget there and returns. + if (!backlog || !sql.includes('push_delivery_batches')) return [{ changes: 0 }] + return [{ changes: PRUNE_BATCH_ROWS }] + }, + transaction: (operation) => operation(database), + lockQuotaScope: async () => {}, + tryLockScope: async () => true, + tryLockSharedScope: async () => true, + close: async () => {} + } + const store = new DurablePushStore(database, Date.now, reserveRequestConnection(database, 2)) + const sweeps: Promise[] = [] + let inFlight = 0 + let concurrentSweeps = 0 + h.deliveries.mockImplementation(() => { + concurrentSweeps = Math.max(concurrentSweeps, ++inFlight) + const sweep = store.prune().finally(() => void inFlight--) + sweeps.push(sweep) + return sweep + }) + try { + await vi.advanceTimersByTimeAsync(10 * 60_000) + const queuedAt = finishedDeletes + const startedAt = Date.now() + let statementsAhead: number | undefined + let claimDelay: number | undefined + const claim = store.claim().then(() => { + statementsAhead = finishedDeletes - queuedAt + claimDelay = Date.now() - startedAt + }) + await vi.advanceTimersByTimeAsync(44_000) + await claim + // Sweeping serially parks one statement ahead of the claim no matter how long the drain runs. + expect({ statementsAhead, concurrentSweeps }).toEqual({ + statementsAhead: 1, + concurrentSweeps: 1 + }) + expect(claimDelay).toBeLessThanOrEqual(4_000) + // Each sweep exhausts its 50-batch budget, so the drain continues instead of idling out the tick. + expect(sweeps.length).toBeGreaterThan(1) + } finally { + await h.stop() + backlog = false + await vi.advanceTimersByTimeAsync(10 * 60_000) + await Promise.all(sweeps) + } +}) + +it('keeps validation mode free of sweeps and timers', async () => { + const h = await fixture('validation') + await vi.advanceTimersByTimeAsync(20 * 60_000) + expect(h.challenges).not.toHaveBeenCalled() + expect(h.sessions).not.toHaveBeenCalled() + expect(h.deliveries).not.toHaveBeenCalled() + expect(vi.getTimerCount()).toBe(0) +}) diff --git a/cloud/apps/push/src/push-background.ts b/cloud/apps/push/src/push-background.ts index 8925fd41fed..9fad4700e7b 100644 --- a/cloud/apps/push/src/push-background.ts +++ b/cloud/apps/push/src/push-background.ts @@ -1,26 +1,55 @@ import type { PushConfig } from './config.js' +import type { PushPruneSweep } from './durable-push-store.js' import type { createPushServer } from './push-server.js' const CHALLENGE_PRUNE_INTERVAL_MS = 60_000 const SESSION_PRUNE_INTERVAL_MS = 10 * 60_000 const DELIVERY_PRUNE_INTERVAL_MS = 60_000 -function prune(label: string, run: () => Promise, intervalMs: number): NodeJS.Timeout { - const timer = setInterval(() => { - void run().catch((error: unknown) => { - console.warn( - JSON.stringify({ - event: 'orca_push_prune_failed', - target: label, - error: error instanceof Error ? error.name : 'unknown' - }) - ) - }) - }, intervalMs) - timer.unref() - return timer +// Chained rather than periodic, so a sweep spanning many bounded statements never overlaps itself and +// one that exhausted its batch budget resumes at once instead of idling out the rest of the interval. +function prune(label: string, run: () => Promise, intervalMs: number): () => void { + let timer: NodeJS.Timeout | undefined + let overdue: NodeJS.Timeout | undefined + let stopped = false + function schedule(delayMs: number): void { + if (stopped) return + timer = setTimeout(tick, delayMs) + timer.unref() + } + function tick(): void { + // Admission waits are untimed, so a lost slot release would otherwise stall retention silently. + overdue = setTimeout(() => { + console.warn(JSON.stringify({ event: 'orca_push_prune_overdue', target: label })) + }, intervalMs) + overdue.unref() + void run() + .then((sweep) => schedule(sweep.saturated ? 0 : intervalMs)) + .catch((error: unknown) => { + console.warn( + JSON.stringify({ + event: 'orca_push_prune_failed', + target: label, + error: error instanceof Error ? error.name : 'unknown' + }) + ) + schedule(intervalMs) + }) + .finally(() => clearTimeout(overdue)) + } + schedule(intervalMs) + return () => { + stopped = true + clearTimeout(timer) + clearTimeout(overdue) + } } +// One DELETE under a statement timeout: there is no batch budget for it to exhaust. +const unbatchedSweep = + (run: () => Promise) => + async (): Promise => ({ deleted: await run(), saturated: false }) + export function startPushBackground( config: Pick, runtime: Pick< @@ -30,14 +59,22 @@ export function startPushBackground( ): () => Promise { if (config.mode === 'validation') return async () => {} const { challenges, sessions, deliveryStore, worker } = runtime - const timers = [ - prune('challenges', () => challenges.pruneExpired(), CHALLENGE_PRUNE_INTERVAL_MS), - prune('sessions', () => sessions.pruneExpired(), SESSION_PRUNE_INTERVAL_MS), + const stops = [ + prune( + 'challenges', + unbatchedSweep(() => challenges.pruneExpired()), + CHALLENGE_PRUNE_INTERVAL_MS + ), + prune( + 'sessions', + unbatchedSweep(() => sessions.pruneExpired()), + SESSION_PRUNE_INTERVAL_MS + ), prune('deliveries', () => deliveryStore.prune(), DELIVERY_PRUNE_INTERVAL_MS) ] worker.start() return async () => { - for (const timer of timers) clearInterval(timer) + for (const stop of stops) stop() await worker.stop() } }