diff --git a/cloud/apps/relay/src/host-session-client-accept.test.ts b/cloud/apps/relay/src/host-session-client-accept.test.ts new file mode 100644 index 00000000000..0155c91f80a --- /dev/null +++ b/cloud/apps/relay/src/host-session-client-accept.test.ts @@ -0,0 +1,304 @@ +import { EventEmitter } from 'node:events' +import { RELAY_CLOSE_CODE } from '@orca-cloud/relay-contract' +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' +import type WebSocket from 'ws' +import type { RelayAssignmentStore } from './assignment-store.js' +import type { RelayConfig } from './config.js' +import type { CredentialReservation, RelayCredentialStore } from './credential-store.js' +import { + CONTROL_LEASE_JITTER_MS, + CONTROL_LEASE_MS, + HostSessionRegistry +} from './host-session-registry.js' +import type { RelayRuntimeObserver } from './relay-observability.js' +import type { RelayTokenClaims } from './relay-token-verifier.js' +import { ProcessQueuedByteBudget } from './splice-forwarder.js' + +// Incident 2026-09-04 ~01:05Z: the phone's dial bound ran out while the cell was +// still inside acceptClient's serialized Postgres phase (cell-inventory lock +// contention). The cell then finished the work for a socket nobody held, leaked +// a 90s activity lease, and logged `host_data_reservation_already_bound`. + +class FakeSocket extends EventEmitter { + readonly OPEN = 1 + readonly CLOSING = 2 + readonly CLOSED = 3 + readyState = this.OPEN + readonly send = vi.fn() + readonly close = vi.fn((code?: number, reason?: string) => { + this.readyState = this.CLOSED + this.emit('close', code, Buffer.from(reason ?? '')) + }) + readonly terminate = vi.fn(() => { + this.readyState = this.CLOSED + this.emit('close') + }) +} + +const config = { + port: 8080, + publicUrl: 'https://relay-c3.example.com', + cellUrl: 'https://relay-c3.example.com', + authIssuer: 'https://auth.example.com', + authAudience: 'orca-relay', + jwksUrl: 'https://auth.example.com/jwks', + assignmentSigningKey: new Uint8Array(32), + role: 'cell', + cellId: 'production-gce-c3', + cells: [{ id: 'production-gce-c3', url: 'https://relay-c3.example.com', capacityRequests: 4_000 }], + adminAudience: 'https://relay-c3.example.com/v1/admin/drain', + deployServiceAccount: 'deploy@example.com', + runtimeServiceAccount: 'runtime@example.com', + adminJwksUrl: 'https://auth.example.com/admin-jwks', + databasePoolMax: 10, + publicAssignmentsEnabled: true, + publicAssignmentConcurrency: 2, + publicAssignmentQueueMax: 128, + publicAssignmentWaitMs: 4_000, + publicResolveConcurrency: 1, + publicResolveWaitMs: 5_000, + publicAssignmentRetryAfterSeconds: 5, + dataDir: './test-data' +} satisfies RelayConfig + +const identity = { + sub: 'user-1', + prof: 'profile-1', + relayHostId: 'abcdefghijklmnop', + purpose: 'host-control', + exp: 4_102_444_800 +} satisfies RelayTokenClaims + +function deferred(): { promise: Promise; resolve: (value: T) => void } { + let resolve!: (value: T) => void + const promise = new Promise((next) => (resolve = next)) + return { promise, resolve } +} + +const reservation: CredentialReservation = { + userId: identity.sub, + relayHostId: identity.relayHostId, + credentialKind: 'resume', + relayDeviceId: 'device-1', + tokenHash: 'hash', + reservationId: 'reservation-1', + leaseExpiresAt: Date.now() + 60_000, + acceptedCredentialVersion: 2, + acceptedAs: 'current' +} + +function harness(options: { random?: () => number; now?: () => number } = {}) { + const acquireActivity = vi.fn().mockResolvedValue(undefined) + const releaseActivity = vi.fn().mockResolvedValue(true) + const assignments = { + activateControl: vi.fn().mockResolvedValue('control:production-gce-c3:1'), + markMigrationTargetRegistered: vi.fn().mockResolvedValue(undefined), + resolve: vi.fn().mockResolvedValue({ cellId: config.cellId }), + acquireActivity, + renewControlActivity: vi.fn().mockResolvedValue(undefined), + releaseActivity + } as unknown as RelayAssignmentStore + const store = { + resolveResume: vi.fn().mockResolvedValue({ userId: identity.sub }), + reserveCredential: vi.fn().mockResolvedValue(reservation), + failReservation: vi.fn().mockResolvedValue(undefined) + } + const observer = { + recordAuth: vi.fn(), + recordForwardedBytes: vi.fn(), + recordHttp: vi.fn(), + recordReconnect: vi.fn(), + recordSql: vi.fn(), + recordClientAcceptAbandoned: vi.fn() + } satisfies RelayRuntimeObserver + const registry = new HostSessionRegistry( + config, + vi.fn(), + store as unknown as RelayCredentialStore, + assignments, + new ProcessQueuedByteBudget(), + observer, + options.now, + options.random + ) + const activate = ( + registry as unknown as { + activate: ( + socket: WebSocket, + identity: RelayTokenClaims, + existing: null, + generation: number, + rebind: boolean, + assignmentEpoch: number, + appVersion: string + ) => Promise + } + ).activate.bind(registry) + return { registry, store, assignments, acquireActivity, releaseActivity, observer, activate } +} + +async function activeHost(h: ReturnType): Promise { + const control = new FakeSocket() + await h.activate(control as unknown as WebSocket, identity, null, 1, false, 1, '1.4.197') + return control +} + +describe('client accept abandoned mid-DB-phase', () => { + beforeEach(() => vi.useFakeTimers()) + afterEach(() => { + vi.clearAllTimers() + vi.useRealTimers() + }) + + it('stops after a slow activity acquire when the phone already hung up', async () => { + const h = harness() + const control = await activeHost(h) + const slowAcquire = deferred() + h.acquireActivity.mockReturnValueOnce(slowAcquire.promise) + const capacity = { bind: vi.fn(), release: vi.fn() } + const client = new FakeSocket() + const warn = vi.spyOn(console, 'warn').mockImplementation(() => undefined) + try { + const accepting = h.registry.acceptClient( + client as unknown as WebSocket, + identity.relayHostId, + 'credential', + capacity + ) + await vi.advanceTimersByTimeAsync(0) + expect(h.acquireActivity).toHaveBeenCalledOnce() + // The phone's 12s bound fires while the cell still waits on Postgres. + client.close(1000, 'client bound') + capacity.release() + slowAcquire.resolve() + await accepting + + // No conn-open reached the desktop; nothing pending; the lease it just took is + // released instead of leaking to expiry cleanup; bind never throws. + expect(control.send).not.toHaveBeenCalledWith(expect.stringContaining('conn-open')) + expect(capacity.bind).not.toHaveBeenCalled() + const session = h.registry.get({ userId: identity.sub, relayHostId: identity.relayHostId }) + expect(session?.pendingConns.size).toBe(0) + expect(h.store.failReservation).toHaveBeenCalledWith(reservation) + expect(h.releaseActivity).toHaveBeenCalledWith( + { userId: identity.sub, relayHostId: identity.relayHostId }, + expect.stringMatching(/^confirmation:/) + ) + expect(h.observer.recordClientAcceptAbandoned).toHaveBeenCalledWith( + 'activity', + expect.any(Number) + ) + const line = warn.mock.calls.map((call) => String(call[0])).find((entry) => + entry.includes('orca_relay_client_accept_abandoned') + ) + expect(line).toBeDefined() + expect(JSON.parse(line!)).toMatchObject({ stage: 'activity' }) + expect(line).not.toContain(identity.relayHostId) + } finally { + warn.mockRestore() + h.registry.drain(0) + vi.advanceTimersByTime(0) + } + }) + + it('stops after a slow credential reservation without acquiring an activity lease', async () => { + const h = harness() + await activeHost(h) + const slowReserve = deferred() + h.store.reserveCredential.mockReturnValueOnce(slowReserve.promise) + const client = new FakeSocket() + const warn = vi.spyOn(console, 'warn').mockImplementation(() => undefined) + try { + const accepting = h.registry.acceptClient( + client as unknown as WebSocket, + identity.relayHostId, + 'credential' + ) + await vi.advanceTimersByTimeAsync(0) + client.close(1000, 'client bound') + slowReserve.resolve(reservation) + await accepting + + expect(h.acquireActivity).not.toHaveBeenCalled() + expect(h.store.failReservation).toHaveBeenCalledWith(reservation) + expect(h.observer.recordClientAcceptAbandoned).toHaveBeenCalledWith( + 'credential', + expect.any(Number) + ) + } finally { + warn.mockRestore() + h.registry.drain(0) + vi.advanceTimersByTime(0) + } + }) + + it('still opens the connection when the phone is holding on', async () => { + const h = harness() + const control = await activeHost(h) + const capacity = { bind: vi.fn(), release: vi.fn() } + const client = new FakeSocket() + await h.registry.acceptClient( + client as unknown as WebSocket, + identity.relayHostId, + 'credential', + capacity + ) + expect(control.send).toHaveBeenCalledWith(expect.stringContaining('"type":"conn-open"')) + expect(capacity.bind).toHaveBeenCalledOnce() + expect(h.observer.recordClientAcceptAbandoned).not.toHaveBeenCalled() + expect(client.close).not.toHaveBeenCalled() + h.registry.drain(0) + vi.advanceTimersByTime(0) + }) +}) + +describe('control lease jitter', () => { + beforeEach(() => vi.useFakeTimers()) + afterEach(() => { + vi.clearAllTimers() + vi.useRealTimers() + }) + + it('grants a lease inside [55min - jitter, 55min] so cohorts drift apart', async () => { + const now = 1_700_000_000_000 + const helloAck = (socket: FakeSocket) => + JSON.parse( + String(socket.send.mock.calls.find((call) => String(call[0]).includes('host-hello-ack'))![0]) + ) as { leaseExpiresAt: number } + + const shortest = harness({ now: () => now, random: () => 0.999999 }) + const shortestAck = helloAck(await activeHost(shortest)) + const longest = harness({ now: () => now, random: () => 0 }) + const longestAck = helloAck(await activeHost(longest)) + + expect(longestAck.leaseExpiresAt).toBe(now + CONTROL_LEASE_MS) + expect(shortestAck.leaseExpiresAt).toBeGreaterThan(now + CONTROL_LEASE_MS - CONTROL_LEASE_JITTER_MS) + expect(shortestAck.leaseExpiresAt).toBeLessThan(longestAck.leaseExpiresAt) + // Nine of ten hosts that connected together now differ by minutes, not zero. + expect(longestAck.leaseExpiresAt - shortestAck.leaseExpiresAt).toBeGreaterThan(9 * 60 * 1000) + shortest.registry.drain(0) + longest.registry.drain(0) + vi.advanceTimersByTime(0) + }) + + it('rebinds re-roll the jitter instead of pinning the cohort phase', async () => { + const now = 1_700_000_000_000 + let roll = 0 + const h = harness({ now: () => now, random: () => roll }) + const first = await activeHost(h) + const session = h.registry.get({ userId: identity.sub, relayHostId: identity.relayHostId })! + const firstLease = session.leaseExpiresAt + roll = 0.5 + const rebind = new FakeSocket() + await ( + h.registry as unknown as { + activate: (...args: unknown[]) => Promise + } + ).activate(rebind as unknown as WebSocket, identity, session, 1, true, 1, '1.4.197') + expect(session.leaseExpiresAt).toBe(now + CONTROL_LEASE_MS - CONTROL_LEASE_JITTER_MS / 2) + expect(session.leaseExpiresAt).not.toBe(firstLease) + expect(first.close).toHaveBeenCalledWith(RELAY_CLOSE_CODE.PEER_DROPPED, 'control rebound') + h.registry.drain(0) + vi.advanceTimersByTime(0) + }) +}) diff --git a/cloud/apps/relay/src/host-session-registry.ts b/cloud/apps/relay/src/host-session-registry.ts index 11b7d1de030..8c5f7fde4e4 100644 --- a/cloud/apps/relay/src/host-session-registry.ts +++ b/cloud/apps/relay/src/host-session-registry.ts @@ -27,7 +27,7 @@ import { } from './credential-store.js' import { relayHostLogDigest } from './relay-host-log-digest.js' import type { RelayTokenClaims } from './relay-token-verifier.js' -import type { RelayRuntimeObserver } from './relay-observability.js' +import type { RelayClientAcceptStage, RelayRuntimeObserver } from './relay-observability.js' import type { PendingHostDataReservation } from './relay-connection-ledger.js' import { closeRelayWebSocket } from './relay-websocket-close.js' import { ProcessQueuedByteBudget, wireSplice } from './splice-forwarder.js' @@ -127,6 +127,13 @@ function send(socket: WebSocket, type: string, message: object): void { // stalled predecessor only accumulates doomed sockets. const ACTIVATION_QUEUE_WAIT_MS = 30_000 +// Why: hosts rebind 1-2 min before this expires, so every host that (re)connected +// in the same minute (a cell recreate dumps hundreds at once) rebinds as one +// cohort every cycle, forever, and each cohort lands on the cell-inventory lock +// as a single wave. Jittering the grant walks the cohort apart across cycles. +export const CONTROL_LEASE_MS = 55 * 60 * 1000 +export const CONTROL_LEASE_JITTER_MS = 10 * 60 * 1000 + export class HostSessionRegistry { private readonly sessions = new Map() private readonly activationQueues = new Map>() @@ -139,9 +146,14 @@ export class HostSessionRegistry { private readonly assignments: RelayAssignmentStore, private readonly queuedByteBudget: ProcessQueuedByteBudget, private readonly observer: RelayRuntimeObserver, - private readonly now: () => number = Date.now + private readonly now: () => number = Date.now, + private readonly random: () => number = Math.random ) {} + private controlLeaseExpiresAt(): number { + return this.now() + CONTROL_LEASE_MS - Math.floor(this.random() * CONTROL_LEASE_JITTER_MS) + } + async acceptClient( socket: WebSocket, hostId: string, @@ -153,6 +165,22 @@ export class HostSessionRegistry { this.rejectClient(socket, RELAY_CLOSE_CODE.DRAINING) return } + // Why: the accept runs several serialized Postgres calls behind the contended + // cell-inventory lock, and phones bound their dial. Finishing the work for a + // phone that already hung up used to acquire (and leak for 90s) an activity + // lease, then fail at bind with host_data_reservation_already_bound. + const acceptStartedAt = this.now() + const abandonedByClient = (stage: RelayClientAcceptStage, cleanup?: () => void): boolean => { + if (socket.readyState === socket.OPEN) return false + capacityReservation?.release() + cleanup?.() + const elapsedMs = this.now() - acceptStartedAt + this.observer.recordClientAcceptAbandoned?.(stage, elapsedMs) + console.warn( + JSON.stringify({ event: 'orca_relay_client_accept_abandoned', stage, elapsedMs }) + ) + return true + } if (this.config.role === 'cell') { const outerIdentity = (await this.store.resolveResume(hostId, credential)) ?? @@ -166,6 +194,7 @@ export class HostSessionRegistry { this.rejectClient(socket, RELAY_CLOSE_CODE.WRONG_CELL) return } + if (abandonedByClient('assignment')) return } const reservation = await this.store.reserveCredential(hostId, credential) if (!reservation) { @@ -175,6 +204,7 @@ export class HostSessionRegistry { return } this.observer.recordAuth(true) + if (abandonedByClient('credential', () => this.failReservationBestEffort(reservation))) return const session = this.sessions.get(this.key(reservation.userId, hostId)) if ( !session || @@ -214,6 +244,14 @@ export class HostSessionRegistry { return } } + if ( + abandonedByClient('activity', () => { + this.failReservationBestEffort(reservation) + if (credentialActivityId) this.releaseActivityBestEffort(identity, credentialActivityId) + }) + ) { + return + } const attachTimer = setTimeout(() => { session.pendingConns.delete(connId) capacityReservation?.release() @@ -727,7 +765,7 @@ export class HostSessionRegistry { existing.socket = socket existing.state = existing.regionalDrainAttemptId ? 'drain-only' : 'active' existing.appVersion = appVersion - existing.leaseExpiresAt = this.now() + 55 * 60 * 1000 + existing.leaseExpiresAt = this.controlLeaseExpiresAt() existing.lastPongAt = this.now() existing.activityRenewalDueAt = this.now() + RELAY_PROTOCOL_LIMITS.controlPingIntervalMs @@ -778,7 +816,7 @@ export class HostSessionRegistry { appVersion, state: 'active', socket, - leaseExpiresAt: this.now() + 55 * 60 * 1000, + leaseExpiresAt: this.controlLeaseExpiresAt(), orphanTimer: null, heartbeatTimer: null, lastPongAt: this.now(), diff --git a/cloud/apps/relay/src/relay-observability.test.ts b/cloud/apps/relay/src/relay-observability.test.ts index 2b9ceb0b72a..fc8a4fcb4af 100644 --- a/cloud/apps/relay/src/relay-observability.test.ts +++ b/cloud/apps/relay/src/relay-observability.test.ts @@ -195,16 +195,23 @@ describe('relay observability', () => { observability.recordControlClose(4402) observability.recordSpliceClose('host-oversize-frame') observability.recordSpliceClose('queue-limit') + observability.recordClientAcceptAbandoned('activity', 14_250.4) + observability.recordClientAcceptAbandoned('activity', 2_000) + observability.recordClientAcceptAbandoned('credential', 3_000) observability.flush(counts) observability.flush(counts) expect(entries[0]).toMatchObject({ controlClosesByCodeDelta: { 1006: 2, 4402: 1 }, - spliceClosesByTriggerDelta: { 'host-oversize-frame': 1, 'queue-limit': 1 } + spliceClosesByTriggerDelta: { 'host-oversize-frame': 1, 'queue-limit': 1 }, + clientAcceptsAbandonedByStageDelta: { activity: 2, credential: 1 }, + clientAcceptAbandonedMsMax: 14_250.4 }) expect(entries[1]).toMatchObject({ controlClosesByCodeDelta: {}, - spliceClosesByTriggerDelta: {} + spliceClosesByTriggerDelta: {}, + clientAcceptsAbandonedByStageDelta: {}, + clientAcceptAbandonedMsMax: 0 }) }) diff --git a/cloud/apps/relay/src/relay-observability.ts b/cloud/apps/relay/src/relay-observability.ts index 2266217d607..59ff437e40b 100644 --- a/cloud/apps/relay/src/relay-observability.ts +++ b/cloud/apps/relay/src/relay-observability.ts @@ -64,8 +64,12 @@ export interface RelayRuntimeObserver { }): void recordControlClose?(code: number): void recordSpliceClose?(trigger: string): void + recordClientAcceptAbandoned?(stage: RelayClientAcceptStage, elapsedMs: number): void } +// Which serialized accept step the phone had already hung up behind. +export type RelayClientAcceptStage = 'assignment' | 'credential' | 'activity' + type RelayMetricDeltas = { forwardedBytes: number authSuccesses: number @@ -87,6 +91,8 @@ type RelayMetricDeltas = { unavailableRegions: Record controlClosesByCode: Record spliceClosesByTrigger: Record + clientAcceptsAbandonedByStage: Record + clientAcceptAbandonedMsMax: number controlRenewalLatenciesMs: number[] controlRenewalsByOutcome: Record controlActivityRecoveries: number @@ -116,6 +122,8 @@ const emptyDeltas = (): RelayMetricDeltas => ({ unavailableRegions: {}, controlClosesByCode: {}, spliceClosesByTrigger: {}, + clientAcceptsAbandonedByStage: {}, + clientAcceptAbandonedMsMax: 0, controlRenewalLatenciesMs: [], controlRenewalsByOutcome: {}, controlActivityRecoveries: 0, @@ -228,6 +236,14 @@ export class RelayObservability implements RelayRuntimeObserver { (this.deltas.spliceClosesByTrigger[trigger] ?? 0) + 1 } + recordClientAcceptAbandoned(stage: RelayClientAcceptStage, elapsedMs: number): void { + increment(this.deltas.clientAcceptsAbandonedByStage, stage) + this.deltas.clientAcceptAbandonedMsMax = Math.max( + this.deltas.clientAcceptAbandonedMsMax, + elapsedMs + ) + } + start(readCounts: () => RelayProcessCounts, intervalMs = 30_000): void { if (this.timer) return this.eventLoop.enable() @@ -289,6 +305,8 @@ export class RelayObservability implements RelayRuntimeObserver { unavailableRegionsDelta: deltas.unavailableRegions, controlClosesByCodeDelta: deltas.controlClosesByCode, spliceClosesByTriggerDelta: deltas.spliceClosesByTrigger, + clientAcceptsAbandonedByStageDelta: deltas.clientAcceptsAbandonedByStage, + clientAcceptAbandonedMsMax: Number(deltas.clientAcceptAbandonedMsMax.toFixed(3)), sqlQueriesDelta: deltas.sqlQueries, sqlFailuresDelta: deltas.sqlFailures, sqlLatencyMsMax: Number(deltas.sqlLatencyMsMax.toFixed(3)), diff --git a/cloud/apps/relay/src/relay-server.ts b/cloud/apps/relay/src/relay-server.ts index a14240cfa6a..0f44e0b44c0 100644 --- a/cloud/apps/relay/src/relay-server.ts +++ b/cloud/apps/relay/src/relay-server.ts @@ -78,6 +78,7 @@ export function createRelayServer( database: RelayDatabase, options: { now?: () => number + random?: () => number connectionLedgerLimits?: { hardCap: number; controlReserve: number } cellIncarnation?: string } = {} @@ -113,7 +114,8 @@ export function createRelayServer( assignments, queuedBytes, observability, - options.now + options.now, + options.random ) const app = createRelayApp(config, { store, diff --git a/mobile/src/transport/mobile-direct-endpoint-probe.test.ts b/mobile/src/transport/mobile-direct-endpoint-probe.test.ts index 69fe8a2ab8a..6526326fe0f 100644 --- a/mobile/src/transport/mobile-direct-endpoint-probe.test.ts +++ b/mobile/src/transport/mobile-direct-endpoint-probe.test.ts @@ -71,4 +71,28 @@ describe('mobile direct endpoint probe', () => { expect(clients.get(host.endpoint)?.close).toHaveBeenCalledOnce() expect(result?.client.close).not.toHaveBeenCalled() }) + + it('fails as soon as every candidate falls into its own reconnect backoff', async () => { + // Incident 2026-09-04: foregrounding on a dead LAN produced an instant 1006 and + // the direct client's 500/1000/2000ms redials, while the probe sat on the + // 'connecting' phase and held the supervisor mutex for the whole 12s bound. + const clients: FakeClient[] = [] + const openDirect = vi.fn(() => { + const client = new FakeClient('connecting') + clients.push(client) + setTimeout(() => client.publishState('reconnecting'), 20) + return client + }) + + const probing = openAuthenticatedDirectEndpoint(host, openDirect, 12_000) + await vi.advanceTimersByTimeAsync(20) + await expect(probing).resolves.toBeNull() + + expect(clients).toHaveLength(2) + for (const client of clients) { + expect(client.close).toHaveBeenCalledOnce() + } + // No 12s timer is left behind to fire into a settled probe. + expect(vi.getTimerCount()).toBe(0) + }) }) diff --git a/mobile/src/transport/mobile-direct-endpoint-probe.ts b/mobile/src/transport/mobile-direct-endpoint-probe.ts index 114a4f29130..bf001564744 100644 --- a/mobile/src/transport/mobile-direct-endpoint-probe.ts +++ b/mobile/src/transport/mobile-direct-endpoint-probe.ts @@ -35,7 +35,11 @@ function waitForAuthenticatedSession(session: RpcClient, timeoutMs: number): Pro if (state === 'connected') { finish() resolve() - } else if (state === 'disconnected' || state === 'auth-failed') { + } else if (state === 'disconnected' || state === 'auth-failed' || state === 'reconnecting') { + // Why: a probe asks whether direct answers NOW. 'reconnecting' is the direct + // client's own backoff loop after a failed dial (an instant 1006 on a dead + // LAN); waiting it out held the supervisor's operation mutex for the full + // 12s bound and blocked relay recovery behind three doomed redials. finish() reject(new Error(`probe session ${state}`)) } diff --git a/mobile/src/transport/mobile-endpoint-supervisor.test.ts b/mobile/src/transport/mobile-endpoint-supervisor.test.ts index aeb9cddef63..1ba5be77f58 100644 --- a/mobile/src/transport/mobile-endpoint-supervisor.test.ts +++ b/mobile/src/transport/mobile-endpoint-supervisor.test.ts @@ -228,6 +228,32 @@ describe('mobile endpoint supervisor', () => { supervisor.stop() }) + it('does not block relay recovery behind a direct probe stuck in its redial loop', 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() + + // Foreground return: the probe dials direct at once, the dead LAN answers with + // an instant 1006, and the direct client enters its 500/1000/2000ms backoff. + supervisor.setForeground(false) + supervisor.setForeground(true) + await vi.advanceTimersByTimeAsync(0) + expect(deps.openDirect).toHaveBeenCalledOnce() + direct.publishState('reconnecting') + logical.publishState('disconnected') + + // Relay recovery must not wait out the probe's 12s bound. + await vi.advanceTimersByTimeAsync(0) + expect(openRelay).toHaveBeenCalledOnce() + expect(direct.close).toHaveBeenCalled() + expect(logical.getState()).toBe('connected') + expect(logical.getActivePath()).toBe('relay') + 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 diff --git a/src/main/runtime/relay/relay-origin-pool.ts b/src/main/runtime/relay/relay-origin-pool.ts index a95d8f21340..7ab87d61643 100644 --- a/src/main/runtime/relay/relay-origin-pool.ts +++ b/src/main/runtime/relay/relay-origin-pool.ts @@ -9,6 +9,15 @@ import { RelayHttpError, requestRelayAssignment, type RelayAssignment } from './ import type { RelayBrokerStatus, RelayIdentity } from './relay-session-broker-contract' import type { RelayRegion } from './relay-region-preference' +// Why: every host whose lease expires in the same minute rebinds in the same +// minute, and each rebind takes the relay's contended cell-inventory lock. A +// cell recreate re-homes hundreds of hosts at once and pins that cohort to one +// phase for the life of the process (observed 2026-09-04: ~1k rebinds in 3s +// every ~54 min). A wide early window spreads each cycle; the floor keeps the +// rebind clear of the relay's expiry sweep even under a slow director. +const CONTROL_ROTATION_EARLY_MIN_MS = 60_000 +const CONTROL_ROTATION_EARLY_JITTER_MS = 5 * 60_000 + type RelayOriginPoolOptions = { directorUrl: string relayHostId: string @@ -244,7 +253,8 @@ export class RelayOriginPool { } const now = (this.options.now ?? Date.now)() const random = this.options.random ?? Math.random - const earlyMs = 60_000 + Math.floor(random() * 60_001) + const earlyMs = + CONTROL_ROTATION_EARLY_MIN_MS + Math.floor(random() * (CONTROL_ROTATION_EARLY_JITTER_MS + 1)) const delay = Math.max(0, origin.controlLeaseExpiresAt - earlyMs - now) this.rotationTimer = setTimeout(() => void this.rebindActiveControl(origin), delay) } diff --git a/src/main/runtime/relay/relay-session-broker.test.ts b/src/main/runtime/relay/relay-session-broker.test.ts index 6f27b4f2e14..7f110568179 100644 --- a/src/main/runtime/relay/relay-session-broker.test.ts +++ b/src/main/runtime/relay/relay-session-broker.test.ts @@ -486,6 +486,59 @@ describe('RelaySessionBroker lifecycle ownership', () => { }) }) +describe('RelaySessionBroker control rotation spreading', () => { + beforeEach(() => { + fakes.controls.length = 0 + fakes.transports.length = 0 + fakes.controlConnect.mockReset() + fakes.exchange.mockReset().mockResolvedValue({ relayToken: 'relay-jwt', expiresAt: 10_000_000 }) + fakes.assign.mockReset().mockResolvedValue({ + cellUrl: 'https://relay.example.test', + assignmentEpoch: 1, + leaseExpiresAt: 10_000_000 + }) + }) + + // Incident 2026-09-04 00:50Z: ~1k hosts re-homed by one cell recreate rebound + // together every ~54 min, each rebind taking the relay's cell-inventory lock. + it('spreads same-lease hosts across a multi-minute window instead of one minute', async () => { + vi.useFakeTimers() + try { + const leaseExpiresAt = 55 * 60_000 + const ack: RelayHostHelloAckMessage = { + type: 'host-hello-ack', + v: 1, + generation: 1, + controlResumeSecret: 'R'.repeat(43), + leaseExpiresAt, + activeConnIds: [], + pendingConns: [] + } + fakes.controlConnect.mockResolvedValue(ack) + const earliest = await RelaySessionBroker.connect(brokerOptions({ random: () => 0.999999 })) + const latest = await RelaySessionBroker.connect(brokerOptions({ random: () => 0 })) + expect(fakes.controls).toHaveLength(2) + + // The widest early roll rebinds ~6 min before expiry; the narrowest at 1 min. + await vi.advanceTimersByTimeAsync(leaseExpiresAt - 6 * 60_000 - 1) + expect(fakes.controls).toHaveLength(2) + await vi.advanceTimersByTimeAsync(2) + expect(fakes.controls).toHaveLength(3) + expect(fakes.controls[2]!.options.previousGeneration).toBe(1) + + await vi.advanceTimersByTimeAsync(4 * 60_000 + 59_000) + expect(fakes.controls).toHaveLength(3) + await vi.advanceTimersByTimeAsync(1_000) + expect(fakes.controls).toHaveLength(4) + + earliest.closeNow() + latest.closeNow() + } finally { + vi.useRealTimers() + } + }) +}) + function brokerBasisIds(broker: RelaySessionBroker): string[] { const pool = (broker as unknown as { originPool: unknown }).originPool return [...(pool as { basisOrigins: Map }).basisOrigins.keys()]