From 261a7b714ca26e23e68a5acfca0e774376be5ee2 Mon Sep 17 00:00:00 2001 From: Brennan Benson <79079362+brennanb2025@users.noreply.github.com> Date: Tue, 21 Jul 2026 18:28:16 -0700 Subject: [PATCH] fix(mobile-relay): back off relay reconnects to stop cellular connect/disconnect churn (#9460) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * fix(mobile-relay): back off relay reconnects to stop cellular connect/disconnect churn On cellular, the relay path re-dialed instantly on every network flap: a NAT rebind / Wi-Fi<->cellular handoff silently kills the socket, the revival trigger treats it as 'link came back' and calls recoverRelay(), and the relay cell answers the overlapping resume with PEER_DROPPED (4408) or LIMIT_EXCEEDED (4429). The session collapsed every close to a plain 'disconnected' and re-dialed with no delay, so the phone ping-ponged connect/disconnect. The documented recovery contract (mobileRelayRecoveryFor, which prescribes fullJitter backoff) had no callers. - Add RelayReconnectBackoff: full-jitter exponential backoff (250ms floor, 30s ceiling) that debounces re-dials via a cooldown window and wires up mobileRelayRecoveryFor. Reset on a successful migrate and on a genuine background->foreground transition (not on repeat foreground nudges). - Extract the lease-rotation timer into RelayLeaseRotationTimer so the supervisor stays under max-lines (the direct-probe path can't be split out — it shares the operationInFlight mutex with recoverRelay). - Add a deterministic test: repeated network-flap nudges re-dial instantly before the fix and are suppressed by the backoff window after. * fix(mobile-relay): recover drops during direct probes * fix(mobile-relay): recover half-open relay sessions * fix(mobile-relay): preserve direct handshakes * fix(mobile-relay): keep recovery retries bounded * fix(mobile-relay): avoid redundant recovery dials * fix(mobile-relay): keep all retries inside cooldown * fix(mobile-relay): close recovery lifecycle races * fix(mobile-relay): preserve in-progress direct auth * fix(mobile-relay): preserve fatal recovery gates * fix(mobile-relay): preserve backoff across unstable resumes * fix(mobile-relay): close remaining recovery lifecycle gaps * fix(mobile-relay): reset backoff only after stable relay --- .../mobile-endpoint-supervisor-contract.ts | 23 + .../mobile-endpoint-supervisor.test.ts | 671 +++++++++++++++++- .../transport/mobile-endpoint-supervisor.ts | 197 ++--- .../mobile-relay-lease-rotation-timer.ts | 59 ++ .../mobile-relay-reconnect-controller.test.ts | 147 ++++ .../mobile-relay-reconnect-controller.ts | 276 +++++++ 6 files changed, 1278 insertions(+), 95 deletions(-) create mode 100644 mobile/src/transport/mobile-endpoint-supervisor-contract.ts create mode 100644 mobile/src/transport/mobile-relay-lease-rotation-timer.ts create mode 100644 mobile/src/transport/mobile-relay-reconnect-controller.test.ts create mode 100644 mobile/src/transport/mobile-relay-reconnect-controller.ts diff --git a/mobile/src/transport/mobile-endpoint-supervisor-contract.ts b/mobile/src/transport/mobile-endpoint-supervisor-contract.ts new file mode 100644 index 00000000000..d26e0706f70 --- /dev/null +++ b/mobile/src/transport/mobile-endpoint-supervisor-contract.ts @@ -0,0 +1,23 @@ +import type { MobileRelayEndpoint } from '../../../src/shared/mobile-relay-credential-contract' +import type { MobileRelayCredentialBundle } from './mobile-relay-credential-bundle' +import type { MobileRelayRpcSession } from './mobile-relay-rpc-session' +import type { resolveMobileRelayEndpoint } from './mobile-relay-resume-director' +import type { RpcClient } from './rpc-client' +import type { HostProfile } from './types' + +export type MobileEndpointSupervisorDependencies = { + openDirect: (endpoint: string) => RpcClient + openRelay: ( + relay: MobileRelayEndpoint, + credential: { token: string; version: number }, + confirmReqId: string + ) => MobileRelayRpcSession + resolveRelay: typeof resolveMobileRelayEndpoint + readBundle: (hostId: string) => Promise + writeBundle: (bundle: MobileRelayCredentialBundle) => Promise + saveHost: (host: HostProfile) => Promise + now: () => number + randomBytes: (length: number) => Uint8Array + setTimer: typeof setTimeout + clearTimer: typeof clearTimeout +} diff --git a/mobile/src/transport/mobile-endpoint-supervisor.test.ts b/mobile/src/transport/mobile-endpoint-supervisor.test.ts index f7b62584168..0bfe75d3b73 100644 --- a/mobile/src/transport/mobile-endpoint-supervisor.test.ts +++ b/mobile/src/transport/mobile-endpoint-supervisor.test.ts @@ -1,5 +1,6 @@ import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' import type { MobileRelayCredentialBundle } from './mobile-relay-credential-bundle' +import { hashMobileRelayCredential } from './mobile-relay-credential-hash' import { RelayOuterError } from './mobile-relay-e2ee-link' import type { MobileRelayRpcSession } from './mobile-relay-rpc-session' import { @@ -16,7 +17,7 @@ vi.mock('expo-crypto', () => ({ getRandomBytes: (length: number) => new Uint8Arr class FakeSession implements RpcClient { readonly sendRequest = vi.fn( - async (): Promise => ({ + async (_method: string, _params?: unknown): Promise => ({ id: 'rpc-1', ok: true, result: {}, @@ -83,6 +84,7 @@ class FakeLogicalClient extends FakeSession implements StableLogicalRpcClient { } this.path = path this.generation += 1 + this.publishState('connected') }) suspendActiveSession = vi.fn(() => this.publishState('disconnected')) getActivePath = () => this.path @@ -141,6 +143,36 @@ function dependencies( } } +function mockCredentialRotation(logical: FakeLogicalClient): void { + let installResult: Record | null = null + logical.sendRequest.mockImplementation(async (method, params) => { + const request = params as { installReqId?: string; reqId?: string } + if (method === 'pairing.provisionRelay') { + installResult = { + v: 1, + reqId: request.reqId, + authorizationMode: 'authenticated-direct', + currentVersion: 3, + resumeExpiresAt: Date.now() + 300_000, + graceExpiresAt: Date.now() + 60_000 + } + return { id: 'rpc-2', ok: true, result: installResult, _meta: { runtimeId: 'runtime-1' } } + } + return { + id: 'rpc-1', + ok: true, + result: { + v: 1, + relay, + installStatus: installResult + ? { v: 1, reqId: request.installReqId, state: 'committed', result: installResult } + : { v: 1, reqId: request.installReqId, state: 'not-found' } + }, + _meta: { runtimeId: 'runtime-1' } + } + }) +} + describe('mobile endpoint supervisor', () => { beforeEach(() => { vi.useFakeTimers() @@ -172,6 +204,14 @@ describe('mobile endpoint supervisor', () => { const supervisor = new MobileEndpointSupervisor(logical, host, deps) await supervisor.start() + logical.publishState('handshaking') + await vi.advanceTimersByTimeAsync(0) + expect(deps.openRelay).not.toHaveBeenCalled() + + supervisor.setForeground(true) + await vi.advanceTimersByTimeAsync(0) + expect(deps.openRelay).not.toHaveBeenCalled() + logical.publishState('reconnecting') await vi.waitFor(() => expect(logical.getActivePath()).toBe('relay')) @@ -191,6 +231,27 @@ describe('mobile endpoint supervisor', () => { supervisor.stop() }) + it('does not spend a queued relay retry while direct authentication is progressing', async () => { + const logical = new FakeLogicalClient('disconnected', 'lan') + const openRelay = vi.fn(() => new FakeRelaySession('disconnected', new RelayOuterError(4408))) + const deps = dependencies({ + openRelay, + randomBytes: () => new Uint8Array([128, 0]) + }) + const supervisor = new MobileEndpointSupervisor(logical, host, deps) + + await supervisor.start() + expect(openRelay).toHaveBeenCalledOnce() + + logical.publishState('handshaking') + await vi.advanceTimersByTimeAsync(250) + expect(openRelay).toHaveBeenCalledOnce() + + logical.publishState('disconnected') + await vi.waitFor(() => expect(openRelay).toHaveBeenCalledTimes(2)) + supervisor.stop() + }) + it('uses POST resolve for wrong-cell recovery and persists the authoritative target', async () => { const logical = new FakeLogicalClient('disconnected', 'lan') const openRelay = vi @@ -228,6 +289,613 @@ describe('mobile endpoint supervisor', () => { supervisor.stop() }) + it('recovers the relay when it drops during an unavailable direct probe', async () => { + const logical = new FakeLogicalClient('connected', 'relay') + const direct = new FakeSession('connecting') + const openRelay = vi.fn(() => new FakeRelaySession('connected')) + const deps = dependencies({ + openDirect: vi.fn(() => direct), + openRelay + }) + const supervisor = new MobileEndpointSupervisor(logical, host, deps) + await supervisor.start() + + // Start the probe, then drop the active relay while the probe owns the + // operation mutex. The failed probe must hand recovery back to the relay. + await vi.advanceTimersByTimeAsync(15_000) + expect(deps.openDirect).toHaveBeenCalledOnce() + logical.publishState('disconnected') + direct.publishState('disconnected') + + await vi.waitFor(() => expect(openRelay).toHaveBeenCalledOnce()) + supervisor.stop() + }) + + it('replaces a half-open relay on a network nudge, then backs off failed resumes', async () => { + const logical = new FakeLogicalClient('disconnected', 'lan') + const openRelay = vi + .fn() + .mockReturnValueOnce(new FakeRelaySession('connected')) + .mockImplementation(() => new FakeRelaySession('disconnected', new RelayOuterError(4408))) + const deps = dependencies({ + openRelay, + // Keep direct unavailable so relay recovery stays the only path under test. + openDirect: vi.fn(() => new FakeSession('disconnected')), + // Deterministic full jitter: fraction 0.5 → half the backoff window. + randomBytes: () => new Uint8Array([128, 0]) + }) + const supervisor = new MobileEndpointSupervisor(logical, host, deps) + + await supervisor.start() + expect(openRelay).toHaveBeenCalledOnce() + + // The OS reports a network handoff, but the dead relay never published onclose. + supervisor.setForeground(true) + await vi.advanceTimersByTimeAsync(0) + expect(logical.suspendActiveSession).toHaveBeenCalledOnce() + expect(openRelay).toHaveBeenCalledTimes(2) + + // The relay cell rejects the replacement with PEER_DROPPED; more flap nudges + // must share the existing cooldown rather than opening more sockets. + for (let i = 0; i < 5; i++) { + supervisor.setForeground(true) + await vi.advanceTimersByTimeAsync(0) + } + expect(openRelay).toHaveBeenCalledTimes(2) + + // Exactly one retry fires at the 250 ms deterministic backoff boundary. + await vi.advanceTimersByTimeAsync(249) + expect(openRelay).toHaveBeenCalledTimes(2) + await vi.advanceTimersByTimeAsync(1) + expect(openRelay).toHaveBeenCalledTimes(3) + supervisor.stop() + }) + + it('backs off a close from the active relay before opening its replacement', async () => { + const logical = new FakeLogicalClient('disconnected', 'lan') + const openRelay = vi + .fn() + .mockReturnValueOnce(new FakeRelaySession('connected', new RelayOuterError(4429))) + .mockImplementation(() => new FakeRelaySession('connected')) + const deps = dependencies({ + openRelay, + randomBytes: () => new Uint8Array([128, 0]) + }) + const supervisor = new MobileEndpointSupervisor(logical, host, deps) + + await supervisor.start() + logical.publishState('disconnected') + + expect(openRelay).toHaveBeenCalledOnce() + await vi.advanceTimersByTimeAsync(249) + expect(openRelay).toHaveBeenCalledOnce() + await vi.advanceTimersByTimeAsync(1) + expect(openRelay).toHaveBeenCalledTimes(2) + supervisor.stop() + }) + + it('escalates backoff when relay sessions connect and then drop repeatedly', async () => { + const logical = new FakeLogicalClient('disconnected', 'lan') + const openRelay = vi.fn(() => new FakeRelaySession('connected', new RelayOuterError(4408))) + const deps = dependencies({ + openRelay, + randomBytes: () => new Uint8Array([128, 0]) + }) + const supervisor = new MobileEndpointSupervisor(logical, host, deps) + + await supervisor.start() + logical.publishState('disconnected') + await vi.advanceTimersByTimeAsync(250) + expect(openRelay).toHaveBeenCalledTimes(2) + + logical.publishState('disconnected') + await vi.advanceTimersByTimeAsync(499) + expect(openRelay).toHaveBeenCalledTimes(2) + await vi.advanceTimersByTimeAsync(1) + expect(openRelay).toHaveBeenCalledTimes(3) + + logical.publishState('disconnected') + await vi.advanceTimersByTimeAsync(999) + expect(openRelay).toHaveBeenCalledTimes(3) + await vi.advanceTimersByTimeAsync(1) + expect(openRelay).toHaveBeenCalledTimes(4) + supervisor.stop() + }) + + it('does not try a grace credential for a capacity failure before backing off', async () => { + const logical = new FakeLogicalClient('disconnected', 'lan') + const openRelay = vi.fn(() => new FakeRelaySession('disconnected', new RelayOuterError(4429))) + const deps = dependencies({ + readBundle: vi.fn(async () => ({ + ...bundle, + grace: { ...bundle.current, token: 'C'.repeat(43), hash: 'D'.repeat(43), version: 1 } + })), + openRelay, + randomBytes: () => new Uint8Array([128, 0]) + }) + const supervisor = new MobileEndpointSupervisor(logical, host, deps) + + await supervisor.start() + + expect(openRelay).toHaveBeenCalledOnce() + await vi.advanceTimersByTimeAsync(249) + expect(openRelay).toHaveBeenCalledOnce() + await vi.advanceTimersByTimeAsync(1) + expect(openRelay).toHaveBeenCalledTimes(2) + supervisor.stop() + }) + + it('does not redial a rejected current credential on grace cooldown retries', async () => { + const logical = new FakeLogicalClient('disconnected', 'lan') + const openRelay = vi.fn( + (_relay, credential: { version: number }) => + new FakeRelaySession( + 'disconnected', + new RelayOuterError(credential.version === bundle.current.version ? 4401 : 4429) + ) + ) + const deps = dependencies({ + readBundle: vi.fn(async () => ({ + ...bundle, + grace: { ...bundle.current, token: 'C'.repeat(43), hash: 'D'.repeat(43), version: 1 } + })), + openRelay, + randomBytes: () => new Uint8Array([128, 0]) + }) + const supervisor = new MobileEndpointSupervisor(logical, host, deps) + + await supervisor.start() + expect(openRelay).toHaveBeenCalledTimes(2) + + await vi.advanceTimersByTimeAsync(250) + expect(openRelay).toHaveBeenCalledTimes(3) + expect(openRelay.mock.calls[2]?.[1]).toEqual(expect.objectContaining({ version: 1 })) + supervisor.stop() + }) + + it('rotates a rejected current credential after grace keeps relay recovery alive', async () => { + const logical = new FakeLogicalClient('disconnected', 'lan') + const current = { + ...bundle.current, + hash: hashMobileRelayCredential(bundle.current.token) + } + const writeBundle = vi.fn(async () => {}) + const deps = dependencies({ + readBundle: vi.fn(async () => ({ + ...bundle, + current, + grace: { ...current, token: 'C'.repeat(43), hash: 'D'.repeat(43), version: 1 } + })), + openRelay: vi.fn( + (_relay, credential: { version: number }) => + new FakeRelaySession( + credential.version === current.version ? 'disconnected' : 'connected', + credential.version === current.version ? new RelayOuterError(4401) : null + ) + ), + writeBundle + }) + const supervisor = new MobileEndpointSupervisor(logical, host, deps) + + await supervisor.start() + writeBundle.mockClear() + await vi.advanceTimersByTimeAsync(60_000) + + expect(writeBundle).toHaveBeenCalledWith( + expect.objectContaining({ pending: expect.any(Object) }) + ) + supervisor.stop() + }) + + it('does not duplicate transport failures across current and grace credentials', async () => { + const logical = new FakeLogicalClient('disconnected', 'lan') + const openRelay = vi.fn(() => new FakeRelaySession('disconnected', new Error('network down'))) + const resolveRelay = vi.fn(async () => { + throw new Error('director unreachable') + }) + const deps = dependencies({ + readBundle: vi.fn(async () => ({ + ...bundle, + grace: { ...bundle.current, token: 'C'.repeat(43), hash: 'D'.repeat(43), version: 1 } + })), + openRelay, + resolveRelay + }) + const supervisor = new MobileEndpointSupervisor(logical, host, deps) + + await supervisor.start() + + expect(openRelay).toHaveBeenCalledOnce() + expect(resolveRelay).toHaveBeenCalledOnce() + supervisor.stop() + }) + + it('keeps an authenticated relay off the backoff path when persistence fails', async () => { + const logical = new FakeLogicalClient('disconnected', 'lan') + const openRelay = vi.fn(() => new FakeRelaySession('connected')) + const deps = dependencies({ + readBundle: vi.fn(async () => ({ + ...bundle, + grace: { ...bundle.current, token: 'C'.repeat(43), hash: 'D'.repeat(43), version: 1 } + })), + openRelay, + writeBundle: vi.fn(async () => { + throw new Error('secure store unavailable') + }) + }) + const supervisor = new MobileEndpointSupervisor(logical, host, deps) + + await supervisor.start() + expect(openRelay).toHaveBeenCalledOnce() + + logical.publishState('disconnected') + await vi.waitFor(() => expect(openRelay).toHaveBeenCalledTimes(2)) + supervisor.stop() + }) + + it('cancels a pending relay retry when the original direct path reconnects', async () => { + const logical = new FakeLogicalClient('disconnected', 'lan') + const deps = dependencies({ + openRelay: vi.fn(() => new FakeRelaySession('disconnected', new RelayOuterError(4408))), + randomBytes: () => new Uint8Array([128, 0]) + }) + const supervisor = new MobileEndpointSupervisor(logical, host, deps) + + await supervisor.start() + expect(vi.getTimerCount()).toBe(1) + + logical.publishState('connected') + expect(vi.getTimerCount()).toBe(0) + supervisor.stop() + }) + + it('recovers a relay drop while post-migration persistence owns the mutex', async () => { + const logical = new FakeLogicalClient('disconnected', 'lan') + let finishWrite: (() => void) | undefined + const writePending = new Promise((resolve) => { + finishWrite = resolve + }) + const openRelay = vi + .fn() + .mockReturnValueOnce(new FakeRelaySession('connected', new RelayOuterError(4408))) + .mockImplementation(() => new FakeRelaySession('connected')) + const deps = dependencies({ + openRelay, + writeBundle: vi.fn(() => writePending), + randomBytes: () => new Uint8Array([128, 0]) + }) + const supervisor = new MobileEndpointSupervisor(logical, host, deps) + + const starting = supervisor.start() + await vi.waitFor(() => expect(deps.writeBundle).toHaveBeenCalledOnce()) + logical.publishState('disconnected') + finishWrite?.() + await starting + + expect(openRelay).toHaveBeenCalledOnce() + await vi.advanceTimersByTimeAsync(250) + expect(openRelay).toHaveBeenCalledTimes(2) + supervisor.stop() + }) + + it('waits for an external signal instead of polling a host-offline relay', async () => { + const logical = new FakeLogicalClient('disconnected', 'lan') + const openRelay = vi.fn(() => new FakeRelaySession('disconnected', new RelayOuterError(4404))) + const deps = dependencies({ + openRelay, + randomBytes: () => new Uint8Array([128, 0]) + }) + const supervisor = new MobileEndpointSupervisor(logical, host, deps) + + await supervisor.start() + + expect(openRelay).toHaveBeenCalledOnce() + logical.publishState('disconnected') + expect(vi.getTimerCount()).toBe(0) + + supervisor.setForeground(true) + expect(vi.getTimerCount()).toBe(1) + await vi.advanceTimersByTimeAsync(250) + expect(openRelay).toHaveBeenCalledTimes(2) + expect(vi.getTimerCount()).toBe(0) + supervisor.stop() + }) + + it('waits for direct connectivity before replacing a rejected relay credential', async () => { + const logical = new FakeLogicalClient('disconnected', 'lan') + const openRelay = vi.fn(() => new FakeRelaySession('disconnected', new RelayOuterError(4401))) + const deps = dependencies({ openRelay }) + const supervisor = new MobileEndpointSupervisor(logical, host, deps) + + await supervisor.start() + supervisor.setForeground(true) + await vi.advanceTimersByTimeAsync(1000) + + expect(openRelay).toHaveBeenCalledOnce() + expect(deps.writeBundle).not.toHaveBeenCalled() + + logical.publishState('connected') + await vi.waitFor(() => expect(deps.writeBundle).toHaveBeenCalledOnce()) + supervisor.stop() + }) + + it('keeps rejected relay credentials gated until their replacement is durable', async () => { + const logical = new FakeLogicalClient('disconnected', 'lan') + const openRelay = vi + .fn() + .mockReturnValueOnce(new FakeRelaySession('disconnected', new RelayOuterError(4401))) + .mockImplementation(() => new FakeRelaySession('connected')) + let finishCredentialWrite: (() => void) | undefined + const credentialWritePending = new Promise((resolve) => { + finishCredentialWrite = resolve + }) + const writeBundle = vi + .fn<(value: MobileRelayCredentialBundle) => Promise>() + .mockResolvedValue() + .mockResolvedValueOnce() + .mockReturnValueOnce(credentialWritePending) + mockCredentialRotation(logical) + const deps = dependencies({ openRelay, writeBundle }) + const supervisor = new MobileEndpointSupervisor(logical, host, deps) + + await supervisor.start() + logical.publishState('connected') + await vi.waitFor(() => expect(writeBundle).toHaveBeenCalledTimes(2)) + + // The direct socket can disappear after the server commits but before the + // replacement credential finishes its durable write. + logical.publishState('disconnected') + await vi.advanceTimersByTimeAsync(0) + expect(openRelay).toHaveBeenCalledOnce() + + finishCredentialWrite?.() + await vi.waitFor(() => expect(openRelay).toHaveBeenCalledTimes(2)) + expect(openRelay).toHaveBeenLastCalledWith( + relay, + expect.objectContaining({ version: 3 }), + expect.any(String) + ) + supervisor.stop() + }) + + it('uses a scheduled credential rotation that finishes after relay rejection', async () => { + const logical = new FakeLogicalClient('connected', 'lan') + let finishCredentialWrite: (() => void) | undefined + const credentialWritePending = new Promise((resolve) => { + finishCredentialWrite = resolve + }) + const writeBundle = vi + .fn<(value: MobileRelayCredentialBundle) => Promise>() + .mockResolvedValue() + .mockResolvedValueOnce() + .mockReturnValueOnce(credentialWritePending) + mockCredentialRotation(logical) + const openRelay = vi.fn( + (_relay, credential: { version: number }) => + new FakeRelaySession( + credential.version === bundle.current.version ? 'disconnected' : 'connected', + credential.version === bundle.current.version ? new RelayOuterError(4401) : null + ) + ) + const deps = dependencies({ + readBundle: vi.fn(async () => ({ + ...bundle, + current: { ...bundle.current, expiresAt: Date.now() + 60_000 } + })), + openRelay, + writeBundle + }) + const supervisor = new MobileEndpointSupervisor(logical, host, deps) + + await supervisor.start() + logical.publishState('connected') + await vi.waitFor(() => expect(writeBundle).toHaveBeenCalledTimes(2)) + + // The expiring credential can be rejected while its replacement is waiting on SecureStore. + logical.publishState('disconnected') + await vi.waitFor(() => expect(openRelay).toHaveBeenCalledOnce()) + finishCredentialWrite?.() + + await vi.waitFor(() => expect(openRelay).toHaveBeenCalledTimes(2)) + expect(openRelay).toHaveBeenLastCalledWith( + relay, + expect.objectContaining({ version: 3 }), + expect.any(String) + ) + supervisor.stop() + }) + + it('does not open a resolved relay replacement after backgrounding', async () => { + const logical = new FakeLogicalClient('disconnected', 'lan') + let finishResolve: ((value: typeof relay) => void) | undefined + const resolvePending = new Promise((resolve) => { + finishResolve = resolve + }) + const openRelay = vi + .fn() + .mockReturnValueOnce(new FakeRelaySession('disconnected', new RelayOuterError(4409))) + .mockReturnValueOnce(new FakeRelaySession('disconnected', new RelayOuterError(4408))) + const deps = dependencies({ + openRelay, + resolveRelay: vi.fn(() => resolvePending) + }) + const supervisor = new MobileEndpointSupervisor(logical, host, deps) + + const starting = supervisor.start() + await vi.waitFor(() => expect(deps.resolveRelay).toHaveBeenCalledOnce()) + supervisor.setForeground(false) + finishResolve?.(relay) + await starting + + expect(openRelay).toHaveBeenCalledOnce() + expect(vi.getTimerCount()).toBe(0) + supervisor.stop() + }) + + it('does not recreate a lease retry after forced replacement is backgrounded', async () => { + const logical = new FakeLogicalClient('disconnected', 'lan') + let finishResolve: ((value: typeof relay) => void) | undefined + const resolvePending = new Promise((resolve) => { + finishResolve = resolve + }) + const openRelay = vi + .fn() + .mockReturnValueOnce(new FakeRelaySession('connected', null, Date.now() + 31_000)) + .mockReturnValueOnce(new FakeRelaySession('disconnected', new RelayOuterError(4409))) + .mockReturnValueOnce(new FakeRelaySession('disconnected', new RelayOuterError(4408))) + const deps = dependencies({ + openRelay, + resolveRelay: vi.fn(() => resolvePending) + }) + const supervisor = new MobileEndpointSupervisor(logical, host, deps) + + await supervisor.start() + await vi.advanceTimersByTimeAsync(1000) + await vi.waitFor(() => expect(deps.resolveRelay).toHaveBeenCalledOnce()) + supervisor.setForeground(false) + finishResolve?.(relay) + await vi.waitFor(() => expect(deps.saveHost).toHaveBeenCalledOnce()) + await vi.advanceTimersByTimeAsync(0) + + expect(openRelay).toHaveBeenCalledTimes(2) + expect(vi.getTimerCount()).toBe(0) + supervisor.stop() + }) + + it('does not recreate a lease timer after stop races relay persistence', async () => { + const logical = new FakeLogicalClient('disconnected', 'lan') + let finishWrite: (() => void) | undefined + const writePending = new Promise((resolve) => { + finishWrite = resolve + }) + const deps = dependencies({ writeBundle: vi.fn(() => writePending) }) + const supervisor = new MobileEndpointSupervisor(logical, host, deps) + + const starting = supervisor.start() + await vi.waitFor(() => expect(deps.writeBundle).toHaveBeenCalledOnce()) + supervisor.stop() + finishWrite?.() + await starting + + expect(vi.getTimerCount()).toBe(0) + }) + + it('does not poll a host-offline relay through forced lease retries', async () => { + const logical = new FakeLogicalClient('disconnected', 'lan') + const openRelay = vi + .fn() + .mockReturnValueOnce(new FakeRelaySession('connected', null, Date.now() + 31_000)) + .mockImplementation(() => new FakeRelaySession('disconnected', new RelayOuterError(4404))) + const deps = dependencies({ openRelay }) + const supervisor = new MobileEndpointSupervisor(logical, host, deps) + + await supervisor.start() + await vi.advanceTimersByTimeAsync(1000) + expect(openRelay).toHaveBeenCalledTimes(2) + + await vi.advanceTimersByTimeAsync(5000) + expect(openRelay).toHaveBeenCalledTimes(2) + supervisor.stop() + }) + + it('keeps a fatal lease-replacement gate after the active relay later drops', async () => { + const logical = new FakeLogicalClient('disconnected', 'lan') + const openRelay = vi + .fn() + .mockReturnValueOnce( + new FakeRelaySession('connected', new RelayOuterError(4408), Date.now() + 31_000) + ) + .mockReturnValueOnce(new FakeRelaySession('disconnected', new RelayOuterError(4401))) + .mockImplementation(() => new FakeRelaySession('connected')) + const deps = dependencies({ + openRelay, + randomBytes: () => new Uint8Array([128, 0]) + }) + const supervisor = new MobileEndpointSupervisor(logical, host, deps) + + await supervisor.start() + await vi.advanceTimersByTimeAsync(1000) + expect(openRelay).toHaveBeenCalledTimes(2) + + // The old relay can outlive its rejected lease replacement, then close separately. + logical.publishState('disconnected') + await vi.advanceTimersByTimeAsync(30_000) + + expect(openRelay).toHaveBeenCalledTimes(2) + supervisor.stop() + }) + + it('keeps revival nudges inside a failed lease rotation cooldown', async () => { + const logical = new FakeLogicalClient('disconnected', 'lan') + const openRelay = vi + .fn() + .mockReturnValueOnce(new FakeRelaySession('connected', null, Date.now() + 31_000)) + .mockReturnValueOnce(new FakeRelaySession('disconnected', new RelayOuterError(4429))) + .mockImplementation(() => new FakeRelaySession('connected')) + const deps = dependencies({ + openRelay, + randomBytes: () => new Uint8Array([128, 0]) + }) + const supervisor = new MobileEndpointSupervisor(logical, host, deps) + + await supervisor.start() + await vi.advanceTimersByTimeAsync(1000) + expect(openRelay).toHaveBeenCalledTimes(2) + + supervisor.setForeground(true) + await vi.advanceTimersByTimeAsync(249) + expect(openRelay).toHaveBeenCalledTimes(2) + await vi.advanceTimersByTimeAsync(1) + expect(openRelay).toHaveBeenCalledTimes(3) + supervisor.stop() + }) + + it('keeps lease rotation inside an active relay failure cooldown', async () => { + const logical = new FakeLogicalClient('disconnected', 'lan') + const openRelay = vi + .fn() + .mockReturnValueOnce( + new FakeRelaySession('connected', new RelayOuterError(4429), Date.now() + 31_000) + ) + .mockImplementation(() => new FakeRelaySession('connected')) + const deps = dependencies({ + openRelay, + randomBytes: () => new Uint8Array([128, 0]) + }) + const supervisor = new MobileEndpointSupervisor(logical, host, deps) + + await supervisor.start() + await vi.advanceTimersByTimeAsync(900) + logical.publishState('disconnected') + + await vi.advanceTimersByTimeAsync(100) + expect(openRelay).toHaveBeenCalledOnce() + await vi.advanceTimersByTimeAsync(150) + expect(openRelay).toHaveBeenCalledTimes(2) + supervisor.stop() + }) + + it('clears relay backoff on a genuine foreground so the retry is immediate', async () => { + const logical = new FakeLogicalClient('disconnected', 'lan') + const openRelay = vi.fn(() => new FakeRelaySession('disconnected', new RelayOuterError(4408))) + const deps = dependencies({ + openRelay, + openDirect: vi.fn(() => new FakeSession('disconnected')), + randomBytes: () => new Uint8Array([128, 0]) + }) + const supervisor = new MobileEndpointSupervisor(logical, host, deps) + + await supervisor.start() + const afterStart = openRelay.mock.calls.length + + // Background → foreground is a fresh signal: dial now, not after the cooldown. + supervisor.setForeground(false) + supervisor.setForeground(true) + await vi.advanceTimersByTimeAsync(0) + expect(openRelay.mock.calls.length).toBeGreaterThan(afterStart) + supervisor.stop() + }) + it('releases a background relay session and reconnects it on foreground', async () => { const logical = new FakeLogicalClient('connected', 'relay') const deps = dependencies() @@ -237,6 +905,7 @@ describe('mobile endpoint supervisor', () => { supervisor.setForeground(false) expect(logical.suspendActiveSession).toHaveBeenCalledOnce() expect(logical.getState()).toBe('disconnected') + expect(vi.getTimerCount()).toBe(0) supervisor.setForeground(true) await vi.waitFor(() => expect(logical.migrateTo).toHaveBeenCalled()) diff --git a/mobile/src/transport/mobile-endpoint-supervisor.ts b/mobile/src/transport/mobile-endpoint-supervisor.ts index 9815a1a3ed4..c8b3cdbcfb1 100644 --- a/mobile/src/transport/mobile-endpoint-supervisor.ts +++ b/mobile/src/transport/mobile-endpoint-supervisor.ts @@ -1,5 +1,7 @@ -import type { MobileRelayEndpoint } from '../../../src/shared/mobile-relay-credential-contract' import { openAuthenticatedDirectEndpoint } from './mobile-direct-endpoint-probe' +import type { MobileEndpointSupervisorDependencies } from './mobile-endpoint-supervisor-contract' +import { RelayReconnectController } from './mobile-relay-reconnect-controller' +import { RelayLeaseRotationTimer } from './mobile-relay-lease-rotation-timer' import { MobileEndpointHysteresis } from './mobile-endpoint-hysteresis' import { encodeBase64Url, @@ -13,37 +15,17 @@ import { rotateMobileRelayCredential } from './mobile-relay-credential-rotation' import type { MobileRelayCredentialBundle } from './mobile-relay-credential-bundle' -import type { MobileRelayRpcSession } from './mobile-relay-rpc-session' -import { resolveMobileRelayEndpoint } from './mobile-relay-resume-director' -import type { RpcClient } from './rpc-client' import type { StableLogicalRpcClient } from './stable-logical-rpc-client' import type { HostProfile } from './types' +export type { MobileEndpointSupervisorDependencies } from './mobile-endpoint-supervisor-contract' + const DIRECT_PROBE_INTERVAL_MS = 15_000 const DIRECT_OBSERVATION_MS = 30_000 const MINIMUM_DWELL_MS = 60_000 const FAILURE_COOLDOWN_MS = 60_000 -const LEASE_ROTATION_MARGIN_MS = 30_000 - -export type MobileEndpointSupervisorDependencies = { - openDirect: (endpoint: string) => RpcClient - openRelay: ( - relay: MobileRelayEndpoint, - credential: { token: string; version: number }, - confirmReqId: string - ) => MobileRelayRpcSession - resolveRelay: typeof resolveMobileRelayEndpoint - readBundle: (hostId: string) => Promise - writeBundle: (bundle: MobileRelayCredentialBundle) => Promise - saveHost: (host: HostProfile) => Promise - now: () => number - randomBytes: (length: number) => Uint8Array - setTimer: typeof setTimeout - clearTimer: typeof clearTimeout -} export class MobileEndpointSupervisor { - private host: HostProfile private bundle: MobileRelayCredentialBundle | null = null private stopped = false private foreground = true @@ -51,22 +33,27 @@ export class MobileEndpointSupervisor { private credentialRotationInFlight = false private relayRotationPending = false private probeTimer: ReturnType | null = null - private leaseTimer: ReturnType | null = null private unsubscribeState: (() => void) | null = null private readonly hysteresis: MobileEndpointHysteresis + private readonly relayReconnect: RelayReconnectController + private readonly leaseRotation: RelayLeaseRotationTimer constructor( private readonly logical: StableLogicalRpcClient, - host: HostProfile, + private host: HostProfile, private readonly dependencies: MobileEndpointSupervisorDependencies ) { - this.host = host this.hysteresis = new MobileEndpointHysteresis(dependencies.now(), { directSuccessesRequired: 3, directObservationMs: DIRECT_OBSERVATION_MS, failureCooldownMs: FAILURE_COOLDOWN_MS, minimumDwellMs: MINIMUM_DWELL_MS }) + this.relayReconnect = new RelayReconnectController(dependencies, this.recoverRelay.bind(this)) + this.leaseRotation = new RelayLeaseRotationTimer(dependencies, () => { + this.relayRotationPending = true + void this.recoverRelay(true) + }) } async start(): Promise { @@ -77,21 +64,16 @@ export class MobileEndpointSupervisor { this.unsubscribeState = this.logical.onStateChange((state) => { if (state === 'connected') { if (this.logical.getActivePath() !== 'relay') { - void this.rotateCredentialIfNeeded() + void this.rotateCredentialIfNeeded(this.relayReconnect.resetForDirectConnection()) } this.scheduleDirectProbe() - } else if (state === 'reconnecting' || state === 'disconnected' || state === 'auth-failed') { + } else { // Why: the direct client enters reconnecting after its first failed // dial and may never publish disconnected while its retry loop lives. - void this.recoverRelay() + this.relayReconnect.handleStateFailure(this.logical, state) } }) - const initialState = this.logical.getState() - if ( - initialState === 'reconnecting' || - initialState === 'disconnected' || - initialState === 'auth-failed' - ) { + if (this.relayReconnect.needsRecovery(this.logical.getState())) { // Why: the first direct dial can fail while encrypted relay credentials // are still loading, before the supervisor subscribes to state changes. await this.recoverRelay() @@ -101,20 +83,17 @@ export class MobileEndpointSupervisor { } setForeground(foreground: boolean): void { + const wasForeground = this.foreground this.foreground = foreground if (foreground) { - void this.recoverRelay(this.relayRotationPending) + this.relayReconnect.handleForeground(this.logical, wasForeground) this.scheduleDirectProbe(0) } else { - if (this.logical.getActivePath() === 'relay') { - // Why: background phones must not hold billed relay data splices; the - // stable client keeps subscriptions for authenticated foreground replay. - this.logical.suspendActiveSession() - } - if (this.probeTimer) { - this.dependencies.clearTimer(this.probeTimer) - this.probeTimer = null - } + // Why: background phones must not hold billed relay data splices. + this.relayReconnect.suspendActiveRelay(this.logical) + this.clearDirectProbeTimer() + this.relayReconnect.clear() + this.leaseRotation.clear() } } @@ -122,42 +101,64 @@ export class MobileEndpointSupervisor { this.stopped = true this.unsubscribeState?.() this.unsubscribeState = null - if (this.probeTimer) { - this.dependencies.clearTimer(this.probeTimer) - this.probeTimer = null - } - this.clearLeaseTimer() + this.clearDirectProbeTimer() + this.relayReconnect.clear() + this.leaseRotation.clear() } private async recoverRelay(forceReplacement = false): Promise { + // Why: connecting/handshaking is live direct progress; a relay dial would race it. if ( this.stopped || !this.foreground || this.operationInFlight || !this.bundle || !this.host.relay || - (!forceReplacement && this.logical.getState() === 'connected') + (!forceReplacement && !this.relayReconnect.needsRecovery(this.logical.getState())) ) { return } + // Why: revival and lease timers can overlap resume failures; one shared cooldown + // prevents PEER_DROPPED/LIMIT_EXCEEDED reconnect churn. + if (this.relayReconnect.shouldDefer()) { + return + } this.operationInFlight = true + let lastError: Error | null = null + let retryAfterOperation = false try { - const credentials = [this.bundle.current, this.bundle.grace].filter( - (credential): credential is NonNullable => - Boolean(credential && credential.expiresAt > this.dependencies.now()) + const credentials = this.relayReconnect.eligibleCredentials( + this.bundle.current, + this.bundle.grace ) for (const credential of credentials) { - if (await this.tryRelayCredential(credential)) { + const result = await this.tryRelayCredential(credential) + if (result.ok) { + retryAfterOperation = this.logical.getState() !== 'connected' return } + lastError = result.error + if (this.relayReconnect.shouldTryGraceAfterRelayFailure(result.error)) { + // Why: a rejected version stays invalid; retry only the grace credential. + this.relayReconnect.recordRejectedCredential(credential.version) + } else { + break + } + } + if (credentials.length > 0) { + // Why: cleanup may happen while a relay dial is awaiting the network; + // record its outcome without recreating a foreground retry timer. + const scheduleRetry = !forceReplacement && this.foreground && !this.stopped + this.relayReconnect.registerFailure(lastError, scheduleRetry) } } finally { this.operationInFlight = false - if (forceReplacement && this.relayRotationPending && !this.stopped && !this.leaseTimer) { - this.leaseTimer = this.dependencies.setTimer(() => { - this.leaseTimer = null - void this.recoverRelay(true) - }, 5000) + if (forceReplacement && this.relayRotationPending && !this.stopped && this.foreground) { + this.leaseRotation.armRetry(this.relayReconnect.retryDelayMs(5000)) + } + // Why: the active relay can drop while migration follow-up still owns the mutex. + if (retryAfterOperation && !this.stopped && this.foreground) { + void this.recoverRelay() } } } @@ -165,13 +166,13 @@ export class MobileEndpointSupervisor { private async tryRelayCredential(credential: { token: string version: number - }): Promise { + }): Promise<{ ok: true } | { ok: false; error: Error }> { const first = await this.openAndMigrateRelay(credential) if (first.ok) { - return true + return first } if (!isDirectorResolutionFailure(first.error) || !this.host.relay) { - return false + return first } try { const resolved = await this.dependencies.resolveRelay({ @@ -179,9 +180,9 @@ export class MobileEndpointSupervisor { resumeToken: credential.token }) this.host = await persistRelayHost(this.host, resolved, this.dependencies.saveHost) - return (await this.openAndMigrateRelay(credential)).ok - } catch { - return false + return await this.openAndMigrateRelay(credential) + } catch (error) { + return { ok: false, error: toError(error) } } } @@ -189,7 +190,8 @@ export class MobileEndpointSupervisor { token: string version: number }): Promise<{ ok: true } | { ok: false; error: Error }> { - if (!this.host.relay || !this.bundle) { + // Why: director resolution and grace fallback can finish after background/stop. + if (this.stopped || !this.foreground || !this.host.relay || !this.bundle) { return { ok: false, error: new Error('relay state missing') } } const session = this.dependencies.openRelay( @@ -199,17 +201,23 @@ export class MobileEndpointSupervisor { ) try { await this.logical.migrateTo(session, 'relay') + this.relayReconnect.setActiveSession(session) if (!this.foreground) { - this.logical.suspendActiveSession() + this.relayReconnect.suspendActiveRelay(this.logical) } this.relayRotationPending = false this.hysteresis.recordMigration(this.dependencies.now()) const confirmation = session.getResumeConfirmation() if (confirmation) { this.bundle = applyResumeConfirmation(this.bundle, credential.version, confirmation) - await this.dependencies.writeBundle(this.bundle) + // Why: the relay is already authenticated; a SecureStore failure must + // not open another socket or count against transport recovery backoff. + await this.dependencies.writeBundle(this.bundle).catch(() => {}) } - this.scheduleLeaseRotation(session) + // Why: async persistence can finish after stop/background; never recreate a stale timer. + this.leaseRotation.scheduleFromLease( + this.stopped || !this.foreground ? null : session.getLeaseExpiresAt() + ) this.scheduleDirectProbe() return { ok: true } } catch (error) { @@ -258,30 +266,32 @@ export class MobileEndpointSupervisor { await this.logical.migrateTo(successful.client, successful.path) successful = null this.hysteresis.recordMigration(this.dependencies.now()) - this.clearLeaseTimer() + this.leaseRotation.clear() this.relayRotationPending = false - await this.rotateCredentialIfNeeded() + await this.rotateCredentialIfNeeded(this.relayReconnect.resetForDirectConnection()) } finally { successful?.client.close() this.operationInFlight = false - if (this.relayRotationPending) { - void this.recoverRelay(true) + // Why: a relay drop or backoff timer can arrive while the direct probe owns the mutex. + if (this.relayRotationPending || this.logical.getState() !== 'connected') { + void this.recoverRelay(this.relayRotationPending) } this.scheduleDirectProbe() } } - private async rotateCredentialIfNeeded(): Promise { + private async rotateCredentialIfNeeded(force = false): Promise { if ( this.stopped || this.credentialRotationInFlight || !this.bundle || this.logical.getActivePath() === 'relay' || - !mobileRelayCredentialNeedsRotation(this.bundle, this.dependencies.now()) + (!force && !mobileRelayCredentialNeedsRotation(this.bundle, this.dependencies.now())) ) { return } this.credentialRotationInFlight = true + let credentialRefreshed = false try { const result = await rotateMobileRelayCredential({ client: this.logical, @@ -290,33 +300,32 @@ export class MobileEndpointSupervisor { randomBytes: this.dependencies.randomBytes }) this.bundle = result.bundle + // Why: a scheduled rotation can finish after the old credential enters the rejection gate. + credentialRefreshed = true this.host = await persistRelayHost(this.host, result.relay, this.dependencies.saveHost) } catch { // Why: pending material remains durable; the next authenticated direct // opportunity must reconcile it before creating another install key. } finally { + if (credentialRefreshed) { + this.relayReconnect.completeCredentialRefresh() + } this.credentialRotationInFlight = false + if ( + credentialRefreshed && + !this.stopped && + this.foreground && + this.relayReconnect.needsRecovery(this.logical.getState()) + ) { + void this.recoverRelay() + } } } - private scheduleLeaseRotation(session: MobileRelayRpcSession): void { - this.clearLeaseTimer() - const deadline = session.getLeaseExpiresAt() - if (!deadline) { - return - } - const delay = Math.max(1000, deadline - this.dependencies.now() - LEASE_ROTATION_MARGIN_MS) - this.leaseTimer = this.dependencies.setTimer(() => { - this.leaseTimer = null - this.relayRotationPending = true - void this.recoverRelay(true) - }, delay) - } - - private clearLeaseTimer(): void { - if (this.leaseTimer) { - this.dependencies.clearTimer(this.leaseTimer) - this.leaseTimer = null + private clearDirectProbeTimer(): void { + if (this.probeTimer) { + this.dependencies.clearTimer(this.probeTimer) + this.probeTimer = null } } } diff --git a/mobile/src/transport/mobile-relay-lease-rotation-timer.ts b/mobile/src/transport/mobile-relay-lease-rotation-timer.ts new file mode 100644 index 00000000000..bf2a3dd4d1c --- /dev/null +++ b/mobile/src/transport/mobile-relay-lease-rotation-timer.ts @@ -0,0 +1,59 @@ +// Why: the relay resume lease expires; the phone must proactively re-resume a +// little before the deadline (and retry shortly if a forced rotation didn't land) +// so the session never lapses. Owns the single lease/rotation timer slot. +const LEASE_ROTATION_MARGIN_MS = 30_000 + +export type RelayLeaseRotationDependencies = { + now: () => number + setTimer: typeof setTimeout + clearTimer: typeof clearTimeout +} + +export class RelayLeaseRotationTimer { + private timer: ReturnType | null = null + + constructor( + private readonly dependencies: RelayLeaseRotationDependencies, + private readonly onRotate: () => void + ) {} + + // Arm rotation a margin before the lease deadline. No-op with no deadline. + scheduleFromLease(leaseExpiresAt: number | null): void { + this.clear() + if (!leaseExpiresAt) { + return + } + const delay = Math.max( + 1000, + leaseExpiresAt - this.dependencies.now() - LEASE_ROTATION_MARGIN_MS + ) + this.arm(delay) + } + + // Re-arm a short retry when a forced rotation's recovery did not complete. + // Ignored while a timer is already pending so retries never stack. + armRetry(delayMs: number | null): void { + if (delayMs == null || this.timer) { + return + } + this.arm(delayMs) + } + + get pending(): boolean { + return this.timer != null + } + + clear(): void { + if (this.timer) { + this.dependencies.clearTimer(this.timer) + this.timer = null + } + } + + private arm(delayMs: number): void { + this.timer = this.dependencies.setTimer(() => { + this.timer = null + this.onRotate() + }, delayMs) + } +} diff --git a/mobile/src/transport/mobile-relay-reconnect-controller.test.ts b/mobile/src/transport/mobile-relay-reconnect-controller.test.ts new file mode 100644 index 00000000000..c68ed48a4f8 --- /dev/null +++ b/mobile/src/transport/mobile-relay-reconnect-controller.test.ts @@ -0,0 +1,147 @@ +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' +import type { MobileRelayRpcSession } from './mobile-relay-rpc-session' +import { MobileE2EEAuthenticationError } from './mobile-e2ee-v2-physical-channel' +import { RelayOuterError } from './mobile-relay-e2ee-link' +import { RelayReconnectController } from './mobile-relay-reconnect-controller' +import type { StableLogicalRpcClient } from './stable-logical-rpc-client' + +vi.mock('react-native', () => ({ Platform: { OS: 'ios' } })) +vi.mock('expo-crypto', () => ({ getRandomBytes: (length: number) => new Uint8Array(length) })) + +describe('relay reconnect controller', () => { + beforeEach(() => { + vi.useFakeTimers() + }) + + afterEach(() => { + vi.useRealTimers() + }) + + it('keeps one retry timer and cancels it when recovery needs an external signal', () => { + const onRetry = vi.fn() + const reconnect = createController(onRetry) + + reconnect.registerFailure(new RelayOuterError(4429)) + reconnect.registerFailure(new RelayOuterError(4408)) + expect(vi.getTimerCount()).toBe(1) + + reconnect.registerFailure(new RelayOuterError(4404)) + expect(vi.getTimerCount()).toBe(0) + expect(reconnect.shouldDefer()).toBe(true) + expect(vi.getTimerCount()).toBe(0) + vi.runAllTimers() + expect(onRetry).not.toHaveBeenCalled() + }) + + it('drops a pending relay retry after direct connectivity wins', () => { + const onRetry = vi.fn() + const reconnect = createController(onRetry) + + reconnect.registerFailure(new RelayOuterError(4408)) + expect(vi.getTimerCount()).toBe(1) + + reconnect.resetForDirectConnection() + expect(vi.getTimerCount()).toBe(0) + vi.runAllTimers() + expect(onRetry).not.toHaveBeenCalled() + }) + + it('waits for an external signal after rejected E2EE authentication', () => { + const onRetry = vi.fn() + const reconnect = createController(onRetry) + + reconnect.registerFailure(new MobileE2EEAuthenticationError()) + + expect(reconnect.shouldDefer()).toBe(true) + expect(vi.getTimerCount()).toBe(0) + vi.advanceTimersByTime(60_000) + expect(onRetry).not.toHaveBeenCalled() + }) + + it('upgrades host-revival gating to fresh credentials without later downgrading it', () => { + const reconnect = createController(vi.fn()) + + reconnect.registerFailure(new RelayOuterError(4404)) + reconnect.registerFailure(new RelayOuterError(4401)) + reconnect.registerFailure(new RelayOuterError(4408)) + + expect(reconnect.resetForDirectConnection()).toBe(true) + expect(vi.getTimerCount()).toBe(0) + }) + + it('extends forced-rotation retries to the exponential cooldown', () => { + const reconnect = createController(vi.fn()) + + for (let failure = 0; failure < 6; failure++) { + reconnect.registerFailure(new RelayOuterError(4429), false) + } + + expect(reconnect.retryDelayMs(5000)).toBe(8000) + expect(vi.getTimerCount()).toBe(0) + }) + + it('does not reset backoff merely because a failed attempt took one ceiling', () => { + const onRetry = vi.fn() + const reconnect = createController(onRetry) + + reconnect.registerFailure(new RelayOuterError(4429)) + vi.advanceTimersByTime(250) + expect(onRetry).toHaveBeenCalledOnce() + + vi.advanceTimersByTime(30_000) + reconnect.registerFailure(new RelayOuterError(4429)) + vi.advanceTimersByTime(249) + expect(onRetry).toHaveBeenCalledOnce() + vi.advanceTimersByTime(1) + expect(onRetry).toHaveBeenCalledOnce() + vi.advanceTimersByTime(249) + expect(onRetry).toHaveBeenCalledOnce() + vi.advanceTimersByTime(1) + expect(onRetry).toHaveBeenCalledTimes(2) + }) + + it('resets backoff after an authenticated relay remains stable', () => { + const onRetry = vi.fn() + const reconnect = createController(onRetry) + const session = { + getFailure: () => new RelayOuterError(4408) + } as MobileRelayRpcSession + const logical = { + getActivePath: () => 'relay' + } as StableLogicalRpcClient + + reconnect.registerFailure(new RelayOuterError(4429)) + vi.advanceTimersByTime(250) + reconnect.setActiveSession(session) + vi.advanceTimersByTime(30_000) + reconnect.registerActiveFailure(logical) + + vi.advanceTimersByTime(249) + expect(onRetry).toHaveBeenCalledOnce() + vi.advanceTimersByTime(1) + expect(onRetry).toHaveBeenCalledTimes(2) + }) + + it('uses grace only when the outer relay credential was rejected', () => { + const reconnect = createController(vi.fn()) + + expect(reconnect.shouldTryGraceAfterRelayFailure(new RelayOuterError(4401))).toBe(true) + expect(reconnect.shouldTryGraceAfterRelayFailure(new Error('relay transport error'))).toBe( + false + ) + expect(reconnect.shouldTryGraceAfterRelayFailure(new RelayOuterError(4408))).toBe(false) + expect(reconnect.shouldTryGraceAfterRelayFailure(new RelayOuterError(4429))).toBe(false) + }) +}) + +function createController(onRetry: () => void): RelayReconnectController { + return new RelayReconnectController( + { + now: Date.now, + randomBytes: () => new Uint8Array([128, 0]), + setTimer: setTimeout, + clearTimer: clearTimeout + }, + onRetry + ) +} diff --git a/mobile/src/transport/mobile-relay-reconnect-controller.ts b/mobile/src/transport/mobile-relay-reconnect-controller.ts new file mode 100644 index 00000000000..7e3e1ec0e13 --- /dev/null +++ b/mobile/src/transport/mobile-relay-reconnect-controller.ts @@ -0,0 +1,276 @@ +import { + isMobileRelayCloseCode, + MOBILE_RELAY_CLOSE_CODE, + mobileRelayRecoveryFor +} from '../../../src/shared/mobile-relay-close-codes' +import type { MobileRelayRpcSession } from './mobile-relay-rpc-session' +import { MobileE2EEAuthenticationError } from './mobile-e2ee-v2-physical-channel' +import { RelayOuterError } from './mobile-relay-e2ee-link' +import type { StableLogicalRpcClient } from './stable-logical-rpc-client' +import type { ConnectionState } from './types' + +// Why: relay resume closes and silent cellular NAT rebinds otherwise cause +// immediate re-dials that ping-pong the phone between connected and disconnected. +const RELAY_BACKOFF_MIN_MS = 250 +const RELAY_BACKOFF_BASE_MS = 500 +const RELAY_BACKOFF_CEILING_MS = 30_000 +const RELAY_STABLE_CONNECTION_MS = RELAY_BACKOFF_CEILING_MS + +export type RelayReconnectDependencies = { + now: () => number + randomBytes: (length: number) => Uint8Array + setTimer: typeof setTimeout + clearTimer: typeof clearTimeout +} + +type RecoveryGate = 'external-signal' | 'fresh-credential' + +export class RelayReconnectController { + private consecutiveFailures = 0 + private activeRelayConnectedAt: number | null = null + private nextAttemptAt = 0 + private timer: ReturnType | null = null + private activeSession: MobileRelayRpcSession | null = null + private recoveryGate: RecoveryGate | null = null + private readonly rejectedCredentialVersions = new Set() + + constructor( + private readonly dependencies: RelayReconnectDependencies, + private readonly onRetry: (forceReplacement?: boolean) => void + ) {} + + handleForeground(logical: StableLogicalRpcClient, wasForeground: boolean): void { + if (!wasForeground) { + // Why: an app resume is a fresh signal, unlike repeated network-flap nudges. + if (this.recoveryGate !== 'fresh-credential') { + this.reset() + } + } else if (this.recoveryGate === 'external-signal') { + this.recoveryGate = null + } + if ( + wasForeground && + this.recoveryGate !== 'fresh-credential' && + logical.getState() === 'connected' + ) { + // Why: a network handoff can leave the relay half-open without publishing a close. + this.suspendActiveRelay(logical) + } + // Why: revival nudges must honor failure cooldowns even when lease rotation is pending. + this.onRetry() + } + + handleStateFailure(logical: StableLogicalRpcClient, state: ConnectionState): void { + if (!this.needsRecovery(state)) { + return + } + this.registerActiveFailure(logical) + this.onRetry() + } + + needsRecovery(state: ConnectionState): boolean { + return state !== 'connected' && state !== 'connecting' && state !== 'handshaking' + } + + suspendActiveRelay(logical: StableLogicalRpcClient): void { + if (logical.getActivePath() !== 'relay') { + return + } + this.activeSession = null + this.activeRelayConnectedAt = null + logical.suspendActiveSession() + } + + setActiveSession(session: MobileRelayRpcSession): void { + this.activeSession = session + this.activeRelayConnectedAt = this.dependencies.now() + this.nextAttemptAt = 0 + this.recoveryGate = null + this.clearTimer() + } + + resetForDirectConnection(): boolean { + const needsCredentialRefresh = + this.recoveryGate === 'fresh-credential' || this.rejectedCredentialVersions.size > 0 + this.activeSession = null + this.activeRelayConnectedAt = null + if (needsCredentialRefresh) { + // Why: the rejected credential stays unusable until its replacement is durable. + this.consecutiveFailures = 0 + this.nextAttemptAt = 0 + this.recoveryGate = 'fresh-credential' + this.clearTimer() + } else { + this.reset() + } + return needsCredentialRefresh + } + + completeCredentialRefresh(): void { + if (this.recoveryGate === 'fresh-credential') { + this.rejectedCredentialVersions.clear() + this.reset() + } + } + + eligibleCredentials( + ...credentials: Array + ): T[] { + const eligible = credentials.filter((credential): credential is T => + Boolean( + credential && + credential.expiresAt > this.dependencies.now() && + !this.rejectedCredentialVersions.has(credential.version) + ) + ) + if (eligible.length === 0 && this.rejectedCredentialVersions.size > 0) { + this.recoveryGate = 'fresh-credential' + this.clearTimer() + } + return eligible + } + + recordRejectedCredential(version: number): void { + this.rejectedCredentialVersions.add(version) + } + + registerActiveFailure(logical: StableLogicalRpcClient): void { + if (logical.getActivePath() !== 'relay') { + return + } + const failure = this.activeSession?.getFailure() + this.activeSession = null + if (failure) { + // Why: active relay closes need the same cooldown as failed replacement dials. + this.registerFailure(failure) + } else { + this.activeRelayConnectedAt = null + } + } + + // True when the caller is still inside the cooldown window and must not + // re-dial. Arms the self-scheduled retry so recovery still happens on its own. + shouldDefer(): boolean { + if (this.recoveryGate) { + return true + } + if (this.dependencies.now() < this.nextAttemptAt) { + this.scheduleRetry() + return true + } + return false + } + + registerFailure(error: Error | null, scheduleRetry = true): void { + const code = error instanceof RelayOuterError ? error.code : null + const recovery = + code != null && isMobileRelayCloseCode(code) + ? mobileRelayRecoveryFor(code, 'phone-resume') + : null + if ( + this.recoveryGate === 'fresh-credential' || + (this.recoveryGate === 'external-signal' && recovery?.kind !== 'disable-relay-credential') + ) { + // Why: only the gate's external signal can make a known-fatal recovery retryable. + this.clearTimer() + return + } + const now = this.dependencies.now() + if ( + this.activeRelayConnectedAt != null && + now - this.activeRelayConnectedAt >= RELAY_STABLE_CONNECTION_MS + ) { + this.consecutiveFailures = 0 + } + // Why: elapsed time inside a slow failed dial is not evidence of recovery; + // only an authenticated relay that survived the stability window resets the streak. + this.activeRelayConnectedAt = null + this.consecutiveFailures += 1 + const delay = this.delayMs() + this.nextAttemptAt = now + delay + if (error instanceof MobileE2EEAuthenticationError) { + // Why: pairing state cannot change on a timer; polling only wakes the radio. + this.recoveryGate = 'external-signal' + this.clearTimer() + return + } + if (recovery?.kind === 'wait-for-host-revival') { + // Why: retrying HOST_OFFLINE without a revival signal is polling a known-negative state. + this.recoveryGate = 'external-signal' + this.clearTimer() + return + } + if (recovery?.kind === 'disable-relay-credential') { + // Why: a rejected outer credential cannot recover until direct connectivity refreshes it. + this.recoveryGate = 'fresh-credential' + this.clearTimer() + return + } + this.recoveryGate = null + if (!scheduleRetry) { + this.clearTimer() + return + } + this.scheduleRetry(delay) + } + + shouldTryGraceAfterRelayFailure(error: Error): boolean { + // Why: only a rejected outer credential can be repaired by the grace token; + // retrying session/capacity close codes immediately recreates relay churn. + return ( + error instanceof RelayOuterError && + error.code === MOBILE_RELAY_CLOSE_CODE.BAD_OUTER_CREDENTIAL + ) + } + + retryDelayMs(minimumMs: number): number | null { + if (this.recoveryGate) { + return null + } + return Math.max(minimumMs, this.nextAttemptAt - this.dependencies.now()) + } + + reset(): void { + this.consecutiveFailures = 0 + this.activeRelayConnectedAt = null + this.nextAttemptAt = 0 + this.recoveryGate = null + this.clearTimer() + } + + clear(): void { + this.clearTimer() + this.activeSession = null + this.activeRelayConnectedAt = null + } + + private clearTimer(): void { + if (this.timer) { + this.dependencies.clearTimer(this.timer) + this.timer = null + } + } + + private scheduleRetry(delayMs?: number): void { + if (this.timer) { + return + } + const delay = delayMs ?? Math.max(0, this.nextAttemptAt - this.dependencies.now()) + this.timer = this.dependencies.setTimer(() => { + this.timer = null + this.onRetry() + }, delay) + } + + private delayMs(): number { + const exponent = Math.max(0, this.consecutiveFailures - 1) + const cap = Math.min(RELAY_BACKOFF_CEILING_MS, RELAY_BACKOFF_BASE_MS * 2 ** exponent) + // Full jitter (uniform in [0, cap)), floored so retries never busy-loop. + return Math.max(RELAY_BACKOFF_MIN_MS, Math.floor(cap * this.jitterFraction())) + } + + private jitterFraction(): number { + const [high, low] = this.dependencies.randomBytes(2) + return (((high ?? 0) << 8) | (low ?? 0)) / 0x1_00_00 + } +}