diff --git a/cloud/apps/push/src/durable-push-store.ts b/cloud/apps/push/src/durable-push-store.ts index 02d26f20b3e..b84ebb17893 100644 --- a/cloud/apps/push/src/durable-push-store.ts +++ b/cloud/apps/push/src/durable-push-store.ts @@ -2,13 +2,14 @@ import { isDismissedAlert, reconcileQueuedDismissal } from './push-queued-dismis import { parsePushDeliveryPayload } from './push-delivery-payload.js' import { createHash, randomUUID } from 'node:crypto' import { PUSH_LIMITS, type PushNotification } from '@orca-cloud/push-contract' +import { WORKER_DRAINS } from './push-worker-concurrency.js' import type { PushDatabase, SqlRow } from './push-database.js' const RETENTION_MS = 24 * 60 * 60_000 // accept() caps expires_at at due_at + TTL, so older due_at is expired: scans skip an unpruned backlog. const TTL_MS = PUSH_LIMITS.notificationTtlSeconds * 1000 -// Covers one claimer per worker drain skipping a device another drain holds. -const CLAIM_CANDIDATE_ATTEMPTS = 4 +// One winner plus every other drain holding a device head, so a full set of peers cannot exhaust it. +const CLAIM_CANDIDATE_ATTEMPTS = WORKER_DRAINS export const PRUNE_BATCH_ROWS = 2_000 export const PRUNE_MAX_BATCHES = 50 export const DELIVERY_LEASE_MS = 30_000 diff --git a/cloud/apps/push/src/durable-push-worker.ts b/cloud/apps/push/src/durable-push-worker.ts index 2a1a22f7a33..f5bcc20fc2c 100644 --- a/cloud/apps/push/src/durable-push-worker.ts +++ b/cloud/apps/push/src/durable-push-worker.ts @@ -1,10 +1,7 @@ import { buildPushDelivery } from './push-delivery-message.js' import type { PushDispatcher } from './push-dispatcher.js' import type { DurablePushStore } from './durable-push-store.js' - -// Four drains capped delivery near 30/s at ~120 ms per item, below the 2026-09 inflow; each drain -// holds one delivery in flight, so this is the worker's concurrency, not its DB draw. -const WORKER_DRAINS = 12 +import { WORKER_DRAINS } from './push-worker-concurrency.js' export class DurablePushWorker { private timer?: NodeJS.Timeout @@ -33,7 +30,8 @@ export class DurablePushWorker { return } if (this.stopped) return - const pending = Promise.allSettled(Array.from({ length: WORKER_DRAINS }, () => this.drain())).then( + 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 diff --git a/cloud/apps/push/src/push-worker-concurrency.ts b/cloud/apps/push/src/push-worker-concurrency.ts new file mode 100644 index 00000000000..daec5659d21 --- /dev/null +++ b/cloud/apps/push/src/push-worker-concurrency.ts @@ -0,0 +1,3 @@ +// Twelve drains lift the ~30/s ceiling four drains hit at ~120 ms per item; each drain holds one +// delivery in flight, so this is the worker's concurrency, not its database draw. +export const WORKER_DRAINS = 12