refactor(push): centralize consent and unregister drain ownership

This commit is contained in:
Jinwoo-H
2026-09-10 01:11:37 -04:00
parent e3b63ed063
commit 317680a9cc
13 changed files with 429 additions and 74 deletions
+7 -3
View File
@@ -289,8 +289,9 @@ Secret Manager names (already exist in `onorca-cloud`): `orca-cloud-push-apns-ke
- 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
drain re-reads the queue as it goes, so a delete queued mid-drain lands in the same pass, and a pass
that leaves retryable items schedules an unref'd backoff retry (30 s, doubling, capped at 10 min)
drain processes one queue snapshot per pass; every enqueue requests a flush, so the outer loop
takes another snapshot for work arriving during a pass. Retryable failures schedule an unref'd
backoff retry (30 s, doubling, capped at 10 min)
instead of waiting for the next launch.
- Both RPCs added to `runtime-rpc-mobile-method-allowlist.ts`.
- Push client `src/main/runtime/push/push-gateway-client.ts`: challenge/proof/session with token cache,
@@ -326,7 +327,10 @@ Secret Manager names (already exist in `onorca-cloud`): `orca-cloud-push-apns-ke
controls native push registration. Hint: “Get agent alerts even when the app is closed.
Delivered through Orca’s push service and Apple or Google.” Desktop category controls are
authoritative and are not duplicated as phone overrides. Phone sound and viewing controls remain
independent. Consent is stored only in `orca:pushNotificationsEnabled`; a missing preference
independent. Settings, onboarding and permission-based opt-in use one consent mutation function.
It persists consent and cleanup intent, then schedules reconciliation once without waiting for
network completion. A cleanup-intent write failure still schedules reconciliation and reaches the caller.
Consent is stored only in `orca:pushNotificationsEnabled`; a missing preference
remains off, and obsolete test-build push keys do not grant consent. **Only when away from desktop**
defaults on (180 seconds of OS input idle, or locked). Unknown/headless presence does not suppress; it is never inferred from remote CPU
activity. The detailed payload disclosure remains in the notification documentation.
+2 -2
View File
@@ -22,7 +22,7 @@ import {
saveDefaultSessionView,
type MobileSessionView
} from '../src/storage/session-view-preferences'
import { savePushNotificationsEnabled } from '../src/storage/preferences'
import { setRemotePushEnabled } from '../src/notifications/push-registration'
const SLIDE_DURATION_MS = 280
@@ -127,7 +127,7 @@ function MobileOnboardingFlow({
setError(null)
try {
const enabled = choice === 'enable' ? await ensureNotificationPermissions() : false
await savePushNotificationsEnabled(enabled)
await setRemotePushEnabled(enabled)
advanceOrContinue()
} catch {
setError('Notification settings could not be updated. Try again.')
@@ -1,12 +0,0 @@
const listeners = new Set<() => void>()
export function notifyNotificationConsentChanged(): void {
for (const listener of listeners) {
listener()
}
}
export function subscribeNotificationConsent(listener: () => void): () => void {
listeners.add(listener)
return () => {
listeners.delete(listener)
}
}
@@ -0,0 +1,267 @@
import { createElement } from 'react'
import { act, create, type ReactTestRenderer } from 'react-test-renderer'
import { afterEach, beforeEach, expect, it, vi } from 'vitest'
import AsyncStorage from '@react-native-async-storage/async-storage'
import NotificationsScreen from '../../app/notifications'
import MobileOnboardingScreen from '../../app/mobile-onboarding'
import { shouldPresentNotificationOptIn } from './notification-opt-in-gate'
import {
attachPushRegistration,
NOTIFICATIONS_REMOTE_PUSH_CAPABILITY,
resetPushRegistrationForTests,
setRemotePushEnabled,
startPushTokenSync
} from './push-registration'
import { getDevicePushToken } from './push-token'
const mocks = vi.hoisted(() => ({ storage: new Map<string, string>(), replace: vi.fn() }))
vi.mock('@react-native-async-storage/async-storage', () => ({
default: {
getItem: vi.fn(async (key: string) => mocks.storage.get(key) ?? null),
setItem: vi.fn(async (key: string, value: string) => {
mocks.storage.set(key, value)
})
}
}))
vi.mock('react-native', () => ({
AppState: { currentState: 'active', addEventListener: () => ({ remove: vi.fn() }) },
AccessibilityInfo: {
addEventListener: () => ({ remove: vi.fn() }),
isReduceMotionEnabled: async () => false
},
Animated: { Value: class {}, View: 'View', multiply: () => 0 },
BackHandler: { addEventListener: () => ({ remove: vi.fn() }) },
StyleSheet: { create: (styles: unknown) => styles },
Text: 'Text',
View: 'View',
Switch: 'Switch',
ScrollView: 'ScrollView',
Pressable: 'Pressable',
Alert: { alert: vi.fn() },
Linking: { openSettings: vi.fn() },
useWindowDimensions: () => ({ width: 390, height: 844 })
}))
vi.mock('expo-router', () => ({
useFocusEffect: vi.fn(),
useLocalSearchParams: () => ({ hostId: 'host', steps: 'notifications' }),
useRouter: () => ({ replace: mocks.replace })
}))
vi.mock('react-native-safe-area-context', () => ({
SafeAreaView: 'View',
useSafeAreaInsets: () => ({ top: 0, bottom: 0 })
}))
vi.mock('lucide-react-native', () => ({ ChevronLeft: 'Icon' }))
vi.mock('../components/OrcaLogo', () => ({ OrcaLogo: 'Logo' }))
vi.mock('../onboarding/MobileOnboardingPage', () => ({ MobileOnboardingPage: 'Page' }))
vi.mock('./NotificationDeliverySection', () => ({ NotificationDeliverySection: 'Delivery' }))
vi.mock('./use-remote-push-capable-hosts', () => ({ useRemotePushCapableHosts: () => [] }))
vi.mock('./mobile-notifications', () => ({
ensureNotificationPermissions: async () => true,
getNotificationPermissionState: async () => ({
granted: true,
status: 'granted',
canAskAgain: true,
authorizationReflectsUserChoice: true
})
}))
vi.mock('./desktop-notification-channel', () => ({
ensureDesktopNotificationChannel: async () => {}
}))
vi.mock('./push-token', () => ({
getDevicePushToken: vi.fn(),
addPushTokenListener: () => () => {}
}))
const token = {
platform: 'ios' as const,
token: 'a'.repeat(64),
apnsEnvironment: 'sandbox' as const
}
let renderer: ReactTestRenderer | undefined
let stopSync: () => void
const records = () => JSON.parse(mocks.storage.get('orca:remotePushHostRegistrations') ?? '{}')
const drain = () => vi.advanceTimersByTimeAsync(0)
function deferred<T>() {
let resolve!: (value: T) => void
const promise = new Promise<T>((done) => {
resolve = done
})
return { promise, resolve }
}
function connection() {
return {
sendRequest: vi.fn(async (method: string): Promise<unknown> => ({
ok: true,
result:
method === 'status.get'
? { capabilities: [NOTIFICATIONS_REMOTE_PUSH_CAPABILITY] }
: { registered: true, unregistered: true }
}))
}
}
async function connectedHost() {
const client = connection()
attachPushRegistration('host', client as never)
await drain()
client.sendRequest.mockClear()
return client
}
async function choose(entry: string) {
if (entry === 'gate') {
await expect(shouldPresentNotificationOptIn()).resolves.toBe(false)
return
}
await act(async () => {
renderer = create(
createElement(entry === 'settings' ? NotificationsScreen : MobileOnboardingScreen)
)
})
await act(async () => {
if (entry === 'settings') {
renderer!.root.findByType('Switch').props.onValueChange(true)
} else {
renderer!.root.findByType('Page').props.onNotificationChoice('enable')
}
})
}
function expectChoiceComplete(entry: string) {
expect(mocks.storage.get('orca:pushNotificationsEnabled')).toBe('true')
if (entry === 'settings') {
expect(renderer!.root.findByType('Switch').props).toMatchObject({
value: true,
disabled: false
})
}
if (entry === 'onboarding') {
expect(mocks.replace).toHaveBeenCalledExactlyOnceWith('/h/host')
}
}
beforeEach(() => {
vi.useFakeTimers()
vi.clearAllMocks()
mocks.storage.clear()
resetPushRegistrationForTests()
vi.mocked(getDevicePushToken).mockResolvedValue(token)
stopSync = startPushTokenSync()
})
afterEach(async () => {
await act(async () => renderer?.unmount())
renderer = undefined
stopSync()
resetPushRegistrationForTests()
vi.useRealTimers()
})
it.each(['settings', 'onboarding', 'gate'])(
'%s schedules exactly one registration with token sync running',
async (entry) => {
const client = await connectedHost()
await choose(entry)
await drain()
expectChoiceComplete(entry)
expect(client.sendRequest.mock.calls.map(([method]) => method)).toEqual([
'notifications.registerPush'
])
expect(records().registeredHostIds).toEqual(['host'])
}
)
it.each(['settings', 'onboarding', 'gate'])(
'%s finishes local consent while native token acquisition is pending',
async (entry) => {
const client = await connectedHost()
const pending = deferred<typeof token>()
vi.mocked(getDevicePushToken).mockReturnValue(pending.promise)
await choose(entry)
await drain()
expectChoiceComplete(entry)
expect(getDevicePushToken).toHaveBeenCalledOnce()
expect(client.sendRequest).not.toHaveBeenCalled()
pending.resolve(token)
await drain()
expect(client.sendRequest.mock.calls.map(([method]) => method)).toEqual([
'notifications.registerPush'
])
}
)
it.each(['settings', 'onboarding', 'gate'])(
'%s finishes local consent while registration RPC is pending',
async (entry) => {
const client = await connectedHost()
const pending = deferred<unknown>()
client.sendRequest.mockImplementationOnce(() => pending.promise)
await choose(entry)
await drain()
expectChoiceComplete(entry)
expect(client.sendRequest).toHaveBeenCalledOnce()
expect(records().registeredHostIds).toEqual([])
pending.resolve({ ok: true, result: { registered: true } })
await drain()
expect(records().registeredHostIds).toEqual(['host'])
expect(client.sendRequest).toHaveBeenCalledOnce()
}
)
it('waits for durable local records and schedules one unregister without waiting for its RPC', async () => {
const client = await connectedHost()
await setRemotePushEnabled(true)
await drain()
client.sendRequest.mockClear()
const write = deferred<void>()
vi.mocked(AsyncStorage.setItem)
.mockImplementationOnce(async (key, value) => {
mocks.storage.set(key, value)
})
.mockImplementationOnce(async (key, value) => {
await write.promise
mocks.storage.set(key, value)
})
const rpc = deferred<unknown>()
client.sendRequest.mockImplementationOnce(() => rpc.promise)
const completed = vi.fn()
const disable = setRemotePushEnabled(false).then(completed)
await drain()
expect(completed).not.toHaveBeenCalled()
expect(client.sendRequest).not.toHaveBeenCalled()
write.resolve()
await disable
await drain()
expect(completed).toHaveBeenCalledOnce()
expect(records().pendingUnregisterHostIds).toEqual(['host'])
expect(client.sendRequest.mock.calls.map(([method]) => method)).toEqual([
'notifications.unregisterPush'
])
rpc.resolve({ ok: true })
await drain()
expect(records()).toEqual({ registeredHostIds: [], pendingUnregisterHostIds: [] })
expect(client.sendRequest).toHaveBeenCalledOnce()
})
it('exposes a failed consent write without scheduling or changing durable consent', async () => {
const client = await connectedHost()
vi.mocked(AsyncStorage.setItem).mockRejectedValueOnce(new Error('consent write failed'))
await expect(setRemotePushEnabled(true)).rejects.toThrow('consent write failed')
await drain()
expect(mocks.storage.has('orca:pushNotificationsEnabled')).toBe(false)
expect(client.sendRequest).not.toHaveBeenCalled()
})
it('exposes a failed records write and still schedules exactly one cleanup', async () => {
const client = await connectedHost()
await setRemotePushEnabled(true)
await drain()
client.sendRequest.mockClear()
vi.mocked(AsyncStorage.setItem)
.mockImplementationOnce(async (key, value) => {
mocks.storage.set(key, value)
})
.mockRejectedValueOnce(new Error('records write failed'))
await expect(setRemotePushEnabled(false)).rejects.toThrow('records write failed')
await drain()
expect(mocks.storage.get('orca:pushNotificationsEnabled')).toBe('false')
expect(client.sendRequest.mock.calls.map(([method]) => method)).toEqual([
'notifications.unregisterPush'
])
expect(records()).toEqual({ registeredHostIds: [], pendingUnregisterHostIds: [] })
})
@@ -1,16 +1,15 @@
import { beforeEach, describe, expect, it, vi } from 'vitest'
import {
readPushNotificationsPreference,
savePushNotificationsEnabled
} from '../storage/preferences'
import { readPushNotificationsPreference } from '../storage/preferences'
import { setRemotePushEnabled } from './push-registration'
import { getNotificationPermissionState } from './mobile-notifications'
import { shouldPresentNotificationOptIn } from './notification-opt-in-gate'
vi.mock('../storage/preferences', () => ({
readPushNotificationsPreference: vi.fn(),
savePushNotificationsEnabled: vi.fn()
readPushNotificationsPreference: vi.fn()
}))
vi.mock('./push-registration', () => ({ setRemotePushEnabled: vi.fn() }))
vi.mock('./mobile-notifications', () => ({
getNotificationPermissionState: vi.fn()
}))
@@ -18,7 +17,7 @@ vi.mock('./mobile-notifications', () => ({
describe('notification opt-in gate', () => {
beforeEach(() => {
vi.mocked(readPushNotificationsPreference).mockReset()
vi.mocked(savePushNotificationsEnabled).mockReset()
vi.mocked(setRemotePushEnabled).mockReset()
vi.mocked(getNotificationPermissionState).mockReset()
})
@@ -32,7 +31,7 @@ describe('notification opt-in gate', () => {
})
await expect(shouldPresentNotificationOptIn()).resolves.toBe(true)
expect(savePushNotificationsEnabled).not.toHaveBeenCalled()
expect(setRemotePushEnabled).not.toHaveBeenCalled()
})
it.each([true, false])('preserves an existing %s mobile preference', async (value) => {
@@ -52,7 +51,7 @@ describe('notification opt-in gate', () => {
})
await expect(shouldPresentNotificationOptIn()).resolves.toBe(false)
expect(savePushNotificationsEnabled).toHaveBeenCalledWith(true)
expect(setRemotePushEnabled).toHaveBeenCalledExactlyOnceWith(true)
})
it('still presents when a pre-Android 13 default grant is not an opt-in decision', async () => {
@@ -65,7 +64,7 @@ describe('notification opt-in gate', () => {
})
await expect(shouldPresentNotificationOptIn()).resolves.toBe(true)
expect(savePushNotificationsEnabled).not.toHaveBeenCalled()
expect(setRemotePushEnabled).not.toHaveBeenCalled()
})
it('skips the gate when iOS has already denied permission', async () => {
@@ -78,7 +77,7 @@ describe('notification opt-in gate', () => {
})
await expect(shouldPresentNotificationOptIn()).resolves.toBe(false)
expect(savePushNotificationsEnabled).toHaveBeenCalledWith(false)
expect(setRemotePushEnabled).toHaveBeenCalledExactlyOnceWith(false)
})
it('does not block startup when storage or permission checks fail', async () => {
@@ -1,7 +1,5 @@
import {
readPushNotificationsPreference,
savePushNotificationsEnabled
} from '../storage/preferences'
import { readPushNotificationsPreference } from '../storage/preferences'
import { setRemotePushEnabled } from './push-registration'
import { getNotificationPermissionState } from './mobile-notifications'
export async function shouldPresentNotificationOptIn(): Promise<boolean> {
@@ -18,13 +16,13 @@ export async function shouldPresentNotificationOptIn(): Promise<boolean> {
}
// Why: an already-authorized device should inherit the useful default
// without seeing an onboarding decision it has effectively made.
await savePushNotificationsEnabled(true)
await setRemotePushEnabled(true)
return false
}
if (permission.status === 'denied' || !permission.canAskAgain) {
// Why: iOS cannot show its authorization prompt again, so a blocking
// onboarding screen would be a dead end; Settings remains the recovery.
await savePushNotificationsEnabled(false)
await setRemotePushEnabled(false)
return false
}
return permission.status === 'undetermined'
@@ -116,6 +116,11 @@ it('does not register with stale consent after the user disables notifications d
const disabled = setRemotePushEnabled(false)
pending.resolve(token)
await disabled
await vi.waitFor(() =>
expect(connection.sendRequest.mock.calls.map(([method]) => method)).toContain(
'notifications.unregisterPush'
)
)
expect(connection.sendRequest.mock.calls.map(([method]) => method)).not.toContain(
'notifications.registerPush'
)
@@ -155,9 +160,12 @@ it('completes disable while native token acquisition remains unresolved, and rej
attachPushRegistration('host', connection as never)
await vi.advanceTimersByTimeAsync(0)
expect(getDevicePushToken).toHaveBeenCalledOnce()
const disabled = setRemotePushEnabled(false)
await setRemotePushEnabled(false)
expect(records().pendingUnregisterHostIds).toEqual(['host'])
expect(connection.sendRequest.mock.calls.map(([method]) => method)).not.toContain(
'notifications.unregisterPush'
)
await vi.advanceTimersByTimeAsync(2_000)
await disabled
expect(storage.get('orca:pushNotificationsEnabled')).toBe('false')
expect(records()).toEqual({ registeredHostIds: [], pendingUnregisterHostIds: [] })
expect(connection.sendRequest.mock.calls.map(([method]) => method)).toContain(
@@ -191,6 +199,7 @@ it('restores registration without reconnect after metadata removal fails, retain
detach()
connection.sendRequest.mockClear()
await setRemotePushEnabled(true)
await new Promise((resolve) => setTimeout(resolve, 0))
expect(connection.sendRequest).not.toHaveBeenCalled()
})
@@ -212,6 +221,7 @@ it('does not revive a connection detached while metadata removal was pending', a
commit.resolve()
await removal
await setRemotePushEnabled(true)
await new Promise((resolve) => setTimeout(resolve, 0))
expect(connection.sendRequest).not.toHaveBeenCalled()
})
@@ -1,5 +1,4 @@
import { ensureDesktopNotificationChannel } from './desktop-notification-channel'
import { subscribeNotificationConsent } from './notification-consent-events'
import { AppState } from 'react-native'
import { startMobilePushLeaseRenewal } from './mobile-push-lease-renewal'
import {
@@ -256,6 +255,7 @@ export function attachPushRegistration(hostId: string, client: PushClient): () =
}
}
// Consent completion covers local persistence; host reconciliation runs in the background.
export async function setRemotePushEnabled(enabled: boolean): Promise<void> {
consentGeneration++
await savePushNotificationsEnabled(enabled)
@@ -270,7 +270,7 @@ export async function setRemotePushEnabled(enabled: boolean): Promise<void> {
current.pending.clear()
})
} finally {
await reconcileAllHosts()
void reconcileAllHosts()
}
}
@@ -306,16 +306,12 @@ export async function unregisterPushForRemovedHost(hostId: string): Promise<() =
/** A rolled token stops delivering, so re-register every connected host at once. */
export function startPushTokenSync(): () => void {
const stopConsent = subscribeNotificationConsent(() => {
void reconcileAllHosts()
})
const stopLease = startMobilePushLeaseRenewal(reconcileAllHosts)
const stopToken = addPushTokenListener((token) => {
tokenPromise = Promise.resolve(token)
void reconcileAllHosts()
})
return () => {
stopConsent()
stopLease()
stopToken()
}
@@ -10,7 +10,7 @@ const mocks = vi.hoisted(() => ({
animatedTiming: vi.fn(),
ensureNotificationPermissions: vi.fn(),
saveDefaultSessionView: vi.fn(),
savePushNotificationsEnabled: vi.fn()
setRemotePushEnabled: vi.fn()
}))
vi.mock('react-native', () => ({
@@ -48,8 +48,8 @@ vi.mock('../notifications/mobile-notifications', () => ({
vi.mock('../storage/session-view-preferences', () => ({
saveDefaultSessionView: mocks.saveDefaultSessionView
}))
vi.mock('../storage/preferences', () => ({
savePushNotificationsEnabled: mocks.savePushNotificationsEnabled
vi.mock('../notifications/push-registration', () => ({
setRemotePushEnabled: mocks.setRemotePushEnabled
}))
describe('MobileOnboardingScreen', () => {
@@ -64,7 +64,7 @@ describe('MobileOnboardingScreen', () => {
})
mocks.ensureNotificationPermissions.mockReset().mockResolvedValue(true)
mocks.saveDefaultSessionView.mockReset().mockResolvedValue(undefined)
mocks.savePushNotificationsEnabled.mockReset().mockResolvedValue(undefined)
mocks.setRemotePushEnabled.mockReset().mockResolvedValue(undefined)
})
afterEach(() => {
@@ -100,7 +100,30 @@ describe('MobileOnboardingScreen', () => {
await act(async () => pages()[1].props.onNotificationChoice('skip'))
expect(mocks.ensureNotificationPermissions).not.toHaveBeenCalled()
expect(mocks.savePushNotificationsEnabled).toHaveBeenCalledWith(false)
expect(mocks.setRemotePushEnabled).toHaveBeenCalledWith(false)
expect(mocks.replace).toHaveBeenCalledWith('/h/paired-host')
})
it.each([true, false])(
'saves permission result %s through the consent owner once',
async (granted) => {
mocks.params = { hostId: 'paired-host', steps: 'notifications' }
mocks.ensureNotificationPermissions.mockResolvedValue(granted)
await renderScreen()
await act(async () => pages()[0].props.onNotificationChoice('enable'))
expect(mocks.setRemotePushEnabled).toHaveBeenCalledExactlyOnceWith(granted)
expect(mocks.replace).toHaveBeenCalledWith('/h/paired-host')
}
)
it('keeps notification consent retryable when its local write fails', async () => {
mocks.params = { hostId: 'paired-host', steps: 'notifications' }
mocks.setRemotePushEnabled.mockRejectedValueOnce(new Error('disk full'))
await renderScreen()
await act(async () => pages()[0].props.onNotificationChoice('enable'))
expect(pages()[0].props.error).toBe('Notification settings could not be updated. Try again.')
expect(mocks.replace).not.toHaveBeenCalled()
await act(async () => pages()[0].props.onNotificationChoice('enable'))
expect(mocks.replace).toHaveBeenCalledWith('/h/paired-host')
})
+5 -13
View File
@@ -1,4 +1,3 @@
import { subscribeNotificationConsent } from '../notifications/notification-consent-events'
import AsyncStorage from '@react-native-async-storage/async-storage'
import { beforeEach, describe, expect, it, vi } from 'vitest'
import {
@@ -306,24 +305,17 @@ describe('push notification preference', () => {
await expect(loadPushNotificationsEnabled()).resolves.toBe(false)
})
it('persists and reloads master consent and notifies listeners after each choice', async () => {
it('persists and reloads master consent', async () => {
const storage = new Map<string, string>()
vi.mocked(AsyncStorage.getItem).mockImplementation(async (key) => storage.get(key) ?? null)
vi.mocked(AsyncStorage.setItem).mockImplementation(async (key, value) => {
storage.set(key, value)
})
const changed = vi.fn()
const unsubscribe = subscribeNotificationConsent(changed)
try {
for (const enabled of [true, false]) {
await savePushNotificationsEnabled(enabled)
await expect(loadPushNotificationsEnabled()).resolves.toBe(enabled)
}
expect([...storage]).toEqual([['orca:pushNotificationsEnabled', 'false']])
expect(changed).toHaveBeenCalledTimes(2)
} finally {
unsubscribe()
for (const enabled of [true, false]) {
await savePushNotificationsEnabled(enabled)
await expect(loadPushNotificationsEnabled()).resolves.toBe(enabled)
}
expect([...storage]).toEqual([['orca:pushNotificationsEnabled', 'false']])
})
})
-2
View File
@@ -1,4 +1,3 @@
import { notifyNotificationConsentChanged } from '../notifications/notification-consent-events'
import AsyncStorage from '@react-native-async-storage/async-storage'
const PINS_PREFIX = 'orca:pins:'
@@ -29,7 +28,6 @@ export async function loadPushNotificationsEnabled(): Promise<boolean> {
export async function savePushNotificationsEnabled(enabled: boolean): Promise<void> {
await AsyncStorage.setItem(NOTIF_KEY, String(enabled))
notifyNotificationConsentChanged()
}
const REMOTE_PUSH_HOST_REGISTRATIONS_KEY = 'orca:remotePushHostRegistrations'
@@ -225,16 +225,9 @@ export class DesktopPushService {
/** Returns true when the pass left behind an item the gateway may still accept. */
private async drainPending(): Promise<boolean> {
const attempted = new Set<string>()
let retryable = false
for (;;) {
// Re-read per item: a snapshot taken at loop entry misses anything queued
// while an await was in flight, and the outbox swaps arrays on every write.
const item = this.outbox.pending().find((candidate) => !attempted.has(candidate.reqId))
if (!item) {
return retryable
}
attempted.add(item.reqId)
// Every enqueue requests a flush; the outer loop owns work added during this pass.
for (const item of this.outbox.pending()) {
try {
const deleted = await runKeyedSerializedOperation(
this.deviceOperations,
@@ -250,6 +243,7 @@ export class DesktopPushService {
retryable = true
}
}
return retryable
}
private async deleteQueued(reqId: string, registrationId: string): Promise<boolean> {
@@ -27,6 +27,7 @@ function harness() {
const registry = new DeviceRegistry(path)
const deviceId = registry.addDevice('phone', 'mobile').deviceId
const outbox = new PushUnregisterOutbox(path)
const retries: { run: () => void; delayMs: number }[] = []
let live = false
let reachable = true
const client = {
@@ -34,7 +35,7 @@ function harness() {
live = true
return { ok: true, registrationId: 'stable-id' }
}),
deleteDevice: vi.fn(async () => {
deleteDevice: vi.fn(async (_registrationId: string) => {
if (!reachable) {
return false
}
@@ -46,7 +47,7 @@ function harness() {
const service = DesktopPushService.create({
gatewayUrl: 'https://push.example.test',
client: client as never,
scheduleRetry: () => {},
scheduleRetry: (run, delayMs) => retries.push({ run, delayMs }),
runtime: {
setMobilePushRegistrar: () => {},
onNotificationDispatched: () => () => {}
@@ -60,6 +61,8 @@ function harness() {
})!
service.start()
return {
path,
retries,
registry,
deviceId,
outbox,
@@ -100,7 +103,7 @@ it('waits for an already-running delete before re-registering', async () => {
await new Promise<void>((resolve) => {
release = resolve
})
return normalDelete()
return normalDelete('stable-id')
})
await h.service.unregister(h.deviceId)
await tick()
@@ -181,3 +184,86 @@ it('drains durable deletion even if clearing the local registration fails', asyn
expect(h.client.deleteDevice).toHaveBeenCalledWith('stable-id')
expect(h.outbox.pending()).toEqual([])
})
it('retries an old failure before mid-drain work, then waits for the armed backoff', async () => {
const h = harness()
await h.service.flushUnregisterOutbox()
const deletes: string[] = []
h.client.deleteDevice.mockImplementation(async (registrationId) => {
deletes.push(registrationId)
if (deletes.length === 1) {
h.outbox.enqueue({ registrationId: 'new', deviceId: 'new-phone' })
void h.service.flushUnregisterOutbox()
}
return registrationId === 'new'
})
h.outbox.enqueue({ registrationId: 'old', deviceId: h.deviceId })
await h.service.flushUnregisterOutbox()
expect(deletes).toEqual(['old', 'old', 'new'])
expect(h.outbox.pending().map((item) => item.registrationId)).toEqual(['old'])
expect(h.retries.map((retry) => retry.delayMs)).toEqual([30_000])
await tick()
expect(deletes).toHaveLength(3)
h.retries[0].run()
await tick()
expect(deletes).toEqual(['old', 'old', 'new', 'old'])
expect(h.retries.map((retry) => retry.delayMs)).toEqual([30_000, 60_000])
h.service.stop()
})
it('skips a snapshot delete consumed by same-device registration cleanup', async () => {
const h = harness()
await h.service.flushUnregisterOutbox()
let release!: () => void
h.client.deleteDevice.mockImplementationOnce(
() =>
new Promise<boolean>((resolve) => {
release = () => resolve(true)
})
)
h.outbox.enqueue({ registrationId: 'blocker', deviceId: 'other-phone' })
h.outbox.enqueue({ registrationId: 'stable-id', deviceId: h.deviceId })
const flush = h.service.flushUnregisterOutbox()
await tick()
expect(await h.service.register({ ...input, deviceId: h.deviceId })).toMatchObject({
registered: true
})
expect(h.live()).toBe(true)
release()
await flush
expect(h.client.deleteDevice.mock.calls).toEqual([['blocker'], ['stable-id']])
expect(h.outbox.pending()).toEqual([])
expect(h.live()).toBe(true)
h.service.stop()
})
it('finishes the current snapshot on stop and leaves later work durable for restart', async () => {
const h = harness()
await h.service.flushUnregisterOutbox()
let release!: () => void
h.client.deleteDevice.mockImplementationOnce(
() =>
new Promise<boolean>((resolve) => {
release = () => resolve(false)
})
)
h.outbox.enqueue({ registrationId: 'blocked', deviceId: h.deviceId })
h.outbox.enqueue({ registrationId: 'in-snapshot', deviceId: 'other-phone' })
const flush = h.service.flushUnregisterOutbox()
await tick()
h.outbox.enqueue({ registrationId: 'late', deviceId: 'late-phone' })
void h.service.flushUnregisterOutbox()
h.service.stop()
release()
await flush
expect(h.client.deleteDevice.mock.calls).toEqual([['blocked'], ['in-snapshot']])
expect(h.retries).toEqual([])
const recovered = new PushUnregisterOutbox(h.path)
expect(recovered.pending().map((item) => item.registrationId)).toEqual(['blocked', 'late'])
await h.service.flushUnregisterOutbox()
expect(h.client.deleteDevice).toHaveBeenCalledTimes(2)
h.service.start()
await h.service.flushUnregisterOutbox()
expect(h.outbox.pending()).toEqual([])
h.service.stop()
})