From 54ded3bc185fb3fec7b3ad55cfa2ca9c0596cfd7 Mon Sep 17 00:00:00 2001 From: Jinwoo Hong <73622457+Jinwoo-H@users.noreply.github.com> Date: Tue, 6 Oct 2026 00:53:51 -0400 Subject: [PATCH] fix(relay): commit the cell counter in one round trip; cells boot without the database (#25765) * fix(relay): commit the cell counter in one round trip; cells boot without the DB Step 2 (option B) cell image: - One-round-trip counter commit at acquireActivity, releaseActivity and activateControl: the final counter UPDATE and COMMIT go as one simple-query message. Server errors mean COMMIT never ran (retry as today; 22012 = no row, rolled back and disambiguated outside the transaction); a lost connection is never retried. - Cells skip the schema apply and region backfill, so they listen while the database is down and turn ready on their first successful query. - G13: rehome target connection headroom folded into the existing NOWAIT UPDATE, excluding the host's own reservation by key. - fixLevel on every runtime metrics line, plus declared (not applied) cell fix-level metrics and alert. - Per-desktop drain disconnect-gap measurement from existing log lines. - Census test that fails on floating database promises; fixes two shutdown sites. Lock-wait sample keeps the combined role. * fix(relay): make the outdated-image alert creatable: one PromQL condition, 1 h lookback, fixed floor A PromQL condition must be the only condition in its policy, and alerts on log-based metrics may look back at most 25 h. Replace the 6-day/7-day design with relay_cell_min_fix_level (tfvars, raised by a targeted apply after each wave) and one query: a serving cell below the floor or reporting no level, sustained 6 h. Drops the separate without-level metric. * fix(relay): review fixes: gap-script ordering, wider promise census, fused-path guard, row-busy as scheduled - Drain gap script: sort closes by time (gcloud exports newest first) and refuse an invalid drain start. - Census: any floating promise in relay src, including callback-discarded and never-read ones, with a reviewed never-rejects list. - Test the fused counter commit through the store the server builds, so a wrapper that stops forwarding commitWithFinal fails CI. - Same-cap shadow gate: a row-busy refusal (the host's own release still holds its row) is a scheduled 503, like an own early retry. No client change. * test(relay): judge drain redials by no host refused twice, not a refusal count The row-busy count tracks how many releases are still in flight at the dial (80 of 180 every run at 1 s, against a bar of 90). What matters is that the release has finished by the next dial: assert no host is refused twice, keep the time-to-placed p95 bound. * fix(relay): cap row-busy as scheduled at the drain-return admissions; bound the gap script's window Shadow gate: a row-busy refusal of a drained host follows its drain-return lane admission, so per minute only that many (plus a rounding margin of 2) are scheduled; the rest stay non-drain, so row contention the drain does not explain still fails the budget. Gap script: --drain-ended-at excludes the new container's closes after the roll; later grants still close a gap. --- cloud/apps/relay/src/assignment-store.ts | 107 +++++++-- ...ell-boot-without-database-postgres.test.ts | 121 +++++++++++ .../src/cell-startup-environment.test.ts | 55 +++++ ...trol-accept-cell-row-lock-postgres.test.ts | 69 +++++- cloud/apps/relay/src/database.ts | 86 +++++++- ...lease-cell-row-contention-postgres.test.ts | 19 +- .../relay/src/fault-injection-test-entry.ts | 2 +- .../relay/src/floating-promise-census.test.ts | 186 ++++++++++++++++ cloud/apps/relay/src/index.ts | 5 +- .../apps/relay/src/observed-relay-database.ts | 10 + ...s-commit-with-final-write-postgres.test.ts | 203 ++++++++++++++++++ .../src/postgres-lock-wait-sample.test.ts | 28 +++ .../relay/src/postgres-lock-wait-sample.ts | 2 +- ...al-rehome-target-row-lock-postgres.test.ts | 60 ++++++ cloud/apps/relay/src/relay-fix-level.ts | 5 + .../relay/src/relay-observability.test.ts | 15 ++ cloud/apps/relay/src/relay-observability.ts | 2 + .../terraform-root-partition/families.json | 2 + .../measure-relay-drain-disconnect-gap.mjs | 161 ++++++++++++++ ...easure-relay-drain-disconnect-gap.test.mjs | 158 ++++++++++++++ .../relay-same-cap-shadow-gate-verdict.mjs | 20 ++ .../scripts/relay-same-cap-shadow-gate.mjs | 3 +- .../relay-same-cap-shadow-gate.test.mjs | 38 ++++ .../terraform/environments/production.tfvars | 4 + cloud/infra/terraform/relay-observability.tf | 65 ++++++ cloud/infra/terraform/variables.tf | 6 + cloud/package.json | 2 +- 27 files changed, 1397 insertions(+), 37 deletions(-) create mode 100644 cloud/apps/relay/src/cell-boot-without-database-postgres.test.ts create mode 100644 cloud/apps/relay/src/cell-startup-environment.test.ts create mode 100644 cloud/apps/relay/src/floating-promise-census.test.ts create mode 100644 cloud/apps/relay/src/postgres-commit-with-final-write-postgres.test.ts create mode 100644 cloud/apps/relay/src/postgres-lock-wait-sample.test.ts create mode 100644 cloud/apps/relay/src/relay-fix-level.ts create mode 100644 cloud/dev/scripts/measure-relay-drain-disconnect-gap.mjs create mode 100644 cloud/dev/scripts/measure-relay-drain-disconnect-gap.test.mjs diff --git a/cloud/apps/relay/src/assignment-store.ts b/cloud/apps/relay/src/assignment-store.ts index 2eaf143940b..e082aaa72a5 100644 --- a/cloud/apps/relay/src/assignment-store.ts +++ b/cloud/apps/relay/src/assignment-store.ts @@ -66,6 +66,7 @@ import { type ControlRenewalRequest } from './control-renewal-statement.js' import { + commitWithFinalWrite, REGIONAL_REHOME_DEFAULT_HOST_COOLDOWN_MS } from './database.js' import type { RelayCellConfig } from './config.js' @@ -356,6 +357,13 @@ const ACTIVITY_REQUEST_UNITS: Record = { } const ASSIGNMENT_LOCK_RETRY_DEADLINE_MS = 15_000 +// Clamped at zero on release; an increase that would pass capacity changes no row. +const CELL_RESERVATION_UPDATE = `UPDATE relay_cells SET reserved_requests = + CASE WHEN reserved_requests + ? < 0 THEN 0 ELSE reserved_requests + ? END, + updated_at = ? + WHERE cell_id = ? + AND (? <= 0 OR reserved_requests + ? <= capacity_requests) + RETURNING cell_id` // About one release round trip from Asia, with margin; see assignStickyOnce. const STICKY_ASSIGNMENT_ROW_WAIT_MS = 1_000 @@ -3724,7 +3732,7 @@ export class RelayAssignmentStore { ) // Keep the contended cell row locked for only the final write and commit. if (!existing) { - await this.adjustCellReservationAtomically(transaction, input.cellId, units) + await this.commitCellReservationAtomically(transaction, input.cellId, units) } }) }) @@ -3886,9 +3894,13 @@ export class RelayAssignmentStore { WHERE attempt_id = ?`, [now, request.attemptId] ) // Must stay the last statement: the row is held from here to COMMIT. - const target = await this.reserveRegionalRehomeTargetRow( - transaction, attempt.targetCellId, attempt.targetReservedUnits, now - ) + const target = await this.reserveRegionalRehomeTargetRow(transaction, { + identity: request, + assignmentEpoch: attempt.assignmentEpoch, + cellId: attempt.targetCellId, + units: attempt.targetReservedUnits, + now + }) if (target !== 'reserved') { targetDeferral = target throw new RegionalRehomeTargetDeferred() @@ -4204,7 +4216,7 @@ export class RelayAssignmentStore { now ) // Release paths can safely defer the cell-row lock until their final write. - await this.adjustCellReservationAtomically( + await this.commitCellReservationAtomically( transaction, text(existing, 'cell_id'), -integer(existing, 'request_units') @@ -4337,7 +4349,7 @@ export class RelayAssignmentStore { // trips. Holding its write lock from the first of them capped a // far-from-Postgres cell at a couple of accepts a second. if (reservationDelta !== 0) { - await this.adjustCellReservationAtomically( + await this.commitCellReservationAtomically( transaction, input.cellId, reservationDelta @@ -6249,12 +6261,21 @@ export class RelayAssignmentStore { // the statement, because an admission flip no longer serialises against this // transaction any other way. Locking the admission row too makes a flip that // committed after the statement's snapshot re-evaluate on its new version. + // Connection headroom is re-checked here too: the candidate read happened before this + // transaction's own reservation insert, and a placement may have reserved since. Every + // other reservation inserter holds this row first, so NOWAIT defers rather than races. + // The host's own reservation is excluded by key because the insert may have added none. private async reserveRegionalRehomeTargetRow( database: RelayDatabase, - cellId: string, - units: number, - now: number + input: { + identity: AssignmentIdentity + assignmentEpoch: number + cellId: string + units: number + now: number + } ): Promise<'reserved' | RegionalRehomeTargetDeferral> { + const { identity, cellId, units, now } = input const lockClause = database.dialect === 'sqlite' ? '' : 'FOR UPDATE OF cell, admission NOWAIT' try { const rows = await database.queryLocked( @@ -6268,8 +6289,39 @@ export class RelayAssignmentStore { UPDATE relay_cells SET reserved_requests = reserved_requests + ?, updated_at = ? WHERE cell_id IN (SELECT cell_id FROM target) AND reserved_requests + ? <= capacity_requests + AND ( + NOT EXISTS (SELECT 1 FROM relay_cell_connection_limits WHERE cell_id = ?) + OR EXISTS ( + SELECT 1 FROM relay_cell_connection_limits limits + JOIN relay_cell_connection_snapshots snapshot ON snapshot.cell_id = limits.cell_id + JOIN relay_cell_runtime runtime ON runtime.cell_id = limits.cell_id + WHERE limits.cell_id = ? + AND snapshot.snapshot_at > ? + AND snapshot.cell_incarnation = runtime.cell_incarnation + AND snapshot.enforced_connection_units + + (SELECT COUNT(*) FROM relay_control_connection_reservations reservation + WHERE reservation.cell_id = limits.cell_id + AND reservation.state IN ('reserved', 'late-arrival-debt', 'claimed') + AND NOT (reservation.user_id = ? AND reservation.relay_host_id = ? + AND reservation.assignment_epoch = ?)) + + limits.unobserved_bound + < limits.hard_cap - ? + ) + ) RETURNING cell_id`, - [cellId, units, now, units], + [ + cellId, + units, + now, + units, + cellId, + cellId, + now - this.heartbeatTtlMs, + identity.userId, + identity.relayHostId, + input.assignmentEpoch, + RELAY_ADMISSION_BUDGETS.reservedHostControls + ], { failIfUnavailable: true, lockClauseInStatement: true, @@ -8492,16 +8544,31 @@ export class RelayAssignmentStore { cellId: string, delta: number ): Promise { - const rows = await database.query( - `UPDATE relay_cells SET reserved_requests = - CASE WHEN reserved_requests + ? < 0 THEN 0 ELSE reserved_requests + ? END, - updated_at = ? - WHERE cell_id = ? - AND (? <= 0 OR reserved_requests + ? <= capacity_requests) - RETURNING cell_id`, - [delta, delta, this.now(), cellId, delta, delta] - ) - if (rows.length > 0) return + const rows = await database.query(CELL_RESERVATION_UPDATE, [ + delta, + delta, + this.now(), + cellId, + delta, + delta + ]) + if (rows.length === 0) await this.refuseCellReservation(database, cellId) + } + + // The same write as the last statement of its transaction, committed in the same round + // trip so the shared cell row is not held while the reply crosses to a far cell. + private async commitCellReservationAtomically( + transaction: RelayDatabase, + cellId: string, + delta: number + ): Promise { + const params = [delta, delta, this.now(), cellId, delta, delta] + if (!(await commitWithFinalWrite(transaction, CELL_RESERVATION_UPDATE, params))) { + await this.refuseCellReservation(transaction, cellId) + } + } + + private async refuseCellReservation(database: RelayDatabase, cellId: string): Promise { const cell = ( await database.query(`SELECT cell_id FROM relay_cells WHERE cell_id = ?`, [cellId]) )[0] diff --git a/cloud/apps/relay/src/cell-boot-without-database-postgres.test.ts b/cloud/apps/relay/src/cell-boot-without-database-postgres.test.ts new file mode 100644 index 00000000000..03df1728aa4 --- /dev/null +++ b/cloud/apps/relay/src/cell-boot-without-database-postgres.test.ts @@ -0,0 +1,121 @@ +import { createServer as createHttpServer, type Server } from 'node:http' +import { createServer as createTcpServer, connect, type AddressInfo } from 'node:net' +import { afterEach, describe, expect, it } from 'vitest' +import type { RelayConfig } from './config.js' +import { openRelayDatabase, type RelayDatabase } from './database.js' +import { createRelayReadiness } from './relay-readiness.js' +import { createRelayServer } from './relay-server.js' + +const databaseUrl = process.env.ORCA_RELAY_TEST_POSTGRES_URL +const describePostgres = databaseUrl ? describe : describe.skip + +// c25 exited after 45 s of boot retries when its database was unreachable, and the MIG +// recreated it into the same loop. A cell now opens a lazy pool and listens at once. +describePostgres('cell boot while PostgreSQL is unreachable', () => { + const cleanups: Array<() => Promise> = [] + + afterEach(async () => { + for (const cleanup of cleanups.splice(0).reverse()) await cleanup() + }) + + it('listens with /health 200 and /ready 503, then turns ready once the database answers', async () => { + const proxyPort = await unusedPort() + const target = new URL(databaseUrl!) + const viaProxy = new URL(databaseUrl!) + viaProxy.hostname = '127.0.0.1' + viaProxy.port = String(proxyPort) + + const startedAt = performance.now() + const database = await openRelayDatabase({ + databaseUrl: viaProxy.toString(), + dataDir: '', + appliesPostgresSchema: false + }) + cleanups.push(async () => await database.close()) + // Nothing dialled: the 2 s connect timeout would show here otherwise. + expect(performance.now() - startedAt).toBeLessThan(500) + await expect(database.query('SELECT 1')).rejects.toThrow() + + const jwks = await listen( + createHttpServer((_request, response) => { + response.setHeader('content-type', 'application/json') + response.end('{"keys":[]}') + }) + ) + cleanups.push(async () => await close(jwks)) + const jwksUrl = `http://127.0.0.1:${(jwks.address() as AddressInfo).port}/jwks` + const relay = createRelayServer(cellConfig(jwksUrl), database) + const server = await listen(relay.server) + cleanups.push(async () => await close(server)) + const base = `http://127.0.0.1:${(server.address() as AddressInfo).port}` + + expect((await fetch(`${base}/health`)).status).toBe(200) + expect((await fetch(`${base}/ready`)).status).toBe(503) + + const proxy = await listenOn( + createTcpServer((socket) => { + const upstream = connect(Number(target.port), target.hostname) + socket.pipe(upstream).pipe(socket) + socket.on('error', () => upstream.destroy()) + upstream.on('error', () => socket.destroy()) + }), + proxyPort + ) + cleanups.push(async () => await close(proxy)) + const readiness = createRelayReadiness(database, jwksUrl, { cacheMs: 0 }) + await expect(readiness.check()).resolves.toBe(true) + }, 20_000) +}) + +async function unusedPort(): Promise { + const probe = await listen(createTcpServer()) + const { port } = probe.address() as AddressInfo + await close(probe) + return port +} + +async function listen>(server: T): Promise { + return await listenOn(server, 0) +} + +async function listenOn>( + server: T, + port: number +): Promise { + await new Promise((resolve) => server.listen(port, '127.0.0.1', resolve)) + return server +} + +async function close(server: Server | ReturnType): Promise { + if ('closeAllConnections' in server) server.closeAllConnections() + await new Promise((resolve) => server.close(() => resolve())) +} + +function cellConfig(jwksUrl: string): RelayConfig { + return { + port: 0, + publicUrl: 'https://c25.relay.example.test', + cellUrl: 'https://c25.relay.example.test', + region: 'asia-east2', + authIssuer: 'https://auth.example.test', + authAudience: 'orca-relay', + jwksUrl, + assignmentSigningKey: new Uint8Array(32), + role: 'cell', + cellId: 'production-gce-c25', + cells: [], + adminAudience: 'https://relay.example.test/v1/admin/drain', + deployServiceAccount: 'deploy@example.test', + runtimeServiceAccount: 'relay-cell@example.test', + adminJwksUrl: jwksUrl, + databasePoolMax: 10, + publicAssignmentsEnabled: true, + publicAssignmentConcurrency: 2, + publicAssignmentQueueMax: 128, + publicAssignmentWaitMs: 4_000, + publicResolveConcurrency: 1, + publicResolveWaitMs: 5_000, + publicAssignmentRetryAfterSeconds: 5, + dataDir: './data' + } +} diff --git a/cloud/apps/relay/src/cell-startup-environment.test.ts b/cloud/apps/relay/src/cell-startup-environment.test.ts new file mode 100644 index 00000000000..f85493dd6d0 --- /dev/null +++ b/cloud/apps/relay/src/cell-startup-environment.test.ts @@ -0,0 +1,55 @@ +import { readFileSync } from 'node:fs' +import { describe, expect, it } from 'vitest' +import { loadRelayConfig } from './config.js' + +// A same-cap roll swaps only the image: the cell keeps the env its startup script already +// wrote. So every image must boot on exactly the variables that script renders, and a +// config change that needs a new variable is a template change, not an image-only roll. +const template = readFileSync( + new URL('../../../infra/terraform/relay-gce-startup.sh.tftpl', import.meta.url), + 'utf8' +) + +// One plausible production value per variable the template can write. +const RENDERED: Record = { + DATABASE_URL: 'postgres://relay@127.0.0.1:5432/orca_relay', + ORCA_RELAY_ASSIGNMENT_SIGNING_KEY: 'assignment-key-with-at-least-thirty-two-bytes', + ORCA_RELAY_PUBLIC_URL: 'https://c25.relay.onorca.dev', + ORCA_RELAY_CELL_URL: 'https://c25.relay.onorca.dev', + ORCA_RELAY_AUTH_ISSUER: 'https://auth.onorca.dev', + ORCA_RELAY_AUTH_AUDIENCE: 'orca-relay', + ORCA_RELAY_JWKS_URL: 'https://auth.onorca.dev/.well-known/jwks.json', + ORCA_RELAY_ROLE: 'cell', + ORCA_RELAY_CELL_ID: 'production-gce-c25', + ORCA_RELAY_REGION: 'asia-east2', + ORCA_RELAY_CELL_CAPACITY: '3000', + ORCA_RELAY_DATABASE_POOL_MAX: '16', + ORCA_RELAY_CELL_CONNECTION_HARD_CAP: '3000', + ORCA_RELAY_CELL_CONNECTION_UNOBSERVED_BOUND: '60', + ORCA_RELAY_CELLS_JSON: '[]', + ORCA_RELAY_ADMIN_AUDIENCE: 'https://relay.onorca.dev/v1/admin/drain', + ORCA_RELAY_DEPLOY_SERVICE_ACCOUNT: 'deploy@example.iam.gserviceaccount.com', + ORCA_RELAY_CAPACITY_SERVICE_ACCOUNT: 'capacity@example.iam.gserviceaccount.com', + ORCA_RELAY_ASIA_PROOF_SERVICE_ACCOUNT: 'asia-proof@example.iam.gserviceaccount.com', + ORCA_RELAY_RUNTIME_SERVICE_ACCOUNT: 'relay-cell@example.iam.gserviceaccount.com', + ORCA_RELAY_REHOME_DIRECTOR_SERVICE_ACCOUNT: 'relay-director@example.iam.gserviceaccount.com', + ORCA_RELAY_REHOME_AUDIENCE: 'https://relay.onorca.dev/v1/admin/host-drain', + ORCA_RELAY_DIRECTOR_URL: 'https://relay.onorca.dev', + ORCA_RELAY_HEARTBEAT_AUDIENCE: 'https://relay.onorca.dev/v1/admin/cell-heartbeat', + ORCA_RELAY_IMAGE_DIGEST: `sha256:${'a'.repeat(64)}` +} + +const TEMPLATE_VARIABLES = [...template.matchAll(/printf '([A-Z0-9_]+)=/g)].map(([, name]) => name!) + +describe('cell startup environment', () => { + it('boots a cell on exactly the variables the startup template writes', () => { + expect(TEMPLATE_VARIABLES.length).toBeGreaterThan(20) + expect(TEMPLATE_VARIABLES.filter((name) => RENDERED[name] === undefined)).toEqual([]) + const env = Object.fromEntries(TEMPLATE_VARIABLES.map((name) => [name, RENDERED[name]])) + expect(loadRelayConfig(env)).toMatchObject({ + role: 'cell', + cellId: 'production-gce-c25', + databaseUrl: RENDERED.DATABASE_URL + }) + }) +}) diff --git a/cloud/apps/relay/src/control-accept-cell-row-lock-postgres.test.ts b/cloud/apps/relay/src/control-accept-cell-row-lock-postgres.test.ts index 11a7b4d1b0e..fea3b3ebc09 100644 --- a/cloud/apps/relay/src/control-accept-cell-row-lock-postgres.test.ts +++ b/cloud/apps/relay/src/control-accept-cell-row-lock-postgres.test.ts @@ -1,6 +1,9 @@ -import { afterAll, beforeAll, describe, expect, it } from 'vitest' +import pg from 'pg' +import { afterAll, beforeAll, describe, expect, it, vi } from 'vitest' import { RelayAssignmentStore } from './assignment-store.js' +import type { RelayConfig } from './config.js' import { openRelayDatabase, type RelayDatabase } from './database.js' +import { createRelayServer } from './relay-server.js' const databaseUrl = process.env.ORCA_RELAY_TEST_POSTGRES_URL const describePostgres = databaseUrl ? describe : describe.skip @@ -177,6 +180,42 @@ describePostgres('PostgreSQL control accept without a held cell row', () => { expect(Number(cells[0]!.reserved_requests)).toBe(Number(units[0]!.units)) }, 20_000) + // The fused commit is optional on RelayDatabase, so a wrapper that stopped forwarding it + // would silently fall back to a separate COMMIT. Drive the store the server builds. + it('commits every counter write in one message through the production store', async () => { + await removeTestRows(databases[0]!) + const store = new RelayAssignmentStore(databases[0]!, () => 100) + await prepareCell(store) + const { assignments } = createRelayServer(cellConfig(), databases[0]!, { now: () => 100 }) + const assignment = await store.assign(first) + const firstControl = await assignments.activateControl(first, { + cellId: cell.id, + assignmentEpoch: assignment.assignmentEpoch, + generation: 1 + }) + + const sent = vi.spyOn(pg.Client.prototype, 'query') + let texts: string[] = [] + try { + await assignments.releaseActivity(first, firstControl) + await assignments.activateControl(first, { + cellId: cell.id, + assignmentEpoch: assignment.assignmentEpoch, + generation: 2 + }) + await assignments.acquireActivity(first, { + activityId: 'invite:fused', + kind: 'invite', + cellId: cell.id + }) + } finally { + texts = sent.mock.calls.map(([text]) => String(text)) + sent.mockRestore() + } + expect(texts.filter((text) => text.endsWith('; COMMIT'))).toHaveLength(3) + expect(texts.filter((text) => text === 'COMMIT')).toEqual([]) + }) + async function prepareCell(store: RelayAssignmentStore): Promise { await store.reconcileCells([cell]) await store.recordCellHeartbeat({ @@ -197,3 +236,31 @@ describePostgres('PostgreSQL control accept without a held cell row', () => { } }) + +function cellConfig(): RelayConfig { + return { + port: 0, + publicUrl: cell.url, + cellUrl: cell.url, + authIssuer: 'https://auth.example.test', + authAudience: 'orca-relay', + jwksUrl: 'https://auth.example.test/jwks', + assignmentSigningKey: new Uint8Array(32), + role: 'cell', + cellId: cell.id, + cells: [], + adminAudience: 'https://relay.example.test/v1/admin/drain', + deployServiceAccount: 'deploy@example.test', + runtimeServiceAccount: 'relay-cell@example.test', + adminJwksUrl: 'https://auth.example.test/jwks', + databasePoolMax: 10, + publicAssignmentsEnabled: true, + publicAssignmentConcurrency: 2, + publicAssignmentQueueMax: 128, + publicAssignmentWaitMs: 4_000, + publicResolveConcurrency: 1, + publicResolveWaitMs: 5_000, + publicAssignmentRetryAfterSeconds: 5, + dataDir: './data' + } +} diff --git a/cloud/apps/relay/src/database.ts b/cloud/apps/relay/src/database.ts index dbc5e88df8c..122dd3fc1c0 100644 --- a/cloud/apps/relay/src/database.ts +++ b/cloud/apps/relay/src/database.ts @@ -80,9 +80,23 @@ export interface RelayDatabase { operation: (transaction: RelayDatabase) => Promise, options?: RelayTransactionOptions ): Promise + // Only on a PostgreSQL transaction handle; see commitWithFinalWrite. + commitWithFinal?(sql: string, params?: unknown[]): Promise close(): Promise } +// Runs a single-row write with RETURNING as the transaction's last statement and reports +// whether it changed a row. On PostgreSQL the write and COMMIT go as one message, so the +// row lock is held for no round trip; a false result has already rolled the transaction back. +export async function commitWithFinalWrite( + database: RelayDatabase, + sql: string, + params: unknown[] = [] +): Promise { + if (database.commitWithFinal) return await database.commitWithFinal(sql, params) + return (await database.query(sql, params)).length > 0 +} + // RULE - no new index and no new column on `relay_control_connection_reservations`, // `relay_confirm_results`, `relay_audit_events`, `relay_connection_bases`, or any other large // table may be added to SCHEMA or to POSTGRES_SCHEMA_MIGRATIONS. The catalog pre-check skips a @@ -844,14 +858,62 @@ class SqliteDatabase extends SqliteTransaction { } } +// Literals for a simple-query message, which carries no bind parameters. Only safe +// integers and strings: anything else is a caller bug, not something to stringify. +function inlinePostgresParameters(sql: string, params: unknown[], client: pg.PoolClient): string { + let index = 0 + const inlined = sql.replace(/\?/g, () => { + const value = params[index++] + if (typeof value === 'number' && Number.isSafeInteger(value)) return String(value) + if (typeof value === 'string') return client.escapeLiteral(value) + throw new Error('unsupported_inline_parameter') + }) + if (index !== params.length) throw new Error('inline_parameter_count_mismatch') + return inlined +} + class PostgresTransaction implements RelayDatabase { readonly dialect = 'postgres' as const private held: { fromMs: number; site: CellLockHoldSite } | undefined private lockUnavailable = 0 private lockTimeouts = 0 + private state: 'open' | 'committed' | 'rolled-back' = 'open' constructor(protected readonly client: pg.PoolClient) {} + get open(): boolean { + return this.state === 'open' + } + + async commitWithFinal(sql: string, params: unknown[] = []): Promise { + this.assertNotCommitted() + // Zero rows divides by zero, so the message stops before COMMIT exactly when the write missed. + const message = + `WITH final_write AS (${inlinePostgresParameters(sql, params, this.client)}) ` + + 'SELECT 1 / (SELECT count(*)::int FROM final_write); COMMIT' + try { + await this.client.query(message) + } catch (error) { + // Simple query stops at the first error, so any server error means COMMIT never ran: + // retryable codes take the caller's normal rollback-and-retry path. A lost connection + // leaves the outcome unknown, and nothing retries it. + if (String((error as { code?: unknown }).code) !== '22012') { + rememberPostgresTransactionPhase(error, sql) + throw error + } + await this.client.query('ROLLBACK') + this.state = 'rolled-back' + return false + } + this.state = 'committed' + return true + } + + private assertNotCommitted(): void { + // A later statement would run in autocommit, outside the work it belongs to. + if (this.state === 'committed') throw new Error('postgres_transaction_already_committed') + } + consumeHold(): MeasuredHold | undefined { if (this.held === undefined) return undefined const hold = { holdMs: performance.now() - this.held.fromMs, site: this.held.site } @@ -874,6 +936,7 @@ class PostgresTransaction implements RelayDatabase { } async query(sql: string, params: unknown[] = []): Promise { + this.assertNotCommitted() try { const result = await this.client.query(postgresSql(sql), params) return returnsRows(sql) ? (result.rows as SqlRow[]) : [{ changes: result.rowCount ?? 0 }] @@ -1051,16 +1114,22 @@ export class PostgresDatabase implements RelayDatabase { try { await client.query('BEGIN') const result = await operation(transaction) - await client.query('COMMIT') + if (transaction.open) await client.query('COMMIT') recordMeasuredHold(this.holds, transaction) this.holds.recordUnavailable(transaction.consumeLockUnavailable()) this.holds.recordLockTimeout(transaction.consumeLockTimeouts()) return result } catch (error) { - await client.query('ROLLBACK').catch(() => undefined) + const open = transaction.open + if (open) await client.query('ROLLBACK').catch(() => undefined) this.holds.recordUnavailable(transaction.consumeLockUnavailable()) this.holds.recordLockTimeout(transaction.consumeLockTimeouts()) - if (!retryablePostgresTransactionError(error) || attempt === POSTGRES_TRANSACTION_ATTEMPTS) { + // Once the fused commit has ended the transaction, a retry would apply the work twice. + if ( + !open || + !retryablePostgresTransactionError(error) || + attempt === POSTGRES_TRANSACTION_ATTEMPTS + ) { if (retryablePostgresTransactionError(error) && options.reportRetries !== false) { console.warn( JSON.stringify({ @@ -1204,12 +1273,19 @@ export type RelayDatabaseOpenInput = { poolMax?: number applicationName?: string statementTimeoutMs?: number + // Directors own the PostgreSQL schema. A cell skips it and never touches the database + // at boot, so it starts listening while the database is down and stays unready until + // its first successful query. + appliesPostgresSchema?: boolean } export async function openRelayDatabase(input: RelayDatabaseOpenInput): Promise { let database: RelayDatabase + const appliesPostgresSchema = input.appliesPostgresSchema !== false if (input.databaseUrl) { - await applySchemaOnUntimedPool(input.databaseUrl, input.applicationName) + if (appliesPostgresSchema) { + await applySchemaOnUntimedPool(input.databaseUrl, input.applicationName) + } const pool = new pg.Pool({ connectionString: input.databaseUrl, max: input.poolMax ?? 10, @@ -1229,7 +1305,7 @@ export async function openRelayDatabase(input: RelayDatabaseOpenInput): Promise< } try { if (!input.databaseUrl) await applySchema(database) - await backfillRelayCellRegions(database) + if (!input.databaseUrl || appliesPostgresSchema) await backfillRelayCellRegions(database) return database } catch (error) { await database.close().catch(() => undefined) diff --git a/cloud/apps/relay/src/drain-release-cell-row-contention-postgres.test.ts b/cloud/apps/relay/src/drain-release-cell-row-contention-postgres.test.ts index 23f7877ac7f..a484fc81246 100644 --- a/cloud/apps/relay/src/drain-release-cell-row-contention-postgres.test.ts +++ b/cloud/apps/relay/src/drain-release-cell-row-contention-postgres.test.ts @@ -120,6 +120,8 @@ type DepartingReport = { firstAttempt: { placed: number; rejected: Tally; failed: Tally } // Row-busy refusals on any attempt, and errors no client would retry. busyRefusals: number + // Hosts refused more than once, for any reason, before being placed. + hostsRefusedTwice: number unexpected: Tally redials: number // Release start to the grant that re-placed the host, redials included. @@ -508,6 +510,7 @@ describePostgres('PostgreSQL drain releases against director placement', () => { const timeToPlaced: number[] = [] let redials = 0 let busyRefusals = 0 + let hostsRefusedTwice = 0 const unexpected: Tally = {} let placedInWindow = 0 const activationWork: Promise[] = [] @@ -576,6 +579,7 @@ describePostgres('PostgreSQL drain releases against director placement', () => { break } if (!outcome.startsWith('rejected:')) break + if (attempt === 1) hostsRefusedTwice += 1 const nextAt = sentAt + HOST_ASSIGN_MIN_INTERVAL_MS + @@ -598,6 +602,7 @@ describePostgres('PostgreSQL drain releases against director placement', () => { placed: timeToPlaced.length, firstAttempt, busyRefusals, + hostsRefusedTwice, unexpected, redials, timeToPlacedMs: { @@ -698,9 +703,10 @@ describePostgres('PostgreSQL drain releases against director placement', () => { expect(report.director.lockWaitingMean).toBeLessThan(0.5) expect(report.dials.placedCrossRegion).toBe(neighboursCapped ? report.dials.placed : 0) if (rate === 18) { - // Still cell-side: releases on one row serialise at ~1/RTT (~5.8/s) - // and the excess sheds to lease expiry. The director no longer waits. - expect(report.releases.ok).toBeLessThan(0.6 * report.releases.attempted) + // The release commits with its counter write, so the source row is no longer + // held for a round trip and releases stop shedding to lease expiry. With the + // separate COMMIT, 79 of 180 landed and the rest timed out waiting for the pool. + expect(report.releases.ok).toBe(report.releases.attempted) } } }, 240_000) @@ -718,8 +724,11 @@ describePostgres('PostgreSQL drain releases against director placement', () => { report.firstAttempt.rejected expect(total(slotRejections)).toBeLessThan(0.1 * report.hosts) expect(report.lockWaitingMean).toBeLessThan(0.5) - // Measured 16-20 refusals of 180 and a p95 of 5.7-6.3s; before, p95 was 11-16s. - expect(report.busyRefusals).toBeLessThanOrEqual(0.2 * report.hosts) + // A redial that beats its own release meets its own row (16-80 of 180, by how many + // releases are still in flight at the dial) and is answered at once. That release + // has finished by the next dial, so no host is refused twice. Measured p95 5.6-6.5s; + // before the row-busy answer, p95 was 11-16s. + expect(report.hostsRefusedTwice).toBe(0) expect(report.timeToPlacedMs.p95).toBeLessThanOrEqual(8_000) } } diff --git a/cloud/apps/relay/src/fault-injection-test-entry.ts b/cloud/apps/relay/src/fault-injection-test-entry.ts index 84c7e34f93a..37dea42cc30 100644 --- a/cloud/apps/relay/src/fault-injection-test-entry.ts +++ b/cloud/apps/relay/src/fault-injection-test-entry.ts @@ -53,7 +53,7 @@ server.listen(config.port, () => { const shutdown = (): void => { sessions.drain(0) - server.close(() => void realDatabase.close()) + server.close(() => void realDatabase.close().catch(() => undefined)) } process.once('SIGTERM', shutdown) process.once('SIGINT', shutdown) diff --git a/cloud/apps/relay/src/floating-promise-census.test.ts b/cloud/apps/relay/src/floating-promise-census.test.ts new file mode 100644 index 00000000000..d127409aa1c --- /dev/null +++ b/cloud/apps/relay/src/floating-promise-census.test.ts @@ -0,0 +1,186 @@ +import { readdirSync } from 'node:fs' +import { basename, join } from 'node:path' +import { fileURLToPath } from 'node:url' +import ts from 'typescript' +import { describe, expect, it } from 'vitest' + +// A promise nobody awaits rejects into Node's unhandledRejection, which ends the process: +// one lost database reply would drop every host on the cell. Production code must await, +// return, store and use, or `.catch` every promise it starts. +// Exempt: functions written to settle every failure themselves. Keyed file#name. +const NEVER_REJECTS = new Map([ + ['relay-background-operation.ts#runRelayBackgroundOperation', 'catches and logs every failure'], + ['assignment-cleanup-steps.ts#runAssignmentCleanup', 'runs each step through the above'], + ['cell-heartbeat-client.ts#send', 'one try/catch around the whole send'], + ['control-renewal-batch.ts#flush', 'rejects the waiters, never itself'], + ['regional-rehome-worker.ts#run', 'one try/catch around the whole poll'] +]) + +const sourceDirectory = fileURLToPath(new URL('.', import.meta.url)) +const COMPILER_OPTIONS: ts.CompilerOptions = { + target: ts.ScriptTarget.ES2022, + module: ts.ModuleKind.NodeNext, + moduleResolution: ts.ModuleResolutionKind.NodeNext, + strict: true, + noEmit: true, + skipLibCheck: true +} + +function productionFiles(): string[] { + return readdirSync(sourceDirectory) + .filter((entry) => entry.endsWith('.ts') && !entry.endsWith('.test.ts')) + .map((entry) => join(sourceDirectory, entry)) +} + +function isPromise(checker: ts.TypeChecker, node: ts.Node): boolean { + return checker.getTypeAtLocation(node).getSymbol()?.getName() === 'Promise' +} + +function isChainStep(node: ts.Node): node is ts.PropertyAccessExpression { + return ( + ts.isPropertyAccessExpression(node) && ['then', 'catch', 'finally'].includes(node.name.text) + ) +} + +// Walks a `.then/.catch/.finally` chain up to the expression that consumes it. +function consumer(call: ts.CallExpression): { node: ts.Node; caught: boolean } { + let node: ts.Node = call + let caught = false + while (isChainStep(node.parent) && ts.isCallExpression(node.parent.parent)) { + const step = node.parent.parent + const method = node.parent.name.text + if (method === 'catch' || (method === 'then' && step.arguments.length > 1)) caught = true + node = step + } + while (ts.isParenthesizedExpression(node.parent)) node = node.parent + return { node, caught } +} + +// A callback whose contextual type returns void (setTimeout, an event listener) drops the promise. +function discardedByCallback(checker: ts.TypeChecker, fn: ts.ArrowFunction): boolean { + const signatures = checker.getContextualType(fn)?.getCallSignatures() ?? [] + return ( + signatures.length > 0 && + signatures.every( + (signature) => (checker.getReturnTypeOfSignature(signature).flags & ts.TypeFlags.Void) !== 0 + ) + ) +} + +function neverRead(checker: ts.TypeChecker, declaration: ts.VariableDeclaration): boolean { + if (!ts.isIdentifier(declaration.name)) return false + const symbol = checker.getSymbolAtLocation(declaration.name) + let read = false + const visit = (node: ts.Node): void => { + if (read) return + if (ts.isIdentifier(node) && node !== declaration.name) { + read = checker.getSymbolAtLocation(node) === symbol + } + ts.forEachChild(node, visit) + } + visit(declaration.getSourceFile()) + return !read +} + +function floatingShape(checker: ts.TypeChecker, call: ts.CallExpression): string | null { + const { node, caught } = consumer(call) + if (caught) return null + const parent = node.parent + if (ts.isExpressionStatement(parent)) return 'statement' + if (ts.isVoidExpression(parent)) return 'void' + if (ts.isArrowFunction(parent) && parent.body === node && discardedByCallback(checker, parent)) { + return 'callback' + } + if (ts.isVariableDeclaration(parent) && parent.initializer === node) { + return neverRead(checker, parent) ? 'unread' : null + } + return null +} + +function exempt(checker: ts.TypeChecker, call: ts.CallExpression): boolean { + const declaration = checker.getResolvedSignature(call)?.getDeclaration() + if (!declaration) return false + const name = ts.getNameOfDeclaration(declaration)?.getText() + return NEVER_REJECTS.has(`${basename(declaration.getSourceFile().fileName)}#${name}`) +} + +function census(program: ts.Program, files: string[]): { floating: string[]; seen: number } { + const checker = program.getTypeChecker() + const floating: string[] = [] + let seen = 0 + for (const file of files) { + const source = program.getSourceFile(file) + if (!source) throw new Error(`not in program: ${file}`) + const visit = (node: ts.Node): void => { + // A chain step is judged through the call at its base. + if (ts.isCallExpression(node) && !isChainStep(node.expression) && isPromise(checker, node)) { + seen += 1 + const shape = floatingShape(checker, node) + if (shape && !exempt(checker, node)) { + const { line } = source.getLineAndCharacterOfPosition(node.getStart(source)) + floating.push(`${shape} ${basename(file)}:${line + 1} ${node.expression.getText(source)}`) + } + } + ts.forEachChild(node, visit) + } + visit(source) + } + return { floating, seen } +} + +function probeCensus(text: string): string[] { + const name = join(sourceDirectory, 'floating-promise-probe.ts') + const host = ts.createCompilerHost(COMPILER_OPTIONS) + const getSourceFile = host.getSourceFile.bind(host) + host.getSourceFile = (fileName, language) => + fileName === name + ? ts.createSourceFile(fileName, text, language) + : getSourceFile(fileName, language) + const fileExists = host.fileExists.bind(host) + host.fileExists = (fileName) => fileName === name || fileExists(fileName) + return census(ts.createProgram([name], COMPILER_OPTIONS, host), [name]).floating.map( + (entry) => entry.split(' ')[0]! + ) +} + +describe('floating promises', () => { + it('finds none in production code', () => { + const files = productionFiles() + const result = census(ts.createProgram(files, COMPILER_OPTIONS), files) + // Resolution worked: a broken program would see no promises and pass vacuously. + expect(result.seen).toBeGreaterThan(500) + expect(result.floating).toEqual([]) + }, 120_000) + + it('flags each floating shape and accepts the handled ones', () => { + const prelude = `declare const db: { query(sql: string): Promise } + async function wrapper(): Promise { await db.query('x') }\n` + const floating = [ + `db.query('x')`, + `void db.query('x')`, + `void wrapper()`, + `db.query('x').then(() => 1)`, + `setTimeout(() => db.query('x'), 1)`, + `export function f() { const pending = db.query('x') }` + ] + for (const statement of floating) { + expect({ statement, shapes: probeCensus(prelude + statement) }).toEqual({ + statement, + shapes: [expect.any(String)] + }) + } + const handled = [ + `export async function f() { await db.query('x') }`, + `void db.query('x').catch(() => undefined)`, + `void db.query('x').then(() => 1, () => 2)`, + `export async function f() { const pending = db.query('x'); await pending }`, + `export const g = () => db.query('x')` + ] + for (const statement of handled) { + expect({ statement, shapes: probeCensus(prelude + statement) }).toEqual({ + statement, + shapes: [] + }) + } + }) +}) diff --git a/cloud/apps/relay/src/index.ts b/cloud/apps/relay/src/index.ts index 9c34dfb7673..0121619878c 100644 --- a/cloud/apps/relay/src/index.ts +++ b/cloud/apps/relay/src/index.ts @@ -33,7 +33,8 @@ const database = await openRelayDatabaseAtBoot({ databaseUrl: config.databaseUrl, dataDir: config.dataDir, poolMax: config.databasePoolMax, - applicationName: `orca-relay/${config.role}/${config.cellId}` + applicationName: `orca-relay/${config.role}/${config.cellId}`, + appliesPostgresSchema: config.role !== 'cell' }) await reconcileCellAdmissionAtStartup(config, new RelayAssignmentStore(database)) const { @@ -158,7 +159,7 @@ const shutdown = (): void => { heartbeat?.stop() regionalRehomeWorker?.stop() sessions.drain(0) - server.close(() => void database.close()) + server.close(() => void database.close().catch(() => undefined)) } process.once('SIGTERM', shutdown) process.once('SIGINT', shutdown) diff --git a/cloud/apps/relay/src/observed-relay-database.ts b/cloud/apps/relay/src/observed-relay-database.ts index 4720cae6c98..5149c1c2dd9 100644 --- a/cloud/apps/relay/src/observed-relay-database.ts +++ b/cloud/apps/relay/src/observed-relay-database.ts @@ -29,10 +29,20 @@ export function observeRelayDatabase( error instanceof Error && error.message === 'database_lock_unavailable' ) + const commitWithFinal = database.commitWithFinal?.bind(database) return { dialect: database.dialect, query, queryLocked, + ...(commitWithFinal + ? { + commitWithFinal: (sql: string, params?: unknown[]): Promise => + timedRelayOperation( + () => commitWithFinal(sql, params), + (durationMs, success) => observer.recordSql(durationMs, success) + ) + } + : {}), transaction: async ( operation: (transaction: RelayDatabase) => Promise, options?: RelayTransactionOptions diff --git a/cloud/apps/relay/src/postgres-commit-with-final-write-postgres.test.ts b/cloud/apps/relay/src/postgres-commit-with-final-write-postgres.test.ts new file mode 100644 index 00000000000..43fa340ca35 --- /dev/null +++ b/cloud/apps/relay/src/postgres-commit-with-final-write-postgres.test.ts @@ -0,0 +1,203 @@ +import { afterAll, beforeAll, beforeEach, describe, expect, it, vi } from 'vitest' +import pg from 'pg' +import { + commitWithFinalWrite, + openRelayDatabase, + type RelayDatabase +} from './database.js' + +const databaseUrl = process.env.ORCA_RELAY_TEST_POSTGRES_URL +const describePostgres = databaseUrl ? describe : describe.skip + +const table = 'relay_commit_final_write_test' +const counterUpdate = `UPDATE ${table} SET n = n + ? + WHERE id = ? AND (? <= 0 OR n + ? <= cap) + RETURNING id` + +function counterParams(id: string, delta: number): unknown[] { + return [delta, id, delta, delta] +} + +// The counter write and COMMIT travel as one simple-query message. These pin the +// outcomes callers rely on: a server error means COMMIT never ran, and only a +// lost connection leaves the outcome unknown. +describePostgres('PostgreSQL counter write committed in the same round trip', () => { + let database: RelayDatabase + let other: RelayDatabase + let admin: pg.Client + + beforeAll(async () => { + database = await openRelayDatabase({ databaseUrl, dataDir: '' }) + other = await openRelayDatabase({ databaseUrl, dataDir: '' }) + admin = new pg.Client({ connectionString: databaseUrl }) + await admin.connect() + await admin.query( + `CREATE TABLE IF NOT EXISTS ${table} (id TEXT PRIMARY KEY, n BIGINT NOT NULL, cap BIGINT NOT NULL)` + ) + await admin.query( + `CREATE TABLE IF NOT EXISTS ${table}_log (id TEXT NOT NULL, note TEXT NOT NULL)` + ) + }) + + beforeEach(async () => { + await admin.query(`DELETE FROM ${table}`) + await admin.query(`DELETE FROM ${table}_log`) + await admin.query(`INSERT INTO ${table} (id, n, cap) VALUES ('a', 0, 2), ('b', 0, 2), ('it''s', 0, 2)`) + }) + + afterAll(async () => { + await admin.query(`DROP TABLE IF EXISTS ${table}`) + await admin.query(`DROP TABLE IF EXISTS ${table}_log`) + await admin.end() + await database.close() + await other.close() + }) + + async function counter(id: string): Promise { + return Number((await admin.query(`SELECT n FROM ${table} WHERE id = $1`, [id])).rows[0].n) + } + + async function logged(): Promise { + return Number((await admin.query(`SELECT count(*) AS c FROM ${table}_log`)).rows[0].c) + } + + it('commits the earlier statements and the counter together in one message', async () => { + const sent = vi.spyOn(pg.Client.prototype, 'query') + let messagesForCommit = 0 + const result = await database.transaction(async (transaction) => { + await transaction.query(`INSERT INTO ${table}_log (id, note) VALUES (?, ?)`, ['a', 'x']) + const before = sent.mock.calls.length + const committed = await commitWithFinalWrite(transaction, counterUpdate, counterParams("it's", 1)) + messagesForCommit = sent.mock.calls.length - before + return committed + }) + const texts = sent.mock.calls.map(([text]) => String(text)) + sent.mockRestore() + expect(result).toBe(true) + expect(messagesForCommit).toBe(1) + expect(texts.filter((text) => text === 'COMMIT')).toEqual([]) + expect(texts.at(-1)).toMatch(/; COMMIT$/) + expect(await counter("it's")).toBe(1) + expect(await logged()).toBe(1) + }) + + it('rolls back over cap and lets the caller read the row outside the transaction', async () => { + await admin.query(`UPDATE ${table} SET n = 2 WHERE id = 'a'`) + let attempts = 0 + const read = await database.transaction(async (transaction) => { + attempts += 1 + await transaction.query(`INSERT INTO ${table}_log (id, note) VALUES (?, ?)`, ['a', 'x']) + if (await commitWithFinalWrite(transaction, counterUpdate, counterParams('a', 1))) { + return 'committed' + } + // Already rolled back, so this runs outside the aborted transaction. + const rows = await transaction.query(`SELECT id FROM ${table} WHERE id = ?`, ['a']) + return rows.length > 0 ? 'over-cap' : 'missing' + }) + expect(read).toBe('over-cap') + expect(attempts).toBe(1) + expect(await counter('a')).toBe(2) + expect(await logged()).toBe(0) + }) + + it('reports a missing row the same way', async () => { + const read = await database.transaction(async (transaction) => { + await transaction.query(`INSERT INTO ${table}_log (id, note) VALUES (?, ?)`, ['z', 'x']) + if (await commitWithFinalWrite(transaction, counterUpdate, counterParams('z', -1))) { + return 'committed' + } + const rows = await transaction.query(`SELECT id FROM ${table} WHERE id = ?`, ['z']) + return rows.length > 0 ? 'over-cap' : 'missing' + }) + expect(read).toBe('missing') + expect(await logged()).toBe(0) + }) + + it('retries a deadlock or lock timeout raised by the fused message and commits once', async () => { + let attempts = 0 + let theirAttempts = 0 + let otherHolds!: () => void + const otherHolding = new Promise((resolve) => (otherHolds = resolve)) + let ourHold!: () => void + const weHold = new Promise((resolve) => (ourHold = resolve)) + const theirs = other.transaction(async (transaction) => { + theirAttempts += 1 + await transaction.query(`UPDATE ${table} SET n = n WHERE id = 'b'`) + otherHolds() + await weHold + // Waits for 'a', which the first attempt below holds: one side deadlocks. + await transaction.query(`UPDATE ${table} SET n = n WHERE id = 'a'`) + }, { reportRetries: false }) + const ours = database.transaction(async (transaction) => { + attempts += 1 + await transaction.query(`UPDATE ${table} SET n = n WHERE id = 'a'`) + await otherHolding + ourHold() + return await commitWithFinalWrite(transaction, counterUpdate, counterParams('b', 1)) + }, { reportRetries: false }) + const [mine, their] = await Promise.allSettled([ours, theirs]) + // Whichever side PostgreSQL picks as the victim retries and then succeeds. + expect(mine.status).toBe('fulfilled') + expect(their.status).toBe('fulfilled') + expect(await counter('b')).toBe(1) + expect(attempts + theirAttempts).toBeGreaterThanOrEqual(3) + }, 20_000) + + it('retries a lock timeout raised by the fused message', async () => { + let attempts = 0 + const blocker = new pg.Client({ connectionString: databaseUrl }) + await blocker.connect() + try { + await blocker.query('BEGIN') + await blocker.query(`SELECT n FROM ${table} WHERE id = 'a' FOR UPDATE`) + const ours = database.transaction(async (transaction) => { + attempts += 1 + if (attempts === 2) await blocker.query('COMMIT') + return await commitWithFinalWrite(transaction, counterUpdate, counterParams('a', 1)) + }, { reportRetries: false }) + await expect(ours).resolves.toBe(true) + } finally { + await blocker.end() + } + expect(attempts).toBe(2) + expect(await counter('a')).toBe(1) + }, 20_000) + + it('never retries when the connection is lost under the fused message', async () => { + let attempts = 0 + const blocker = new pg.Client({ connectionString: databaseUrl }) + await blocker.connect() + try { + await blocker.query('BEGIN') + await blocker.query(`SELECT n FROM ${table} WHERE id = 'a' FOR UPDATE`) + const ours = database.transaction(async (transaction) => { + attempts += 1 + const pid = Number((await transaction.query('SELECT pg_backend_pid() AS pid'))[0]!.pid) + // Ends the backend while the fused message waits on the row lock. + setTimeout(() => { + void admin.query('SELECT pg_terminate_backend($1)', [pid]) + }, 200) + return await commitWithFinalWrite(transaction, counterUpdate, counterParams('a', 1)) + }) + await expect(ours).rejects.toThrow() + } finally { + await blocker.query('ROLLBACK').catch(() => undefined) + await blocker.end() + } + expect(attempts).toBe(1) + expect(await counter('a')).toBe(0) + // The pool still serves after dropping the dead client. + await expect(database.query('SELECT 1 AS one')).resolves.toEqual([{ one: 1 }]) + }, 20_000) + + it('refuses a statement after the fused commit', async () => { + await expect( + database.transaction(async (transaction) => { + await commitWithFinalWrite(transaction, counterUpdate, counterParams('a', 1)) + await transaction.query(`INSERT INTO ${table}_log (id, note) VALUES (?, ?)`, ['a', 'late']) + }) + ).rejects.toThrow('postgres_transaction_already_committed') + expect(await counter('a')).toBe(1) + expect(await logged()).toBe(0) + }) +}) diff --git a/cloud/apps/relay/src/postgres-lock-wait-sample.test.ts b/cloud/apps/relay/src/postgres-lock-wait-sample.test.ts new file mode 100644 index 00000000000..9b3d5627c04 --- /dev/null +++ b/cloud/apps/relay/src/postgres-lock-wait-sample.test.ts @@ -0,0 +1,28 @@ +import { describe, expect, it } from 'vitest' +import type { RelayDatabase } from './database.js' +import { readPostgresLockWaitSample } from './postgres-lock-wait-sample.js' + +function sampled(rows: Array>): RelayDatabase { + const database: RelayDatabase = { + query: async () => rows, + queryLocked: async () => rows, + transaction: async (operation) => await operation(database), + close: async () => undefined + } + return database +} + +describe('lock-wait sample roles', () => { + it('keeps combined-role sessions as their own class and folds unknown roles into other', async () => { + const sample = await readPostgresLockWaitSample( + sampled([ + { waiter_role: 'combined', waited_table: 'relay_cells', holder_role: 'combined', waiters: 2 }, + { waiter_role: 'psql', waited_table: 'relay_cells', holder_role: null, waiters: 1 } + ]) + ) + expect(sample).toEqual([ + { waiterRole: 'combined', table: 'relay_cells', holderRole: 'combined', waiters: 2 }, + { waiterRole: 'other', table: 'relay_cells', holderRole: 'other', waiters: 1 } + ]) + }) +}) diff --git a/cloud/apps/relay/src/postgres-lock-wait-sample.ts b/cloud/apps/relay/src/postgres-lock-wait-sample.ts index 0d9a46d69af..f26c26d9b86 100644 --- a/cloud/apps/relay/src/postgres-lock-wait-sample.ts +++ b/cloud/apps/relay/src/postgres-lock-wait-sample.ts @@ -37,7 +37,7 @@ LEFT JOIN pg_stat_activity holder ON holder.pid = root.pid WHERE w.application_name LIKE 'orca-relay/%' GROUP BY 1, 2, 3` -const RELAY_ROLES = new Set(['director', 'cell']) +const RELAY_ROLES = new Set(['director', 'cell', 'combined']) export async function readPostgresLockWaitSample( database: RelayDatabase diff --git a/cloud/apps/relay/src/regional-rehome-target-row-lock-postgres.test.ts b/cloud/apps/relay/src/regional-rehome-target-row-lock-postgres.test.ts index c61c6b3ff0f..ba7c0fde0d5 100644 --- a/cloud/apps/relay/src/regional-rehome-target-row-lock-postgres.test.ts +++ b/cloud/apps/relay/src/regional-rehome-target-row-lock-postgres.test.ts @@ -263,6 +263,66 @@ describePostgres('PostgreSQL regional rehome target-row lock', () => { await expectNothingCommitted(context) }) + // A placement that reserves on the target between the candidate read and the target-row + // statement: modelled by raising the target's enforced units right before that statement. + async function fillTargetBefore( + context: Awaited>, + freeSeats: number + ): Promise { + control.beforeTrip = async (sql) => { + if (!sql.includes('WITH target AS')) return + const foreign = await observer.query( + `SELECT COUNT(*) AS count FROM relay_control_connection_reservations + WHERE cell_id = ? AND state IN ('reserved', 'late-arrival-debt', 'claimed') + AND user_id <> ?`, + [context.target.id, context.identity.userId] + ) + // Headroom: enforced + outstanding + unobserved < hard cap - reserved host controls. + const enforced = 1_000 - 100 - 60 - Number(foreign[0]!.count) - freeSeats + await observer.query( + `UPDATE relay_cell_connection_snapshots SET enforced_connection_units = ? WHERE cell_id = ?`, + [enforced, context.target.id] + ) + } + } + + it('defers when a placement takes the target last connection seat after selection', async () => { + const context = await fixture() + const request = await context.select() + await fillTargetBefore(context, 0) + control.enabled = true + + const result = await context.delayedStore.commitIdleRegionalRehome(request, context.safety()) + control.enabled = false + + expect(result).toEqual({ outcome: 'deferred', reason: 'candidate-ineligible' }) + await expectNothingCommitted(context) + }) + + it('admits into the last connection seat, not counting its own reservation', async () => { + const context = await fixture() + const request = await context.select() + await fillTargetBefore(context, 1) + control.enabled = true + + const result = await context.delayedStore.commitIdleRegionalRehome(request, context.safety()) + control.enabled = false + + expect(result).toEqual({ outcome: 'committed' }) + }) + + it('admits a target with no connection limits row', async () => { + const context = await fixture() + const request = await context.select() + await primary.query(`DELETE FROM relay_cell_connection_limits WHERE cell_id = ?`, [ + context.target.id + ]) + + const result = await context.delayedStore.commitIdleRegionalRehome(request, context.safety()) + + expect(result).toEqual({ outcome: 'committed' }) + }) + async function expectNothingCommitted(context: Awaited>) { expect(await commitCounts(context.identity.userId)).toEqual({ attempts: 0, migrations: 0 }) expect(await reservedRequests(context.target.id)).toBe(context.targetReservedBefore) diff --git a/cloud/apps/relay/src/relay-fix-level.ts b/cloud/apps/relay/src/relay-fix-level.ts new file mode 100644 index 00000000000..8deb9b3d81b --- /dev/null +++ b/cloud/apps/relay/src/relay-fix-level.ts @@ -0,0 +1,5 @@ +// Bump by one in any change that fixes a cell crash or a cell safety bug. Every runtime +// metrics line carries it, and an alert pages on a serving cell left below the newest +// level any cell reports, so a fleet that silently kept an unfixed image gets caught. +// Images from before this constant report nothing, which the alert reads as below. +export const RELAY_FIX_LEVEL = 1 diff --git a/cloud/apps/relay/src/relay-observability.test.ts b/cloud/apps/relay/src/relay-observability.test.ts index 04051519074..a993225acfe 100644 --- a/cloud/apps/relay/src/relay-observability.test.ts +++ b/cloud/apps/relay/src/relay-observability.test.ts @@ -9,6 +9,7 @@ import { RelayObservability, type RelayProcessCounts } from './relay-observability.js' +import { RELAY_FIX_LEVEL } from './relay-fix-level.js' const counts: RelayProcessCounts = { totalConnections: 9, @@ -57,6 +58,20 @@ function renameStageKeys(bucket: unknown): unknown { } describe('relay observability', () => { + it('stamps every runtime metrics line with the fix level', () => { + const entries: Array> = [] + const observability = new RelayObservability( + { role: 'cell', cellId: 'production-gce-c25', region: 'asia-east2' }, + (entry) => entries.push(entry) + ) + observability.flush(counts) + expect(entries[0]).toMatchObject({ + event: 'orca_relay_runtime_metrics', + fixLevel: RELAY_FIX_LEVEL + }) + expect(RELAY_FIX_LEVEL).toBeGreaterThanOrEqual(1) + }) + it('emits safe readiness dependency outcomes', () => { const entries: Array> = [] const observability = new RelayObservability( diff --git a/cloud/apps/relay/src/relay-observability.ts b/cloud/apps/relay/src/relay-observability.ts index bc7f03cc1d5..3638de91549 100644 --- a/cloud/apps/relay/src/relay-observability.ts +++ b/cloud/apps/relay/src/relay-observability.ts @@ -5,6 +5,7 @@ import type { ControlRenewalFlush } from './control-renewal-batch.js' import type { CellInventoryHoldCounts } from './cell-inventory-hold-samples.js' import type { PostgresPoolPressureCounts } from './postgres-pool-pressure.js' import type { RelayReadinessGraceEvent, RelayReadinessObservation } from './relay-readiness.js' +import { RELAY_FIX_LEVEL } from './relay-fix-level.js' export type RelayRuntimeCounts = { totalConnections: number @@ -480,6 +481,7 @@ export class RelayObservability implements RelayRuntimeObserver { message: 'Orca Relay runtime metrics', event: 'orca_relay_runtime_metrics', metricVersion: 2, + fixLevel: RELAY_FIX_LEVEL, role: this.identity.role, cellId: this.identity.cellId, region: this.identity.region, diff --git a/cloud/dev/fixtures/terraform-root-partition/families.json b/cloud/dev/fixtures/terraform-root-partition/families.json index 57d7043aeff..11e9db811a3 100644 --- a/cloud/dev/fixtures/terraform-root-partition/families.json +++ b/cloud/dev/fixtures/terraform-root-partition/families.json @@ -134,12 +134,14 @@ "google_iam_workload_identity_pool_provider.github_relay_asia_topology", "google_iam_workload_identity_pool_provider.github_staging_relay_capacity", "google_iam_workload_identity_pool_provider.github_staging_relay_deploy", + "google_logging_metric.relay_cell_fix_level", "google_logging_metric.relay_incident", "google_logging_metric.relay_long_cell_inventory_hold", "google_logging_metric.relay_snapshot", "google_monitoring_alert_policy.relay_assignment_5xx", "google_monitoring_alert_policy.relay_assignment_edge_429", "google_monitoring_alert_policy.relay_cell_control_rtt", + "google_monitoring_alert_policy.relay_cell_outdated_fix_level", "google_monitoring_alert_policy.relay_cell_process_exit", "google_monitoring_alert_policy.relay_cloud_nat_port_drops", "google_monitoring_alert_policy.relay_cloud_sql_backends", diff --git a/cloud/dev/scripts/measure-relay-drain-disconnect-gap.mjs b/cloud/dev/scripts/measure-relay-drain-disconnect-gap.mjs new file mode 100644 index 00000000000..9487c24b21f --- /dev/null +++ b/cloud/dev/scripts/measure-relay-drain-disconnect-gap.mjs @@ -0,0 +1,161 @@ +import { readFileSync } from 'node:fs' +import { pathToFileURL } from 'node:url' + +// Per-desktop disconnect gap for one drained cell, read from log lines that already ship: +// the source cell's `control closed` line and the director's reconnect `assignment granted` +// line (drains send `reconnect: true`, so every drained desktop's grant carries `hinted=true`). +// Read-only: the operator exports both logs; this file never queries anything. + +const CLOSE = /^\[orca-relay\] control closed host=(\S+) .*?\bsplices=(\d+) .*?\bcode=(\d+) reason=("(?:[^"\\]|\\.)*")/ +const GRANT = /^\[orca-relay\] assignment granted lane=\S+ hinted=true host=(\S+) cell=(\S+)$/ +// The cell's own drain close; a desktop closed this way had not moved yet, idle or not. +const DRAIN_CLOSE_REASON = 'resolve configured director' + +function entryText(entry) { + if (typeof entry.textPayload === 'string') return entry.textPayload + if (typeof entry.jsonPayload?.message === 'string') return entry.jsonPayload.message + return null +} + +function entryTime(entry) { + const at = Date.parse(entry.timestamp) + if (!Number.isFinite(at)) throw new Error('log entry has no timestamp') + return at +} + +export function readDrainCloses(entries) { + const closes = [] + for (const entry of entries) { + const match = CLOSE.exec(entryText(entry) ?? '') + if (!match) continue + closes.push({ + host: match[1], + splices: Number(match[2]), + code: Number(match[3]), + reason: JSON.parse(match[4]), + at: entryTime(entry) + }) + } + // gcloud exports newest first; the measurement needs each host's earliest close. + return closes.sort((left, right) => left.at - right.at) +} + +export function readReconnectGrants(entries) { + const grants = [] + for (const entry of entries) { + const match = GRANT.exec(entryText(entry) ?? '') + if (match) grants.push({ host: match[1], cellId: match[2], at: entryTime(entry) }) + } + return grants.sort((left, right) => left.at - right.at) +} + +function percentile(sorted, fraction) { + if (sorted.length === 0) return null + return sorted[Math.min(sorted.length - 1, Math.ceil(fraction * sorted.length) - 1)] +} + +function summary(values) { + const sorted = [...values].sort((left, right) => left - right) + return { + count: sorted.length, + p50Ms: percentile(sorted, 0.5), + p95Ms: percentile(sorted, 0.95), + maxMs: sorted.at(-1) ?? null + } +} + +// `controlsAtDrainStart` is the denominator: hosts the cell held when the drain began, +// read from its runtime metrics line, so a desktop that never logged a close still counts. +// `drainEndedAt` closes the window: after it the rolled cell's new container serves new sessions, +// and their closes are not part of the drain. Grants after it still count, as the gap's far end. +export function measureDrainDisconnectGap({ + sourceCellId, + drainStartedAt, + drainEndedAt, + controlsAtDrainStart, + closes, + grants, + cellRegions = {} +}) { + if (!Number.isFinite(drainStartedAt)) throw new Error('drain start time is invalid') + if (!Number.isFinite(drainEndedAt) || drainEndedAt <= drainStartedAt) { + throw new Error('drain end time is invalid') + } + const sourceRegion = cellRegions[sourceCellId] + // First close per host after the drain began; later closes are the host's new sessions. + const firstClose = new Map() + for (const close of closes) { + if (close.at < drainStartedAt || close.at > drainEndedAt || firstClose.has(close.host)) continue + firstClose.set(close.host, close) + } + const grantsByHost = new Map() + for (const grant of grants) { + if (grant.at < drainStartedAt || grant.cellId === sourceCellId) continue + grantsByHost.set(grant.host, [...(grantsByHost.get(grant.host) ?? []), grant]) + } + const cutOffGapsMs = [] + const counts = { movedFirst: 0, cutOff: 0, cutOffUnresolved: 0, otherClose: 0 } + const leftRegion = [] + let phoneSessionsDropped = 0 + for (const close of firstClose.values()) { + phoneSessionsDropped += close.splices + const hostGrants = grantsByHost.get(close.host) ?? [] + const first = hostGrants[0] + if (first && first.at <= close.at) counts.movedFirst += 1 + else if (close.reason === DRAIN_CLOSE_REASON) { + if (first) { + counts.cutOff += 1 + cutOffGapsMs.push(first.at - close.at) + } else counts.cutOffUnresolved += 1 + } else counts.otherClose += 1 + if (first && sourceRegion && cellRegions[first.cellId] && cellRegions[first.cellId] !== sourceRegion) { + const back = hostGrants.find( + (grant) => grant.at > first.at && cellRegions[grant.cellId] === sourceRegion + ) + leftRegion.push(back ? back.at - first.at : null) + } + } + const returned = leftRegion.filter((value) => value !== null) + return { + sourceCellId, + controlsAtDrainStart, + closedHosts: firstClose.size, + // Below 1 means some drained desktops left no close line in the export; widen it. + closeCoverage: controlsAtDrainStart > 0 ? firstClose.size / controlsAtDrainStart : null, + ...counts, + // Lower bound until the target cells run an image that logs control activation. + cutOffGap: summary(cutOffGapsMs), + phoneSessionsDropped, + leftSourceRegion: leftRegion.length, + leftSourceRegionStillAway: leftRegion.length - returned.length, + timeUntilBack: summary(returned) + } +} + +function argument(name) { + const index = process.argv.indexOf(`--${name}`) + return index === -1 ? undefined : process.argv[index + 1] +} + +function required(name) { + const value = argument(name) + if (value === undefined) throw new Error(`--${name} is required`) + return value +} + +function main() { + const readJson = (path) => JSON.parse(readFileSync(path, 'utf8')) + const cellRegionsPath = argument('cell-regions') + const result = measureDrainDisconnectGap({ + sourceCellId: required('source-cell'), + drainStartedAt: Date.parse(required('drain-started-at')), + drainEndedAt: Date.parse(required('drain-ended-at')), + controlsAtDrainStart: Number(required('controls')), + closes: readDrainCloses(readJson(required('cell-log'))), + grants: readReconnectGrants(readJson(required('director-log'))), + cellRegions: cellRegionsPath ? readJson(cellRegionsPath) : {} + }) + console.log(JSON.stringify(result, null, 2)) +} + +if (import.meta.url === pathToFileURL(process.argv[1] ?? '').href) main() diff --git a/cloud/dev/scripts/measure-relay-drain-disconnect-gap.test.mjs b/cloud/dev/scripts/measure-relay-drain-disconnect-gap.test.mjs new file mode 100644 index 00000000000..900cd58d7dc --- /dev/null +++ b/cloud/dev/scripts/measure-relay-drain-disconnect-gap.test.mjs @@ -0,0 +1,158 @@ +import assert from 'node:assert/strict' +import { describe, it } from 'node:test' +import { + measureDrainDisconnectGap, + readDrainCloses, + readReconnectGrants +} from './measure-relay-drain-disconnect-gap.mjs' + +const start = Date.parse('2026-10-06T10:00:00Z') +const at = (seconds) => new Date(start + seconds * 1000).toISOString() + +function close(host, seconds, reason, splices = 0) { + return { + timestamp: at(seconds), + jsonPayload: { + message: + `[orca-relay] control closed host=${host} gen=3 state=closed ageMs=100 app="1.4.0"` + + ` splices=${splices} pending=0 code=4001 reason=${JSON.stringify(reason)}` + } + } +} + +function grant(host, seconds, cellId) { + return { + timestamp: at(seconds), + textPayload: `[orca-relay] assignment granted lane=drain-return hinted=true host=${host} cell=${cellId}` + } +} + +describe('drain disconnect gap', () => { + it('splits moved-first, cut-off and unresolved desktops and sums dropped phone sessions', () => { + const closes = readDrainCloses([ + close('h-moved', 30, 'migration completed', 2), + close('h-cut', 300, 'resolve configured director', 1), + close('h-idle', 310, 'resolve configured director'), + close('h-lost', 320, 'resolve configured director'), + close('h-early', -5, 'resolve configured director'), + // A later close of the same host is its next session, not the drain. + close('h-cut', 900, 'resolve configured director', 7) + ]) + const grants = readReconnectGrants([ + grant('h-moved', 20, 'c2'), + grant('h-cut', 304, 'c2'), + grant('h-idle', 330, 'c3'), + grant('h-other', 5, 'c2'), + grant('h-cut', 290, 'c1') + ]) + const result = measureDrainDisconnectGap({ + sourceCellId: 'c1', + drainStartedAt: start, + drainEndedAt: start + 1_200_000, + controlsAtDrainStart: 5, + closes, + grants + }) + assert.equal(result.closedHosts, 4) + assert.equal(result.closeCoverage, 0.8) + assert.equal(result.movedFirst, 1) + assert.equal(result.cutOff, 2) + assert.equal(result.cutOffUnresolved, 1) + assert.equal(result.otherClose, 0) + assert.deepEqual(result.cutOffGap, { count: 2, p50Ms: 4000, p95Ms: 20000, maxMs: 20000 }) + assert.equal(result.phoneSessionsDropped, 3) + }) + + it('times desktops that left the source region until a grant brings them back', () => { + const result = measureDrainDisconnectGap({ + sourceCellId: 'a1', + drainStartedAt: start, + drainEndedAt: start + 1_200_000, + controlsAtDrainStart: 2, + closes: readDrainCloses([ + close('h-away', 100, 'resolve configured director'), + close('h-stuck', 100, 'resolve configured director') + ]), + grants: readReconnectGrants([ + grant('h-away', 101, 'u1'), + grant('h-away', 3701, 'a2'), + grant('h-stuck', 102, 'u1') + ]), + cellRegions: { a1: 'asia-east2', a2: 'asia-east2', u1: 'us-central1' } + }) + assert.equal(result.leftSourceRegion, 2) + assert.equal(result.leftSourceRegionStillAway, 1) + assert.equal(result.timeUntilBack.maxMs, 3_600_000) + }) + + it('keeps each host earliest close when the export is newest first', () => { + const result = measureDrainDisconnectGap({ + sourceCellId: 'c1', + drainStartedAt: start, + drainEndedAt: start + 1_200_000, + controlsAtDrainStart: 1, + closes: readDrainCloses([ + close('h', 300, 'resolve configured director'), + close('h', 60, 'resolve configured director', 3) + ]), + grants: readReconnectGrants([grant('h', 120, 'c2')]) + }) + assert.equal(result.movedFirst, 0) + assert.equal(result.cutOff, 1) + assert.deepEqual(result.cutOffGap, { count: 1, p50Ms: 60000, p95Ms: 60000, maxMs: 60000 }) + assert.equal(result.phoneSessionsDropped, 3) + }) + + it('refuses an invalid drain start time', () => { + assert.throws( + () => + measureDrainDisconnectGap({ + sourceCellId: 'c1', + drainStartedAt: Date.parse('not a time'), + drainEndedAt: start, + controlsAtDrainStart: 1, + closes: [], + grants: [] + }), + /drain start time is invalid/ + ) + }) + + it('leaves out closes after the drain ended and refuses a bad end time', () => { + const input = { + sourceCellId: 'c1', + drainStartedAt: start, + drainEndedAt: start + 600_000, + controlsAtDrainStart: 1, + closes: readDrainCloses([ + close('h-drained', 100, 'resolve configured director'), + // The new container's session after the roll. + close('h-new', 700, '', 4) + ]), + grants: readReconnectGrants([grant('h-drained', 650, 'c2')]) + } + const result = measureDrainDisconnectGap(input) + assert.equal(result.closedHosts, 1) + assert.equal(result.otherClose, 0) + assert.equal(result.phoneSessionsDropped, 0) + // A grant after the window still ends that host's gap. + assert.equal(result.cutOffGap.maxMs, 550_000) + for (const drainEndedAt of [Number.NaN, start]) { + assert.throws( + () => measureDrainDisconnectGap({ ...input, drainEndedAt }), + /drain end time is invalid/ + ) + } + }) + + it('ignores lines that are not the two it reads', () => { + assert.deepEqual(readDrainCloses([{ timestamp: at(0), textPayload: 'unrelated' }]), []) + assert.deepEqual( + readReconnectGrants([grant('h', 0, 'c1')].map((entry) => ({ + ...entry, + textPayload: entry.textPayload.replace('hinted=true', 'hinted=false') + }))), + [] + ) + }) +}) diff --git a/cloud/dev/scripts/relay-same-cap-shadow-gate-verdict.mjs b/cloud/dev/scripts/relay-same-cap-shadow-gate-verdict.mjs index bb40afde202..23ba31a222b 100644 --- a/cloud/dev/scripts/relay-same-cap-shadow-gate-verdict.mjs +++ b/cloud/dev/scripts/relay-same-cap-shadow-gate-verdict.mjs @@ -178,6 +178,17 @@ function ownRetries(sample) { ) } +// A drained host whose redial beats its own release meets its own row. That host was admitted to +// the drain-return lane first (the admission is counted before the assign that is refused), so +// row-busy refusals up to that minute's drain-return admissions are scheduled. Anything beyond is +// row contention the drain does not explain, and stays in the budget. The margin absorbs rounding +// from splitting each 30 s sample across clock minutes. +const ROW_BUSY_MARGIN_PER_MINUTE = 2 + +function rowBusy(sample) { + return sample.assign503sByCauseDelta?.relay_assignment_row_busy ?? 0 +} + /** * The director's scheduled 503s per clock minute, from its runtime-metrics samples: drain-return * deferrals and answers to a host's own early retry, plus the re-placements. Each sample's count is @@ -187,6 +198,7 @@ function ownRetries(sample) { export function drainReturnByMinute(reads, limit) { const deferrals = new Map() const retries = new Map() + const busy = new Map() const assignments = new Map() let retryAfterSecondsMax = 0 let truncated = false @@ -209,6 +221,7 @@ export function drainReturnByMinute(reads, limit) { const endedAt = Date.parse(sample.timestamp) charge(deferrals, endedAt, sample.drainReturnDeferralsDelta ?? 0) charge(retries, endedAt, ownRetries(sample)) + charge(busy, endedAt, rowBusy(sample)) charge(assignments, endedAt, sample.drainReturnAssignmentsDelta ?? 0) retryAfterSecondsMax = Math.max( retryAfterSecondsMax, @@ -217,6 +230,12 @@ export function drainReturnByMinute(reads, limit) { } } const sum = (map) => Math.round([...map.values()].reduce((total, count) => total + count, 0)) + let rowBusyBeyondDrain = 0 + for (const [minute, count] of busy) { + const scheduled = Math.min(count, (assignments.get(minute) ?? 0) + ROW_BUSY_MARGIN_PER_MINUTE) + rowBusyBeyondDrain += count - scheduled + retries.set(minute, (retries.get(minute) ?? 0) + scheduled) + } return { deferralsPerMinute: Object.fromEntries(deferrals), ownRetriesPerMinute: Object.fromEntries(retries), @@ -224,6 +243,7 @@ export function drainReturnByMinute(reads, limit) { deferralsPeakPerMinute: Math.round(Math.max(0, ...deferrals.values())), assignmentsTotal: sum(assignments), assignmentsPeakPerMinute: Math.round(Math.max(0, ...assignments.values())), + rowBusyBeyondDrainTotal: Math.round(rowBusyBeyondDrain), retryAfterSecondsMax, truncated } diff --git a/cloud/dev/scripts/relay-same-cap-shadow-gate.mjs b/cloud/dev/scripts/relay-same-cap-shadow-gate.mjs index 294090639af..9332017f39a 100644 --- a/cloud/dev/scripts/relay-same-cap-shadow-gate.mjs +++ b/cloud/dev/scripts/relay-same-cap-shadow-gate.mjs @@ -199,7 +199,8 @@ const DIRECTOR_DRAIN_FIELDS = [ 'placementRejectionsByReasonDelta', 'drainReturnDeferralsDelta', 'drainReturnAssignmentsDelta', - 'drainReturnRetryAfterSecondsMax' + 'drainReturnRetryAfterSecondsMax', + 'assign503sByCauseDelta' ] // Five instances at one sample per 30 s is ~100 per 10-min sub-window; this many is truncation. diff --git a/cloud/dev/scripts/relay-same-cap-shadow-gate.test.mjs b/cloud/dev/scripts/relay-same-cap-shadow-gate.test.mjs index 09e7539beef..adcf0edfebe 100644 --- a/cloud/dev/scripts/relay-same-cap-shadow-gate.test.mjs +++ b/cloud/dev/scripts/relay-same-cap-shadow-gate.test.mjs @@ -611,6 +611,44 @@ function director503s(minute, count) { return [{ minute503: minute, count }] } +test('row-busy 503s are scheduled only up to the drain-return admissions they ride on', () => { + const drain = drainReturnByMinute([{ + failed: false, + samples: [{ + timestamp: '2026-10-05T20:01:00Z', + drainReturnAssignmentsDelta: 10, + assign503sByCauseDelta: { relay_assignment_row_busy: 8, 'placement-lane': 5 } + }] + }], 1000) + assert.deepEqual(drain.ownRetriesPerMinute, { '2026-10-05T20:00': 8 }) + assert.equal(drain.rowBusyBeyondDrainTotal, 0) + const split = withoutDrainDeferrals( + { perMinute: { '2026-10-05T20:00': 13 } }, + drain, + ['2026-10-05T20:00'] + ) + // The placement-lane refusals stay: only the row-busy ones were scheduled. + assert.deepEqual(split.series, [5]) +}) + +test('row-busy 503s beyond the drain stay in the non-drain budget and fail it', () => { + const minutes = Array.from({ length: 10 }, (_, index) => `2026-10-05T20:0${index}`) + // No drain in the background, and 3 drain-return admissions a minute in the window against 60 + // row-busy refusals: row contention the drain does not explain. + const samples = minutes.map((minute, index) => ({ + timestamp: new Date(Date.parse(`${minute}:30Z`) + 30_000).toISOString(), + drainReturnAssignmentsDelta: index < 5 ? 0 : 3, + assign503sByCauseDelta: { relay_assignment_row_busy: index < 5 ? 0 : 60 } + })) + const drain = drainReturnByMinute([{ failed: false, samples }], 1000) + assert.equal(drain.rowBusyBeyondDrainTotal, 5 * (60 - 3 - 2)) + const perMinute = Object.fromEntries(minutes.map((minute, index) => [minute, index < 5 ? 2 : 60])) + const background = backgroundOf(withoutDrainDeferrals({ perMinute }, drain, minutes.slice(0, 5))) + const observed = withoutDrainDeferrals({ perMinute }, drain, minutes.slice(5)) + assert.deepEqual(observed.series, [55, 55, 55, 55, 55]) + assert.equal(judgeNonDrain503Budget({ observed, background }).status, 'would-block') +}) + test('scheduled 503s come out of the count, split across the minutes they cover', () => { const drain = drainReturnByMinute([{ failed: false, diff --git a/cloud/infra/terraform/environments/production.tfvars b/cloud/infra/terraform/environments/production.tfvars index a3521914b20..2081472fa82 100644 --- a/cloud/infra/terraform/environments/production.tfvars +++ b/cloud/infra/terraform/environments/production.tfvars @@ -499,6 +499,10 @@ relay_region_rehome_source_cell_ids = [ # was otherwise going to strip it from every policy, leaving the alerts firing at nobody. relay_alert_notification_channels = ["projects/onorca-cloud/notificationChannels/4879431412695417284"] +# Cells below this RELAY_FIX_LEVEL page after 6 hours. Raise it with a targeted apply of the +# outdated-image alert once a wave has rolled every serving cell, never mid-wave. +relay_cell_min_fix_level = 1 + # Mobile push gateway. Production is the only environment that runs one; the runtime account, # the three Apple secrets, and their accessor bindings already exist and are imported once # (see docs/push-gateway.md). diff --git a/cloud/infra/terraform/relay-observability.tf b/cloud/infra/terraform/relay-observability.tf index 91601de9805..a0a91a64add 100644 --- a/cloud/infra/terraform/relay-observability.tf +++ b/cloud/infra/terraform/relay-observability.tf @@ -1086,6 +1086,71 @@ resource "google_monitoring_alert_policy" "relay_region_hint_skew" { depends_on = [google_logging_metric.relay_snapshot] } +# Declared for a targeted apply after the step-2 wave. Cells only: a director-only fix must not +# raise the level the cells are held to. sum / count of the distribution is the exact level. +resource "google_logging_metric" "relay_cell_fix_level" { + project = var.project_id + name = "orca_relay_cell_fix_level" + description = "RELAY_FIX_LEVEL reported by each cell's runtime metrics line." + filter = "${local.relay_runtime_log_filter} AND jsonPayload.role=\"cell\" AND jsonPayload.fixLevel:*" + value_extractor = "EXTRACT(jsonPayload.fixLevel)" + label_extractors = { + cell_id = "EXTRACT(jsonPayload.cellId)" + } + + metric_descriptor { + metric_kind = "DELTA" + value_type = "DISTRIBUTION" + unit = "1" + + labels { + key = "cell_id" + value_type = "STRING" + description = "Durable relay cell identifier." + } + } + + bucket_options { + linear_buckets { + num_finite_buckets = 64 + width = 1 + offset = 0 + } + } +} + +locals { + relay_cell_fix_level_hourly = "(sum by (cell_id) (increase(logging_googleapis_com:user_orca_relay_cell_fix_level_sum{monitored_resource=\"gce_instance\"}[1h])) / sum by (cell_id) (increase(logging_googleapis_com:user_orca_relay_cell_fix_level_count{monitored_resource=\"gce_instance\"}[1h])))" + relay_cell_serving = "(sum by (cell_id) (increase(logging_googleapis_com:user_orca_relay_controls_sum{monitored_resource=\"gce_instance\",role=\"cell\"}[1h])) > 0)" +} + +# One PromQL condition per policy, and log-based metrics allow at most 25 h of lookback, so the +# floor is a fixed variable rather than a fleet maximum: raise it with a targeted apply once a +# wave has rolled every serving cell. Images from before the field report no level at all. +resource "google_monitoring_alert_policy" "relay_cell_outdated_fix_level" { + project = var.project_id + display_name = "Orca Relay: cell left on an outdated image" + combiner = "OR" + enabled = true + notification_channels = var.relay_alert_notification_channels + + conditions { + display_name = "Serving cell below fix level ${var.relay_cell_min_fix_level} for 6 hours" + + condition_prometheus_query_language { + query = "((${local.relay_cell_fix_level_hourly} < ${var.relay_cell_min_fix_level}) or (${local.relay_cell_serving} unless on (cell_id) ${local.relay_cell_fix_level_hourly})) and on (cell_id) ${local.relay_cell_serving}" + duration = "21600s" + } + } + + documentation { + content = "A cell holding desktops runs an image below `relay_cell_min_fix_level`, or one too old to report `fixLevel`. On 2026-09-28, 18 cells still ran images without the pg connection-error fix and crashed in a two-minute database failover, dropping ~16.3k hosts. Roll the named cell with a same-capacity roll to the current digest. Empty cells do not fire because they hold no controls. Raise the floor only after a wave has rolled every serving cell." + mime_type = "text/markdown" + } + + depends_on = [google_logging_metric.relay_cell_fix_level] +} + # Why: the four signals that had to be assembled by hand during the 2026-09-04 incident. resource "google_monitoring_dashboard" "relay_incident" { project = var.project_id diff --git a/cloud/infra/terraform/variables.tf b/cloud/infra/terraform/variables.tf index b4360bb6a8a..6be05489c3f 100644 --- a/cloud/infra/terraform/variables.tf +++ b/cloud/infra/terraform/variables.tf @@ -362,6 +362,12 @@ variable "relay_alert_notification_channels" { default = [] } +variable "relay_cell_min_fix_level" { + type = number + description = "Lowest RELAY_FIX_LEVEL a serving relay cell may run before the outdated-image alert fires. Raise it after a wave rolls every serving cell." + default = 1 +} + variable "relay_gce_domain" { type = string description = "Parent DNS name for GCE relay cells; each cell is one exact host below it." diff --git a/cloud/package.json b/cloud/package.json index bb9a853d9a2..272c532a10e 100644 --- a/cloud/package.json +++ b/cloud/package.json @@ -22,7 +22,7 @@ "load:relay:recovery-gate": "node dev/scripts/run-relay-recovery-wave-gate.mjs", "ops:relay": "pnpm --filter @orca-cloud/relay-ops dev", "pretest": "node --test dev/scripts/capture-terraform-plan-baseline.test.mjs dev/scripts/operate-relay-asia-admission.test.mjs dev/scripts/prepare-relay-asia-director-cells.test.mjs dev/scripts/prepare-relay-asia-topology-input.test.mjs dev/scripts/production-cloud-sql-rollout-lock.test.mjs dev/scripts/read-relay-serving-regional-placement-version.test.mjs dev/scripts/relay-asia-rollout-evidence.test.mjs dev/scripts/relay-asia-topology-workflow.test.mjs dev/scripts/relay-cloud-sql-connection-budget.test.mjs dev/scripts/relay-load-reader-evidence.test.mjs dev/scripts/relay-lock-contention-alerts.test.mjs dev/scripts/relay-region-hint-metrics.test.mjs dev/scripts/relay-staging-deploy-identity.test.mjs dev/scripts/sanitize-relay-asia-admission-result.test.mjs dev/scripts/terraform-root-partition.test.mjs dev/scripts/validate-relay-asia-topology-plan.test.mjs ../.github/actions/cloud-sql-rollout-lease/action-contract.test.mjs ../.github/actions/cloud-sql-rollout-lease/storage-lease.test.mjs", - "test": "pnpm -r test && node --test dev/scripts/check-relay-same-cap-headroom.test.mjs dev/scripts/classify-relay-production-capacity-director.test.mjs dev/scripts/classify-relay-staging-bootstrap.test.mjs dev/scripts/deploy-relay-blue-green.test.mjs dev/scripts/deploy-relay-gce-candidate.test.mjs dev/scripts/deploy-relay-gce-multi-target.test.mjs dev/scripts/drive-relay-director-deploy.test.mjs dev/scripts/github-smoke-token.test.mjs dev/scripts/infra.test.mjs dev/scripts/operate-relay-regional-rehome.test.mjs dev/scripts/power-staging-relay.test.mjs dev/scripts/prepare-relay-capacity-canary.test.mjs dev/scripts/prepare-relay-production-capacity-canary.test.mjs dev/scripts/probe-relay-legacy-admission.test.mjs dev/scripts/probe-relay-rehome-trust.test.mjs dev/scripts/production-cell-image-digest-consistency.test.mjs dev/scripts/push-gateway-workflow.test.mjs dev/scripts/push-gateway-recovery.test.mjs dev/scripts/read-relay-production-capacity-identity.test.mjs dev/scripts/relay-admin-endpoint-retry-workflow.test.mjs dev/scripts/relay-admin-transient-retry.test.mjs dev/scripts/relay-admission-selector.test.mjs dev/scripts/relay-gce-terraform-fence.test.mjs dev/scripts/relay-load-connection-failure.test.mjs dev/scripts/relay-load-control-peer.test.mjs dev/scripts/relay-load-director-capacity-gate.test.mjs dev/scripts/relay-load-model.test.mjs dev/scripts/relay-load-phase-barrier.test.mjs dev/scripts/relay-load-placement-boundary.test.mjs dev/scripts/relay-load-profile.test.mjs dev/scripts/relay-load-rebind-boundary.test.mjs dev/scripts/relay-load-region-behavior.test.mjs dev/scripts/relay-load-request-unit-boundary.test.mjs dev/scripts/relay-load-run-lifecycle.test.mjs dev/scripts/relay-monitor-evidence.test.mjs dev/scripts/relay-production-capacity-wave.test.mjs dev/scripts/relay-production-capacity-workflow.test.mjs dev/scripts/relay-production-identity-boundaries.test.mjs dev/scripts/relay-production-same-cap-wave.test.mjs dev/scripts/relay-public-workflow-contract.test.mjs dev/scripts/relay-recovery-wave-gate.test.mjs dev/scripts/relay-region-observation-evidence.test.mjs dev/scripts/relay-rehome-aggregate-evidence.test.mjs dev/scripts/relay-repository.test.mjs dev/scripts/relay-same-cap-job-mode-conditions.test.mjs dev/scripts/relay-same-cap-script-census.test.mjs dev/scripts/relay-same-cap-shadow-gate.test.mjs dev/scripts/relay-staging-c4-refresh-workflow.test.mjs dev/scripts/relay-staging-capacity-identity.test.mjs dev/scripts/staging-relay-apply-guard.test.mjs dev/scripts/validate-relay-capacity-plan.test.mjs dev/scripts/verify-relay-capacity-transition.test.mjs dev/scripts/verify-relay-legacy-bootstrap.test.mjs dev/scripts/workload-identity-attribute-conditions.test.mjs", + "test": "pnpm -r test && node --test dev/scripts/check-relay-same-cap-headroom.test.mjs dev/scripts/classify-relay-production-capacity-director.test.mjs dev/scripts/classify-relay-staging-bootstrap.test.mjs dev/scripts/deploy-relay-blue-green.test.mjs dev/scripts/deploy-relay-gce-candidate.test.mjs dev/scripts/deploy-relay-gce-multi-target.test.mjs dev/scripts/drive-relay-director-deploy.test.mjs dev/scripts/github-smoke-token.test.mjs dev/scripts/infra.test.mjs dev/scripts/measure-relay-drain-disconnect-gap.test.mjs dev/scripts/operate-relay-regional-rehome.test.mjs dev/scripts/power-staging-relay.test.mjs dev/scripts/prepare-relay-capacity-canary.test.mjs dev/scripts/prepare-relay-production-capacity-canary.test.mjs dev/scripts/probe-relay-legacy-admission.test.mjs dev/scripts/probe-relay-rehome-trust.test.mjs dev/scripts/production-cell-image-digest-consistency.test.mjs dev/scripts/push-gateway-workflow.test.mjs dev/scripts/push-gateway-recovery.test.mjs dev/scripts/read-relay-production-capacity-identity.test.mjs dev/scripts/relay-admin-endpoint-retry-workflow.test.mjs dev/scripts/relay-admin-transient-retry.test.mjs dev/scripts/relay-admission-selector.test.mjs dev/scripts/relay-gce-terraform-fence.test.mjs dev/scripts/relay-load-connection-failure.test.mjs dev/scripts/relay-load-control-peer.test.mjs dev/scripts/relay-load-director-capacity-gate.test.mjs dev/scripts/relay-load-model.test.mjs dev/scripts/relay-load-phase-barrier.test.mjs dev/scripts/relay-load-placement-boundary.test.mjs dev/scripts/relay-load-profile.test.mjs dev/scripts/relay-load-rebind-boundary.test.mjs dev/scripts/relay-load-region-behavior.test.mjs dev/scripts/relay-load-request-unit-boundary.test.mjs dev/scripts/relay-load-run-lifecycle.test.mjs dev/scripts/relay-monitor-evidence.test.mjs dev/scripts/relay-production-capacity-wave.test.mjs dev/scripts/relay-production-capacity-workflow.test.mjs dev/scripts/relay-production-identity-boundaries.test.mjs dev/scripts/relay-production-same-cap-wave.test.mjs dev/scripts/relay-public-workflow-contract.test.mjs dev/scripts/relay-recovery-wave-gate.test.mjs dev/scripts/relay-region-observation-evidence.test.mjs dev/scripts/relay-rehome-aggregate-evidence.test.mjs dev/scripts/relay-repository.test.mjs dev/scripts/relay-same-cap-job-mode-conditions.test.mjs dev/scripts/relay-same-cap-script-census.test.mjs dev/scripts/relay-same-cap-shadow-gate.test.mjs dev/scripts/relay-staging-c4-refresh-workflow.test.mjs dev/scripts/relay-staging-capacity-identity.test.mjs dev/scripts/staging-relay-apply-guard.test.mjs dev/scripts/validate-relay-capacity-plan.test.mjs dev/scripts/verify-relay-capacity-transition.test.mjs dev/scripts/verify-relay-legacy-bootstrap.test.mjs dev/scripts/workload-identity-attribute-conditions.test.mjs", "typecheck": "pnpm -r typecheck" }, "devDependencies": {