Files
orca/cloud/apps/push/src/durable-push-worker.ts
T
Jinwoo Hong e8e09eed46 fix(push): size the claim-attempt budget from the drain count (#24040)
* fix(push): size the claim-attempt budget from the drain count

With twelve drains, up to eleven peers can hold device heads, so a
four-attempt claim budget can run out while claimable rows remain and the
drain exits idle for a tick. Move the drain count into one module and derive
the attempt budget from it.

Claude-Session: ced32ebb-7155-4413-adad-1eccd14c2010

* style(push): keep the worker and store in repo formatting

Claude-Session: ced32ebb-7155-4413-adad-1eccd14c2010

* style(push): drop the stray semicolon in the concurrency constant

Claude-Session: ced32ebb-7155-4413-adad-1eccd14c2010
2026-09-30 02:59:48 -04:00

92 lines
2.8 KiB
TypeScript

import { buildPushDelivery } from './push-delivery-message.js'
import type { PushDispatcher } from './push-dispatcher.js'
import type { DurablePushStore } from './durable-push-store.js'
import { WORKER_DRAINS } from './push-worker-concurrency.js'
export class DurablePushWorker {
private timer?: NodeJS.Timeout
private running: Promise<void> | null = null
private stopped = false
constructor(
private readonly store: DurablePushStore,
private readonly dispatcher: PushDispatcher,
private readonly options: { now?: () => number; onRetry?: () => void } = {}
) {}
start(): void {
if (this.timer) return
this.stopped = false
this.timer = setInterval(() => {
void this.runDue().catch(() => {
console.warn(JSON.stringify({ event: 'orca_push_worker_failed' }))
})
}, 1000)
this.timer.unref()
}
async runDue(): Promise<void> {
if (this.running) {
await this.running
return
}
if (this.stopped) return
const drains = Array.from({ length: WORKER_DRAINS }, () => this.drain())
const pending = Promise.allSettled(drains).then(
(results) => {
const failure = results.find((result) => result.status === 'rejected')
if (failure?.status === 'rejected') throw failure.reason
}
)
this.running = pending
try {
await pending
} finally {
this.running = null
}
}
private async drain(): Promise<void> {
for (let count = 0; count < 25 && !this.stopped; count++) {
const queued = await this.store.claim()
if (!queued) return
const delivery = buildPushDelivery({
expiresAt: queued.expiresAt,
registrationId: queued.registrationId,
hostFingerprint: queued.hostFingerprint,
notification: queued.notification
})
if ((this.options.now ?? Date.now)() >= queued.expiresAt) {
await this.store.finish(queued)
continue
}
const heartbeat = setInterval(() => {
void this.store.renew(queued).catch(() => {})
}, 10_000)
heartbeat.unref()
try {
if (queued.attempts > 1) this.options.onRetry?.()
const outcome = await this.dispatcher.sendOnce(delivery)
const retryAfterMs =
outcome.status === 'error' && outcome.retryable
? Math.max(
outcome.retryAfterMs ?? 0,
Math.min(30_000, 1000 * 2 ** Math.min(queued.attempts, 5))
)
: undefined
await this.store.finish(queued, retryAfterMs)
} catch {
await this.store.finish(queued, 5000)
} finally {
clearInterval(heartbeat)
}
}
}
async stop(): Promise<void> {
this.stopped = true
if (this.timer) clearInterval(this.timer)
this.timer = undefined
await this.running
}
}