diff --git a/README.md b/README.md
index 7a3cbe2360c..2ae59035da8 100644
--- a/README.md
+++ b/README.md
@@ -238,9 +238,9 @@ Pair with your desktop app to monitor and steer your agents from your phone.
- **Discord:** Join the community on **[Discord](https://discord.gg/fzjDKHxv8Q)**.
- **Twitter / X:** Follow **[@orca_build](https://x.com/orca_build)** for updates and announcements.
-- **WeChat:** Scan to join the Orca community WeChat group 8.
+- **WeChat:** Scan to join the Orca community WeChat group 8. Group 8 may be full; if so, scan the Group 9 QR code instead.
-
+
- **Feedback & Ideas:** We ship fast. Missing something? [Request a new feature](https://github.com/stablyai/orca/issues).
- **Privacy:** See the [privacy & telemetry docs](https://www.onorca.dev/docs/telemetry) for what anonymous usage data Orca collects and how to opt out.
diff --git a/cloud/apps/relay-ops/src/incident-monitor.test.ts b/cloud/apps/relay-ops/src/incident-monitor.test.ts
index 61a73b64dbe..4e1da9fab26 100644
--- a/cloud/apps/relay-ops/src/incident-monitor.test.ts
+++ b/cloud/apps/relay-ops/src/incident-monitor.test.ts
@@ -111,14 +111,19 @@ describe('incident monitor evaluator', () => {
})
})
- it('freezes when postgres retries exceed the recalibrated ceiling', () => {
- const sample = healthySample()
- sample.sources['relay-logs']!.signals['relay.postgres_retries'] =
- signal(INCIDENT_MONITOR_THRESHOLDS.relayPostgresRetries + 1)
- expect(evaluateIncidentSample(sample, startedAt)).toMatchObject({
+ // Why: the global relay_cells lock made retries a steady-state rate (24 h p99
+ // 1320/5min on 2026-09-04); the bar fences only unbounded growth beyond that.
+ it('tolerates the measured healthy retry rate and freezes above the bar', () => {
+ const healthy = healthySample()
+ healthy.sources['relay-logs']!.signals['relay.postgres_retries'] = signal(1504)
+ expect(evaluateIncidentSample(healthy, startedAt).status).toBe('green')
+
+ const incident = healthySample()
+ incident.sources['relay-logs']!.signals['relay.postgres_retries'] = signal(2001)
+ expect(evaluateIncidentSample(incident, startedAt)).toMatchObject({
status: 'freeze',
failures: [
- expect.objectContaining({ signal: 'relay.postgres_retries', threshold: 300 })
+ expect.objectContaining({ signal: 'relay.postgres_retries', threshold: 2000 })
]
})
})
diff --git a/cloud/apps/relay-ops/src/incident-monitor.ts b/cloud/apps/relay-ops/src/incident-monitor.ts
index 868bb86fb93..a121568d918 100644
--- a/cloud/apps/relay-ops/src/incident-monitor.ts
+++ b/cloud/apps/relay-ops/src/incident-monitor.ts
@@ -32,11 +32,20 @@ export const INCIDENT_MONITOR_THRESHOLDS = {
relayPoolWaiting: 800,
relayPoolWaitMs: 2_500,
// Why: successful lock retries are the contention machinery working, not harm.
- // Healthy 2026-08-26 baseline bursts to 234/5min (26% of windows crossed the old
- // bar of 20, set unmeasured at the monitor's 2026-07-28 birth); the 2026-08-23
- // incident ran ~2,200-3,000/5min. 300 clears healthy bursts with ~10x incident
- // margin; relayPostgresRetryExhausted below bounds the terminally failed share.
- relayPostgresRetries: 300,
+ // Recalibrated 2026-09-04 from 300, which was set 2026-08-26 when healthy bursts
+ // reached 234/5min. The global relay_cells FOR UPDATE lock has since become the
+ // fleet's steady state: measured fleet-wide (director + cells, summed per five
+ // minutes) 2026-09-03T05Z..2026-09-04T05Z p50 430 / p90 924 / p99 1320 / max
+ // 1504, with 55% of windows over 300 and only 22% of 15-minute gates clean, so
+ // the bar blocked the very cell roll that carries the 500 ms lock wait (#18521)
+ // and the beginProof crash guard to the cells. The 2026-08-23 lock incident on
+ // this same metric peaked at 1510 in one window and 646 in the next, so it is
+ // not separable from today's contention by retries alone; it is caught by
+ // relayPostgresRetryExhausted (467 at the peak vs a 300 bar), director
+ // concurrency, and the pool bars. 2000 passes every healthy 15-minute window
+ // measured in the last 24 h and still fences unbounded growth. Re-tighten once
+ // the fleet is on the 500 ms lock wait and the baseline is re-measured.
+ relayPostgresRetries: 2000,
// Why: 300 per five minutes, recalibrated 2026-09-04 from a bar of zero that no
// production window has cleared since #18521 shipped to the director. That
// change cut the request-path cell-inventory wait from the 1 s pool lock_timeout
@@ -48,7 +57,7 @@ export const INCIDENT_MONITOR_THRESHOLDS = {
// quiet hours p50 2 / max 36; pre-#18521 daytime p50 10 / p90 25 / max 87;
// post-#18521 p50 42 / p90 147 / max 220. The 2026-08-23 lock incident peaked
// at 467. 300 clears every measured healthy window and still sits below the
- // incident shape; relayPostgresRetries above stays the ~10x discriminator.
+ // incident shape; retries above fence only unbounded growth.
// User-facing /v1/assign 503 share did not move with #18521 (13.9% old image
// vs 12.3% new, same evening), so exhaustion is not a proxy for user harm.
relayPostgresRetryExhausted: 300,
diff --git a/cloud/apps/relay-ops/src/resource-inventory.test.ts b/cloud/apps/relay-ops/src/resource-inventory.test.ts
index 6d8b3070c98..4bf43fe9a7a 100644
--- a/cloud/apps/relay-ops/src/resource-inventory.test.ts
+++ b/cloud/apps/relay-ops/src/resource-inventory.test.ts
@@ -23,8 +23,10 @@ describe('readResourceInventory', () => {
calls += 1
return new Response(null, { status: 200 })
},
- async () => {
- waits += 1
+ {
+ wait: async () => {
+ waits += 1
+ }
}
)
@@ -45,8 +47,10 @@ describe('readResourceInventory', () => {
calls.set(path, call)
return new Response(null, { status: path === '/ready' && call === 1 ? 503 : 200 })
},
- async (ms) => {
- waits.push(ms)
+ {
+ wait: async (ms) => {
+ waits.push(ms)
+ }
}
)
@@ -65,17 +69,130 @@ describe('readResourceInventory', () => {
calls += 1
return new Response(null, { status: 503 })
},
- async (ms) => {
- waits.push(ms)
+ {
+ wait: async (ms) => {
+ waits.push(ms)
+ }
}
)
expect(result.health).toBe(false)
expect(result.ready).toBe(false)
expect(calls).toBe(4)
+ // A refusing endpoint is a reading, so only the independent retry runs.
expect(waits).toEqual([11_000])
})
+ it('treats a thrown fetch as no reading and re-asks that path once', async () => {
+ const calls: string[] = []
+ const waits: number[] = []
+ const result = await probeEndpointHealth(
+ 'https://c9.relay.onorca.dev',
+ async (input) => {
+ const path = new URL(String(input)).pathname
+ calls.push(path)
+ if (path === '/health' && calls.filter((call) => call === '/health').length === 1) {
+ throw new TypeError('fetch failed')
+ }
+ return new Response(null, { status: 200 })
+ },
+ {
+ wait: async (ms) => {
+ waits.push(ms)
+ }
+ }
+ )
+
+ expect(result.health).toBe(true)
+ expect(result.ready).toBe(true)
+ expect(calls.filter((call) => call === '/health')).toEqual(['/health', '/health'])
+ expect(waits).toEqual([1_000])
+ })
+
+ it('fails closed when both attempts of a path throw', async () => {
+ const calls: string[] = []
+ const waits: number[] = []
+ const result = await probeEndpointHealth(
+ 'https://c9.relay.onorca.dev',
+ async (input) => {
+ const path = new URL(String(input)).pathname
+ calls.push(path)
+ if (path === '/health') throw new TypeError('fetch failed')
+ return new Response(null, { status: 200 })
+ },
+ {
+ wait: async (ms) => {
+ waits.push(ms)
+ }
+ }
+ )
+
+ expect(result.health).toBe(false)
+ expect(calls.filter((call) => call === '/health')).toHaveLength(4)
+ expect(waits).toEqual([1_000, 11_000, 1_000])
+ })
+
+ it('accepts an auth-shaped endpoint that serves no readiness path', async () => {
+ const calls: string[] = []
+ const waits: number[] = []
+ const result = await probeEndpointHealth(
+ 'https://login.onorca.dev',
+ async (input) => {
+ const path = new URL(String(input)).pathname
+ calls.push(path)
+ return new Response(null, { status: path === '/ready' ? 404 : 200 })
+ },
+ {
+ requiresReady: false,
+ wait: async (ms) => {
+ waits.push(ms)
+ }
+ }
+ )
+
+ expect(result.health).toBe(true)
+ expect(result.ready).toBeNull()
+ expect(calls).toEqual(['/health'])
+ expect(waits).toEqual([])
+ })
+
+ it('still requires readiness for the director and cells', async () => {
+ const waits: number[] = []
+ const result = await probeEndpointHealth(
+ 'https://relay.onorca.dev',
+ async (input) => new Response(null, {
+ status: new URL(String(input)).pathname === '/ready' ? 503 : 200
+ }),
+ {
+ wait: async (ms) => {
+ waits.push(ms)
+ }
+ }
+ )
+
+ expect(result.health).toBe(true)
+ expect(result.ready).toBe(false)
+ expect(waits).toEqual([11_000])
+ })
+
+ it('measures latency as the answering round trip, not the retry delay', async () => {
+ let healthCalls = 0
+ const result = await probeEndpointHealth(
+ 'https://c9.relay.onorca.dev',
+ async (input) => {
+ if (new URL(String(input)).pathname !== '/health') return new Response(null, { status: 200 })
+ healthCalls += 1
+ if (healthCalls === 1) throw new TypeError('fetch failed')
+ return new Response(null, { status: 200 })
+ },
+ { wait: async (ms) => await new Promise((resolve) => setTimeout(resolve, Math.min(ms, 60))) }
+ )
+
+ expect(result.health).toBe(true)
+ expect(result.latencyMs).not.toBeNull()
+ expect(result.latencyMs!).toBeLessThan(60)
+ })
+
it('uses aggregate REST inventory without probing sleeping staging endpoints', async () => {
const gcloud: GcloudClient = { accessToken: async () => 'a'.repeat(40) }
let publicProbeCalls = 0
diff --git a/cloud/apps/relay-ops/src/resource-inventory.ts b/cloud/apps/relay-ops/src/resource-inventory.ts
index 62ed3fd862b..c1d01191baf 100644
--- a/cloud/apps/relay-ops/src/resource-inventory.ts
+++ b/cloud/apps/relay-ops/src/resource-inventory.ts
@@ -102,6 +102,7 @@ export type ResourceInventory = {
const unavailableEndpoint = (): EndpointHealth => ({ health: null, ready: null, latencyMs: null })
const independentEndpointRetryDelayMs = 11_000
+const transientProbeRetryDelayMs = 1_000
function finalSegment(value: string): string {
return value.split('/').at(-1) ?? value
@@ -138,41 +139,80 @@ async function googleRequest(
return await response.json()
}
-async function endpointProbe(origin: string, fetchImpl: typeof fetch): Promise {
- const startedAt = performance.now()
- const check = async (path: '/health' | '/ready'): Promise => {
+// A reading the endpoint actually produced: ok is its answer, latencyMs is that answer's round trip.
+type PathReading = { ok: boolean; latencyMs: number | null }
+
+async function probePath(
+ origin: string,
+ path: '/health' | '/ready',
+ fetchImpl: typeof fetch,
+ wait: (ms: number) => Promise
+): Promise {
+ // null means the request never produced an answer (DNS/TCP/TLS failure or the 8s abort).
+ const attempt = async (): Promise => {
+ const startedAt = performance.now()
try {
const response = await fetchImpl(`${origin}${path}`, {
redirect: 'error',
signal: AbortSignal.timeout(8_000)
})
- return response.ok
+ return { ok: response.ok, latencyMs: Math.round(performance.now() - startedAt) }
} catch {
- return false
+ return null
}
}
- const [health, ready] = await Promise.all([check('/health'), check('/ready')])
- return { health, ready, latencyMs: Math.round(performance.now() - startedAt) }
+ const first = await attempt()
+ if (first) return first
+ // A thrown fetch is the absence of a reading, not an unhealthy answer, so re-ask before concluding.
+ await wait(transientProbeRetryDelayMs)
+ return (await attempt()) ?? { ok: false, latencyMs: null }
+}
+
+async function endpointProbe(
+ origin: string,
+ fetchImpl: typeof fetch,
+ requiresReady: boolean,
+ wait: (ms: number) => Promise
+): Promise {
+ const [health, ready] = await Promise.all([
+ probePath(origin, '/health', fetchImpl, wait),
+ requiresReady ? probePath(origin, '/ready', fetchImpl, wait) : null
+ ])
+ // Latency is the slowest answering round trip in this probe; retry delays are not serving latency.
+ const latencies = [health.latencyMs, ready?.latencyMs ?? null].filter(
+ (value): value is number => value !== null
+ )
+ return {
+ health: health.ok,
+ ready: ready ? ready.ok : null,
+ latencyMs: latencies.length > 0 ? Math.max(...latencies) : null
+ }
+}
+
+export type EndpointProbeOptions = {
+ // Auth serves no /ready by design, so it is judged on /health and latency alone.
+ requiresReady?: boolean
+ wait?: (ms: number) => Promise
}
export async function probeEndpointHealth(
origin: string,
fetchImpl: typeof fetch,
- wait: (ms: number) => Promise = async (ms) =>
- await new Promise((resolvePromise) => setTimeout(resolvePromise, ms))
+ options: EndpointProbeOptions = {}
): Promise {
- const first = await endpointProbe(origin, fetchImpl)
- if (
- first.health &&
- first.ready &&
- first.latencyMs !== null &&
- first.latencyMs <= INCIDENT_MONITOR_THRESHOLDS.endpointLatencyMs
- ) {
- return first
- }
+ const requiresReady = options.requiresReady ?? true
+ const wait = options.wait ??
+ (async (ms: number) => await new Promise((resolvePromise) => setTimeout(resolvePromise, ms)))
+ const accepted = (probe: EndpointHealth): boolean =>
+ probe.health === true &&
+ (!requiresReady || probe.ready === true) &&
+ probe.latencyMs !== null &&
+ probe.latencyMs <= INCIDENT_MONITOR_THRESHOLDS.endpointLatencyMs
+ const first = await endpointProbe(origin, fetchImpl, requiresReady, wait)
+ if (accepted(first)) return first
// Outwait Relay's ten-second readiness cache before treating the retry as independent.
await wait(independentEndpointRetryDelayMs)
- return await endpointProbe(origin, fetchImpl)
+ return await endpointProbe(origin, fetchImpl, requiresReady, wait)
}
function imageDigest(template: z.infer): string | null {
@@ -338,7 +378,8 @@ export async function readResourceInventory(
? [unavailableEndpoint(), unavailableEndpoint()]
: await Promise.all([
probeEndpointHealth(environment.directorOrigin, fetchImpl),
- probeEndpointHealth(environment.authOrigin, fetchImpl)
+ // The auth service exposes no /ready, so requiring it would fail every first probe.
+ probeEndpointHealth(environment.authOrigin, fetchImpl, { requiresReady: false })
])
const cells = await Promise.all(environment.cells.map((cell, index) =>
readCell(environment, cell, migValues[index] ?? null, token, fetchImpl)
diff --git a/cloud/apps/relay/src/assignment-connection-headroom-postgres.test.ts b/cloud/apps/relay/src/assignment-connection-headroom-postgres.test.ts
index 6ac9521c3d6..80a74a47eeb 100644
--- a/cloud/apps/relay/src/assignment-connection-headroom-postgres.test.ts
+++ b/cloud/apps/relay/src/assignment-connection-headroom-postgres.test.ts
@@ -44,6 +44,12 @@ describePostgres('PostgreSQL assignment connection headroom', () => {
`DELETE FROM relay_assignments
WHERE user_id LIKE 'connection-headroom-postgres-%'`
)
+ // A snapshot left by an aborted run rejects the replayed watermark
+ // with stale_connection_snapshot.
+ await database.query(
+ `DELETE FROM relay_cell_connection_snapshots WHERE cell_id = ?`,
+ [cell.id]
+ )
await database.query(
`DELETE FROM relay_cell_connection_runtime WHERE cell_id = ?`,
[cell.id]
diff --git a/cloud/apps/relay/src/assignment-control-supersession-postgres.test.ts b/cloud/apps/relay/src/assignment-control-supersession-postgres.test.ts
index 10193b78cc6..cf8819686b5 100644
--- a/cloud/apps/relay/src/assignment-control-supersession-postgres.test.ts
+++ b/cloud/apps/relay/src/assignment-control-supersession-postgres.test.ts
@@ -38,6 +38,10 @@ describePostgres('PostgreSQL control supersession', () => {
[identity.userId]
)
await database.query(`DELETE FROM relay_assignments WHERE user_id = ?`, [identity.userId])
+ // A snapshot left by an aborted run rejects the replayed watermark with stale_connection_snapshot.
+ await database.query(`DELETE FROM relay_cell_connection_snapshots WHERE cell_id = ?`, [
+ cell.id
+ ])
await database.query(`DELETE FROM relay_cell_connection_runtime WHERE cell_id = ?`, [cell.id])
await database.query(`DELETE FROM relay_cell_connection_limits WHERE cell_id = ?`, [cell.id])
await database.query(`DELETE FROM relay_cell_runtime WHERE cell_id = ?`, [cell.id])
diff --git a/cloud/apps/relay/src/assignment-store.ts b/cloud/apps/relay/src/assignment-store.ts
index d0517d46746..226df9b3984 100644
--- a/cloud/apps/relay/src/assignment-store.ts
+++ b/cloud/apps/relay/src/assignment-store.ts
@@ -645,9 +645,17 @@ export class RelayAssignmentStore {
): Promise {
const now = this.now()
return await this.database.transaction(async (transaction) => {
- const lockedCells = inventoryFirst
- ? await this.lockCellInventory(transaction, lockMode)
+ // Why: the retry exists to take a cell row before the assignment row, the
+ // order placement uses. It only ever needs the one cell this host is
+ // pinned to, so read the pin unlocked and lock that row alone; taking all
+ // 23 queued every sticky refresh in the fleet behind every other one.
+ const pinnedCellId = inventoryFirst
+ ? await this.pinnedCellId(transaction, identity)
: undefined
+ const lockedCells =
+ pinnedCellId === undefined
+ ? undefined
+ : await this.lockCellRows(transaction, [pinnedCellId], lockMode)
const existing = await this.assignmentRow(transaction, identity, inventoryFirst)
if (!existing) return null
const activityLeases = await this.lockAssignmentActivities(transaction, identity, true)
@@ -661,6 +669,11 @@ export class RelayAssignmentStore {
}
const currentCellId = text(existing, 'cell_id')
+ // The pin moved between the unlocked read and the assignment lock, so the
+ // row held is the wrong one. Same recovery as losing the lock: retry.
+ if (pinnedCellId !== undefined && pinnedCellId !== currentCellId) {
+ throw new Error('database_lock_unavailable')
+ }
const hadControl = holdsControlLease(
activityLeases,
currentCellId,
@@ -701,14 +714,9 @@ export class RelayAssignmentStore {
if (hadControl) {
await this.touchAssignment(transaction, identity, leaseExpiresAt, now)
} else {
- const nextReservation = integer(currentRow, 'reserved_requests') + 1
- if (nextReservation > integer(currentRow, 'capacity_requests')) {
- throw new Error('relay_capacity_exhausted')
- }
- await transaction.query(
- `UPDATE relay_cells SET reserved_requests = ?, updated_at = ? WHERE cell_id = ?`,
- [nextReservation, now, currentCellId]
- )
+ // Delta, not the value read from the snapshot: an absolute write here
+ // would clobber any concurrent movement of the same counter.
+ await this.adjustCellReservationAtomically(transaction, currentCellId, 1)
await this.adjustActivityCount(transaction, identity, 'control', 1, leaseExpiresAt, now)
await this.insertPendingControlLease(
transaction,
@@ -3202,8 +3210,7 @@ export class RelayAssignmentStore {
)
const requestDelta = ACTIVITY_REQUEST_UNITS[kind] * (after - before)
if (requestDelta !== 0) {
- await this.lockCellInventory(transaction, 'request')
- await this.adjustCellReservation(transaction, text(row, 'cell_id'), requestDelta)
+ await this.adjustCellReservationAtomically(transaction, text(row, 'cell_id'), requestDelta)
}
})
})
@@ -3263,9 +3270,12 @@ export class RelayAssignmentStore {
}
const units = ACTIVITY_REQUEST_UNITS[input.kind]
if (existing) {
- await this.lockCellInventory(transaction, 'request')
+ // Why: a client-chosen activity id can move between cells, so lock the
+ // one or two rows this path touches in cell_id order, the same order
+ // placement takes the inventory in, and no cycle can form.
+ await this.lockCellRows(transaction, [text(existing, 'cell_id'), input.cellId])
await this.removeActivityLease(transaction, identity, existing, now)
- await this.adjustCellReservation(transaction, input.cellId, units)
+ await this.adjustCellReservationAtomically(transaction, input.cellId, units)
}
await this.adjustActivityCount(transaction, identity, input.kind, 1, expiresAt, now)
await transaction.query(
@@ -3580,8 +3590,7 @@ export class RelayAssignmentStore {
)
await this.touchAssignment(transaction, identity, expiresAt, now)
} else {
- await this.lockCellInventory(transaction, 'request')
- await this.adjustCellReservation(transaction, input.cellId, 1)
+ await this.adjustCellReservationAtomically(transaction, input.cellId, 1)
await this.adjustActivityCount(transaction, identity, 'control', 1, expiresAt, now)
await transaction.query(
`INSERT INTO relay_assignment_activity_leases
@@ -6863,13 +6872,24 @@ export class RelayAssignmentStore {
targetCellId
]
)
- const cells = await this.lockCellInventory(transaction, 'pool-default')
+ // Only the two cells this repairs need holding. The id set below is an
+ // existence check against a table that only reconcileCells writes, so it
+ // reads unlocked instead of dragging the other 21 rows into the section.
+ const cellIds = new Set(
+ (await transaction.query(`SELECT cell_id FROM relay_cells`)).map((row) =>
+ text(row, 'cell_id')
+ )
+ )
+ const cells = await this.lockCellRows(
+ transaction,
+ [sourceCellId, targetCellId],
+ 'pool-default'
+ )
const assignmentKeys = new Set(
assignments.map((row) =>
assignmentKey(text(row, 'user_id'), text(row, 'relay_host_id'))
)
)
- const cellIds = new Set(cells.map((row) => text(row, 'cell_id')))
const assignmentCounts = new Map<
string,
{ counts: Record; leaseExpiresAt: number }
@@ -6923,9 +6943,7 @@ export class RelayAssignmentStore {
)
}
- for (const row of cells.filter((cell) =>
- [sourceCellId, targetCellId].includes(text(cell, 'cell_id'))
- )) {
+ for (const row of cells) {
const cellId = text(row, 'cell_id')
const expected = cellUnits.get(cellId) ?? 0
if (expected > integer(row, 'capacity_requests')) {
@@ -6954,6 +6972,43 @@ export class RelayAssignmentStore {
return rows
}
+ // Per-connection paths touch one or two cells. Locking exactly those rows,
+ // in the same ascending order the inventory lock uses (ORDER BY fixes the
+ // row-lock order), keeps them off the fleet-wide lock without a cycle.
+ // The wait policy follows the caller for the same reason the inventory lock's
+ // does: a sweep must not fail terminally on ordinary contention. Hold time is
+ // deliberately not sampled here — the metric tracks the fleet-wide lock these
+ // rows replace, and mixing in short single-row holds would flatter it.
+ private async lockCellRows(
+ database: RelayDatabase,
+ cellIds: string[],
+ mode: CellInventoryLockMode = 'request'
+ ): Promise {
+ const distinct = [...new Set(cellIds)]
+ const { measureHoldMs: _sampled, ...wait } = cellInventoryLockOptions(mode)
+ return await database.queryLocked(
+ `SELECT * FROM relay_cells WHERE cell_id IN (${distinct.map(() => '?').join(', ')})
+ ORDER BY cell_id ASC`,
+ distinct,
+ wait
+ )
+ }
+
+ // Unlocked on purpose: this only names the row to lock next, and the caller
+ // re-checks the pin once the assignment row is held.
+ private async pinnedCellId(
+ database: RelayDatabase,
+ identity: AssignmentIdentity
+ ): Promise {
+ const row = (
+ await database.query(
+ `SELECT cell_id FROM relay_assignments WHERE user_id = ? AND relay_host_id = ?`,
+ [identity.userId, identity.relayHostId]
+ )
+ )[0]
+ return row ? text(row, 'cell_id') : undefined
+ }
+
private async lockGeneralCellInventory(
database: RelayDatabase,
mode: CellInventoryLockMode
@@ -6972,10 +7027,11 @@ export class RelayAssignmentStore {
private async leastLoadedCell(
database: RelayDatabase,
- lockedCells: SqlRow[] | undefined,
+ // Required: the one caller has already locked the inventory it selects from,
+ // and an optional parameter left a second fleet-wide lock reachable here.
+ rows: SqlRow[],
preferredRegion: RelayRegion
): Promise {
- const rows = lockedCells ?? (await this.lockCellInventory(database, 'pool-default'))
const regions = new Map(
(await database.query(`SELECT cell_id, region FROM relay_cell_regions`)).map((row) => [
text(row, 'cell_id'),
@@ -7590,7 +7646,10 @@ export class RelayAssignmentStore {
) {
throw new Error('activity_lease_shape_mismatch')
}
- const cells = await this.lockCellInventory(database, 'request')
+ // Why: this recomputes one cell's reservation from its leases, so only that
+ // row needs to be held; the 23-row inventory lock here serialised every
+ // desktop control rebind in the fleet behind every other one.
+ const cellRow = (await this.lockCellRows(database, [cellId]))[0]
await database.query(
`DELETE FROM relay_assignment_activity_leases
WHERE user_id = ? AND relay_host_id = ? AND activity_kind = 'control'
@@ -7611,7 +7670,6 @@ export class RelayAssignmentStore {
[cellId]
)
)[0]!
- const cellRow = cells.find((cell) => text(cell, 'cell_id') === cellId)
const cellUnits = integer(cellUnitsRow, 'request_units')
if (!cellRow) throw new Error('assigned_cell_missing')
if (cellUnits > integer(cellRow, 'capacity_requests')) {
diff --git a/cloud/apps/relay/src/cell-inventory-lock-census.test.ts b/cloud/apps/relay/src/cell-inventory-lock-census.test.ts
index 8ca7cee55f5..0ac4c8225e3 100644
--- a/cloud/apps/relay/src/cell-inventory-lock-census.test.ts
+++ b/cloud/apps/relay/src/cell-inventory-lock-census.test.ts
@@ -18,17 +18,21 @@ type CensusEntry = { method: string; mode: CensusMode; reach: Reachability }
// assignment-store.ts, in source order. A new site fails this test until it is
// classified here, which is the point.
const CENSUS: CensusEntry[] = [
- { method: 'assignStickyOnce', mode: 'caller', reach: 'both' },
+ // assignStickyOnce is gone from this list: its retry now locks only the row
+ // the host is pinned to (lockCellRows), which is what a sticky refresh
+ // touches. Placement below is the one genuinely fleet-wide decision left.
{ method: 'assignOnce', mode: 'caller', reach: 'both' },
{ method: 'assignOnce', mode: 'caller', reach: 'both' },
{ method: 'assignOnce', mode: 'nowait', reach: 'both' },
{ method: 'assignOnce', mode: 'nowait', reach: 'both' },
{ method: 'assignOnce', mode: 'nowait', reach: 'both' },
{ method: 'refreshDrainMigrationLeasesOnce', mode: 'request', reach: 'request' },
- // Reachable from neither: changeActivity has no production callers, only tests.
- { method: 'changeActivity', mode: 'request', reach: 'orphan' },
- { method: 'acquireActivity', mode: 'request', reach: 'request' },
- { method: 'activateControl', mode: 'request', reach: 'request' },
+ // changeActivity, acquireActivity, activateControl and
+ // removeSupersededSameCellControls no longer take the inventory: they lock
+ // only the one or two cell rows they touch, in cell_id order (lockCellRows),
+ // so they cannot cycle with placement's ordered inventory lock, and the
+ // 23-row lock there had serialised every reconnect in the fleet behind every
+ // other one.
{ method: 'startEvacuation', mode: 'request', reach: 'request' },
{ method: 'completeEvacuationFromDeadSourceOnce', mode: 'request', reach: 'request' },
{ method: 'completeEvacuationFromDeadSourceOnce', mode: 'nowait', reach: 'request' },
@@ -47,9 +51,33 @@ const CENSUS: CensusEntry[] = [
{ method: 'abortExpiredEvacuations', mode: 'nowait', reach: 'sweep' },
{ method: 'releaseExpiredActivityLeases', mode: 'nowait', reach: 'sweep' },
{ method: 'releaseExpiredActivity', mode: 'nowait', reach: 'sweep' },
- { method: 'reconcileReservationAccounting', mode: 'pool-default', reach: 'both' },
- { method: 'leastLoadedCell', mode: 'pool-default', reach: 'both' },
- { method: 'removeSupersededSameCellControls', mode: 'request', reach: 'request' }
+ // reconcileReservationAccounting and leastLoadedCell are gone too: the first
+ // repairs exactly two cells' counters and now holds only those rows, and the
+ // second selects from the inventory its single caller has already locked.
+]
+
+// Every inline `FROM relay_cells ... FOR UPDATE` outside the named lock helpers,
+// in source order: whole-table locks in reconciliation and sticky placement,
+// and single-row locks for a cell the method is already scoped to (heartbeat,
+// fence, drain generation, configuration, or a reservation adjust that runs
+// under a lock its caller already holds). A new inline lock fails the census
+// below until it is listed here; per-connection paths that touch more than one
+// cell go through lockCellRows so the order is fixed.
+const NAMED_LOCK_HELPERS = ['lockCellInventory', 'lockGeneralCellInventory', 'lockCellRows']
+
+const INLINE_CELL_LOCK_SITES = [
+ 'reconcileCellsWithOptions',
+ 'assignStickyOnce',
+ 'recordCellHeartbeat',
+ 'attestCellFence',
+ 'adoptLegacyCellFence',
+ 'commitLegacyCellFenceAdoption',
+ 'prepareCellFenceAttempt',
+ 'attestCellFenceAttempt',
+ 'attestCellFenceAttempt',
+ 'configureCell',
+ 'assertDrainCellGeneration',
+ 'adjustCellReservation'
]
// The background sweeps, and nothing else. A method reachable from one of these
@@ -151,6 +179,42 @@ describe('cell inventory lock call-site census', () => {
)
})
+ // Why: the census only sees lockCellInventory calls, so a hand-written
+ // `relay_cells ... FOR UPDATE` would escape classification entirely.
+ it('routes every relay_cells row lock through a named lock helper', () => {
+ const lines = storeSource()
+ const rawSites: string[] = []
+ // Whole statements, not a fixed window: a wide column list or a raw
+ // FOR UPDATE inside query() must not slip past.
+ const source = lines.join('\n')
+ const bounds: { name: string; start: number }[] = []
+ lines.forEach((line, index) => {
+ const declaration = DECLARATION.exec(line)
+ if (declaration) bounds.push({ name: declaration[1]!, start: index })
+ })
+ const methodAt = (offset: number): string => {
+ const lineIndex = source.slice(0, offset).split('\n').length - 1
+ let name = ''
+ for (const bound of bounds) if (bound.start <= lineIndex) name = bound.name
+ return name
+ }
+ const tick = String.fromCharCode(96)
+ const statementCall = new RegExp(
+ '\\.(queryLocked|query)\\(\\s*' + tick + '([^' + tick + ']*)' + tick,
+ 'g'
+ )
+ for (const call of source.matchAll(statementCall)) {
+ const statement = call[2]!
+ if (!/\bFROM\s+relay_cells\b/.test(statement)) continue
+ const locks = call[1] === 'queryLocked' || /\bFOR\s+UPDATE\b/.test(statement)
+ if (!locks) continue
+ const method = methodAt(call.index)
+ if (NAMED_LOCK_HELPERS.includes(method)) continue
+ rawSites.push(method)
+ }
+ expect(rawSites).toEqual(INLINE_CELL_LOCK_SITES)
+ })
+
it('leaves no call site taking the inventory without naming a mode', () => {
const source = readFileSync(new URL('./assignment-store.ts', import.meta.url), 'utf8')
const unclassified = source
diff --git a/cloud/apps/relay/src/cell-inventory-per-cell-locking-postgres.test.ts b/cloud/apps/relay/src/cell-inventory-per-cell-locking-postgres.test.ts
new file mode 100644
index 00000000000..8c0ebd73f32
--- /dev/null
+++ b/cloud/apps/relay/src/cell-inventory-per-cell-locking-postgres.test.ts
@@ -0,0 +1,206 @@
+import { afterAll, beforeAll, describe, expect, it } from 'vitest'
+import { RelayAssignmentStore } from './assignment-store.js'
+import { openRelayDatabase, type RelayDatabase } from './database.js'
+
+const databaseUrl = process.env.ORCA_RELAY_TEST_POSTGRES_URL
+const describePostgres = databaseUrl ? describe : describe.skip
+
+// Sorted ascending, and the host is pinned to the LAST id on purpose: the
+// fleet-wide lock is one ordered scan, so it holds every earlier row while it
+// waits on the pinned one. Pinning to the first id would make the two locking
+// models indistinguishable.
+const cells = ['a', 'b', 'c'].map((suffix) => ({
+ id: `percell-postgres-${suffix}`,
+ url: `https://percell-postgres-${suffix}.example.com`,
+ capacityRequests: 1_000,
+ connectionHardCap: 600 as const,
+ connectionUnobservedBound: 50
+}))
+const [cellA, cellB, cellC] = cells as [(typeof cells)[0], (typeof cells)[0], (typeof cells)[0]]
+const identity = { userId: 'percell-postgres-user', relayHostId: 'percellhost00001' }
+
+function heartbeat(cell: (typeof cells)[number]) {
+ return {
+ cellId: cell.id,
+ cellUrl: cell.url,
+ cellIncarnation: '11111111-1111-4111-8111-111111111111',
+ startedAt: 50,
+ ready: true,
+ observedRequests: 0,
+ totalConnections: 0,
+ inFlightConnections: 0,
+ reservedConnectionUnits: 0,
+ enforcedConnectionUnits: 0,
+ connectionInclusionWatermark: 1,
+ connectionHardCap: 600 as const,
+ connectionUnobservedBound: 50
+ }
+}
+
+describePostgres('PostgreSQL per-cell inventory locking', () => {
+ const databases: RelayDatabase[] = []
+
+ beforeAll(async () => {
+ for (let index = 0; index < 3; index++) {
+ databases.push(await openRelayDatabase({ databaseUrl, dataDir: '' }))
+ }
+ })
+
+ async function removeTestRows(database: RelayDatabase): Promise {
+ await database.query(
+ `DELETE FROM relay_control_connection_reservations WHERE user_id LIKE 'percell-postgres-%'`
+ )
+ for (const table of [
+ 'relay_assignment_activity_leases',
+ 'relay_post_drain_migration_pins',
+ 'relay_assignment_migration_incarnations',
+ 'relay_assignment_migrations',
+ 'relay_assignment_region_preferences',
+ 'relay_assignments'
+ ]) {
+ await database.query(`DELETE FROM ${table} WHERE user_id LIKE 'percell-postgres-%'`)
+ }
+ for (const cell of cells) {
+ for (const table of [
+ 'relay_cell_connection_snapshots',
+ 'relay_cell_connection_runtime',
+ 'relay_cell_connection_limits',
+ 'relay_cell_runtime',
+ 'relay_cells'
+ ]) {
+ await database.query(`DELETE FROM ${table} WHERE cell_id = ?`, [cell.id])
+ }
+ }
+ }
+
+ afterAll(async () => {
+ if (databases[0]) await removeTestRows(databases[0])
+ for (const connection of databases) await connection.close()
+ })
+
+ async function pinHostToLastCell(store: RelayAssignmentStore): Promise {
+ await store.reconcileCells(cells)
+ for (const cell of cells) await store.recordCellHeartbeat(heartbeat(cell))
+ await store.setCellEnabled(cellA.id, false)
+ await store.setCellEnabled(cellB.id, false)
+ const assignment = await store.assign(identity)
+ expect(assignment.cellId).toBe(cellC.id)
+ await store.setCellEnabled(cellA.id, true)
+ await store.setCellEnabled(cellB.id, true)
+ }
+
+ async function lockWaiterAppeared(database: RelayDatabase): Promise {
+ const deadline = Date.now() + 4_000
+ while (Date.now() < deadline) {
+ const rows = await database.query(
+ `SELECT count(*) AS waiting FROM pg_stat_activity
+ WHERE datname = current_database() AND wait_event_type = 'Lock'`
+ )
+ if (Number(rows[0]!.waiting) > 0) return true
+ await new Promise((resolve) => setTimeout(resolve, 10))
+ }
+ return false
+ }
+
+ // Why: a sticky refresh whose first NOWAIT probe loses retries by taking a
+ // cell row before the assignment row. That retry used to take the whole
+ // inventory, so one busy cell stalled every other cell's reconnects.
+ it('waits only on the pinned cell row while refreshing a sticky assignment', async () => {
+ await removeTestRows(databases[0]!)
+ const store = new RelayAssignmentStore(databases[0]!, () => 100)
+ await pinHostToLastCell(store)
+ // A host whose control lease was already reaped still holds its pin; that
+ // is the shape that reaches the cell-row probe instead of touchAssignment.
+ await databases[0]!.query(
+ `DELETE FROM relay_assignment_activity_leases WHERE user_id = ?`,
+ [identity.userId]
+ )
+
+ let releaseRow!: () => void
+ const rowReleased = new Promise((resolve) => {
+ releaseRow = resolve
+ })
+ let rowHeld!: () => void
+ const rowHeldPromise = new Promise((resolve) => {
+ rowHeld = resolve
+ })
+ const holder = databases[1]!.transaction(async (transaction) => {
+ await transaction.queryLocked(`SELECT * FROM relay_cells WHERE cell_id = ?`, [cellC.id])
+ rowHeld()
+ await rowReleased
+ })
+ await rowHeldPromise
+
+ const refresh = store.assign(identity)
+ expect(await lockWaiterAppeared(databases[2]!)).toBe(true)
+ // The refresh is blocked on cell C. Every earlier row must still be free:
+ // the ordered fleet-wide scan would be holding both of them by now.
+ const heldWhileRefreshWaits: string[] = []
+ await databases[2]!.transaction(async (transaction) => {
+ for (const cell of [cellA, cellB]) {
+ try {
+ await transaction.queryLocked(
+ `SELECT * FROM relay_cells WHERE cell_id = ?`,
+ [cell.id],
+ { failIfUnavailable: true }
+ )
+ } catch {
+ heldWhileRefreshWaits.push(cell.id)
+ }
+ }
+ })
+ releaseRow()
+ await holder
+
+ expect(heldWhileRefreshWaits).toEqual([])
+ expect((await refresh).cellId).toBe(cellC.id)
+ }, 15_000)
+
+ // Why: the counter moves by a delta now instead of an absolute value read
+ // from a snapshot, so concurrent movement on the same cell must still sum.
+ it('keeps a cell reservation exact under concurrent same-cell activity', async () => {
+ await removeTestRows(databases[0]!)
+ const store = new RelayAssignmentStore(databases[0]!, () => 100)
+ await store.reconcileCells(cells)
+ for (const cell of cells) await store.recordCellHeartbeat(heartbeat(cell))
+ await store.setCellEnabled(cellA.id, false)
+ await store.setCellEnabled(cellB.id, false)
+
+ const hosts = Array.from({ length: 6 }, (_, index) => ({
+ userId: `percell-postgres-user-${index}`,
+ relayHostId: `percellhost0000${index}`
+ }))
+ const stores = databases.map((database) => new RelayAssignmentStore(database, () => 100))
+ await Promise.all(hosts.map((host, index) => stores[index % stores.length]!.assign(host)))
+
+ // One splice each (2 units) on the same cell, from three connections at once.
+ await Promise.all(
+ hosts.map((host, index) =>
+ stores[index % stores.length]!.acquireActivity(host, {
+ activityId: `splice:percell-${index}`,
+ kind: 'splice',
+ cellId: cellC.id
+ })
+ )
+ )
+ const afterAcquire = await databases[0]!.query(
+ `SELECT reserved_requests FROM relay_cells WHERE cell_id = ?`,
+ [cellC.id]
+ )
+ // 6 pending control grants + 6 splices at 2 units each.
+ expect(Number(afterAcquire[0]!.reserved_requests)).toBe(6 + 12)
+
+ await Promise.all(
+ hosts.map((host, index) =>
+ stores[index % stores.length]!.releaseActivity(host, `splice:percell-${index}`)
+ )
+ )
+ const afterRelease = await databases[0]!.query(
+ `SELECT reserved_requests FROM relay_cells WHERE cell_id = ?`,
+ [cellC.id]
+ )
+ expect(Number(afterRelease[0]!.reserved_requests)).toBe(6)
+ await store.setCellEnabled(cellA.id, true)
+ await store.setCellEnabled(cellB.id, true)
+ }, 15_000)
+})
diff --git a/cloud/apps/relay/src/control-rebind-inventory-lock-postgres.test.ts b/cloud/apps/relay/src/control-rebind-inventory-lock-postgres.test.ts
new file mode 100644
index 00000000000..e990ac1ed1a
--- /dev/null
+++ b/cloud/apps/relay/src/control-rebind-inventory-lock-postgres.test.ts
@@ -0,0 +1,260 @@
+import { afterAll, beforeAll, describe, expect, it } from 'vitest'
+import { RelayAssignmentStore } from './assignment-store.js'
+import { openRelayDatabase, type RelayDatabase } from './database.js'
+
+const databaseUrl = process.env.ORCA_RELAY_TEST_POSTGRES_URL
+const describePostgres = databaseUrl ? describe : describe.skip
+
+// Three cells: the inventory lock covers more than the rows a move touches, and
+// a high-to-low move exposes any lock taken out of cell_id order.
+const cells = [
+ {
+ id: 'rebind-inventory-postgres-a',
+ url: 'https://rebind-inventory-postgres-a.example.com',
+ capacityRequests: 1_000,
+ connectionHardCap: 600 as const,
+ connectionUnobservedBound: 50
+ },
+ {
+ id: 'rebind-inventory-postgres-b',
+ url: 'https://rebind-inventory-postgres-b.example.com',
+ capacityRequests: 1_000,
+ connectionHardCap: 600 as const,
+ connectionUnobservedBound: 50
+ },
+ {
+ id: 'rebind-inventory-postgres-c',
+ url: 'https://rebind-inventory-postgres-c.example.com',
+ capacityRequests: 1_000,
+ connectionHardCap: 600 as const,
+ connectionUnobservedBound: 50
+ }
+]
+const identity = { userId: 'rebind-inventory-postgres-user', relayHostId: 'rebindinvhost001' }
+
+function heartbeat(cell: (typeof cells)[number]) {
+ return {
+ cellId: cell.id,
+ cellUrl: cell.url,
+ cellIncarnation: '11111111-1111-4111-8111-111111111111',
+ startedAt: 50,
+ ready: true,
+ observedRequests: 0,
+ totalConnections: 0,
+ inFlightConnections: 0,
+ reservedConnectionUnits: 0,
+ enforcedConnectionUnits: 0,
+ connectionInclusionWatermark: 1,
+ connectionHardCap: 600 as const,
+ connectionUnobservedBound: 50
+ }
+}
+
+// Why: every desktop control rebind used to take the fleet-wide relay_cells
+// FOR UPDATE lock, so a rebind on one cell queued behind whatever held any
+// other cell's row, until COMMIT (55P03 at the request bound). A rebind only
+// touches its own cell row, so it must proceed while another cell's row is
+// held elsewhere.
+describePostgres('PostgreSQL control rebind under a held cell row', () => {
+ const databases: RelayDatabase[] = []
+
+ beforeAll(async () => {
+ databases.push(
+ await openRelayDatabase({ databaseUrl, dataDir: '' }),
+ await openRelayDatabase({ databaseUrl, dataDir: '' })
+ )
+ })
+
+ async function removeTestRows(database: RelayDatabase): Promise {
+ await database.query(
+ `DELETE FROM relay_control_connection_reservations WHERE user_id = ?`,
+ [identity.userId]
+ )
+ for (const table of [
+ 'relay_assignment_activity_leases',
+ 'relay_post_drain_migration_pins',
+ 'relay_assignment_migration_incarnations',
+ 'relay_assignment_migrations',
+ 'relay_assignments'
+ ]) {
+ await database.query(`DELETE FROM ${table} WHERE user_id = ?`, [identity.userId])
+ }
+ for (const cell of cells) {
+ for (const table of [
+ 'relay_cell_connection_snapshots',
+ 'relay_cell_connection_runtime',
+ 'relay_cell_connection_limits',
+ 'relay_cell_runtime',
+ 'relay_cells'
+ ]) {
+ await database.query(`DELETE FROM ${table} WHERE cell_id = ?`, [cell.id])
+ }
+ }
+ }
+
+ afterAll(async () => {
+ if (databases[0]) await removeTestRows(databases[0])
+ for (const connection of databases) await connection.close()
+ })
+
+ it("rebinds and supersedes a control while another cell's row is held", async () => {
+ // A prior aborted run leaves connection snapshots that reject a replayed watermark.
+ await removeTestRows(databases[0]!)
+ const store = new RelayAssignmentStore(databases[0]!, () => 100)
+ await store.reconcileCells(cells)
+ for (const cell of cells) await store.recordCellHeartbeat(heartbeat(cell))
+ // Pin the host to cell A so placement is deterministic.
+ await store.setCellEnabled(cells[1]!.id, false)
+ await store.setCellEnabled(cells[2]!.id, false)
+ const assignment = await store.assign(identity)
+ expect(assignment.cellId).toBe(cells[0]!.id)
+ await store.setCellEnabled(cells[1]!.id, true)
+ await store.setCellEnabled(cells[2]!.id, true)
+ await store.activateControl(identity, {
+ cellId: cells[0]!.id,
+ assignmentEpoch: assignment.assignmentEpoch,
+ generation: 1,
+ connectionInclusionWatermark: 10
+ })
+
+ // Hold only cell B's row on a second connection, the way a rebind on B
+ // does, for longer than the request-path lock bound.
+ let releaseInventory!: () => void
+ const inventoryReleased = new Promise((resolve) => {
+ releaseInventory = resolve
+ })
+ let inventoryHeld!: () => void
+ const inventoryHeldPromise = new Promise((resolve) => {
+ inventoryHeld = resolve
+ })
+ const holder = databases[1]!.transaction(async (transaction) => {
+ await transaction.queryLocked(`SELECT * FROM relay_cells WHERE cell_id = ?`, [cells[1]!.id])
+ inventoryHeld()
+ await inventoryReleased
+ })
+ await inventoryHeldPromise
+
+ // A generation-2 rebind on cell A supersedes generation 1. It must not
+ // wait on cell B's row.
+ const startedAt = Date.now()
+ const blockedStatement = async (): Promise => {
+ const rows = await databases[1]!.query(
+ `SELECT left(query, 160) AS q FROM pg_stat_activity
+ WHERE datname = current_database() AND wait_event_type = 'Lock'`
+ )
+ return rows.map((row) => String(row.q)).join(' | ')
+ }
+ const timeout = new Promise((_, reject) =>
+ setTimeout(
+ () =>
+ void blockedStatement().then((statement) =>
+ reject(new Error(`rebind on cell A blocked behind cell B's row: ${statement}`))
+ ),
+ 2_000
+ )
+ )
+ const rebound = await Promise.race([
+ store.activateControl(identity, {
+ cellId: cells[0]!.id,
+ assignmentEpoch: assignment.assignmentEpoch,
+ generation: 2,
+ connectionInclusionWatermark: 11
+ }),
+ timeout
+ ])
+ const elapsedMs = Date.now() - startedAt
+ releaseInventory()
+ await holder
+
+ expect(rebound).toBe(`control:${cells[0]!.id}:2`)
+ expect(elapsedMs).toBeLessThan(2_000)
+ const controls = await databases[0]!.query(
+ `SELECT activity_id FROM relay_assignment_activity_leases
+ WHERE user_id = ? AND activity_kind = 'control' ORDER BY activity_id`,
+ [identity.userId]
+ )
+ expect(controls).toEqual([{ activity_id: `control:${cells[0]!.id}:2` }])
+ const reserved = await databases[0]!.query(
+ `SELECT reserved_requests FROM relay_cells WHERE cell_id = ?`,
+ [cells[0]!.id]
+ )
+ expect(Number(reserved[0]!.reserved_requests)).toBe(1)
+ }, 15_000)
+
+ // Why: a phone's activity id is client-chosen and can follow the host across
+ // a migration, so acquireActivity may touch two cell rows. Moving from the
+ // higher cell to the lower one is where an unordered lock cycles with
+ // placement's ascending inventory lock (reproduced live before this fix).
+ it('moves an activity from a higher cell to a lower one in cell_id order', async () => {
+ await removeTestRows(databases[0]!)
+ const [cellA, cellB, cellC] = cells as [typeof cells[0], typeof cells[0], typeof cells[0]]
+ const store = new RelayAssignmentStore(databases[0]!, () => 100)
+ await store.reconcileCells(cells)
+ for (const cell of cells) await store.recordCellHeartbeat(heartbeat(cell))
+ await store.setCellEnabled(cellA.id, false)
+ await store.setCellEnabled(cellB.id, false)
+ const assignment = await store.assign(identity)
+ expect(assignment.cellId).toBe(cellC.id)
+ await store.setCellEnabled(cellA.id, true)
+ await store.setCellEnabled(cellB.id, true)
+ const activityId = 'splice:rebind-inventory-postgres'
+ await store.acquireActivity(identity, { activityId, kind: 'splice', cellId: cellC.id })
+ // The migration makes B authoritative; the lease still sits on C.
+ const migration = await store.startEvacuation(identity, cellB.id)
+ expect(migration.targetCellId).toBe(cellB.id)
+
+ // Hold B elsewhere. An ordered move locks B first and queues here holding
+ // nothing else. Locking C first (the old lease's row, as an unordered move
+ // does) or the whole inventory (which takes A) shows up as a held row.
+ let releaseRow!: () => void
+ const rowReleased = new Promise((resolve) => {
+ releaseRow = resolve
+ })
+ let rowHeld!: () => void
+ const rowHeldPromise = new Promise((resolve) => {
+ rowHeld = resolve
+ })
+ const heldWhileMoverWaits: string[] = []
+ const holder = databases[1]!.transaction(async (transaction) => {
+ await transaction.queryLocked(`SELECT * FROM relay_cells WHERE cell_id = ?`, [cellB.id])
+ rowHeld()
+ await rowReleased
+ for (const cell of [cellA, cellC]) {
+ try {
+ await transaction.queryLocked(`SELECT * FROM relay_cells WHERE cell_id = ?`, [cell.id], {
+ failIfUnavailable: true
+ })
+ } catch {
+ heldWhileMoverWaits.push(cell.id)
+ }
+ }
+ })
+ await rowHeldPromise
+ const move = store.acquireActivity(identity, { activityId, kind: 'splice', cellId: cellB.id })
+ let moved = false
+ void move.then(() => {
+ moved = true
+ })
+ await new Promise((resolve) => setTimeout(resolve, 250))
+ expect(moved).toBe(false)
+ releaseRow()
+ await holder
+ await move
+ expect(heldWhileMoverWaits).toEqual([])
+
+ const reservations = await databases[0]!.query(
+ `SELECT cell_id, reserved_requests FROM relay_cells
+ WHERE cell_id IN (?, ?, ?) ORDER BY cell_id ASC`,
+ [cellA.id, cellB.id, cellC.id]
+ )
+ const reserved = reservations.map((row) => [String(row.cell_id), Number(row.reserved_requests)])
+ expect(reserved).toEqual([
+ [cellA.id, 0],
+ // Migration grant plus the moved splice, as in the SQLite origin-scoped
+ // reservation case: the lock change did not alter accounting.
+ [cellB.id, 6],
+ // The sticky grant stays on the source until the migration completes.
+ [cellC.id, 1]
+ ])
+ }, 15_000)
+})
diff --git a/cloud/apps/relay/src/database-postgres-timeout.test.ts b/cloud/apps/relay/src/database-postgres-timeout.test.ts
index 7aba1234f7f..c9021a9ef18 100644
--- a/cloud/apps/relay/src/database-postgres-timeout.test.ts
+++ b/cloud/apps/relay/src/database-postgres-timeout.test.ts
@@ -2,7 +2,10 @@ import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
const fakes = vi.hoisted(() => ({
configs: [] as Array>,
- query: vi.fn(async () => ({ rows: [], rowCount: 0 })),
+ // Pool construction and pool shutdown interleaved, so "the schema pool is
+ // gone before the serving pool opens" is checkable rather than assumed.
+ lifecycle: [] as string[],
+ query: vi.fn(async (_sql: string) => ({ rows: [], rowCount: 0 })),
release: vi.fn(),
end: vi.fn(async () => undefined)
}))
@@ -13,20 +16,37 @@ vi.mock('pg', () => ({
totalCount = 1
idleCount = 1
waitingCount = 0
- end = fakes.end
on = vi.fn()
connect = vi.fn(async () => ({ query: fakes.query, release: fakes.release }))
+ private readonly label: string
constructor(config: Record) {
fakes.configs.push(config)
+ this.label = `max=${String(config.max)} statement_timeout=${String(config.statement_timeout)}`
+ fakes.lifecycle.push(`open ${this.label}`)
+ }
+
+ async end(): Promise {
+ fakes.lifecycle.push(`end ${this.label}`)
+ await fakes.end()
}
}
}
}))
-import { openRelayDatabase } from './database.js'
+import { openRelayDatabase, relayPostgresStatementTimeoutMs } from './database.js'
import { applyPostgresSchema } from './postgres-schema-startup.js'
+const SCHEMA_POOL = {
+ max: 1,
+ application_name: 'orca-relay/director/director/schema',
+ connectionTimeoutMillis: 2_000,
+ // Why: DDL must not inherit the request deadline.
+ statement_timeout: 0,
+ lock_timeout: 1_000,
+ idle_in_transaction_session_timeout: 5_000
+}
+
afterEach(() => {
vi.restoreAllMocks()
})
@@ -34,9 +54,11 @@ afterEach(() => {
describe('PostgreSQL relay deadlines', () => {
beforeEach(() => {
fakes.configs.length = 0
+ fakes.lifecycle.length = 0
fakes.query.mockClear()
fakes.release.mockClear()
fakes.end.mockClear()
+ delete process.env.ORCA_RELAY_POSTGRES_STATEMENT_TIMEOUT_MS
})
it('bounds pool acquisition, statements, locks, and abandoned transactions', async () => {
@@ -48,6 +70,7 @@ describe('PostgreSQL relay deadlines', () => {
})
expect(fakes.configs).toEqual([
+ expect.objectContaining(SCHEMA_POOL),
expect.objectContaining({
max: 3,
application_name: 'orca-relay/director/director',
@@ -59,6 +82,110 @@ describe('PostgreSQL relay deadlines', () => {
])
await database.close()
})
+
+ // Why: an untimed session left open would be a standing way for request work
+ // to escape the deadline this whole pool config exists to enforce.
+ it('closes the untimed schema pool before the serving pool opens', async () => {
+ const database = await openRelayDatabase({
+ databaseUrl: 'postgresql://relay:secret@127.0.0.1:5432/relay',
+ dataDir: './unused',
+ poolMax: 3,
+ applicationName: 'orca-relay/director/director'
+ })
+
+ expect(fakes.lifecycle).toEqual([
+ 'open max=1 statement_timeout=0',
+ 'end max=1 statement_timeout=0',
+ 'open max=3 statement_timeout=5000'
+ ])
+ await database.close()
+ })
+
+ it('applies the schema on the untimed pool, never on the serving one', async () => {
+ fakes.query.mockClear()
+ const ddl: string[] = []
+ fakes.query.mockImplementation(async (sql: string) => {
+ // Every statement issued before the serving pool exists is schema work.
+ if (fakes.lifecycle.length === 1) ddl.push(sql)
+ return { rows: [], rowCount: 0 }
+ })
+ const database = await openRelayDatabase({
+ databaseUrl: 'postgresql://relay:secret@127.0.0.1:5432/relay',
+ dataDir: './unused'
+ })
+
+ expect(ddl.length).toBeGreaterThan(0)
+ // Statements can open with a leading `--` rationale comment.
+ const body = (statement: string): string =>
+ statement.replace(/^(?:\s*--[^\n]*\n)*\s*/, '')
+ expect(ddl.every((statement) => /^CREATE\b/i.test(body(statement)))).toBe(true)
+ // The backfill is DML, so it stays on the deadline-bearing serving pool.
+ expect(ddl.some((statement) => statement.includes('INSERT INTO'))).toBe(false)
+ await database.close()
+ })
+
+ it('takes the serving statement deadline from the environment', async () => {
+ process.env.ORCA_RELAY_POSTGRES_STATEMENT_TIMEOUT_MS = '2500'
+
+ const database = await openRelayDatabase({
+ databaseUrl: 'postgresql://relay:secret@127.0.0.1:5432/relay',
+ dataDir: './unused'
+ })
+
+ expect(fakes.configs).toEqual([
+ expect.objectContaining({ statement_timeout: 0 }),
+ expect.objectContaining({ statement_timeout: 2_500 })
+ ])
+ await database.close()
+ })
+
+ it.each(['0', '-1', '2.5', 'soon', ' '])(
+ 'refuses %s as a statement deadline instead of running unbounded',
+ (value) => {
+ expect(() =>
+ relayPostgresStatementTimeoutMs({ ORCA_RELAY_POSTGRES_STATEMENT_TIMEOUT_MS: value })
+ ).toThrow('invalid_statement_timeout')
+ }
+ )
+
+ it.each([undefined, ''])('defaults to 5s when the environment says %s', (value) => {
+ expect(
+ relayPostgresStatementTimeoutMs(
+ value === undefined ? {} : { ORCA_RELAY_POSTGRES_STATEMENT_TIMEOUT_MS: value }
+ )
+ ).toBe(5_000)
+ })
+
+ // Why: a statement deadline that reaches the caller as a crash converts a
+ // transient stall into a failed assignment. It aborts the transaction exactly
+ // as a lock timeout does, so it belongs on the same bounded retry.
+ it('retries a statement timeout on a fresh client', async () => {
+ vi.spyOn(console, 'warn').mockImplementation(() => undefined)
+ const database = await openRelayDatabase({
+ databaseUrl: 'postgresql://relay:secret@127.0.0.1:5432/relay',
+ dataDir: './unused'
+ })
+ let attempts = 0
+
+ const result = await database.transaction(async (transaction) => {
+ attempts += 1
+ if (attempts === 1) {
+ await transaction.query('SELECT 1')
+ throw Object.assign(new Error('canceling statement due to statement timeout'), {
+ code: '57014'
+ })
+ }
+ return 'committed'
+ })
+
+ expect(result).toBe('committed')
+ expect(attempts).toBe(2)
+ expect(console.warn).toHaveBeenCalledWith(
+ expect.stringContaining('"event":"orca_relay_postgres_transaction_retry"')
+ )
+ expect(console.warn).toHaveBeenCalledWith(expect.stringContaining('"code":"57014"'))
+ await database.close()
+ })
})
describe('PostgreSQL schema startup', () => {
diff --git a/cloud/apps/relay/src/database-statement-timeout-postgres.test.ts b/cloud/apps/relay/src/database-statement-timeout-postgres.test.ts
new file mode 100644
index 00000000000..6b7ebb0334e
--- /dev/null
+++ b/cloud/apps/relay/src/database-statement-timeout-postgres.test.ts
@@ -0,0 +1,98 @@
+import { afterAll, beforeAll, describe, expect, it } from 'vitest'
+import { openRelayDatabase, type RelayDatabase } from './database.js'
+
+const databaseUrl = process.env.ORCA_RELAY_TEST_POSTGRES_URL
+const describePostgres = databaseUrl ? describe : describe.skip
+const applicationName = 'orca-relay/statement-timeout-postgres'
+
+describePostgres('PostgreSQL statement deadline', () => {
+ const databases: RelayDatabase[] = []
+
+ beforeAll(async () => {
+ databases.push(await openRelayDatabase({ databaseUrl, dataDir: '' }))
+ })
+
+ afterAll(async () => {
+ for (const database of databases) await database.close()
+ })
+
+ it('serves requests under the configured deadline', async () => {
+ const database = await openRelayDatabase({ databaseUrl, dataDir: '', statementTimeoutMs: 300 })
+ databases.push(database)
+
+ expect(await database.query(`SELECT current_setting('statement_timeout') AS statement_timeout`)).toEqual([
+ { statement_timeout: '300ms' }
+ ])
+ })
+
+ // Why: a real 57014 aborts the transaction exactly as a lock timeout does. If
+ // it escapes the bounded retry it becomes a failed assignment instead of a
+ // slow one.
+ it('retries a real statement timeout on a fresh client', async () => {
+ const database = await openRelayDatabase({ databaseUrl, dataDir: '', statementTimeoutMs: 300 })
+ databases.push(database)
+ let attempts = 0
+
+ const result = await database.transaction(async (transaction) => {
+ attempts += 1
+ if (attempts === 1) await transaction.query(`SELECT pg_sleep(2)`)
+ return attempts
+ })
+
+ expect(result).toBe(2)
+ }, 15_000)
+
+ // Why: DDL runs on its own untimed connection. relay_invites carries a
+ // CREATE INDEX IF NOT EXISTS, which (unlike CREATE TABLE IF NOT EXISTS)
+ // really does queue behind an ACCESS EXCLUSIVE lock on the table.
+ it('applies the schema behind a held ACCESS EXCLUSIVE lock', async () => {
+ let releaseTable!: () => void
+ const tableReleased = new Promise((resolve) => {
+ releaseTable = resolve
+ })
+ let tableHeld!: () => void
+ const tableHeldPromise = new Promise((resolve) => {
+ tableHeld = resolve
+ })
+ const holder = databases[0]!.transaction(async (transaction) => {
+ await transaction.query(`LOCK TABLE relay_invites IN ACCESS EXCLUSIVE MODE`)
+ tableHeld()
+ await tableReleased
+ })
+ await tableHeldPromise
+
+ const opening = openRelayDatabase({
+ databaseUrl,
+ dataDir: '',
+ applicationName,
+ // Far too short for a blocked DDL; the serving pool wears it, the schema
+ // connection must not.
+ statementTimeoutMs: 200
+ })
+ const blockedOnSchemaConnection = async (): Promise => {
+ const deadline = Date.now() + 4_000
+ while (Date.now() < deadline) {
+ const rows = await databases[0]!.query(
+ `SELECT count(*) AS waiting FROM pg_stat_activity
+ WHERE datname = current_database() AND wait_event_type = 'Lock'
+ AND application_name = ?`,
+ [`${applicationName}/schema`]
+ )
+ if (Number(rows[0]!.waiting) > 0) return true
+ await new Promise((resolve) => setTimeout(resolve, 10))
+ }
+ return false
+ }
+ const blocked = await blockedOnSchemaConnection()
+ releaseTable()
+ await holder
+
+ const database = await opening
+ databases.push(database)
+ expect(blocked).toBe(true)
+ // The serving pool still carries the short deadline it was opened with.
+ expect(await database.query(`SELECT current_setting('statement_timeout') AS statement_timeout`)).toEqual([
+ { statement_timeout: '200ms' }
+ ])
+ }, 15_000)
+})
diff --git a/cloud/apps/relay/src/database.ts b/cloud/apps/relay/src/database.ts
index 326ab010ccb..2558831ca64 100644
--- a/cloud/apps/relay/src/database.ts
+++ b/cloud/apps/relay/src/database.ts
@@ -796,12 +796,34 @@ class PostgresTransaction implements RelayDatabase {
const POSTGRES_TRANSACTION_ATTEMPTS = 3
const POSTGRES_RETRY_MAX_DELAY_MS = 25
const POSTGRES_CONNECTION_TIMEOUT_MS = 2_000
-const POSTGRES_STATEMENT_TIMEOUT_MS = 5_000
+// Derivation: a control renewal must land inside its own 30s tick
+// (RELAY_PROTOCOL_LIMITS.controlPingIntervalMs * 2), and a transaction gets
+// POSTGRES_TRANSACTION_ATTEMPTS tries, so the worst case a renewal can spend in
+// Postgres is attempts * timeout. 5s keeps that at 15s, half the tick, and still
+// leaves room for the connect timeout above.
+export const POSTGRES_STATEMENT_TIMEOUT_MS = 5_000
const POSTGRES_IDLE_TRANSACTION_TIMEOUT_MS = 5_000
+export function relayPostgresStatementTimeoutMs(
+ env: NodeJS.ProcessEnv = process.env
+): number {
+ const configured = env.ORCA_RELAY_POSTGRES_STATEMENT_TIMEOUT_MS
+ if (configured === undefined || configured === '') return POSTGRES_STATEMENT_TIMEOUT_MS
+ const milliseconds = Number(configured)
+ // 0 is PostgreSQL's "no timeout"; refusing it keeps the deadline this exists
+ // to enforce from being disabled by a typo in an environment variable.
+ if (!Number.isInteger(milliseconds) || milliseconds < 1) {
+ throw new Error('invalid_statement_timeout')
+ }
+ return milliseconds
+}
+
function retryablePostgresTransactionError(error: unknown): boolean {
const code = String((error as { code?: unknown }).code)
- return code === '40P01' || code === '40001' || code === '55P03'
+ // 57014 is the pool statement_timeout firing. It aborts the transaction the
+ // same way a lock timeout does, so it belongs on the bounded retry path
+ // rather than surfacing as a terminal failure to the caller.
+ return code === '40P01' || code === '40001' || code === '55P03' || code === '57014'
}
export function isRelayDatabaseTransientError(error: unknown): boolean {
@@ -963,11 +985,36 @@ async function applySchema(database: RelayDatabase): Promise {
}
}
-async function applySchemaWithPostgresRetries(database: RelayDatabase): Promise {
- await applyPostgresSchema(
- SCHEMA.split(';').filter((statement) => statement.trim()),
- async (statement) => await database.query(statement)
- )
+// Why: DDL is not a request. A CREATE INDEX on a grown table legitimately runs
+// longer than the request statement_timeout, and inheriting that timeout would
+// make every startup fail at the same statement instead of finishing once. One
+// short-lived connection of its own, ended before the serving pool opens, keeps
+// the untimed session off the request path entirely.
+async function applySchemaOnUntimedPool(
+ databaseUrl: string,
+ applicationName: string | undefined
+): Promise {
+ const pool = new pg.Pool({
+ connectionString: databaseUrl,
+ max: 1,
+ application_name: applicationName ? `${applicationName}/schema` : undefined,
+ connectionTimeoutMillis: POSTGRES_CONNECTION_TIMEOUT_MS,
+ statement_timeout: 0,
+ // Kept: a DDL blocked behind another director's ACCESS EXCLUSIVE lock must
+ // yield to the bounded schema retry instead of holding the connection.
+ lock_timeout: POSTGRES_LOCK_TIMEOUT_MS,
+ idle_in_transaction_session_timeout: POSTGRES_IDLE_TRANSACTION_TIMEOUT_MS
+ })
+ absorbPostgresIdleClientErrors(pool)
+ const database = new PostgresDatabase(pool)
+ try {
+ await applyPostgresSchema(
+ SCHEMA.split(';').filter((statement) => statement.trim()),
+ async (statement) => await database.query(statement)
+ )
+ } finally {
+ await database.close().catch(() => undefined)
+ }
}
async function backfillRelayCellRegions(database: RelayDatabase): Promise {
@@ -983,15 +1030,17 @@ export async function openRelayDatabase(input: {
dataDir: string
poolMax?: number
applicationName?: string
+ statementTimeoutMs?: number
}): Promise {
let database: RelayDatabase
if (input.databaseUrl) {
+ await applySchemaOnUntimedPool(input.databaseUrl, input.applicationName)
const pool = new pg.Pool({
connectionString: input.databaseUrl,
max: input.poolMax ?? 10,
application_name: input.applicationName,
connectionTimeoutMillis: POSTGRES_CONNECTION_TIMEOUT_MS,
- statement_timeout: POSTGRES_STATEMENT_TIMEOUT_MS,
+ statement_timeout: input.statementTimeoutMs ?? relayPostgresStatementTimeoutMs(),
lock_timeout: POSTGRES_LOCK_TIMEOUT_MS,
idle_in_transaction_session_timeout: POSTGRES_IDLE_TRANSACTION_TIMEOUT_MS
})
@@ -1004,8 +1053,7 @@ export async function openRelayDatabase(input: {
database = new SqliteDatabase(sqlite)
}
try {
- if (input.databaseUrl) await applySchemaWithPostgresRetries(database)
- else await applySchema(database)
+ if (!input.databaseUrl) await applySchema(database)
await backfillRelayCellRegions(database)
return database
} catch (error) {
diff --git a/cloud/apps/relay/src/host-close-reason-memory.test.ts b/cloud/apps/relay/src/host-close-reason-memory.test.ts
new file mode 100644
index 00000000000..2985e6f1a1d
--- /dev/null
+++ b/cloud/apps/relay/src/host-close-reason-memory.test.ts
@@ -0,0 +1,82 @@
+import { ASSIGNMENT_LIMITS, RELAY_HOST_CLOSE_REASON } from '@orca-cloud/relay-contract'
+import { describe, expect, it } from 'vitest'
+import { HostCloseReasonMemory } from './host-close-reason-memory.js'
+
+function memoryAt(clock: { now: number }): HostCloseReasonMemory {
+ return new HostCloseReasonMemory(() => clock.now)
+}
+
+describe('HostCloseReasonMemory', () => {
+ it('remembers only reasons it knows', () => {
+ const clock = { now: 1_000 }
+ const memory = memoryAt(clock)
+
+ memory.record('a', RELAY_HOST_CLOSE_REASON.SIGNED_OUT)
+ memory.record('b', 'quitting')
+ memory.record('c', Buffer.alloc(0))
+ memory.record('d', undefined)
+
+ expect(memory.read('a')).toBe(RELAY_HOST_CLOSE_REASON.SIGNED_OUT)
+ expect(memory.read('b')).toBeNull()
+ expect(memory.read('c')).toBeNull()
+ expect(memory.read('d')).toBeNull()
+ })
+
+ it('accepts the reason as the Buffer a ws close delivers', () => {
+ const clock = { now: 1_000 }
+ const memory = memoryAt(clock)
+
+ memory.record('a', Buffer.from(RELAY_HOST_CLOSE_REASON.SIGNED_OUT))
+
+ expect(memory.read('a')).toBe(RELAY_HOST_CLOSE_REASON.SIGNED_OUT)
+ })
+
+ it('expires an entry once its host may have been rebalanced away', () => {
+ const clock = { now: 1_000 }
+ const memory = memoryAt(clock)
+ memory.record('a', RELAY_HOST_CLOSE_REASON.SIGNED_OUT)
+
+ clock.now += ASSIGNMENT_LIMITS.dormantTtlMs - 1
+ expect(memory.read('a')).toBe(RELAY_HOST_CLOSE_REASON.SIGNED_OUT)
+
+ clock.now += 1
+ expect(memory.read('a')).toBeNull()
+ expect(memory.size()).toBe(0)
+ })
+
+ it('forgets on demand', () => {
+ const clock = { now: 1_000 }
+ const memory = memoryAt(clock)
+ memory.record('a', RELAY_HOST_CLOSE_REASON.SIGNED_OUT)
+
+ memory.forget('a')
+
+ expect(memory.read('a')).toBeNull()
+ })
+
+ it('drops the oldest survivors rather than growing without bound', () => {
+ const clock = { now: 1_000 }
+ const memory = memoryAt(clock)
+ for (let index = 0; index < 50_050; index++) {
+ memory.record(`host-${index}`, RELAY_HOST_CLOSE_REASON.SIGNED_OUT)
+ }
+
+ expect(memory.size()).toBe(50_000)
+ expect(memory.read('host-0')).toBeNull()
+ expect(memory.read('host-50049')).toBe(RELAY_HOST_CLOSE_REASON.SIGNED_OUT)
+ })
+
+ it('re-recording refreshes recency so a live host is not evicted first', () => {
+ const clock = { now: 1_000 }
+ const memory = memoryAt(clock)
+ memory.record('a', RELAY_HOST_CLOSE_REASON.SIGNED_OUT)
+ memory.record('b', RELAY_HOST_CLOSE_REASON.SIGNED_OUT)
+ memory.record('a', RELAY_HOST_CLOSE_REASON.SIGNED_OUT)
+
+ expect([...['a', 'b'].map((key) => memory.read(key))]).toEqual([
+ RELAY_HOST_CLOSE_REASON.SIGNED_OUT,
+ RELAY_HOST_CLOSE_REASON.SIGNED_OUT
+ ])
+ expect(memory.size()).toBe(2)
+ })
+})
diff --git a/cloud/apps/relay/src/host-close-reason-memory.ts b/cloud/apps/relay/src/host-close-reason-memory.ts
new file mode 100644
index 00000000000..ed01aacd666
--- /dev/null
+++ b/cloud/apps/relay/src/host-close-reason-memory.ts
@@ -0,0 +1,72 @@
+import {
+ ASSIGNMENT_LIMITS,
+ relayHostCloseReasonFrom,
+ type RelayHostCloseReason
+} from '@orca-cloud/relay-contract'
+
+// Retention matches the dormant assignment TTL: past it the host may have been
+// rebalanced onto another cell, so this cell is no longer the one a phone asks.
+const RETENTION_MS = ASSIGNMENT_LIMITS.dormantTtlMs
+// A fleet-wide auth outage signs out every host at once; the cap bounds that
+// burst well above any single cell's host count without becoming a leak.
+const MAX_ENTRIES = 50_000
+
+// Why in-memory and not Postgres: a phone reaches the cell its host's assignment
+// row already names, which is the same cell that watched the control socket
+// close. Losing this on a cell restart degrades to the pre-existing generic
+// verdict, so the failure mode is the old behaviour rather than a wrong one.
+export class HostCloseReasonMemory {
+ private readonly entries = new Map()
+
+ constructor(private readonly now: () => number = Date.now) {}
+
+ // Silently ignores anything that is not a known reason, which is every close
+ // from a host that predates the field and every abrupt 1006.
+ record(key: string, reason: unknown): void {
+ const parsed = relayHostCloseReasonFrom(reason)
+ if (!parsed) {
+ return
+ }
+ this.entries.delete(key)
+ this.entries.set(key, { reason: parsed, expiresAt: this.now() + RETENTION_MS })
+ this.evict()
+ }
+
+ forget(key: string): void {
+ this.entries.delete(key)
+ }
+
+ read(key: string): RelayHostCloseReason | null {
+ const entry = this.entries.get(key)
+ if (!entry) {
+ return null
+ }
+ if (entry.expiresAt <= this.now()) {
+ this.entries.delete(key)
+ return null
+ }
+ return entry.reason
+ }
+
+ size(): number {
+ return this.entries.size
+ }
+
+ private evict(): void {
+ const now = this.now()
+ for (const [key, entry] of this.entries) {
+ if (entry.expiresAt > now) {
+ break
+ }
+ this.entries.delete(key)
+ }
+ // Insertion order is recency order (record deletes before setting), so the
+ // head is always the oldest survivor.
+ for (const key of this.entries.keys()) {
+ if (this.entries.size <= MAX_ENTRIES) {
+ break
+ }
+ this.entries.delete(key)
+ }
+ }
+}
diff --git a/cloud/apps/relay/src/host-session-registry.ts b/cloud/apps/relay/src/host-session-registry.ts
index 11b7d1de030..5c53041e7ff 100644
--- a/cloud/apps/relay/src/host-session-registry.ts
+++ b/cloud/apps/relay/src/host-session-registry.ts
@@ -14,7 +14,8 @@ import {
HostHelloSchema,
InviteCreateSchema,
RELAY_PROTOCOL_LIMITS,
- RELAY_CLOSE_CODE
+ RELAY_CLOSE_CODE,
+ type RelayHostCloseReason
} from '@orca-cloud/relay-contract'
import nacl from 'tweetnacl'
import type WebSocket from 'ws'
@@ -25,6 +26,7 @@ import {
RelayCredentialStore,
type CredentialReservation
} from './credential-store.js'
+import { HostCloseReasonMemory } from './host-close-reason-memory.js'
import { relayHostLogDigest } from './relay-host-log-digest.js'
import type { RelayTokenClaims } from './relay-token-verifier.js'
import type { RelayRuntimeObserver } from './relay-observability.js'
@@ -130,6 +132,10 @@ const ACTIVATION_QUEUE_WAIT_MS = 30_000
export class HostSessionRegistry {
private readonly sessions = new Map()
private readonly activationQueues = new Map>()
+ // Why it outlives `sessions`: the orphan grace deletes the session within 30s,
+ // but a signed-out desktop never comes back, so the phone that asks minutes
+ // later would otherwise find nothing to explain its rejection with.
+ private readonly hostCloseReasons = new HostCloseReasonMemory(() => this.now())
private draining = false
constructor(
@@ -175,7 +181,8 @@ export class HostSessionRegistry {
return
}
this.observer.recordAuth(true)
- const session = this.sessions.get(this.key(reservation.userId, hostId))
+ const sessionKey = this.key(reservation.userId, hostId)
+ const session = this.sessions.get(sessionKey)
if (
!session ||
session.state !== 'active' ||
@@ -184,7 +191,13 @@ export class HostSessionRegistry {
) {
capacityReservation?.release()
await this.store.failReservation(reservation)
- this.rejectClient(socket, RELAY_CLOSE_CODE.HOST_OFFLINE)
+ // The only rejection that can name a cause: the host is genuinely absent.
+ // The attach-deadline 4404 below fires while control is still connected.
+ this.rejectClient(
+ socket,
+ RELAY_CLOSE_CODE.HOST_OFFLINE,
+ this.hostCloseReasons.read(sessionKey)
+ )
return
}
if (session.activeConnIds.size + session.pendingConns.size >= 8) {
@@ -793,7 +806,10 @@ export class HostSessionRegistry {
regionalDrainTimer: null,
regionalDrainExpiresAt: null
}
- this.sessions.set(this.key(identity.sub, identity.relayHostId), session)
+ const sessionKey = this.key(identity.sub, identity.relayHostId)
+ // A host that proved itself again is not signed out, whatever it said last.
+ this.hostCloseReasons.forget(sessionKey)
+ this.sessions.set(sessionKey, session)
this.wireActiveControl(session)
this.sendHelloAck(session)
}
@@ -813,6 +829,11 @@ export class HostSessionRegistry {
})
socket.once('close', (code, reason) => {
this.observer.recordControlClose?.(code)
+ // Guarded on identity: a predecessor retired by a rebind must not stamp a
+ // cause onto the live session that replaced it.
+ if (session.socket === socket) {
+ this.hostCloseReasons.record(this.key(session.identity.sub, session.relayHostId), reason)
+ }
// One line per control close makes reconnect churners attributable by
// host digest without exposing the raw relay host id.
console.warn(
@@ -1187,9 +1208,16 @@ export class HostSessionRegistry {
if (session.socket) send(session.socket, 'control-error', { ...(reqId ? { reqId } : {}), code })
}
- private rejectClient(socket: WebSocket, code: number): void {
+ // hostCloseReason rides the WebSocket close reason, never relay-hello: every
+ // shipped phone parses relay-hello with a strict schema that rejects an
+ // unknown key, and none of them read the close reason at all.
+ private rejectClient(
+ socket: WebSocket,
+ code: number,
+ hostCloseReason?: RelayHostCloseReason | null
+ ): void {
send(socket, 'relay-hello', { ok: false, code })
- closeRelayWebSocket(socket, code, 'relay connection rejected')
+ closeRelayWebSocket(socket, code, hostCloseReason ?? 'relay connection rejected')
}
private releaseControlActivity(session: HostSession): void {
diff --git a/cloud/apps/relay/src/host-signed-out-rejection.test.ts b/cloud/apps/relay/src/host-signed-out-rejection.test.ts
new file mode 100644
index 00000000000..0f8542c960f
--- /dev/null
+++ b/cloud/apps/relay/src/host-signed-out-rejection.test.ts
@@ -0,0 +1,206 @@
+import { EventEmitter } from 'node:events'
+import {
+ CONTROL_CONTINUITY_LIMITS,
+ RELAY_CLOSE_CODE,
+ RELAY_HOST_CLOSE_REASON
+} 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 { RelayCredentialStore } from './credential-store.js'
+import { 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'
+
+class FakeSocket extends EventEmitter {
+ readonly OPEN = 1
+ 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', 1006, Buffer.alloc(0))
+ })
+}
+
+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: []
+} as unknown as RelayConfig
+
+const identity = {
+ sub: 'user-1',
+ prof: 'profile-1',
+ org: 'org-1',
+ relayHostId: 'AbCdEf0123_-xyZ9'
+} as unknown as RelayTokenClaims
+
+const reservation = {
+ userId: identity.sub,
+ relayHostId: identity.relayHostId,
+ credentialKind: 'resume',
+ relayDeviceId: 'device-1',
+ leaseExpiresAt: Date.now() + 60_000
+}
+
+function createRegistry() {
+ const store = {
+ resolveResume: vi.fn().mockResolvedValue({ userId: identity.sub }),
+ reserveCredential: vi.fn().mockResolvedValue(reservation),
+ failReservation: vi.fn().mockResolvedValue(undefined)
+ }
+ const assignments = {
+ activateControl: vi.fn().mockResolvedValue('control:production-gce-c3:1'),
+ markMigrationTargetRegistered: vi.fn().mockResolvedValue(undefined),
+ resolve: vi.fn().mockResolvedValue({ cellId: config.cellId }),
+ acquireActivity: vi.fn().mockResolvedValue(undefined),
+ renewControlActivity: vi.fn().mockResolvedValue(undefined),
+ releaseActivity: vi.fn().mockResolvedValue(true)
+ } as unknown as RelayAssignmentStore
+ const observer = {
+ recordAuth: vi.fn(),
+ recordForwardedBytes: vi.fn(),
+ recordHttp: vi.fn(),
+ recordReconnect: vi.fn(),
+ recordSql: vi.fn(),
+ recordControlClose: vi.fn(),
+ recordSpliceClose: vi.fn()
+ } satisfies RelayRuntimeObserver
+ const registry = new HostSessionRegistry(
+ config,
+ vi.fn(),
+ store as unknown as RelayCredentialStore,
+ assignments,
+ new ProcessQueuedByteBudget(),
+ observer
+ )
+ const activate = (socket: WebSocket, generation: number): Promise =>
+ (
+ registry as unknown as {
+ activate: (
+ socket: WebSocket,
+ identity: RelayTokenClaims,
+ existing: null,
+ generation: number,
+ rebind: boolean,
+ assignmentEpoch: number,
+ appVersion: string
+ ) => Promise
+ }
+ ).activate(socket, identity, null, generation, false, 1, '1.4.173')
+ return { registry, activate }
+}
+
+async function dialPhone(registry: HostSessionRegistry): Promise {
+ const phone = new FakeSocket()
+ await registry.acceptClient(phone as unknown as WebSocket, identity.relayHostId, 'credential')
+ return phone
+}
+
+// The 4404 hello body is unchanged: every shipped phone parses it with a strict
+// schema, so the cause has to ride the close frame instead.
+const HOST_OFFLINE_HELLO = JSON.stringify({
+ type: 'relay-hello',
+ ok: false,
+ code: RELAY_CLOSE_CODE.HOST_OFFLINE
+})
+
+describe('host sign-out reason on phone rejection', () => {
+ beforeEach(() => vi.useFakeTimers())
+ afterEach(() => {
+ vi.clearAllTimers()
+ vi.useRealTimers()
+ })
+
+ it('names the sign-out to a phone that arrives after the host is gone', async () => {
+ const { registry, activate } = createRegistry()
+ const control = new FakeSocket()
+ await activate(control as unknown as WebSocket, 1)
+
+ control.close(1000, RELAY_HOST_CLOSE_REASON.SIGNED_OUT)
+ vi.advanceTimersByTime(CONTROL_CONTINUITY_LIMITS.orphanGraceMs + 1)
+
+ const phone = await dialPhone(registry)
+ expect(phone.send).toHaveBeenCalledWith(HOST_OFFLINE_HELLO)
+ expect(phone.close).toHaveBeenCalledWith(
+ RELAY_CLOSE_CODE.HOST_OFFLINE,
+ RELAY_HOST_CLOSE_REASON.SIGNED_OUT
+ )
+ })
+
+ it('says nothing when the host died without naming a cause', async () => {
+ const { registry, activate } = createRegistry()
+ const control = new FakeSocket()
+ await activate(control as unknown as WebSocket, 1)
+
+ control.terminate()
+ vi.advanceTimersByTime(CONTROL_CONTINUITY_LIMITS.orphanGraceMs + 1)
+
+ const phone = await dialPhone(registry)
+ expect(phone.close).toHaveBeenCalledWith(
+ RELAY_CLOSE_CODE.HOST_OFFLINE,
+ 'relay connection rejected'
+ )
+ })
+
+ it('ignores a close reason the host invented', async () => {
+ const { registry, activate } = createRegistry()
+ const control = new FakeSocket()
+ await activate(control as unknown as WebSocket, 1)
+
+ control.close(1000, 'signed-out-ish')
+ vi.advanceTimersByTime(CONTROL_CONTINUITY_LIMITS.orphanGraceMs + 1)
+
+ const phone = await dialPhone(registry)
+ expect(phone.close).toHaveBeenCalledWith(
+ RELAY_CLOSE_CODE.HOST_OFFLINE,
+ 'relay connection rejected'
+ )
+ })
+
+ it('forgets the sign-out once the host proves itself again', async () => {
+ const { registry, activate } = createRegistry()
+ const control = new FakeSocket()
+ await activate(control as unknown as WebSocket, 1)
+ control.close(1000, RELAY_HOST_CLOSE_REASON.SIGNED_OUT)
+ vi.advanceTimersByTime(CONTROL_CONTINUITY_LIMITS.orphanGraceMs + 1)
+
+ const reconnected = new FakeSocket()
+ await activate(reconnected as unknown as WebSocket, 2)
+ // Drop it abruptly, as a network death would, so only the stale memory
+ // could still name a cause.
+ reconnected.terminate()
+ vi.advanceTimersByTime(CONTROL_CONTINUITY_LIMITS.orphanGraceMs + 1)
+
+ const phone = await dialPhone(registry)
+ expect(phone.close).toHaveBeenCalledWith(
+ RELAY_CLOSE_CODE.HOST_OFFLINE,
+ 'relay connection rejected'
+ )
+ })
+
+ // A live host is present: the 4404 there is an attach deadline, not absence.
+ it('never names a cause while the host control is connected', async () => {
+ const { registry, activate } = createRegistry()
+ const control = new FakeSocket()
+ await activate(control as unknown as WebSocket, 1)
+
+ const phone = await dialPhone(registry)
+ expect(phone.close).not.toHaveBeenCalled()
+ expect(control.send).toHaveBeenCalledWith(expect.stringContaining('"type":"conn-open"'))
+ })
+})
diff --git a/cloud/apps/relay/src/postgres-transaction-recovery.test.ts b/cloud/apps/relay/src/postgres-transaction-recovery.test.ts
index a49c07d2e7d..ae9d52a7c86 100644
--- a/cloud/apps/relay/src/postgres-transaction-recovery.test.ts
+++ b/cloud/apps/relay/src/postgres-transaction-recovery.test.ts
@@ -503,12 +503,19 @@ describePostgres('PostgreSQL transaction recovery', () => {
const directorLockOrder: string[] = []
const assignmentDatabase = new TransactionProbeDatabase(database, async (phase, sql) => {
if (phase === 'before') {
- if (sql.includes('FROM relay_assignments WHERE user_id = ?')) {
+ // Only locked statements reach this hook, so classifying the pin read
+ // is what proves it stays unlocked: if it ever grows a FOR UPDATE it
+ // shows up in the order below instead of silently joining the queue.
+ if (sql.includes('SELECT cell_id FROM relay_assignments')) {
+ directorLockOrder.push('pin-read')
+ } else if (sql.includes('FROM relay_assignments WHERE user_id = ?')) {
directorLockOrder.push('assignment')
} else if (sql.includes('FROM relay_assignment_activity_leases')) {
directorLockOrder.push('activity')
} else if (sql.includes('FROM relay_cells ORDER BY')) {
directorLockOrder.push('cell-inventory')
+ } else if (sql.includes('FROM relay_cells WHERE cell_id IN')) {
+ directorLockOrder.push('cell-rows')
} else if (sql.includes('FROM relay_cells WHERE cell_id = ?')) {
directorLockOrder.push('cell')
}
@@ -530,11 +537,14 @@ describePostgres('PostgreSQL transaction recovery', () => {
})
await expect(legacyTransaction).resolves.toBeUndefined()
expect(assignmentDatabase.attempts).toBe(2)
+ // The retry still takes a cell row before the assignment row — the order
+ // that avoids the legacy cycle — but only the pinned row, never the
+ // inventory.
expect(directorLockOrder).toEqual([
'assignment',
'activity',
'cell',
- 'cell-inventory',
+ 'cell-rows',
'assignment',
'activity'
])
diff --git a/cloud/dev/fixtures/terraform-root-partition/families.json b/cloud/dev/fixtures/terraform-root-partition/families.json
index 9664c1eb299..dfe100fd2dd 100644
--- a/cloud/dev/fixtures/terraform-root-partition/families.json
+++ b/cloud/dev/fixtures/terraform-root-partition/families.json
@@ -133,10 +133,15 @@
"google_logging_metric.relay_snapshot",
"google_monitoring_alert_policy.relay_assignment_5xx",
"google_monitoring_alert_policy.relay_assignment_edge_429",
+ "google_monitoring_alert_policy.relay_cell_process_exit",
+ "google_monitoring_alert_policy.relay_cloud_nat_port_drops",
"google_monitoring_alert_policy.relay_cloud_sql_backends",
+ "google_monitoring_alert_policy.relay_cloud_sql_checkpoint_loop",
+ "google_monitoring_alert_policy.relay_cloud_sql_disk",
"google_monitoring_alert_policy.relay_custom",
"google_monitoring_alert_policy.relay_gce_connection_headroom",
"google_monitoring_alert_policy.relay_postgres_retry_exhausted",
+ "google_monitoring_dashboard.relay_incident",
"google_project_iam_custom_role.github_production_relay_capacity_mutation",
"google_project_iam_custom_role.github_relay_asia_topology_mutation",
"google_project_iam_custom_role.github_relay_asia_topology_read",
diff --git a/cloud/docs/relay-incident-monitor.md b/cloud/docs/relay-incident-monitor.md
index 337d3f1b20f..870c95dd413 100644
--- a/cloud/docs/relay-incident-monitor.md
+++ b/cloud/docs/relay-incident-monitor.md
@@ -99,7 +99,7 @@ durably marked consumed before mutation and cannot authorize another run.
| Cloud SQL deadlocks | over 0 |
| Relay pool waiters | over 800 |
| Relay pool wait | over 2,500 ms |
-| PostgreSQL retries in five minutes | over 300 |
+| PostgreSQL retries in five minutes | over 2,000 |
| Exhausted PostgreSQL retries in five minutes | over 300 |
| Director instances | outside 5–6 |
| Director CPU or memory | over 80% |
@@ -138,7 +138,23 @@ heartbeats, and matching live admission.
`jsonPayload.event="orca_relay_postgres_transaction_retry"` in production
logs: healthy-day bursts reach 234/5min with zero exhausted retries and 26%
of five-minute windows over 20, while the 2026-08-23 lock-contention
- incident ran roughly 2,200–3,000/5min.
+ incident ran roughly 2,200–3,000/5min by raw log-line count (the gate's
+ own `orca_relay_postgres_retries` metric read 1,510 for that window; see the
+ 2026-09-04 entry).
+- Recalibrated the PostgreSQL-retry freeze from 300 to 2,000 per five minutes
+ (2026-09-04). Basis: the global `relay_cells FOR UPDATE` lock made
+ successful retries a steady-state rate. Measured fleet-wide (director +
+ cells, summed per five minutes from the `orca_relay_postgres_retries`
+ log metric) over 2026-09-03T05Z..2026-09-04T05Z: p50 430 / p90 924 /
+ p99 1,320 / max 1,504; 55% of windows over 300; only 22% of 15-minute gates
+ clean at 300 versus 100% at 2,000. Three read-only dry-runs on 2026-09-04
+ froze on this bar (runs 33836470590, 33838698725) or on a genuine six-cell
+ crash storm (33837160275), blocking the same-cap roll that carries #18521
+ and the `beginProof` crash guard to the 23 cells. The 2026-08-23 incident
+ on this metric peaked at 1,510 then 646, so retries alone no longer
+ separate it from today's baseline; the exhausted-retry bar (incident peak
+ 467 vs bar 300), director concurrency, and the pool bars carry that role.
+ Re-tighten after the fleet is on the 500 ms lock wait.
- Recalibrated the exhausted-PostgreSQL-retry freeze from 0 to 300 per five
minutes (2026-09-04). Basis: #18521 cut the request-path cell-inventory
lock wait from the 1 s pool `lock_timeout` to 500 ms, so contended waiters
diff --git a/cloud/infra/terraform/relay-gce-cells.tf b/cloud/infra/terraform/relay-gce-cells.tf
index d6b7f3351f9..a4505ba2e37 100644
--- a/cloud/infra/terraform/relay-gce-cells.tf
+++ b/cloud/infra/terraform/relay-gce-cells.tf
@@ -242,6 +242,7 @@ resource "google_compute_instance_template" "relay_gce_cell" {
artifact_registry_host = "${var.region}-docker.pkg.dev"
relay_image = each.value.image
cloud_sql_proxy_image = var.relay_gce_cloud_sql_proxy_image
+ cloud_sql_private_ip = var.relay_cloud_sql_private_ip
# Keep cell-only plans independent from unrelated database configuration drift.
cloud_sql_connection_name = local.relay_database_connection_name
})
diff --git a/cloud/infra/terraform/relay-gce-foundation.tf b/cloud/infra/terraform/relay-gce-foundation.tf
index aab3b4579eb..a8d64b3fcea 100644
--- a/cloud/infra/terraform/relay-gce-foundation.tf
+++ b/cloud/infra/terraform/relay-gce-foundation.tf
@@ -42,6 +42,12 @@ resource "google_compute_router_nat" "relay_gce" {
router = google_compute_router.relay_gce[0].name
nat_ip_allocate_option = "AUTO_ONLY"
source_subnetwork_ip_ranges_to_nat = "LIST_OF_SUBNETWORKS"
+ # Cells reach Cloud SQL's public IP through this NAT. The static default of 64 ports per VM
+ # filled during the 2026-09-04 incident and every cell's proxy dial timed out at once.
+ enable_dynamic_port_allocation = true
+ enable_endpoint_independent_mapping = false
+ min_ports_per_vm = 64
+ max_ports_per_vm = 4096
subnetwork {
name = google_compute_subnetwork.relay_gce[0].id
@@ -85,6 +91,12 @@ resource "google_compute_router_nat" "relay_gce_additional" {
router = google_compute_router.relay_gce_additional[each.key].name
nat_ip_allocate_option = "AUTO_ONLY"
source_subnetwork_ip_ranges_to_nat = "LIST_OF_SUBNETWORKS"
+ # Cells reach Cloud SQL's public IP through this NAT. The static default of 64 ports per VM
+ # filled during the 2026-09-04 incident and every cell's proxy dial timed out at once.
+ enable_dynamic_port_allocation = true
+ enable_endpoint_independent_mapping = false
+ min_ports_per_vm = 64
+ max_ports_per_vm = 4096
subnetwork {
name = google_compute_subnetwork.relay_gce_additional[each.key].id
diff --git a/cloud/infra/terraform/relay-gce-startup.sh.tftpl b/cloud/infra/terraform/relay-gce-startup.sh.tftpl
index f593d94e9e5..a77466f2169 100644
--- a/cloud/infra/terraform/relay-gce-startup.sh.tftpl
+++ b/cloud/infra/terraform/relay-gce-startup.sh.tftpl
@@ -109,6 +109,9 @@ docker run --detach \
--user 0:0 \
--volume "$${cloudsql_dir}:/cloudsql" \
'${cloud_sql_proxy_image}' \
+%{ if cloud_sql_private_ip ~}
+ --private-ip \
+%{ endif ~}
--unix-socket=/cloudsql \
'${cloud_sql_connection_name}'
diff --git a/cloud/infra/terraform/relay-observability.tf b/cloud/infra/terraform/relay-observability.tf
index a8fa4276eb2..964fc1060e6 100644
--- a/cloud/infra/terraform/relay-observability.tf
+++ b/cloud/infra/terraform/relay-observability.tf
@@ -37,6 +37,16 @@ locals {
description = "Relay PostgreSQL transactions that exhausted bounded retry."
filter = "((resource.type=\"cloud_run_revision\" AND (${local.relay_service_log_filter})) OR resource.type=\"gce_instance\") AND jsonPayload.event=\"orca_relay_postgres_transaction_exhausted\""
}
+ cell_process_exit = {
+ # The docker event stream is the only per-exit line: the relay's own crash footer only
+ # appears for unhandled rejections, and `container start` also counts healthy first boots.
+ description = "Relay cell container exits, one Docker `container die` event per process exit."
+ filter = "resource.type=\"gce_instance\" AND logName=\"projects/${var.project_id}/logs/cos_system\" AND jsonPayload.SYSLOG_IDENTIFIER=\"docker\" AND jsonPayload.MESSAGE:\"container die\" AND jsonPayload.MESSAGE:\"name=orca-relay)\""
+ }
+ cloud_sql_wal_checkpoint = {
+ description = "Cloud SQL checkpoints triggered by WAL volume instead of the timed schedule; a sustained run is the fsync loop that stalled every relay process at once on 2026-09-04."
+ filter = "resource.type=\"cloudsql_database\" AND resource.labels.database_id=\"${var.project_id}:${local.relay_database_instance_name}\" AND textPayload:\"checkpoint starting: wal\""
+ }
}
relay_runtime_metrics = {
@@ -523,3 +533,280 @@ resource "google_monitoring_alert_policy" "relay_cloud_sql_backends" {
mime_type = "text/markdown"
}
}
+
+resource "google_monitoring_alert_policy" "relay_cloud_sql_checkpoint_loop" {
+ project = var.project_id
+ display_name = "Orca Relay: Cloud SQL checkpoint loop"
+ combiner = "OR"
+ enabled = true
+ notification_channels = var.relay_alert_notification_channels
+
+ conditions {
+ display_name = "WAL-triggered checkpoints above 3 in 5 minutes"
+
+ condition_threshold {
+ filter = "resource.type=\"cloudsql_database\" AND metric.type=\"logging.googleapis.com/user/orca_relay_cloud_sql_wal_checkpoint\""
+ comparison = "COMPARISON_GT"
+ threshold_value = 3
+ duration = "300s"
+
+ aggregations {
+ alignment_period = "300s"
+ per_series_aligner = "ALIGN_SUM"
+ cross_series_reducer = "REDUCE_SUM"
+ }
+
+ trigger {
+ count = 1
+ }
+ }
+ }
+
+ documentation {
+ content = "Healthy operation is one timed checkpoint every 5 minutes. Repeated `checkpoint starting: wal` lines mean WAL is outrunning `max_wal_size` and every checkpoint fsync stalls all relay SQL for seconds. Check `checkpoint complete` sync= times and disk write throughput against the PD-SSD ceiling; the fix is disk size and `max_wal_size` in the Terraform root that owns the instance (orca-cloud `infra/terraform-foundation`)."
+ mime_type = "text/markdown"
+ }
+
+ depends_on = [google_logging_metric.relay_incident]
+}
+
+resource "google_monitoring_alert_policy" "relay_cloud_sql_disk" {
+ project = var.project_id
+ display_name = "Orca Relay: Cloud SQL disk utilization"
+ combiner = "OR"
+ enabled = true
+ notification_channels = var.relay_alert_notification_channels
+
+ conditions {
+ display_name = "Cloud SQL disk above 70%"
+
+ condition_threshold {
+ filter = "resource.type=\"cloudsql_database\" AND resource.label.\"database_id\"=\"${var.project_id}:${local.relay_database_instance_name}\" AND metric.type=\"cloudsql.googleapis.com/database/disk/utilization\""
+ comparison = "COMPARISON_GT"
+ threshold_value = 0.7
+ duration = "600s"
+
+ aggregations {
+ alignment_period = "300s"
+ per_series_aligner = "ALIGN_MAX"
+ }
+
+ trigger {
+ count = 1
+ }
+ }
+ }
+
+ documentation {
+ content = "The shared auth/relay Cloud SQL disk is filling. `refresh_tokens` is the largest table and grows without pruning; grow the disk (IOPS scale with size) before it reaches the WAL checkpoint loop, and prune revoked token rows."
+ mime_type = "text/markdown"
+ }
+}
+
+resource "google_monitoring_alert_policy" "relay_cloud_nat_port_drops" {
+ count = local.relay_gce_configured ? 1 : 0
+
+ project = var.project_id
+ display_name = "Orca Relay: Cloud NAT port exhaustion"
+ combiner = "OR"
+ enabled = true
+ notification_channels = var.relay_alert_notification_channels
+
+ conditions {
+ display_name = "NAT packets dropped for lack of ports"
+
+ condition_threshold {
+ filter = "resource.type=\"nat_gateway\" AND resource.label.\"gateway_name\"=monitoring.regex.full_match(\"${local.relay_gce_name}(-.*)?\") AND metric.type=\"router.googleapis.com/nat/dropped_sent_packets_count\" AND metric.label.\"reason\"=\"OUT_OF_RESOURCES\""
+ comparison = "COMPARISON_GT"
+ threshold_value = 0
+ duration = "120s"
+
+ aggregations {
+ alignment_period = "60s"
+ per_series_aligner = "ALIGN_SUM"
+ cross_series_reducer = "REDUCE_SUM"
+ group_by_fields = ["resource.label.\"gateway_name\""]
+ }
+
+ trigger {
+ count = 1
+ }
+ }
+ }
+
+ documentation {
+ content = "Relay cells reach Cloud SQL's public IP through this NAT. Port exhaustion makes every cell's Cloud SQL Auth Proxy dial time out at once, which reads as a fleet-wide SQL stall with a healthy database. Check `nat/port_usage` per VM and raise `max_ports_per_vm` in `relay-gce-foundation.tf`, or move the database to a private IP."
+ mime_type = "text/markdown"
+ }
+}
+
+resource "google_monitoring_alert_policy" "relay_cell_process_exit" {
+ project = var.project_id
+ display_name = "Orca Relay: cell process exits"
+ combiner = "OR"
+ enabled = true
+ notification_channels = var.relay_alert_notification_channels
+
+ conditions {
+ display_name = "Cell container exits above 3 in 15 minutes"
+
+ condition_threshold {
+ filter = "resource.type=\"gce_instance\" AND metric.type=\"logging.googleapis.com/user/orca_relay_cell_process_exit\""
+ comparison = "COMPARISON_GT"
+ threshold_value = 3
+ duration = "0s"
+
+ aggregations {
+ alignment_period = "900s"
+ per_series_aligner = "ALIGN_SUM"
+ cross_series_reducer = "REDUCE_SUM"
+ group_by_fields = ["resource.label.\"instance_id\""]
+ }
+
+ trigger {
+ count = 1
+ }
+ }
+ }
+
+ documentation {
+ content = "A Relay GCE cell restarted its container more than three times in 15 minutes. Each exit drops every host and phone on that cell, and 201 exits went unpaged over 48 h on 2026-09-04. The instance hostname is `relay--`; read `jsonPayload.MESSAGE` on `cos_system` for the exit code and the container's own stderr for the stack before blaming MIG autoheal or load. A same-capacity roll is the remedy when the running image is behind."
+ mime_type = "text/markdown"
+ }
+
+ depends_on = [google_logging_metric.relay_incident]
+}
+
+# Why: the four signals that had to be assembled by hand during the 2026-09-04 incident.
+resource "google_monitoring_dashboard" "relay_incident" {
+ project = var.project_id
+
+ dashboard_json = jsonencode({
+ displayName = "Orca Relay: incident overview"
+ mosaicLayout = {
+ columns = 12
+ tiles = [
+ {
+ xPos = 0
+ yPos = 0
+ width = 3
+ height = 4
+ widget = {
+ title = "Cloud SQL WAL checkpoints"
+ xyChart = {
+ dataSets = [{
+ plotType = "LINE"
+ targetAxis = "Y1"
+ timeSeriesQuery = {
+ timeSeriesFilter = {
+ filter = "metric.type=\"logging.googleapis.com/user/orca_relay_cloud_sql_wal_checkpoint\" AND resource.type=\"cloudsql_database\""
+ aggregation = {
+ alignmentPeriod = "300s"
+ perSeriesAligner = "ALIGN_SUM"
+ crossSeriesReducer = "REDUCE_SUM"
+ }
+ }
+ }
+ }]
+ yAxis = {
+ label = "checkpoints"
+ scale = "LINEAR"
+ }
+ }
+ }
+ },
+ {
+ xPos = 3
+ yPos = 0
+ width = 3
+ height = 4
+ widget = {
+ title = "Cloud NAT dropped packets"
+ xyChart = {
+ dataSets = [{
+ plotType = "LINE"
+ targetAxis = "Y1"
+ timeSeriesQuery = {
+ timeSeriesFilter = {
+ filter = "metric.type=\"router.googleapis.com/nat/dropped_sent_packets_count\" AND resource.type=\"nat_gateway\" AND resource.label.\"gateway_name\"=monitoring.regex.full_match(\"${local.relay_gce_name}(-.*)?\")"
+ aggregation = {
+ alignmentPeriod = "60s"
+ perSeriesAligner = "ALIGN_SUM"
+ crossSeriesReducer = "REDUCE_SUM"
+ groupByFields = ["resource.label.\"gateway_name\"", "metric.label.\"reason\""]
+ }
+ }
+ }
+ }]
+ yAxis = {
+ label = "packets"
+ scale = "LINEAR"
+ }
+ }
+ }
+ },
+ {
+ xPos = 6
+ yPos = 0
+ width = 3
+ height = 4
+ widget = {
+ title = "Auth refresh 401s"
+ xyChart = {
+ dataSets = [{
+ plotType = "LINE"
+ targetAxis = "Y1"
+ timeSeriesQuery = {
+ timeSeriesFilter = {
+ filter = "metric.type=\"logging.googleapis.com/user/orca_auth_refresh_401\""
+ aggregation = {
+ alignmentPeriod = "300s"
+ perSeriesAligner = "ALIGN_SUM"
+ crossSeriesReducer = "REDUCE_SUM"
+ }
+ }
+ }
+ }]
+ yAxis = {
+ label = "rejections"
+ scale = "LINEAR"
+ }
+ }
+ }
+ },
+ {
+ xPos = 9
+ yPos = 0
+ width = 3
+ height = 4
+ widget = {
+ title = "Standing desktop controls (fleet sum)"
+ xyChart = {
+ dataSets = [{
+ plotType = "LINE"
+ targetAxis = "Y1"
+ timeSeriesQuery = {
+ timeSeriesFilter = {
+ # ALIGN_MEAN, not ALIGN_SUM: each process reports its standing control count once per interval.
+ filter = "metric.type=\"logging.googleapis.com/user/orca_relay_controls\""
+ aggregation = {
+ alignmentPeriod = "300s"
+ perSeriesAligner = "ALIGN_MEAN"
+ crossSeriesReducer = "REDUCE_SUM"
+ }
+ }
+ }
+ }]
+ yAxis = {
+ label = "controls"
+ scale = "LINEAR"
+ }
+ }
+ }
+ }
+ ]
+ }
+ })
+
+ depends_on = [google_logging_metric.relay_incident, google_logging_metric.relay_snapshot]
+}
diff --git a/cloud/infra/terraform/variables.tf b/cloud/infra/terraform/variables.tf
index 68f6c555bd3..91f67e8ebe0 100644
--- a/cloud/infra/terraform/variables.tf
+++ b/cloud/infra/terraform/variables.tf
@@ -468,6 +468,12 @@ variable "relay_gce_fenced_cells" {
default = []
}
+variable "relay_cloud_sql_private_ip" {
+ type = bool
+ description = "Dial Cloud SQL over its private IP inside this VPC instead of its public IP through Cloud NAT. Requires the foundation root's private services access peering to be applied first; a cell that cannot reach the private IP never becomes ready."
+ default = false
+}
+
variable "relay_gce_cloud_sql_proxy_image" {
type = string
description = "Digest-pinned Cloud SQL Auth Proxy image used by private relay workers."
diff --git a/cloud/packages/relay-contract/src/host-close-reason.ts b/cloud/packages/relay-contract/src/host-close-reason.ts
new file mode 100644
index 00000000000..3a5abde3f00
--- /dev/null
+++ b/cloud/packages/relay-contract/src/host-close-reason.ts
@@ -0,0 +1,18 @@
+// Mirror of src/shared/relay-host-close-reason.ts in the Orca app repo half.
+// A host control socket may close with one of these as its WebSocket close
+// reason; the cell records it so a later phone rejection can name the cause.
+// Anything else (including the empty reason of an abrupt 1006) means "unknown",
+// which is what every peer that predates this file sends.
+export const RELAY_HOST_CLOSE_REASON = {
+ SIGNED_OUT: 'signed-out'
+} as const
+
+export type RelayHostCloseReason =
+ (typeof RELAY_HOST_CLOSE_REASON)[keyof typeof RELAY_HOST_CLOSE_REASON]
+
+const REASONS: readonly string[] = Object.values(RELAY_HOST_CLOSE_REASON)
+
+export function relayHostCloseReasonFrom(value: unknown): RelayHostCloseReason | null {
+ const text = typeof value === 'string' ? value : (value?.toString() ?? '')
+ return REASONS.includes(text) ? (text as RelayHostCloseReason) : null
+}
diff --git a/cloud/packages/relay-contract/src/index.ts b/cloud/packages/relay-contract/src/index.ts
index 2a7d7d0feda..aab3b53b5f3 100644
--- a/cloud/packages/relay-contract/src/index.ts
+++ b/cloud/packages/relay-contract/src/index.ts
@@ -5,6 +5,7 @@ export * from './control-messages.js'
export * from './control-continuity.js'
export * from './credential-messages.js'
export * from './director-messages.js'
+export * from './host-close-reason.js'
export * from './host-proof-transcript.js'
export * from './persistence-invariants.js'
export * from './protocol-limits.js'
diff --git a/config/patches/node-pty@1.1.0.patch b/config/patches/node-pty@1.1.0.patch
index ff474f7d95e..8f5045b932a 100644
--- a/config/patches/node-pty@1.1.0.patch
+++ b/config/patches/node-pty@1.1.0.patch
@@ -603,7 +603,7 @@ index 7b4b9e1f990fbf95b51528bb56dc9717f5b87532..2ae787c5bd4f3eba470584dc658a01a5
}
#endif
diff --git a/src/win/conpty.cc b/src/win/conpty.cc
-index 7b286d3d644c26141df516929703aa6e129df4b2..ec6bf3932c65b89c013ff133dc6bf46a6a4082ce 100644
+index 7b286d3d644c26141df516929703aa6e129df4b2..4aed260dd68e6a171dcfd349e9a7c5c97209248e 100644
--- a/src/win/conpty.cc
+++ b/src/win/conpty.cc
@@ -18,6 +18,7 @@
@@ -614,7 +614,7 @@ index 7b286d3d644c26141df516929703aa6e129df4b2..ec6bf3932c65b89c013ff133dc6bf46a
#include
#include
#include
-@@ -44,12 +45,29 @@ struct pty_baton {
+@@ -44,12 +45,39 @@ struct pty_baton {
HANDLE hOut;
HPCON hpc;
@@ -630,22 +630,32 @@ index 7b286d3d644c26141df516929703aa6e129df4b2..ec6bf3932c65b89c013ff133dc6bf46a
+ // refused to create or assign one (an outer job without breakaway rights),
+ // in which case callers fall back to their pre-job behaviour.
+ HANDLE hJob = nullptr;
++
++ // Orca: teardown needs BOTH the shell's death and an explicit kill() before
++ // the baton can be freed, so each side records that it has run. Whichever
++ // arrives second frees it. Freeing on the shell's death alone -- what this
++ // file did before -- destroyed the only record of `hpc` while
++ // ClosePseudoConsole was still owed, which is why a self-exiting shell
++ // leaked its pseudoconsole and the console host it reaps (#18601 / F24).
++ bool shellExited = false;
++ bool consoleClosed = false;
pty_baton(int _id, HANDLE _hIn, HANDLE _hOut, HPCON _hpc) : id(_id), hIn(_hIn), hOut(_hOut), hpc(_hpc) {};
};
static std::vector> ptyHandles;
-+// Orca: guards the job accessors below against the exit watcher thread. It does
-+// NOT make the whole table safe -- PtyResize/PtyClear/PtyKill read it unlocked,
-+// as they always have -- but it closes the window this patch opened, where the
-+// watcher can close hShell/hJob and free the baton between a lookup and its use.
++// Orca: guards the job accessors below, and PtyKill, against the exit watcher
++// thread. It does NOT make the whole table safe -- PtyResize and PtyClear still
++// read it unlocked, as they always have -- but it closes the window this patch
++// opened, where the watcher can close hShell/hJob and free the baton between a
++// lookup and its use.
+// Handle VALUES are recycled aggressively, so an unguarded read could pass the
+// shell-pid check against an unrelated process and terminate the wrong job.
+static std::mutex ptyJobMutex;
static volatile LONG ptyCounter;
static pty_baton* get_pty_baton(int id) {
-@@ -102,8 +120,27 @@ void SetupExitCallback(Napi::Env env, Napi::Function cb, pty_baton* baton) {
+@@ -102,8 +130,31 @@ void SetupExitCallback(Napi::Env env, Napi::Function cb, pty_baton* baton) {
// Get process exit code.
GetExitCodeProcess(baton->hShell, (LPDWORD)(&exit_event->exit_code));
// Clean up handles
@@ -665,9 +675,13 @@ index 7b286d3d644c26141df516929703aa6e129df4b2..ec6bf3932c65b89c013ff133dc6bf46a
+ // Why inside the lock: erasing frees the baton the job accessors hold a
+ // pointer to. Note remove_pty_baton must not be an assert() argument --
+ // NDEBUG would compile the call away and leak every baton.
-+ const bool removed = remove_pty_baton(baton->id);
-+ assert(removed);
-+ (void)removed;
++ baton->shellExited = true;
++ if (baton->consoleClosed) {
++ const bool removed = remove_pty_baton(baton->id);
++ assert(removed);
++ (void)removed;
++ }
++ // Else PtyKill has not run yet and still owns hpc. It frees the baton.
+ }
+ // Why the lock ends here: BlockingCall below waits on the JS thread, and the
+ // JS thread can be waiting on ptyJobMutex inside PtyTerminateJob. Holding
@@ -675,7 +689,7 @@ index 7b286d3d644c26141df516929703aa6e129df4b2..ec6bf3932c65b89c013ff133dc6bf46a
auto status = tsfn.BlockingCall(exit_event, callback); // In main thread
switch (status) {
-@@ -409,6 +446,15 @@ static Napi::Value PtyConnect(const Napi::CallbackInfo& info) {
+@@ -409,6 +460,15 @@ static Napi::Value PtyConnect(const Napi::CallbackInfo& info) {
throw errorWithCode(info, "UpdateProcThreadAttribute failed");
}
@@ -691,7 +705,7 @@ index 7b286d3d644c26141df516929703aa6e129df4b2..ec6bf3932c65b89c013ff133dc6bf46a
PROCESS_INFORMATION piClient{};
fSuccess = !!CreateProcessW(
nullptr,
-@@ -416,7 +462,10 @@ static Napi::Value PtyConnect(const Napi::CallbackInfo& info) {
+@@ -416,7 +476,10 @@ static Napi::Value PtyConnect(const Napi::CallbackInfo& info) {
nullptr, // lpProcessAttributes
nullptr, // lpThreadAttributes
false, // bInheritHandles VERY IMPORTANT that this is false
@@ -703,7 +717,7 @@ index 7b286d3d644c26141df516929703aa6e129df4b2..ec6bf3932c65b89c013ff133dc6bf46a
envArg, // lpEnvironment
mutableCwd.get(), // lpCurrentDirectory
&siEx.StartupInfo, // lpStartupInfo
-@@ -426,8 +475,47 @@ static Napi::Value PtyConnect(const Napi::CallbackInfo& info) {
+@@ -426,8 +489,47 @@ static Napi::Value PtyConnect(const Napi::CallbackInfo& info) {
throw errorWithCode(info, "Cannot create process");
}
@@ -753,7 +767,7 @@ index 7b286d3d644c26141df516929703aa6e129df4b2..ec6bf3932c65b89c013ff133dc6bf46a
if (useConptyDll && fLoadedDll)
{
PFNRELEASEPSEUDOCONSOLE const pfnReleasePseudoConsole = (PFNRELEASEPSEUDOCONSOLE)GetProcAddress(
-@@ -440,6 +528,8 @@ static Napi::Value PtyConnect(const Napi::CallbackInfo& info) {
+@@ -440,6 +542,8 @@ static Napi::Value PtyConnect(const Napi::CallbackInfo& info) {
// Update handle
handle->hShell = piClient.hProcess;
@@ -762,7 +776,91 @@ index 7b286d3d644c26141df516929703aa6e129df4b2..ec6bf3932c65b89c013ff133dc6bf46a
// Close the thread handle to avoid resource leak
CloseHandle(piClient.hThread);
-@@ -567,6 +657,143 @@ static Napi::Value PtyKill(const Napi::CallbackInfo& info) {
+@@ -544,29 +648,215 @@ static Napi::Value PtyKill(const Napi::CallbackInfo& info) {
+ int id = info[0].As().Int32Value();
+ const bool useConptyDll = info[1].As().Value();
+
+- const pty_baton* handle = get_pty_baton(id);
++ // Orca: resolve the DLL BEFORE touching any baton state, for the same reason
++ // PtyConnect does it before creating anything. LoadConptyDll throws when
++ // conpty.dll is missing, and a throw after consoleClosed was set would strand
++ // the pseudoconsole permanently: the retry would find the work already
++ // claimed and do nothing. Only the useConptyDll path can throw here; the
++ // other returns kernel32.
++ HANDLE hLibrary = LoadConptyDll(info, useConptyDll);
++ PFNCLOSEPSEUDOCONSOLE pfnClosePseudoConsole = nullptr;
++ if (hLibrary != nullptr) {
++ pfnClosePseudoConsole = (PFNCLOSEPSEUDOCONSOLE)GetProcAddress(
++ (HMODULE)hLibrary,
++ useConptyDll ? "ConptyClosePseudoConsole" : "ClosePseudoConsole");
++ }
+
+- if (handle != nullptr) {
+- HANDLE hLibrary = LoadConptyDll(info, useConptyDll);
+- bool fLoadedDll = hLibrary != nullptr;
+- if (fLoadedDll)
+- {
+- PFNCLOSEPSEUDOCONSOLE const pfnClosePseudoConsole = (PFNCLOSEPSEUDOCONSOLE)GetProcAddress(
+- (HMODULE)hLibrary,
+- useConptyDll ? "ConptyClosePseudoConsole" : "ClosePseudoConsole");
+- if (pfnClosePseudoConsole)
+- {
+- pfnClosePseudoConsole(handle->hpc);
++ // Orca: the baton now outlives the shell, so this runs on a self-exited pty
++ // too -- that is the whole point. Take what we need under the lock: the
++ // watcher thread nulls hShell the moment the shell dies, and TerminateProcess
++ // on a handle it just closed is an invalid-handle operation. Duplicating
++ // rather than reordering keeps upstream's close-then-terminate sequence.
++ HPCON hpc = nullptr;
++ HANDLE hShellDup = nullptr;
++ bool owed = false;
++ {
++ std::lock_guard guard(ptyJobMutex);
++ pty_baton* handle = get_pty_baton(id);
++ // Why the consoleClosed check: a second kill() would otherwise close the
++ // same pseudoconsole twice. Upstream relied on the baton being gone.
++ if (handle != nullptr && !handle->consoleClosed) {
++ hpc = handle->hpc;
++ owed = true;
++ handle->consoleClosed = true;
++ // Null hShell means a self-exited pty, where there is nothing to kill.
++ if (useConptyDll && handle->hShell != nullptr) {
++ if (!DuplicateHandle(GetCurrentProcess(), handle->hShell, GetCurrentProcess(),
++ &hShellDup, 0, FALSE, DUPLICATE_SAME_ACCESS)) {
++ // Why terminate here instead of skipping: a failed duplication leaves
++ // hShellDup null, which is indistinguishable from the self-exit case,
++ // and skipping would leave the shell RUNNING after its pane closed --
++ // a worse outcome than the leak this all exists to fix. hShell is
++ // valid under this lock and TerminateProcess does not block, so the
++ // only cost is that this rare path kills before the console closes.
++ hShellDup = nullptr;
++ TerminateProcess(handle->hShell, 1);
++ }
++ }
++ if (handle->shellExited) {
++ const bool removed = remove_pty_baton(id);
++ assert(removed);
++ (void)removed;
+ }
++ // Else the shell is still running and the watcher frees the baton.
+ }
+- if (useConptyDll) {
+- TerminateProcess(handle->hShell, 1);
++ }
++
++ // Why outside the lock: ClosePseudoConsole blocks until the conout side has
++ // drained, and the watcher must be able to take the lock while it does.
++ if (owed) {
++ if (pfnClosePseudoConsole)
++ {
++ pfnClosePseudoConsole(hpc);
++ }
++ if (hShellDup != nullptr) {
++ TerminateProcess(hShellDup, 1);
++ CloseHandle(hShellDup);
+ }
+ }
+
return env.Undefined();
}
@@ -808,9 +906,11 @@ index 7b286d3d644c26141df516929703aa6e129df4b2..ec6bf3932c65b89c013ff133dc6bf46a
+ * Orca: the pids still alive in this pty's tree, straight from the kernel.
+ *
+ * Descendant liveness for a tree that is still tracked, including children that
-+ * detached from the console. Once the shell exits the baton is gone, so this
-+ * returns null rather than an empty list -- null means "no answer", never
-+ * "they died". Also returns null when no job was assigned.
++ * detached from the console. Once the shell exits the watcher nulls hJob, which
++ * ownsShell rejects, so this returns null rather than an empty list -- null
++ * means "no answer", never "they died". (The baton itself now outlives the
++ * shell, until kill() runs; hJob is what makes the answer null.) Also returns
++ * null when no job was assigned.
+ *
+ * Does not include the ConPTY console host: CreatePseudoConsole spawns it
+ * before this job exists, so it is not a member and ClosePseudoConsole is what
@@ -906,7 +1006,7 @@ index 7b286d3d644c26141df516929703aa6e129df4b2..ec6bf3932c65b89c013ff133dc6bf46a
/**
* Init
*/
-@@ -577,6 +804,9 @@ Napi::Object init(Napi::Env env, Napi::Object exports) {
+@@ -577,6 +867,9 @@ Napi::Object init(Napi::Env env, Napi::Object exports) {
exports.Set("resize", Napi::Function::New(env, PtyResize));
exports.Set("clear", Napi::Function::New(env, PtyClear));
exports.Set("kill", Napi::Function::New(env, PtyKill));
@@ -917,7 +1017,7 @@ index 7b286d3d644c26141df516929703aa6e129df4b2..ec6bf3932c65b89c013ff133dc6bf46a
};
diff --git a/lib/windowsPtyAgent.js b/lib/windowsPtyAgent.js
-index a358ffb..fb3a96f 100644
+index a358ffb177357e177661033c1b092f9c9d0e5f5a..26c2a4c58799ce649f5113131e4c52f7ed2d87ad 100644
--- a/lib/windowsPtyAgent.js
+++ b/lib/windowsPtyAgent.js
@@ -136,6 +136,9 @@ var WindowsPtyAgent = /** @class */ (function () {
@@ -930,6 +1030,20 @@ index a358ffb..fb3a96f 100644
this._outSocket.readable = false;
this._getConsoleProcessList().then(function (consoleProcessList) {
consoleProcessList.forEach(function (pid) {
+@@ -154,9 +157,10 @@ var WindowsPtyAgent = /** @class */ (function () {
+ // Close the input write handle to signal the end of session.
+ this._inSocket.destroy();
+ this._ptyNative.kill(this._pty, this._useConptyDll);
+- this._outSocket.on('data', function () {
+- _this._conoutSocketWorker.dispose();
+- });
++ // Orca: dispose unconditionally, as the non-DLL branch above does.
++ // Waiting for another 'data' event leaks the conout worker on every
++ // self-exiting shell, because no more data ever arrives (F24).
++ this._conoutSocketWorker.dispose();
+ }
+ }
+ else {
diff --git a/lib/windowsTerminal.js b/lib/windowsTerminal.js
index 3c38f89..e20b3e6 100644
--- a/lib/windowsTerminal.js
@@ -1015,7 +1129,7 @@ index 3c38f89..e20b3e6 100644
\ No newline at end of file
+//# sourceMappingURL=windowsTerminal.js.map
diff --git a/src/windowsPtyAgent.ts b/src/windowsPtyAgent.ts
-index d705444..ce611b8 100644
+index d7054449516f0c9a62af351c2caa17331206d530..0c28a32e2e1db2b3f208ddde8443cd4e67bb1ad6 100644
--- a/src/windowsPtyAgent.ts
+++ b/src/windowsPtyAgent.ts
@@ -143,6 +143,9 @@ export class WindowsPtyAgent {
@@ -1028,6 +1142,20 @@ index d705444..ce611b8 100644
this._outSocket.readable = false;
this._getConsoleProcessList().then(consoleProcessList => {
consoleProcessList.forEach((pid: number) => {
+@@ -159,9 +162,10 @@ export class WindowsPtyAgent {
+ // Close the input write handle to signal the end of session.
+ this._inSocket.destroy();
+ (this._ptyNative as IConptyNative).kill(this._pty, this._useConptyDll);
+- this._outSocket.on('data', () => {
+- this._conoutSocketWorker.dispose();
+- });
++ // Orca: dispose unconditionally, as the non-DLL branch above does.
++ // Waiting for another 'data' event leaks the conout worker on every
++ // self-exiting shell, because no more data ever arrives (F24).
++ this._conoutSocketWorker.dispose();
+ }
+ } else {
+ // Because pty.kill closes the handle, it will kill most processes by itself.
diff --git a/src/windowsTerminal.ts b/src/windowsTerminal.ts
index 13f6c6d..eda63c8 100644
--- a/src/windowsTerminal.ts
diff --git a/config/relay-assets/node-pty-1.1.0-windows-pty-teardown-patch.cjs b/config/relay-assets/node-pty-1.1.0-windows-pty-teardown-patch.cjs
new file mode 100644
index 00000000000..dd26784ee46
--- /dev/null
+++ b/config/relay-assets/node-pty-1.1.0-windows-pty-teardown-patch.cjs
@@ -0,0 +1,158 @@
+const { createHash } = require('node:crypto')
+const { readFileSync, renameSync, rmSync, writeFileSync } = require('node:fs')
+const { join, resolve } = require('node:path')
+
+/**
+ * Release the ConPTY teardown handles a relay's npm-installed node-pty never releases.
+ *
+ * Two files, and the ORDER of one of the edits is the whole fix.
+ *
+ * `windowsPtyAgent.js` -- `kill()` flips `readable` on both sockets and destroys neither.
+ * `_cleanUpProcess` destroys `_outSocket`, so the conout handle comes back; nothing ever destroys
+ * `_inSocket`, and it wraps a real Windows named-pipe handle from `fs.openSync(term.conin, 'w')`.
+ * Every terminal leaks one File handle for the life of the host process.
+ *
+ * The obvious fix -- and the one the desktop patch ships -- releases it at the TOP of the branch,
+ * before `_getConsoleProcessList()` forks and before the native kill. That is measurably worse than
+ * leaving the leak alone: teardown aborts partway, the forked console-list agent is never reaped,
+ * and both pipe handles stay alive instead of one. This asset releases it at the END of the branch
+ * instead, after the fork and the kill have already happened.
+ *
+ * Measured on a Windows SSH host, 20 spawn/kill cycles, handles bucketed by NT object type
+ * (identical numbers standalone and through a real relay):
+ *
+ * published node-pty File +1/terminal, Process flat
+ * desktop patch placement File +2/terminal, Process +1/terminal <-- 3x WORSE
+ * released last (here) File flat, Process flat
+ *
+ * `windowsTerminal.js` carries the desktop's error-listener hunks verbatim. The conin listener is
+ * what keeps a pipe error retiring one terminal instead of the host -- its own comment names the
+ * failure mode: "Without a listener, Node promotes errors such as write EAGAIN to uncaughtException".
+ * It is not what fixes the leak (adding it changed nothing on its own), but it is the guard that
+ * makes destroying conin safe at all.
+ *
+ * Why this ships as a relay asset rather than only in config/patches/node-pty@1.1.0.patch: pnpm
+ * patches do not cross the SSH boundary -- a relay host runs the tree `npm install` put there.
+ *
+ * DELIBERATE DIVERGENCE FROM THE DESKTOP: the desktop patch has the early placement and therefore
+ * the +2 File / +1 Process regression, measured against its exact installed tree. Correcting it
+ * there is a separate change with its own verification, so the two trees differ on this one hunk on
+ * purpose, and the test pins that so a future "sync the patches" does not copy the bug back.
+ *
+ * NOT ADDRESSED, AND A SEPARATE DEFECT THAT IS STILL OPEN: a terminal that exits on its own is
+ * still torn down through `kill()` -- both hosts call `destroy()` on natural exit and
+ * `WindowsTerminal.destroy()` is `kill()` -- but the shell is already gone by then, and the
+ * ordering this patch relies on does not hold. Measured over 20 self-exit cycles with that
+ * `destroy()` issued: published +3 File/+1 Process per terminal, desktop-patched +2/+1, this tree
+ * +2/+1. So this patch does not close it and the desktop patch does not either. It is reachable
+ * for every Windows user, local and relay, on every terminal closed by typing `exit`.
+ */
+
+const EXPECTED_NODE_PTY_VERSION = '1.1.0'
+
+/** Each entry is one published file, its patched form, and the edits between them. */
+const PATCH_TARGETS = [
+ {
+ relativePath: ['lib', 'windowsPtyAgent.js'],
+ originalSha256: '8636d16b38266112204061a22b135734177c242837982fd3a4055be726efa64a',
+ patchedSha256: '1e23ef480569e73706e3ab4f5482c7e553c76f51414ae8e7b0bdcc2fd75f7280',
+ replacements: [
+ [
+ ' this._ptyNative.kill(this._pty, this._useConptyDll);\n this._conoutSocketWorker.dispose();\n',
+ ' this._ptyNative.kill(this._pty, this._useConptyDll);\n this._conoutSocketWorker.dispose();\n // Orca: released AFTER the console-list fork and the native kill, not before them.\n // Destroying conin first aborts teardown partway -- measured on a Windows SSH relay\n // as +2 File and +1 Process handles per terminal, against +1 File unpatched.\n this._inSocket.destroy();\n'
+ ]
+ ]
+ },
+ {
+ relativePath: ['lib', 'windowsTerminal.js'],
+ originalSha256: 'c3a65716f53fed0135a8a633373d5f9c2ab092544d651f27ef0a67096dd3bcd9',
+ patchedSha256: '8247ecd69be8b18257050fb026b290024612c5ffc6d492ff1d46f81e613be2cf',
+ replacements: [
+ [
+ ' _this._agent = new windowsPtyAgent_1.WindowsPtyAgent(file, args, parsedEnv, cwd, _this._cols, _this._rows, false, opt.useConpty, opt.useConptyDll, opt.conptyInheritCursor);\n _this._socket = _this._agent.outSocket;\n // Not available until `ready` event emitted.\n _this._pid = _this._agent.innerPid;',
+ " _this._agent = new windowsPtyAgent_1.WindowsPtyAgent(file, args, parsedEnv, cwd, _this._cols, _this._rows, false, opt.useConpty, opt.useConptyDll, opt.conptyInheritCursor);\n _this._socket = _this._agent.outSocket;\n // Attach before readiness so a broken ConPTY output pipe cannot be unhandled.\n _this._socket.on('error', function (err) {\n var code = err && err.code;\n // PTY output can report EPIPE before `_close()` wins the race.\n _this._close();\n if (code === 'EPIPE' || code === 'ERR_STREAM_PUSH_AFTER_EOF' || code === 'ERR_STREAM_DESTROYED') {\n return;\n }\n // EIO, happens when someone closes our child process: the only process\n // in the terminal.\n // node < 0.6.14: errno 5\n // node >= 0.6.14: read EIO\n if (typeof code === 'string') {\n if (~code.indexOf('errno 5') || ~code.indexOf('EIO'))\n return;\n }\n // Throw anything else.\n if (_this.listeners('error').length < 2) {\n throw err;\n }\n });\n // Not available until `ready` event emitted.\n _this._pid = _this._agent.innerPid;"
+ ],
+ [
+ " }\n });\n // Shutdown if `error` event is emitted.\n _this._socket.on('error', function (err) {\n // Close terminal session.\n _this._close();\n // EIO, happens when someone closes our child process: the only process\n // in the terminal.\n // node < 0.6.14: errno 5\n // node >= 0.6.14: read EIO\n if (err.code) {\n if (~err.code.indexOf('errno 5') || ~err.code.indexOf('EIO'))\n return;\n }\n // Throw anything else.\n if (_this.listeners('error').length < 2) {\n throw err;\n }\n });\n // Cleanup after the socket is closed.\n _this._socket.on('close', function () {",
+ " }\n });\n // Cleanup after the socket is closed.\n _this._socket.on('close', function () {"
+ ],
+ [
+ ' _this._readable = true;\n _this._writable = true;\n _this._forwardEvents();\n return _this;',
+ " _this._readable = true;\n _this._writable = true;\n // A ConPTY input-pipe error must retire only this terminal. Without a listener, Node promotes\n // errors such as write EAGAIN to uncaughtException and kills every PTY in the daemon.\n _this._agent.inSocket.on('error', function () {\n if (!_this._writable) {\n return;\n }\n _this._close();\n try {\n _this._agent.kill();\n }\n catch (_a) {\n // The failing pipe may have raced process exit; the terminal is already unwritable.\n }\n });\n _this._forwardEvents();\n return _this;"
+ ],
+ [
+ 'exports.WindowsTerminal = WindowsTerminal;\n//# sourceMappingURL=windowsTerminal.js.map',
+ 'exports.WindowsTerminal = WindowsTerminal;\n//# sourceMappingURL=windowsTerminal.js.map\n'
+ ]
+ ]
+ }
+]
+
+function inspectTarget(relayDir, target) {
+ const nodePtyDir = resolve(relayDir, 'node_modules', 'node-pty')
+ const packageJson = JSON.parse(readFileSync(join(nodePtyDir, 'package.json'), 'utf8'))
+ if (packageJson.version !== EXPECTED_NODE_PTY_VERSION) {
+ throw new Error(
+ `Refusing to patch node-pty ${packageJson.version}; expected ${EXPECTED_NODE_PTY_VERSION}`
+ )
+ }
+ const filePath = join(nodePtyDir, ...target.relativePath)
+ return { filePath, source: readFileSync(filePath, 'utf8') }
+}
+
+function assertPatchedNodePtyWindowsTeardown(relayDir = process.cwd()) {
+ for (const target of PATCH_TARGETS) {
+ const inspected = inspectTarget(relayDir, target)
+ if (sourceSha256(inspected.source) !== target.patchedSha256) {
+ throw new Error(
+ `node-pty ConPTY teardown release is not installed in ${target.relativePath.join('/')}`
+ )
+ }
+ }
+}
+
+function patchNodePtyWindowsTeardown(relayDir = process.cwd()) {
+ for (const target of PATCH_TARGETS) {
+ const inspected = inspectTarget(relayDir, target)
+ const sourceHash = sourceSha256(inspected.source)
+ if (sourceHash === target.patchedSha256) {
+ continue
+ }
+ if (sourceHash !== target.originalSha256) {
+ throw new Error(
+ `Refusing to patch unexpected node-pty source in ${target.relativePath.join('/')}`
+ )
+ }
+ let patchedSource = inspected.source
+ for (const [from, to] of target.replacements) {
+ // Why the count check: an anchor that matched twice would patch the wrong site silently, and
+ // the hash below would then reject a tree this script had already rewritten.
+ if (patchedSource.split(from).length - 1 !== 1) {
+ throw new Error(`Refusing to patch ${target.relativePath.join('/')}; anchor is not unique`)
+ }
+ patchedSource = patchedSource.replace(from, to)
+ }
+ const temporaryPath = `${inspected.filePath}.orca-patch-${process.pid}`
+ // Why: a terminated remote install must leave either known source version recoverable on reconnect.
+ try {
+ writeFileSync(temporaryPath, patchedSource)
+ renameSync(temporaryPath, inspected.filePath)
+ } finally {
+ rmSync(temporaryPath, { force: true })
+ }
+ }
+ assertPatchedNodePtyWindowsTeardown(relayDir)
+}
+
+function sourceSha256(source) {
+ return createHash('sha256').update(source).digest('hex')
+}
+
+if (require.main === module) {
+ patchNodePtyWindowsTeardown()
+}
+
+module.exports = {
+ assertPatchedNodePtyWindowsTeardown,
+ patchNodePtyWindowsTeardown
+}
diff --git a/config/scripts/build-relay.mjs b/config/scripts/build-relay.mjs
index 289c7a957bd..4d408712f97 100644
--- a/config/scripts/build-relay.mjs
+++ b/config/scripts/build-relay.mjs
@@ -57,6 +57,13 @@ const NODE_PTY_CONSOLE_LIST_PATCH_SOURCE = join(
'relay-assets',
NODE_PTY_CONSOLE_LIST_PATCH_FILENAME
)
+const NODE_PTY_WINDOWS_TEARDOWN_PATCH_FILENAME = 'node-pty-1.1.0-windows-pty-teardown-patch.cjs'
+const NODE_PTY_WINDOWS_TEARDOWN_PATCH_SOURCE = join(
+ ROOT,
+ 'config',
+ 'relay-assets',
+ NODE_PTY_WINDOWS_TEARDOWN_PATCH_FILENAME
+)
const NODE_PTY_MASTER_CLOEXEC_PATCH_FILENAME = 'node-pty-1.1.0-master-cloexec-patch.cjs'
const NODE_PTY_MASTER_CLOEXEC_PATCH_SOURCE = join(
ROOT,
@@ -132,6 +139,10 @@ for (const platform of RELAY_BUILD_PLATFORMS) {
NODE_PTY_CONSOLE_LIST_PATCH_SOURCE,
join(outDir, NODE_PTY_CONSOLE_LIST_PATCH_FILENAME)
)
+ copyFileSync(
+ NODE_PTY_WINDOWS_TEARDOWN_PATCH_SOURCE,
+ join(outDir, NODE_PTY_WINDOWS_TEARDOWN_PATCH_FILENAME)
+ )
}
copyFileSync(
NODE_PTY_MASTER_CLOEXEC_PATCH_SOURCE,
diff --git a/config/scripts/locale-ko-key-overrides.json b/config/scripts/locale-ko-key-overrides.json
index f368ecc3cbc..bf5f62d1fa5 100644
--- a/config/scripts/locale-ko-key-overrides.json
+++ b/config/scripts/locale-ko-key-overrides.json
@@ -492,7 +492,7 @@
"ko": "agent CLI를 찾지 못했습니다. 하나를 설치하거나 설정에서 기본 agent를 선택하세요."
},
"auto.components.Terminal.7958465754": {
- "ko": "실행 중인 프로세스가 있는 로컬 terminals이 있습니다. 그래도 창을 닫으시겠습니까?"
+ "ko": "실행 중인 프로세스가 있는 terminals이 있습니다. 그래도 창을 닫으시겠습니까?"
},
"auto.components.Terminal.cdc9ac4b2d": {
"ko": "편집기"
diff --git a/config/scripts/node-pty-windows-pty-teardown-patch.test.mjs b/config/scripts/node-pty-windows-pty-teardown-patch.test.mjs
new file mode 100644
index 00000000000..64fb1b056b8
--- /dev/null
+++ b/config/scripts/node-pty-windows-pty-teardown-patch.test.mjs
@@ -0,0 +1,213 @@
+// The relay's copy of the ConPTY teardown release, and the guard that keeps it in lockstep with the
+// desktop's own node-pty patch. pnpm patches do not cross the SSH boundary, so a relay runs the tree
+// `npm install` put there; the desktop had this fix and the relay did not, and every terminal on a
+// Windows SSH host leaked one File handle for the life of the relay process.
+//
+// The ORDER of the conin release is the fix. Releasing it at the top of the branch -- what the
+// desktop patch does -- was measured at 3x WORSE than shipping nothing (File +2/terminal and a new
+// Process +1/terminal); releasing it after the console-list fork and the native kill is flat.
+import { createRequire } from 'node:module'
+import { existsSync, mkdirSync, mkdtempSync, readFileSync, rmSync, writeFileSync } from 'node:fs'
+import { join, resolve } from 'node:path'
+import { afterEach, describe, expect, it } from 'vitest'
+
+const require = createRequire(import.meta.url)
+const {
+ assertPatchedNodePtyWindowsTeardown,
+ patchNodePtyWindowsTeardown
+} = require('../relay-assets/node-pty-1.1.0-windows-pty-teardown-patch.cjs')
+const projectDir = resolve(import.meta.dirname, '..', '..')
+const cleanupDirs = []
+
+const PATCHED_FILES = ['windowsPtyAgent.js', 'windowsTerminal.js']
+
+/** The hunks config/patches/node-pty@1.1.0.patch adds to the installed desktop tree. */
+const DESKTOP_HUNKS = {
+ 'windowsPtyAgent.js': [
+ [
+ [
+ ' this._inSocket.readable = false;',
+ ' // The non-DLL path previously only flipped `readable`, leaving the',
+ ' // conin PipeWrap alive until the host exited (#947).',
+ ' this._inSocket.destroy();',
+ ' this._outSocket.readable = false;',
+ ''
+ ].join('\n'),
+ [
+ ' this._inSocket.readable = false;',
+ ' this._outSocket.readable = false;',
+ ''
+ ].join('\n')
+ ],
+ // The useConptyDll branch, which only the DESKTOP runs -- the relay takes the
+ // non-DLL branch above, where the dispose is already unconditional. Listed here
+ // so un-applying still yields published; the relay asset needs no counterpart.
+ [
+ [
+ ' // Orca: dispose unconditionally, as the non-DLL branch above does.',
+ " // Waiting for another 'data' event leaks the conout worker on every",
+ ' // self-exiting shell, because no more data ever arrives (F24).',
+ ' this._conoutSocketWorker.dispose();',
+ ''
+ ].join('\n'),
+ [
+ " this._outSocket.on('data', function () {",
+ ' _this._conoutSocketWorker.dispose();',
+ ' });',
+ ''
+ ].join('\n')
+ ]
+ ],
+ 'windowsTerminal.js': [
+ [
+ ' // Attach before readiness so a broken ConPTY output pipe cannot be unhandled.',
+ null
+ ],
+ [' // A ConPTY input-pipe error must retire only this terminal.', null]
+ ]
+}
+
+function desktopPath(file) {
+ return join(projectDir, 'node_modules', 'node-pty', 'lib', file)
+}
+
+afterEach(() => {
+ for (const dir of cleanupDirs.splice(0)) {
+ rmSync(dir, { recursive: true, force: true })
+ }
+})
+
+describe('Windows SSH relay node-pty ConPTY teardown patch', () => {
+ // Why reconstruct rather than vendor upstream: the installed tree IS the published file plus the
+ // desktop's hunks, so un-applying them yields upstream exactly -- and pinning that against this
+ // asset's own hashes is what fails loudly if either side of the pair moves.
+ it('takes the desktop error listeners verbatim', () => {
+ const fixture = writeNodePtyFixture('1.1.0')
+ patchNodePtyWindowsTeardown(fixture.root)
+
+ expect(readFileSync(join(fixture.libDir, 'windowsTerminal.js'), 'utf8')).toBe(
+ readFileSync(desktopPath('windowsTerminal.js'), 'utf8')
+ )
+ })
+
+ // The one hunk that must NOT match the desktop, and the reason is measured, not stylistic:
+ // releasing conin before `_getConsoleProcessList()` forks aborts teardown partway.
+ it('releases conin after the console-list fork, not before it like the desktop patch', () => {
+ const fixture = writeNodePtyFixture('1.1.0')
+ patchNodePtyWindowsTeardown(fixture.root)
+ const patched = readFileSync(join(fixture.libDir, 'windowsPtyAgent.js'), 'utf8')
+
+ const branch = patched.slice(
+ patched.indexOf('if (!this._useConptyDll) {'),
+ patched.indexOf('else {', patched.indexOf('if (!this._useConptyDll) {'))
+ )
+ expect(branch).toContain('this._inSocket.destroy();')
+ expect(branch.indexOf('this._inSocket.destroy();')).toBeGreaterThan(
+ branch.indexOf('this._conoutSocketWorker.dispose();')
+ )
+ expect(branch.indexOf('this._inSocket.destroy();')).toBeGreaterThan(
+ branch.indexOf('this._getConsoleProcessList()')
+ )
+ // Pinned so a future "sync the relay asset to config/patches" cannot copy the regression back.
+ expect(patched).not.toBe(readFileSync(desktopPath('windowsPtyAgent.js'), 'utf8'))
+ })
+
+ it('installs and verifies idempotently', () => {
+ const fixture = writeNodePtyFixture('1.1.0')
+
+ patchNodePtyWindowsTeardown(fixture.root)
+ const once = PATCHED_FILES.map((file) => readFileSync(join(fixture.libDir, file), 'utf8'))
+ for (const file of PATCHED_FILES) {
+ expect(existsSync(`${join(fixture.libDir, file)}.orca-patch-${process.pid}`)).toBe(false)
+ }
+ expect(() => assertPatchedNodePtyWindowsTeardown(fixture.root)).not.toThrow()
+
+ patchNodePtyWindowsTeardown(fixture.root)
+ expect(PATCHED_FILES.map((file) => readFileSync(join(fixture.libDir, file), 'utf8'))).toEqual(
+ once
+ )
+ })
+
+ it('refuses a different package version or unexpected source', () => {
+ const wrongVersion = writeNodePtyFixture('1.2.0-beta.11')
+ expect(() => patchNodePtyWindowsTeardown(wrongVersion.root)).toThrow('expected 1.1.0')
+
+ for (const file of PATCHED_FILES) {
+ const drifted = writeNodePtyFixture('1.1.0')
+ const path = join(drifted.libDir, file)
+ writeFileSync(path, `${readFileSync(path, 'utf8')}\n// drift`)
+ expect(() => patchNodePtyWindowsTeardown(drifted.root)).toThrow('unexpected node-pty')
+ }
+ })
+
+ it('refuses a half-applied tree, so one file cannot pass for both', () => {
+ for (const file of PATCHED_FILES) {
+ const partial = writeNodePtyFixture('1.1.0')
+ const fixture = writeNodePtyFixture('1.1.0')
+ patchNodePtyWindowsTeardown(fixture.root)
+ writeFileSync(join(partial.libDir, file), readFileSync(join(fixture.libDir, file), 'utf8'))
+ expect(() => assertPatchedNodePtyWindowsTeardown(partial.root)).toThrow('is not installed')
+ }
+ })
+})
+
+/** A published node-pty tree, rebuilt by un-applying the desktop hunks from the installed one. */
+function writeNodePtyFixture(version) {
+ const root = mkdtempSync(join(projectDir, '.node-pty-teardown-patch-test-'))
+ cleanupDirs.push(root)
+ const libDir = join(root, 'node_modules', 'node-pty', 'lib')
+ mkdirSync(libDir, { recursive: true })
+ writeFileSync(join(root, 'node_modules', 'node-pty', 'package.json'), JSON.stringify({ version }))
+ for (const file of PATCHED_FILES) {
+ const desktop = readFileSync(desktopPath(file), 'utf8')
+ for (const [marker] of DESKTOP_HUNKS[file]) {
+ expect(desktop).toContain(marker)
+ }
+ writeFileSync(join(libDir, file), unapplyDesktopHunks(file, desktop))
+ }
+ return { root, libDir }
+}
+
+/**
+ * Reverse of the published-to-desktop transform.
+ *
+ * `windowsTerminal.js` is taken verbatim from the desktop, so the asset's own replacement table is
+ * the transform and reversing it is exact. `windowsPtyAgent.js` deliberately diverges, so its
+ * published form is rebuilt from the desktop hunk instead -- which is also what makes this file the
+ * place that notices if the desktop hunk itself ever moves.
+ */
+function unapplyDesktopHunks(file, desktop) {
+ if (file === 'windowsPtyAgent.js') {
+ let published = desktop
+ for (const [patched, original] of DESKTOP_HUNKS[file]) {
+ expect(published.split(patched).length - 1).toBe(1)
+ published = published.replace(patched, original)
+ }
+ return published
+ }
+ const asset = readFileSync(
+ join(projectDir, 'config', 'relay-assets', 'node-pty-1.1.0-windows-pty-teardown-patch.cjs'),
+ 'utf8'
+ )
+ const { PATCH_TARGETS } = loadPatchTargets(asset)
+ const target = PATCH_TARGETS.find((entry) => entry.relativePath.at(-1) === file)
+ expect(target).toBeDefined()
+ let published = desktop
+ for (const [from, to] of target.replacements.toReversed()) {
+ expect(published.split(to).length - 1).toBe(1)
+ published = published.replace(to, from)
+ }
+ return published
+}
+
+function loadPatchTargets(assetSource) {
+ const module = { exports: {} }
+ const factory = new Function(
+ 'module',
+ 'exports',
+ 'require',
+ `${assetSource}\nmodule.exports.PATCH_TARGETS = PATCH_TARGETS`
+ )
+ factory(module, module.exports, require)
+ return module.exports
+}
diff --git a/config/scripts/skill-description-length.test.mjs b/config/scripts/skill-description-length.test.mjs
new file mode 100644
index 00000000000..e7a9db79541
--- /dev/null
+++ b/config/scripts/skill-description-length.test.mjs
@@ -0,0 +1,39 @@
+import { readdirSync, readFileSync } from 'node:fs'
+import { join, resolve } from 'node:path'
+import { describe, expect, it } from 'vitest'
+import { parse } from 'yaml'
+
+const skillsDir = resolve(import.meta.dirname, '../../skills')
+// Why: the Agent Skills spec caps `description` at 1024 chars and conforming installers
+// reject the whole skill (#17935); the frontmatter is what the installer parses, so check it.
+const MAX_DESCRIPTION_LENGTH = 1024
+
+function readDescription(skillName) {
+ const skillMarkdown = readFileSync(join(skillsDir, skillName, 'SKILL.md'), 'utf8')
+ const frontmatter = /^---\r?\n([\s\S]*?)\r?\n---\r?\n/u.exec(skillMarkdown)?.[1]
+
+ expect(frontmatter, `${skillName}: missing frontmatter`).toBeDefined()
+
+ return parse(frontmatter ?? '').description
+}
+
+describe('bundled skill descriptions', () => {
+ const skillNames = readdirSync(skillsDir, { withFileTypes: true })
+ .filter((entry) => entry.isDirectory())
+ .map((entry) => entry.name)
+
+ it('discovers the bundled skills', () => {
+ expect(skillNames).toContain('orchestration')
+ })
+
+ it.each(skillNames)('%s keeps description within the Agent Skills spec limit', (name) => {
+ const description = readDescription(name)
+
+ expect(typeof description, `${name}: description must be a string`).toBe('string')
+ expect(description.trim().length, `${name}: description is empty`).toBeGreaterThan(0)
+ expect(
+ description.length,
+ `${name}: description is ${description.length} chars`
+ ).toBeLessThanOrEqual(MAX_DESCRIPTION_LENGTH)
+ })
+})
diff --git a/docs/assets/readme-downloads.svg b/docs/assets/readme-downloads.svg
index ef8ebb61bb4..c240c965fee 100644
--- a/docs/assets/readme-downloads.svg
+++ b/docs/assets/readme-downloads.svg
@@ -1,5 +1,5 @@
- |