fix: harden mobile push delivery and deployment recovery

This commit is contained in:
Jinwoo-H
2026-09-06 18:42:05 -04:00
parent 41ace23a51
commit b71ffe09fe
56 changed files with 1439 additions and 308 deletions
+11 -6
View File
@@ -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
+1
View File
@@ -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
+5 -1
View File
@@ -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
+2 -1
View File
@@ -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",
+3 -3
View File
@@ -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 })
})
})
+11 -2
View File
@@ -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 })
}
}
}
+5 -2
View File
@@ -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
}
@@ -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<typeof import('node:http2')>()),
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<typeof vi.fn>
close: ReturnType<typeof vi.fn>
}
> = []
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()
})
+11 -2
View File
@@ -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<string, unknown>) => {
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.
+17 -4
View File
@@ -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<Promise<void>>()
private stopped = false
private readonly windows = new Map<string, PendingWindow>()
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<void> {
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()
}
+10 -3
View File
@@ -169,10 +169,17 @@ export class PushDeviceRegistryStore {
return row ? toRegistration(row) : null
}
async markDead(registrationId: string): Promise<void> {
async markDead(registrationId: string, observed?: PushDeviceRegistration): Promise<void> {
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 ?? ''] : [])
]
)
}
}
+16 -8
View File
@@ -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<string, string> }
message: {
android: { collapse_key: string; notification: { tag: string } }
data: Record<string, string>
}
}
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
})
})
})
+19 -4
View File
@@ -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<FcmResponse>
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)
}
}
}
@@ -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
])
+31 -9
View File
@@ -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<number>, 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<void>((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)
@@ -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
}
@@ -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()
})
+7 -1
View File
@@ -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<void> {
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)
}
@@ -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<void>((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()
})
+23 -1
View File
@@ -9,6 +9,9 @@ export type PushDispatcherOptions = {
devices: PushDeviceRegistryStore
apns?: ApnsClient
fcm?: FcmClient
wait?: (ms: number) => Promise<void>
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<void> {
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
}
}
+3 -1
View File
@@ -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,
+1 -1
View File
@@ -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 }
+28
View File
@@ -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<void> {
this.draining = true
return this.active === 0
? Promise.resolve()
: new Promise((resolve) => this.waiters.add(resolve))
}
}
@@ -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<ReturnType<typeof createPushServerHarness>>[] = []
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)
})
@@ -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
+25 -4
View File
@@ -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<void>
apnsTransport?: ApnsTransport
fcmTransport?: FcmTransport
fcmAccessToken?: () => Promise<string>
@@ -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()
}
}
}
@@ -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 }])
})
})
@@ -0,0 +1,23 @@
import type { PushDatabase } from './push-database.js'
export async function ensurePushSessionIndex(database: PushDatabase): Promise<void> {
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)'
)
})
}
@@ -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')
})
})
+28 -7
View File
@@ -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<PushQuotaDecision> {
async reserve(
hostFingerprint: string,
registrationId: string,
event?: { notificationEpoch: string; notificationSeq: number }
): Promise<PushQuotaDecision> {
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<PushQuotaDecision>(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'
})
+5 -1
View File
@@ -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
+2 -1
View File
@@ -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",
+1 -105
View File
@@ -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<void>
}
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<void> {
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<unknown>,
options: SchemaStartupOptions = {}
): Promise<void> {
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'
@@ -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
`)
})
@@ -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/)
+22
View File
@@ -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.
+1 -1
View File
@@ -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": {
@@ -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"
}
}
+103
View File
@@ -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<void>
}
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<void> {
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<unknown>,
options: SchemaStartupOptions = {}
): Promise<void> {
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
}
}
}
}
@@ -0,0 +1,11 @@
{
"extends": "./tsconfig.json",
"compilerOptions": {
"declaration": true,
"emitDeclarationOnly": false,
"noEmit": false,
"outDir": "dist",
"rootDir": "src"
},
"exclude": ["src/**/*.test.ts"]
}
@@ -0,0 +1,5 @@
{
"extends": "../../tsconfig.base.json",
"compilerOptions": { "noEmit": true },
"include": ["src/**/*.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
)
})
@@ -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({
+18
View File
@@ -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:
+22 -9
View File
@@ -99,12 +99,12 @@ host restart can re-read it. iOS token is 64 hex chars; Android token is the FCM
{ "v": 1,
"registrationIds": ["<id>", "..."],
"notification": {
"notificationId": "<string, may be absent for terminal-bell>",
"notificationId": "<max 2048 chars, may be absent for terminal-bell>",
"notificationSeq": <int>, "notificationEpoch": "<uuid>",
"source": "agent-task-complete" | "terminal-bell" | "plugin",
"agentState": "needs-input" | "finished" | null,
"title": "<max 80 chars>", "body": "<max 180 chars>",
"worktreeId": "<string|absent>" } }
"worktreeId": "<max 2048 chars|absent>" } }
```
→ 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 `<N> agents need attention` (or `<N> updates` when n
carries the latest event's fields plus `coalescedCount`. Collapse id for a summary is
`host:<hostFingerprint>` 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).
@@ -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([])
+48 -7
View File
@@ -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<string, Promise<void>>()
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<MobilePushRegisterResult> {
// 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<boolean> {
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
@@ -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()
})
})
@@ -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)
})
@@ -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> = {}
): MobilePushRegistration {
return {
registrationId: 'reg-1',
platform: 'ios',
filter: ALL_SOURCES,
registeredAt: 1,
...overrides
}
}
export type SendCall = Parameters<PushGatewayClient['send']>[0]
export function createHarness(options: {
devices: { deviceId: string; pushRegistration?: MobilePushRegistration }[]
results?: PushSendResult[]
sendImpl?: () => Promise<never>
}): {
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> = {}
): 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<void> => new Promise((resolve) => setImmediate(resolve))
+11 -109
View File
@@ -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> = {}): MobilePushRegistration {
return {
registrationId: 'reg-1',
platform: 'ios',
filter: ALL_SOURCES,
registeredAt: 1,
...overrides
}
}
type SendCall = Parameters<PushGatewayClient['send']>[0]
function createHarness(options: {
devices: { deviceId: string; pushRegistration?: MobilePushRegistration }[]
results?: PushSendResult[]
sendImpl?: () => Promise<never>
}): {
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> = {}): 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<void> => 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) => {
+48 -4
View File
@@ -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<void> {
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 {
+1 -1
View File
@@ -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.
@@ -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()
}
})
@@ -0,0 +1,27 @@
type PushOutcome = 'error' | 'rate_limited' | 'rejected' | 'unreachable'
export class PushOutcomeCounters {
private readonly counts = new Map<PushOutcome, number>()
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()
}
}
@@ -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<void>((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<void>((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)
})