mirror of
https://github.com/stablyai/orca.git
synced 2026-09-29 08:03:20 +00:00
fix(relay): stop dead accept work, spread control rotations, fail direct probes fast
Incident 2026-09-04 ~01:05Z: after a background/foreground cycle the phone's
relay dial timed out inside the cell's acceptClient DB phase while the fleet
was in a cell-inventory lock storm (55P03 retries ~7.5k/h vs a ~1k/h floor).
Relay (cloud/apps/relay)
- acceptClient checks socket.readyState after each serialized Postgres call
and abandons the accept once the phone has hung up, releasing the capacity
reservation, failing the credential reservation, and releasing the activity
lease it just acquired instead of leaking it to expiry cleanup and then
throwing host_data_reservation_already_bound at bind.
- New structured event orca_relay_client_accept_abandoned {stage, elapsedMs}
and runtime-metric fields clientAcceptsAbandonedByStageDelta /
clientAcceptAbandonedMsMax so the "phone gave up behind the lock" rate is
quantifiable per cell.
- Control lease grants are jittered: 55 min minus [0, 10 min). Every host
that (re)connected in the same minute rebound as one cohort every ~54 min
(c27 autoheal recreate at 23:23Z re-homed ~420 controls; ~1.1k 1006 +
~1k 4408 "control rebound" closes landed in a 3 s window at 00:50:14Z),
and each rebind is an activateControl transaction on the inventory lock.
No wire change: leaseExpiresAt was always a server-chosen absolute time.
Desktop (src/main/runtime/relay)
- Control rotation rebinds 1-6 min early instead of 1-2 min, so a re-homed
cohort spreads across cycles rather than pinning one phase for the life of
the process.
Phone (mobile/src/transport)
- openAuthenticatedDirectEndpoint treats 'reconnecting' as a failed probe. On
a dead LAN the foreground direct dial dies with an instant 1006 and the
direct client enters its own 500/1000/2000 ms backoff; the probe used to
wait out its full 12 s bound holding the supervisor mutex, so relay recovery
queued behind three doomed redials. The stage-aware bound from #18518 is
unaffected.
Not done here: rolling the 23 GCE cells onto the post-#18521 image (500 ms
lock_timeout) is a deploy owned by cloud-deploy-relay-production-same-cap.
This commit is contained in:
@@ -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<T>(): { promise: Promise<T>; resolve: (value: T) => void } {
|
||||
let resolve!: (value: T) => void
|
||||
const promise = new Promise<T>((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<void>
|
||||
}
|
||||
).activate.bind(registry)
|
||||
return { registry, store, assignments, acquireActivity, releaseActivity, observer, activate }
|
||||
}
|
||||
|
||||
async function activeHost(h: ReturnType<typeof harness>): Promise<FakeSocket> {
|
||||
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<void>()
|
||||
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<CredentialReservation>()
|
||||
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<void>
|
||||
}
|
||||
).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)
|
||||
})
|
||||
})
|
||||
@@ -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<string, HostSession>()
|
||||
private readonly activationQueues = new Map<string, Promise<void>>()
|
||||
@@ -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(),
|
||||
|
||||
@@ -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
|
||||
})
|
||||
})
|
||||
|
||||
|
||||
@@ -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<string, number>
|
||||
controlClosesByCode: Record<string, number>
|
||||
spliceClosesByTrigger: Record<string, number>
|
||||
clientAcceptsAbandonedByStage: Record<string, number>
|
||||
clientAcceptAbandonedMsMax: number
|
||||
controlRenewalLatenciesMs: number[]
|
||||
controlRenewalsByOutcome: Record<string, number>
|
||||
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)),
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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)
|
||||
})
|
||||
})
|
||||
|
||||
@@ -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}`))
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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<string, unknown> }).basisOrigins.keys()]
|
||||
|
||||
Reference in New Issue
Block a user