mirror of
https://github.com/stablyai/orca.git
synced 2026-10-01 16:02:03 +00:00
* 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>
81 lines
2.6 KiB
TypeScript
81 lines
2.6 KiB
TypeScript
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
|
|
|
|
// 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<
|
|
ReturnType<typeof createPushServer>,
|
|
'challenges' | 'sessions' | 'deliveryStore' | 'worker'
|
|
>
|
|
): () => Promise<void> {
|
|
if (config.mode === 'validation') return async () => {}
|
|
const { challenges, sessions, deliveryStore, worker } = runtime
|
|
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 stop of stops) stop()
|
|
await worker.stop()
|
|
}
|
|
}
|