From 10bea3a8faefc7adf8771f5bfd57a37435ccede4 Mon Sep 17 00:00:00 2001 From: Jinwoo Hong <73622457+Jinwoo-H@users.noreply.github.com> Date: Mon, 5 Oct 2026 17:55:52 -0400 Subject: [PATCH] feat(relay): report lane service times, 503 causes, per-site hold p99 and a director-vs-cell lock-wait split (#25645) * feat(relay): report lane service times, 503 causes, per-site hold p99 and lock-wait split Adds the step-1 observability: drain-return and sticky slot service times, every /v1/assign 503 by cause with a non-drain total, a p99 per lock hold site (drain-return regional rows get their own site), a director-side pg_stat_activity sample that splits lock waiters by director vs cell, and a log line for each reserved_requests drift reconciliation corrects. Log-based metrics for the new fields are declared in Terraform, not applied. * fix(relay): classify a lock waiter by the first relay table its statement names * fix(relay): attribute lock waiters to the root holder; document the director hold alert change Every waiter after the first in a row-lock convoy is blocked by the first waiter, so the sample now walks pg_blocking_pids to the root, reading it once per waiter. The sampler is single-flight. Drain-return regional holds feeding the non-paging director hold policy is documented as expected. * docs(relay): note the lock-wait sample undercounts director waiters when the pool is full --- cloud/apps/relay/src/app.ts | 63 ++++++++++---- ...signment-isolated-cell-replacement.test.ts | 18 ++++ cloud/apps/relay/src/assignment-store.test.ts | 33 ++++++++ cloud/apps/relay/src/assignment-store.ts | 26 ++++-- .../src/cell-inventory-hold-samples.test.ts | 20 +++++ .../relay/src/cell-inventory-hold-samples.ts | 26 +++++- .../relay/src/drain-return-admission.test.ts | 26 +++++- .../apps/relay/src/drain-return-admission.ts | 3 + cloud/apps/relay/src/index.ts | 17 ++++ ...postgres-lock-wait-sample-postgres.test.ts | 76 +++++++++++++++++ .../relay/src/postgres-lock-wait-sample.ts | 57 +++++++++++++ ...ublic-assignment-drain-return-lane.test.ts | 36 +++++++- .../relay/src/relay-observability.test.ts | 78 +++++++++++++++++ cloud/apps/relay/src/relay-observability.ts | 83 +++++++++++++++++++ cloud/apps/relay/src/relay-server.ts | 3 + cloud/docs/orca-relay-operations.md | 22 +++++ cloud/docs/relay-incident-monitor.md | 10 +++ cloud/infra/terraform/relay-observability.tf | 18 +++- 18 files changed, 586 insertions(+), 29 deletions(-) create mode 100644 cloud/apps/relay/src/postgres-lock-wait-sample-postgres.test.ts create mode 100644 cloud/apps/relay/src/postgres-lock-wait-sample.ts diff --git a/cloud/apps/relay/src/app.ts b/cloud/apps/relay/src/app.ts index c3db6291287..e88162dfd88 100644 --- a/cloud/apps/relay/src/app.ts +++ b/cloud/apps/relay/src/app.ts @@ -54,6 +54,7 @@ import type { RelayReadinessDependency } from './relay-readiness.js' import type { AssignmentAdmissionLane, AssignmentAdmissionOutcome, + AssignmentUnavailableCause, RegionalRehomeSafetySnapshot, RelayRuntimeCounts } from './relay-observability.js' @@ -121,6 +122,8 @@ export function createRelayApp( reason: AssignmentAdmissionRejection ) => void recordDrainReturnRetryAfter?: (seconds: number) => void + recordAdmissionServiceMs?: (lane: 'sticky' | 'drain-return', durationMs: number) => void + recordAssignmentUnavailable?: (cause: AssignmentUnavailableCause) => void recordRegionRequest?: (region: RelayRegion | undefined) => void recordRegionSelection?: (input: { targetRegion: RelayRegion @@ -202,11 +205,25 @@ export function createRelayApp( context.header('Retry-After', String(stickyRetryAfterSeconds)) return context.json({ error: 'assignments_temporarily_unavailable' }, 503) } + // The slot's hold, not the request's latency: herd recovery is slots / this. + const timeStickySlot = (lease: { release(): void }): { release(): void } => { + const startedAt = performance.now() + let released = false + return { + release: () => { + if (released) return + released = true + operations.recordAdmissionServiceMs?.('sticky', performance.now() - startedAt) + lease.release() + } + } + } // Classified server-side from the host's own assignment row, never from a // client claim: only a host whose home is roll-isolated reaches this lane. const drainReturnAdmission = new RelayDrainReturnAdmission(publicAssignmentAdmission, { maxConcurrent: Math.max(1, drainReturnConcurrency), - maxRetryAfterSeconds: config.drainReturnMaxRetryAfterSeconds ?? 300 + maxRetryAfterSeconds: config.drainReturnMaxRetryAfterSeconds ?? 300, + onServiceMs: (durationMs) => operations.recordAdmissionServiceMs?.('drain-return', durationMs) }) const deferDrainReturn = (context: Context, retryAfterSeconds: number): Response => { context.header('Retry-After', String(retryAfterSeconds)) @@ -287,7 +304,10 @@ export function createRelayApp( }) app.post('/v1/assign', async (context) => { if (config.role === 'cell') return context.json({ error: 'director_only' }, 404) - if (!config.publicAssignmentsEnabled) return rejectPublicAssignment(context) + if (!config.publicAssignmentsEnabled) { + operations.recordAssignmentUnavailable?.('disabled') + return rejectPublicAssignment(context) + } const bearer = readBearer(context.req.header('authorization')) if (!bearer) return context.json({ error: 'invalid_token' }, 401) const claims = await verifyRelayToken(bearer) @@ -312,11 +332,12 @@ export function createRelayApp( let lane: AssignmentAdmissionLane = 'placement' if (body.data.reconnect) { let rejection: AssignmentAdmissionRejection | undefined - const fastLane = await stickyAssignmentAdmission.acquire(claims.relayHostId, (reason) => { + const stickySlot = await stickyAssignmentAdmission.acquire(claims.relayHostId, (reason) => { rejection = reason }) - if (!fastLane) { + if (!stickySlot) { operations.recordAssignmentAdmission?.('sticky-rejected') + operations.recordAssignmentUnavailable?.('sticky-lane') logAdmissionRejection({ route: 'assign', lane: 'sticky', @@ -326,6 +347,7 @@ export function createRelayApp( }) return rejectStickyAssignment(context) } + const fastLane = timeStickySlot(stickySlot) let verified: ResolvedRelayAssignment | null = null try { verified = await operations.assignments.resolve(identity, { @@ -342,6 +364,7 @@ export function createRelayApp( relayHostId: claims.relayHostId, reason: operationError(error) }) + operations.recordAssignmentUnavailable?.('sticky-verify-database') return rejectStickyAssignment(context) } throw error @@ -356,6 +379,7 @@ export function createRelayApp( operations.recordAssignmentAdmission?.('drain-return-deferred') operations.recordAssignmentRejectionReason?.('drain-return', drainReturn.reason) operations.recordDrainReturnRetryAfter?.(drainReturn.retryAfterSeconds) + operations.recordAssignmentUnavailable?.('drain-return-deferred') logAdmissionRejection({ route: 'assign', lane, @@ -391,6 +415,7 @@ export function createRelayApp( relayHostId: claims.relayHostId, reason: rejection }) + operations.recordAssignmentUnavailable?.('placement-lane') return rejectPublicAssignment(context) } } @@ -430,6 +455,7 @@ export function createRelayApp( }) // The host's own release is settling; its next dial finds the row free. context.header('Retry-After', String(ASSIGNMENT_ROW_BUSY_RETRY_AFTER_SECONDS)) + operations.recordAssignmentUnavailable?.('relay_assignment_row_busy') return context.json({ error: 'assignment_row_busy' }, 503) } if (isRelayAssignmentUnavailableError(error) || isRelayDatabaseTransientError(error)) { @@ -442,13 +468,16 @@ export function createRelayApp( ...homeCellRejectionDetail(error) }) } - if (isRelayAssignmentUnavailableError(error)) { + const unavailable = assignmentUnavailableCause(error) + if (unavailable) { if (lane === 'placement') { operations.recordRegionSelection?.({ targetRegion, fallback: false }) } + operations.recordAssignmentUnavailable?.(unavailable) return context.json({ error: operationError(error) }, 503) } if (isRelayDatabaseTransientError(error)) { + operations.recordAssignmentUnavailable?.('database') return lane === 'placement' ? rejectPublicAssignment(context) : rejectStickyAssignment(context) @@ -2054,16 +2083,22 @@ function logAssignmentRejection(input: { // The home-cell reason is not capacity, but it is the same answer to the client: // retry, the director cannot place you right now. +const RELAY_ASSIGNMENT_UNAVAILABLE_ERRORS = [ + 'relay_capacity_exhausted', + 'relay_connection_headroom_exhausted', + 'relay_home_cell_unavailable', + 'relay_assignment_row_busy' +] as const satisfies readonly AssignmentUnavailableCause[] + +function assignmentUnavailableCause( + error: unknown +): (typeof RELAY_ASSIGNMENT_UNAVAILABLE_ERRORS)[number] | undefined { + if (!(error instanceof Error)) return undefined + return RELAY_ASSIGNMENT_UNAVAILABLE_ERRORS.find((message) => message === error.message) +} + function isRelayAssignmentUnavailableError(error: unknown): boolean { - return ( - error instanceof Error && - [ - 'relay_capacity_exhausted', - 'relay_connection_headroom_exhausted', - 'relay_home_cell_unavailable', - 'relay_assignment_row_busy' - ].includes(error.message) - ) + return assignmentUnavailableCause(error) !== undefined } function homeCellRejectionDetail( diff --git a/cloud/apps/relay/src/assignment-isolated-cell-replacement.test.ts b/cloud/apps/relay/src/assignment-isolated-cell-replacement.test.ts index 21cfcc3cbe6..8c1e9088cf8 100644 --- a/cloud/apps/relay/src/assignment-isolated-cell-replacement.test.ts +++ b/cloud/apps/relay/src/assignment-isolated-cell-replacement.test.ts @@ -8,6 +8,7 @@ import { } from './cell-admission-selector.js' import type { RelayCellConfig } from './config.js' import { + consumeRelayCellInventoryHold, openInMemoryRelayDatabase, type RelayDatabase, type RelayLockOptions, @@ -365,6 +366,23 @@ describe('re-placing a host off a cell isolated for a roll', () => { } }) + // Why its own site: the drain-return lane's service time is either this lock or + // the inventory, and only a separate p99 can say which. + it('samples the regional target rows it locks under their own hold site', async () => { + const { store, database, isolateForRoll } = await setup() + const first = await store.assign(IDENTITY, 'us-central1') + await isolateForRoll(first.cellId) + consumeRelayCellInventoryHold(database) + + await store.assign(IDENTITY, 'us-central1') + + expect(consumeRelayCellInventoryHold(database)).toMatchObject({ + cellInventoryHolds: 1, + cellInventoryHoldMaxSite: 'isolated-replacement', + isolatedReplacementHolds: 1 + }) + }) + it('stops re-placing once restore clears the stamp', async () => { const { store, isolateForRoll, restore, rollIsolatedAt } = await setup() const first = await store.assign(IDENTITY, 'us-central1') diff --git a/cloud/apps/relay/src/assignment-store.test.ts b/cloud/apps/relay/src/assignment-store.test.ts index 9b889a133a8..39eeaa0f924 100644 --- a/cloud/apps/relay/src/assignment-store.test.ts +++ b/cloud/apps/relay/src/assignment-store.test.ts @@ -2476,6 +2476,39 @@ describe('RelayAssignmentStore', () => { }) }) + it('logs the drift reconciliation corrects, once, after it commits', async () => { + const store = await setup(() => 100, [ + { id: 'cell-a', url: 'https://relay-a.example.com', capacityRequests: 10 }, + { id: 'cell-b', url: 'https://relay-b.example.com', capacityRequests: 10 } + ]) + const grant = await store.assign({ userId: 'user-a', relayHostId: 'host000000000001' }) + await database!.query(`UPDATE relay_cells SET reserved_requests = 5 WHERE cell_id = ?`, [ + grant.cellId + ]) + const warn = vi.spyOn(console, 'warn').mockImplementation(() => {}) + + let lines: string[] + try { + await store.cellEvacuationStatus('cell-a', 'cell-b', true) + await store.cellEvacuationStatus('cell-a', 'cell-b', true) + } finally { + lines = warn.mock.calls.map((call) => String(call[0])) + warn.mockRestore() + } + + const drift = lines + .filter((line) => line.includes('orca_relay_reservation_drift')) + .map((line) => JSON.parse(line) as unknown) + expect(drift).toEqual([ + { + event: 'orca_relay_reservation_drift', + cellId: grant.cellId, + reservedRequests: 5, + leaseUnits: 1 + } + ]) + }) + it('blocks aggregate fenced completion on assignment accounting mismatch', async () => { let now = 100 const cells = [ diff --git a/cloud/apps/relay/src/assignment-store.ts b/cloud/apps/relay/src/assignment-store.ts index fac9c514916..2eaf143940b 100644 --- a/cloud/apps/relay/src/assignment-store.ts +++ b/cloud/apps/relay/src/assignment-store.ts @@ -69,6 +69,7 @@ import { REGIONAL_REHOME_DEFAULT_HOST_COOLDOWN_MS } from './database.js' import type { RelayCellConfig } from './config.js' +import type { CellLockHoldSite } from './cell-inventory-hold-samples.js' import type { RelayDatabase, RelayLockOptions, @@ -7374,7 +7375,7 @@ export class RelayAssignmentStore { targetCellId: string ): Promise { const now = this.now() - await this.database.transaction(async (transaction) => { + const drift = await this.database.transaction(async (transaction) => { // Reconciliation takes the same assignment→activity→cell order as live // mutations so correcting drift never races a credential or socket lease. const assignments = await transaction.queryLocked( @@ -7494,6 +7495,7 @@ export class RelayAssignmentStore { ) } + const corrected: { cellId: string; reservedRequests: number; leaseUnits: number }[] = [] for (const row of cells) { const cellId = text(row, 'cell_id') const expected = cellUnits.get(cellId) ?? 0 @@ -7505,8 +7507,18 @@ export class RelayAssignmentStore { `UPDATE relay_cells SET reserved_requests = ?, updated_at = ? WHERE cell_id = ?`, [expected, now, cellId] ) + corrected.push({ + cellId, + reservedRequests: integer(row, 'reserved_requests'), + leaseUnits: expected + }) } + return corrected }) + // Logged after COMMIT so a retried transaction reports each correction once. + for (const sample of drift) { + console.warn(JSON.stringify({ event: 'orca_relay_reservation_drift', ...sample })) + } } private async lockCellInventory( @@ -7528,12 +7540,13 @@ export class RelayAssignmentStore { // row-lock order), keeps them off the fleet-wide lock without a cycle. // The wait policy follows the caller for the same reason the inventory lock's // does: a sweep must not fail terminally on ordinary contention. Hold time is - // deliberately not sampled here — the metric tracks the fleet-wide lock these - // rows replace, and mixing in short single-row holds would flatter it. + // sampled only for a caller that names its site: short single-row holds under + // the inventory label would flatter the fleet-wide lock they replace. private async lockCellRows( database: RelayDatabase, cellIds: string[], - mode: CellInventoryLockMode = 'request' + mode: CellInventoryLockMode = 'request', + holdSite?: CellLockHoldSite ): Promise { const distinct = [...new Set(cellIds)] const { measureHoldMs: _sampled, ...wait } = cellInventoryLockOptions(mode) @@ -7541,7 +7554,7 @@ export class RelayAssignmentStore { `SELECT * FROM relay_cells WHERE cell_id IN (${distinct.map(() => '?').join(', ')}) ORDER BY cell_id ASC`, distinct, - wait + holdSite ? { ...wait, measureHoldMs: true, holdSite } : wait ) } @@ -7618,7 +7631,8 @@ export class RelayAssignmentStore { return await this.lockCellRows( database, candidates.map((row) => text(row, 'cell_id')), - mode + mode, + 'isolated-replacement' ) } diff --git a/cloud/apps/relay/src/cell-inventory-hold-samples.test.ts b/cloud/apps/relay/src/cell-inventory-hold-samples.test.ts index 61c7ef1ba9d..1e6fd86d2a9 100644 --- a/cloud/apps/relay/src/cell-inventory-hold-samples.test.ts +++ b/cloud/apps/relay/src/cell-inventory-hold-samples.test.ts @@ -102,6 +102,26 @@ describe('cell inventory hold samples', () => { expect(samples.consumeCounts()).toEqual(emptyCellInventoryHoldCounts()) }) + // Why: the shared max and p95 mix every lock, so only the per-site p99 says + // whether a drain return waits on its own regional rows or on the inventory. + it('reports a p99 per hold site', () => { + const samples = new CellInventoryHoldSamples() + for (let ms = 1; ms <= 100; ms++) { + samples.record(ms) + samples.record(ms * 2, 'isolated-replacement') + } + samples.record(7, 'rehome-target-row') + + expect(samples.readCounts()).toMatchObject({ + cellInventoryHolds: 201, + cellInventoryHoldMaxSite: 'isolated-replacement', + inventoryHoldMsP99: 99, + isolatedReplacementHoldMsP99: 198, + isolatedReplacementHolds: 100, + rehomeTargetRowHoldMsP99: 7 + }) + }) + // Why: the alert reads one max across every lock, so the label is the only // thing that says whether a long hold was the inventory or a rehome target row. it('names the site of the longest hold and reports rehome target rows apart', () => { diff --git a/cloud/apps/relay/src/cell-inventory-hold-samples.ts b/cloud/apps/relay/src/cell-inventory-hold-samples.ts index 3454c97b55a..f8fed1bb349 100644 --- a/cloud/apps/relay/src/cell-inventory-hold-samples.ts +++ b/cloud/apps/relay/src/cell-inventory-hold-samples.ts @@ -3,7 +3,8 @@ // hold distribution, and no runtime metric carried it before this change. // Which lock a hold sample came from. Every site feeds the same max, so one // alert on cellInventoryHoldMsMax covers them all; the label names the holder. -export type CellLockHoldSite = 'inventory' | 'rehome-target-row' +// 'isolated-replacement' is a drain return's regional target rows. +export type CellLockHoldSite = 'inventory' | 'rehome-target-row' | 'isolated-replacement' export type CellInventoryHoldCounts = { cellInventoryHoldMsMax: number @@ -14,6 +15,11 @@ export type CellInventoryHoldCounts = { // reads its bound and its presence here. rehomeTargetRowHoldMsMax: number rehomeTargetRowHolds: number + // Per-site p99: the shared max and p95 cannot say which lock a drain waits on. + inventoryHoldMsP99: number + rehomeTargetRowHoldMsP99: number + isolatedReplacementHoldMsP99: number + isolatedReplacementHolds: number // Why: a failed acquisition produces no hold sample, so the hold fields alone // read healthy while the lock is saturated. Split by wait policy, not by // caller: fail-fast covers background sweeps that step aside by design AND @@ -36,6 +42,10 @@ export function emptyCellInventoryHoldCounts(): CellInventoryHoldCounts { cellInventoryHoldMaxSite: 'none', rehomeTargetRowHoldMsMax: 0, rehomeTargetRowHolds: 0, + inventoryHoldMsP99: 0, + rehomeTargetRowHoldMsP99: 0, + isolatedReplacementHoldMsP99: 0, + isolatedReplacementHolds: 0, cellInventoryLockUnavailable: 0, cellInventoryLockTimeouts: 0 } @@ -79,19 +89,29 @@ export class CellInventoryHoldSamples { if (this.samples.length === 0) return { ...emptyCellInventoryHoldCounts(), ...failures } const sorted = [...this.samples].sort((left, right) => left.holdMs - right.holdMs) const max = sorted[sorted.length - 1]! - const rehome = sorted.filter((sample) => sample.site === 'rehome-target-row') + const site = (name: CellLockHoldSite) => sorted.filter((sample) => sample.site === name) + const rehome = site('rehome-target-row') + const isolated = site('isolated-replacement') return { cellInventoryHoldMsMax: round(max.holdMs), - cellInventoryHoldMsP95: round(sorted[Math.ceil(0.95 * sorted.length) - 1]?.holdMs ?? 0), + cellInventoryHoldMsP95: nearestRank(sorted, 0.95), cellInventoryHolds: sorted.length, cellInventoryHoldMaxSite: max.site, rehomeTargetRowHoldMsMax: round(rehome[rehome.length - 1]?.holdMs ?? 0), rehomeTargetRowHolds: rehome.length, + inventoryHoldMsP99: nearestRank(site('inventory'), 0.99), + rehomeTargetRowHoldMsP99: nearestRank(rehome, 0.99), + isolatedReplacementHoldMsP99: nearestRank(isolated, 0.99), + isolatedReplacementHolds: isolated.length, ...failures } } } +function nearestRank(sorted: { holdMs: number }[], rank: number): number { + return round(sorted[Math.ceil(rank * sorted.length) - 1]?.holdMs ?? 0) +} + function round(value: number): number { return Number(value.toFixed(3)) } diff --git a/cloud/apps/relay/src/drain-return-admission.test.ts b/cloud/apps/relay/src/drain-return-admission.test.ts index 33ee05eace1..560d4515ca0 100644 --- a/cloud/apps/relay/src/drain-return-admission.test.ts +++ b/cloud/apps/relay/src/drain-return-admission.test.ts @@ -9,7 +9,13 @@ import { RelayPublicAssignmentAdmission } from './public-assignment-admission.js type Timer = { at: number; callback: () => void; cancelled: boolean } -function harness(overrides: { maxQueued?: number; maxRetryAfterSeconds?: number } = {}) { +function harness( + overrides: { + maxQueued?: number + maxRetryAfterSeconds?: number + onServiceMs?: (durationMs: number) => void + } = {} +) { let now = 0 const timers: Timer[] = [] const placement = new RelayPublicAssignmentAdmission({ @@ -33,7 +39,8 @@ function harness(overrides: { maxQueued?: number; maxRetryAfterSeconds?: number const lane = new RelayDrainReturnAdmission(placement, { maxConcurrent: 1, maxRetryAfterSeconds: overrides.maxRetryAfterSeconds ?? 300, - now: () => now + now: () => now, + onServiceMs: overrides.onServiceMs }) return { lane, @@ -133,6 +140,21 @@ describe('drain-return admission', () => { expect(Math.max(...retries)).toBe(3) }) + // The unclamped sample: the EWMA's floor and ceiling would hide the real tail. + it('reports each slot hold once, as measured', async () => { + const samples: number[] = [] + const { lane, advance } = harness({ onServiceMs: (ms) => samples.push(ms) }) + const fast = admitted(await lane.acquire(host(1))) + advance(5) + fast.lease.release() + fast.lease.release() + const slow = admitted(await lane.acquire(host(2))) + advance(20_000) + slow.lease.release() + + expect(samples).toEqual([5, 20_000]) + }) + it('answers a host’s own early retry with its interval, not a place behind the cohort', async () => { const { lane } = harness() admitted(await lane.acquire(host(0))).lease.release() diff --git a/cloud/apps/relay/src/drain-return-admission.ts b/cloud/apps/relay/src/drain-return-admission.ts index 955b7a3b640..5b9989b6b5f 100644 --- a/cloud/apps/relay/src/drain-return-admission.ts +++ b/cloud/apps/relay/src/drain-return-admission.ts @@ -44,6 +44,8 @@ export class RelayDrainReturnAdmission { maxConcurrent: number maxRetryAfterSeconds: number now?: () => number + // The raw sample behind the EWMA, so the lane's service time is reported. + onServiceMs?: (durationMs: number) => void } ) {} @@ -108,6 +110,7 @@ export class RelayDrainReturnAdmission { } private recordService(durationMs: number): void { + this.options.onServiceMs?.(durationMs) const sample = Math.min(SERVICE_MS_CEILING, Math.max(SERVICE_MS_FLOOR, durationMs)) this.serviceMs += SERVICE_EWMA_WEIGHT * (sample - this.serviceMs) } diff --git a/cloud/apps/relay/src/index.ts b/cloud/apps/relay/src/index.ts index d7d37ce209b..9c34dfb7673 100644 --- a/cloud/apps/relay/src/index.ts +++ b/cloud/apps/relay/src/index.ts @@ -18,6 +18,7 @@ import { } from './database.js' import { runAssignmentCleanup } from './assignment-cleanup-steps.js' import { runRelayBackgroundOperation } from './relay-background-operation.js' +import { readPostgresLockWaitSample } from './postgres-lock-wait-sample.js' import { jitteredSweepIntervalMs } from './relay-sweep-schedule.js' import { observedRelayRequests } from './relay-observability.js' import { startRegionalRehomeWorker } from './regional-rehome-worker.js' @@ -87,8 +88,23 @@ const migrationInventoryTimer = roleOwnsAssignmentMaintenance(config.role) }, '[orca-relay] migration inventory failed') }, 5 * 60_000) : null +// Directors only: one role's view covers every backend, and cells roll separately. +// Single-flight, so a slow database never stacks samples on the 3-slot pool. +let lockWaitSampling = false +const lockWaitSampleTimer = roleOwnsAssignmentMaintenance(config.role) + ? setInterval(() => { + if (database.dialect !== 'postgres' || lockWaitSampling) return + lockWaitSampling = true + void runRelayBackgroundOperation(async () => { + observability.recordDatabaseLockWaitSample(await readPostgresLockWaitSample(database)) + }, '[orca-relay] lock wait sample failed').finally(() => { + lockWaitSampling = false + }) + }, 5_000) + : null cleanupTimer?.unref() assignmentCleanupTimer?.unref() +lockWaitSampleTimer?.unref() inventorySnapshotTimer?.unref() migrationInventoryTimer?.unref() observability.start(() => ({ @@ -137,6 +153,7 @@ const shutdown = (): void => { if (assignmentCleanupTimer) clearInterval(assignmentCleanupTimer) if (inventorySnapshotTimer) clearInterval(inventorySnapshotTimer) if (migrationInventoryTimer) clearInterval(migrationInventoryTimer) + if (lockWaitSampleTimer) clearInterval(lockWaitSampleTimer) observability.stop() heartbeat?.stop() regionalRehomeWorker?.stop() diff --git a/cloud/apps/relay/src/postgres-lock-wait-sample-postgres.test.ts b/cloud/apps/relay/src/postgres-lock-wait-sample-postgres.test.ts new file mode 100644 index 00000000000..072dfcf1988 --- /dev/null +++ b/cloud/apps/relay/src/postgres-lock-wait-sample-postgres.test.ts @@ -0,0 +1,76 @@ +import { afterAll, beforeAll, describe, expect, it } from 'vitest' +import { RelayAssignmentStore } from './assignment-store.js' +import { openRelayDatabase, type RelayDatabase } from './database.js' +import { readPostgresLockWaitSample } from './postgres-lock-wait-sample.js' + +const databaseUrl = process.env.ORCA_RELAY_TEST_POSTGRES_URL +const describePostgres = databaseUrl ? describe : describe.skip + +const cell = { + id: 'lock-wait-sample-cell', + url: 'https://lock-wait-sample.example.com', + capacityRequests: 10 +} + +// The split rests on application_name surviving into pg_stat_activity for both +// the waiter and its blocker, which only a real server shows. +describePostgres('PostgreSQL lock-wait sample', () => { + let cellDatabase: RelayDatabase + let directorDatabase: RelayDatabase + + beforeAll(async () => { + cellDatabase = await openRelayDatabase({ + databaseUrl, + dataDir: '', + applicationName: `orca-relay/cell/${cell.id}` + }) + directorDatabase = await openRelayDatabase({ + databaseUrl, + dataDir: '', + applicationName: 'orca-relay/director/director' + }) + await new RelayAssignmentStore(directorDatabase).reconcileCells([cell]) + }) + + afterAll(async () => { + await directorDatabase?.query(`DELETE FROM relay_cells WHERE cell_id = ?`, [cell.id]) + await cellDatabase?.close() + await directorDatabase?.close() + }) + + // Every waiter after the first is blocked by the first waiter's tuple lock, so + // only the root of the chain names the cell transaction that holds the row. + it('attributes a convoy of directors to the cell holding the relay_cells row', async () => { + expect(await readPostgresLockWaitSample(directorDatabase)).toEqual([]) + + let release!: () => void + const released = new Promise((resolve) => (release = resolve)) + let held!: () => void + const holding = new Promise((resolve) => (held = resolve)) + const holder = cellDatabase.transaction(async (transaction) => { + await transaction.query( + `UPDATE relay_cells SET reserved_requests = reserved_requests WHERE cell_id = ?`, + [cell.id] + ) + held() + await released + }) + await holding + const waiters = Array.from({ length: 3 }, () => + directorDatabase.transaction(async (transaction) => { + await transaction.queryLocked(`SELECT * FROM relay_cells WHERE cell_id = ?`, [cell.id]) + }) + ) + + try { + await expect + .poll(async () => await readPostgresLockWaitSample(directorDatabase), { timeout: 900 }) + .toEqual([ + { waiterRole: 'director', table: 'relay_cells', holderRole: 'cell', waiters: 3 } + ]) + } finally { + release() + await Promise.all([holder, ...waiters]) + } + }) +}) diff --git a/cloud/apps/relay/src/postgres-lock-wait-sample.ts b/cloud/apps/relay/src/postgres-lock-wait-sample.ts new file mode 100644 index 00000000000..0d9a46d69af --- /dev/null +++ b/cloud/apps/relay/src/postgres-lock-wait-sample.ts @@ -0,0 +1,57 @@ +import type { RelayDatabase } from './database.js' +import type { DatabaseLockWaitSample } from './relay-observability.js' + +// Every relay process connects as the same user through a socket, so neither +// pg_stat_statements nor Query Insights can tell a director's lock wait from a +// cell's. application_name (`orca-relay//`) can, so this samples it. +// The holder is the root of the wait chain: in a row-lock convoy every later +// waiter is blocked by the first waiter, not by the transaction holding the row. +// Blockers are read once per waiter; non-relay waiters stay in so chains through +// them resolve, and the depth cap bounds a cycle. The table is the first relay +// table the waiting statement names, which may be one it only references. +// It shares the director's 3-slot pool, so when every slot is a lock waiter the +// sample queues and undercounts director waiters; a dedicated connection fixes that. +const LOCK_WAIT_SAMPLE_SQL = ` +WITH RECURSIVE waiting AS MATERIALIZED ( + SELECT pid, application_name, query, (pg_blocking_pids(pid))[1] AS blocker + FROM pg_stat_activity + WHERE datname = current_database() AND wait_event_type = 'Lock' +), chain AS ( + SELECT pid AS waiter, blocker AS pid, 1 AS depth FROM waiting + UNION ALL + SELECT chain.waiter, waiting.blocker, chain.depth + 1 + FROM chain JOIN waiting ON waiting.pid = chain.pid + WHERE chain.depth < 8 +), root AS ( + SELECT DISTINCT ON (waiter) waiter, pid FROM chain ORDER BY waiter, depth DESC +) +SELECT split_part(w.application_name, '/', 2) AS waiter_role, + COALESCE( + substring(w.query FROM '\\m(relay_cells|relay_assignments)\\M'), + 'other' + ) AS waited_table, + split_part(holder.application_name, '/', 2) AS holder_role, + COUNT(*) AS waiters +FROM waiting w JOIN root ON root.waiter = w.pid +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']) + +export async function readPostgresLockWaitSample( + database: RelayDatabase +): Promise { + const rows = await database.query(LOCK_WAIT_SAMPLE_SQL) + return rows.map((row) => ({ + waiterRole: relayRole(row['waiter_role']), + table: String(row['waited_table']), + holderRole: relayRole(row['holder_role']), + waiters: Number(row['waiters']) + })) +} + +// Anything else (an operator session, a finished holder) stays one bounded key. +function relayRole(value: unknown): string { + return typeof value === 'string' && RELAY_ROLES.has(value) ? value : 'other' +} diff --git a/cloud/apps/relay/src/public-assignment-drain-return-lane.test.ts b/cloud/apps/relay/src/public-assignment-drain-return-lane.test.ts index 11ded0fa0e8..cd1344dc5b6 100644 --- a/cloud/apps/relay/src/public-assignment-drain-return-lane.test.ts +++ b/cloud/apps/relay/src/public-assignment-drain-return-lane.test.ts @@ -37,12 +37,14 @@ describe('drain-return lane', () => { relayHostId === drained ? isolatedHome(relayHostId) : assignment('cell-o', relayHostId) ) const outcomes: string[] = [] + const servedLanes: string[] = [] const app = createRelayApp(config(), { store: {} as never, assignments: { assign, resolve } as never, drain: vi.fn(), ready: vi.fn(async () => true), - recordAssignmentAdmission: (outcome) => outcomes.push(outcome) + recordAssignmentAdmission: (outcome) => outcomes.push(outcome), + recordAdmissionServiceMs: (lane) => servedLanes.push(lane) }) // The re-placement holds the drain lane; the sticky lane's only slot is free. @@ -61,6 +63,8 @@ describe('drain-return lane', () => { replacement.resolve(assignment('cell-new', drained)) expect((await moving).status).toBe(200) expect(outcomes).toEqual(['drain-return', 'sticky']) + // The drained host held a sticky slot for its verification read too. + expect(servedLanes).toEqual(['sticky', 'sticky', 'drain-return']) }) it('defers an overflowing drain return with a paced Retry-After and says why', async () => { @@ -73,6 +77,7 @@ describe('drain-return lane', () => { const outcomes: string[] = [] const reasons: string[] = [] const retryAfters: number[] = [] + const causes: string[] = [] const app = createRelayApp(config({ drainReturnQueueMax: 1, drainReturnWaitMs: 50 }), { store: {} as never, assignments: { assign, resolve } as never, @@ -80,7 +85,8 @@ describe('drain-return lane', () => { ready: vi.fn(async () => true), recordAssignmentAdmission: (outcome) => outcomes.push(outcome), recordAssignmentRejectionReason: (lane, reason) => reasons.push(`${lane}:${reason}`), - recordDrainReturnRetryAfter: (seconds) => retryAfters.push(seconds) + recordDrainReturnRetryAfter: (seconds) => retryAfters.push(seconds), + recordAssignmentUnavailable: (cause) => causes.push(cause) }) const warn = vi.spyOn(console, 'warn').mockImplementation(() => {}) @@ -103,6 +109,7 @@ describe('drain-return lane', () => { expect(reasons).toEqual(['drain-return:queue-full', 'drain-return:wait-timeout']) expect(retryAfters).toEqual([2, Number(timedOut.headers.get('retry-after'))]) expect(outcomes.filter((outcome) => outcome === 'drain-return-deferred')).toHaveLength(2) + expect(causes).toEqual(['drain-return-deferred', 'drain-return-deferred']) expect(warn.mock.calls.flat().join('\n')).toMatch( /lane=drain-return hinted=true reason=queue-full .*retryAfter=2/ ) @@ -115,12 +122,14 @@ describe('drain-return lane', () => { const assign = vi.fn(async () => assignment('cell-s', host)) const resolve = vi.fn(async () => assignment('cell-s', host)) const outcomes: string[] = [] + const causes: string[] = [] const app = createRelayApp(config(), { store: {} as never, assignments: { assign, resolve } as never, drain: vi.fn(), ready: vi.fn(async () => true), - recordAssignmentAdmission: (outcome) => outcomes.push(outcome) + recordAssignmentAdmission: (outcome) => outcomes.push(outcome), + recordAssignmentUnavailable: (cause) => causes.push(cause) }) expect( @@ -131,6 +140,27 @@ describe('drain-return lane', () => { expect(repeat.status).toBe(503) expect(repeat.headers.get('retry-after')).toBe('2') expect(outcomes).toEqual(['sticky', 'sticky-rejected']) + expect(causes).toEqual(['sticky-lane']) + }) + + it('names the cause of a placement that found no capacity', async () => { + vi.spyOn(console, 'warn').mockImplementation(() => {}) + const assign = vi.fn(async () => { + throw new Error('relay_capacity_exhausted') + }) + const causes: string[] = [] + const app = createRelayApp(config(), { + store: {} as never, + assignments: { assign, resolve: vi.fn() } as never, + drain: vi.fn(), + ready: vi.fn(async () => true), + recordAssignmentUnavailable: (cause) => causes.push(cause) + }) + + const response = await app.request('/v1/assign', assignmentRequest('cccccccccccccccc')) + + expect(response.status).toBe(503) + expect(causes).toEqual(['relay_capacity_exhausted']) }) it('never classifies an unhinted request, which keeps the placement lane', async () => { diff --git a/cloud/apps/relay/src/relay-observability.test.ts b/cloud/apps/relay/src/relay-observability.test.ts index 9881c54808d..04051519074 100644 --- a/cloud/apps/relay/src/relay-observability.test.ts +++ b/cloud/apps/relay/src/relay-observability.test.ts @@ -200,6 +200,84 @@ describe('relay observability', () => { }) }) + it('splits non-drain 503s from scheduled drain-return deferrals', () => { + const entries: Array> = [] + const observability = new RelayObservability( + { role: 'director', cellId: 'director', region: 'us-central1' }, + (entry) => entries.push(entry) + ) + for (let index = 0; index < 40; index++) { + observability.recordAssignmentUnavailable('drain-return-deferred') + } + observability.recordAssignmentUnavailable('sticky-lane') + observability.recordAssignmentUnavailable('relay_capacity_exhausted') + observability.recordAssignmentUnavailable('relay_capacity_exhausted') + observability.flush(counts) + observability.flush(counts) + + expect(entries[0]).toMatchObject({ + assign503sByCauseDelta: { + 'drain-return-deferred': 40, + 'sticky-lane': 1, + relay_capacity_exhausted: 2 + }, + assignNonDrain503sDelta: 3 + }) + expect(entries[1]).toMatchObject({ assign503sByCauseDelta: {}, assignNonDrain503sDelta: 0 }) + }) + + it('reports lane slot service times only for windows that served the lane', () => { + const entries: Array> = [] + const observability = new RelayObservability( + { role: 'director', cellId: 'director', region: 'us-central1' }, + (entry) => entries.push(entry) + ) + for (let ms = 1; ms <= 100; ms++) { + observability.recordAdmissionServiceMs('sticky', ms) + observability.recordAdmissionServiceMs('drain-return', ms * 10) + } + observability.flush(counts) + observability.flush(counts) + + expect(entries[0]).toMatchObject({ + stickyServiceMsP50: 50, + stickyServiceMsP99: 99, + drainReturnServiceMsP50: 500, + drainReturnServiceMsP95: 950 + }) + expect(entries[1]).not.toHaveProperty('stickyServiceMsP50') + expect(entries[1]).not.toHaveProperty('drainReturnServiceMsP50') + }) + + it('sums lock waiters per role and table over the samples taken', () => { + const entries: Array> = [] + const observability = new RelayObservability( + { role: 'director', cellId: 'director', region: 'us-central1' }, + (entry) => entries.push(entry) + ) + observability.recordDatabaseLockWaitSample([ + { waiterRole: 'cell', table: 'relay_cells', holderRole: 'cell', waiters: 3 }, + { waiterRole: 'director', table: 'relay_cells', holderRole: 'cell', waiters: 1 }, + { waiterRole: 'cell', table: 'relay_assignments', holderRole: 'director', waiters: 2 } + ]) + observability.recordDatabaseLockWaitSample([]) + observability.recordDatabaseLockWaitSample([ + { waiterRole: 'cell', table: 'relay_cells', holderRole: 'cell', waiters: 1 } + ]) + observability.flush(counts) + + expect(entries[0]).toMatchObject({ + dbLockWaitSamplesDelta: 3, + dbLockWaitersByKeyDelta: { + 'cell:relay_cells:cell': 4, + 'director:relay_cells:cell': 1, + 'cell:relay_assignments:director': 2 + }, + dbCellRowLockWaitersDirectorDelta: 1, + dbCellRowLockWaitersCellDelta: 4 + }) + }) + it('aggregates coarse region requests, selections, fallbacks, and outages', () => { 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 16ee785a3f6..bc7f03cc1d5 100644 --- a/cloud/apps/relay/src/relay-observability.ts +++ b/cloud/apps/relay/src/relay-observability.ts @@ -49,6 +49,29 @@ export type AssignmentAdmissionOutcome = export type AssignmentAdmissionLane = 'sticky' | 'placement' | 'drain-return' +// Why every /v1/assign 503 has a cause: drain-return deferrals are scheduled +// 503s, so a gate on all 503s trips on every fast drain. +export type AssignmentUnavailableCause = + | 'disabled' + | 'sticky-lane' + | 'sticky-verify-database' + | 'placement-lane' + | 'drain-return-deferred' + | 'database' + | 'relay_capacity_exhausted' + | 'relay_connection_headroom_exhausted' + | 'relay_home_cell_unavailable' + | 'relay_assignment_row_busy' + +// Row-lock waiters seen in one pg_stat_activity sample, by who waits, on which +// table, behind whom. Roles come from each pool's application_name. +export type DatabaseLockWaitSample = { + waiterRole: string + table: string + holderRole: string + waiters: number +}[] + export interface RelayRuntimeObserver { recordAuth(success: boolean): void recordForwardedBytes(bytes: number): void @@ -61,6 +84,9 @@ export interface RelayRuntimeObserver { recordAssignmentAdmission?(outcome: AssignmentAdmissionOutcome): void recordAssignmentRejectionReason?(lane: AssignmentAdmissionLane, reason: string): void recordDrainReturnRetryAfter?(seconds: number): void + recordAdmissionServiceMs?(lane: 'sticky' | 'drain-return', durationMs: number): void + recordAssignmentUnavailable?(cause: AssignmentUnavailableCause): void + recordDatabaseLockWaitSample?(sample: DatabaseLockWaitSample): void recordRegionRequest?(region: RelayRegion | undefined): void recordRegionSelection?(input: { targetRegion: RelayRegion @@ -115,6 +141,11 @@ type RelayMetricDeltas = { drainReturnRejectionsByReason: Record drainReturnRetryAfterSecondsMax: number drainReturnRetryAfterSecondsSum: number + drainReturnServiceMs: number[] + stickyServiceMs: number[] + assign503sByCause: Record + dbLockWaitSamples: number + dbLockWaitersByKey: Record requestedRegions: Record selectedRegions: Record regionFallbacks: Record @@ -161,6 +192,11 @@ const emptyDeltas = (): RelayMetricDeltas => ({ drainReturnRejectionsByReason: {}, drainReturnRetryAfterSecondsMax: 0, drainReturnRetryAfterSecondsSum: 0, + drainReturnServiceMs: [], + stickyServiceMs: [], + assign503sByCause: {}, + dbLockWaitSamples: 0, + dbLockWaitersByKey: {}, requestedRegions: {}, selectedRegions: {}, regionFallbacks: {}, @@ -279,6 +315,23 @@ export class RelayObservability implements RelayRuntimeObserver { this.deltas.drainReturnRetryAfterSecondsSum += seconds } + recordAdmissionServiceMs(lane: 'sticky' | 'drain-return', durationMs: number): void { + if (lane === 'sticky') this.deltas.stickyServiceMs.push(durationMs) + else this.deltas.drainReturnServiceMs.push(durationMs) + } + + recordAssignmentUnavailable(cause: AssignmentUnavailableCause): void { + increment(this.deltas.assign503sByCause, cause) + } + + recordDatabaseLockWaitSample(sample: DatabaseLockWaitSample): void { + this.deltas.dbLockWaitSamples++ + for (const row of sample) { + const key = `${row.waiterRole}:${row.table}:${row.holderRole}` + this.deltas.dbLockWaitersByKey[key] = (this.deltas.dbLockWaitersByKey[key] ?? 0) + row.waiters + } + } + recordRegionRequest(region: RelayRegion | undefined): void { increment(this.deltas.requestedRegions, region ?? 'unhinted') } @@ -446,6 +499,30 @@ export class RelayObservability implements RelayRuntimeObserver { drainReturnRejectionsByReasonDelta: deltas.drainReturnRejectionsByReason, drainReturnRetryAfterSecondsMax: deltas.drainReturnRetryAfterSecondsMax, drainReturnRetryAfterSecondsSum: deltas.drainReturnRetryAfterSecondsSum, + // Slot hold times; a lane serves concurrency / service time per second. + // Omitted when empty for the same reason as the accept percentiles below. + ...(deltas.drainReturnServiceMs.length === 0 + ? {} + : { + drainReturnServiceMsP50: roundMs(percentile(deltas.drainReturnServiceMs, 0.5)), + drainReturnServiceMsP95: roundMs(percentile(deltas.drainReturnServiceMs, 0.95)) + }), + ...(deltas.stickyServiceMs.length === 0 + ? {} + : { + stickyServiceMsP50: roundMs(percentile(deltas.stickyServiceMs, 0.5)), + stickyServiceMsP99: roundMs(percentile(deltas.stickyServiceMs, 0.99)) + }), + assign503sByCauseDelta: deltas.assign503sByCause, + assignNonDrain503sDelta: Object.entries(deltas.assign503sByCause) + .filter(([cause]) => cause !== 'drain-return-deferred') + .reduce((sum, [, count]) => sum + count, 0), + // Waiters summed over samples; divided by samples it is the mean number of + // backends waiting, i.e. lock-wait seconds per second. + dbLockWaitSamplesDelta: deltas.dbLockWaitSamples, + dbLockWaitersByKeyDelta: deltas.dbLockWaitersByKey, + dbCellRowLockWaitersDirectorDelta: lockWaiters(deltas.dbLockWaitersByKey, 'director'), + dbCellRowLockWaitersCellDelta: lockWaiters(deltas.dbLockWaitersByKey, 'cell'), requestedRegionsDelta: deltas.requestedRegions, selectedRegionsDelta: deltas.selectedRegions, ...regionCounterFields('requestedRegion', deltas.requestedRegions), @@ -526,6 +603,12 @@ function regionCounterFields( ) } +function lockWaiters(byKey: Record, waiterRole: string): number { + return Object.entries(byKey) + .filter(([key]) => key.startsWith(`${waiterRole}:relay_cells:`)) + .reduce((sum, [, waiters]) => sum + waiters, 0) +} + function increment(counts: Record, key: string): void { counts[key] = (counts[key] ?? 0) + 1 } diff --git a/cloud/apps/relay/src/relay-server.ts b/cloud/apps/relay/src/relay-server.ts index 57e4eac2c99..b957443a35f 100644 --- a/cloud/apps/relay/src/relay-server.ts +++ b/cloud/apps/relay/src/relay-server.ts @@ -168,6 +168,9 @@ export function createRelayServer( recordAssignmentRejectionReason: (lane, reason) => observability.recordAssignmentRejectionReason?.(lane, reason), recordDrainReturnRetryAfter: (seconds) => observability.recordDrainReturnRetryAfter?.(seconds), + recordAdmissionServiceMs: (lane, durationMs) => + observability.recordAdmissionServiceMs?.(lane, durationMs), + recordAssignmentUnavailable: (cause) => observability.recordAssignmentUnavailable?.(cause), recordRegionRequest: (region) => observability.recordRegionRequest?.(region), recordRegionSelection: (input) => observability.recordRegionSelection?.(input) }) diff --git a/cloud/docs/orca-relay-operations.md b/cloud/docs/orca-relay-operations.md index fc56af8656e..bab64731cee 100644 --- a/cloud/docs/orca-relay-operations.md +++ b/cloud/docs/orca-relay-operations.md @@ -513,6 +513,28 @@ pause through `applyRegionalRehomeControl` waits 1 s for that row, 3 attempts, so it can fail during one commit: retry a pause that fails once, and do not treat that as a fault. +Each hold site also reports its own p99: `inventoryHoldMsP99`, +`rehomeTargetRowHoldMsP99` and `isolatedReplacementHoldMsP99`, the last being +the regional target rows a drain return locks; a site with no holds in the +interval reports 0, so read these filtered to values above 0. Directors also +report lane slot times (`stickyServiceMsP50/P99`, `drainReturnServiceMsP50/P95`, +omitted when the lane served nothing; a sticky slot held only for a drain-return +or unverified host's verification read counts too, so its p50 moves with that +share), every `/v1/assign` 503 by cause in +`assign503sByCauseDelta`, and `assignNonDrain503sDelta`, which leaves out the +scheduled `drain-return-deferred` answers a fast drain produces by design. +Every 5 s each director samples `pg_stat_activity` for relay backends waiting +on a lock, keyed by waiter role, table and holder role +(`dbLockWaitersByKeyDelta`). The holder is the root of the wait chain, and the +table is the first relay table the waiting statement names. The sample runs +on the director's own 3-connection pool, so when every connection is waiting on +a lock it queues behind them and undercounts director waiters; a dedicated +sampler connection would fix that. Summed waiters divided by `dbLockWaitSamplesDelta` +is the mean number waiting, i.e. lock-wait seconds per second; directors and +cells share one database user, so Query Insights cannot make this split. +Reconciliation logs `orca_relay_reservation_drift` for each cell whose +`reserved_requests` it corrected. + `host-cooldown-ms` is the minimum gap between two rehomes of one host. It bounds the damage from a desktop whose region probe flips: without it the host would be dragged back across the ocean on every flip, since the preference age never expires while the host keeps reconnecting. diff --git a/cloud/docs/relay-incident-monitor.md b/cloud/docs/relay-incident-monitor.md index f193991e3ac..699835b5d0e 100644 --- a/cloud/docs/relay-incident-monitor.md +++ b/cloud/docs/relay-incident-monitor.md @@ -252,6 +252,16 @@ paused. Do not drain or restart cells for a director hold. A cell that has stopped reporting emits no hold sample, so these policies catch lock convoys, not outages. +Since the step-1 observability deploy, director holds also include the +regional target rows a drain return locks (`cellInventoryHoldMaxSite: +isolated-replacement`), so the director policy can fire during a roll's drain. +That is expected and needs no action; `isolatedReplacementHoldMsP99` is the +per-site view. The same deploy adds that lock's NOWAIT refusals and bounded-wait +timeouts to the director's `cellInventoryLockUnavailable` and +`cellInventoryLockTimeouts`, so compare those fields across the deploy only with +that site subtracted; `orca_relay_cloud_sql_lock_timeouts` is unaffected. Drain +returns run only on directors, so the paging cell policy is unchanged. + ## Implementation log - Gave `collector_failed` the same two-consecutive-sample tolerance as an unread diff --git a/cloud/infra/terraform/relay-observability.tf b/cloud/infra/terraform/relay-observability.tf index ba124407564..91601de9805 100644 --- a/cloud/infra/terraform/relay-observability.tf +++ b/cloud/infra/terraform/relay-observability.tf @@ -53,6 +53,10 @@ locals { description = "Relay statements Postgres cancelled after waiting out their lock timeout; a burst means one transaction is holding rows every other relay process needs." filter = "resource.type=\"cloudsql_database\" AND resource.labels.database_id=\"${var.project_id}:${local.relay_database_instance_name}\" AND textPayload:\"db=orca_relay,\" AND textPayload:\"canceling statement due to lock timeout\"" } + reservation_drift = { + description = "Cells whose reserved_requests disagreed with their lease units when reconciliation corrected them; should trend to zero." + filter = "((resource.type=\"cloud_run_revision\" AND (${local.relay_service_log_filter})) OR resource.type=\"gce_instance\") AND jsonPayload.event=\"orca_relay_reservation_drift\"" + } } relay_runtime_metrics = { @@ -108,6 +112,18 @@ locals { cell_inventory_holds = { field = "cellInventoryHolds", description = "Cell-inventory locks acquired in the interval; the percentiles above summarise these." } cell_inventory_lock_unavailable = { field = "cellInventoryLockUnavailable", description = "Fail-fast cell-inventory acquisitions that found the lock held. Includes background sweeps, which step aside by design, so this is contention pressure rather than user-visible failure." } cell_inventory_lock_timeouts = { field = "cellInventoryLockTimeouts", description = "Bounded cell-inventory waits that expired, counted per attempt rather than per request. This is the user-visible lane." } + inventory_hold_ms_p99 = { field = "inventoryHoldMsP99", description = "Cell-inventory lock (every row, or every general row) hold p99 in the interval, this site only." } + rehome_target_row_hold_ms_p99 = { field = "rehomeTargetRowHoldMsP99", description = "Rehome target-row lock hold p99 in the interval, this site only." } + isolated_replacement_hold_ms_p99 = { field = "isolatedReplacementHoldMsP99", description = "Drain-return regional target-row lock hold p99 in the interval, this site only." } + isolated_replacement_holds = { field = "isolatedReplacementHolds", description = "Drain-return regional target-row locks acquired in the interval." } + drain_return_service_ms_p50 = { field = "drainReturnServiceMsP50", description = "Drain-return lane slot hold p50; lane capacity is concurrency / this. Omitted when no drain return ran." } + drain_return_service_ms_p95 = { field = "drainReturnServiceMsP95", description = "Drain-return lane slot hold p95 in the interval." } + sticky_service_ms_p50 = { field = "stickyServiceMsP50", description = "Sticky lane slot hold p50, verification included; herd recovery is concurrency / this." } + sticky_service_ms_p99 = { field = "stickyServiceMsP99", description = "Sticky lane slot hold p99 in the interval." } + assign_non_drain_503s = { field = "assignNonDrain503sDelta", description = "Director /v1/assign 503s other than scheduled drain-return deferrals; the per-cause split is assign503sByCauseDelta." } + db_lock_wait_samples = { field = "dbLockWaitSamplesDelta", description = "Relay-database lock-wait samples taken by this director in the interval; divide the waiter sums below by this." } + db_cell_row_lock_waiters_director = { field = "dbCellRowLockWaitersDirectorDelta", description = "Director backends waiting on a relay_cells row lock, summed over samples." } + db_cell_row_lock_waiters_cell = { field = "dbCellRowLockWaitersCellDelta", description = "Cell backends waiting on a relay_cells row lock, summed over samples." } } # Regions the director can hint or select. Pinned to relay-contract's RELAY_REGIONS by @@ -323,7 +339,7 @@ resource "google_logging_metric" "relay_snapshot" { metric_descriptor { metric_kind = "DELTA" value_type = "DISTRIBUTION" - unit = contains(["sql_latency_ms", "control_rtt_ms_p50", "control_rtt_ms_p95", "control_rtt_ms_max", "client_accept_total_ms_p50", "client_accept_total_ms_p95", "client_accept_total_ms_max", "client_accept_assignment_ms_p95", "client_accept_credential_ms_p95", "client_accept_activity_ms_p95", "client_accept_attach_ms_p95", "client_accept_basis_ms_p95", "control_renewal_latency_ms_p50", "control_renewal_latency_ms_p95", "control_renewal_latency_ms_max", "http_latency_ms", "event_loop_ms_p99", "db_oldest_wait_ms", "db_wait_ms_max"], each.key) ? "ms" : each.key == "queued_bytes" || each.key == "heap_used_bytes" || each.key == "forwarded_bytes" ? "By" : "1" + unit = contains(["sql_latency_ms", "control_rtt_ms_p50", "control_rtt_ms_p95", "control_rtt_ms_max", "client_accept_total_ms_p50", "client_accept_total_ms_p95", "client_accept_total_ms_max", "client_accept_assignment_ms_p95", "client_accept_credential_ms_p95", "client_accept_activity_ms_p95", "client_accept_attach_ms_p95", "client_accept_basis_ms_p95", "control_renewal_latency_ms_p50", "control_renewal_latency_ms_p95", "control_renewal_latency_ms_max", "http_latency_ms", "event_loop_ms_p99", "db_oldest_wait_ms", "db_wait_ms_max", "inventory_hold_ms_p99", "rehome_target_row_hold_ms_p99", "isolated_replacement_hold_ms_p99", "drain_return_service_ms_p50", "drain_return_service_ms_p95", "sticky_service_ms_p50", "sticky_service_ms_p99"], each.key) ? "ms" : each.key == "queued_bytes" || each.key == "heap_used_bytes" || each.key == "forwarded_bytes" ? "By" : "1" labels { key = "role"