mirror of
https://github.com/stablyai/orca.git
synced 2026-09-28 00:02:41 +00:00
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 <noreply@anthropic.com>
---------
Co-authored-by: Claude <noreply@anthropic.com>
This commit is contained in:
@@ -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)
|
||||
})
|
||||
|
||||
@@ -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<number> {
|
||||
async prune(): Promise<PushPruneSweep> {
|
||||
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<number> {
|
||||
): Promise<PushPruneSweep> {
|
||||
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 }
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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<void>)[] = []
|
||||
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<PushPruneSweep>((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<PushPruneSweep>((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<PushPruneSweep>[] = []
|
||||
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)
|
||||
})
|
||||
@@ -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<number>, 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<PushPruneSweep>, 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<number>) =>
|
||||
async (): Promise<PushPruneSweep> => ({ deleted: await run(), saturated: false })
|
||||
|
||||
export function startPushBackground(
|
||||
config: Pick<PushConfig, 'mode'>,
|
||||
runtime: Pick<
|
||||
@@ -30,14 +59,22 @@ export function startPushBackground(
|
||||
): () => Promise<void> {
|
||||
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()
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user