From e8e09eed46e313cefee64b1165912b60d1bf81b7 Mon Sep 17 00:00:00 2001 From: Jinwoo Hong <73622457+Jinwoo-H@users.noreply.github.com> Date: Wed, 30 Sep 2026 02:59:48 -0400 Subject: [PATCH] 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 --- cloud/apps/push/src/durable-push-store.ts | 5 +++-- cloud/apps/push/src/durable-push-worker.ts | 8 +++----- cloud/apps/push/src/push-worker-concurrency.ts | 3 +++ 3 files changed, 9 insertions(+), 7 deletions(-) create mode 100644 cloud/apps/push/src/push-worker-concurrency.ts 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