diff --git a/src/main/runtime/relay/relay-auth-coordinator-contract.ts b/src/main/runtime/relay/relay-auth-coordinator-contract.ts index 2e31ef899eb..ed15aad3b89 100644 --- a/src/main/runtime/relay/relay-auth-coordinator-contract.ts +++ b/src/main/runtime/relay/relay-auth-coordinator-contract.ts @@ -33,6 +33,12 @@ export type RelayAuthCoordinatorOptions = { refreshAccessToken: () => Promise }) => Promise onStatus: (status: RelayBrokerStatus, cellUrl?: string) => void + /** + * Test seam, not a policy knob. The only production construction is in `desktop-relay-service.ts` + * and it passes neither this nor `random`, so every shipped build lingers for the hardcoded + * ten-minute default in `relay-auth-coordinator.ts`. Read it as such before tuning it: changing + * the default is the only thing a user can feel. + */ lingerMs?: number random?: () => number } diff --git a/src/main/runtime/relay/relay-auth-coordinator-wait-budget.test.ts b/src/main/runtime/relay/relay-auth-coordinator-wait-budget.test.ts new file mode 100644 index 00000000000..e46431f42f6 --- /dev/null +++ b/src/main/runtime/relay/relay-auth-coordinator-wait-budget.test.ts @@ -0,0 +1,57 @@ +import { afterEach, describe, expect, it, vi } from 'vitest' +import { RelayAuthCoordinator, type RelayAuthContext } from './relay-auth-coordinator' + +const context: RelayAuthContext = { + identity: { userId: 'user-1', profileId: 'profile-1', organizationId: 'org-1' }, + accessToken: 'access-1', + relayEntitled: true +} + +/** + * What the live-broker wait budget does and does not bound. + * + * `LIVE_BROKER_WAIT_BUDGET_MS` is 20s and the deadline is checked only AFTER `await pending`, which + * is deliberately unbounded — cutting a slow-but-succeeding open short would fail a pairing that + * was about to work. So the budget bounds the armed-retry chain and nothing else. + * + * That is worth a test rather than a comment because the consequence lands on the phone: + * `pairing.provisionRelay` reaches this through `requireActiveBroker`, and `readContext`'s own + * ceiling is the cloud refresh timeout (60s, with one retry for a definitive 5xx) — several times + * the phone's request budget. The phone gives up and retries while the desktop is still holding a + * transient demand ref for the call it abandoned. + */ +describe('live-broker wait budget', () => { + afterEach(() => { + vi.useRealTimers() + }) + + it('does not bound a reconcile already in flight, so it can outlast the phone request budget', async () => { + vi.useFakeTimers() + let releaseContext = (): void => {} + const coordinator = new RelayAuthCoordinator({ + readContext: () => + new Promise((resolve) => { + releaseContext = (): void => resolve(context) + }), + openBroker: async () => ({ closeNow: vi.fn() }), + onStatus: vi.fn() + }) + coordinator.reconcile() + await Promise.resolve() + + let settled = false + const wait = coordinator.waitForLiveBrokerResult().then((result) => { + settled = true + return result + }) + + // Well past both the 20s budget and the phone's 30s request budget. + await vi.advanceTimersByTimeAsync(45_000) + expect(settled).toBe(false) + + releaseContext() + await vi.advanceTimersByTimeAsync(0) + await expect(wait).resolves.toEqual({ broker: expect.anything() }) + expect(settled).toBe(true) + }) +}) diff --git a/src/main/runtime/relay/relay-auth-coordinator.ts b/src/main/runtime/relay/relay-auth-coordinator.ts index 80cff8f878e..f26281f7a54 100644 --- a/src/main/runtime/relay/relay-auth-coordinator.ts +++ b/src/main/runtime/relay/relay-auth-coordinator.ts @@ -7,7 +7,11 @@ import type { RelayBrokerStatus } from './relay-session-broker' import { RelayHttpError, shouldRetryRelayConnectionError } from './relay-http-client' import { relayOfflineReasonForOpenFailure, type RelayOfflineReason } from './relay-offline-reason' import { RelayRetrySchedule } from './relay-retry-schedule' -import { withTimeout } from '../../../shared/promise-timeout-fallback' +import { + runLiveBrokerWait, + LIVE_BROKER_WAIT_BUDGET_MS, + type LiveBrokerWaitSource +} from './relay-live-broker-wait' import { relayAuthIdentityKey as identityKey, type CoordinatedRelayBroker, @@ -31,17 +35,25 @@ type BrokerOwnership = { } export class RelayAuthCoordinator { - // Why 20s: bounds only how long a waiter sits through armed retries, never - // an open already in flight. It spans the first few rungs of the backoff - // ladder and stays inside the phone's 30s request budget, so a sustained - // outage fails the caller with its cause instead of parking the demand ref. - private static readonly LIVE_BROKER_WAIT_BUDGET_MS = 20_000 private readonly options: RelayAuthCoordinatorOptions private authEpoch = 0 private offlineReason: RelayOfflineReason | null = null private ownership: BrokerOwnership | null = null private readonly pendingOwnerships = new Set() private latestReconcile: Promise = Promise.resolve() + // Why waiters are woken on every authority turnover (fresh reconcile, fence) + // rather than left on the reconcile they joined: that reconcile's result is + // discarded once it is superseded, so parking on it holds the caller behind + // an open nobody will use — even after a newer one registered a broker. + private authorityChange = Promise.withResolvers() + private readonly waitSource: LiveBrokerWaitSource = { + stopped: () => this.stopped, + liveBroker: () => this.getLiveBroker(), + reconcile: () => this.latestReconcile, + authorityChange: () => this.authorityChange.promise, + armedRetry: () => this.retry.settled, + offlineReason: () => this.offlineReason + } private lingerTimer: ReturnType | null = null private readonly retry: RelayRetrySchedule private stopped = false @@ -71,9 +83,16 @@ export class RelayAuthCoordinator { this.invalidatePendingOwnerships() const reconcile = this.reconcileEpoch(epoch, expectedIdentityKey, options) this.latestReconcile = reconcile + this.wakeWaiters() void reconcile } + private wakeWaiters(): void { + const change = this.authorityChange + this.authorityChange = Promise.withResolvers() + change.resolve() + } + // hostCloseReason names an auth loss the phone should be told about. Quit, // relaunch and every other fence pass nothing, so the control socket dies // abruptly exactly as before and the cell records no cause. @@ -85,6 +104,7 @@ export class RelayAuthCoordinator { this.invalidatePendingOwnerships() this.invalidateOwnership(hostCloseReason) this.publish('offline', hostCloseReason) + this.wakeWaiters() } // Why derived rather than passed in: the coordinator republishes `registered` @@ -136,38 +156,9 @@ export class RelayAuthCoordinator { // Why the caller maps offlineReason to a code instead of the coordinator // throwing: the mint-failure vocabulary belongs to the RPC layer. async waitForLiveBrokerResult( - budgetMs = RelayAuthCoordinator.LIVE_BROKER_WAIT_BUDGET_MS + budgetMs = LIVE_BROKER_WAIT_BUDGET_MS ): Promise { - const deadline = Date.now() + budgetMs - while (!this.stopped) { - const broker = this.getLiveBroker() - if (broker) { - return { broker } - } - const pending = this.latestReconcile - // Why unbounded: a reconcile always settles (opens carry HTTP deadlines), - // and cutting a slow-but-succeeding open short would fail a pairing that - // was about to work. The budget bounds only the retry chain below. - await pending - if (pending !== this.latestReconcile) { - continue - } - // Why: a reconcile that failed transiently has already armed its own - // retry; returning now would surface a hiccup fixed moments later. A - // terminal outcome (signed out, unentitled, rejected) arms nothing, so - // its cause returns without waiting. - const armed = this.retry.settled - if (!armed || Date.now() >= deadline) { - return this.settledLiveBrokerResult() - } - await withTimeout(armed, deadline - Date.now(), undefined) - } - return this.settledLiveBrokerResult() - } - - private settledLiveBrokerResult(): LiveBrokerWaitResult { - const broker = this.getLiveBroker() - return broker ? { broker } : { broker: null, offlineReason: this.offlineReason } + return await runLiveBrokerWait(this.waitSource, budgetMs) } stop(): void { diff --git a/src/main/runtime/relay/relay-concurrency-broker-wait.test.ts b/src/main/runtime/relay/relay-concurrency-broker-wait.test.ts new file mode 100644 index 00000000000..ba782a03862 --- /dev/null +++ b/src/main/runtime/relay/relay-concurrency-broker-wait.test.ts @@ -0,0 +1,256 @@ +import { afterEach, describe, expect, it, vi } from 'vitest' +import { + RelayAuthCoordinator, + type CoordinatedRelayBroker, + type LiveBrokerWaitResult, + type RelayAuthContext +} from './relay-auth-coordinator' +import { RelayHttpError } from './relay-http-client' + +const context: RelayAuthContext = { + identity: { userId: 'user-1', profileId: 'profile-1', organizationId: 'org-1' }, + accessToken: 'access-1', + relayEntitled: true +} + +function deferred() { + let resolve!: (value: T) => void + const promise = new Promise((resolvePromise) => { + resolve = resolvePromise + }) + return { promise, resolve } +} + +// Reads a wait without ever awaiting it, so a waiter that never settles fails an +// assertion instead of hanging out to the runner timeout — which is what every +// leak in this file looks like. +function observe(wait: Promise): () => LiveBrokerWaitResult | 'unsettled' { + let seen: LiveBrokerWaitResult | 'unsettled' = 'unsettled' + void wait.then((result) => { + seen = result + }) + return () => seen +} + +afterEach(() => { + vi.useRealTimers() +}) + +describe('relay live-broker wait under interleaving', () => { + it('hands a waiter the broker a newer reconcile registered instead of the abandoned open', async () => { + // A reconcile settles mid-wait. The waiter joined the previous one, whose + // result is discarded — parking on it held the caller behind an open nobody + // would use, for as long as that open took to time out. + vi.useFakeTimers() + const abandoned = deferred() + const fresh = { closeNow: vi.fn() } + let opens = 0 + const coordinator = new RelayAuthCoordinator({ + readContext: async () => context, + openBroker: async () => { + opens += 1 + return opens === 1 ? await abandoned.promise : fresh + }, + onStatus: vi.fn(), + random: () => 0.5 + }) + + coordinator.reconcile() + await vi.advanceTimersByTimeAsync(0) + expect(opens).toBe(1) + + const seen = observe(coordinator.waitForLiveBrokerResult(20_000)) + await vi.advanceTimersByTimeAsync(0) + expect(seen()).toBe('unsettled') + + coordinator.reconcile() + await vi.advanceTimersByTimeAsync(0) + expect(coordinator.getLiveBroker()).toBe(fresh) + + expect(seen()).toEqual({ broker: fresh }) + coordinator.stop() + }) + + it('follows a superseding reconcile that is still opening rather than answering from the old one', async () => { + // The other half of a mid-wait supersession: waking the waiter is only + // right if it then joins the newer reconcile. Answering as soon as it wakes + // would report "no broker, no cause" while the open that will supply one is + // still in flight. + vi.useFakeTimers() + const opens = [deferred(), deferred()] + const broker = { closeNow: vi.fn() } + let call = 0 + const coordinator = new RelayAuthCoordinator({ + readContext: async () => context, + openBroker: () => opens[call++]!.promise, + onStatus: vi.fn(), + random: () => 0.5 + }) + + coordinator.reconcile() + await vi.advanceTimersByTimeAsync(0) + const seen = observe(coordinator.waitForLiveBrokerResult(20_000)) + await vi.advanceTimersByTimeAsync(0) + + coordinator.reconcile() + await vi.advanceTimersByTimeAsync(0) + expect(call).toBe(2) + expect(seen()).toBe('unsettled') + + opens[1]!.resolve(broker) + await vi.advanceTimersByTimeAsync(0) + expect(seen()).toEqual({ broker }) + coordinator.stop() + }) + + it('does not let the first waiter off an armed retry strand or answer for the second', async () => { + // Two waiters share one armed schedule. The short-budget one giving up must + // neither cancel the retry nor leave the long-budget one on a signal that + // has already fired. + vi.useFakeTimers() + const broker = { closeNow: vi.fn() } + const openBroker = vi + .fn() + .mockRejectedValueOnce(new RelayHttpError('assignment', 500)) + .mockResolvedValue(broker) + const coordinator = new RelayAuthCoordinator({ + readContext: async () => context, + openBroker, + onStatus: vi.fn(), + random: () => 0.5 + }) + coordinator.reconcile() + await vi.advanceTimersByTimeAsync(0) + + const short = observe(coordinator.waitForLiveBrokerResult(100)) + const long = observe(coordinator.waitForLiveBrokerResult(20_000)) + await vi.advanceTimersByTimeAsync(100) + expect(short()).toEqual({ broker: null, offlineReason: 'broker_unavailable' }) + expect(long()).toBe('unsettled') + + // The retry the short waiter walked away from still fires, and its result + // reaches the waiter that stayed. + await vi.advanceTimersByTimeAsync(400) + expect(openBroker).toHaveBeenCalledTimes(2) + expect(long()).toEqual({ broker }) + coordinator.stop() + }) + + it('lets the armed retry take the tie when it fires on the tick the budget expires', async () => { + // attempt 0 with random 0.5 is a 500ms delay; the budget matches it exactly. + // The retry wins the tick, so the caller is told the fresh attempt's cause + // rather than the generic no-cause result. + vi.useFakeTimers() + const openBroker = vi.fn().mockRejectedValue(new RelayHttpError('assignment', 500)) + const coordinator = new RelayAuthCoordinator({ + readContext: async () => context, + openBroker, + onStatus: vi.fn(), + random: () => 0.5 + }) + coordinator.reconcile() + await vi.advanceTimersByTimeAsync(0) + + const seen = observe(coordinator.waitForLiveBrokerResult(500)) + await vi.advanceTimersByTimeAsync(500) + expect(openBroker).toHaveBeenCalledTimes(2) + expect(seen()).toEqual({ broker: null, offlineReason: 'broker_unavailable' }) + coordinator.stop() + }) + + it('settles a waiter parked on an armed retry when the coordinator stops, leaving no timer', async () => { + vi.useFakeTimers() + const openBroker = vi.fn().mockRejectedValue(new RelayHttpError('assignment', 500)) + const coordinator = new RelayAuthCoordinator({ + readContext: async () => context, + openBroker, + onStatus: vi.fn(), + random: () => 0.5 + }) + coordinator.reconcile() + await vi.advanceTimersByTimeAsync(0) + const seen = observe(coordinator.waitForLiveBrokerResult(20_000)) + await vi.advanceTimersByTimeAsync(0) + expect(vi.getTimerCount()).toBeGreaterThan(0) + + coordinator.stop() + await vi.advanceTimersByTimeAsync(0) + expect(seen()).toEqual({ broker: null, offlineReason: null }) + expect(vi.getTimerCount()).toBe(0) + }) + + it('settles a waiter parked on an open that never returns when the coordinator stops', async () => { + // The wait is deliberately unbounded across the open it joined, so nothing + // else can release this waiter; a teardown that did not wake it leaked the + // promise for the life of the process. + vi.useFakeTimers() + const stuck = deferred() + const coordinator = new RelayAuthCoordinator({ + readContext: async () => context, + openBroker: () => stuck.promise, + onStatus: vi.fn(), + random: () => 0.5 + }) + coordinator.reconcile() + await vi.advanceTimersByTimeAsync(0) + const seen = observe(coordinator.waitForLiveBrokerResult(1_000)) + await vi.advanceTimersByTimeAsync(0) + expect(seen()).toBe('unsettled') + + coordinator.stop() + await vi.advanceTimersByTimeAsync(0) + expect(seen()).toEqual({ broker: null, offlineReason: null }) + expect(vi.getTimerCount()).toBe(0) + }) + + it('gives a waiter the terminal cause that landed while its retry was armed', async () => { + vi.useFakeTimers() + let current: RelayAuthContext | null = context + const openBroker = vi.fn().mockRejectedValue(new RelayHttpError('assignment', 500)) + const coordinator = new RelayAuthCoordinator({ + readContext: async () => current, + openBroker, + onStatus: vi.fn(), + random: () => 0.5 + }) + coordinator.reconcile() + await vi.advanceTimersByTimeAsync(0) + const seen = observe(coordinator.waitForLiveBrokerResult(20_000)) + await vi.advanceTimersByTimeAsync(0) + expect(seen()).toBe('unsettled') + + current = null + coordinator.reconcile() + await vi.advanceTimersByTimeAsync(0) + expect(seen()).toEqual({ broker: null, offlineReason: 'signed-out' }) + + // The retryable schedule the terminal cause superseded must arm nothing. + await vi.advanceTimersByTimeAsync(10 * 60_000) + expect(openBroker).toHaveBeenCalledTimes(1) + coordinator.stop() + }) + + it('keeps the budget over a chain of superseding reconciles, not just over armed retries', async () => { + // Every transient-demand acquire and release reconciles, so a busy host can + // supersede a waiter's reconcile faster than opens settle. The budget has to + // survive that or a caller rides the chain indefinitely holding its demand ref. + vi.useFakeTimers() + const coordinator = new RelayAuthCoordinator({ + readContext: async () => context, + openBroker: () => + new Promise((_resolve, reject) => + setTimeout(() => reject(new Error('slow open')), 400) + ), + onStatus: vi.fn(), + random: () => 0.5 + }) + coordinator.reconcile() + const seen = observe(coordinator.waitForLiveBrokerResult(1_000)) + const churn = setInterval(() => coordinator.reconcile(), 200) + + await vi.advanceTimersByTimeAsync(2_000) + clearInterval(churn) + expect(seen()).not.toBe('unsettled') + coordinator.stop() + }) +}) diff --git a/src/main/runtime/relay/relay-concurrency-policy-flip-mid-mint.test.ts b/src/main/runtime/relay/relay-concurrency-policy-flip-mid-mint.test.ts new file mode 100644 index 00000000000..17760c23ecb --- /dev/null +++ b/src/main/runtime/relay/relay-concurrency-policy-flip-mid-mint.test.ts @@ -0,0 +1,180 @@ +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' +import type { OrcaCloudAuthConfig } from '../../orca-profiles/profile-cloud-auth-config' +import type { MobilePairingConnectionMode } from '../../../shared/mobile-pairing-connection-mode' +import type { OrcaRuntimeRpcServer } from '../runtime-rpc' + +// The policy-flip fix is normally exercised with the flip sequenced before the +// mint. These cover the interleaving it was written for: the flip lands while +// the mint is parked at an await, so the coordinator reaches `standby` and +// clears the offline reason out from under the waiter. +const fakes = vi.hoisted(() => ({ + readRelayAuthContext: vi.fn(), + connect: vi.fn(), + brokers: [] as { closeNow: ReturnType }[] +})) + +vi.mock('./relay-auth-context', () => ({ readRelayAuthContext: fakes.readRelayAuthContext })) + +vi.mock('./relay-session-broker', () => { + class RelaySessionBroker { + readonly hostId = 'relay-host-1' + readonly ownerIdentityKey = 'user-1\0profile-1\0org-1' + readonly endpoint = { v: 1 as const, relayHostId: 'relay-host-1' } + readonly closeNow = vi.fn() + createPairingRelay = vi.fn(async (relayDeviceId: string) => ({ + v: 1, + relayHostId: this.hostId, + relayDeviceId, + inviteExpiresAt: 0 + })) + isLive(): boolean { + return this.closeNow.mock.calls.length === 0 + } + static connect = fakes.connect + } + return { RelaySessionBroker } +}) + +import { DesktopRelayService } from './desktop-relay-service' +import { RelaySessionBroker } from './relay-session-broker' + +function deferred() { + let resolve!: (value: T) => void + const promise = new Promise((resolvePromise) => { + resolve = resolvePromise + }) + return { promise, resolve } +} + +// The real broker's constructor is private, so the mocked class is the only way +// to produce something `instanceof RelaySessionBroker` for the service's gate. +type FakeBroker = { + readonly hostId: string + readonly closeNow: ReturnType + createPairingRelay: (relayDeviceId: string) => Promise +} + +function newBroker(): FakeBroker { + const broker = new (RelaySessionBroker as unknown as new () => FakeBroker)() + fakes.brokers.push(broker) + return broker +} + +function service(mode: { current: MobilePairingConnectionMode }): DesktopRelayService { + fakes.readRelayAuthContext.mockResolvedValue({ + identity: { userId: 'user-1', profileId: 'profile-1', organizationId: 'org-1' }, + accessToken: 'access-1', + relayEntitled: true + }) + const runtimeRpc = { + getE2EEKeypair: () => ({ + publicKey: new Uint8Array(32).fill(7), + secretKey: new Uint8Array(32).fill(9), + publicKeyB64: 'x' + }), + getMobileSocketWiring: () => ({ attachTransport: () => () => {} }), + getRelayRevokeOutbox: () => ({ pendingFor: () => [], remove: vi.fn() }), + getDeviceRegistry: () => ({ + listDevices: () => [], + getDevice: () => ({ deviceId: 'device-1', scope: 'mobile' }), + getMobilePairingConnectionMode: () => 'automatic' + }) + } as unknown as OrcaRuntimeRpcServer + return new DesktopRelayService({ + authConfig: { + relayDirectorUrl: 'https://relay.example.test', + relayTokenEndpoint: 'https://login.example.test/relay-token' + } as OrcaCloudAuthConfig, + userDataPath: '/tmp/orca-relay-policy-flip-interleaving', + appVersion: '1.4.188', + runtimeRpc, + onStatus: () => {}, + hostMobilePairingConnectionMode: () => mode.current + }) +} + +beforeEach(() => { + vi.useFakeTimers() + fakes.brokers.length = 0 + fakes.connect.mockReset() +}) + +afterEach(() => { + vi.useRealTimers() +}) + +describe('DesktopRelayService LAN flip landing mid-mint', () => { + it('names the flip when it lands while the mint is parked on an in-flight open', async () => { + const mode = { current: 'automatic' as MobilePairingConnectionMode } + const open = deferred() + fakes.connect.mockImplementation(() => open.promise) + const relayService = service(mode) + try { + const mint = relayService.createPairingRelay('device-1') + const rejection = expect(mint).rejects.toThrow('relay_disabled_for_device') + await vi.advanceTimersByTimeAsync(0) + expect(fakes.connect).toHaveBeenCalled() + + mode.current = 'local-only' + relayService.pairingPolicyChanged() + await vi.advanceTimersByTimeAsync(0) + await rejection + + // The abandoned open must still be discarded rather than left registered. + open.resolve(newBroker()) + await vi.advanceTimersByTimeAsync(0) + expect(fakes.brokers[0]!.closeNow).toHaveBeenCalled() + } finally { + relayService.stop() + } + }) + + it('names the flip when it lands after the broker is live, mid create_pairing_relay', async () => { + const mode = { current: 'automatic' as MobilePairingConnectionMode } + const broker = newBroker() + const mintCall = deferred() + const relayService = service(mode) + broker.createPairingRelay = vi.fn(async () => { + mode.current = 'local-only' + relayService.pairingPolicyChanged() + await mintCall.promise + throw new Error('unreachable') + }) + fakes.connect.mockResolvedValue(broker) + try { + const mint = relayService.createPairingRelay('device-1') + const rejection = expect(mint).rejects.toThrow('relay_disabled_for_device') + await vi.advanceTimersByTimeAsync(0) + expect(broker.createPairingRelay).toHaveBeenCalled() + + // The control the flip closed under it fails the call that was in flight. + mintCall.resolve(undefined as never) + await vi.advanceTimersByTimeAsync(0) + await rejection + } finally { + relayService.stop() + } + }) + + it('leaves the real failure alone when the flip is reverted before the mint fails', async () => { + // The rewrite reads the live policy at the moment of failure, so a flip that + // is undone in the same window must not claim a failure it did not cause. + const mode = { current: 'automatic' as MobilePairingConnectionMode } + const broker = newBroker() + broker.createPairingRelay = vi.fn(async () => { + mode.current = 'local-only' + await Promise.resolve() + mode.current = 'automatic' + throw new Error('relay_assignment_failed_500') + }) + fakes.connect.mockResolvedValue(broker) + const relayService = service(mode) + try { + await expect(relayService.createPairingRelay('device-1')).rejects.toThrow( + 'relay_assignment_failed_500' + ) + } finally { + relayService.stop() + } + }) +}) diff --git a/src/main/runtime/relay/relay-demand-ledger-lifecycle.test.ts b/src/main/runtime/relay/relay-demand-ledger-lifecycle.test.ts new file mode 100644 index 00000000000..1f03a859bbc --- /dev/null +++ b/src/main/runtime/relay/relay-demand-ledger-lifecycle.test.ts @@ -0,0 +1,248 @@ +import { mkdtempSync } from 'node:fs' +import { tmpdir } from 'node:os' +import { join } from 'node:path' +import { afterEach, describe, expect, it, vi } from 'vitest' +import { DeviceRegistry } from '../device-registry' +import type { MobilePairingConnectionMode } from '../../../shared/mobile-pairing-connection-mode' +import { isMobileRelayAllowed } from '../../../shared/mobile-relay-policy' +import { RelayDemandLedger } from './relay-demand-ledger' +import { RelayRevokeOutbox, type RelayDeviceBinding } from './relay-revoke-outbox' + +const ownerIdentityKey = 'user-1\0profile-1\0org-1' +const otherOwnerIdentityKey = 'user-2\0profile-2\0org-2' +const relayHostId = 'relay-host-1' + +/** Mirrors DesktopRelayService.isRelayAllowedForDevice: a live pull, not a snapshot. */ +function fixture(now = 1_000) { + const userDataPath = mkdtempSync(join(tmpdir(), 'orca-relay-demand-lifecycle-')) + const deviceRegistry = new DeviceRegistry(userDataPath) + const revokeOutbox = new RelayRevokeOutbox(userDataPath) + const host = { mode: 'automatic' as MobilePairingConnectionMode } + const ledger = new RelayDemandLedger({ + deviceRegistry, + revokeOutbox, + relayHostId, + now: () => now, + isRelayAllowedForDevice: (deviceId) => + isMobileRelayAllowed({ + hostConnectionMode: host.mode, + deviceConnectionMode: deviceRegistry.getMobilePairingConnectionMode(deviceId) ?? null + }) + }) + return { userDataPath, deviceRegistry, revokeOutbox, ledger, host } +} + +function binding(relayDeviceId: string, inviteExpiresAt?: number): RelayDeviceBinding { + return { relayDeviceId, relayHostId, ownerIdentityKey, inviteExpiresAt } +} + +function pairedPhone(fx: ReturnType, name = 'Phone') { + const phone = fx.deviceRegistry.addDevice(name) + fx.deviceRegistry.setMobilePairingConnectionMode(phone.deviceId, 'automatic') + fx.deviceRegistry.setRelayBinding(phone.deviceId, binding(phone.deviceId)) + fx.deviceRegistry.updateLastSeen(phone.deviceId) + return phone.deviceId +} + +describe('RelayDemandLedger transient-ref lifecycle', () => { + afterEach(() => vi.useRealTimers()) + + // Hazard 1: the live policy flips between acquire and release. + it('releases a ref acquired before a policy flip instead of leaking it', () => { + const fx = fixture() + const deviceId = pairedPhone(fx) + const release = fx.ledger.acquireTransient(`endpoints:${deviceId}`, deviceId) + expect(fx.ledger.hasDemand(ownerIdentityKey)).toBe(true) + + fx.host.mode = 'local-only' + expect(fx.ledger.hasDemand(ownerIdentityKey)).toBe(false) + + // The release path must be policy-blind, or the ref is unreachable forever. + release() + expect(fx.ledger.transientRefsForTests().size).toBe(0) + + fx.host.mode = 'automatic' + // Standing binding demand returns; the released ref must not add a second pin. + expect(fx.ledger.hasDemand(ownerIdentityKey)).toBe(true) + fx.deviceRegistry.removeDevice(deviceId) + expect(fx.ledger.hasDemand(ownerIdentityKey)).toBe(false) + }) + + it('restores a held ref as demand when the policy flips back', () => { + const fx = fixture() + const deviceId = pairedPhone(fx) + fx.deviceRegistry.removeDevice(deviceId) + const release = fx.ledger.acquireTransient(`pairing:${deviceId}`, deviceId) + + fx.host.mode = 'local-only' + expect(fx.ledger.hasDemand(ownerIdentityKey)).toBe(false) + // The filter must be re-evaluated per call, not baked in at acquire time. + fx.host.mode = 'automatic' + expect(fx.ledger.hasDemand(ownerIdentityKey)).toBe(true) + release() + expect(fx.ledger.hasDemand(ownerIdentityKey)).toBe(false) + }) + + // Hazard 2: the same release closure invoked repeatedly. + it('a repeated release does not decrement a count another holder still owns', () => { + const fx = fixture() + const deviceId = pairedPhone(fx) + fx.deviceRegistry.removeDevice(deviceId) + const key = `endpoints:${deviceId}` + const releaseFirst = fx.ledger.acquireTransient(key, deviceId) + const releaseSecond = fx.ledger.acquireTransient(key, deviceId) + expect(fx.ledger.transientRefsForTests().get(key)?.count).toBe(2) + + releaseFirst() + releaseFirst() + releaseFirst() + expect(fx.ledger.transientRefsForTests().get(key)?.count).toBe(1) + expect(fx.ledger.hasDemand(ownerIdentityKey)).toBe(true) + + releaseSecond() + expect(fx.ledger.transientRefsForTests().size).toBe(0) + expect(fx.ledger.hasDemand(ownerIdentityKey)).toBe(false) + }) + + it('a stale release from a finished ref does not cancel a later acquire of the same key', () => { + const fx = fixture() + const deviceId = pairedPhone(fx) + fx.deviceRegistry.removeDevice(deviceId) + const key = `provision:${deviceId}` + const stale = fx.ledger.acquireTransient(key, deviceId) + stale() + expect(fx.ledger.transientRefsForTests().size).toBe(0) + + const fresh = fx.ledger.acquireTransient(key, deviceId) + stale() + expect(fx.ledger.transientRefsForTests().get(key)?.count).toBe(1) + expect(fx.ledger.hasDemand(ownerIdentityKey)).toBe(true) + fresh() + expect(fx.ledger.transientRefsForTests().size).toBe(0) + }) + + // Hazard 3: an operation throws and its ref is never released. + it('an unreleased ref pins demand and retains its map entry', () => { + const fx = fixture() + const deviceId = pairedPhone(fx) + fx.deviceRegistry.removeDevice(deviceId) + for (let index = 0; index < 5; index += 1) { + fx.ledger.acquireTransient(`pairing:abandoned-${index}`, `abandoned-${index}`) + } + expect(fx.ledger.transientRefsForTests().size).toBe(5) + expect(fx.ledger.hasDemand(ownerIdentityKey)).toBe(true) + }) + + // Hazard 4: two refs for one device released in the opposite order. + it('is release-order independent for refs sharing a key', () => { + const fx = fixture() + const deviceId = pairedPhone(fx) + fx.deviceRegistry.removeDevice(deviceId) + const key = `endpoints:${deviceId}` + const outer = fx.ledger.acquireTransient(key, deviceId) + const inner = fx.ledger.acquireTransient(key, deviceId) + inner() + expect(fx.ledger.hasDemand(ownerIdentityKey)).toBe(true) + const middle = fx.ledger.acquireTransient(key, deviceId) + outer() + expect(fx.ledger.transientRefsForTests().get(key)?.count).toBe(1) + middle() + expect(fx.ledger.transientRefsForTests().size).toBe(0) + }) + + // Hazard 5: a release for device B must not disturb device A's ref. + it('keeps per-device refs independent across a concurrent release', () => { + const fx = fixture() + const deviceA = fx.deviceRegistry.addDevice('Phone A').deviceId + const deviceB = fx.deviceRegistry.addDevice('Phone B').deviceId + fx.deviceRegistry.setMobilePairingConnectionMode(deviceA, 'local-only') + fx.deviceRegistry.setMobilePairingConnectionMode(deviceB, 'automatic') + + const releaseA = fx.ledger.acquireTransient(`pairing:${deviceA}`, deviceA) + const releaseB = fx.ledger.acquireTransient(`pairing:${deviceB}`, deviceB) + // A is excluded by its own pair-time mode; B alone carries the demand. + expect(fx.ledger.hasDemand(ownerIdentityKey)).toBe(true) + + releaseB() + expect(fx.ledger.transientRefsForTests().size).toBe(1) + expect(fx.ledger.transientRefsForTests().get(`pairing:${deviceA}`)?.deviceId).toBe(deviceA) + expect(fx.ledger.hasDemand(ownerIdentityKey)).toBe(false) + + releaseA() + expect(fx.ledger.transientRefsForTests().size).toBe(0) + }) + + // Hazard 6: the ledger owns no timers, so teardown is the caller's map only. + it('arms no timers of its own across acquire, query and release', () => { + vi.useFakeTimers() + const fx = fixture() + const deviceId = pairedPhone(fx) + fx.deviceRegistry.setRelayBinding(deviceId, binding(deviceId, 2_000)) + const release = fx.ledger.acquireTransient(`pairing:${deviceId}`, deviceId) + expect(fx.ledger.hasDemand(ownerIdentityKey)).toBe(true) + expect(fx.ledger.nextPendingExpiry()).toBe(2_000) + release() + expect(vi.getTimerCount()).toBe(0) + }) +}) + +describe('RelayDemandLedger nextPendingExpiry scoping', () => { + it('ignores a pending invite bound to a different relay host', () => { + const fx = fixture() + const phone = fx.deviceRegistry.addDevice('Foreign-host phone') + fx.deviceRegistry.setRelayBinding(phone.deviceId, { + ...binding(phone.deviceId, 2_000), + relayHostId: 'relay-host-2' + }) + expect(fx.ledger.hasDemand(ownerIdentityKey)).toBe(false) + expect(fx.ledger.nextPendingExpiry()).toBeNull() + }) + + it('ignores a pending invite on a non-mobile device', () => { + const fx = fixture() + const runtime = fx.deviceRegistry.addDevice('Runtime', 'runtime') + fx.deviceRegistry.setRelayBinding(runtime.deviceId, binding(runtime.deviceId, 2_000)) + expect(fx.ledger.hasDemand(ownerIdentityKey)).toBe(false) + expect(fx.ledger.nextPendingExpiry()).toBeNull() + }) + + it('ignores a pending invite the live policy excludes, and re-arms when it flips back', () => { + const fx = fixture() + const phone = fx.deviceRegistry.addDevice('LAN phone') + fx.deviceRegistry.setMobilePairingConnectionMode(phone.deviceId, 'automatic') + fx.deviceRegistry.setRelayBinding(phone.deviceId, binding(phone.deviceId, 2_000)) + expect(fx.ledger.nextPendingExpiry()).toBe(2_000) + + fx.host.mode = 'local-only' + expect(fx.ledger.hasDemand(ownerIdentityKey)).toBe(false) + expect(fx.ledger.nextPendingExpiry()).toBeNull() + + fx.host.mode = 'automatic' + expect(fx.ledger.nextPendingExpiry()).toBe(2_000) + }) + + it('still reports the nearest in-scope pending invite', () => { + const fx = fixture() + const near = fx.deviceRegistry.addDevice('Near phone') + const far = fx.deviceRegistry.addDevice('Far phone') + fx.deviceRegistry.setRelayBinding(far.deviceId, binding(far.deviceId, 9_000)) + fx.deviceRegistry.setRelayBinding(near.deviceId, binding(near.deviceId, 4_000)) + expect(fx.ledger.nextPendingExpiry()).toBe(4_000) + }) +}) + +describe('RelayDemandLedger owner scoping of transient refs', () => { + // Characterisation: transient refs carry no owner identity, so they answer + // hasDemand for any signed-in identity. See the report for the risk window. + it('counts a transient ref for an identity that did not request it', () => { + const fx = fixture() + const phone = fx.deviceRegistry.addDevice('Phone') + fx.deviceRegistry.setMobilePairingConnectionMode(phone.deviceId, 'automatic') + fx.deviceRegistry.setRelayBinding(phone.deviceId, binding(phone.deviceId)) + const release = fx.ledger.acquireTransient(`pairing:${phone.deviceId}`, phone.deviceId) + expect(fx.ledger.hasDemand(ownerIdentityKey)).toBe(true) + expect(fx.ledger.hasDemand(otherOwnerIdentityKey)).toBe(true) + release() + expect(fx.ledger.hasDemand(otherOwnerIdentityKey)).toBe(false) + }) +}) diff --git a/src/main/runtime/relay/relay-demand-ledger.ts b/src/main/runtime/relay/relay-demand-ledger.ts index 817fbc37e37..2e44e5497f5 100644 --- a/src/main/runtime/relay/relay-demand-ledger.ts +++ b/src/main/runtime/relay/relay-demand-ledger.ts @@ -1,5 +1,5 @@ -import type { DeviceRegistry } from '../device-registry' -import type { RelayRevokeOutbox } from './relay-revoke-outbox' +import type { DeviceEntry, DeviceRegistry } from '../device-registry' +import type { RelayDeviceBinding, RelayRevokeOutbox } from './relay-revoke-outbox' type RelayDemandLedgerOptions = { deviceRegistry: DeviceRegistry @@ -35,6 +35,10 @@ export class RelayDemandLedger { } released = true const current = this.transientRefs.get(key) + // Keep this guard even though nothing can reach it today: `released` makes the lookup happen + // at most once per closure and no public path empties the map, so mutating it away survives + // the suite. It goes load-bearing the moment the ledger grows a dispose()/clear(), and + // nothing would catch that. if (!current) { return } @@ -47,6 +51,20 @@ export class RelayDemandLedger { } hasDemand(ownerIdentityKey: string): boolean { + // KNOWN GAP, deliberately not fixed here: a transient ref carries no owner identity, so this + // loop answers true for ANY signed-in identity. The device-binding and revoke-outbox branches + // below both filter on `ownerIdentityKey`; this one cannot. A profile/org switch mid-pairing + // therefore keeps the coordinator holding a relay control session for an identity that never + // asked for one. + // + // Nothing defends either the current behaviour or that regression: scoping this loop by owner + // fails exactly one test across the whole relay suite, the characterisation test written for + // it. The fix has to thread identity through `acquireTransient`, whose call site is + // desktop-relay-service.ts. Inferring the owner here instead — skipping a ref whose device + // carries a different owner's binding — is UNSAFE: `withTransientDemand('provision')` calls + // `setMobileRelayBinding` inside the operation, so during a re-pair the device still holds the + // old owner's binding and that filter would drop demand mid-provision, tearing the broker down + // under the very operation holding the ref. for (const ref of this.transientRefs.values()) { if (this.isRelayAllowed(ref.deviceId)) { return true @@ -57,14 +75,8 @@ export class RelayDemandLedger { } const now = (this.options.now ?? Date.now)() return this.options.deviceRegistry.listDevices().some((device) => { - const binding = device.relayBinding - if ( - device.scope !== 'mobile' || - !binding || - binding.ownerIdentityKey !== ownerIdentityKey || - binding.relayHostId !== this.options.relayHostId || - !this.isRelayAllowed(device.deviceId) - ) { + const binding = this.demandCandidateBinding(device) + if (!binding || binding.ownerIdentityKey !== ownerIdentityKey) { return false } // Why: E2EE authentication marks a scanned DeviceEntry as seen before @@ -78,7 +90,9 @@ export class RelayDemandLedger { const now = (this.options.now ?? Date.now)() let next: number | null = null for (const device of this.options.deviceRegistry.listDevices()) { - const expiresAt = device.relayBinding?.inviteExpiresAt + // Why: an expiry that hasDemand would never look at still armed a wake + // timer — another host's binding, a runtime device, a LAN-excluded phone. + const expiresAt = this.demandCandidateBinding(device)?.inviteExpiresAt if (expiresAt && expiresAt > now && (next === null || expiresAt < next)) { next = expiresAt } @@ -86,6 +100,25 @@ export class RelayDemandLedger { return next } + /** Test-only view of the outstanding transient refs. */ + transientRefsForTests(): ReadonlyMap> { + return this.transientRefs + } + + /** The binding when this device could stand as demand on this relay host, else null. */ + private demandCandidateBinding(device: DeviceEntry): RelayDeviceBinding | null { + const binding = device.relayBinding + if ( + device.scope !== 'mobile' || + !binding || + binding.relayHostId !== this.options.relayHostId || + !this.isRelayAllowed(device.deviceId) + ) { + return null + } + return binding + } + private isRelayAllowed(deviceId: string): boolean { return this.options.isRelayAllowedForDevice?.(deviceId) ?? true } diff --git a/src/main/runtime/relay/relay-live-broker-wait.ts b/src/main/runtime/relay/relay-live-broker-wait.ts new file mode 100644 index 00000000000..44b3db78dff --- /dev/null +++ b/src/main/runtime/relay/relay-live-broker-wait.ts @@ -0,0 +1,80 @@ +import { withTimeout } from '../../../shared/promise-timeout-fallback' +import type { + CoordinatedRelayBroker, + LiveBrokerWaitResult +} from './relay-auth-coordinator-contract' +import type { RelayOfflineReason } from './relay-offline-reason' + +// Why 20s: bounds only how long a waiter sits through armed retries and superseded opens, never +// the open it arrived on. It spans the first few rungs of the backoff ladder, so a sustained +// outage fails the caller with its cause instead of parking the demand ref. +// +// It does NOT bound the call, and an earlier version of this comment claimed it "stays inside the +// phone's 30s request budget". It cannot: the deadline is consulted only after the await below, +// and a reconcile's own ceiling is `readContext`'s cloud-refresh timeout (60s, plus one retry for +// a definitive 5xx) followed by the broker open. Measured: with a reconcile in flight, a wait at +// this default budget had not settled at 45s. `pairing.provisionRelay` reaches here through +// `requireActiveBroker`, so the phone gives up and retries while the desktop still holds a +// transient demand ref for the call it abandoned. relay-auth-coordinator-wait-budget.test.ts pins +// that so the claim cannot drift back. +export const LIVE_BROKER_WAIT_BUDGET_MS = 20_000 + +// Every member is a function because the wait re-reads all of it after each +// await; a snapshot would answer for a coordinator state that is already gone. +export type LiveBrokerWaitSource = { + stopped: () => boolean + liveBroker: () => CoordinatedRelayBroker | null + reconcile: () => Promise + // Resolves when a fresh reconcile or a fence turns the authority over, so a + // waiter is not left parked on a reconcile whose result is already discarded. + authorityChange: () => Promise + armedRetry: () => Promise | null + offlineReason: () => RelayOfflineReason | null +} + +export async function runLiveBrokerWait( + source: LiveBrokerWaitSource, + budgetMs: number +): Promise { + const deadline = Date.now() + budgetMs + let joinedReconcile = false + while (!source.stopped()) { + const broker = source.liveBroker() + if (broker) { + return { broker } + } + // Why the budget applies only once a reconcile has been joined: the one open + // the waiter arrived on is never cut short, but a chain of superseding opens + // must not outlive the budget the caller asked for. + if (joinedReconcile && Date.now() >= deadline) { + return settledResult(source) + } + const pending = source.reconcile() + const superseded = source.authorityChange() + joinedReconcile = true + // Why unbounded on `pending`: a reconcile always settles — BOTH its awaits carry a deadline, + // the broker open and `readContext`'s cloud refresh (the larger of the two, and the one an + // earlier parenthetical here left out) — and cutting a slow-but-succeeding open short would + // fail a pairing that was about to work. "Settles" is not "settles soon": see the ceiling + // noted on LIVE_BROKER_WAIT_BUDGET_MS. + await Promise.race([pending, superseded]) + if (pending !== source.reconcile()) { + continue + } + // Why: a reconcile that failed transiently has already armed its own retry; + // returning now would surface a hiccup fixed moments later. A terminal + // outcome (signed out, unentitled, rejected) arms nothing, so its cause + // returns without waiting. + const armed = source.armedRetry() + if (!armed || Date.now() >= deadline) { + return settledResult(source) + } + await withTimeout(armed, deadline - Date.now(), undefined) + } + return settledResult(source) +} + +function settledResult(source: LiveBrokerWaitSource): LiveBrokerWaitResult { + const broker = source.liveBroker() + return broker ? { broker } : { broker: null, offlineReason: source.offlineReason() } +}