From b71ffe09feec77290330cd168780463727fe4fd8 Mon Sep 17 00:00:00 2001 From: Jinwoo-H Date: Sun, 6 Sep 2026 18:42:05 -0400 Subject: [PATCH] fix: harden mobile push delivery and deployment recovery --- .github/workflows/cloud-push-deploy.yml | 17 +- .github/workflows/cloud-verify.yml | 1 + cloud/apps/push/Dockerfile | 6 +- cloud/apps/push/package.json | 3 +- cloud/apps/push/src/apns-client.test.ts | 6 +- cloud/apps/push/src/apns-client.ts | 13 +- cloud/apps/push/src/apns-http2-transport.ts | 7 +- .../push/src/apns-session-replacement.test.ts | 45 +++++ cloud/apps/push/src/apns-stream-response.ts | 13 +- cloud/apps/push/src/coalescer.ts | 21 ++- cloud/apps/push/src/device-registry-store.ts | 13 +- cloud/apps/push/src/fcm-client.test.ts | 24 ++- cloud/apps/push/src/fcm-client.ts | 23 ++- cloud/apps/push/src/host-session-store.ts | 1 + cloud/apps/push/src/index.ts | 40 ++++- cloud/apps/push/src/provider-retry-delay.ts | 9 + .../push-database-postgres-startup.test.ts | 4 +- cloud/apps/push/src/push-database.ts | 8 +- .../push/src/push-delivery-lifecycle.test.ts | 155 +++++++++++++++++ cloud/apps/push/src/push-dispatcher.ts | 24 ++- cloud/apps/push/src/push-observability.ts | 4 +- cloud/apps/push/src/push-provider-outcome.ts | 2 +- cloud/apps/push/src/push-request-drain.ts | 28 +++ .../push/src/push-send-idempotency.test.ts | 34 ++++ .../src/push-server-harness.test-fixture.ts | 1 + cloud/apps/push/src/push-server.ts | 29 +++- .../push/src/push-session-concurrency.test.ts | 73 ++++++++ cloud/apps/push/src/push-session-schema.ts | 23 +++ .../apps/push/src/send-quota-postgres.test.ts | 15 +- cloud/apps/push/src/send-quota.ts | 35 +++- cloud/apps/relay/Dockerfile | 6 +- cloud/apps/relay/package.json | 3 +- .../apps/relay/src/postgres-schema-startup.ts | 106 +----------- .../scripts/push-gateway-recovery.test.mjs | 93 ++++++++++ .../scripts/push-gateway-workflow.test.mjs | 10 +- cloud/docs/push-gateway.md | 22 +++ cloud/package.json | 2 +- cloud/packages/postgres-schema/package.json | 20 +++ cloud/packages/postgres-schema/src/index.ts | 103 +++++++++++ .../postgres-schema/tsconfig.build.json | 11 ++ cloud/packages/postgres-schema/tsconfig.json | 5 + .../src/notification-identity-limits.test.ts | 32 ++++ .../push-contract/src/send-messages.ts | 10 +- cloud/pnpm-lock.yaml | 18 ++ docs/reference/mobile-push-contract.md | 31 +++- .../runtime/push/desktop-push-service.test.ts | 1 + src/main/runtime/push/desktop-push-service.ts | 55 +++++- .../runtime/push/push-agent-state.test.ts | 21 +++ .../push/push-cleanup-auth-expiry.test.ts | 41 +++++ .../push/push-dispatcher.test-fixture.ts | 94 ++++++++++ src/main/runtime/push/push-dispatcher.test.ts | 120 ++----------- src/main/runtime/push/push-dispatcher.ts | 52 +++++- src/main/runtime/push/push-gateway-client.ts | 2 +- .../push/push-outcome-counters.test.ts | 25 +++ .../runtime/push/push-outcome-counters.ts | 27 +++ .../push/push-registration-races.test.ts | 160 ++++++++++++++++++ 56 files changed, 1439 insertions(+), 308 deletions(-) create mode 100644 cloud/apps/push/src/apns-session-replacement.test.ts create mode 100644 cloud/apps/push/src/provider-retry-delay.ts create mode 100644 cloud/apps/push/src/push-delivery-lifecycle.test.ts create mode 100644 cloud/apps/push/src/push-request-drain.ts create mode 100644 cloud/apps/push/src/push-send-idempotency.test.ts create mode 100644 cloud/apps/push/src/push-session-concurrency.test.ts create mode 100644 cloud/apps/push/src/push-session-schema.ts create mode 100644 cloud/dev/scripts/push-gateway-recovery.test.mjs create mode 100644 cloud/packages/postgres-schema/package.json create mode 100644 cloud/packages/postgres-schema/src/index.ts create mode 100644 cloud/packages/postgres-schema/tsconfig.build.json create mode 100644 cloud/packages/postgres-schema/tsconfig.json create mode 100644 cloud/packages/push-contract/src/notification-identity-limits.test.ts create mode 100644 src/main/runtime/push/push-agent-state.test.ts create mode 100644 src/main/runtime/push/push-cleanup-auth-expiry.test.ts create mode 100644 src/main/runtime/push/push-dispatcher.test-fixture.ts create mode 100644 src/main/runtime/push/push-outcome-counters.test.ts create mode 100644 src/main/runtime/push/push-outcome-counters.ts create mode 100644 src/main/runtime/push/push-registration-races.test.ts diff --git a/.github/workflows/cloud-push-deploy.yml b/.github/workflows/cloud-push-deploy.yml index 53af9f5d755..9290b4ab2ce 100644 --- a/.github/workflows/cloud-push-deploy.yml +++ b/.github/workflows/cloud-push-deploy.yml @@ -130,11 +130,14 @@ jobs: run: | set -euo pipefail tag="c${GITHUB_RUN_ID}-${GITHUB_RUN_ATTEMPT}" + echo "CANDIDATE_TAG=${tag}" >> "${GITHUB_ENV}" + echo "CANDIDATE_REVISION=${SERVICE_NAME}-${tag}" >> "${GITHUB_ENV}" gcloud run deploy "${SERVICE_NAME}" \ --project "${GCP_PROJECT_ID}" \ --region "${GCP_REGION}" \ --image "${IMAGE}" \ --tag "${tag}" \ + --revision-suffix "${tag}" \ --no-traffic \ --quiet candidate="$(gcloud run services describe "${SERVICE_NAME}" \ @@ -142,8 +145,7 @@ jobs: | jq -er --arg tag "${tag}" \ '[.status.traffic[] | select(.tag == $tag)] | if length == 1 then .[0] else error("tagged candidate is not unique") end')" - echo "CANDIDATE_TAG=${tag}" >> "${GITHUB_ENV}" - echo "CANDIDATE_REVISION=$(jq -r '.revisionName' <<< "${candidate}")" >> "${GITHUB_ENV}" + test "$(jq -r '.revisionName' <<< "${candidate}")" = "${SERVICE_NAME}-${tag}" echo "CANDIDATE_URL=$(jq -r '.url' <<< "${candidate}")" >> "${GITHUB_ENV}" # A tagged revision is directly addressable and sits outside the service-wide cap, so the @@ -224,6 +226,7 @@ jobs: shell: bash run: | set -euo pipefail + echo "TRAFFIC_SHIFT_ATTEMPTED=true" >> "${GITHUB_ENV}" gcloud run services update-traffic "${SERVICE_NAME}" \ --project "${GCP_PROJECT_ID}" \ --region "${GCP_REGION}" \ @@ -240,11 +243,12 @@ jobs: # moved, the rollback target is the single thing an operator needs, and a summary that only # appeared on success would be missing in exactly the run that needs it. - name: Publish the rollout summary + if: ${{ always() && env.CANDIDATE_REVISION != '' && env.ROLLBACK_REVISION != '' }} shell: bash run: | set -euo pipefail { - echo '### Push gateway deployed' + echo '### Push gateway rollout' echo echo "Revision: \`${CANDIDATE_REVISION}\`" echo @@ -275,7 +279,7 @@ jobs: # is not a failure to deploy, it is a live gateway that has to go back, so the traffic move # is undone here rather than left to whoever reads the run. - name: Roll traffic back to the previous revision - if: ${{ failure() && env.TRAFFIC_SHIFTED == 'true' }} + if: ${{ (failure() || cancelled()) && env.TRAFFIC_SHIFT_ATTEMPTED == 'true' }} shell: bash run: | set -euo pipefail @@ -290,6 +294,7 @@ jobs: | jq -r '[.status.traffic[] | select((.percent // 0) > 0)] | if length == 1 and .[0].percent == 100 then .[0].revisionName else empty end')" test "${serving}" = "${ROLLBACK_REVISION}" + echo "TRAFFIC_ROLLED_BACK=true" >> "${GITHUB_ENV}" { echo echo '### Push gateway rolled back' @@ -302,8 +307,8 @@ jobs: # SQL pool for nothing. Its tag comes off first, because Cloud Run refuses to delete a # revision a traffic target still names, and clearing CANDIDATE_TAG makes the always() tag # step below a no-op rather than a second failure. - - name: Delete the candidate revision that never took traffic - if: ${{ failure() && env.TRAFFIC_SHIFTED != 'true' }} + - name: Delete the rejected candidate revision + if: ${{ (failure() || cancelled()) && (env.TRAFFIC_SHIFT_ATTEMPTED != 'true' || env.TRAFFIC_ROLLED_BACK == 'true') }} shell: bash run: | set -euo pipefail diff --git a/.github/workflows/cloud-verify.yml b/.github/workflows/cloud-verify.yml index e2ba9407ac4..5e24cae76cc 100644 --- a/.github/workflows/cloud-verify.yml +++ b/.github/workflows/cloud-verify.yml @@ -90,6 +90,7 @@ jobs: --health-timeout 5s --health-retries 10 env: + ORCA_PUSH_TEST_DATABASE_URL: postgres://relay_test:relay_test@127.0.0.1:5432/orca_relay_test ORCA_RELAY_TEST_POSTGRES_URL: postgres://relay_test:relay_test@127.0.0.1:5432/orca_relay_test steps: - uses: actions/checkout@v4 diff --git a/cloud/apps/push/Dockerfile b/cloud/apps/push/Dockerfile index 48537c5ce07..efdc85fc404 100644 --- a/cloud/apps/push/Dockerfile +++ b/cloud/apps/push/Dockerfile @@ -3,11 +3,13 @@ WORKDIR /app RUN corepack enable COPY package.json pnpm-lock.yaml pnpm-workspace.yaml tsconfig.base.json ./ COPY packages/push-contract/package.json packages/push-contract/package.json +COPY packages/postgres-schema/package.json packages/postgres-schema/package.json COPY apps/push/package.json apps/push/package.json RUN pnpm install --frozen-lockfile COPY packages/push-contract packages/push-contract COPY apps/push apps/push -RUN pnpm --filter @orca-cloud/push-contract build && pnpm --filter @orca-cloud/push build +COPY packages/postgres-schema packages/postgres-schema +RUN pnpm --filter @orca-cloud/postgres-schema build && pnpm --filter @orca-cloud/push-contract build && pnpm --filter @orca-cloud/push build FROM node:24-alpine AS runtime ENV NODE_ENV=production @@ -16,8 +18,10 @@ WORKDIR /app RUN corepack enable COPY package.json pnpm-lock.yaml pnpm-workspace.yaml ./ COPY packages/push-contract/package.json packages/push-contract/package.json +COPY packages/postgres-schema/package.json packages/postgres-schema/package.json COPY apps/push/package.json apps/push/package.json COPY --from=build /app/packages/push-contract/dist packages/push-contract/dist +COPY --from=build /app/packages/postgres-schema/dist packages/postgres-schema/dist COPY --from=build /app/apps/push/dist apps/push/dist RUN pnpm install --prod --frozen-lockfile --filter @orca-cloud/push... USER node diff --git a/cloud/apps/push/package.json b/cloud/apps/push/package.json index ba59f6fae86..d84d0af8b25 100644 --- a/cloud/apps/push/package.json +++ b/cloud/apps/push/package.json @@ -9,13 +9,14 @@ "clean": "node -e \"require('fs').rmSync('dist', { recursive: true, force: true })\"", "dev": "tsx watch src/index.ts", "lint": "tsc -p tsconfig.json --noEmit", - "pretest": "pnpm --filter @orca-cloud/push-contract build", + "pretest": "pnpm --filter @orca-cloud/postgres-schema build && pnpm --filter @orca-cloud/push-contract build", "start": "node dist/index.js", "test": "vitest run", "typecheck": "tsc -p tsconfig.json --noEmit" }, "dependencies": { "@hono/node-server": "^1.19.14", + "@orca-cloud/postgres-schema": "workspace:*", "@orca-cloud/push-contract": "workspace:*", "google-auth-library": "^10.5.0", "hono": "^4.12.27", diff --git a/cloud/apps/push/src/apns-client.test.ts b/cloud/apps/push/src/apns-client.test.ts index 1e7a384f6ab..c0f312e46e6 100644 --- a/cloud/apps/push/src/apns-client.test.ts +++ b/cloud/apps/push/src/apns-client.test.ts @@ -147,7 +147,7 @@ describe('apns client', () => { [400, 'PayloadTooLarge'], [429, 'TooManyRequests'], [500, 'InternalServerError'] - ])('treats %i %s as a retryable error, not a dead token', async (status, reason) => { + ])('treats %i %s with the appropriate retry policy', async (status, reason) => { const fake = fakeTransport({ status, body: JSON.stringify({ reason }) }) const client = new ApnsClient({ topic: 'com.stably.orca.mobile', @@ -156,7 +156,7 @@ describe('apns client', () => { }) await expect( client.send(delivery(), { token: 'a'.repeat(64), apnsEnvironment: 'production' }) - ).resolves.toEqual({ status: 'error', reason }) + ).resolves.toEqual({ status: 'error', reason, retryable: status === 429 || status >= 500 }) }) it('reports a transport failure as an error rather than throwing', async () => { @@ -169,6 +169,6 @@ describe('apns client', () => { }) await expect( client.send(delivery(), { token: 'a'.repeat(64), apnsEnvironment: 'production' }) - ).resolves.toEqual({ status: 'error', reason: 'Error' }) + ).resolves.toEqual({ status: 'error', reason: 'Error', retryable: true }) }) }) diff --git a/cloud/apps/push/src/apns-client.ts b/cloud/apps/push/src/apns-client.ts index 9bca65c6d67..7e6eb813865 100644 --- a/cloud/apps/push/src/apns-client.ts +++ b/cloud/apps/push/src/apns-client.ts @@ -69,7 +69,11 @@ export class ApnsClient { body: apnsBody(delivery) }) } catch (error) { - return { status: 'error', reason: error instanceof Error ? error.name : 'transport_failed' } + return { + status: 'error', + reason: error instanceof Error ? error.name : 'transport_failed', + retryable: true + } } if (response.status === 200) return { status: 'sent' } const reason = readReason(response.body) @@ -77,6 +81,11 @@ export class ApnsClient { if (response.status === 400 && DEAD_TOKEN_REASONS.has(reason)) { return { status: 'dead', reason } } - return { status: 'error', reason } + return { + status: 'error', + reason, + retryable: response.status === 429 || response.status >= 500, + ...(response.retryAfterMs === undefined ? {} : { retryAfterMs: response.retryAfterMs }) + } } } diff --git a/cloud/apps/push/src/apns-http2-transport.ts b/cloud/apps/push/src/apns-http2-transport.ts index 5bcec7049f8..167b4d14e38 100644 --- a/cloud/apps/push/src/apns-http2-transport.ts +++ b/cloud/apps/push/src/apns-http2-transport.ts @@ -20,8 +20,11 @@ export function createApnsHttp2Transport(): ApnsTransport & { close(): void } { const existing = sessions.get(host) if (existing && !existing.closed && !existing.destroyed) return existing const session = connect(`https://${host}`) - session.on('error', () => sessions.delete(host)) - session.on('close', () => sessions.delete(host)) + const forget = (): void => { + if (sessions.get(host) === session) sessions.delete(host) + } + session.on('error', forget) + session.on('close', forget) sessions.set(host, session) return session } diff --git a/cloud/apps/push/src/apns-session-replacement.test.ts b/cloud/apps/push/src/apns-session-replacement.test.ts new file mode 100644 index 00000000000..2678732ca94 --- /dev/null +++ b/cloud/apps/push/src/apns-session-replacement.test.ts @@ -0,0 +1,45 @@ +import { EventEmitter } from 'node:events' +import { expect, it, vi } from 'vitest' +const mocks = vi.hoisted(() => ({ + connect: vi.fn(), + read: vi.fn(async () => ({ status: 200, body: '' })) +})) +vi.mock('node:http2', async (original) => ({ + ...(await original()), + connect: mocks.connect +})) +vi.mock('./apns-stream-response.js', () => ({ readApnsStreamResponse: mocks.read })) +import { createApnsHttp2Transport } from './apns-http2-transport.js' + +it('keeps the replacement cached when the draining session closes later', async () => { + const sessions: Array< + EventEmitter & { + closed: boolean + destroyed: boolean + request: ReturnType + close: ReturnType + } + > = [] + mocks.connect.mockImplementation(() => { + const session = Object.assign(new EventEmitter(), { + closed: false, + destroyed: false, + request: vi.fn(() => ({})), + close: vi.fn() + }) + sessions.push(session) + return session + }) + const transport = createApnsHttp2Transport() + const request = { host: 'api.push.apple.com', path: '/synthetic', headers: {}, body: '{}' } + await transport(request) + sessions[0]!.closed = true + await transport(request) + sessions[0]!.emit('close') + sessions[0]!.emit('error', new Error('old-session')) + await transport(request) + expect(sessions).toHaveLength(2) + expect(sessions[1]!.request).toHaveBeenCalledTimes(2) + transport.close() + expect(sessions[1]!.close).toHaveBeenCalledOnce() +}) diff --git a/cloud/apps/push/src/apns-stream-response.ts b/cloud/apps/push/src/apns-stream-response.ts index 5bb452fa3b3..da001a5df31 100644 --- a/cloud/apps/push/src/apns-stream-response.ts +++ b/cloud/apps/push/src/apns-stream-response.ts @@ -1,7 +1,8 @@ import type { EventEmitter } from 'node:events' +import { providerRetryAfter } from './provider-retry-delay.js' import { constants } from 'node:http2' -export type ApnsResponse = { status: number; body: string } +export type ApnsResponse = { status: number; body: string; retryAfterMs?: number } // The subset of ClientHttp2Stream this module drives, so a fake emitter can // stand in for a real APNs stream in tests. @@ -26,15 +27,23 @@ export function readApnsStreamResponse( run() } let status = 0 + let retryAfterMs: number | undefined const chunks: Buffer[] = [] stream.setTimeout(timeoutMs, () => stream.destroy(new Error('apns_timeout'))) stream.on('response', (headers: Record) => { status = Number(headers[constants.HTTP2_HEADER_STATUS] ?? 0) + retryAfterMs = providerRetryAfter(String(headers['retry-after'] ?? '')) }) stream.on('data', (chunk: Buffer) => chunks.push(chunk)) stream.on('error', (error: Error) => settle(() => reject(error))) stream.on('end', () => - settle(() => resolve({ status, body: Buffer.concat(chunks).toString('utf8') })) + settle(() => + resolve({ + status, + body: Buffer.concat(chunks).toString('utf8'), + ...(retryAfterMs === undefined ? {} : { retryAfterMs }) + }) + ) ) // A peer reset with NGHTTP2_NO_ERROR emits neither 'end' nor 'error', which // would leave the coalescer's delivery pending for the life of the process. diff --git a/cloud/apps/push/src/coalescer.ts b/cloud/apps/push/src/coalescer.ts index 45104b8d7cd..f55b6757418 100644 --- a/cloud/apps/push/src/coalescer.ts +++ b/cloud/apps/push/src/coalescer.ts @@ -37,6 +37,8 @@ export function summaryBody(notifications: readonly PushNotification[]): string // Holds sends per registration for one window so a burst of desktop events // reaches the phone as a single banner instead of a stack of near-duplicates. export class PushCoalescer { + private readonly deliveries = new Set>() + private stopped = false private readonly windows = new Map() private readonly windowMs: number private readonly setTimer: (callback: () => void, delayMs: number) => CoalescerTimer @@ -53,6 +55,7 @@ export class PushCoalescer { hostFingerprint: string notification: PushNotification }): void { + if (this.stopped) throw new Error('push_coalescer_stopped') const existing = this.windows.get(input.registrationId) if (existing) { existing.notifications.push(input.notification) @@ -86,18 +89,28 @@ export class PushCoalescer { body: coalescedCount > 1 ? summaryBody(window.notifications) : latest.body, coalescedCount }) + const pending = Promise.resolve() + .then(() => this.options.deliver(delivery)) + .catch((error) => { + this.options.onDeliveryFailed?.(error) + }) + this.deliveries.add(pending) try { - await this.options.deliver(delivery) - } catch (error) { - this.options.onDeliveryFailed?.(error) + await pending + } finally { + this.deliveries.delete(pending) } } async flushAll(): Promise { - await Promise.all([...this.windows.keys()].map((id) => this.flush(id))) + do { + await Promise.all([...this.windows.keys()].map((id) => this.flush(id))) + await Promise.all([...this.deliveries]) + } while (this.windows.size || this.deliveries.size) } stop(): void { + this.stopped = true for (const window of this.windows.values()) this.clearTimer(window.timer) this.windows.clear() } diff --git a/cloud/apps/push/src/device-registry-store.ts b/cloud/apps/push/src/device-registry-store.ts index 45af7bbfea0..9aac22dd25c 100644 --- a/cloud/apps/push/src/device-registry-store.ts +++ b/cloud/apps/push/src/device-registry-store.ts @@ -169,10 +169,17 @@ export class PushDeviceRegistryStore { return row ? toRegistration(row) : null } - async markDead(registrationId: string): Promise { + async markDead(registrationId: string, observed?: PushDeviceRegistration): Promise { await this.database.query( - 'UPDATE push_devices SET dead_at = ?, updated_at = ? WHERE registration_id = ?', - [this.now(), this.now(), registrationId] + `UPDATE push_devices SET dead_at = ?, updated_at = ? WHERE registration_id = ?${ + observed ? " AND token = ? AND platform = ? AND COALESCE(apns_environment, '') = ?" : '' + }`, + [ + this.now(), + this.now(), + registrationId, + ...(observed ? [observed.token, observed.platform, observed.apnsEnvironment ?? ''] : []) + ] ) } } diff --git a/cloud/apps/push/src/fcm-client.test.ts b/cloud/apps/push/src/fcm-client.test.ts index e1c61c4cf62..3069c62032b 100644 --- a/cloud/apps/push/src/fcm-client.test.ts +++ b/cloud/apps/push/src/fcm-client.test.ts @@ -54,9 +54,7 @@ describe('fcm client', () => { const { fake, client: fcm } = client({ status: 200, body: '{"name":"projects/x/messages/1"}' }) await expect(fcm.send(delivery(), { token: TOKEN })).resolves.toEqual({ status: 'sent' }) const request = fake.requests[0]! - expect(request.url).toBe( - 'https://fcm.googleapis.com/v1/projects/onorca-cloud/messages:send' - ) + expect(request.url).toBe('https://fcm.googleapis.com/v1/projects/onorca-cloud/messages:send') expect(request.accessToken).toBe('access-token') expect(JSON.parse(request.body)).toEqual({ message: { @@ -86,9 +84,14 @@ describe('fcm client', () => { const { fake, client: fcm } = client({ status: 200, body: '{}' }) await fcm.send(delivery(3, null), { token: TOKEN }) const message = JSON.parse(fake.requests[0]!.body) as { - message: { android: { collapse_key: string; notification: { tag: string } }; data: Record } + message: { + android: { collapse_key: string; notification: { tag: string } } + data: Record + } } - expect(Object.values(message.message.data).every((value) => typeof value === 'string')).toBe(true) + expect(Object.values(message.message.data).every((value) => typeof value === 'string')).toBe( + true + ) expect(message.message.data.agentState).toBeUndefined() expect(message.message.data.coalescedCount).toBe('3') expect(message.message.android.notification.tag).toBe(`host:${HOST}`) @@ -146,7 +149,9 @@ describe('fcm client', () => { }) await expect(unnamed.client.send(delivery(), { token: TOKEN })).resolves.toEqual({ status: 'error', - reason: 'INVALID_ARGUMENT' + reason: 'INVALID_ARGUMENT', + retryable: false, + retryAfterMs: 10000 }) }) @@ -157,7 +162,9 @@ describe('fcm client', () => { }) await expect(faulted.client.send(delivery(), { token: TOKEN })).resolves.toEqual({ status: 'error', - reason: 'UNAVAILABLE' + reason: 'UNAVAILABLE', + retryable: true, + retryAfterMs: 10000 }) const broken = new FcmClient({ projectId: 'onorca-cloud', @@ -168,7 +175,8 @@ describe('fcm client', () => { }) await expect(broken.send(delivery(), { token: TOKEN })).resolves.toEqual({ status: 'error', - reason: 'Error' + reason: 'Error', + retryable: true }) }) }) diff --git a/cloud/apps/push/src/fcm-client.ts b/cloud/apps/push/src/fcm-client.ts index 5c86d7b9736..3e9329e2d34 100644 --- a/cloud/apps/push/src/fcm-client.ts +++ b/cloud/apps/push/src/fcm-client.ts @@ -1,3 +1,4 @@ +import { providerRetryAfter } from './provider-retry-delay.js' import { createHash } from 'node:crypto' import { PUSH_DEFAULTS, PUSH_LIMITS } from '@orca-cloud/push-contract' import { orcaDataStrings, type PushDelivery } from './push-delivery-message.js' @@ -6,7 +7,7 @@ import type { PushProviderOutcome } from './push-provider-outcome.js' export const FCM_SCOPE = 'https://www.googleapis.com/auth/firebase.messaging' export type FcmRequest = { url: string; accessToken: string; body: string } -export type FcmResponse = { status: number; body: string } +export type FcmResponse = { status: number; body: string; retryAfterMs?: number } export type FcmTransport = (request: FcmRequest) => Promise export type FcmClientOptions = { @@ -89,7 +90,11 @@ export class FcmClient { }) }) } catch (error) { - return { status: 'error', reason: error instanceof Error ? error.name : 'transport_failed' } + return { + status: 'error', + reason: error instanceof Error ? error.name : 'transport_failed', + retryable: true + } } if (response.status >= 200 && response.status < 300) return { status: 'sent' } const failure = readFcmError(response.body) @@ -100,7 +105,12 @@ export class FcmClient { if (failure.status === 'INVALID_ARGUMENT' && /\btoken\b/i.test(failure.message)) { return { status: 'dead', reason: 'INVALID_ARGUMENT' } } - return { status: 'error', reason: failure.status } + return { + status: 'error', + reason: failure.status, + retryable: response.status === 429 || response.status >= 500, + retryAfterMs: Math.max(response.status === 429 ? 60_000 : 10_000, response.retryAfterMs ?? 0) + } } } @@ -113,8 +123,13 @@ export function createFcmFetchTransport(fetchImpl: typeof fetch = fetch): FcmTra 'content-type': 'application/json' }, body: request.body, + redirect: 'error', signal: AbortSignal.timeout(10_000) }) - return { status: response.status, body: await response.text() } + return { + status: response.status, + body: await response.text(), + retryAfterMs: providerRetryAfter(response.headers.get('retry-after') ?? undefined) + } } } diff --git a/cloud/apps/push/src/host-session-store.ts b/cloud/apps/push/src/host-session-store.ts index d400ce9b926..899bacabcc8 100644 --- a/cloud/apps/push/src/host-session-store.ts +++ b/cloud/apps/push/src/host-session-store.ts @@ -30,6 +30,7 @@ export class PushHostSessionStore { // Why: a desktop holds one session at a time and only re-proves once it is // gone, so an earlier row is dead weight. It also bounds the table to one // row per host however many proofs a self-minted identity answers. + await transaction.lockQuotaScope(`orca-push-session:${hostFingerprint}`) await transaction.query('DELETE FROM push_sessions WHERE host_fingerprint = ?', [ hostFingerprint ]) diff --git a/cloud/apps/push/src/index.ts b/cloud/apps/push/src/index.ts index 04d259801fd..c3415dc307a 100644 --- a/cloud/apps/push/src/index.ts +++ b/cloud/apps/push/src/index.ts @@ -14,8 +14,16 @@ const database = await openPushDatabase({ poolMax: config.databasePoolMax, applicationName: 'orca-push' }) -const { server, challenges, sessions, quota, coalescer, observability, closeTransports } = - createPushServer(config, database) +const { + server, + challenges, + sessions, + quota, + coalescer, + observability, + closeTransports, + requestDrain +} = createPushServer(config, database) function prune(label: string, run: () => Promise, intervalMs: number): NodeJS.Timeout { const timer = setInterval(() => { @@ -45,15 +53,29 @@ server.listen(config.port, () => { console.log(`[orca-push] listening on ${config.publicUrl} (port ${config.port})`) }) +let stopping = false const shutdown = (): void => { + if (stopping) return + stopping = true for (const timer of timers) clearInterval(timer) - observability.stop() - // Drain the coalescing windows so an in-flight burst still reaches the phone. - void coalescer.flushAll().finally(() => { - coalescer.stop() - closeTransports() - server.close(() => void database.close()) - }) + // Cloud Run sends SIGKILL after ten seconds; leave time for explicit cleanup. + const deadline = setTimeout(() => process.exit(1), 9_000) + deadline.unref() + const requests = requestDrain.begin() + const connections = new Promise((resolve) => server.close(() => resolve())) + void Promise.all([requests, connections]) + .then(async () => { + await coalescer.flushAll() + coalescer.stop() + closeTransports() + await database.close() + observability.stop() + clearTimeout(deadline) + }) + .catch(() => { + console.warn(JSON.stringify({ event: 'orca_push_shutdown_failed' })) + process.exitCode = 1 + }) } process.once('SIGTERM', shutdown) process.once('SIGINT', shutdown) diff --git a/cloud/apps/push/src/provider-retry-delay.ts b/cloud/apps/push/src/provider-retry-delay.ts new file mode 100644 index 00000000000..4c77b3c6dc7 --- /dev/null +++ b/cloud/apps/push/src/provider-retry-delay.ts @@ -0,0 +1,9 @@ +export function providerRetryAfter( + value: string | undefined, + now = Date.now() +): number | undefined { + if (!value) return undefined + const seconds = Number(value) + const delay = Number.isFinite(seconds) ? seconds * 1000 : Date.parse(value) - now + return Number.isFinite(delay) ? Math.max(0, delay) : undefined +} diff --git a/cloud/apps/push/src/push-database-postgres-startup.test.ts b/cloud/apps/push/src/push-database-postgres-startup.test.ts index 97736ef25a6..181016d062a 100644 --- a/cloud/apps/push/src/push-database-postgres-startup.test.ts +++ b/cloud/apps/push/src/push-database-postgres-startup.test.ts @@ -60,7 +60,9 @@ describe('PostgreSQL push gateway startup', () => { lock_timeout: 1_000, idle_in_transaction_session_timeout: 5_000 }) - expect(fakes.query.mock.calls.map(([sql]) => sql)).toEqual(pushSchemaStatements()) + expect( + fakes.query.mock.calls.map(([sql]) => sql).slice(0, pushSchemaStatements().length) + ).toEqual(pushSchemaStatements()) await database.close() }) diff --git a/cloud/apps/push/src/push-database.ts b/cloud/apps/push/src/push-database.ts index a396bdde9f5..6f8ba88ed1d 100644 --- a/cloud/apps/push/src/push-database.ts +++ b/cloud/apps/push/src/push-database.ts @@ -2,6 +2,8 @@ import { mkdirSync } from 'node:fs' import { join } from 'node:path' import { DatabaseSync } from 'node:sqlite' import pg from 'pg' +import { applyPostgresSchema } from '@orca-cloud/postgres-schema' +import { ensurePushSessionIndex } from './push-session-schema.js' import { pushSchemaStatements } from './push-schema.js' const POSTGRES_LOCK_TIMEOUT_MS = 1_000 @@ -187,6 +189,7 @@ class PostgresDatabase implements PushDatabase { async function applySchema(database: PushDatabase): Promise { for (const statement of pushSchemaStatements()) await database.query(statement) + await ensurePushSessionIndex(database) } // Why: DDL is not a request. A CREATE INDEX on a grown table can legitimately @@ -210,7 +213,10 @@ async function applySchemaOnUntimedPool( absorbPostgresIdleClientErrors(pool) const database = new PostgresDatabase(pool) try { - await applySchema(database) + await applyPostgresSchema(pushSchemaStatements(), (statement) => database.query(statement), { + eventPrefix: 'orca_push_postgres_schema' + }) + await ensurePushSessionIndex(database) } finally { await database.close().catch(() => undefined) } diff --git a/cloud/apps/push/src/push-delivery-lifecycle.test.ts b/cloud/apps/push/src/push-delivery-lifecycle.test.ts new file mode 100644 index 00000000000..88d95081515 --- /dev/null +++ b/cloud/apps/push/src/push-delivery-lifecycle.test.ts @@ -0,0 +1,155 @@ +import { afterEach, expect, it, vi } from 'vitest' +import { Hono } from 'hono' +import { PushRequestDrain } from './push-request-drain.js' +import { PushCoalescer } from './coalescer.js' +import { PushDispatcher } from './push-dispatcher.js' +import { PushDeviceRegistryStore } from './device-registry-store.js' +import { openInMemoryPushDatabase, type PushDatabase } from './push-database.js' +import { buildPushDelivery } from './push-delivery-message.js' +import { PushNotificationSchema } from '@orca-cloud/push-contract' +import { notification } from './push-server-harness.test-fixture.js' + +const databases: PushDatabase[] = [] +afterEach(async () => { + await Promise.all(databases.splice(0).map((db) => db.close())) + vi.restoreAllMocks() +}) +const note = PushNotificationSchema.parse(notification()) +const tick = () => new Promise((resolve) => setImmediate(resolve)) +function deferred() { + let resolve!: () => void + const promise = new Promise((done) => { + resolve = done + }) + return { promise, resolve } +} +async function registered() { + const db = await openInMemoryPushDatabase() + databases.push(db) + const devices = new PushDeviceRegistryStore(db) + const input = { + hostFingerprint: 'abcdefghijklmnop', + deviceId: 'device', + platform: 'android' as const, + token: 'old-token', + filter: { sources: [], agentStates: [] } + } + const row = await devices.upsert(input) + if (!row.ok) throw new Error('registration failed') + const delivery = buildPushDelivery({ + registrationId: row.registrationId, + hostFingerprint: input.hostFingerprint, + notification: note, + title: note.title, + body: note.body, + coalescedCount: 1 + }) + return { db, devices, input, delivery } +} + +it('does not retire a refreshed token after the old token fails', async () => { + const h = await registered() + const gate = deferred() + const send = vi.fn(async () => { + await gate.promise + return { status: 'dead', reason: 'UNREGISTERED' } + }) + vi.spyOn(console, 'warn').mockImplementation(() => {}) + const dispatcher = new PushDispatcher({ devices: h.devices, fcm: { send } as never }) + const pending = dispatcher.deliver(h.delivery) + await tick() + await h.devices.upsert({ ...h.input, token: 'replacement-token' }) + gate.resolve() + await pending + expect(await h.devices.findById(h.delivery.registrationId)).toMatchObject({ + token: 'replacement-token', + dead: false + }) +}) + +it('drains timer-triggered deliveries that already left the window map', async () => { + const gate = deferred() + const deliver = vi.fn(() => gate.promise) + const coalescer = new PushCoalescer({ + deliver, + setTimer: () => ({ handle: null }), + clearTimer: () => {} + }) + coalescer.enqueue({ + registrationId: 'reg', + hostFingerprint: 'abcdefghijklmnop', + notification: note + }) + const pending = coalescer.flush('reg') + let drained = false + const drain = coalescer.flushAll().then(() => { + drained = true + }) + await tick() + expect(deliver).toHaveBeenCalledOnce() + expect(drained).toBe(false) + gate.resolve() + await Promise.all([pending, drain]) + expect(drained).toBe(true) +}) + +it('rejects new requests during drain and waits for an admitted handler', async () => { + const gate = deferred() + const requests = new PushRequestDrain() + const app = new Hono().use('*', requests.middleware).post('/send', async (c) => { + await gate.promise + return c.json({ queued: true }) + }) + const pending = app.request('/send', { method: 'POST' }) + await tick() + let drained = false + const drain = requests.begin().then(() => { + drained = true + }) + expect((await app.request('/send', { method: 'POST' })).status).toBe(503) + expect(drained).toBe(false) + gate.resolve() + expect((await pending).status).toBe(200) + await drain + expect(drained).toBe(true) +}) + +it('retries transient failures with the provider delay and stops after success', async () => { + const h = await registered() + vi.spyOn(console, 'warn').mockImplementation(() => {}) + const send = vi + .fn() + .mockResolvedValueOnce({ + status: 'error', + reason: 'UNAVAILABLE', + retryable: true, + retryAfterMs: 10000 + }) + .mockResolvedValue({ status: 'sent' }) + const wait = vi.fn(async (_ms: number) => {}) + await new PushDispatcher({ devices: h.devices, fcm: { send } as never, wait }).deliver(h.delivery) + expect(send).toHaveBeenCalledTimes(2) + expect(wait).toHaveBeenCalledExactlyOnceWith(expect.any(Number)) + expect(wait.mock.calls[0]![0]).toBeGreaterThanOrEqual(10000) +}) + +it('bounds retries and rechecks registration after waiting', async () => { + const h = await registered() + vi.spyOn(console, 'warn').mockImplementation(() => {}) + const send = vi.fn().mockResolvedValue({ status: 'error', reason: 'timeout', retryable: true }) + await new PushDispatcher({ + devices: h.devices, + fcm: { send } as never, + wait: async () => {} + }).deliver(h.delivery) + expect(send).toHaveBeenCalledTimes(3) + send.mockClear() + await new PushDispatcher({ + devices: h.devices, + fcm: { send } as never, + wait: async () => { + await h.devices.deleteOwned(h.input.hostFingerprint, h.delivery.registrationId) + } + }).deliver(h.delivery) + expect(send).toHaveBeenCalledOnce() +}) diff --git a/cloud/apps/push/src/push-dispatcher.ts b/cloud/apps/push/src/push-dispatcher.ts index d2bd4b1a41c..39d17c92d11 100644 --- a/cloud/apps/push/src/push-dispatcher.ts +++ b/cloud/apps/push/src/push-dispatcher.ts @@ -9,6 +9,9 @@ export type PushDispatcherOptions = { devices: PushDeviceRegistryStore apns?: ApnsClient fcm?: FcmClient + wait?: (ms: number) => Promise + now?: () => number + onRetry?: () => void onOutcome?: (outcome: PushProviderOutcome['status']) => void } @@ -18,6 +21,22 @@ export class PushDispatcher { constructor(private readonly options: PushDispatcherOptions) {} async deliver(delivery: PushDelivery): Promise { + const now = this.options.now ?? Date.now + const deadline = now() + 120_000 + for (let attempt = 0; attempt < 3; attempt++) { + if (now() >= deadline) return + const retry = await this.deliverAttempt(delivery) + if (!retry || attempt === 2) return + const delay = Math.max(retry.delayMs, 1000 * 2 ** attempt) + Math.floor(Math.random() * 250) + if (now() + delay >= deadline) return + this.options.onRetry?.() + await (this.options.wait ?? ((ms) => new Promise((resolve) => setTimeout(resolve, ms))))( + delay + ) + } + } + + private async deliverAttempt(delivery: PushDelivery): Promise<{ delayMs: number } | undefined> { const device = await this.options.devices.findById(delivery.registrationId) if (!device || device.dead) return let outcome: PushProviderOutcome @@ -35,7 +54,7 @@ export class PushDispatcher { } this.options.onOutcome?.(outcome.status) if (outcome.status === 'dead') { - await this.options.devices.markDead(delivery.registrationId) + await this.options.devices.markDead(delivery.registrationId, device) } if (outcome.status !== 'sent') { console.warn( @@ -48,5 +67,8 @@ export class PushDispatcher { }) ) } + if (outcome.status === 'error' && outcome.retryable) + return { delayMs: outcome.retryAfterMs ?? 0 } + return undefined } } diff --git a/cloud/apps/push/src/push-observability.ts b/cloud/apps/push/src/push-observability.ts index c5a8c6d233e..4840723b7ec 100644 --- a/cloud/apps/push/src/push-observability.ts +++ b/cloud/apps/push/src/push-observability.ts @@ -15,6 +15,7 @@ type PushCounterName = | 'delivery_sent' | 'delivery_dead' | 'delivery_error' + | 'delivery_retry' const COUNTER_NAMES: PushCounterName[] = [ 'ip_rate_limited', @@ -32,7 +33,8 @@ const COUNTER_NAMES: PushCounterName[] = [ 'send_error', 'delivery_sent', 'delivery_dead', - 'delivery_error' + 'delivery_error', + 'delivery_retry' ] // Aggregate counters only. Nothing here may accept a token, a title, a body, diff --git a/cloud/apps/push/src/push-provider-outcome.ts b/cloud/apps/push/src/push-provider-outcome.ts index aadf49f15c6..bc65d10c175 100644 --- a/cloud/apps/push/src/push-provider-outcome.ts +++ b/cloud/apps/push/src/push-provider-outcome.ts @@ -3,4 +3,4 @@ export type PushProviderOutcome = | { status: 'sent' } | { status: 'dead'; reason: string } - | { status: 'error'; reason: string } + | { status: 'error'; reason: string; retryable?: boolean; retryAfterMs?: number } diff --git a/cloud/apps/push/src/push-request-drain.ts b/cloud/apps/push/src/push-request-drain.ts new file mode 100644 index 00000000000..4acaf09ca67 --- /dev/null +++ b/cloud/apps/push/src/push-request-drain.ts @@ -0,0 +1,28 @@ +import type { MiddlewareHandler } from 'hono' + +export class PushRequestDrain { + private draining = false + private active = 0 + private readonly waiters = new Set<() => void>() + + readonly middleware: MiddlewareHandler = async (context, next) => { + if (this.draining) return context.json({ error: 'shutting_down' }, 503) + this.active++ + try { + await next() + } finally { + this.active-- + if (this.active === 0) { + for (const resolve of this.waiters) resolve() + this.waiters.clear() + } + } + } + + begin(): Promise { + this.draining = true + return this.active === 0 + ? Promise.resolve() + : new Promise((resolve) => this.waiters.add(resolve)) + } +} diff --git a/cloud/apps/push/src/push-send-idempotency.test.ts b/cloud/apps/push/src/push-send-idempotency.test.ts new file mode 100644 index 00000000000..ec79512f70e --- /dev/null +++ b/cloud/apps/push/src/push-send-idempotency.test.ts @@ -0,0 +1,34 @@ +import { afterEach, expect, it } from 'vitest' +import { createPushServerHarness, notification } from './push-server-harness.test-fixture.js' +import { createPushHostKeypair } from './host-challenge-answering.test-fixture.js' +const harnesses: Awaited>[] = [] +afterEach(async () => { + await Promise.all(harnesses.splice(0).map((h) => h.close())) +}) + +it('returns queued for concurrent retries without double quota or a false summary', async () => { + const h = await createPushServerHarness() + harnesses.push(h) + const token = await h.signIn(createPushHostKeypair(2)) + const registrationId = await h.registerAndroid(token) + const body = { v: 1, registrationIds: [registrationId], notification: notification() } + const responses = await Promise.all( + Array.from({ length: 10 }, () => h.post('/v1/send', body, token)) + ) + for (const response of responses) + expect(await response.json()).toEqual({ results: [{ registrationId, status: 'queued' }] }) + expect(h.server.coalescer.pendingCount(registrationId)).toBe(1) + await h.server.coalescer.flushAll() + await h.post('/v1/send', body, token) + await h.server.coalescer.flushAll() + expect(h.fcmRequests).toHaveLength(1) + expect(JSON.parse(h.fcmRequests[0]!.body).message.data.coalescedCount).toBe('1') + expect((await h.database.query('SELECT COUNT(*) AS count FROM push_send_log'))[0]?.count).toBe(1) + await h.post( + '/v1/send', + { ...body, notification: notification({ notificationEpoch: 'new-epoch' }) }, + token + ) + await h.server.coalescer.flushAll() + expect(h.fcmRequests).toHaveLength(2) +}) diff --git a/cloud/apps/push/src/push-server-harness.test-fixture.ts b/cloud/apps/push/src/push-server-harness.test-fixture.ts index 60d8a6102ad..4b955fcf68a 100644 --- a/cloud/apps/push/src/push-server-harness.test-fixture.ts +++ b/cloud/apps/push/src/push-server-harness.test-fixture.ts @@ -67,6 +67,7 @@ export async function createPushServerHarness() { let fcmResponse: FcmResponse = { status: 200, body: '{}' } const server = createPushServer(testPushConfig(), database, { now: () => clock, + providerRetryWait: async () => undefined, apnsTransport: async (request) => { apnsRequests.push(request) return apnsResponse diff --git a/cloud/apps/push/src/push-server.ts b/cloud/apps/push/src/push-server.ts index 5bc58337476..1201748095b 100644 --- a/cloud/apps/push/src/push-server.ts +++ b/cloud/apps/push/src/push-server.ts @@ -23,10 +23,12 @@ import type { PushDatabase } from './push-database.js' import { PushDispatcher } from './push-dispatcher.js' import { PushObservability } from './push-observability.js' import { createPushReadiness } from './push-readiness.js' +import { PushRequestDrain } from './push-request-drain.js' import { PushSendQuota } from './send-quota.js' export type PushServerOptions = { now?: () => number + providerRetryWait?: (ms: number) => Promise apnsTransport?: ApnsTransport fcmTransport?: FcmTransport fcmAccessToken?: () => Promise @@ -69,6 +71,9 @@ export function createPushServer( const apnsTransport = options.apnsTransport ?? (config.apns ? createApnsHttp2Transport() : null) const dispatcher = new PushDispatcher({ devices, + now, + ...(options.providerRetryWait ? { wait: options.providerRetryWait } : {}), + onRetry: () => observability.record('delivery_retry'), ...(config.apns && apnsTransport ? { apns: new ApnsClient({ @@ -114,6 +119,8 @@ export function createPushServer( onLimited: () => observability.record('ip_rate_limited') }) const app = new Hono<{ Variables: PushVariables }>() + const requestDrain = new PushRequestDrain() + app.use('*', requestDrain.middleware) // Hono's default handler prints the whole error, and a pg error carries the // offending row in `detail`. Only the error's name may reach the logs. app.onError((error, context) => { @@ -129,7 +136,9 @@ export function createPushServer( app.get('/health', (context) => context.json({ ok: true, pushProtocol: 1 })) app.get('/ready', async (context) => - (await ready()) ? context.json({ ok: true }) : context.json({ error: 'dependency_unavailable' }, 503) + (await ready()) + ? context.json({ ok: true }) + : context.json({ error: 'dependency_unavailable' }, 503) ) const bearerSession: MiddlewareHandler<{ Variables: PushVariables }> = async (context, next) => { @@ -173,7 +182,9 @@ export function createPushServer( if (!verification.ok) { observability.record('session_rejected') return context.json( - { error: verification.reason === 'unknown_challenge' ? 'invalid_challenge' : 'invalid_proof' }, + { + error: verification.reason === 'unknown_challenge' ? 'invalid_challenge' : 'invalid_proof' + }, 401 ) } @@ -236,7 +247,16 @@ export function createPushServer( results.push({ registrationId, status: 'dead' }) continue } - if ((await quota.reserve(hostFingerprint, registrationId)) === 'rate_limited') { + const reservation = await quota.reserve( + hostFingerprint, + registrationId, + body.data.notification + ) + if (reservation === 'duplicate') { + results.push({ registrationId, status: 'queued' }) + continue + } + if (reservation === 'rate_limited') { observability.record('send_rate_limited') results.push({ registrationId, status: 'rate_limited' }) continue @@ -250,6 +270,7 @@ export function createPushServer( return { app, + requestDrain, server: createAdaptorServer(app), challenges, sessions, @@ -261,7 +282,7 @@ export function createPushServer( ready, closeTransports: (): void => { if (apnsTransport && 'close' in apnsTransport) { - (apnsTransport as { close: () => void }).close() + ;(apnsTransport as { close: () => void }).close() } } } diff --git a/cloud/apps/push/src/push-session-concurrency.test.ts b/cloud/apps/push/src/push-session-concurrency.test.ts new file mode 100644 index 00000000000..a43daf0f07b --- /dev/null +++ b/cloud/apps/push/src/push-session-concurrency.test.ts @@ -0,0 +1,73 @@ +import { randomUUID } from 'node:crypto' +import { tmpdir } from 'node:os' +import { afterEach, describe, expect, it } from 'vitest' +import { openInMemoryPushDatabase, openPushDatabase, type PushDatabase } from './push-database.js' +import { PushHostSessionStore } from './host-session-store.js' +import { ensurePushSessionIndex } from './push-session-schema.js' +const databases: PushDatabase[] = [] +afterEach(async () => { + await Promise.all(databases.splice(0).map((db) => db.close())) +}) + +async function concurrentSessions(db: PushDatabase) { + databases.push(db) + const host = randomUUID() + const store = new PushHostSessionStore(db) + try { + const sessions = await Promise.all(Array.from({ length: 20 }, () => store.create(host))) + const decisions = await Promise.all( + sessions.map((session) => store.resolve(session.sessionToken)) + ) + expect(decisions.filter((decision) => decision.ok)).toHaveLength(1) + const [row] = await db.query( + 'SELECT COUNT(*) AS count FROM push_sessions WHERE host_fingerprint = ?', + [host] + ) + expect(Number(row?.count)).toBe(1) + } finally { + await db.query('DELETE FROM push_sessions WHERE host_fingerprint = ?', [host]) + } +} +it('serializes sessions on SQLite', async () => { + await concurrentSessions(await openInMemoryPushDatabase()) +}) + +it('migrates existing duplicate hosts to the newest session and enforces uniqueness', async () => { + const db = await openInMemoryPushDatabase() + databases.push(db) + await db.query('DROP INDEX push_sessions_host') + for (const [token, created] of [ + ['old', 1], + ['new', 2] + ] as const) { + await db.query('INSERT INTO push_sessions VALUES (?, ?, ?, ?)', [token, 'host', 100, created]) + } + await ensurePushSessionIndex(db) + expect(await db.query('SELECT token_hash FROM push_sessions')).toEqual([{ token_hash: 'new' }]) + await expect( + db.query('INSERT INTO push_sessions VALUES (?, ?, ?, ?)', ['third', 'host', 100, 3]) + ).rejects.toThrow() +}) + +describe.skipIf(!process.env.ORCA_PUSH_TEST_DATABASE_URL)('PostgreSQL push sessions', () => { + it('leaves exactly one live token after concurrent creates', async () => { + await concurrentSessions( + await openPushDatabase({ + databaseUrl: process.env.ORCA_PUSH_TEST_DATABASE_URL!, + dataDir: tmpdir() + }) + ) + }) + it('allows concurrent schema startup', async () => { + const opened = await Promise.all( + Array.from({ length: 4 }, () => + openPushDatabase({ + databaseUrl: process.env.ORCA_PUSH_TEST_DATABASE_URL!, + dataDir: tmpdir() + }) + ) + ) + databases.push(...opened) + for (const db of opened) expect(await db.query('SELECT 1 AS ok')).toEqual([{ ok: 1 }]) + }) +}) diff --git a/cloud/apps/push/src/push-session-schema.ts b/cloud/apps/push/src/push-session-schema.ts new file mode 100644 index 00000000000..aeb690ce048 --- /dev/null +++ b/cloud/apps/push/src/push-session-schema.ts @@ -0,0 +1,23 @@ +import type { PushDatabase } from './push-database.js' + +export async function ensurePushSessionIndex(database: PushDatabase): Promise { + await database.transaction(async (transaction) => { + await transaction.lockQuotaScope('orca-push-session-schema') + const indexQuery = + database.dialect === 'postgres' + ? "SELECT indexname FROM pg_indexes WHERE schemaname = current_schema() AND tablename = 'push_sessions' AND indexname = 'push_sessions_host'" + : "SELECT name FROM sqlite_master WHERE type = 'index' AND name = 'push_sessions_host'" + if ((await transaction.query(indexQuery)).length) return + // Retain the newest session when upgrading a database with duplicate hosts. + await transaction.query(`DELETE FROM push_sessions WHERE token_hash IN ( + SELECT token_hash FROM ( + SELECT token_hash, ROW_NUMBER() OVER ( + PARTITION BY host_fingerprint ORDER BY created_at DESC, token_hash DESC + ) AS position FROM push_sessions + ) AS ranked WHERE position > 1 + )`) + await transaction.query( + 'CREATE UNIQUE INDEX IF NOT EXISTS push_sessions_host ON push_sessions(host_fingerprint)' + ) + }) +} diff --git a/cloud/apps/push/src/send-quota-postgres.test.ts b/cloud/apps/push/src/send-quota-postgres.test.ts index 5630ae371d4..9ccdf176f46 100644 --- a/cloud/apps/push/src/send-quota-postgres.test.ts +++ b/cloud/apps/push/src/send-quota-postgres.test.ts @@ -6,8 +6,7 @@ import { PushDeviceRegistryStore } from './device-registry-store.js' import { openPushDatabase, type PushDatabase } from './push-database.js' import { PushSendQuota } from './send-quota.js' -// CI stays SQLite-only. Point this at a throwaway PostgreSQL to prove the -// advisory lock, because SQLite serializes writers and cannot show the race. +// Cloud Verify supplies a disposable PostgreSQL; SQLite cannot expose these races. const DATABASE_URL = process.env.ORCA_PUSH_TEST_DATABASE_URL const CONCURRENT_RESERVES = 80 @@ -86,4 +85,16 @@ describe.skipIf(!DATABASE_URL)('push send quota on postgres', () => { await database.query('DELETE FROM push_send_log WHERE host_fingerprint = ?', [otherHost]) } }) + it('reserves a retried event once under concurrent PostgreSQL transactions', async () => { + const quota = new PushSendQuota(database) + const event = { notificationEpoch: 'epoch', notificationSeq: 1 } + const results = await Promise.all( + Array.from({ length: 40 }, () => quota.reserve(hostFingerprint, 'reg-dedupe', event)) + ) + expect(results.filter((result) => result === 'allowed')).toHaveLength(1) + expect(results.filter((result) => result === 'duplicate')).toHaveLength(39) + expect( + await quota.reserve(hostFingerprint, 'reg-dedupe', { ...event, notificationEpoch: 'next' }) + ).toBe('allowed') + }) }) diff --git a/cloud/apps/push/src/send-quota.ts b/cloud/apps/push/src/send-quota.ts index 8cf34c3705d..3049cb312b1 100644 --- a/cloud/apps/push/src/send-quota.ts +++ b/cloud/apps/push/src/send-quota.ts @@ -1,4 +1,4 @@ -import { randomUUID } from 'node:crypto' +import { createHash, randomUUID } from 'node:crypto' import { PUSH_LIMITS } from '@orca-cloud/push-contract' import type { PushDatabase } from './push-database.js' @@ -6,7 +6,7 @@ const QUOTA_LOCK_PREFIX = 'orca-push-send-quota:' const ROLLING_HOUR_MS = 60 * 60 * 1000 const ROLLING_DAY_MS = 24 * ROLLING_HOUR_MS -export type PushQuotaDecision = 'allowed' | 'rate_limited' +export type PushQuotaDecision = 'allowed' | 'rate_limited' | 'duplicate' export class PushSendQuota { constructor( @@ -18,10 +18,33 @@ export class PushSendQuota { // COMMITTED, so concurrent reserves would each see the same under-quota count // and all be admitted. The host lock serializes them. The registration count // rides the same lock because a registration belongs to exactly one host. - async reserve(hostFingerprint: string, registrationId: string): Promise { + async reserve( + hostFingerprint: string, + registrationId: string, + event?: { notificationEpoch: string; notificationSeq: number } + ): Promise { const now = this.now() + const sendId = event + ? createHash('sha256') + .update( + JSON.stringify([ + hostFingerprint, + registrationId, + event.notificationEpoch, + event.notificationSeq + ]) + ) + .digest('hex') + : randomUUID() return await this.database.transaction(async (transaction) => { await transaction.lockQuotaScope(`${QUOTA_LOCK_PREFIX}${hostFingerprint}`) + if ( + event && + (await transaction.query('SELECT send_id FROM push_send_log WHERE send_id = ?', [sendId])) + .length + ) { + return 'duplicate' + } const [hostRow] = await transaction.query( 'SELECT COUNT(*) AS sends FROM push_send_log WHERE host_fingerprint = ? AND sent_at > ?', [hostFingerprint, now - ROLLING_HOUR_MS] @@ -31,15 +54,13 @@ export class PushSendQuota { 'SELECT COUNT(*) AS sends FROM push_send_log WHERE registration_id = ? AND sent_at > ?', [registrationId, now - ROLLING_DAY_MS] ) - if ( - Number(registrationRow?.sends ?? 0) >= PUSH_LIMITS.registrationSendsPerRollingDay - ) { + if (Number(registrationRow?.sends ?? 0) >= PUSH_LIMITS.registrationSendsPerRollingDay) { return 'rate_limited' } await transaction.query( `INSERT INTO push_send_log (send_id, host_fingerprint, registration_id, sent_at) VALUES (?, ?, ?, ?)`, - [randomUUID(), hostFingerprint, registrationId, now] + [sendId, hostFingerprint, registrationId, now] ) return 'allowed' }) diff --git a/cloud/apps/relay/Dockerfile b/cloud/apps/relay/Dockerfile index 12516cbf749..f0abcf9f5b3 100644 --- a/cloud/apps/relay/Dockerfile +++ b/cloud/apps/relay/Dockerfile @@ -3,11 +3,13 @@ WORKDIR /app RUN corepack enable COPY package.json pnpm-lock.yaml pnpm-workspace.yaml tsconfig.base.json ./ COPY packages/relay-contract/package.json packages/relay-contract/package.json +COPY packages/postgres-schema/package.json packages/postgres-schema/package.json COPY apps/relay/package.json apps/relay/package.json RUN pnpm install --frozen-lockfile COPY packages/relay-contract packages/relay-contract COPY apps/relay apps/relay -RUN pnpm --filter @orca-cloud/relay-contract build && pnpm --filter @orca-cloud/relay build +COPY packages/postgres-schema packages/postgres-schema +RUN pnpm --filter @orca-cloud/postgres-schema build && pnpm --filter @orca-cloud/relay-contract build && pnpm --filter @orca-cloud/relay build FROM node:24-alpine AS runtime ENV NODE_ENV=production @@ -16,8 +18,10 @@ WORKDIR /app RUN corepack enable COPY package.json pnpm-lock.yaml pnpm-workspace.yaml ./ COPY packages/relay-contract/package.json packages/relay-contract/package.json +COPY packages/postgres-schema/package.json packages/postgres-schema/package.json COPY apps/relay/package.json apps/relay/package.json COPY --from=build /app/packages/relay-contract/dist packages/relay-contract/dist +COPY --from=build /app/packages/postgres-schema/dist packages/postgres-schema/dist COPY --from=build /app/apps/relay/dist apps/relay/dist RUN pnpm install --prod --frozen-lockfile --filter @orca-cloud/relay... USER node diff --git a/cloud/apps/relay/package.json b/cloud/apps/relay/package.json index 4c2b2e4269c..ea69572b4f6 100644 --- a/cloud/apps/relay/package.json +++ b/cloud/apps/relay/package.json @@ -9,13 +9,14 @@ "clean": "node -e \"require('fs').rmSync('dist', { recursive: true, force: true })\"", "dev": "tsx watch src/index.ts", "lint": "tsc -p tsconfig.json --noEmit", - "pretest": "pnpm --filter @orca-cloud/relay-contract build", + "pretest": "pnpm --filter @orca-cloud/postgres-schema build && pnpm --filter @orca-cloud/relay-contract build", "start": "node dist/index.js", "test": "vitest run", "typecheck": "tsc -p tsconfig.json --noEmit" }, "dependencies": { "@hono/node-server": "^1.19.14", + "@orca-cloud/postgres-schema": "workspace:*", "@orca-cloud/relay-contract": "workspace:*", "hono": "^4.12.27", "jose": "^6.1.3", diff --git a/cloud/apps/relay/src/postgres-schema-startup.ts b/cloud/apps/relay/src/postgres-schema-startup.ts index ba9efc6a792..3a3428eda32 100644 --- a/cloud/apps/relay/src/postgres-schema-startup.ts +++ b/cloud/apps/relay/src/postgres-schema-startup.ts @@ -1,105 +1 @@ -const RETRYABLE_SCHEMA_CODES = new Set(['55P03', '57014']) -const DEFAULT_RETRY_DEADLINE_MS = 30_000 -const RETRY_BASE_DELAY_MS = 250 -const RETRY_MAX_DELAY_MS = 2_000 - -type SchemaStartupOptions = { - now?: () => number - random?: () => number - retryDeadlineMs?: number - wait?: (delayMs: number) => Promise -} - -function retryDelayMs(attempt: number, random: () => number): number { - const ceiling = Math.min( - RETRY_BASE_DELAY_MS * 2 ** (attempt - 1), - RETRY_MAX_DELAY_MS - ) - return Math.ceil(ceiling * (0.5 + random() * 0.5)) -} - -function wait(delayMs: number): Promise { - return new Promise((resolve) => setTimeout(resolve, delayMs)) -} - -const CREATE_TABLE_IF_NOT_EXISTS = /^\s*CREATE\s+TABLE\s+IF\s+NOT\s+EXISTS\b/i -const CREATE_INDEX_IF_NOT_EXISTS = /^\s*CREATE\s+(?:UNIQUE\s+)?INDEX\s+IF\s+NOT\s+EXISTS\b/i - -// `IF NOT EXISTS` only checks the name before the catalog inserts, so the loser of a concurrent -// CREATE can fail on the catalog unique index (23505) or, when the winner has already committed by -// the time the loser reaches TypeCreate/heap_create_with_catalog, on the name check those routines -// repeat (42710 duplicate type, 42P07 duplicate relation). Each is a no-op on the next attempt. -function concurrentCreateCollision( - value: { code?: unknown; constraint?: unknown }, - statement: string -): boolean { - if (CREATE_TABLE_IF_NOT_EXISTS.test(statement)) { - return ( - (value.code === '23505' && value.constraint === 'pg_type_typname_nsp_index') || - value.code === '42710' || - value.code === '42P07' - ) - } - if (CREATE_INDEX_IF_NOT_EXISTS.test(statement)) { - return ( - (value.code === '23505' && value.constraint === 'pg_class_relname_nsp_index') || - value.code === '42P07' - ) - } - return false -} - -function retryableSchemaError(error: unknown, statement: string): boolean { - const value = error as { code?: unknown; constraint?: unknown } - return ( - RETRYABLE_SCHEMA_CODES.has(String(value.code)) || concurrentCreateCollision(value, statement) - ) -} - -export async function applyPostgresSchema( - statements: string[], - query: (statement: string) => Promise, - options: SchemaStartupOptions = {} -): Promise { - const now = options.now ?? Date.now - const random = options.random ?? Math.random - const pause = options.wait ?? wait - const deadlineAt = now() + (options.retryDeadlineMs ?? DEFAULT_RETRY_DEADLINE_MS) - - for (const statement of statements) { - let attempt = 1 - while (true) { - try { - await query(statement) - break - } catch (error) { - const code = String((error as { code?: unknown }).code) - const remainingMs = deadlineAt - now() - const retryable = retryableSchemaError(error, statement) - if (!retryable || remainingMs <= 0) { - if (retryable) { - console.warn( - JSON.stringify({ - event: 'orca_relay_postgres_schema_retry_exhausted', - code, - attempts: attempt - }) - ) - } - throw error - } - const delayMs = Math.min(remainingMs, retryDelayMs(attempt, random)) - console.warn( - JSON.stringify({ - event: 'orca_relay_postgres_schema_retry', - code, - attempt, - delayMs - }) - ) - await pause(delayMs) - attempt += 1 - } - } - } -} +export { applyPostgresSchema } from '@orca-cloud/postgres-schema' diff --git a/cloud/dev/scripts/push-gateway-recovery.test.mjs b/cloud/dev/scripts/push-gateway-recovery.test.mjs new file mode 100644 index 00000000000..abed4bc6885 --- /dev/null +++ b/cloud/dev/scripts/push-gateway-recovery.test.mjs @@ -0,0 +1,93 @@ +import assert from 'node:assert/strict' +import { mkdtempSync, rmSync } from 'node:fs' +import { tmpdir } from 'node:os' +import { join } from 'node:path' +import { spawnSync } from 'node:child_process' +import test from 'node:test' +import { readRelayWorkflow } from './relay-repository.mjs' + +const workflow = readRelayWorkflow('push-deploy.yml') +function step(name) { + const start = workflow.indexOf(` - name: ${name}\n`) + assert.notEqual(start, -1) + const end = workflow.indexOf('\n - name:', start + 1) + const block = workflow.slice(start, end === -1 ? undefined : end) + return block.slice(block.indexOf(' run: |\n') + ' run: |\n'.length) + .split('\n').filter((line) => line.startsWith(' ')).map((line) => line.slice(10)).join('\n') +} +const candidate = step('Deploy the candidate revision with no traffic') +const shift = step('Shift all traffic to the verified candidate') +const rollback = step('Roll traffic back to the previous revision') +const cleanup = step('Delete the rejected candidate revision') +const env = { SERVICE_NAME: 'push-test', GCP_PROJECT_ID: 'test', GCP_REGION: 'test', + GITHUB_RUN_ID: '123', GITHUB_RUN_ATTEMPT: '1', IMAGE: 'synthetic-image', + CANDIDATE_REVISION: 'push-test-c123-1', ROLLBACK_REVISION: 'push-test-old' } + +function exercise(body) { + const dir = mkdtempSync(join(tmpdir(), 'push-workflow-')) + try { + const run = spawnSync('bash', ['-c', body], { encoding: 'utf8', timeout: 10000, + env: { ...process.env, ...env, GITHUB_ENV: join(dir, 'env'), GITHUB_STEP_SUMMARY: join(dir, 'summary'), + TRACE: join(dir, 'trace'), STATE: join(dir, 'state') } }) + assert.equal(run.status, 0, run.stderr) + } finally { rmSync(dir, { recursive: true, force: true }) } +} + +// Workflow shell behavior is Linux-specific; these tests never call a real cloud CLI. +test('failed candidate discovery retains enough state to remove tag and revision', { skip: process.platform === 'win32' }, () => { + exercise(` + gcloud() { + case "$*" in + 'run deploy '*) echo deployed > "$STATE" ;; + 'run services describe '*) return 1 ;; + *) echo "$*" >> "$TRACE" ;; + esac + } + jq() { return 1; } + ( ${candidate} ) + test "$?" != 0 || exit 1 + source "$GITHUB_ENV" + test "$CANDIDATE_TAG" = c123-1 || exit 1 + test "$CANDIDATE_REVISION" = push-test-c123-1 || exit 1 + ( ${cleanup} ) || exit 1 + grep -q -- '--remove-tags c123-1' "$TRACE" || exit 1 + grep -q 'run revisions delete push-test-c123-1' "$TRACE" || exit 1 + `) +}) + +test('failed post-promotion read retains intent and restores previous traffic', { skip: process.platform === 'win32' }, () => { + exercise(` + gcloud() { + case "$*" in + 'run services update-traffic '*) echo "$*" >> "$TRACE" ;; + 'run services describe '*) return 1 ;; + esac + } + jq() { return 1; } + ( ${shift} ) + test "$?" != 0 || exit 1 + source "$GITHUB_ENV" + test "$TRAFFIC_SHIFT_ATTEMPTED" = true || exit 1 + gcloud() { + case "$*" in + 'run services update-traffic '*) echo "$*" >> "$TRACE" ;; + 'run services describe '*) echo '{}' ;; + esac + } + jq() { echo "$ROLLBACK_REVISION"; } + ( ${rollback} ) || exit 1 + source "$GITHUB_ENV" + test "$TRAFFIC_ROLLED_BACK" = true || exit 1 + grep -q -- '--to-revisions push-test-old=100' "$TRACE" || exit 1 + `) +}) + +test('ambiguous promotion failure also leaves rollback intent', { skip: process.platform === 'win32' }, () => { + exercise(` + gcloud() { return 1; } + ( ${shift} ) + test "$?" != 0 || exit 1 + source "$GITHUB_ENV" + test "$TRAFFIC_SHIFT_ATTEMPTED" = true + `) +}) diff --git a/cloud/dev/scripts/push-gateway-workflow.test.mjs b/cloud/dev/scripts/push-gateway-workflow.test.mjs index 947a9133554..b7c8c7db3fe 100644 --- a/cloud/dev/scripts/push-gateway-workflow.test.mjs +++ b/cloud/dev/scripts/push-gateway-workflow.test.mjs @@ -252,15 +252,15 @@ test('a failure after the shift rolls production back automatically', () => { const shift = workflow.indexOf('- name: Shift all traffic to the verified candidate') assert.ok( workflow.indexOf('echo "TRAFFIC_SHIFTED=true"') > shift, - 'the marker must be set only once the shift has been verified' + 'the success marker follows the shift step' ) const body = workflow.slice( workflow.indexOf('- name: Roll traffic back to the previous revision'), - workflow.indexOf('- name: Delete the candidate revision that never took traffic') + workflow.indexOf('- name: Delete the rejected candidate revision') ) assert.match( body, - /if: \$\{\{ failure\(\) && env\.TRAFFIC_SHIFTED == 'true' \}\}/, + /if: \$\{\{ \(failure\(\) \|\| cancelled\(\)\) && env\.TRAFFIC_SHIFT_ATTEMPTED == 'true' \}\}/, 'the rollback must be conditioned on both failure and the shift marker' ) assert.match(body, /test -n "\$\{ROLLBACK_REVISION:-\}"/) @@ -273,12 +273,12 @@ test('a failure after the shift rolls production back automatically', () => { // tag comes off first, because Cloud Run refuses to delete a revision a traffic target names. test('a failure before the shift deletes the candidate it created', () => { const body = workflow.slice( - workflow.indexOf('- name: Delete the candidate revision that never took traffic'), + workflow.indexOf('- name: Delete the rejected candidate revision'), workflow.indexOf('- name: Drop the candidate traffic tag') ) assert.match( body, - /if: \$\{\{ failure\(\) && env\.TRAFFIC_SHIFTED != 'true' \}\}/, + /env\.TRAFFIC_SHIFT_ATTEMPTED != 'true' \|\| env\.TRAFFIC_ROLLED_BACK == 'true'/, 'the cleanup must be conditioned on both failure and the absence of the shift marker' ) assert.match(body, /test -n "\$\{CANDIDATE_REVISION:-\}" \|\| exit 0/) diff --git a/cloud/docs/push-gateway.md b/cloud/docs/push-gateway.md index b2518093608..f373c7a2bfb 100644 --- a/cloud/docs/push-gateway.md +++ b/cloud/docs/push-gateway.md @@ -313,3 +313,25 @@ push.onorca.dev. CNAME ghs.googlehosted.com. (DNS only, not proxied) `terraform -chdir=infra/terraform output push_dns_record` prints the same three fields. If the record is ever lost, recreate it exactly like that; Cloudflare proxying blocks certificate issuance and breaks Cloud Run host routing. + + +### Recovery and delivery guarantees + +Candidate tags and deterministic revision names are recorded before deployment. Promotion intent is +recorded before changing traffic, so a failed verification or ambiguous mutation result still triggers +rollback. Failed candidates are deleted only before attempted promotion or after verified rollback. +The summary runs even if candidate discovery or traffic verification fails. + +Push uses the relay's schema-startup retry implementation through `@orca-cloud/postgres-schema`. +Session replacement is serialized per host and a unique host index upgrades older databases by +retaining their newest session. Cloud Verify runs push concurrency tests against PostgreSQL. + +Accepted sends deduplicate by host, registration, epoch, and sequence for the quota ledger's 25-hour +retention period. Provider failures retry at most three times within two minutes, respecting provider +retry delays. Queues remain in memory; a crash or the nine-second shutdown deadline can still lose work. +Graceful shutdown first refuses new requests, waits for admitted handlers, and drains pending and active +deliveries before closing transports and SQL. `delivery_retry` counters accompany existing outcomes. + +Notification and worktree IDs allow 2048 characters each, subject to a combined notification JSON +budget of 3000 UTF-8 bytes. This preserves normal long and Unicode paths without exceeding provider +envelope space. No identity is truncated to meet this budget. diff --git a/cloud/package.json b/cloud/package.json index c460b32a0bc..3e33f245527 100644 --- a/cloud/package.json +++ b/cloud/package.json @@ -21,7 +21,7 @@ "load:relay:recovery-gate": "node dev/scripts/run-relay-recovery-wave-gate.mjs", "ops:relay": "pnpm --filter @orca-cloud/relay-ops dev", "pretest": "node --test dev/scripts/capture-terraform-plan-baseline.test.mjs dev/scripts/operate-relay-asia-admission.test.mjs dev/scripts/prepare-relay-asia-director-cells.test.mjs dev/scripts/prepare-relay-asia-topology-input.test.mjs dev/scripts/production-cloud-sql-rollout-lock.test.mjs dev/scripts/read-relay-serving-regional-placement-version.test.mjs dev/scripts/relay-asia-admission-workflow.test.mjs dev/scripts/relay-asia-rollout-evidence.test.mjs dev/scripts/relay-asia-topology-workflow.test.mjs dev/scripts/relay-cloud-sql-connection-budget.test.mjs dev/scripts/relay-load-reader-evidence.test.mjs dev/scripts/relay-staging-deploy-identity.test.mjs dev/scripts/sanitize-relay-asia-admission-result.test.mjs dev/scripts/terraform-root-partition.test.mjs dev/scripts/validate-relay-asia-topology-plan.test.mjs ../.github/actions/cloud-sql-rollout-lease/action-contract.test.mjs ../.github/actions/cloud-sql-rollout-lease/storage-lease.test.mjs", - "test": "pnpm -r test && node --test dev/scripts/classify-relay-production-capacity-director.test.mjs dev/scripts/classify-relay-staging-bootstrap.test.mjs dev/scripts/deploy-relay-blue-green.test.mjs dev/scripts/deploy-relay-gce-candidate.test.mjs dev/scripts/deploy-relay-gce-multi-target.test.mjs dev/scripts/github-smoke-token.test.mjs dev/scripts/infra.test.mjs dev/scripts/operate-relay-regional-rehome.test.mjs dev/scripts/power-staging-relay.test.mjs dev/scripts/prepare-relay-capacity-canary.test.mjs dev/scripts/prepare-relay-production-capacity-canary.test.mjs dev/scripts/probe-relay-legacy-admission.test.mjs dev/scripts/probe-relay-rehome-trust.test.mjs dev/scripts/production-cell-image-digest-consistency.test.mjs dev/scripts/push-gateway-workflow.test.mjs dev/scripts/read-relay-production-capacity-identity.test.mjs dev/scripts/relay-admin-endpoint-retry-workflow.test.mjs dev/scripts/relay-admin-transient-retry.test.mjs dev/scripts/relay-admission-selector.test.mjs dev/scripts/relay-gce-terraform-fence.test.mjs dev/scripts/relay-load-connection-failure.test.mjs dev/scripts/relay-load-control-peer.test.mjs dev/scripts/relay-load-director-capacity-gate.test.mjs dev/scripts/relay-load-model.test.mjs dev/scripts/relay-load-phase-barrier.test.mjs dev/scripts/relay-load-placement-boundary.test.mjs dev/scripts/relay-load-profile.test.mjs dev/scripts/relay-load-rebind-boundary.test.mjs dev/scripts/relay-load-region-behavior.test.mjs dev/scripts/relay-load-request-unit-boundary.test.mjs dev/scripts/relay-load-run-lifecycle.test.mjs dev/scripts/relay-monitor-evidence.test.mjs dev/scripts/relay-production-capacity-wave.test.mjs dev/scripts/relay-production-capacity-workflow.test.mjs dev/scripts/relay-production-identity-boundaries.test.mjs dev/scripts/relay-production-same-cap-wave.test.mjs dev/scripts/relay-public-workflow-contract.test.mjs dev/scripts/relay-recovery-wave-gate.test.mjs dev/scripts/relay-region-observation-evidence.test.mjs dev/scripts/relay-regional-rehome-workflow.test.mjs dev/scripts/relay-rehome-aggregate-evidence.test.mjs dev/scripts/relay-repository.test.mjs dev/scripts/relay-same-cap-script-census.test.mjs dev/scripts/relay-staging-c4-refresh-workflow.test.mjs dev/scripts/relay-staging-capacity-identity.test.mjs dev/scripts/staging-relay-apply-guard.test.mjs dev/scripts/validate-relay-capacity-plan.test.mjs dev/scripts/verify-relay-capacity-transition.test.mjs dev/scripts/verify-relay-legacy-bootstrap.test.mjs dev/scripts/workload-identity-attribute-conditions.test.mjs", + "test": "pnpm -r test && node --test dev/scripts/classify-relay-production-capacity-director.test.mjs dev/scripts/classify-relay-staging-bootstrap.test.mjs dev/scripts/deploy-relay-blue-green.test.mjs dev/scripts/deploy-relay-gce-candidate.test.mjs dev/scripts/deploy-relay-gce-multi-target.test.mjs dev/scripts/github-smoke-token.test.mjs dev/scripts/infra.test.mjs dev/scripts/operate-relay-regional-rehome.test.mjs dev/scripts/power-staging-relay.test.mjs dev/scripts/prepare-relay-capacity-canary.test.mjs dev/scripts/prepare-relay-production-capacity-canary.test.mjs dev/scripts/probe-relay-legacy-admission.test.mjs dev/scripts/probe-relay-rehome-trust.test.mjs dev/scripts/production-cell-image-digest-consistency.test.mjs dev/scripts/push-gateway-workflow.test.mjs dev/scripts/push-gateway-recovery.test.mjs dev/scripts/read-relay-production-capacity-identity.test.mjs dev/scripts/relay-admin-endpoint-retry-workflow.test.mjs dev/scripts/relay-admin-transient-retry.test.mjs dev/scripts/relay-admission-selector.test.mjs dev/scripts/relay-gce-terraform-fence.test.mjs dev/scripts/relay-load-connection-failure.test.mjs dev/scripts/relay-load-control-peer.test.mjs dev/scripts/relay-load-director-capacity-gate.test.mjs dev/scripts/relay-load-model.test.mjs dev/scripts/relay-load-phase-barrier.test.mjs dev/scripts/relay-load-placement-boundary.test.mjs dev/scripts/relay-load-profile.test.mjs dev/scripts/relay-load-rebind-boundary.test.mjs dev/scripts/relay-load-region-behavior.test.mjs dev/scripts/relay-load-request-unit-boundary.test.mjs dev/scripts/relay-load-run-lifecycle.test.mjs dev/scripts/relay-monitor-evidence.test.mjs dev/scripts/relay-production-capacity-wave.test.mjs dev/scripts/relay-production-capacity-workflow.test.mjs dev/scripts/relay-production-identity-boundaries.test.mjs dev/scripts/relay-production-same-cap-wave.test.mjs dev/scripts/relay-public-workflow-contract.test.mjs dev/scripts/relay-recovery-wave-gate.test.mjs dev/scripts/relay-region-observation-evidence.test.mjs dev/scripts/relay-regional-rehome-workflow.test.mjs dev/scripts/relay-rehome-aggregate-evidence.test.mjs dev/scripts/relay-repository.test.mjs dev/scripts/relay-same-cap-script-census.test.mjs dev/scripts/relay-staging-c4-refresh-workflow.test.mjs dev/scripts/relay-staging-capacity-identity.test.mjs dev/scripts/staging-relay-apply-guard.test.mjs dev/scripts/validate-relay-capacity-plan.test.mjs dev/scripts/verify-relay-capacity-transition.test.mjs dev/scripts/verify-relay-legacy-bootstrap.test.mjs dev/scripts/workload-identity-attribute-conditions.test.mjs", "typecheck": "pnpm -r typecheck" }, "devDependencies": { diff --git a/cloud/packages/postgres-schema/package.json b/cloud/packages/postgres-schema/package.json new file mode 100644 index 00000000000..e170973cf2b --- /dev/null +++ b/cloud/packages/postgres-schema/package.json @@ -0,0 +1,20 @@ +{ + "name": "@orca-cloud/postgres-schema", + "version": "0.0.0", + "private": true, + "type": "module", + "main": "dist/index.js", + "types": "dist/index.d.ts", + "scripts": { + "build": "tsc -p tsconfig.build.json", + "clean": "node -e \"require('fs').rmSync('dist', { recursive: true, force: true })\"", + "lint": "tsc -p tsconfig.json --noEmit", + "test": "pnpm build", + "typecheck": "tsc -p tsconfig.json --noEmit" + }, + "devDependencies": { + "@types/node": "^24.10.0", + "typescript": "^5.9.3", + "vitest": "^4.0.8" + } +} diff --git a/cloud/packages/postgres-schema/src/index.ts b/cloud/packages/postgres-schema/src/index.ts new file mode 100644 index 00000000000..10c144b0ad3 --- /dev/null +++ b/cloud/packages/postgres-schema/src/index.ts @@ -0,0 +1,103 @@ +const RETRYABLE_SCHEMA_CODES = new Set(['55P03', '57014']) +const DEFAULT_RETRY_DEADLINE_MS = 30_000 +const RETRY_BASE_DELAY_MS = 250 +const RETRY_MAX_DELAY_MS = 2_000 + +type SchemaStartupOptions = { + eventPrefix?: string + now?: () => number + random?: () => number + retryDeadlineMs?: number + wait?: (delayMs: number) => Promise +} + +function retryDelayMs(attempt: number, random: () => number): number { + const ceiling = Math.min(RETRY_BASE_DELAY_MS * 2 ** (attempt - 1), RETRY_MAX_DELAY_MS) + return Math.ceil(ceiling * (0.5 + random() * 0.5)) +} + +function wait(delayMs: number): Promise { + return new Promise((resolve) => setTimeout(resolve, delayMs)) +} + +const CREATE_TABLE_IF_NOT_EXISTS = /^\s*CREATE\s+TABLE\s+IF\s+NOT\s+EXISTS\b/i +const CREATE_INDEX_IF_NOT_EXISTS = /^\s*CREATE\s+(?:UNIQUE\s+)?INDEX\s+IF\s+NOT\s+EXISTS\b/i + +// `IF NOT EXISTS` only checks the name before the catalog inserts, so the loser of a concurrent +// CREATE can fail on the catalog unique index (23505) or, when the winner has already committed by +// the time the loser reaches TypeCreate/heap_create_with_catalog, on the name check those routines +// repeat (42710 duplicate type, 42P07 duplicate relation). Each is a no-op on the next attempt. +function concurrentCreateCollision( + value: { code?: unknown; constraint?: unknown }, + statement: string +): boolean { + if (CREATE_TABLE_IF_NOT_EXISTS.test(statement)) { + return ( + (value.code === '23505' && value.constraint === 'pg_type_typname_nsp_index') || + value.code === '42710' || + value.code === '42P07' + ) + } + if (CREATE_INDEX_IF_NOT_EXISTS.test(statement)) { + return ( + (value.code === '23505' && value.constraint === 'pg_class_relname_nsp_index') || + value.code === '42P07' + ) + } + return false +} + +function retryableSchemaError(error: unknown, statement: string): boolean { + const value = error as { code?: unknown; constraint?: unknown } + return ( + RETRYABLE_SCHEMA_CODES.has(String(value.code)) || concurrentCreateCollision(value, statement) + ) +} + +export async function applyPostgresSchema( + statements: string[], + query: (statement: string) => Promise, + options: SchemaStartupOptions = {} +): Promise { + const now = options.now ?? Date.now + const random = options.random ?? Math.random + const pause = options.wait ?? wait + const deadlineAt = now() + (options.retryDeadlineMs ?? DEFAULT_RETRY_DEADLINE_MS) + + for (const statement of statements) { + let attempt = 1 + while (true) { + try { + await query(statement) + break + } catch (error) { + const code = String((error as { code?: unknown }).code) + const remainingMs = deadlineAt - now() + const retryable = retryableSchemaError(error, statement) + if (!retryable || remainingMs <= 0) { + if (retryable) { + console.warn( + JSON.stringify({ + event: `${options.eventPrefix ?? 'orca_relay_postgres_schema'}_retry_exhausted`, + code, + attempts: attempt + }) + ) + } + throw error + } + const delayMs = Math.min(remainingMs, retryDelayMs(attempt, random)) + console.warn( + JSON.stringify({ + event: `${options.eventPrefix ?? 'orca_relay_postgres_schema'}_retry`, + code, + attempt, + delayMs + }) + ) + await pause(delayMs) + attempt += 1 + } + } + } +} diff --git a/cloud/packages/postgres-schema/tsconfig.build.json b/cloud/packages/postgres-schema/tsconfig.build.json new file mode 100644 index 00000000000..94c84b60803 --- /dev/null +++ b/cloud/packages/postgres-schema/tsconfig.build.json @@ -0,0 +1,11 @@ +{ + "extends": "./tsconfig.json", + "compilerOptions": { + "declaration": true, + "emitDeclarationOnly": false, + "noEmit": false, + "outDir": "dist", + "rootDir": "src" + }, + "exclude": ["src/**/*.test.ts"] +} diff --git a/cloud/packages/postgres-schema/tsconfig.json b/cloud/packages/postgres-schema/tsconfig.json new file mode 100644 index 00000000000..a552e34dbe9 --- /dev/null +++ b/cloud/packages/postgres-schema/tsconfig.json @@ -0,0 +1,5 @@ +{ + "extends": "../../tsconfig.base.json", + "compilerOptions": { "noEmit": true }, + "include": ["src/**/*.ts"] +} diff --git a/cloud/packages/push-contract/src/notification-identity-limits.test.ts b/cloud/packages/push-contract/src/notification-identity-limits.test.ts new file mode 100644 index 00000000000..e19fd93140a --- /dev/null +++ b/cloud/packages/push-contract/src/notification-identity-limits.test.ts @@ -0,0 +1,32 @@ +import { expect, it } from 'vitest' +import { PushNotificationSchema } from './send-messages.js' +const base = { + source: 'agent-task-complete', + agentState: 'finished', + notificationSeq: 1, + notificationEpoch: 'epoch', + title: 'Done', + body: '' +} +it.each([ + 'repo::/Users/developer/orca/workspaces/monorepo/packages/desktop/integrations/feature-mobile-background-notifications', + 'repo::C:\\Users\\developer\\Documents\\projects\\monorepo\\packages\\desktop\\feature-mobile-notifications', + 'folder::/home/developer/projects/通知/作業ディレクトリ/機能', + 'ssh:host::/home/developer/workspaces/monorepo/packages/desktop/feature-mobile-background-notifications' +])('preserves long desktop identities: %s', (path) => { + const worktreeId = `12345678-1234-1234-1234-123456789012::${path}` + const notificationId = [ + 'agent', + encodeURIComponent(worktreeId), + encodeURIComponent('12345678-1234-1234-1234-123456789012:87654321-4321-4321-4321-210987654321'), + '1780000000123' + ].join(':') + const result = PushNotificationSchema.parse({ ...base, worktreeId, notificationId }) + expect(result.worktreeId).toBe(worktreeId) + expect(result.notificationId).toBe(notificationId) +}) +it('rejects oversized provider data by UTF-8 bytes instead of truncating identities', () => { + expect(PushNotificationSchema.safeParse({ ...base, worktreeId: '界'.repeat(1100) }).success).toBe( + false + ) +}) diff --git a/cloud/packages/push-contract/src/send-messages.ts b/cloud/packages/push-contract/src/send-messages.ts index e0e08786423..63e8fc84c1f 100644 --- a/cloud/packages/push-contract/src/send-messages.ts +++ b/cloud/packages/push-contract/src/send-messages.ts @@ -14,7 +14,7 @@ export const PushNotificationSchema = z notificationId: z .string() .min(1) - .max(256) + .max(2048) .regex(/^[\x20-\x7e]+$/) .optional(), notificationSeq: SequenceSchema, @@ -23,9 +23,15 @@ export const PushNotificationSchema = z agentState: PushAgentStateSchema.nullable(), title: z.string().min(1).max(PUSH_LIMITS.titleMaxChars), body: z.string().max(PUSH_LIMITS.bodyMaxChars), - worktreeId: z.string().min(1).max(256).optional() + worktreeId: z.string().min(1).max(2048).optional() }) .strict() + .refine( + (notification) => new TextEncoder().encode(JSON.stringify(notification)).byteLength <= 3000, + { + message: 'notification exceeds provider payload budget' + } + ) export const PushSendRequestSchema = z .object({ diff --git a/cloud/pnpm-lock.yaml b/cloud/pnpm-lock.yaml index 4d8e70e3260..6011b2f62d5 100644 --- a/cloud/pnpm-lock.yaml +++ b/cloud/pnpm-lock.yaml @@ -26,6 +26,9 @@ importers: '@hono/node-server': specifier: ^1.19.14 version: 1.19.14(hono@4.12.27) + '@orca-cloud/postgres-schema': + specifier: workspace:* + version: link:../../packages/postgres-schema '@orca-cloud/push-contract': specifier: workspace:* version: link:../../packages/push-contract @@ -66,6 +69,9 @@ importers: '@hono/node-server': specifier: ^1.19.14 version: 1.19.14(hono@4.12.27) + '@orca-cloud/postgres-schema': + specifier: workspace:* + version: link:../../packages/postgres-schema '@orca-cloud/relay-contract': specifier: workspace:* version: link:../../packages/relay-contract @@ -157,6 +163,18 @@ importers: specifier: ^4.0.8 version: 4.1.9(@types/node@24.13.2)(vite@8.0.16(@types/node@24.13.2)(esbuild@0.28.1)(tsx@4.22.4)) + packages/postgres-schema: + devDependencies: + '@types/node': + specifier: ^24.10.0 + version: 24.13.2 + typescript: + specifier: ^5.9.3 + version: 5.9.3 + vitest: + specifier: ^4.0.8 + version: 4.1.9(@types/node@24.13.2)(vite@8.0.16(@types/node@24.13.2)(esbuild@0.28.1)(tsx@4.22.4)) + packages/push-contract: dependencies: zod: diff --git a/docs/reference/mobile-push-contract.md b/docs/reference/mobile-push-contract.md index d78e3ff0776..30dbf22b9b3 100644 --- a/docs/reference/mobile-push-contract.md +++ b/docs/reference/mobile-push-contract.md @@ -99,12 +99,12 @@ host restart can re-read it. iOS token is 64 hex chars; Android token is the FCM { "v": 1, "registrationIds": ["", "..."], "notification": { - "notificationId": "", + "notificationId": "", "notificationSeq": , "notificationEpoch": "", "source": "agent-task-complete" | "terminal-bell" | "plugin", "agentState": "needs-input" | "finished" | null, "title": "", "body": "", - "worktreeId": "" } } + "worktreeId": "" } } ``` → 200 ```json @@ -117,6 +117,11 @@ host restart can re-read it. iOS token is 64 hex chars; Android token is the FCM The cap counts the ids as sent; the gateway then dedupes them, so a repeated id spends quota once, yields one result, and counts once toward `coalescedCount`. `results` may therefore be shorter than `registrationIds`, and callers must match a result by its `registrationId`, never by position. +- Notification JSON is limited to 3000 UTF-8 bytes to leave provider envelope space; identities + are preserved exactly, including long filesystem paths. Oversized payloads fail validation. +- Gateway retries are deduplicated by host, registration, notification epoch, and sequence in the + quota ledger for its 25-hour retention window. Duplicates return `queued` without reserving + quota or enqueueing another delivery. - Both quota counters are reserved under a per-host lock held for the whole transaction. PostgreSQL reads at READ COMMITTED, so a concurrent count-then-insert would otherwise admit a whole burst. @@ -151,7 +156,11 @@ summary: title `Orca`, body ` agents need attention` (or ` updates` when n carries the latest event's fields plus `coalescedCount`. Collapse id for a summary is `host:` so a later summary replaces it. The window is held in memory per gateway instance, so with more than one instance a burst can produce up to one summary per instance; accepted -for this release, and the collapse id keeps the phone showing one banner. +for this release, and the collapse id keeps the phone showing one banner. Transient provider errors +retry at most three attempts within two minutes, honoring Retry-After and FCM minimum delays. Permanent failures +are not retried. Unregister/dead-token state is re-read before every attempt. Shutdown stops admission +and drains admitted requests, pending windows, and active deliveries before closing resources; +a nine-second hard deadline remains below Cloud Run's termination grace. Delivery remains in memory. ### Provider payloads @@ -176,8 +185,8 @@ metadata server or `GOOGLE_APPLICATION_CREDENTIALS` locally): - `push_hosts(host_fingerprint pk, host_public_key, created_at, last_seen_at)`, written only on a verified proof and pruned after 1 h of no contact when no `push_devices` row still names the host. Nothing reads it, and any keypair mints a host for free, so it is not allowed to accumulate. -- `push_sessions` holds one row per host: minting a session deletes the host's earlier one, since a - desktop holds a single session and only re-proves once it is gone. +- `push_sessions` holds one row per host, enforced by a unique index and transaction lock. Minting a + session deletes the host's earlier one, since a desktop holds a single session and only re-proves once it is gone. - `push_challenges(challenge_id pk, host_fingerprint, host_public_key, secret_hash, transcript, expires_at, consumed_at)` - `push_sessions(token_hash pk, host_fingerprint, expires_at, created_at)` @@ -216,8 +225,11 @@ Secret Manager names (already exist in `onorca-cloud`): `orca-cloud-push-apns-ke optional field, tolerated by old registries). When the gateway accepted the token but the host could not store it — the device left mobile scope mid-call (`not_mobile`) or the registry write threw (`registration_storage_failed`) — the host queues the gateway delete in the unregister outbox rather - than leaking a registration nothing will ever push to. Phones must treat any `registered: false` as - "retry later", so an unknown reason string is safe to add. + than leaking a registration nothing will ever push to. Registration, unregister, and outbox deletes + are serialized per device; re-registration first settles earlier cleanup. Authentication failure + never drops a durable delete. Stale send responses only clear the exact local registration observed, + while provider dead-token updates match the token/platform/environment that was sent. Phones must + treat any `registered: false` as "retry later", so an unknown reason string is safe to add. - RPC `notifications.unregisterPush` params null → `{ unregistered: boolean }`. Removes the field and enqueues a gateway delete in a durable outbox (`src/main/runtime/push/push-unregister-outbox.ts`, modelled on `relay-revoke-outbox.ts`). Unpair/revoke (`revokeMobileDevice`) enqueues the same. The @@ -236,8 +248,9 @@ Secret Manager names (already exist in `onorca-cloud`): `orca-cloud-push-apns-ke events, maps `agentState` to `needs-input | finished` (blocked/waiting → needs-input, else finished), batches matching registrationIds into `POST /v1/send` requests of at most 20 registrations each (the gateway's per-request cap; extra devices get their own request rather than being dropped), and drops - registrations the gateway reports `dead`. Fire-and-forget with one retry after 2 s per request; never - throws into dispatch. + unchanged registrations the gateway reports `dead`. Failure categories are counted without payload + values and logged at most once per minute (with a final flush on shutdown). Fire-and-forget with + one retry after 2 s per request; never throws into dispatch. - Add `agentState` to `MobileNotificationDispatchEvent` and set it in `src/main/ipc/notifications.ts` from `args.agentState`. Fix `buildAgentTaskCompleteNotificationOptions` so `working|running|busy` never yields "finished" (title says "working" and the dispatcher treats it as not-final, i.e. no push). diff --git a/src/main/runtime/push/desktop-push-service.test.ts b/src/main/runtime/push/desktop-push-service.test.ts index db34230d42b..9177bcbc18f 100644 --- a/src/main/runtime/push/desktop-push-service.test.ts +++ b/src/main/runtime/push/desktop-push-service.test.ts @@ -135,6 +135,7 @@ describe('DesktopPushService', () => { await harness.service.register({ deviceId: harness.deviceId, ...REGISTER_INPUT }) expect(await harness.service.unregister(harness.deviceId)).toEqual({ unregistered: true }) + await harness.service.flushUnregisterOutbox() expect(harness.registry.getDevice(harness.deviceId)?.pushRegistration).toBeUndefined() expect(harness.deletes).toEqual(['reg-1']) expect(harness.outbox.pending()).toEqual([]) diff --git a/src/main/runtime/push/desktop-push-service.ts b/src/main/runtime/push/desktop-push-service.ts index 751aa3c53e4..a459798625a 100644 --- a/src/main/runtime/push/desktop-push-service.ts +++ b/src/main/runtime/push/desktop-push-service.ts @@ -6,6 +6,7 @@ import type { MobilePushRegisterInput, MobilePushRegisterResult } from '../../../shared/mobile-push-contract' +import { runKeyedSerializedOperation } from '../../cli/keyed-promise-queue' import type { DeviceRegistry } from '../device-registry' import type { OrcaRuntimeService } from '../orca-runtime' import type { OrcaRuntimeRpcServer } from '../runtime-rpc' @@ -46,6 +47,7 @@ export class DesktopPushService { private retryArmed = false private retryDelayMs = OUTBOX_RETRY_BASE_MS private stopped = false + private readonly deviceOperations = new Map>() private constructor( options: DesktopPushServiceOptions, @@ -81,6 +83,7 @@ export class DesktopPushService { start(): void { this.stopped = false + this.dispatcher.start() this.runtime.setMobilePushRegistrar(this) this.unsubscribe = this.runtime.onNotificationDispatched((event) => { this.dispatcher.enqueue(event) @@ -95,6 +98,7 @@ export class DesktopPushService { stop(): void { this.stopped = true + this.dispatcher.stop() this.unsubscribe?.() this.unsubscribe = null this.runtimeRpc.setOnPushUnregisterQueued(null) @@ -110,6 +114,24 @@ export class DesktopPushService { if (!this.registerThrottle.allow(input.deviceId)) { return { registered: false, reason: 'throttled' } } + return runKeyedSerializedOperation(this.deviceOperations, input.deviceId, () => + this.registerAfterCleanup(input) + ) + } + + private async registerAfterCleanup( + input: MobilePushRegisterInput + ): Promise { + // A stable gateway ID must not inherit a delete from an earlier registration. + for (const item of this.outbox.pending().filter((entry) => entry.deviceId === input.deviceId)) { + if (!(await this.deleteQueued(item.reqId, item.registrationId))) { + this.scheduleFlushRetry() + return { registered: false, reason: 'gateway_unreachable' } + } + } + if (this.registry.getDevice(input.deviceId)?.scope !== 'mobile' || this.stopped) { + return { registered: false, reason: 'not_mobile' } + } const result = await this.client.registerDevice(input) if (!result.ok) { return { @@ -130,14 +152,19 @@ export class DesktopPushService { } async unregister(deviceId: string): Promise<{ unregistered: boolean }> { + return runKeyedSerializedOperation(this.deviceOperations, deviceId, async () => + this.unregisterCurrent(deviceId) + ) + } + + private unregisterCurrent(deviceId: string): { unregistered: boolean } { const registrationId = this.registry.getDevice(deviceId)?.pushRegistration?.registrationId if (!registrationId) { return { unregistered: false } } - // Why: drop the local registration first. The phone asked to stop being pushed - // to, and that must hold even if the gateway delete has to wait in the outbox. - this.registry.setPushRegistration(deviceId, null) + // Persist cleanup before forgetting its ID; neither write waits on the gateway. this.outbox.enqueue({ registrationId, deviceId }) + this.registry.setPushRegistration(deviceId, null) void this.flushUnregisterOutbox() return { unregistered: true } } @@ -197,10 +224,12 @@ export class DesktopPushService { } attempted.add(item.reqId) try { - const result = await this.client.deleteDevice(item.registrationId) - if (result.deleted || !result.retryable) { - this.outbox.remove(item.reqId) - } else { + const deleted = await runKeyedSerializedOperation( + this.deviceOperations, + item.deviceId, + () => this.deleteQueued(item.reqId, item.registrationId) + ) + if (!deleted) { retryable = true } } catch (error) { @@ -211,6 +240,18 @@ export class DesktopPushService { } } + private async deleteQueued(reqId: string, registrationId: string): Promise { + if (!this.outbox.pending().some((item) => item.reqId === reqId)) { + return true + } + const result = await this.client.deleteDevice(registrationId) + if (!result.deleted) { + return false + } + this.outbox.remove(reqId) + return true + } + private scheduleFlushRetry(): void { if (this.retryArmed || this.stopped) { return diff --git a/src/main/runtime/push/push-agent-state.test.ts b/src/main/runtime/push/push-agent-state.test.ts new file mode 100644 index 00000000000..e56d39ffb01 --- /dev/null +++ b/src/main/runtime/push/push-agent-state.test.ts @@ -0,0 +1,21 @@ +import { describe, expect, it } from 'vitest' +import { mapPushAgentState } from './push-dispatcher' + +describe('mapPushAgentState', () => { + it.each([ + ['blocked', 'needs-input'], + ['waiting', 'needs-input'], + ['done', 'finished'], + [undefined, 'finished'] + ] as const)('maps agent-task-complete %s to %s', (agentState, expected) => { + expect(mapPushAgentState('agent-task-complete', agentState)).toBe(expected) + }) + + it('suppresses a still-working agent', () => { + expect(mapPushAgentState('agent-task-complete', 'working')).toBeUndefined() + }) + + it('leaves non-agent sources without a state', () => { + expect(mapPushAgentState('terminal-bell', undefined)).toBeNull() + }) +}) diff --git a/src/main/runtime/push/push-cleanup-auth-expiry.test.ts b/src/main/runtime/push/push-cleanup-auth-expiry.test.ts new file mode 100644 index 00000000000..b746a03a002 --- /dev/null +++ b/src/main/runtime/push/push-cleanup-auth-expiry.test.ts @@ -0,0 +1,41 @@ +import { createHash } from 'node:crypto' +import { expect, it } from 'vitest' +import { PushGatewayClient } from './push-gateway-client' +import { buildPushChallengeFixture, createPushHostKeypair } from './push-host-challenge-fixtures' + +it('retains a delete when its session proof expires before the DELETE is attempted', async () => { + const keypair = createPushHostKeypair() + const hostFingerprint = createHash('sha256') + .update(keypair.publicKey) + .digest('base64url') + .slice(0, 16) + let now = 1_770_000_000_000 + let deletes = 0 + const client = new PushGatewayClient({ + gatewayUrl: 'https://push.example.test', + keypair, + now: () => now, + fetch: (async (url, init) => { + if (String(url).endsWith('/challenge')) { + const fixture = buildPushChallengeFixture({ + hostKeypair: keypair, + hostFingerprint, + gatewayOrigin: 'https://push.example.test', + issuedAt: now, + challengeId: 'challenge-1' + }) + now += 11_000 + return Response.json(fixture.challenge) + } + if (String(url).endsWith('/session')) { + return Response.json({ error: 'invalid_proof' }, { status: 401 }) + } + if (init?.method === 'DELETE') { + deletes++ + } + return new Response(null, { status: 204 }) + }) as typeof fetch + }) + expect(await client.deleteDevice('registration-1')).toEqual({ deleted: false, retryable: true }) + expect(deletes).toBe(0) +}) diff --git a/src/main/runtime/push/push-dispatcher.test-fixture.ts b/src/main/runtime/push/push-dispatcher.test-fixture.ts new file mode 100644 index 00000000000..9137ed8ea9f --- /dev/null +++ b/src/main/runtime/push/push-dispatcher.test-fixture.ts @@ -0,0 +1,94 @@ +import { vi } from 'vitest' +import type { MobilePushFilter, MobilePushRegistration } from '../../../shared/mobile-push-contract' +import type { MobileNotificationEvent } from '../runtime-mobile-notification-controller' +import type { PushGatewayClient, PushSendResult } from './push-gateway-client' +import { PushDispatcher, type PushDispatcherRegistry } from './push-dispatcher' + +const ALL_SOURCES: MobilePushFilter = { + sources: ['agent-task-complete', 'terminal-bell', 'plugin'], + agentStates: ['needs-input', 'finished'] +} + +export function registration( + overrides: Partial = {} +): MobilePushRegistration { + return { + registrationId: 'reg-1', + platform: 'ios', + filter: ALL_SOURCES, + registeredAt: 1, + ...overrides + } +} + +export type SendCall = Parameters[0] + +export function createHarness(options: { + devices: { deviceId: string; pushRegistration?: MobilePushRegistration }[] + results?: PushSendResult[] + sendImpl?: () => Promise +}): { + dispatcher: PushDispatcher + sends: SendCall[] + cleared: (string | null)[] + runRetry: () => void +} { + const sends: SendCall[] = [] + const cleared: (string | null)[] = [] + let retry: (() => void) | null = null + const client = { + send: vi.fn(async (input: SendCall) => { + sends.push(input) + if (options.sendImpl) { + return await options.sendImpl() + } + return { + ok: true as const, + results: + options.results ?? + input.registrationIds.map((registrationId) => ({ + registrationId, + status: 'queued' as const + })) + } + }) + } as unknown as PushGatewayClient + const registry: PushDispatcherRegistry = { + listDevices: () => options.devices, + setPushRegistration: (deviceId, value) => { + cleared.push(value === null ? deviceId : null) + return true + } + } + return { + dispatcher: new PushDispatcher({ + client, + registry, + scheduleRetry: (run) => { + retry = run + } + }), + sends, + cleared, + runRetry: () => retry?.() + } +} + +export function notification( + overrides: Partial = {} +): MobileNotificationEvent { + return { + type: 'notification', + source: 'agent-task-complete', + title: 'feat/x - Claude finished', + body: 'All done.', + worktreeId: 'repo::wt1', + notificationId: 'agent:one', + notificationSeq: 7, + notificationEpoch: 'epoch-1', + agentState: 'done', + ...overrides + } as MobileNotificationEvent +} + +export const flush = (): Promise => new Promise((resolve) => setImmediate(resolve)) diff --git a/src/main/runtime/push/push-dispatcher.test.ts b/src/main/runtime/push/push-dispatcher.test.ts index 5a30115e364..221383a34b1 100644 --- a/src/main/runtime/push/push-dispatcher.test.ts +++ b/src/main/runtime/push/push-dispatcher.test.ts @@ -1,112 +1,13 @@ import { describe, expect, it, vi } from 'vitest' -import type { MobilePushFilter, MobilePushRegistration } from '../../../shared/mobile-push-contract' -import type { MobileNotificationEvent } from '../runtime-mobile-notification-controller' -import type { PushGatewayClient, PushSendResult } from './push-gateway-client' -import { PushDispatcher, mapPushAgentState, type PushDispatcherRegistry } from './push-dispatcher' - -const ALL_SOURCES: MobilePushFilter = { - sources: ['agent-task-complete', 'terminal-bell', 'plugin'], - agentStates: ['needs-input', 'finished'] -} - -function registration(overrides: Partial = {}): MobilePushRegistration { - return { - registrationId: 'reg-1', - platform: 'ios', - filter: ALL_SOURCES, - registeredAt: 1, - ...overrides - } -} - -type SendCall = Parameters[0] - -function createHarness(options: { - devices: { deviceId: string; pushRegistration?: MobilePushRegistration }[] - results?: PushSendResult[] - sendImpl?: () => Promise -}): { - dispatcher: PushDispatcher - sends: SendCall[] - cleared: (string | null)[] - runRetry: () => void -} { - const sends: SendCall[] = [] - const cleared: (string | null)[] = [] - let retry: (() => void) | null = null - const client = { - send: vi.fn(async (input: SendCall) => { - sends.push(input) - if (options.sendImpl) { - return await options.sendImpl() - } - return { - ok: true as const, - results: - options.results ?? - input.registrationIds.map((registrationId) => ({ - registrationId, - status: 'queued' as const - })) - } - }) - } as unknown as PushGatewayClient - const registry: PushDispatcherRegistry = { - listDevices: () => options.devices, - setPushRegistration: (deviceId, value) => { - cleared.push(value === null ? deviceId : null) - return true - } - } - return { - dispatcher: new PushDispatcher({ - client, - registry, - scheduleRetry: (run) => { - retry = run - } - }), - sends, - cleared, - runRetry: () => retry?.() - } -} - -function notification(overrides: Partial = {}): MobileNotificationEvent { - return { - type: 'notification', - source: 'agent-task-complete', - title: 'feat/x - Claude finished', - body: 'All done.', - worktreeId: 'repo::wt1', - notificationId: 'agent:one', - notificationSeq: 7, - notificationEpoch: 'epoch-1', - agentState: 'done', - ...overrides - } as MobileNotificationEvent -} - -const flush = (): Promise => new Promise((resolve) => setImmediate(resolve)) - -describe('mapPushAgentState', () => { - it.each([ - ['blocked', 'needs-input'], - ['waiting', 'needs-input'], - ['done', 'finished'], - [undefined, 'finished'] - ] as const)('maps agent-task-complete %s to %s', (agentState, expected) => { - expect(mapPushAgentState('agent-task-complete', agentState)).toBe(expected) - }) - - it('suppresses a still-working agent', () => { - expect(mapPushAgentState('agent-task-complete', 'working')).toBeUndefined() - }) - - it('leaves non-agent sources without a state', () => { - expect(mapPushAgentState('terminal-bell', undefined)).toBeNull() - }) -}) +import type { PushGatewayClient } from './push-gateway-client' +import { PushDispatcher } from './push-dispatcher' +import { + createHarness, + flush, + notification, + registration, + type SendCall +} from './push-dispatcher.test-fixture' describe('PushDispatcher', () => { it('batches every matching registration into one send', async () => { @@ -270,10 +171,11 @@ describe('PushDispatcher', () => { }) } as unknown as PushGatewayClient const scheduled: (() => void)[] = [] + const devices = [{ deviceId: 'a', pushRegistration: registration() }] const dispatcher = new PushDispatcher({ client, registry: { - listDevices: () => [{ deviceId: 'a', pushRegistration: registration() }], + listDevices: () => devices, setPushRegistration: () => true }, scheduleRetry: (run, delayMs) => { diff --git a/src/main/runtime/push/push-dispatcher.ts b/src/main/runtime/push/push-dispatcher.ts index f3f9d9af78c..d6a83f3241a 100644 --- a/src/main/runtime/push/push-dispatcher.ts +++ b/src/main/runtime/push/push-dispatcher.ts @@ -6,6 +6,7 @@ import type { MobilePushAgentState, MobilePushRegistration } from '../../../shared/mobile-push-contract' +import { PushOutcomeCounters } from './push-outcome-counters' import { MOBILE_PUSH_SOURCES } from '../../../shared/mobile-push-contract' import type { MobileNotificationEvent } from '../runtime-mobile-notification-controller' import type { PushGatewayClient, PushSendNotification } from './push-gateway-client' @@ -29,7 +30,7 @@ type PushDispatcherOptions = { scheduleRetry?: (run: () => void, delayMs: number) => void } -type PushTarget = { deviceId: string; registrationId: string } +type PushTarget = { deviceId: string; registrationId: string; registration: MobilePushRegistration } function clip(value: string, maxLength: number): string { const normalized = value.replace(/\s+/g, ' ').trim() @@ -58,6 +59,8 @@ export function mapPushAgentState( } export class PushDispatcher { + private readonly outcomes = new PushOutcomeCounters() + private stopped = false private readonly client: PushGatewayClient private readonly registry: PushDispatcherRegistry private readonly scheduleRetry: (run: () => void, delayMs: number) => void @@ -73,7 +76,19 @@ export class PushDispatcher { }) } + start(): void { + this.stopped = false + } + + stop(): void { + this.stopped = true + this.outcomes.flush() + } + enqueue(event: MobileNotificationEvent): void { + if (this.stopped) { + return + } try { const plan = this.planSend(event) if (!plan) { @@ -112,7 +127,9 @@ export class PushDispatcher { if (agentState !== null && !registration.filter.agentStates.includes(agentState)) { return [] } - return [{ deviceId: device.deviceId, registrationId: registration.registrationId }] + return [ + { deviceId: device.deviceId, registrationId: registration.registrationId, registration } + ] }) if (targets.length === 0) { return null @@ -137,15 +154,38 @@ export class PushDispatcher { notification: PushSendNotification, attempt: number ): Promise { + if (this.stopped) { + return + } + const currentTargets = targets.filter((target) => + this.registry + .listDevices() + .some( + (device) => + device.deviceId === target.deviceId && device.pushRegistration === target.registration + ) + ) + if (!currentTargets.length) { + return + } try { const result = await this.client.send({ - registrationIds: targets.map((target) => target.registrationId), + registrationIds: currentTargets.map((target) => target.registrationId), notification }) + if (this.stopped) { + return + } if (result.ok) { + for (const entry of result.results) { + if (entry.status === 'error' || entry.status === 'rate_limited') { + this.outcomes.record(entry.status) + } + } this.dropDeadRegistrations(targets, result.results) return } + this.outcomes.record(result.reason) // Only a transport-level miss is worth repeating; a gateway that refused // this payload will refuse the identical retry. if (attempt === 0 && result.reason === 'unreachable') { @@ -167,7 +207,11 @@ export class PushDispatcher { continue } const target = targets.find((entry) => entry.registrationId === result.registrationId) - if (!target) { + if ( + !target || + this.registry.listDevices().find((device) => device.deviceId === target.deviceId) + ?.pushRegistration !== target.registration + ) { continue } try { diff --git a/src/main/runtime/push/push-gateway-client.ts b/src/main/runtime/push/push-gateway-client.ts index 916bca61590..9d2a09ef906 100644 --- a/src/main/runtime/push/push-gateway-client.ts +++ b/src/main/runtime/push/push-gateway-client.ts @@ -103,7 +103,7 @@ export class PushGatewayClient { method: 'DELETE' }) if (!response.ok) { - return { deleted: false, retryable: response.reason === 'unreachable' } + return { deleted: false, retryable: true } } await cancelUnreadResponseBody(response.response) // A gateway that no longer knows the registration is as deleted as it gets. diff --git a/src/main/runtime/push/push-outcome-counters.test.ts b/src/main/runtime/push/push-outcome-counters.test.ts new file mode 100644 index 00000000000..67ccc475cfc --- /dev/null +++ b/src/main/runtime/push/push-outcome-counters.test.ts @@ -0,0 +1,25 @@ +import { expect, it, vi } from 'vitest' +import { PushOutcomeCounters } from './push-outcome-counters' +it('limits failure logs while retaining category counts', () => { + let now = 0 + const log = vi.spyOn(console, 'warn').mockImplementation(() => {}) + try { + const counters = new PushOutcomeCounters(() => now) + counters.record('rejected') + counters.record('error') + counters.record('error') + expect(log).toHaveBeenCalledTimes(1) + now += 60_000 + counters.record('rate_limited') + expect(JSON.parse(String(log.mock.calls[1]![0]))).toEqual({ + event: 'orca_desktop_push_failures', + error: 2, + rate_limited: 1 + }) + counters.record('unreachable') + counters.flush() + expect(log).toHaveBeenCalledTimes(3) + } finally { + log.mockRestore() + } +}) diff --git a/src/main/runtime/push/push-outcome-counters.ts b/src/main/runtime/push/push-outcome-counters.ts new file mode 100644 index 00000000000..6b2507e5a18 --- /dev/null +++ b/src/main/runtime/push/push-outcome-counters.ts @@ -0,0 +1,27 @@ +type PushOutcome = 'error' | 'rate_limited' | 'rejected' | 'unreachable' + +export class PushOutcomeCounters { + private readonly counts = new Map() + private nextLogAt = 0 + + constructor(private readonly now: () => number = Date.now) {} + + record(outcome: PushOutcome): void { + this.counts.set(outcome, (this.counts.get(outcome) ?? 0) + 1) + if (this.now() < this.nextLogAt) { + return + } + this.nextLogAt = this.now() + 60_000 + this.flush() + } + + flush(): void { + if (!this.counts.size) { + return + } + console.warn( + JSON.stringify({ event: 'orca_desktop_push_failures', ...Object.fromEntries(this.counts) }) + ) + this.counts.clear() + } +} diff --git a/src/main/runtime/push/push-registration-races.test.ts b/src/main/runtime/push/push-registration-races.test.ts new file mode 100644 index 00000000000..afdba983a58 --- /dev/null +++ b/src/main/runtime/push/push-registration-races.test.ts @@ -0,0 +1,160 @@ +import { mkdtempSync, rmSync } from 'node:fs' +import { tmpdir } from 'node:os' +import { join } from 'node:path' +import { afterEach, expect, it, vi } from 'vitest' +import { DeviceRegistry } from '../device-registry' +import { DesktopPushService } from './desktop-push-service' +import { PushUnregisterOutbox } from './push-unregister-outbox' +import { createPushHostKeypair } from './push-host-challenge-fixtures' +import { PushDispatcher } from './push-dispatcher' + +const paths: string[] = [] +afterEach(() => { + for (const path of paths.splice(0)) { + rmSync(path, { recursive: true, force: true }) + } +}) +const input = { + platform: 'android' as const, + token: 'synthetic', + filter: { sources: ['plugin'] as const, agentStates: [] } +} +const tick = () => new Promise((resolve) => setImmediate(resolve)) + +function harness() { + const path = mkdtempSync(join(tmpdir(), 'push-races-')) + paths.push(path) + const registry = new DeviceRegistry(path) + const deviceId = registry.addDevice('phone', 'mobile').deviceId + const outbox = new PushUnregisterOutbox(path) + let live = false + let reachable = true + const client = { + registerDevice: vi.fn(async () => { + live = true + return { ok: true, registrationId: 'stable-id' } + }), + deleteDevice: vi.fn(async () => { + if (!reachable) { + return { deleted: false, retryable: true } + } + live = false + return { deleted: true, retryable: false } + }), + send: vi.fn() + } + const service = DesktopPushService.create({ + gatewayUrl: 'https://push.example.test', + client: client as never, + scheduleRetry: () => {}, + runtime: { + setMobilePushRegistrar: () => {}, + onNotificationDispatched: () => () => {} + } as never, + runtimeRpc: { + getE2EEKeypair: createPushHostKeypair, + getDeviceRegistry: () => registry, + getPushUnregisterOutbox: () => outbox, + setOnPushUnregisterQueued: () => {} + } as never + })! + service.start() + return { + registry, + deviceId, + outbox, + client, + service, + live: () => live, + reachable: (value: boolean) => { + reachable = value + } + } +} + +it('deletes obsolete gateway state before reporting successful re-enable', async () => { + const h = harness() + await h.service.register({ ...input, deviceId: h.deviceId }) + h.reachable(false) + await h.service.unregister(h.deviceId) + await h.service.flushUnregisterOutbox() + expect(h.outbox.pending()).toHaveLength(1) + expect(await h.service.register({ ...input, deviceId: h.deviceId })).toMatchObject({ + registered: false + }) + h.reachable(true) + expect(await h.service.register({ ...input, deviceId: h.deviceId })).toMatchObject({ + registered: true + }) + await h.service.flushUnregisterOutbox() + expect(h.live()).toBe(true) + expect(h.outbox.pending()).toEqual([]) +}) + +it('waits for an already-running delete before re-registering', async () => { + const h = harness() + await h.service.register({ ...input, deviceId: h.deviceId }) + let release!: () => void + const normalDelete = h.client.deleteDevice.getMockImplementation()! + h.client.deleteDevice.mockImplementationOnce(async () => { + await new Promise((resolve) => { + release = resolve + }) + return normalDelete() + }) + await h.service.unregister(h.deviceId) + await tick() + const registration = h.service.register({ ...input, deviceId: h.deviceId }) + await tick() + expect(h.client.registerDevice).toHaveBeenCalledTimes(1) + release() + await registration + await h.service.flushUnregisterOutbox() + expect(h.live()).toBe(true) +}) + +it('orders unregister after a register already in flight', async () => { + const h = harness() + let release!: () => void + const normalRegister = h.client.registerDevice.getMockImplementation()! + h.client.registerDevice.mockImplementationOnce(async () => { + await new Promise((resolve) => { + release = resolve + }) + return normalRegister() + }) + const registered = h.service.register({ ...input, deviceId: h.deviceId }) + await tick() + const unregistered = h.service.unregister(h.deviceId) + release() + await Promise.all([registered, unregistered]) + await h.service.flushUnregisterOutbox() + expect(h.registry.getDevice(h.deviceId)?.pushRegistration).toBeUndefined() + expect(h.live()).toBe(false) +}) + +it('does not clear a replacement with the same ID and timestamp after a stale dead response', async () => { + const h = harness() + await h.service.register({ ...input, deviceId: h.deviceId }) + let finish!: (value: unknown) => void + h.client.send.mockImplementation( + () => + new Promise((resolve) => { + finish = resolve + }) + ) + const dispatcher = new PushDispatcher({ registry: h.registry, client: h.client as never }) + dispatcher.enqueue({ + type: 'notification', + source: 'plugin', + title: 'test', + body: '', + notificationEpoch: 'epoch', + notificationSeq: 1 + }) + const original = h.registry.getDevice(h.deviceId)!.pushRegistration! + h.registry.setPushRegistration(h.deviceId, { ...original }) + finish({ ok: true, results: [{ registrationId: 'stable-id', status: 'dead' }] }) + await tick() + expect(h.registry.getDevice(h.deviceId)?.pushRegistration).toEqual(original) +})