mirror of
https://github.com/stablyai/orca.git
synced 2026-10-07 16:02:29 +00:00
* feat(relay): drain pace window as a reviewed same-cap input, with drain-aware 503 gates The same-cap roll drained every cell over a fixed 300 s window, so a US roll re-placed hosts at ~2/s and spent ~10 minutes draining and waiting for quiet. The window is now a dispatch input from a closed set (300000, 60000, 30000), defaulting to today's 300000. - Below the default is refused for anything but US general cells; Asia drains are bound by their targets' accept rate, and migration-only cells hold no hosts. - A non-default window must be named in the confirmation, the canary authority records it (v2), and a batch may run at its canary's window or slower only. - Each cell job re-checks the window, scales the restart-safe timeout with it (15-min lease + window, unchanged at the default), and records what the cell applied and when it settled. - The report-only shadow gate takes the director's drain-return deferrals out of the 503 count (window and baselines), adds the rung's 5-min non-drain 503 budget and a Retry-After check, and reports the measured re-placement rate. - relay-workflows.md documents the pace ladder and what each rung records. * fix(relay): judge paced drains on counted 503s against the pre-drain minutes, and seal the canary's pace verdict Review of #25639 replayed the shadow gate: it read 10-01 c29 as unverified (the log read stopped at 20k entries), false-blocked 10-02 c22, and was blind on four cells whose 24 h/48 h baseline held an incident. - Director 503s now come from Cloud Run's request_count, aligned per minute by Cloud Monitoring, so volume cannot truncate the count. - The background is the median of the 10 same-day minutes before the drain; the 24 h/48 h baselines are gone. - Scheduled 503s come out: drain-return deferrals and sticky/placement answers to a host's own early retry (host-rate-limited, host-in-flight), each split across the minutes its 30 s sample covers. - The rung budget counts only sustained breaches: two straight minutes over max(1.5x, +20) warn, over max(2x, +40) would-block. - The report carries a paceVerdict over the three pace checks. seal_canary downloads the canary cell's report and seals that verdict; a batch below the default pace needs PASS, from a report on the same cell that drained at that pace. - Docs: the step-down rule reads paceVerdict, 30 s waits for the lane service time (#25645), and the staging step is dropped since staging drains unpaced. Replayed read-only: 10-01 c29 would-block (9 minutes over 41.5/min); all nine 10-02 cells and 10-01 c25 paceVerdict PASS. * fix(relay): a partial count already past a block line blocks, in the shadow gate's Cloud SQL and pool checks A truncated FATAL count is a floor, and one runtime sample over the SQL-failure line is a fact, so neither waits for a complete read. The waiter-run rule still needs a complete run, since holes can join two runs into one. * fix(relay): a canary pace PASS needs a real cohort and whole telemetry; one median-based 503 check From the final review of #25639: - canaryPaceVerdict seals PASS only from a report that drained at least 400 hosts (about half a 10-02 US cell), so a near-empty canary cannot authorize a fast batch. - A director-metrics sub-window with fewer samples than one instance emits is unverified, so an empty or short Logging answer is never a calm drain. - seal_canary names the shadow artifact, report path and cell from the gate's normalized cell list, as cell_1 uploads it. - director503 folds into nonDrain503Budget as a single-minute spike rule, max(10x median, 200), dropping the pre-drain peak statistic. All 30 replayed windows keep their verdicts.
324 lines
14 KiB
JavaScript
324 lines
14 KiB
JavaScript
import { readFileSync } from 'node:fs'
|
|
import { pathToFileURL } from 'node:url'
|
|
import { requireSameEvidenceCode } from './relay-evidence-code-provenance.mjs'
|
|
|
|
// Migration-only by policy: zero hosts and no reservation, so a wave rolls one without
|
|
// displacing anybody. It enters and must leave migration-only, never general.
|
|
// C34 is an Asia spare that stays migration-only by policy.
|
|
export const SAME_CAP_MIGRATION_ONLY_CELLS = [
|
|
'production-gce-c17', 'production-gce-c18', 'production-gce-c34'
|
|
]
|
|
|
|
export const SAME_CAP_CELLS = [
|
|
'production-gce-c7', 'production-gce-c8', 'production-gce-c9', 'production-gce-c10',
|
|
'production-gce-c13', 'production-gce-c14', 'production-gce-c15', 'production-gce-c16',
|
|
'production-gce-c19', 'production-gce-c20', 'production-gce-c21', 'production-gce-c22',
|
|
'production-gce-c23', 'production-gce-c24', 'production-gce-c25', 'production-gce-c26',
|
|
'production-gce-c27', 'production-gce-c28', 'production-gce-c29', 'production-gce-c30',
|
|
'production-gce-c31', 'production-gce-c32', 'production-gce-c33',
|
|
...SAME_CAP_MIGRATION_ONLY_CELLS
|
|
]
|
|
|
|
// A general cell's wave isolates and restores it, advancing the selector twice; a
|
|
// migration-only cell's isolate and restore are both no-ops, so its wave advances nothing.
|
|
export function selectorWaveDelta(cellId) {
|
|
return SAME_CAP_MIGRATION_ONLY_CELLS.includes(cellId) ? 0 : 2
|
|
}
|
|
|
|
// The cell spreads its drain sends evenly over this window (host-session-registry.ts drain), so
|
|
// it sets the re-placement arrival rate: hosts / window. A closed set, stepped down one rung at a
|
|
// time per cloud/docs/relay-workflows.md; the first entry is the default every wave used before.
|
|
export const SAME_CAP_DRAIN_PACE_WINDOWS_MS = [300_000, 60_000, 30_000]
|
|
export const DEFAULT_SAME_CAP_DRAIN_PACE_WINDOW_MS = SAME_CAP_DRAIN_PACE_WINDOWS_MS[0]
|
|
|
|
// Only US general cells may drain faster than the default. An Asia drain is bounded by its
|
|
// targets' own accept rate (~4-6 hosts/s per cell over 176 ms round trips), and a migration-only
|
|
// cell carries no hosts, so neither has anything to gain from a shorter window.
|
|
export const SAME_CAP_FAST_DRAIN_PACE_CELLS = [
|
|
'production-gce-c7', 'production-gce-c8', 'production-gce-c9', 'production-gce-c10',
|
|
'production-gce-c13', 'production-gce-c14', 'production-gce-c15', 'production-gce-c16',
|
|
'production-gce-c19', 'production-gce-c20', 'production-gce-c21', 'production-gce-c22',
|
|
'production-gce-c23', 'production-gce-c24', 'production-gce-c25', 'production-gce-c26',
|
|
'production-gce-c32', 'production-gce-c33'
|
|
]
|
|
|
|
export function drainPaceWindowMs(value, cellIds) {
|
|
const parsed = /^[1-9][0-9]*$/.test(value ?? '') ? Number(value) : Number.NaN
|
|
if (!SAME_CAP_DRAIN_PACE_WINDOWS_MS.includes(parsed)) {
|
|
throw new Error(
|
|
`drain pace window must be one of ${SAME_CAP_DRAIN_PACE_WINDOWS_MS.join(', ')} ms`
|
|
)
|
|
}
|
|
if (parsed !== DEFAULT_SAME_CAP_DRAIN_PACE_WINDOW_MS) {
|
|
const slow = cellIds.filter((cell) => !SAME_CAP_FAST_DRAIN_PACE_CELLS.includes(cell))
|
|
if (slow.length > 0) {
|
|
throw new Error(
|
|
`drain pace window ${parsed} ms is for US general cells only, not ${slow.join(',')}`
|
|
)
|
|
}
|
|
}
|
|
return parsed
|
|
}
|
|
|
|
// The default keeps the confirmation every earlier wave typed; any other window must be named
|
|
// in it, so a dispatch cannot run a pace its confirmation did not state.
|
|
function confirmationPaceSuffix(paceWindowMs) {
|
|
return paceWindowMs === DEFAULT_SAME_CAP_DRAIN_PACE_WINDOW_MS
|
|
? ''
|
|
: ` drain-pace-window-ms=${paceWindowMs}`
|
|
}
|
|
|
|
export function entryAdmission(cellId) {
|
|
return SAME_CAP_MIGRATION_ONLY_CELLS.includes(cellId) ? 'migration-only' : 'general'
|
|
}
|
|
|
|
function digest(value, name) {
|
|
if (!/^sha256:[a-f0-9]{64}$/.test(value ?? '')) throw new Error(`${name} is invalid`)
|
|
return value
|
|
}
|
|
|
|
function cells(value) {
|
|
const parsed = value.split(',').map((cell) => cell.trim()).filter(Boolean)
|
|
if (
|
|
parsed.length < 1 ||
|
|
parsed.length > 10 ||
|
|
new Set(parsed).size !== parsed.length ||
|
|
parsed.some((cell) => !SAME_CAP_CELLS.includes(cell))
|
|
) throw new Error('same-cap wave cells are invalid')
|
|
// Every later cell offsets from one per-wave selector delta, and the two classes
|
|
// have different ones, so a mixed wave has no single offset any cell could use.
|
|
if (new Set(parsed.map(selectorWaveDelta)).size > 1) {
|
|
throw new Error('same-cap wave cells must be all general or all migration-only')
|
|
}
|
|
return parsed
|
|
}
|
|
|
|
export function validateSameCapWave(input) {
|
|
if (!['verify', 'canary-apply', 'batch-apply', 'rollback'].includes(input.mode)) {
|
|
throw new Error('same-cap wave mode is invalid')
|
|
}
|
|
const selected = cells(input.cellIds)
|
|
const targetDigest = digest(input.targetDigest, 'target digest')
|
|
const rollbackDigest = digest(input.rollbackDigest, 'rollback digest')
|
|
if (targetDigest === rollbackDigest) throw new Error('target and rollback digests must differ')
|
|
if (input.mode === 'canary-apply' && selected.length !== 1) {
|
|
throw new Error('canary mode requires exactly one cell')
|
|
}
|
|
// Ten is the wave workflow's statically declared serial cell-job chain, cell_1..cell_10.
|
|
if (input.mode === 'batch-apply' && (selected.length < 2 || selected.length > 10)) {
|
|
throw new Error('batch mode requires two to ten cells')
|
|
}
|
|
// Later waves expect the selector to advance by exactly 2 per predecessor,
|
|
// which a resumed rollback cell (isolate skipped, +1) violates.
|
|
if (input.mode === 'rollback' && selected.length !== 1) {
|
|
throw new Error('rollback mode requires exactly one cell')
|
|
}
|
|
const paceWindowMs = drainPaceWindowMs(input.drainPaceWindowMs, selected)
|
|
const mutation = input.mode !== 'verify'
|
|
const expectedConfirmation = (input.mode === 'rollback'
|
|
? `ROLL_BACK_RELAY_SAME_CAP ${rollbackDigest} ${selected.join(',')}`
|
|
: `ROLL_RELAY_SAME_CAP ${targetDigest} ${selected.join(',')}`) +
|
|
confirmationPaceSuffix(paceWindowMs)
|
|
if (mutation && input.confirmation !== expectedConfirmation) {
|
|
throw new Error('same-cap confirmation does not match the exact digest and cells')
|
|
}
|
|
if (!mutation && input.confirmation) throw new Error('verify does not accept confirmation')
|
|
if (input.mode === 'batch-apply' && !/^[1-9][0-9]*$/.test(input.canaryRunId ?? '')) {
|
|
throw new Error('batch mode requires a canary run ID')
|
|
}
|
|
if (input.mode !== 'batch-apply' && input.canaryRunId) {
|
|
throw new Error('only batch mode accepts a canary run ID')
|
|
}
|
|
return { cells: selected, targetDigest, rollbackDigest, drainPaceWindowMs: paceWindowMs }
|
|
}
|
|
|
|
const CANARY_PACE_VERDICTS = ['PASS', 'WARN', 'WOULD_BLOCK', 'UNVERIFIED']
|
|
|
|
// A pace is a host arrival rate (hosts / window), so a canary proves it only with a real cohort:
|
|
// at least about half the 692-782 hosts a US general cell carried on 10-02, keeping any batch
|
|
// cell within ~2x of the rate the canary actually drained at.
|
|
export const CANARY_MIN_DRAINED_HOSTS = 400
|
|
|
|
// The canary cell's own pace checks, trusted only from a report on this cell that drained enough
|
|
// hosts at this pace; a cell whose image fell back to an unpaced drain proved nothing about it.
|
|
export function canaryPaceVerdict(report, cellId, paceWindowMs) {
|
|
if (
|
|
report?.cellId !== cellId ||
|
|
report.drain?.paceWindowMs !== paceWindowMs ||
|
|
report.drain?.appliedPaceWindowMs !== paceWindowMs ||
|
|
!(report.drain?.targetHosts >= CANARY_MIN_DRAINED_HOSTS) ||
|
|
!CANARY_PACE_VERDICTS.includes(report.paceVerdict)
|
|
) return 'UNVERIFIED'
|
|
return report.paceVerdict
|
|
}
|
|
|
|
export function canaryAuthority(input) {
|
|
const wave = validateSameCapWave({ ...input, mode: 'canary-apply', canaryRunId: '' })
|
|
if (!/^[0-9a-f]{40}$/.test(input.commitSha ?? '')) throw new Error('commit SHA is invalid')
|
|
if (!/^[1-9][0-9]*$/.test(input.runId ?? '')) throw new Error('run ID is invalid')
|
|
const selectorGeneration = Number(input.selectorGeneration)
|
|
const rehomeGeneration = Number(input.rehomeGeneration)
|
|
if (!Number.isSafeInteger(selectorGeneration) || selectorGeneration < 0) {
|
|
throw new Error('selector generation is invalid')
|
|
}
|
|
if (!Number.isSafeInteger(rehomeGeneration) || rehomeGeneration < 0) {
|
|
throw new Error('rehome generation is invalid')
|
|
}
|
|
return {
|
|
v: 2,
|
|
commitSha: input.commitSha,
|
|
runId: input.runId,
|
|
cellId: wave.cells[0],
|
|
targetDigest: wave.targetDigest,
|
|
rollbackDigest: wave.rollbackDigest,
|
|
drainPaceWindowMs: wave.drainPaceWindowMs,
|
|
paceVerdict: canaryPaceVerdict(input.shadowReport, wave.cells[0], wave.drainPaceWindowMs),
|
|
selectorGeneration: selectorGeneration + selectorWaveDelta(wave.cells[0]),
|
|
rehomeGeneration
|
|
}
|
|
}
|
|
|
|
export function verifyCanaryAuthority(authority, expected, repositoryRoot) {
|
|
const selectorGeneration = Number(expected.selectorGeneration)
|
|
// A mixed wave is already rejected, so the batch's first cell names the whole batch's class.
|
|
const batchCells = cells(expected.cellIds ?? '')
|
|
const batchAdmission = entryAdmission(batchCells[0])
|
|
const batchPaceWindowMs = drainPaceWindowMs(expected.drainPaceWindowMs, batchCells)
|
|
if (
|
|
authority?.v !== 2 ||
|
|
!/^[0-9a-f]{40}$/.test(authority.commitSha ?? '') ||
|
|
authority.runId !== expected.runId ||
|
|
authority.targetDigest !== expected.targetDigest ||
|
|
authority.rollbackDigest !== expected.rollbackDigest ||
|
|
!Number.isSafeInteger(authority.selectorGeneration) ||
|
|
authority.selectorGeneration < 0 ||
|
|
!Number.isSafeInteger(selectorGeneration) ||
|
|
selectorGeneration < authority.selectorGeneration ||
|
|
authority.rehomeGeneration !== Number(expected.rehomeGeneration) ||
|
|
!SAME_CAP_CELLS.includes(authority.cellId) ||
|
|
!SAME_CAP_DRAIN_PACE_WINDOWS_MS.includes(authority.drainPaceWindowMs) ||
|
|
!CANARY_PACE_VERDICTS.includes(authority.paceVerdict)
|
|
) throw new Error('canary authority does not match this batch')
|
|
// A canary proves its own pace and every slower one; a faster batch needs its own canary, and
|
|
// falling back to a slower pace mid-ladder never does.
|
|
if (batchPaceWindowMs < authority.drainPaceWindowMs) {
|
|
throw new Error(
|
|
`canary authority drained over ${authority.drainPaceWindowMs} ms, ` +
|
|
`so it cannot authorize a batch draining over ${batchPaceWindowMs} ms`
|
|
)
|
|
}
|
|
// A batch rolls up to ten cells back to back, so a faster one needs a canary whose own drain
|
|
// passed; the default is what every wave ran before, and stays available to any canary.
|
|
if (
|
|
batchPaceWindowMs !== DEFAULT_SAME_CAP_DRAIN_PACE_WINDOW_MS &&
|
|
authority.paceVerdict !== 'PASS'
|
|
) {
|
|
throw new Error(
|
|
`canary authority pace checks were ${authority.paceVerdict}, ` +
|
|
`so it cannot authorize a batch draining over ${batchPaceWindowMs} ms`
|
|
)
|
|
}
|
|
// A migration-only cell carries no hosts and a different cap, so rolling it proves nothing
|
|
// about a general batch, and its wave advances a different selector delta.
|
|
if (entryAdmission(authority.cellId) !== batchAdmission) {
|
|
throw new Error(
|
|
`canary authority cell ${authority.cellId} is ${entryAdmission(authority.cellId)}, ` +
|
|
`but this batch is ${batchAdmission}`
|
|
)
|
|
}
|
|
// Each cell checks exact live selector state; later batches may reuse this control epoch's canary.
|
|
requireSameEvidenceCode({
|
|
sealedSha: authority.commitSha,
|
|
currentSha: expected.commitSha,
|
|
label: 'relay same-cap canary authority',
|
|
repositoryRoot
|
|
})
|
|
return authority
|
|
}
|
|
|
|
// A missing or unreadable report seals UNVERIFIED rather than failing the seal: the default pace
|
|
// never needed one.
|
|
function readShadowReport(path) {
|
|
try {
|
|
return JSON.parse(readFileSync(path, 'utf8'))
|
|
} catch {
|
|
return null
|
|
}
|
|
}
|
|
|
|
function values(argv) {
|
|
const result = {}
|
|
for (let index = 0; index < argv.length; index += 2) {
|
|
if (!argv[index]?.startsWith('--') || argv[index + 1] === undefined) {
|
|
throw new Error('invalid arguments')
|
|
}
|
|
result[argv[index].slice(2)] = argv[index + 1]
|
|
}
|
|
return result
|
|
}
|
|
|
|
export function main(argv = process.argv.slice(2)) {
|
|
const command = argv.shift()
|
|
const input = values(argv)
|
|
if (command === 'validate') {
|
|
const wave = validateSameCapWave({
|
|
mode: input.mode,
|
|
cellIds: input['cell-ids'],
|
|
targetDigest: input['target-digest'],
|
|
rollbackDigest: input['rollback-digest'],
|
|
confirmation: input.confirmation,
|
|
canaryRunId: input['canary-run-id'],
|
|
drainPaceWindowMs: input['drain-pace-window-ms']
|
|
})
|
|
process.stdout.write(`${JSON.stringify(wave.cells)}\n`)
|
|
return
|
|
}
|
|
if (command === 'create-canary') {
|
|
process.stdout.write(`${JSON.stringify(canaryAuthority({
|
|
mode: 'canary-apply',
|
|
cellIds: input['cell-id'],
|
|
targetDigest: input['target-digest'],
|
|
rollbackDigest: input['rollback-digest'],
|
|
confirmation: input.confirmation,
|
|
drainPaceWindowMs: input['drain-pace-window-ms'],
|
|
shadowReport: readShadowReport(input['shadow-report']),
|
|
commitSha: input['commit-sha'],
|
|
runId: input['run-id'],
|
|
selectorGeneration: input['selector-generation'],
|
|
rehomeGeneration: input['rehome-generation']
|
|
}))}\n`)
|
|
return
|
|
}
|
|
if (command === 'cell-class') {
|
|
const cellId = input['cell-id']
|
|
if (!SAME_CAP_CELLS.includes(cellId)) throw new Error('same-cap wave cells are invalid')
|
|
process.stdout.write(`${JSON.stringify({
|
|
entryAdmission: entryAdmission(cellId),
|
|
selectorWaveDelta: selectorWaveDelta(cellId),
|
|
drainPaceWindowMs: drainPaceWindowMs(input['drain-pace-window-ms'], [cellId])
|
|
})}\n`)
|
|
return
|
|
}
|
|
if (command === 'verify-canary') {
|
|
verifyCanaryAuthority(JSON.parse(readFileSync(input.file, 'utf8')), {
|
|
commitSha: input['commit-sha'],
|
|
runId: input['run-id'],
|
|
cellIds: input['cell-ids'],
|
|
targetDigest: input['target-digest'],
|
|
rollbackDigest: input['rollback-digest'],
|
|
selectorGeneration: input['selector-generation'],
|
|
rehomeGeneration: input['rehome-generation'],
|
|
drainPaceWindowMs: input['drain-pace-window-ms']
|
|
})
|
|
return
|
|
}
|
|
throw new Error('unknown same-cap wave command')
|
|
}
|
|
|
|
if (import.meta.url === pathToFileURL(process.argv[1]).href) {
|
|
try { main() } catch (error) {
|
|
process.stderr.write(`${error instanceof Error ? error.message : String(error)}\n`)
|
|
process.exitCode = 1
|
|
}
|
|
}
|