mirror of
https://github.com/stablyai/orca.git
synced 2026-09-28 16:02:45 +00:00
merge PR 20061 relay-broker-wait-concurrency
This commit is contained in:
@@ -33,6 +33,12 @@ export type RelayAuthCoordinatorOptions = {
|
||||
refreshAccessToken: () => Promise<string | null>
|
||||
}) => Promise<CoordinatedRelayBroker>
|
||||
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
|
||||
}
|
||||
|
||||
@@ -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<RelayAuthContext>((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)
|
||||
})
|
||||
})
|
||||
@@ -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<BrokerOwnership>()
|
||||
private latestReconcile: Promise<void> = 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<void>()
|
||||
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<typeof setTimeout> | 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<void>()
|
||||
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<LiveBrokerWaitResult> {
|
||||
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 {
|
||||
|
||||
@@ -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<T>() {
|
||||
let resolve!: (value: T) => void
|
||||
const promise = new Promise<T>((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>): () => 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<CoordinatedRelayBroker>()
|
||||
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<CoordinatedRelayBroker>(), deferred<CoordinatedRelayBroker>()]
|
||||
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<CoordinatedRelayBroker>()
|
||||
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<CoordinatedRelayBroker>((_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()
|
||||
})
|
||||
})
|
||||
@@ -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<typeof vi.fn> }[]
|
||||
}))
|
||||
|
||||
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<T>() {
|
||||
let resolve!: (value: T) => void
|
||||
const promise = new Promise<T>((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<typeof vi.fn>
|
||||
createPairingRelay: (relayDeviceId: string) => Promise<unknown>
|
||||
}
|
||||
|
||||
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<FakeBroker>()
|
||||
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<never>()
|
||||
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()
|
||||
}
|
||||
})
|
||||
})
|
||||
@@ -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<typeof fixture>, 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)
|
||||
})
|
||||
})
|
||||
@@ -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<string, Readonly<TransientRef>> {
|
||||
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
|
||||
}
|
||||
|
||||
@@ -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<void>
|
||||
// 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<void>
|
||||
armedRetry: () => Promise<void> | null
|
||||
offlineReason: () => RelayOfflineReason | null
|
||||
}
|
||||
|
||||
export async function runLiveBrokerWait(
|
||||
source: LiveBrokerWaitSource,
|
||||
budgetMs: number
|
||||
): Promise<LiveBrokerWaitResult> {
|
||||
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() }
|
||||
}
|
||||
Reference in New Issue
Block a user