Files
orca/cloud/apps/relay/src/assignment-isolated-cell-replacement-postgres.test.ts
T
Jinwoo Hong c8a5580659 fix(relay): re-place hosts off a cell isolated for a roll (#21911)
* fix(relay): re-place hosts off a cell isolated for a roll

A roll isolates a cell by moving it out of the 'general' admission class; the
cell then refuses every attach with 4503. The director never noticed, because
the only liveness test it applies to a host's current cell reads
`relay_cell_runtime.ready` and the heartbeat, and an isolated cell keeps
heartbeating ready=1 for the whole drain. So every host on that cell was handed
its own dead cell, closed, and handed it back — 500-1,900 hosts looping for
13-16 minutes per cell roll, at ~6 dials each per minute, with no neighbour
absorbing anything.

The sticky lane now treats a live incumbent whose admission is 'migration-only'
— the state a roll's isolate step writes — the same way it treats a dead one:
it returns null, which means "fall through to placement". The placement lane
had the identical hole eleven lines further down, so it takes the same
predicate; without that second swap the sticky change is inert, because
placement would hand the pin straight back (a draining cell has more headroom
than anyone). An isolated incumbent skips the dead-cell fence branch: that
branch exists to prove an unreachable cell stopped serving a host, and this one
is reachable and enforces the epoch itself.

'existing-only' is deliberately untouched — those cells serve the hosts they
already hold, and only `assignmentStrandedOnUnservedCell` may release that pin.
A host with an open `relay_assignment_migrations` row keeps its pin too, so
this stays disjoint from the migration machinery.

Claude-Session: https://claude.ai/code/session_01JNnE9qzUZMMnqpZWCqM3nb

* fix(relay): gate re-placement on a roll-isolation marker, not on admission

Review of the first commit found the predicate wrong. `migration-only` is an
admission class, not a drain signal: an Asia `--mode rollback`, an evacuation or
forward-recovery target awaiting a separate promote dispatch, a failed same-cap
wave's re-isolate, an abandoned migration retired on its target and a rehome
settlement all park loaded cells there durably, with no migration lease and no
open migration row. All five were indistinguishable from a roll's isolate, so
the first commit would have converted `operate-relay-asia-admission --mode
rollback` from a reversible admission flip into a mass move of ~4,000 hosts —
and, because `leastLoadedCell` treated region as a preference, into us-central1.

The signal is now an explicit stamp. `relay_cell_admission` gains a nullable
`roll_isolated_at`, added through the shared schema runner's catalog pre-check
so a migrated database takes no relation lock on boot and an un-migrated one
gets a catalog-only rewrite. The same-cap isolate step is its only writer, via a
new optional `rollIsolatedCells` on the selector apply; the same UPDATE that
writes the state clears the stamp whenever a cell leaves 'migration-only', so a
restore cannot leave one behind and a failed wave's re-isolate keeps the one it
has. Every other admission writer omits the field, so its cells stay unmarked
and their hosts stay pinned. Old directors ignore the field; old callers never
send it.

Region is now a constraint rather than a preference on this path only: a
re-placement must find a general, live cell with connection headroom in the
host's own region, or the pin is kept and one
`orca_relay_sticky_replacement_deferred` event is logged. Cross-region spill is
no longer reachable here.

The fence bypass is narrowed to a live incumbent. It was always a no-op for the
intended case, and for a stamped cell that stops heartbeating while still
holding sockets it reopened split-brain; that cell now takes the dead-cell path
unchanged.

Also: the hot-path admission reader no longer throws on an unrecognised state —
it sits on every sticky dial and the rule it feeds is "move the host", so an
unreadable row has to mean "don't". And the sticky lane reads the admission row
once for both the stranded rule and the stamp instead of twice.

Claude-Session: https://claude.ai/code/session_01JNnE9qzUZMMnqpZWCqM3nb

* fix(relay): emit the re-placement events after the transaction commits

CodeRabbit on assignment-store.ts:1075. Both events were written where they are
decided, which is inside assignOnce's transaction. A reservation or lease write
failing after that point rolls the placement back, but a line already on stdout
cannot be rolled back with it — so the canary this PR asks an operator to read
would count re-placements that never happened, and a Postgres transaction retry
could leave a stale line behind as well.

The transaction now returns its events alongside the RelayAssignment and the
caller flushes them once it has resolved. Returning them rather than setting a
variable in the enclosing scope is what makes the retry case safe too: only the
attempt that committed can carry its events out. assign()'s signature is
unchanged; the extra shape lives entirely inside assignOnce.

orca_relay_sticky_replacement_deferred was moved the same way. It cost one more
push into the array that already existed, and it is decided inside the same
transaction, so leaving it behind would have been the odd case rather than the
cheap one.

The new test injects a failure on the first write after the decision, asserts no
event is emitted, and asserts the assignment is still on its original cell —
without that second assertion the absence would only prove the emit was early,
not that it would have been wrong. A control dial with nothing injected emits
exactly one event, so the case cannot pass on a broken harness. With the emit
put back inside the transaction, it fails.

Claude-Session: https://claude.ai/code/session_01JNnE9qzUZMMnqpZWCqM3nb

* fix(relay): expire the roll stamp, correct the wire note, assert the stamp landed

Delta review findings B, D and E. A (the deferral path's cost) is deliberately
not implemented; it is now written up under Follow-ups in the PR body as
required before any Asia roll, because it cannot fire in a US canary.

B, which also closes C: the stamp was written, carried and never compared to
anything. A roll isolates and restores one cell inside ~15 minutes, so a stamp
older than two hours is not a roll in progress. It is a failed wave whose
failsafe re-isolated a possibly healthy cell and is waiting on an operator — the
postmortem in this tree records gaps of hours — or an orphan left by a director
rollback whose restore wrote 'general' without the clause that clears the stamp,
which the selector's 'keep' branch would then preserve until some later park
reactivated it. Both want the same answer and it is the pre-existing one: keep
the pin. One comparison against a value already on the row.

The bound takes the caller's `now` rather than reading the clock again, so one
assign reasons about one instant; the stamp's age is now a thing that decides
whether a host moves, and two clock reads could disagree across it.

D: the comment beside the new request field claimed an updated caller reaching
an older director "is simply ignored". The schema is .strict(), so it is a 400.
That fails closed — the isolate aborts before MUTATION_STARTED is set and
nothing is written — but it is a deploy ordering constraint, and it was
undocumented. The comment now says so and the PR body's rollout notes carry it.

E: nothing read the `rollIsolated` the script already prints, so an older script
against a newer director would silently produce today's behaviour and the canary
would read as "the fix did nothing" with no way to tell that from a wrong
premise. Both isolate steps now assert it, beside the generation they already
parse.

Claude-Session: https://claude.ai/code/session_01JNnE9qzUZMMnqpZWCqM3nb
2026-09-21 04:18:57 -04:00

365 lines
14 KiB
TypeScript
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
import { createHash } from 'node:crypto'
import { afterAll, beforeAll, describe, expect, it } from 'vitest'
import { RelayAssignmentStore } from './assignment-store.js'
import {
encodeMembership,
type CellAdmissionMembership,
type CellAdmissionState
} from './cell-admission-selector.js'
import type { RelayCellConfig } from './config.js'
import { openRelayDatabase, type RelayDatabase } from './database.js'
const databaseUrl = process.env.ORCA_RELAY_TEST_POSTGRES_URL
const describePostgres = databaseUrl ? describe : describe.skip
const USER_PREFIX = 'isolated-replacement-postgres'
const NOW = 100
const CAPPED = {
capacityRequests: 1_000,
connectionHardCap: 600,
connectionUnobservedBound: 50
} as const
const ISOLATED: RelayCellConfig = {
id: 'isolated-replacement-source',
url: 'https://isolated-replacement-source.example.com',
region: 'us-central1',
...CAPPED
}
const TARGETS: RelayCellConfig[] = [
{
id: 'isolated-replacement-target-a',
url: 'https://isolated-replacement-target-a.example.com',
region: 'us-central1',
...CAPPED
},
{
id: 'isolated-replacement-target-b',
url: 'https://isolated-replacement-target-b.example.com',
region: 'us-central1',
...CAPPED
}
]
const CELLS = [ISOLATED, ...TARGETS]
const HOST_COUNT = 50
function hostIdentity(index: number): { userId: string; relayHostId: string } {
return {
userId: `${USER_PREFIX}-${index}`,
// relay host ids are fixed-width opaque ids.
relayHostId: `isolatedhost${String(index).padStart(4, '0')}`
}
}
describePostgres('PostgreSQL re-placement off a cell isolated for a roll', () => {
const databases: RelayDatabase[] = []
let stores: RelayAssignmentStore[] = []
async function reservedRequests(cellId: string): Promise<number> {
const rows = await databases[0]!.query(
`SELECT reserved_requests FROM relay_cells WHERE cell_id = ?`,
[cellId]
)
return Number(rows[0]!['reserved_requests'])
}
async function heartbeatAll(): Promise<void> {
for (const [index, cell] of CELLS.entries()) {
await stores[0]!.recordCellHeartbeat({
cellId: cell.id,
cellUrl: cell.url,
cellIncarnation: `1111111${index}-1111-4111-8111-111111111111`,
startedAt: 50,
ready: true,
observedRequests: 0,
region: cell.region,
totalConnections: 0,
inFlightConnections: 0,
reservedConnectionUnits: 0,
enforcedConnectionUnits: 0,
connectionHardCap: 600,
connectionUnobservedBound: 50
})
}
}
// The real isolate path: one selector apply, generation CAS and all, naming
// the cells it stamps. Nothing else in the fleet may write the stamp.
async function applySelector(
states: Record<string, CellAdmissionState>,
rollIsolatedCells?: string[]
): Promise<void> {
const current = await stores[0]!.inspectCellAdmissionSelector()
// Built from relay_cells, not from the selector's own membership: the apply
// requires exact coverage of the fleet, and this database is shared.
const fleet = await databases[0]!.query(
`SELECT cell_id FROM relay_cells ORDER BY cell_id ASC`
)
const membership: CellAdmissionMembership = {
existingOnly: [],
migrationOnly: [],
general: []
}
for (const row of fleet) {
const cellId = String(row['cell_id'])
const state =
states[cellId] ??
(current.selector.membership.existingOnly.includes(cellId)
? 'existing-only'
: current.selector.membership.migrationOnly.includes(cellId)
? 'migration-only'
: 'general')
if (state === 'existing-only') membership.existingOnly.push(cellId)
else if (state === 'migration-only') membership.migrationOnly.push(cellId)
else membership.general.push(cellId)
}
await stores[0]!.applyCellAdmissionSelector({
attemptId: `isolated-${current.selector.generation}-${Object.keys(states).join('-')}`,
expectedGeneration: current.selector.generation,
...(current.selector.generation === 0
? {
expectedMembershipSha256: createHash('sha256')
.update(encodeMembership(current.selector.membership))
.digest('hex')
}
: {}),
membership,
...(rollIsolatedCells ? { rollIsolatedCells } : {})
})
}
async function rollIsolatedAt(cellId: string): Promise<number | null> {
const rows = await databases[0]!.query(
`SELECT roll_isolated_at FROM relay_cell_admission WHERE cell_id = ?`,
[cellId]
)
const value = rows[0]?.['roll_isolated_at']
return value === undefined || value === null ? null : Number(value)
}
// Every case starts from the whole fleet general; a case that leaves a cell
// isolated would otherwise starve the next one of placement candidates.
async function resetFleet(): Promise<void> {
await deleteHostRows()
await applySelector(Object.fromEntries(CELLS.map((cell) => [cell.id, 'general'])))
}
// The selector is fleet-wide and this database is shared with every other
// Postgres file in the project, all of which write admission through the
// generation-0 helpers. Advancing the generation and leaving it advanced
// would fail every one of them with admission_selector_boundary_active, so
// this file puts the boundary back exactly as it found it.
async function resetSelectorBoundary(): Promise<void> {
await databases[0]!.query(
`UPDATE relay_admission_selectors SET generation = 0, attempt_id = NULL
WHERE selector_id = 'general'`
)
await databases[0]!.query(
`DELETE FROM relay_admission_selector_intents WHERE attempt_id LIKE 'isolated-%'`
)
// Rewrites membership_json from the live fleet, which generation 0 allows.
await stores[0]!.reconcileCells([], false)
}
async function deleteHostRows(): Promise<void> {
for (const table of [
'relay_control_connection_reservations',
'relay_assignment_activity_leases',
'relay_assignment_migrations',
'relay_assignments'
]) {
await databases[0]!.query(`DELETE FROM ${table} WHERE user_id LIKE '${USER_PREFIX}-%'`)
}
}
beforeAll(async () => {
// Four connections so the concurrent dials below really contend on the
// fleet-wide relay_cells lock rather than queueing in one client.
for (let index = 0; index < 4; index += 1) {
databases.push(await openRelayDatabase({ databaseUrl, dataDir: '' }))
}
stores = databases.map(
(database) =>
new RelayAssignmentStore(database, () => NOW, {
requireLiveCells: true,
heartbeatTtlMs: 45_000
})
)
await deleteHostRows()
// Registers the cell rows without touching admission, so this works at any
// selector generation the shared database happens to be sitting at.
await stores[0]!.reconcileCells(CELLS, false)
await resetSelectorBoundary()
await heartbeatAll()
})
afterAll(async () => {
if (databases[0]) {
await deleteHostRows()
// Order matters: drop the cells first, then rebuild the selector's
// membership from what is left, or it names rows that no longer exist and
// every later read fails admission_selector_membership_drift.
for (const cell of CELLS) {
for (const table of [
'relay_cell_connection_snapshots',
'relay_cell_connection_runtime',
'relay_cell_connection_limits',
'relay_cell_runtime',
'relay_cell_admission',
'relay_cell_regions',
'relay_cells'
]) {
await databases[0].query(`DELETE FROM ${table} WHERE cell_id = ?`, [cell.id])
}
}
await resetSelectorBoundary()
}
for (const connection of databases) await connection.close()
})
it('stamps only the named cell on isolate and clears it on restore', async () => {
await resetFleet()
expect(await rollIsolatedAt(ISOLATED.id)).toBeNull()
await applySelector({ [ISOLATED.id]: 'migration-only' }, [ISOLATED.id])
const stamped = await rollIsolatedAt(ISOLATED.id)
expect(stamped).toBe(NOW)
// Cells the isolate did not name stay unmarked even when parked in the same
// apply: that is every non-roll flow, and its hosts must keep their pin.
await applySelector({ [TARGETS[0]!.id]: 'migration-only' })
expect(await rollIsolatedAt(TARGETS[0]!.id)).toBeNull()
// A failed wave re-isolates rather than restoring; the stamp survives, and
// does not move, so the hosts left behind stay eligible.
await applySelector({ [ISOLATED.id]: 'migration-only' }, [ISOLATED.id])
expect(await rollIsolatedAt(ISOLATED.id)).toBe(stamped)
// Restore writes 'general', and the same statement clears the stamp.
await applySelector({ [ISOLATED.id]: 'general' })
expect(await rollIsolatedAt(ISOLATED.id)).toBeNull()
}, 30_000)
it('ignores a stamp older than the roll it is supposed to describe', async () => {
// A failed wave keeps its stamp on purpose and can sit for hours; past the
// bound the cell stops shedding hosts one dial at a time.
await resetFleet()
const identity = hostIdentity(902)
const first = await stores[0]!.assign(identity, 'us-central1')
await applySelector({ [ISOLATED.id]: 'migration-only' }, [ISOLATED.id])
await databases[0]!.query(
`UPDATE relay_cell_admission SET roll_isolated_at = ? WHERE cell_id = ?`,
[NOW - (2 * 60 * 60_000 + 1), ISOLATED.id]
)
expect(await stores[0]!.assign(identity, 'us-central1')).toMatchObject({
cellId: first.cellId,
assignmentEpoch: first.assignmentEpoch
})
}, 30_000)
it('re-places every host off an isolated cell without leaking a reservation', async () => {
await resetFleet()
const identities = Array.from({ length: HOST_COUNT }, (_, index) => hostIdentity(index))
const sourceBaseline = await reservedRequests(ISOLATED.id)
const targetBaseline =
(await reservedRequests(TARGETS[0]!.id)) + (await reservedRequests(TARGETS[1]!.id))
// Everyone lands on the cell about to be isolated.
await applySelector({ [TARGETS[0]!.id]: 'migration-only', [TARGETS[1]!.id]: 'migration-only' })
const first = new Map<string, { cellId: string; assignmentEpoch: number }>()
for (const identity of identities) {
const grant = await stores[0]!.assign(identity, 'us-central1')
expect(grant.cellId).toBe(ISOLATED.id)
first.set(identity.relayHostId, grant)
}
await applySelector({ [TARGETS[0]!.id]: 'general', [TARGETS[1]!.id]: 'general' })
// The isolate step. The state and the stamp are written by one UPDATE under
// the fleet-wide relay_cells lock, so there is no torn state to race against.
await applySelector({ [ISOLATED.id]: 'migration-only' }, [ISOLATED.id])
expect(await rollIsolatedAt(ISOLATED.id)).not.toBeNull()
const grants = await Promise.all(
identities.map(
async (identity, index) =>
await stores[1 + (index % (stores.length - 1))]!.assign(identity, 'us-central1')
)
)
for (const [index, grant] of grants.entries()) {
const identity = identities[index]!
expect(grant.cellId).not.toBe(ISOLATED.id)
expect(TARGETS.map(({ id }) => id)).toContain(grant.cellId)
// Exactly once: a double bump would mean two transactions both moved it.
expect(grant.assignmentEpoch).toBe(first.get(identity.relayHostId)!.assignmentEpoch + 1)
}
expect(await reservedRequests(ISOLATED.id)).toBe(sourceBaseline)
expect(
(await reservedRequests(TARGETS[0]!.id)) + (await reservedRequests(TARGETS[1]!.id))
).toBe(targetBaseline + HOST_COUNT)
const rows = await databases[0]!.query(
`SELECT COUNT(*) AS count FROM relay_assignments
WHERE user_id LIKE '${USER_PREFIX}-%' AND cell_id = ?`,
[ISOLATED.id]
)
expect(Number(rows[0]!['count'])).toBe(0)
}, 60_000)
it('lets exactly one of a host’s racing dials win the re-placement', async () => {
await resetFleet()
const identity = hostIdentity(900)
await applySelector({ [TARGETS[0]!.id]: 'migration-only', [TARGETS[1]!.id]: 'migration-only' })
const first = await stores[0]!.assign(identity, 'us-central1')
expect(first.cellId).toBe(ISOLATED.id)
await applySelector({ [TARGETS[0]!.id]: 'general', [TARGETS[1]!.id]: 'general' })
const targetBaseline =
(await reservedRequests(TARGETS[0]!.id)) + (await reservedRequests(TARGETS[1]!.id))
const sourceBefore = await reservedRequests(ISOLATED.id)
await applySelector({ [ISOLATED.id]: 'migration-only' }, [ISOLATED.id])
// Two dials from the same host, on separate connections, through the same
// sticky lane. The per-assignment row lock is what must serialise them.
const raced = await Promise.allSettled([
stores[1]!.assign(identity, 'us-central1'),
stores[2]!.assign(identity, 'us-central1')
])
const granted = raced.flatMap((result) =>
result.status === 'fulfilled' ? [result.value] : []
)
expect(granted.length).toBeGreaterThan(0)
// However many dials were granted, only one re-placement may have
// committed: the epoch advances by exactly one and the reservation moves
// exactly one unit.
const settled = await stores[0]!.resolve(identity)
expect(settled?.assignmentEpoch).toBe(first.assignmentEpoch + 1)
expect(settled?.cellId).not.toBe(ISOLATED.id)
for (const grant of granted) expect(grant.cellId).toBe(settled?.cellId)
expect(await reservedRequests(ISOLATED.id)).toBe(sourceBefore - 1)
expect(
(await reservedRequests(TARGETS[0]!.id)) + (await reservedRequests(TARGETS[1]!.id))
).toBe(targetBaseline + 1)
}, 30_000)
it('grants no host the isolated cell after the admission flip commits', async () => {
await resetFleet()
const identity = hostIdentity(901)
const first = await stores[0]!.assign(identity, 'us-central1')
await applySelector({ [ISOLATED.id]: 'migration-only' }, [ISOLATED.id])
for (let dial = 0; dial < 5; dial += 1) {
const grant = await stores[1 + (dial % (stores.length - 1))]!.assign(
identity,
'us-central1'
)
expect(grant.cellId).not.toBe(ISOLATED.id)
}
// Only the first dial re-places; the rest are ordinary sticky re-grants.
expect((await stores[0]!.resolve(identity))?.assignmentEpoch).toBe(
first.assignmentEpoch + 1
)
}, 30_000)
})