mirror of
https://github.com/stablyai/orca.git
synced 2026-10-07 00:02:29 +00:00
feat(relay): one-command director deploy driver (#25642)
* feat(relay): add an operator-local driver for director deploys One command runs the audited director deploy: preflight, pause rehome if enabled, publish, deploy, optional cell configure, a digest-bound inspect, the monitor dry-run, and re-enable. It only dispatches the existing workflows, reads the published digest from the registry and the run log, reads the monitor verdict from its sealed state, and enables with the digests gcloud reports after the deploy. It stops at the first failure, records state, and resumes from it. * fix(relay): report and own the rehome pause; operator types every phrase - Read the control back from pause and enable runs whatever their conclusion, and report PAUSED or UNCONFIRMED loudly. - Resume re-enables only the pause this driver recorded (generation and run). - Ctrl-C and SIGTERM print the same state and resume report. - The operator types every workflow confirmation. A 5-minute soak gated on director 5xx runs before configure. - A dry run keeps no state file. The quiet check pages through all runs. Run IDs come only from the printed URL. One step table drives execute, dry run and resume. The monitor verdict reuses verify-authority. * fix(relay): read back only the driver's own rehome run; anchor the soak at the traffic switch - The pause and enable read-back accepts only the control line its own step's mode prints, at the generation its own dispatch expected. A run adopted after a crash is settled even when green. - The soak window opens a minute before the deploy run completed and is read a minute after it ends, for log ingestion lag. * fix(relay): say what typing ENABLE_REGIONAL_REHOMING commits to * fix(relay): the ENABLE prompt also names the 150 s evidence budget * refactor(relay): derive every deploy decision from live state; no resume machinery The driver keeps no state between runs. Each run reads the serving director, its configured cells and the rehome control, and skips what is already done. - The only rehome fact it owns is the run that paused rehome. A re-run names it (--pause-run), and the driver checks it against that run's log and the live generation. - A recover-enable line counts as the driver's own pause only with recovered: true. A director safety pause is never adopted (F3). - Publish runs before the pause. A fresh run that finds rehome paused stops unless given --pause-run or --rehome-disabled (F2). - Interrupts report a pause or enable still in flight as REHOME IS CHANGING (F1). PAUSE UNCONFIRMED and ENABLE UNCONFIRMED are distinct. - A tripped soak is judged again on fresh traffic (F5). SIGHUP is handled, and a pending signal stops the driver before its next dispatch. - Every stop prints the single command that finishes the deploy. Removes the state file, --resume, step statuses, monitor adoption and interrupted-dispatch adoption. * fix(relay): never report done over an unexplained pause; prove pause ownership by actor - G1: always read rehome; no early DONE. - G2: --pause-run must be a rehome-control run by the same user. - G3: a failed enable run is never an enable. - C1: a pause or enable is reported as changing from the moment it is dispatched. - C2: --leave-rehome-paused (was --rehome-disabled) refuses an enabled switch. Every printed command parses. - G4: the enable-in-flight report prints both finishing commands. - G5: the driver's own runs never block the quiet-lane check. * test(relay): port the round-3 probes: hard kill mid-pause, unexplained disable after a failed enable Claude-Session: 1145a80d-dec4-4a9b-9373-bbbb876b9041
This commit is contained in:
@@ -0,0 +1,987 @@
|
||||
// Operator-local driver for a production relay director deploy. It only dispatches the existing
|
||||
// audited workflows and reads their results back; every safety check stays in those workflows.
|
||||
//
|
||||
// It keeps no state between runs. Every decision comes from live state read at the start: the
|
||||
// serving director, its configured cells, and the rehome control. The one fact live state cannot
|
||||
// show, that a pause is this driver's own, is the run that made it (`--pause-run`), checked
|
||||
// against that run's log and the live generation. On any stop it prints the command that
|
||||
// finishes the deploy from wherever it got to.
|
||||
import { spawnSync } from 'node:child_process'
|
||||
import { appendFileSync, mkdirSync, mkdtempSync, readFileSync, rmSync } from 'node:fs'
|
||||
import { homedir, tmpdir } from 'node:os'
|
||||
import { join, resolve } from 'node:path'
|
||||
import { createInterface } from 'node:readline/promises'
|
||||
import { pathToFileURL } from 'node:url'
|
||||
import { verifyDryRunAuthority } from './relay-monitor-evidence.mjs'
|
||||
import {
|
||||
DIRECTOR_SERVICE,
|
||||
IMAGE_REPOSITORY,
|
||||
PROJECT,
|
||||
REGION,
|
||||
REPOSITORY,
|
||||
WORKFLOW_REF,
|
||||
WORKFLOWS,
|
||||
admissionInspectInputs,
|
||||
admissionInspectResult,
|
||||
blocksDeploy,
|
||||
configureInputs,
|
||||
directorDeployInputs,
|
||||
directorRevisions,
|
||||
logConfirmsPublishedDigest,
|
||||
monitorDryRunInputs,
|
||||
parseConfigureWave,
|
||||
pausedGeneration,
|
||||
publishInputs,
|
||||
rehomeEnableInputs,
|
||||
rehomeInspectInputs,
|
||||
rehomePauseInputs,
|
||||
rehomeResultFromLog,
|
||||
requireCommit,
|
||||
requireDigest,
|
||||
revisionDigest,
|
||||
validateDispatchInputs
|
||||
} from './relay-director-deploy-plan.mjs'
|
||||
|
||||
const ACTIVE_RUN_STATUSES = ['queued', 'in_progress', 'waiting', 'requested', 'pending']
|
||||
// The enable job verifies the monitor within 5 minutes of completion after ~2.5 minutes of setup.
|
||||
export const MONITOR_MAX_AGE_AT_ENABLE_MS = 150_000
|
||||
// The manual procedure watched the new director for 5+ minutes before configuring cells.
|
||||
export const SOAK_MS = 5 * 60_000
|
||||
// Traffic moves about a minute before the deploy run completes (ops-log 05:00:33Z vs 05:01:34Z).
|
||||
const TRAFFIC_SWITCH_LEAD_MS = 60_000
|
||||
// Request logs land up to a minute late; reading at the window's end would undercount it.
|
||||
const LOG_INGESTION_LAG_MS = 60_000
|
||||
const DIRECTOR_5XX_FILTER = [
|
||||
'resource.type="cloud_run_revision"',
|
||||
`resource.labels.service_name="${DIRECTOR_SERVICE}"`,
|
||||
`logName="projects/${PROJECT}/logs/run.googleapis.com%2Frequests"`,
|
||||
'httpRequest.status>=500'
|
||||
].join(' AND ')
|
||||
const LOG_COUNT_LIMIT = 5_000
|
||||
const LOG_ATTEMPTS = 6
|
||||
const LOG_INTERVAL_MS = 10_000
|
||||
const REHOME_HISTORY_RUNS = 5
|
||||
const RUN_ID = /^[1-9][0-9]*$/
|
||||
|
||||
export class DriverStop extends Error {}
|
||||
|
||||
export function parseDriverArguments(argv, home = homedir()) {
|
||||
const config = {
|
||||
dryRun: false,
|
||||
leaveRehomePaused: false,
|
||||
configure: [],
|
||||
logDirectory: join(home, '.orca', 'relay-director-deploy')
|
||||
}
|
||||
const runId = (value, key) => {
|
||||
if (!RUN_ID.test(value)) throw new Error(`${key} must be a run ID`)
|
||||
return Number(value)
|
||||
}
|
||||
for (let index = 0; index < argv.length; index += 1) {
|
||||
const key = argv[index]
|
||||
if (key === '--dry-run' || key === '--leave-rehome-paused') {
|
||||
config[key === '--dry-run' ? 'dryRun' : 'leaveRehomePaused'] = true
|
||||
continue
|
||||
}
|
||||
const value = argv[index + 1]
|
||||
if (value === undefined || value.startsWith('--')) throw new Error(`${key} needs a value`)
|
||||
index += 1
|
||||
if (key === '--commit') config.commit = requireCommit(value, '--commit')
|
||||
else if (key === '--configure') config.configure.push(parseConfigureWave(value))
|
||||
else if (key === '--publish-run') config.publishRun = runId(value, key)
|
||||
else if (key === '--pause-run') config.pauseRun = runId(value, key)
|
||||
else if (key === '--rehome-generation') {
|
||||
if (!/^(0|[1-9][0-9]*)$/.test(value)) throw new Error('--rehome-generation is invalid')
|
||||
config.rehomeGeneration = Number(value)
|
||||
} else if (key === '--log-directory') config.logDirectory = resolve(value)
|
||||
else throw new Error(`unsupported argument ${key}`)
|
||||
}
|
||||
if (!config.commit) throw new Error('missing --commit <40-character main commit to publish>')
|
||||
if (config.pauseRun && config.leaveRehomePaused) {
|
||||
throw new Error('--pause-run and --leave-rehome-paused contradict each other')
|
||||
}
|
||||
return config
|
||||
}
|
||||
|
||||
const timestamp = (ms) => new Date(ms).toISOString()
|
||||
const runUrl = (runId) => `https://github.com/${REPOSITORY}/actions/runs/${runId}`
|
||||
// Argument lists whose values never contain spaces.
|
||||
const words = (text) => text.split(' ')
|
||||
|
||||
export function createDriver(config, deps) {
|
||||
let logPath
|
||||
const typedPhrases = new Map()
|
||||
// Everything learned this run; nothing outlives the process.
|
||||
const known = { director: undefined, selector: undefined, control: undefined }
|
||||
let rollbackPoint
|
||||
let published = config.publishRun ? { runId: config.publishRun } : undefined
|
||||
let owned // { generation, runId }: the pause this driver made, the only rehome state it owns
|
||||
// A pause or enable this run cannot vouch for: `changing` while it may still apply on its own,
|
||||
// `unconfirmed` once it finished without a usable result. { kind, name, runId?, url? }
|
||||
let uncertain
|
||||
let login
|
||||
const ownRunIds = new Set()
|
||||
let enabled = false
|
||||
|
||||
function log(message) {
|
||||
const line = `${timestamp(deps.now())} ${message}`
|
||||
deps.print(line)
|
||||
if (logPath) appendFileSync(logPath, `${line}\n`)
|
||||
}
|
||||
|
||||
function command(program, args, input) {
|
||||
const result = deps.run(program, args, input)
|
||||
if (result.status !== 0) {
|
||||
const detail = String(result.stderr ?? '')
|
||||
.trim()
|
||||
.split('\n')
|
||||
.slice(-3)
|
||||
.join(' | ')
|
||||
throw new Error(`${program} ${args.slice(0, 4).join(' ')} failed: ${detail}`)
|
||||
}
|
||||
return String(result.stdout ?? '')
|
||||
}
|
||||
|
||||
const gh = (args, input) => command('gh', args, input)
|
||||
const ghJson = (args) => JSON.parse(gh(args))
|
||||
const gcloudJson = (args) =>
|
||||
JSON.parse(command('gcloud', [...args, '--project', PROJECT, '--format=json']))
|
||||
|
||||
function readDirector() {
|
||||
const service = gcloudJson(
|
||||
words(`run services describe ${DIRECTOR_SERVICE} --region ${REGION}`)
|
||||
)
|
||||
const { servingRevision, rollbackRevision } = directorRevisions(service)
|
||||
const describe = (revision) =>
|
||||
gcloudJson(words(`run revisions describe ${revision} --region ${REGION}`))
|
||||
const serving = describe(servingRevision)
|
||||
const cellsJson = serving.spec.containers[0].env?.find(
|
||||
(variable) => variable.name === 'ORCA_RELAY_CELLS_JSON'
|
||||
)?.value
|
||||
return {
|
||||
servingRevision,
|
||||
servingDigest: revisionDigest(serving, `revision ${servingRevision}`),
|
||||
servingCreatedAt: serving.metadata?.creationTimestamp,
|
||||
cells: new Set(JSON.parse(cellsJson ?? '[]').map((cell) => cell.id)),
|
||||
rollbackRevision,
|
||||
rollbackDigest: revisionDigest(describe(rollbackRevision), `revision ${rollbackRevision}`)
|
||||
}
|
||||
}
|
||||
|
||||
function describeDirector(director) {
|
||||
return `serving ${director.servingRevision} ${director.servingDigest}; rollback ${director.rollbackRevision} ${director.rollbackDigest}`
|
||||
}
|
||||
|
||||
// Paginated per status: the repository-wide first page can hide an in-flight relay run.
|
||||
function activeRuns() {
|
||||
const active = []
|
||||
for (const status of ACTIVE_RUN_STATUSES) {
|
||||
const lines = gh(
|
||||
words(
|
||||
`api --paginate -X GET repos/${REPOSITORY}/actions/runs -f status=${status} -f per_page=100 --jq .workflow_runs[]|{id,path,name,status}`
|
||||
)
|
||||
)
|
||||
for (const run of lines.split('\n').filter(Boolean).map(JSON.parse)) {
|
||||
if (blocksDeploy(run.path) && !ownRunIds.has(run.id))
|
||||
active.push(`${run.name} ${runUrl(run.id)} (${run.status})`)
|
||||
}
|
||||
}
|
||||
return active
|
||||
}
|
||||
|
||||
function requireQuietLane() {
|
||||
const active = activeRuns()
|
||||
if (active.length > 0) {
|
||||
throw new DriverStop(
|
||||
`relay workflows are in flight; wait for them:\n ${active.join('\n ')}`
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
function viewRun(runId) {
|
||||
return ghJson(
|
||||
words(
|
||||
`run view ${runId} -R ${REPOSITORY} --json status,conclusion,attempt,headSha,headBranch,event,workflowName`
|
||||
)
|
||||
)
|
||||
}
|
||||
|
||||
async function waitForRun(run) {
|
||||
for (;;) {
|
||||
deps.stream(
|
||||
'gh',
|
||||
words(`run watch ${run.runId} -R ${REPOSITORY} --exit-status --interval 10`)
|
||||
)
|
||||
const view = viewRun(run.runId)
|
||||
if (view.status === 'completed') {
|
||||
log(`${run.name}: ${view.conclusion} ${run.url}`)
|
||||
return { ...run, conclusion: view.conclusion, attempt: view.attempt }
|
||||
}
|
||||
await deps.sleep(LOG_INTERVAL_MS)
|
||||
}
|
||||
}
|
||||
|
||||
// Dispatches and watches one workflow run. The run ID comes only from the URL `gh workflow run`
|
||||
// prints, so another dispatch by the same account can never be mistaken for this one.
|
||||
async function dispatch({ name, workflow, inputs: build, changesRehome = false }) {
|
||||
const inputs = validateDispatchInputs(build(view()))
|
||||
// Lets a signal that arrived during synchronous work stop the driver before it dispatches.
|
||||
await new Promise((resolveYield) => setImmediate(resolveYield))
|
||||
requireQuietLane()
|
||||
log(`dispatch ${name}: ${workflow.file} ${JSON.stringify(inputs)}`)
|
||||
if (changesRehome) uncertain = { kind: 'changing', name }
|
||||
const output = gh(
|
||||
words(`workflow run ${workflow.file} -R ${REPOSITORY} --ref ${WORKFLOW_REF} --json`),
|
||||
JSON.stringify(inputs)
|
||||
)
|
||||
const printed = output.match(/\/actions\/runs\/([0-9]+)/)
|
||||
if (!printed) {
|
||||
throw new DriverStop(
|
||||
`gh printed no run URL for ${workflow.file}; upgrade gh. The run may exist: find it before re-running`
|
||||
)
|
||||
}
|
||||
const runId = Number(printed[1])
|
||||
const url = runUrl(runId)
|
||||
ownRunIds.add(runId)
|
||||
if (changesRehome) Object.assign(uncertain, { runId, url })
|
||||
log(`${name}: run ${url}`)
|
||||
const dispatched = viewRun(runId)
|
||||
if (
|
||||
dispatched.workflowName !== workflow.name ||
|
||||
dispatched.event !== 'workflow_dispatch' ||
|
||||
dispatched.headBranch !== WORKFLOW_REF
|
||||
) {
|
||||
throw new DriverStop(
|
||||
`run ${runUrl(runId)} is not a ${workflow.file} dispatch on ${WORKFLOW_REF}`
|
||||
)
|
||||
}
|
||||
// `uncertain` is cleared by the step only once it has read the run's result.
|
||||
return await waitForRun({ name, runId, url, headSha: dispatched.headSha })
|
||||
}
|
||||
|
||||
function requireSuccess(run) {
|
||||
if (run.conclusion !== 'success') {
|
||||
throw new DriverStop(`${run.name} run ${run.url} concluded ${run.conclusion}`)
|
||||
}
|
||||
}
|
||||
|
||||
async function runLog(runId) {
|
||||
for (let attempt = 1; ; attempt += 1) {
|
||||
try {
|
||||
const text = gh(words(`run view ${runId} -R ${REPOSITORY} --log`))
|
||||
if (text.trim()) return text
|
||||
} catch (error) {
|
||||
if (attempt >= LOG_ATTEMPTS)
|
||||
throw new Error(`log for ${runUrl(runId)} is unavailable: ${error.message}`)
|
||||
}
|
||||
if (attempt >= LOG_ATTEMPTS) throw new Error(`log for ${runUrl(runId)} is empty`)
|
||||
await deps.sleep(LOG_INTERVAL_MS)
|
||||
}
|
||||
}
|
||||
|
||||
async function rehomeResult(runId) {
|
||||
try {
|
||||
return rehomeResultFromLog(await runLog(runId))
|
||||
} catch {
|
||||
return undefined
|
||||
}
|
||||
}
|
||||
|
||||
async function withArtifact(runId, artifact, read) {
|
||||
const directory = mkdtempSync(join(tmpdir(), 'relay-director-deploy-'))
|
||||
try {
|
||||
gh(words(`run download ${runId} -R ${REPOSITORY} -n ${artifact} -D ${directory}`))
|
||||
return await read(directory)
|
||||
} finally {
|
||||
rmSync(directory, { recursive: true, force: true })
|
||||
}
|
||||
}
|
||||
|
||||
async function typed(phrase, meaning = '') {
|
||||
if (typedPhrases.has(phrase)) return phrase
|
||||
const answer = (await deps.prompt(`Type ${phrase} to continue${meaning}: `)).trim()
|
||||
if (answer !== phrase)
|
||||
throw new DriverStop(`expected ${phrase}; nothing further was dispatched`)
|
||||
typedPhrases.set(phrase, answer)
|
||||
log(`operator typed ${phrase}`)
|
||||
return answer
|
||||
}
|
||||
|
||||
// The digest the publish run pushed: the registry's digest for the commit's tag, which the push
|
||||
// line in the run's own log must name too.
|
||||
async function publishedDigest(runId) {
|
||||
const view = viewRun(runId)
|
||||
if (view.workflowName !== WORKFLOWS.publish.name || view.conclusion !== 'success') {
|
||||
throw new DriverStop(`${runUrl(runId)} is not a successful ${WORKFLOWS.publish.file} run`)
|
||||
}
|
||||
if (view.headSha !== config.commit) {
|
||||
throw new DriverStop(
|
||||
`publish ${runUrl(runId)} built ${view.headSha}, not the reviewed ${config.commit}; do not deploy it`
|
||||
)
|
||||
}
|
||||
const tag = `${IMAGE_REPOSITORY}:sha-${config.commit}`
|
||||
const digest = command(
|
||||
'gcloud',
|
||||
words(
|
||||
`artifacts docker images describe ${tag} --project ${PROJECT} --format=value(image_summary.digest)`
|
||||
)
|
||||
).trim()
|
||||
requireDigest(digest, 'registry digest of the published tag')
|
||||
if (!logConfirmsPublishedDigest(await runLog(runId), config.commit, digest)) {
|
||||
throw new DriverStop(
|
||||
`registry digest ${digest} is not the digest ${runUrl(runId)} pushed; the tag moved`
|
||||
)
|
||||
}
|
||||
return digest
|
||||
}
|
||||
|
||||
async function lastKnownRehomeGeneration() {
|
||||
const runs = ghJson(
|
||||
words(
|
||||
`run list -R ${REPOSITORY} --workflow ${WORKFLOWS.rehome.file} --status completed --limit ${REHOME_HISTORY_RUNS} --json databaseId`
|
||||
)
|
||||
)
|
||||
for (const run of runs) {
|
||||
const result = await rehomeResult(run.databaseId)
|
||||
if (result) {
|
||||
log(
|
||||
`rehome generation candidate ${result.control.generation} from ${runUrl(run.databaseId)}`
|
||||
)
|
||||
return result.control.generation
|
||||
}
|
||||
}
|
||||
throw new DriverStop(
|
||||
'no recent rehome run printed the control generation; pass --rehome-generation'
|
||||
)
|
||||
}
|
||||
|
||||
// Values every input builder reads. A dry run shows what is not known yet as `<placeholder>`,
|
||||
// which validateDispatchInputs refuses, so a placeholder can never be dispatched.
|
||||
function view() {
|
||||
const unknown = (label) => `<${label}>`
|
||||
const placeholder = (label) => new Proxy({}, { get: () => unknown(label) })
|
||||
const fromAdmission = [unknown('from admission inspect')]
|
||||
const control = known.control
|
||||
return {
|
||||
director: known.director,
|
||||
selector: known.selector ?? {
|
||||
generation: fromAdmission[0],
|
||||
membership: new Proxy({}, { get: () => fromAdmission })
|
||||
},
|
||||
control: control ?? placeholder('inspected'),
|
||||
pausedGeneration:
|
||||
owned?.generation ??
|
||||
(control?.enabled === false
|
||||
? control.generation
|
||||
: unknown('rehome generation after pause')),
|
||||
published: published?.digest ?? unknown('published digest'),
|
||||
rollbackDigest: rollbackPoint?.digest ?? known.director?.servingDigest,
|
||||
monitor: placeholder('monitor run'),
|
||||
afterDeploy: placeholder('digest read from gcloud after the deploy'),
|
||||
notBefore: config.dryRun ? unknown('now, epoch ms') : Math.floor(deps.now() / 1000) * 1000,
|
||||
confirmation: (phrase) =>
|
||||
typedPhrases.has(phrase) ? phrase : unknown(`operator types ${phrase}`)
|
||||
}
|
||||
}
|
||||
|
||||
const deployPending = () => known.director.servingDigest !== published?.digest
|
||||
const pendingWaves = () =>
|
||||
config.configure.filter((wave) => wave.cells.some((cell) => !known.director.cells.has(cell)))
|
||||
const work = () => deployPending() || pendingWaves().length > 0
|
||||
|
||||
function count5xx(fromMs, toMs) {
|
||||
const filter = `${DIRECTOR_5XX_FILTER} AND timestamp>="${timestamp(fromMs)}" AND timestamp<"${timestamp(toMs)}"`
|
||||
return command('gcloud', [
|
||||
'logging',
|
||||
'read',
|
||||
filter,
|
||||
...words(`--project ${PROJECT} --limit ${LOG_COUNT_LIMIT} --format=value(timestamp)`)
|
||||
])
|
||||
.split('\n')
|
||||
.filter(Boolean).length
|
||||
}
|
||||
|
||||
// Director 5xx over SOAK_MS on the new image against the same span before its revision existed.
|
||||
// The window starts at the traffic switch when this run deployed, otherwise now, so a re-run
|
||||
// after a tripped soak judges fresh traffic.
|
||||
async function soak(deployedAt) {
|
||||
const created = Date.parse(known.director.servingCreatedAt)
|
||||
const start =
|
||||
deployedAt === undefined ? deps.now() : Math.max(deployedAt - TRAFFIC_SWITCH_LEAD_MS, created)
|
||||
const end = start + SOAK_MS
|
||||
const readAt = end + LOG_INGESTION_LAG_MS
|
||||
if (deps.now() < readAt) {
|
||||
log(
|
||||
`soak: watching the new director until ${timestamp(end)}, reading at ${timestamp(readAt)}`
|
||||
)
|
||||
await deps.sleep(readAt - deps.now())
|
||||
}
|
||||
const before = count5xx(created - SOAK_MS, created)
|
||||
const after = count5xx(start, end)
|
||||
log(
|
||||
`soak: director 5xx ${after} in ${SOAK_MS / 60_000} min on the new image, ${before} before it`
|
||||
)
|
||||
if (after >= LOG_COUNT_LIMIT || after > 2 * before + 25) {
|
||||
throw new DriverStop(
|
||||
`director 5xx rose from ${before} to ${after} after the deploy; investigate before configuring cells`
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
// The same audited check the enable job runs, on the same sealed files.
|
||||
async function monitorCompletedAt(run) {
|
||||
const incidentId = `relay-${run.runId}-dry-run`
|
||||
return await withArtifact(
|
||||
run.runId,
|
||||
`relay-monitor-dry-run-${run.runId}-${run.attempt}`,
|
||||
async (directory) => {
|
||||
try {
|
||||
const argv = Object.entries({
|
||||
directory,
|
||||
'incident-id': incidentId,
|
||||
'run-id': run.runId,
|
||||
'run-attempt': run.attempt,
|
||||
'commit-sha': run.headSha,
|
||||
mode: 'dry-run',
|
||||
'required-migration-policy': 'strict'
|
||||
}).flatMap(([key, value]) => [`--${key}`, String(value)])
|
||||
const { state } = await verifyDryRunAuthority(argv, deps.now)
|
||||
log(`monitor GREEN, completed ${state.completedAt}`)
|
||||
return state.completedAt
|
||||
} catch (error) {
|
||||
// The freeze reason lives only in the sealed state, never in the run log.
|
||||
const sealed = (() => {
|
||||
try {
|
||||
return JSON.parse(readFileSync(join(directory, `${incidentId}.state.json`), 'utf8'))
|
||||
} catch {
|
||||
return {}
|
||||
}
|
||||
})()
|
||||
const failures = (sealed.failures ?? []).map((failure) =>
|
||||
[
|
||||
failure.source,
|
||||
failure.code,
|
||||
failure.signal,
|
||||
failure.observed,
|
||||
failure.threshold
|
||||
].join(' ')
|
||||
)
|
||||
throw new DriverStop(
|
||||
[
|
||||
`monitor ${run.url} is not usable: ${error.message}`,
|
||||
`frozenAt ${sealed.frozenAt}`,
|
||||
...failures
|
||||
].join('\n ')
|
||||
)
|
||||
}
|
||||
}
|
||||
)
|
||||
}
|
||||
|
||||
async function preflight() {
|
||||
log('preflight: relay workflow lane, target commit, serving director')
|
||||
if (config.dryRun) {
|
||||
const active = activeRuns()
|
||||
if (active.length > 0)
|
||||
log(
|
||||
`WARNING relay workflows are in flight; a real run would stop:\n ${active.join('\n ')}`
|
||||
)
|
||||
} else {
|
||||
requireQuietLane()
|
||||
}
|
||||
if (published) {
|
||||
published.digest = await publishedDigest(published.runId)
|
||||
} else {
|
||||
const main = gh(['api', `repos/${REPOSITORY}/commits/${WORKFLOW_REF}`, '--jq', '.sha']).trim()
|
||||
if (main !== config.commit) {
|
||||
throw new DriverStop(
|
||||
`${WORKFLOW_REF} is at ${main}, not the reviewed ${config.commit}; review the difference and run with --commit ${main}`
|
||||
)
|
||||
}
|
||||
}
|
||||
known.director = readDirector()
|
||||
if (known.director.servingDigest !== published?.digest) {
|
||||
rollbackPoint = {
|
||||
revision: known.director.servingRevision,
|
||||
digest: known.director.servingDigest
|
||||
}
|
||||
}
|
||||
log(describeDirector(known.director))
|
||||
let claim
|
||||
if (config.pauseRun) {
|
||||
// Only this operator's own rehome-control run can prove a pause belongs to this driver.
|
||||
const claimed = ghJson(
|
||||
words(
|
||||
`api repos/${REPOSITORY}/actions/runs/${config.pauseRun} --jq {path,actor:.triggering_actor.login}`
|
||||
)
|
||||
)
|
||||
if (
|
||||
!claimed.path?.split('@')[0].endsWith(`/${WORKFLOWS.rehome.file}`) ||
|
||||
claimed.actor !== login
|
||||
) {
|
||||
throw new DriverStop(
|
||||
`${runUrl(config.pauseRun)} is not a ${WORKFLOWS.rehome.file} run by ${login}, so it cannot prove a pause is this driver's`
|
||||
)
|
||||
}
|
||||
const result = await rehomeResult(config.pauseRun)
|
||||
claim = result && pausedGeneration(result)
|
||||
if (claim === undefined) {
|
||||
throw new DriverStop(
|
||||
`${runUrl(config.pauseRun)} printed no pause of its own, so it cannot prove a pause is this driver's`
|
||||
)
|
||||
}
|
||||
}
|
||||
// A director safety pause moves the generation without a run; the inspect then fails closed.
|
||||
const generation = config.rehomeGeneration ?? (await lastKnownRehomeGeneration())
|
||||
if (config.dryRun) {
|
||||
known.control = undefined
|
||||
return generation
|
||||
}
|
||||
const [admissionStep, inspectStep] = readSteps(generation)
|
||||
const admission = await dispatch(admissionStep)
|
||||
requireSuccess(admission)
|
||||
known.selector = await withArtifact(
|
||||
admission.runId,
|
||||
`relay-asia-admission-result-${admission.runId}-${admission.attempt}`,
|
||||
(directory) =>
|
||||
admissionInspectResult(JSON.parse(readFileSync(join(directory, 'result.json'), 'utf8')))
|
||||
)
|
||||
const inspect = await dispatch(inspectStep)
|
||||
if (inspect.conclusion !== 'success') {
|
||||
throw new DriverStop(
|
||||
`rehome inspect at generation ${generation} failed (${inspect.url}): the generation, selector or a digest moved. A director safety pause bumps the generation. Read the log, then pass --rehome-generation`
|
||||
)
|
||||
}
|
||||
known.control = (await rehomeResult(inspect.runId))?.control
|
||||
if (!known.control) throw new DriverStop(`rehome inspect ${inspect.url} printed no control`)
|
||||
const control = known.control
|
||||
log(
|
||||
`rehome: generation ${control.generation}, ${control.enabled ? 'ENABLED' : 'disabled'}; selector generation ${known.selector.generation}`
|
||||
)
|
||||
if (claim !== undefined) {
|
||||
if (!control.enabled && control.generation === claim)
|
||||
owned = { generation: claim, runId: config.pauseRun }
|
||||
else if (control.enabled && control.generation === claim + 1) enabled = true
|
||||
else
|
||||
throw new DriverStop(
|
||||
`rehome is generation ${control.generation} ${control.enabled ? 'enabled' : 'paused'}, not the pause ${runUrl(config.pauseRun)} made at ${claim}. Something else changed it, such as a director safety pause; this driver will not touch it`
|
||||
)
|
||||
} else if (control.enabled && config.leaveRehomePaused) {
|
||||
throw new DriverStop(
|
||||
'rehome is enabled; --leave-rehome-paused accepts only a paused switch. Drop it'
|
||||
)
|
||||
} else if (!control.enabled && !config.leaveRehomePaused) {
|
||||
throw new DriverStop(
|
||||
`rehome is paused at generation ${control.generation}, and this run did not pause it. If an earlier run of this driver did, re-run with the --pause-run its log printed. Otherwise pass --leave-rehome-paused to deploy and leave it paused`
|
||||
)
|
||||
}
|
||||
if (control.hostCooldownMs === undefined && (control.enabled || owned)) {
|
||||
throw new DriverStop(
|
||||
'the director reports no per-host rehome cooldown, so enable would refuse; not pausing'
|
||||
)
|
||||
}
|
||||
return generation
|
||||
}
|
||||
|
||||
// The two read-only inspects every real run starts with; never resumed, always run fresh.
|
||||
function readSteps(generation) {
|
||||
return [
|
||||
{
|
||||
name: 'preflight-admission',
|
||||
workflow: WORKFLOWS.admission,
|
||||
inputs: (v) => admissionInspectInputs(v.director.servingDigest)
|
||||
},
|
||||
{
|
||||
name: 'preflight-rehome',
|
||||
workflow: WORKFLOWS.rehome,
|
||||
inputs: (v) => rehomeInspectInputs({ ...v, controlGeneration: generation })
|
||||
}
|
||||
]
|
||||
}
|
||||
|
||||
// The one ordered plan. `when` reads live state, so a re-run skips what is already done; a dry run
|
||||
// prints every step whose need it cannot know yet.
|
||||
function plan() {
|
||||
let deployedAt
|
||||
let verified
|
||||
let monitor
|
||||
return [
|
||||
{
|
||||
name: 'publish',
|
||||
workflow: WORKFLOWS.publish,
|
||||
inputs: publishInputs,
|
||||
when: () => !published,
|
||||
run: async (step) => {
|
||||
const run = await dispatch(step)
|
||||
requireSuccess(run)
|
||||
published = { runId: run.runId }
|
||||
published.digest = await publishedDigest(run.runId)
|
||||
log(`published ${IMAGE_REPOSITORY}@${published.digest}`)
|
||||
}
|
||||
},
|
||||
{
|
||||
name: 'pause',
|
||||
workflow: WORKFLOWS.rehome,
|
||||
changesRehome: true,
|
||||
inputs: (v) =>
|
||||
rehomePauseInputs({ ...v, confirmation: v.confirmation('PAUSE_REGIONAL_REHOMING') }),
|
||||
when: () => known.control?.enabled === true && !owned && work(),
|
||||
run: async (step) => {
|
||||
await typed('PAUSE_REGIONAL_REHOMING')
|
||||
const run = await dispatch(step)
|
||||
const result = await rehomeResult(run.runId)
|
||||
const generation = result && pausedGeneration(result)
|
||||
if (generation !== known.control.generation + 1) {
|
||||
uncertain.kind = 'unconfirmed'
|
||||
throw new DriverStop(
|
||||
`${run.url} (${run.conclusion}) printed no pause at generation ${known.control.generation + 1}`
|
||||
)
|
||||
}
|
||||
owned = { generation, runId: run.runId }
|
||||
uncertain = undefined
|
||||
log(
|
||||
`REHOME PAUSED at generation ${generation}. To finish from here after any stop: ${rerunCommand()}`
|
||||
)
|
||||
}
|
||||
},
|
||||
{
|
||||
name: 'deploy',
|
||||
workflow: WORKFLOWS.director,
|
||||
inputs: (v) =>
|
||||
directorDeployInputs({
|
||||
imageDigest: v.published,
|
||||
predecessorDigest: v.rollbackDigest,
|
||||
rehomeGeneration: v.pausedGeneration
|
||||
}),
|
||||
when: () => !published || deployPending(),
|
||||
run: async (step) => {
|
||||
requireSuccess(await dispatch(step))
|
||||
deployedAt = deps.now()
|
||||
known.director = readDirector()
|
||||
if (deployPending())
|
||||
throw new DriverStop(
|
||||
`deploy: ${describeDirector(known.director)}, not ${published.digest}`
|
||||
)
|
||||
}
|
||||
},
|
||||
...config.configure.slice(0, 1).map(() => ({
|
||||
name: 'soak',
|
||||
description: `wait ${SOAK_MS / 60_000} min, then compare director 5xx`,
|
||||
// A wave already configured means an earlier run passed the soak on this image.
|
||||
when: () => pendingWaves().length === config.configure.length,
|
||||
run: () => soak(deployedAt)
|
||||
})),
|
||||
...config.configure.map((wave) => ({
|
||||
name: `configure:${wave.cells.join(',')}`,
|
||||
workflow: WORKFLOWS.admission,
|
||||
inputs: (v) =>
|
||||
configureInputs({
|
||||
...wave,
|
||||
directorDigest: v.published,
|
||||
selectorGeneration: v.selector.generation,
|
||||
confirmation: v.confirmation('CONFIGURE_ASIA_DIRECTOR')
|
||||
}),
|
||||
when: () => pendingWaves().includes(wave),
|
||||
run: async (step) => {
|
||||
await typed('CONFIGURE_ASIA_DIRECTOR')
|
||||
requireSuccess(await dispatch(step))
|
||||
known.director = readDirector()
|
||||
if (deployPending())
|
||||
throw new DriverStop(
|
||||
`${step.name}: ${describeDirector(known.director)}, not ${published.digest}`
|
||||
)
|
||||
}
|
||||
})),
|
||||
// Inspect binds the exact serving and rollback digests, so a wrong one fails here, read-only,
|
||||
// before 15 minutes of monitor evidence is spent on it.
|
||||
{
|
||||
name: 'verify-identities',
|
||||
workflow: WORKFLOWS.rehome,
|
||||
inputs: (v) =>
|
||||
rehomeInspectInputs({
|
||||
director: verified ?? v.afterDeploy,
|
||||
selector: v.selector,
|
||||
controlGeneration: v.pausedGeneration
|
||||
}),
|
||||
when: () => Boolean(owned),
|
||||
run: async (step) => {
|
||||
verified = readDirector()
|
||||
const run = await dispatch(step)
|
||||
const control = (await rehomeResult(run.runId))?.control
|
||||
if (
|
||||
run.conclusion !== 'success' ||
|
||||
control?.enabled !== false ||
|
||||
control.generation !== owned.generation
|
||||
) {
|
||||
throw new DriverStop(
|
||||
`verify-identities ${run.url} (${run.conclusion}) did not find rehome paused at ${owned.generation} with ${describeDirector(verified)}`
|
||||
)
|
||||
}
|
||||
}
|
||||
},
|
||||
{
|
||||
name: 'monitor',
|
||||
workflow: WORKFLOWS.monitor,
|
||||
inputs: (v) => monitorDryRunInputs(v.selector),
|
||||
when: () => Boolean(owned),
|
||||
run: async (step) => {
|
||||
// The operator arms the enable before the 15-minute watch; it still dispatches only on green.
|
||||
await typed(
|
||||
'ENABLE_REGIONAL_REHOMING',
|
||||
` (this arms an automatic enable, sent about 17 min from now and only if the monitor is green and its evidence is at most ${MONITOR_MAX_AGE_AT_ENABLE_MS / 1000} s old)`
|
||||
)
|
||||
const run = await dispatch(step)
|
||||
requireSuccess(run)
|
||||
monitor = {
|
||||
runId: run.runId,
|
||||
attempt: run.attempt,
|
||||
completedAt: await monitorCompletedAt(run)
|
||||
}
|
||||
}
|
||||
},
|
||||
{
|
||||
name: 'enable',
|
||||
workflow: WORKFLOWS.rehome,
|
||||
changesRehome: true,
|
||||
inputs: (v) =>
|
||||
rehomeEnableInputs({
|
||||
...v,
|
||||
director: verified ?? v.afterDeploy,
|
||||
controlGeneration: v.pausedGeneration,
|
||||
monitor: monitor ?? v.monitor,
|
||||
confirmation: v.confirmation('ENABLE_REGIONAL_REHOMING')
|
||||
}),
|
||||
when: () => Boolean(owned),
|
||||
run: async (step) => {
|
||||
const ageMs = deps.now() - Date.parse(monitor.completedAt)
|
||||
if (ageMs > MONITOR_MAX_AGE_AT_ENABLE_MS) {
|
||||
throw new DriverStop(
|
||||
`monitor evidence is ${Math.round(ageMs / 1000)} s old, past the ${MONITOR_MAX_AGE_AT_ENABLE_MS / 1000} s budget; re-run for a fresh monitor`
|
||||
)
|
||||
}
|
||||
const now = readDirector()
|
||||
if (
|
||||
now.servingDigest !== verified.servingDigest ||
|
||||
now.rollbackDigest !== verified.rollbackDigest
|
||||
) {
|
||||
throw new DriverStop(
|
||||
`the director changed after its digests were verified: ${describeDirector(now)}`
|
||||
)
|
||||
}
|
||||
if (known.control.ratePerMinute !== 10)
|
||||
log(
|
||||
`enable starts at the job's fixed 10 hosts/min (was ${known.control.ratePerMinute})`
|
||||
)
|
||||
const run = await dispatch(step)
|
||||
const result = await rehomeResult(run.runId)
|
||||
if (result) known.control = result.control
|
||||
// A failed run is never an enable, whatever it printed last.
|
||||
if (
|
||||
run.conclusion === 'success' &&
|
||||
result?.mode === 'enable' &&
|
||||
result.control.enabled &&
|
||||
result.control.generation === owned.generation + 1
|
||||
) {
|
||||
owned = undefined
|
||||
uncertain = undefined
|
||||
enabled = true
|
||||
log(`REHOME RE-ENABLED at generation ${result.control.generation}`)
|
||||
return
|
||||
}
|
||||
if (result?.mode !== 'recover-enable') {
|
||||
uncertain.kind = 'unconfirmed'
|
||||
throw new DriverStop(`${run.url} (${run.conclusion}) printed no enable result`)
|
||||
}
|
||||
// Recovery that disabled rehome itself is this driver's pause too. Recovery that found it
|
||||
// already disabled at another generation found someone else's pause, such as a director
|
||||
// safety pause: this driver gives up ownership and will never lift it.
|
||||
uncertain = undefined
|
||||
const recovered = pausedGeneration(result)
|
||||
if (recovered !== undefined) owned = { generation: recovered, runId: run.runId }
|
||||
else if (result.control.generation !== owned.generation) owned = undefined
|
||||
throw new DriverStop(
|
||||
`enable ${run.url} failed; the job's recovery left rehome paused at generation ${result.control.generation}`
|
||||
)
|
||||
}
|
||||
}
|
||||
]
|
||||
}
|
||||
|
||||
function rerunCommand(pauseRun = owned?.runId) {
|
||||
return [
|
||||
'node dev/scripts/drive-relay-director-deploy.mjs',
|
||||
`--commit ${config.commit}`,
|
||||
...(published?.digest ? [`--publish-run ${published.runId}`] : []),
|
||||
...(pauseRun ? [`--pause-run ${pauseRun}`] : []),
|
||||
...(config.leaveRehomePaused ? ['--leave-rehome-paused'] : []),
|
||||
...config.configure.map(
|
||||
(wave) => `--configure ${wave.cells.join(',')}=${wave.cellImageDigest}`
|
||||
)
|
||||
].join(' ')
|
||||
}
|
||||
|
||||
function stopReport(reason) {
|
||||
const lines = [`STOPPED: ${reason}`, '', 'State now:']
|
||||
const workflowPage = `https://github.com/${REPOSITORY}/actions/workflows/${WORKFLOWS.rehome.file}`
|
||||
const where = uncertain?.url ?? `(gh printed no URL; find it at ${workflowPage})`
|
||||
if (uncertain?.name === 'pause') {
|
||||
lines.push(
|
||||
uncertain.kind === 'changing'
|
||||
? `- *** REHOME IS CHANGING: the pause run ${where} was dispatched and applies on its own, if it has not already. ***`
|
||||
: `- *** PAUSE UNCONFIRMED: ${where} printed no pause. Rehome MAY BE PAUSED. ***`,
|
||||
' Inspect rehome before walking away.'
|
||||
)
|
||||
if (uncertain.runId) {
|
||||
lines.push(
|
||||
` If it paused rehome at generation ${known.control.generation + 1}, finish with: cd cloud && ${rerunCommand(uncertain.runId)}`
|
||||
)
|
||||
}
|
||||
} else if (uncertain?.name === 'enable') {
|
||||
lines.push(
|
||||
uncertain.kind === 'changing'
|
||||
? `- *** REHOME IS CHANGING: the enable run ${where} applies on its own: it enables rehome, or disables it again if it fails. ***`
|
||||
: `- *** ENABLE UNCONFIRMED: ${where} did not confirm the enable. Rehome may be enabled, or paused. ***`,
|
||||
` If it enabled rehome, or failed before applying, finish with: cd cloud && ${rerunCommand(owned.runId)}`
|
||||
)
|
||||
if (uncertain.runId) {
|
||||
lines.push(
|
||||
` If its recovery paused rehome again, finish with: cd cloud && ${rerunCommand(uncertain.runId)}`
|
||||
)
|
||||
}
|
||||
} else if (owned) {
|
||||
lines.push(
|
||||
`- *** REHOME IS PAUSED by this driver at generation ${owned.generation} (${runUrl(owned.runId)}). It stays paused until the re-run below enables it. ***`
|
||||
)
|
||||
} else if (known.control) {
|
||||
const state =
|
||||
enabled || known.control.enabled
|
||||
? 'enabled'
|
||||
: 'PAUSED, not by this driver; it will not be re-enabled here'
|
||||
lines.push(`- rehome: generation ${known.control.generation} as last read, ${state}`)
|
||||
} else {
|
||||
lines.push('- rehome: not read; this run did not change it')
|
||||
}
|
||||
try {
|
||||
lines.push(`- director now: ${describeDirector(readDirector())}`)
|
||||
} catch (error) {
|
||||
lines.push(`- director: could not re-read (${error.message})`)
|
||||
}
|
||||
if (published?.digest)
|
||||
lines.push(`- published: ${published.digest} (${runUrl(published.runId)})`)
|
||||
if (rollbackPoint) {
|
||||
lines.push(
|
||||
`- rollback point: ${rollbackPoint.revision} ${rollbackPoint.digest}. To undo the deploy, while rehome is paused:`
|
||||
)
|
||||
lines.push(
|
||||
` ${ghCommand(WORKFLOWS.director, directorDeployInputs({ imageDigest: rollbackPoint.digest, predecessorDigest: rollbackPoint.digest, rehomeGeneration: owned?.generation ?? '<paused generation>' }))}`
|
||||
)
|
||||
}
|
||||
if (!config.dryRun && !uncertain) {
|
||||
lines.push(
|
||||
'',
|
||||
`Re-run to finish from here (it re-reads everything and skips what is done):`,
|
||||
` cd cloud && ${rerunCommand()}`
|
||||
)
|
||||
}
|
||||
return lines.join('\n')
|
||||
}
|
||||
|
||||
function ghCommand(workflow, inputs) {
|
||||
const fields = Object.entries(inputs).map(([key, value]) => `-f ${key}=${value}`)
|
||||
return `gh workflow run ${workflow.file} -R ${REPOSITORY} --ref ${WORKFLOW_REF} ${fields.join(' ')}`
|
||||
}
|
||||
|
||||
async function run() {
|
||||
const stamp = timestamp(deps.now()).replace(/[:.]/g, '-')
|
||||
mkdirSync(config.logDirectory, { recursive: true, mode: 0o700 })
|
||||
logPath = join(
|
||||
config.logDirectory,
|
||||
`${stamp}-${config.commit.slice(0, 12)}${config.dryRun ? '-dry-run' : ''}.log`
|
||||
)
|
||||
log(`${config.dryRun ? 'dry run' : 'start'}: ${rerunCommand()}`)
|
||||
try {
|
||||
login = gh(['api', 'user', '--jq', '.login']).trim()
|
||||
const generation = await preflight()
|
||||
const steps = plan()
|
||||
for (const step of [...(config.dryRun ? readSteps(generation) : []), ...steps]) {
|
||||
const needed = config.dryRun || step.when() ? '' : ' (done or not needed)'
|
||||
const detail = step.workflow
|
||||
? `${step.workflow.file} ${JSON.stringify(step.inputs({ ...view(), control: { ...view().control, generation } }))}`
|
||||
: step.description
|
||||
log(`plan ${step.name}${needed}: ${detail}`)
|
||||
}
|
||||
if (config.dryRun) {
|
||||
log('dry run: nothing dispatched')
|
||||
return { dryRun: true }
|
||||
}
|
||||
if (steps.some((step) => step.when())) await typed(`DEPLOY ${config.commit.slice(0, 12)}`)
|
||||
for (const step of steps) if (step.when()) await step.run(step)
|
||||
const rehome =
|
||||
enabled || known.control.enabled ? 'enabled' : 'left paused (--leave-rehome-paused)'
|
||||
log(
|
||||
`DONE: ${describeDirector(known.director)}; rehome ${rehome}${rollbackPoint ? `; rollback point ${rollbackPoint.revision} ${rollbackPoint.digest}` : ''}`
|
||||
)
|
||||
return { done: true, logPath }
|
||||
} catch (error) {
|
||||
for (const line of stopReport(error.message).split('\n')) log(line)
|
||||
throw Object.assign(error instanceof DriverStop ? error : new DriverStop(error.message), {
|
||||
reported: true,
|
||||
logPath
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// Ctrl-C, SIGTERM or a closed terminal still leaves the operator the live state and the re-run.
|
||||
function interrupt(signal) {
|
||||
if (logPath) for (const line of stopReport(`interrupted by ${signal}`).split('\n')) log(line)
|
||||
}
|
||||
|
||||
return { run, interrupt }
|
||||
}
|
||||
|
||||
function defaultDependencies() {
|
||||
return {
|
||||
run: (program, args, input) =>
|
||||
spawnSync(program, args, {
|
||||
encoding: 'utf8',
|
||||
input,
|
||||
maxBuffer: 256 * 1024 * 1024,
|
||||
stdio: ['pipe', 'pipe', 'pipe']
|
||||
}),
|
||||
stream: (program, args) =>
|
||||
spawnSync(program, args, { stdio: ['ignore', 'inherit', 'inherit'] }).status,
|
||||
now: () => Date.now(),
|
||||
sleep: (ms) => new Promise((resolveSleep) => setTimeout(resolveSleep, ms)),
|
||||
print: (line) => process.stdout.write(`${line}\n`),
|
||||
prompt: async (question) => {
|
||||
const reader = createInterface({ input: process.stdin, output: process.stdout })
|
||||
try {
|
||||
return await reader.question(question)
|
||||
} finally {
|
||||
reader.close()
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
export async function main(argv = process.argv.slice(2), deps = defaultDependencies()) {
|
||||
const driver = createDriver(parseDriverArguments(argv), deps)
|
||||
for (const [signal, code] of [
|
||||
['SIGINT', 130],
|
||||
['SIGTERM', 143],
|
||||
['SIGHUP', 129]
|
||||
]) {
|
||||
process.once(signal, () => {
|
||||
driver.interrupt(signal)
|
||||
process.exit(code)
|
||||
})
|
||||
}
|
||||
return await driver.run()
|
||||
}
|
||||
|
||||
if (process.argv[1] && import.meta.url === pathToFileURL(process.argv[1]).href) {
|
||||
main().catch((error) => {
|
||||
if (!error.reported)
|
||||
process.stderr.write(`${error instanceof Error ? error.message : String(error)}\n`)
|
||||
process.exitCode = 1
|
||||
})
|
||||
}
|
||||
File diff suppressed because it is too large
Load Diff
@@ -0,0 +1,333 @@
|
||||
// Pure pieces of the director deploy driver: the exact inputs each existing workflow receives,
|
||||
// and the parsers that read a run's result back. No process, network, or clock access here.
|
||||
|
||||
import {
|
||||
RELAY_GITHUB_REPOSITORY,
|
||||
RELAY_WORKFLOW_FILE_PREFIX,
|
||||
relayWorkflowFile
|
||||
} from './relay-repository.mjs'
|
||||
|
||||
export const REPOSITORY = RELAY_GITHUB_REPOSITORY
|
||||
export const WORKFLOW_REF = 'main'
|
||||
export const PROJECT = 'onorca-cloud'
|
||||
export const REGION = 'us-central1'
|
||||
export const DIRECTOR_SERVICE = 'orca-cloud-relay'
|
||||
export const ROLLBACK_TAG = 'selector-rollback'
|
||||
export const IMAGE_REPOSITORY = 'us-central1-docker.pkg.dev/onorca-cloud/orca-cloud/relay'
|
||||
// Any reviewed registration wave is accepted by admission inspect; the launch wave never changes.
|
||||
const ADMISSION_INSPECT_CELLS = 'production-gce-c27,production-gce-c28,production-gce-c29'
|
||||
|
||||
const workflow = (name, title) => ({ file: relayWorkflowFile(name), name: title })
|
||||
export const WORKFLOWS = {
|
||||
publish: workflow('publish-relay-production.yml', 'Publish Relay Production Image'),
|
||||
director: workflow('deploy-relay-production-director.yml', 'Deploy Relay Production Director'),
|
||||
rehome: workflow('operate-relay-production-rehome.yml', 'Operate Relay Production Rehome'),
|
||||
admission: workflow('operate-relay-asia-admission.yml', 'Operate Relay Asia Admission'),
|
||||
monitor: workflow('monitor-relay-production.yml', 'Monitor Relay Production')
|
||||
}
|
||||
|
||||
// Read-only or unrelated to the production rollout lane, so they never block a deploy.
|
||||
const NON_BLOCKING_WORKFLOWS = new Set(
|
||||
['verify.yml', 'monitor-relay-clock-skew.yml'].map(relayWorkflowFile)
|
||||
)
|
||||
|
||||
const DIGEST = /^sha256:[a-f0-9]{64}$/
|
||||
const COMMIT = /^[a-f0-9]{40}$/
|
||||
const PRODUCTION_CELL = /^production-gce-c[1-9][0-9]?$/
|
||||
const MEMBERSHIP_KEYS = ['existingOnly', 'migrationOnly', 'general']
|
||||
|
||||
export function requireDigest(value, label) {
|
||||
if (typeof value !== 'string' || !DIGEST.test(value))
|
||||
throw new Error(`${label} is not an immutable sha256 digest`)
|
||||
return value
|
||||
}
|
||||
|
||||
export function requireCommit(value, label) {
|
||||
if (typeof value !== 'string' || !COMMIT.test(value))
|
||||
throw new Error(`${label} is not a full commit SHA`)
|
||||
return value
|
||||
}
|
||||
|
||||
function requireGeneration(value, label) {
|
||||
if (!Number.isSafeInteger(value) || value < 0) throw new Error(`${label} is invalid`)
|
||||
return value
|
||||
}
|
||||
|
||||
export function blocksDeploy(path) {
|
||||
// A `@ref` suffix never appears on run paths today; stripping it keeps the check fail-closed.
|
||||
const file = String(path ?? '')
|
||||
.split('@')[0]
|
||||
.split('/')
|
||||
.at(-1)
|
||||
return (
|
||||
file.startsWith(RELAY_WORKFLOW_FILE_PREFIX) &&
|
||||
file.endsWith('.yml') &&
|
||||
!NON_BLOCKING_WORKFLOWS.has(file)
|
||||
)
|
||||
}
|
||||
|
||||
export function membershipInput(cells) {
|
||||
return cells.length === 0 ? 'none' : [...cells].sort().join(',')
|
||||
}
|
||||
|
||||
export function parseSelector(value, label) {
|
||||
const selector = {
|
||||
generation: requireGeneration(value?.generation, `${label} generation`),
|
||||
membership: {}
|
||||
}
|
||||
for (const key of MEMBERSHIP_KEYS) {
|
||||
const cells = value?.membership?.[key]
|
||||
if (!Array.isArray(cells) || cells.some((cell) => !PRODUCTION_CELL.test(cell))) {
|
||||
throw new Error(`${label} ${key} membership is invalid`)
|
||||
}
|
||||
selector.membership[key] = [...cells].sort()
|
||||
}
|
||||
const all = MEMBERSHIP_KEYS.flatMap((key) => selector.membership[key])
|
||||
if (new Set(all).size !== all.length) throw new Error(`${label} membership has duplicates`)
|
||||
return selector
|
||||
}
|
||||
|
||||
function selectorInputs(selector) {
|
||||
return {
|
||||
'expected-selector-generation': String(selector.generation),
|
||||
'expected-existing-only-cells': membershipInput(selector.membership.existingOnly),
|
||||
'expected-migration-only-cells': membershipInput(selector.membership.migrationOnly),
|
||||
'expected-general-cells': membershipInput(selector.membership.general)
|
||||
}
|
||||
}
|
||||
|
||||
export function parseControl(value, label) {
|
||||
const control = value ?? {}
|
||||
const integers = [
|
||||
'generation',
|
||||
'notBefore',
|
||||
'ratePerMinute',
|
||||
'preferenceMaxAgeMs',
|
||||
'drainGraceMs'
|
||||
]
|
||||
if (
|
||||
typeof control.enabled !== 'boolean' ||
|
||||
integers.some((key) => !Number.isSafeInteger(control[key]))
|
||||
) {
|
||||
throw new Error(`${label} is not a complete regional rehome control`)
|
||||
}
|
||||
if (control.hostCooldownMs !== undefined && !Number.isSafeInteger(control.hostCooldownMs)) {
|
||||
throw new Error(`${label} host cooldown is invalid`)
|
||||
}
|
||||
return Object.fromEntries(
|
||||
[...integers, 'enabled', 'hostCooldownMs'].map((key) => [key, control[key]])
|
||||
)
|
||||
}
|
||||
|
||||
const INTEGER_INPUT = /^(0|[1-9][0-9]*)$/
|
||||
|
||||
/**
|
||||
* The last check before `gh workflow run`. Builders also render dry-run placeholders such as
|
||||
* `<published digest>`; this guarantees none of them, or a malformed digest, is ever dispatched.
|
||||
*/
|
||||
export function validateDispatchInputs(inputs) {
|
||||
for (const [key, value] of Object.entries(inputs)) {
|
||||
if (typeof value !== 'string' || value.startsWith('<'))
|
||||
throw new Error(`input ${key} is not resolved`)
|
||||
if (key.endsWith('digest') && !DIGEST.test(value))
|
||||
throw new Error(`input ${key} is not a sha256 digest`)
|
||||
if (
|
||||
(key.includes('generation') ||
|
||||
['not-before', 'monitor-run-id', 'monitor-run-attempt'].includes(key)) &&
|
||||
!INTEGER_INPUT.test(value)
|
||||
) {
|
||||
throw new Error(`input ${key} is not an integer`)
|
||||
}
|
||||
}
|
||||
return inputs
|
||||
}
|
||||
|
||||
export function admissionInspectInputs(servingDigest) {
|
||||
return {
|
||||
environment: 'production',
|
||||
mode: 'inspect',
|
||||
'cell-ids': ADMISSION_INSPECT_CELLS,
|
||||
'image-digest': servingDigest
|
||||
}
|
||||
}
|
||||
|
||||
function rehomeInputs(mode, { director, selector, controlGeneration }) {
|
||||
return {
|
||||
mode,
|
||||
'director-image-digest': director.servingDigest,
|
||||
'rollback-image-digest': director.rollbackDigest,
|
||||
...selectorInputs(selector),
|
||||
'expected-control-generation': String(controlGeneration)
|
||||
}
|
||||
}
|
||||
|
||||
export function rehomeInspectInputs(context) {
|
||||
return rehomeInputs('inspect', context)
|
||||
}
|
||||
|
||||
// Pause keeps every durable field as inspected; only `enabled` and the generation change.
|
||||
// Every `confirmation` is the phrase the operator typed, never filled in here.
|
||||
export function rehomePauseInputs({ director, selector, control, confirmation }) {
|
||||
return {
|
||||
...rehomeInputs('pause', { director, selector, controlGeneration: control.generation }),
|
||||
'not-before': String(control.notBefore),
|
||||
'rate-per-minute': String(control.ratePerMinute),
|
||||
'preference-max-age-ms': String(control.preferenceMaxAgeMs),
|
||||
'host-cooldown-ms': String(control.hostCooldownMs ?? 604_800_000),
|
||||
'drain-grace-ms': String(control.drainGraceMs),
|
||||
confirmation
|
||||
}
|
||||
}
|
||||
|
||||
export function rehomeEnableInputs({
|
||||
director,
|
||||
selector,
|
||||
control,
|
||||
controlGeneration,
|
||||
notBefore,
|
||||
monitor,
|
||||
confirmation
|
||||
}) {
|
||||
return {
|
||||
...rehomeInputs('enable', { director, selector, controlGeneration }),
|
||||
'not-before': String(notBefore),
|
||||
// The enable job accepts exactly 10 per minute.
|
||||
'rate-per-minute': '10',
|
||||
'preference-max-age-ms': String(control.preferenceMaxAgeMs),
|
||||
'host-cooldown-ms': String(control.hostCooldownMs),
|
||||
'drain-grace-ms': String(control.drainGraceMs),
|
||||
'monitor-run-id': String(monitor.runId),
|
||||
'monitor-run-attempt': String(monitor.attempt),
|
||||
confirmation
|
||||
}
|
||||
}
|
||||
|
||||
export function publishInputs() {
|
||||
return { mode: 'publish' }
|
||||
}
|
||||
|
||||
export function directorDeployInputs({ imageDigest, predecessorDigest, rehomeGeneration }) {
|
||||
return {
|
||||
'image-digest': imageDigest,
|
||||
'regional-placement-mode': 'preserve',
|
||||
'region-correction-cohort-percent': 'preserve',
|
||||
'prune-incompatible-revisions': 'false',
|
||||
'expected-rehome-generation': String(rehomeGeneration),
|
||||
'bootstrap-runtime-identity': 'false',
|
||||
// Required by the form even without the bootstrap; it is only format-checked then.
|
||||
'predecessor-image-digest': predecessorDigest
|
||||
}
|
||||
}
|
||||
|
||||
export function configureInputs({
|
||||
cells,
|
||||
cellImageDigest,
|
||||
directorDigest,
|
||||
selectorGeneration,
|
||||
confirmation
|
||||
}) {
|
||||
return {
|
||||
environment: 'production',
|
||||
mode: 'configure',
|
||||
'cell-ids': cells.join(','),
|
||||
'selector-generation': String(selectorGeneration),
|
||||
'image-digest': cellImageDigest,
|
||||
'director-image-digest': directorDigest,
|
||||
confirmation
|
||||
}
|
||||
}
|
||||
|
||||
export function monitorDryRunInputs(selector) {
|
||||
return {
|
||||
mode: 'dry-run',
|
||||
...selectorInputs(selector),
|
||||
'migration-policy': 'strict',
|
||||
'recovery-source-cell-id': 'none',
|
||||
'capacity-cell-id': 'none'
|
||||
}
|
||||
}
|
||||
|
||||
export function parseConfigureWave(value) {
|
||||
const separator = String(value).lastIndexOf('=')
|
||||
const cells = String(value)
|
||||
.slice(0, separator)
|
||||
.split(',')
|
||||
.map((cell) => cell.trim())
|
||||
.filter(Boolean)
|
||||
const cellImageDigest = String(value).slice(separator + 1)
|
||||
if (
|
||||
separator < 1 ||
|
||||
cells.length === 0 ||
|
||||
new Set(cells).size !== cells.length ||
|
||||
cells.some((cell) => !PRODUCTION_CELL.test(cell))
|
||||
) {
|
||||
throw new Error('--configure must be <cell>[,<cell>...]=sha256:<cell image digest>')
|
||||
}
|
||||
return { cells, cellImageDigest: requireDigest(cellImageDigest, '--configure cell image digest') }
|
||||
}
|
||||
|
||||
/** The single 100% revision and the selector-rollback revision of `gcloud run services describe`. */
|
||||
export function directorRevisions(service) {
|
||||
const traffic = Array.isArray(service?.status?.traffic) ? service.status.traffic : []
|
||||
const serving = traffic.filter((entry) => (entry.percent ?? 0) > 0)
|
||||
const rollback = traffic.filter((entry) => entry.tag === ROLLBACK_TAG)
|
||||
if (serving.length !== 1 || serving[0].percent !== 100 || !serving[0].revisionName) {
|
||||
throw new Error('director does not serve one revision at 100%')
|
||||
}
|
||||
if (rollback.length !== 1 || !rollback[0].revisionName) {
|
||||
throw new Error(`director has no single ${ROLLBACK_TAG} revision`)
|
||||
}
|
||||
return { servingRevision: serving[0].revisionName, rollbackRevision: rollback[0].revisionName }
|
||||
}
|
||||
|
||||
export function revisionDigest(revision, label) {
|
||||
const image = revision?.spec?.containers?.[0]?.image
|
||||
const [repository, digest] = String(image ?? '').split('@')
|
||||
if (repository !== IMAGE_REPOSITORY)
|
||||
throw new Error(`${label} does not run the production relay image`)
|
||||
return requireDigest(digest, `${label} image`)
|
||||
}
|
||||
|
||||
/** The relay push line, `sha-<commit>: digest: sha256:… size: …`, must name the registry digest. */
|
||||
export function logConfirmsPublishedDigest(log, commit, digest) {
|
||||
return String(log)
|
||||
.split('\n')
|
||||
.some((line) => line.includes(`sha-${commit}: digest: ${digest} size:`))
|
||||
}
|
||||
|
||||
/** The last `relay_regional_rehome_control` JSON line a rehome run printed, of any mode. */
|
||||
export function rehomeResultFromLog(log) {
|
||||
let found
|
||||
for (const line of String(log).split('\n')) {
|
||||
const start = line.indexOf('{"event":"relay_regional_rehome_control"')
|
||||
if (start < 0) continue
|
||||
try {
|
||||
found = JSON.parse(line.slice(start).trim())
|
||||
} catch {
|
||||
// A truncated or echoed line is not the result.
|
||||
}
|
||||
}
|
||||
if (!found) throw new Error('the rehome run printed no control result')
|
||||
return {
|
||||
mode: found.mode,
|
||||
recovered: found.recovered,
|
||||
control: parseControl(found.control, `rehome ${found.mode} control`)
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* The generation a run paused rehome at, or undefined. Only two lines prove a run paused it: a
|
||||
* `pause`, and a failed enable's `recover-enable` that itself disabled rehome (`recovered: true`).
|
||||
* `recovered: false` means rehome was already disabled, possibly by a director safety pause.
|
||||
*/
|
||||
export function pausedGeneration(result) {
|
||||
const paused =
|
||||
!result.control.enabled &&
|
||||
(result.mode === 'pause' || (result.mode === 'recover-enable' && result.recovered === true))
|
||||
return paused ? result.control.generation : undefined
|
||||
}
|
||||
|
||||
export function admissionInspectResult(result) {
|
||||
if (result?.mode !== 'inspect') throw new Error('admission result is not an inspect')
|
||||
return parseSelector(result, 'admission selector')
|
||||
}
|
||||
Reference in New Issue
Block a user