Files
orca/cloud/dev/scripts/verify-relay-capacity-transition.mjs
T
Jinwoo Hong 7c119465b0 fix(relay): define restart-safe by the cell runtime, and refuse waves without headroom (#24259)
* fix(relay): let a same-cap drain finish when only unplaceable hosts remain

The c28 canary on 2026-10-01 drained the cell to zero live connections, but four
hosts with no free slot anywhere kept redialling and held director leases on it,
so the restart-safe wait timed out and left the cell isolated and empty.

The drain wait now also passes once the runtime has carried nothing for a
sustained quiet window while a small, capped number of leases remain, and logs
the escape. Apply modes also refuse a cell whose hosts exceed 80% of the free
slots on the other general cells, so a wave cannot strand hosts in the first place.

Claude-Session: ced32ebb-7155-4413-adad-1eccd14c2010

* fix(relay): define restart-safe by the cell runtime, not director leases

Replaces the opt-in stranded-host escape with a corrected definition. A
restart is safe when the cell runtime carries nothing live and no migration
is open, sustained for the drain pace window. Director activity leases lag
hosts that already left or cannot be placed, so they are reported in a
progress line and the verified result instead of blocking the restart.

The same-cap drain passes its existing pace window. The headroom script is
added to the trusted evidence code paths with the other production scripts.

Claude-Session: ced32ebb-7155-4413-adad-1eccd14c2010

* fix(relay): print stranded director counts on every restart-safe sample

Each restart-safe poll now prints its sample count and the director's
restart-blocking leases, request units, reserved remainder, and migrations
under `stranded`; the verified line carries the same object. Open migrations
still block because each is pinned to the cell incarnation a restart replaces.

Claude-Session: ced32ebb-7155-4413-adad-1eccd14c2010

* fix(relay): require the pace window for every live restart-safe wait

Pre-auth and total connections no longer reset the restart-safe window:
on drained c28 they flickered with unauthenticated redials in a third of
samples, which a restart does not lose. They stay in the progress output.

Every live restart-safe call must now pass --pace-window-ms. The capacity
job and staging proof drain unpaced, so they pass the production 300000 ms
window, and the calls that relied on the 180000 ms default get 480000 ms.

Headroom free slots now follow the director's placement rule: the admission
pause minus the larger of observed and enforced units, minus outstanding
control reservations.

Claude-Session: ced32ebb-7155-4413-adad-1eccd14c2010
2026-09-30 22:25:36 -04:00

522 lines
21 KiB
JavaScript

import { pathToFileURL } from 'node:url'
import { fetchAdminOnceMore } from './relay-admin-transient-retry.mjs'
const CAPACITY_PROTOCOL = 2
const POLL_INTERVAL_MS = 5_000
function integer(value, name) {
const parsed = Number(value)
if (!Number.isSafeInteger(parsed) || parsed < 0) throw new Error(`${name} is invalid`)
return parsed
}
function signedInteger(value, name) {
if (typeof value !== 'number' || !Number.isSafeInteger(value)) {
throw new Error(`${name} is invalid`)
}
return value
}
export function parseCapacityTransitionArguments(argv) {
const values = {}
for (let index = 0; index < argv.length; index += 2) {
const key = argv[index]
const value = argv[index + 1]
if (!key?.startsWith('--') || value === undefined) throw new Error('invalid arguments')
values[key.slice(2)] = value
}
for (const key of [
'director-origin',
'cell-origin',
'cell-id',
'heartbeat',
'admission',
'draining',
'activity'
]) {
if (!values[key]) throw new Error(`missing --${key}`)
}
if (!['fresh', 'stale', 'either'].includes(values.heartbeat)) {
throw new Error('--heartbeat must be fresh, stale, or either')
}
if (
!['general', 'migration-only', 'general-or-migration-only', 'non-general', 'either'].includes(
values.admission
)
) {
throw new Error(
'--admission must be general, migration-only, general-or-migration-only, non-general, or either'
)
}
if (!['required', 'forbidden', 'either'].includes(values.draining)) {
throw new Error('--draining must be required, forbidden, or either')
}
if (!['quiescent', 'restart-safe', 'allowed'].includes(values.activity)) {
throw new Error('--activity must be quiescent, restart-safe, or allowed')
}
const runtime = values.runtime ?? 'required'
if (!['required', 'unavailable'].includes(runtime)) {
throw new Error('--runtime must be required or unavailable')
}
if (
values.activity === 'restart-safe' &&
runtime === 'required' &&
(values.admission !== 'migration-only' || values.draining !== 'required')
) {
throw new Error('restart-safe activity requires migration-only admission and draining')
}
if (
runtime === 'unavailable' &&
(values.heartbeat !== 'stale' ||
values.admission !== 'migration-only' ||
values.draining !== 'either' ||
values.activity !== 'restart-safe')
) {
throw new Error('unavailable runtime requires stale migration-only durable state')
}
const origin = new URL(values['director-origin'])
const cellOrigin = new URL(values['cell-origin'])
if (
origin.protocol !== 'https:' ||
origin.origin !== values['director-origin'] ||
cellOrigin.protocol !== 'https:' ||
cellOrigin.origin !== values['cell-origin']
) {
throw new Error('origins must be canonical HTTPS origins')
}
const hardCap = values['hard-cap'] === undefined
? undefined
: integer(values['hard-cap'], '--hard-cap')
const unobservedBound = values['unobserved-bound'] === undefined
? undefined
: integer(values['unobserved-bound'], '--unobserved-bound')
if ((hardCap === undefined) !== (unobservedBound === undefined)) {
throw new Error('capacity expectations must be paired')
}
if (runtime === 'unavailable' && hardCap !== undefined) {
throw new Error('unavailable runtime cannot prove live capacity')
}
const expectedImageDigests = values['expected-image-digests']?.split(',') ?? []
if (
new Set(expectedImageDigests).size !== expectedImageDigests.length ||
expectedImageDigests.some((digest) => !/^sha256:[a-f0-9]{64}$/.test(digest))
) {
throw new Error('--expected-image-digests is invalid')
}
if (runtime === 'unavailable' && expectedImageDigests.length > 0) {
throw new Error('unavailable runtime cannot prove a live image')
}
const regionalRehomeProtocol = values['regional-rehome-protocol'] === undefined
? undefined
: integer(values['regional-rehome-protocol'], '--regional-rehome-protocol')
if (regionalRehomeProtocol !== undefined && ![0, 1, 3].includes(regionalRehomeProtocol)) {
throw new Error('--regional-rehome-protocol must be 0, 1, or 3')
}
if (runtime === 'unavailable' && regionalRehomeProtocol !== undefined) {
throw new Error('unavailable runtime cannot prove the regional rehome protocol')
}
const paceWindowMs = values['pace-window-ms'] === undefined
? undefined
: integer(values['pace-window-ms'], '--pace-window-ms')
const liveRestartSafe = values.activity === 'restart-safe' && runtime === 'required'
if ((paceWindowMs !== undefined) !== liveRestartSafe) {
throw new Error('--pace-window-ms is required exactly for restart-safe activity on a live runtime')
}
return {
directorOrigin: origin.origin,
cellOrigin: cellOrigin.origin,
cellId: values['cell-id'],
heartbeat: values.heartbeat,
admission: values.admission,
draining: values.draining,
activity: values.activity,
runtime,
expectedImageDigests,
...(regionalRehomeProtocol === undefined ? {} : { regionalRehomeProtocol }),
...(paceWindowMs === undefined ? {} : { paceWindowMs }),
hardCap,
unobservedBound,
timeoutMs: integer(values['timeout-ms'] ?? 180_000, '--timeout-ms')
}
}
async function responseJson(response, label) {
const body = await response.json().catch(() => ({}))
if (!response.ok) throw new Error(`${label} returned ${response.status}`)
return body
}
async function cellRuntime(fetchImpl, config, token) {
let response
try {
response = await fetchImpl(`${config.cellOrigin}/v1/admin/runtime-status`, {
method: 'POST',
headers: { authorization: `Bearer ${token}`, 'content-type': 'application/json' },
body: JSON.stringify({ v: 1 }),
signal: AbortSignal.timeout(30_000)
})
} catch (error) {
if (config.runtime === 'unavailable') return null
throw error
}
if ([502, 503, 504].includes(response.status)) {
await response.arrayBuffer().catch(() => undefined)
return null
}
return await responseJson(response, 'cell runtime status')
}
function offlineRollbackMatches(status, config) {
if (status.admissionState !== 'migration-only') {
throw new Error('capacity transition admission does not match the required state')
}
const durableCounts = [
status.activityLeases,
status.activityRequestUnits,
status.reservedRequests,
status.restartBlockingActivityLeases,
status.restartBlockingActivityRequestUnits,
status.outgoingMigrations,
status.incomingMigrations,
status.connectionCapacity?.pendingControlReservations
]
const restartBlockingReservedRequests = signedInteger(
status.restartBlockingReservedRequests,
'restart-blocking reserved requests'
)
// Only a positive remainder is unexplained; the other gates reject real work.
return heartbeatMatches(status, config.heartbeat) &&
durableCounts.every((value) => integer(value, 'durable activity count') === 0) &&
restartBlockingReservedRequests <= 0
}
function directorActivityMatches(status, config) {
const durable = [
status.activityLeases,
status.reservedRequests,
status.outgoingMigrations,
status.incomingMigrations
]
// Reconnect reservations survive replacement; draining prevents activation on this process.
const transient = [
status.connectionCapacity?.observedConnections,
status.connectionCapacity?.inFlightConnections,
status.connectionCapacity?.reservedConnectionUnits,
status.connectionCapacity?.enforcedConnectionUnits,
status.connectionCapacity?.pendingControlReservations
]
const restartSafe = config.activity !== 'restart-safe' || (() => {
// Leases lag hosts that left or cannot be placed, so the runtime's counters gate instead;
// still validated so a malformed director fails closed.
integer(status.restartBlockingActivityLeases, 'restart-blocking activity leases')
integer(
status.restartBlockingActivityRequestUnits,
'restart-blocking activity request units'
)
// Open migrations are pinned to this cell's incarnation, which a restart replaces.
return integer(status.outgoingMigrations, 'outgoing migrations') === 0 &&
integer(status.incomingMigrations, 'incoming migrations') === 0
})()
const quiescent = [...durable, ...transient]
.filter((value) => value !== undefined)
.every((value) => integer(value, 'activity count') === 0)
if (
(config.admission === 'general' && status.admissionState !== 'general') ||
(config.admission === 'migration-only' && status.admissionState !== 'migration-only') ||
// A failed same-cap canary leaves its cell migration-only; the documented
// rollback recovery must accept that state alongside a completed general roll.
(config.admission === 'general-or-migration-only' &&
!['general', 'migration-only'].includes(status.admissionState)) ||
(config.admission === 'non-general' &&
!['existing-only', 'migration-only'].includes(status.admissionState))
) {
throw new Error('capacity transition admission does not match the required state')
}
return config.activity === 'allowed' ||
(config.activity === 'restart-safe' ? restartSafe : quiescent)
}
function capacityMatches(status, config) {
if (config.hardCap === undefined) return true
const capacity = status.connectionCapacity
return (
capacity?.hardCap === config.hardCap &&
capacity.unobservedBound === config.unobservedBound &&
capacity.controlRebindReserve === 100 &&
capacity.ordinaryConnectionLimit === config.hardCap - 100 &&
capacity.normalAdmissionPause === config.hardCap - 100 - config.unobservedBound
)
}
function heartbeatMatches(status, expectation) {
if (expectation === 'either') return true
const fresh = status.connectionCapacity?.heartbeatFresh ?? status.runtime?.heartbeatFresh
return fresh === (expectation === 'fresh')
}
function aggregateCount(value) {
const parsed = Number(value)
return Number.isSafeInteger(parsed) && parsed >= 0 ? parsed : null
}
function signedAggregateCount(value) {
return typeof value === 'number' && Number.isSafeInteger(value) ? value : null
}
function capacityObservation(capacity) {
if (capacity === null || capacity === undefined) return null
return {
hardCap: aggregateCount(capacity.hardCap),
controlRebindReserve: aggregateCount(capacity.controlRebindReserve),
ordinaryConnectionLimit: aggregateCount(capacity.ordinaryConnectionLimit),
unobservedBound: aggregateCount(capacity.unobservedBound),
normalAdmissionPause: aggregateCount(capacity.normalAdmissionPause),
observedConnections: aggregateCount(capacity.observedConnections),
inFlightConnections: aggregateCount(capacity.inFlightConnections),
reservedConnectionUnits: aggregateCount(capacity.reservedConnectionUnits),
enforcedConnectionUnits: aggregateCount(capacity.enforcedConnectionUnits),
pendingControlReservations: aggregateCount(capacity.pendingControlReservations),
heartbeatFresh: typeof capacity.heartbeatFresh === 'boolean'
? capacity.heartbeatFresh
: null
}
}
function transitionObservation(runtime, status) {
return {
runtimeAvailable: runtime !== null,
admissionState: ['general', 'migration-only', 'existing-only'].includes(status.admissionState)
? status.admissionState
: null,
draining: typeof runtime?.draining === 'boolean' ? runtime.draining : null,
runtime: runtime === null
? null
: {
totalConnections: aggregateCount(runtime.runtime?.totalConnections),
preAuthConnections: aggregateCount(runtime.runtime?.preAuthConnections),
inFlightConnections: aggregateCount(runtime.runtime?.inFlightConnections),
reservedConnectionUnits: aggregateCount(runtime.runtime?.reservedConnectionUnits),
enforcedConnectionUnits: aggregateCount(runtime.runtime?.enforcedConnectionUnits),
controls: aggregateCount(runtime.runtime?.controls),
splices: aggregateCount(runtime.runtime?.splices),
pendingSplices: aggregateCount(runtime.runtime?.pendingSplices),
queuedBytes: aggregateCount(runtime.runtime?.queuedBytes)
},
director: {
activityLeases: aggregateCount(status.activityLeases),
activityRequestUnits: aggregateCount(status.activityRequestUnits),
reservedRequests: aggregateCount(status.reservedRequests),
restartBlockingActivityLeases: aggregateCount(status.restartBlockingActivityLeases),
restartBlockingActivityRequestUnits:
aggregateCount(status.restartBlockingActivityRequestUnits),
restartBlockingReservedRequests:
signedAggregateCount(status.restartBlockingReservedRequests),
outgoingMigrations: aggregateCount(status.outgoingMigrations),
incomingMigrations: aggregateCount(status.incomingMigrations)
},
runtimeCapacity: capacityObservation(runtime?.connectionCapacity),
directorCapacity: capacityObservation(status.connectionCapacity),
runtimeHeartbeatFresh: typeof status.runtime?.heartbeatFresh === 'boolean'
? status.runtime.heartbeatFresh
: null
}
}
// Director records of hosts that left or cannot be placed; reported, never gated on.
function strandedObservation(observation) {
const director = observation.director
return director === undefined
? null
: {
restartBlockingActivityLeases: director.restartBlockingActivityLeases,
restartBlockingActivityRequestUnits: director.restartBlockingActivityRequestUnits,
restartBlockingReservedRequests: director.restartBlockingReservedRequests,
outgoingMigrations: director.outgoingMigrations,
incomingMigrations: director.incomingMigrations
}
}
function runtimeQuiescent(runtime, config) {
if (
runtime.role !== 'cell' ||
runtime.cellId !== config.cellId ||
runtime.cellUrl !== config.cellOrigin ||
// Legacy pre-rehome images omit the field; the exact digest binds absence to protocol 0.
(config.regionalRehomeProtocol !== undefined &&
(runtime.regionalRehomeProtocol ?? 0) !== config.regionalRehomeProtocol) ||
(config.expectedImageDigests?.length > 0 &&
!config.expectedImageDigests.includes(runtime.imageDigest))
) {
throw new Error('capacity transition runtime does not match the cell')
}
if (
(config.draining === 'required' && runtime.draining !== true) ||
(config.draining === 'forbidden' && runtime.draining === true)
) {
return false
}
const counts = [runtime.runtime?.totalConnections, runtime.runtime?.preAuthConnections]
if (counts.some((value) => value === undefined)) {
throw new Error('capacity transition runtime is incomplete')
}
if (runtime.connectionCapacity !== null && runtime.connectionCapacity !== undefined) {
if (runtime.runtime?.enforcedConnectionUnits === undefined) {
throw new Error('capacity transition runtime is incomplete')
}
counts.push(runtime.runtime.enforcedConnectionUnits)
}
const quiescent = counts.every((value) => integer(value, 'runtime connection count') === 0)
// Total and pre-auth connections also count unauthenticated redials, which never stop while
// hosts have nowhere else to go and lose nothing on a restart.
const restartSafe = config.activity !== 'restart-safe' ||
[
runtime.runtime?.inFlightConnections,
runtime.runtime?.reservedConnectionUnits,
runtime.runtime?.controls,
runtime.runtime?.splices,
runtime.runtime?.pendingSplices,
runtime.runtime?.queuedBytes
].every((value) => integer(value, 'live runtime count') === 0)
if (
config.heartbeat === 'fresh' &&
!capacityMatches({ connectionCapacity: runtime.connectionCapacity }, config)
) {
return false
}
return config.activity === 'allowed' ||
(config.activity === 'restart-safe' ? restartSafe : quiescent)
}
export async function verifyCapacityTransition(config, overrides = {}) {
if (
config.activity === 'restart-safe' &&
config.runtime !== 'unavailable' &&
(config.admission !== 'migration-only' || config.draining !== 'required')
) {
throw new Error('restart-safe activity requires migration-only admission and draining')
}
const fetchImpl = overrides.fetch ?? fetch
const progress = overrides.progress ??
((line) => process.stdout.write(`${JSON.stringify(line)}\n`))
const wait = overrides.wait ?? ((ms) => new Promise((resolve) => setTimeout(resolve, ms)))
const now = overrides.now ?? Date.now
const token = overrides.token ?? process.env.ORCA_RELAY_ADMIN_ID_TOKEN
if (!token || token.length > 8_192) throw new Error('admin identity token is unavailable')
const health = await responseJson(
await fetchAdminOnceMore(
fetchImpl,
`${config.directorOrigin}/health`,
{},
{ wait, timeoutMs: 15_000 }
),
'director health'
)
if (health.ok !== true || health.connectionCapacityProtocol !== CAPACITY_PROTOCOL) {
throw new Error('director is not capacity-protocol compatible')
}
const deadline = now() + config.timeoutMs
// A paced drain can refill the runtime until its window ends, so the quiet must outlast it.
// Each poll waits a full interval, so consecutive samples bound the quiet time from below.
const requiredRestartSafeSamples =
Math.max(Math.ceil((config.paceWindowMs ?? 0) / POLL_INTERVAL_MS) + 1, 2)
let restartSafeSamples = 0
let lastObservation = { runtimeAvailable: false }
for (;;) {
const runtime = await cellRuntime(fetchImpl, config, token)
lastObservation = { runtimeAvailable: runtime !== null }
if ((runtime === null) === (config.runtime === 'unavailable')) {
const result = await responseJson(
await fetchAdminOnceMore(
fetchImpl,
`${config.directorOrigin}/v1/admin/cell-status`,
{
method: 'POST',
headers: { authorization: `Bearer ${token}`, 'content-type': 'application/json' },
body: JSON.stringify({ v: 1, cellId: config.cellId })
},
{ wait }
),
'cell status'
)
const status = result.status
if (
status?.cellId !== config.cellId ||
status.cellUrl !== config.cellOrigin ||
status.runtime?.cellUrl !== config.cellOrigin
) {
throw new Error('capacity transition director status does not match the cell')
}
lastObservation = transitionObservation(runtime, status)
const restartBlockingReservedRequests =
config.activity === 'restart-safe'
? signedInteger(
status.restartBlockingReservedRequests,
'restart-blocking reserved requests'
)
: null
const matches = runtime === null
? offlineRollbackMatches(status, config)
: runtimeQuiescent(runtime, config) &&
directorActivityMatches(status, config) &&
capacityMatches(status, config) &&
heartbeatMatches(status, config.heartbeat)
if (
matches &&
(config.activity !== 'restart-safe' ||
restartSafeSamples + 1 >= requiredRestartSafeSamples)
) {
return {
cellId: status.cellId,
admissionState: status.admissionState,
assignments: integer(status.assignments, 'assignments'),
hardCap: runtime === null ? null : status.connectionCapacity?.hardCap ?? null,
unobservedBound:
runtime === null ? null : status.connectionCapacity?.unobservedBound ?? null,
heartbeatFresh:
status.connectionCapacity?.heartbeatFresh ?? status.runtime?.heartbeatFresh ?? false,
imageDigest: runtime?.imageDigest ?? null,
...(config.activity === 'restart-safe'
? { restartBlockingReservedRequests, stranded: strandedObservation(lastObservation) }
: {})
}
}
restartSafeSamples = matches ? restartSafeSamples + 1 : 0
} else {
restartSafeSamples = 0
}
if (config.activity === 'restart-safe') {
lastObservation = {
...lastObservation,
restartSafeSamples,
requiredRestartSafeSamples
}
progress({
event: 'relay_capacity_transition_restart_progress',
restartSafeSamples,
requiredRestartSafeSamples,
totalConnections: lastObservation.runtime?.totalConnections ?? null,
preAuthConnections: lastObservation.runtime?.preAuthConnections ?? null,
stranded: strandedObservation(lastObservation)
})
}
if (now() >= deadline) {
throw new Error(
`capacity transition verification timed out: ${JSON.stringify(lastObservation)}`
)
}
await wait(POLL_INTERVAL_MS)
}
}
export async function main(argv = process.argv.slice(2)) {
const result = await verifyCapacityTransition(parseCapacityTransitionArguments(argv))
process.stdout.write(`${JSON.stringify({ event: 'relay_capacity_transition_verified', ...result })}\n`)
}
if (import.meta.url === pathToFileURL(process.argv[1]).href) {
main().catch((error) => {
process.stderr.write(`${error instanceof Error ? error.message : String(error)}\n`)
process.exitCode = 1
})
}