diff --git a/cloud/apps/relay/src/assignment-cleanup-steps.test.ts b/cloud/apps/relay/src/assignment-cleanup-steps.test.ts index 0e5eac13e31..38fbb105580 100644 --- a/cloud/apps/relay/src/assignment-cleanup-steps.test.ts +++ b/cloud/apps/relay/src/assignment-cleanup-steps.test.ts @@ -16,6 +16,7 @@ function stubStore(overrides: Partial = {}) { abortExpiredRegionalRehomes: method('abortExpiredRegionalRehomes'), reapRegionalRehomeAttempts: method('reapRegionalRehomeAttempts'), releaseExpiredActivityLeases: method('releaseExpiredActivityLeases'), + pruneReleasedControlReservations: method('pruneReleasedControlReservations'), releaseExpiredActivity: method('releaseExpiredActivity'), releaseExpiredRegionPreferences: method('releaseExpiredRegionPreferences'), evacuateDeadCells: method('evacuateDeadCells'), @@ -43,6 +44,7 @@ describe('assignment cleanup steps', () => { 'abortExpiredRegionalRehomes', 'reapRegionalRehomeAttempts', 'releaseExpiredActivityLeases', + 'pruneReleasedControlReservations', 'releaseExpiredActivity', 'releaseExpiredRegionPreferences', 'evacuateDeadCells' diff --git a/cloud/apps/relay/src/assignment-cleanup-steps.ts b/cloud/apps/relay/src/assignment-cleanup-steps.ts index 31e25be3688..9edd7044313 100644 --- a/cloud/apps/relay/src/assignment-cleanup-steps.ts +++ b/cloud/apps/relay/src/assignment-cleanup-steps.ts @@ -1,9 +1,9 @@ import { runRelayBackgroundOperation } from './relay-background-operation.js' -// The eleven periodic assignment sweeps the director runs every 30s. Each step +// The twelve periodic assignment sweeps the director runs every 30s. Each step // re-derives its state from the database and is idempotent, so they carry no // intra-tick ordering dependency — which is what makes per-step isolation -// sound: one failing sweep costs one tick of itself, never the other ten. +// sound: one failing sweep costs one tick of itself, never the other eleven. // (A single poisoned rehome row once silenced the whole chained form // fleet-wide.) Sweep failures are logged, never fed into the rehome worker's // dispatch-failure budget: a sweep exception is not a dispatch failure and @@ -17,6 +17,7 @@ export type AssignmentCleanupStore = { abortExpiredRegionalRehomes(): Promise reapRegionalRehomeAttempts(): Promise releaseExpiredActivityLeases(): Promise + pruneReleasedControlReservations(): Promise releaseExpiredActivity(): Promise releaseExpiredRegionPreferences(): Promise evacuateDeadCells(): Promise @@ -34,6 +35,10 @@ function assignmentCleanupSteps( ['abort-expired-regional-rehomes', () => assignments.abortExpiredRegionalRehomes()], ['reap-regional-rehome-attempts', () => assignments.reapRegionalRehomeAttempts()], ['release-expired-activity-leases', () => assignments.releaseExpiredActivityLeases()], + [ + 'prune-released-control-reservations', + () => assignments.pruneReleasedControlReservations() + ], ['release-expired-activity', () => assignments.releaseExpiredActivity()], ['release-expired-region-preferences', () => assignments.releaseExpiredRegionPreferences()], ['evacuate-dead-cells', () => assignments.evacuateDeadCells()] diff --git a/cloud/apps/relay/src/assignment-store.ts b/cloud/apps/relay/src/assignment-store.ts index 11cf72c6892..14c3e50a0dc 100644 --- a/cloud/apps/relay/src/assignment-store.ts +++ b/cloud/apps/relay/src/assignment-store.ts @@ -1,4 +1,5 @@ import { createDrainMigrationRowLookup } from './drain-migration-row-lookup.js' +import { HeapWindowReaper } from './heap-window-reaper.js' import { selectIdleRegionalRehomes, type IdleRegionalRehomeCandidate, @@ -498,6 +499,19 @@ export class RelayHomeCellUnavailableError extends Error { // never return otherwise starves connection headroom fleet-wide and turns // every placement into relay_capacity_exhausted. const LATE_ARRIVAL_DEBT_RETENTION_MS = 10 * 60 * 1_000 +// A released reservation is read by nothing: every reader filters it out by state, and the only +// statement that still touches it is the per-host lock, which just makes that lock set longer. A day +// is margin for forensics, not for reads. +export const RELEASED_CONTROL_RESERVATION_RETENTION_MS = 24 * 60 * 60 * 1_000 +// ~27 rows per page, so a statement deletes a few hundred rows at most. With ~9 director ticks a +// minute the row cap is ~5M rows a day: the 13.7M-row backlog drains over about three days, and a +// walk that finds nothing reads ~1 MB a tick. +const RELEASED_CONTROL_RESERVATION_REAP_BUDGET = { + pagesPerStatement: 16, + maxPagesPerTick: 128, + maxRowsPerTick: 400, + budgetMs: 250 +} const CELL_FENCE_TTL_MS = 5 * 60 * 1_000 const CELL_FENCE_ATTEMPT_TTL_MS = 60 * 60 * 1_000 const CELL_DRAIN_SEND_PERMIT_MS = 30_000 @@ -541,6 +555,11 @@ export class RelayAssignmentStore { private readonly admissionSelector: RelayCellAdmissionSelector private readonly migrationCellRegistrar: RelayMigrationCellRegistrar private readonly activityQueue = new AssignmentIdentityQueue() + private readonly releasedReservationReaper = new HeapWindowReaper( + 'relay_control_connection_reservations', + `state = 'released' AND released_at <= ?`, + RELEASED_CONTROL_RESERVATION_REAP_BUDGET + ) private assignmentTail: Promise = Promise.resolve() constructor( @@ -7200,6 +7219,12 @@ export class RelayAssignmentStore { return aborted } + async pruneReleasedControlReservations(): Promise { + return await this.releasedReservationReaper.reap(this.database, [ + this.now() - RELEASED_CONTROL_RESERVATION_RETENTION_MS + ]) + } + async releaseExpiredActivityLeases(): Promise { const now = this.now() await this.database.query( diff --git a/cloud/apps/relay/src/heap-window-reaper.ts b/cloud/apps/relay/src/heap-window-reaper.ts new file mode 100644 index 00000000000..3e7a1dd74cf --- /dev/null +++ b/cloud/apps/relay/src/heap-window-reaper.ts @@ -0,0 +1,90 @@ +import { performance } from 'node:perf_hooks' +import type { RelayDatabase, SqlRow } from './database.js' + +export type HeapWindowReapBudget = { + // Pages one DELETE may visit, so its row locks and WAL stay a few hundred rows. + pagesPerStatement: number + // Pages one tick may visit while it finds nothing to delete. + maxPagesPerTick: number + // Rows one tick deletes before it stops; this is what paces a backlog over days. + maxRowsPerTick: number + budgetMs: number +} + +// Why a TID range and not `WHERE LIMIT n`: these tables have no index on their +// retention column, so the planner answers LIMIT with a sequential scan from page 0 (production +// EXPLAIN, 2026-10-04). That scan gets longer every tick as the reaped head of the heap empties. +// A TID range bounds each statement to its own pages whatever the table holds. +export class HeapWindowReaper { + private nextPage: number | undefined + + constructor( + private readonly table: string, + private readonly predicate: string, + private readonly budget: HeapWindowReapBudget, + private readonly random: () => number = Math.random, + private readonly clock: () => number = () => performance.now() + ) {} + + async reap(database: RelayDatabase, params: unknown[]): Promise { + if (database.dialect !== 'postgres') { + return changes( + await database.query( + `DELETE FROM ${this.table} WHERE rowid IN ( + SELECT rowid FROM ${this.table} WHERE ${this.predicate} + LIMIT ${this.budget.maxRowsPerTick} + )`, + params + ) + ) + } + const pages = await heapPages(database, this.table) + if (pages === 0) return 0 + // A random first page: every director runs this sweep, and walks that all start at page 0 + // after a rollout would read the same pages in lockstep. + let page = this.nextPage ?? Math.floor(this.random() * pages) + const startedAt = this.clock() + let scanned = 0 + let deleted = 0 + while ( + scanned < this.budget.maxPagesPerTick && + deleted < this.budget.maxRowsPerTick && + this.clock() - startedAt < this.budget.budgetMs + ) { + if (page >= pages) page = 0 + const end = Math.min(page + this.budget.pagesPerStatement, pages) + // SKIP LOCKED: a row some request holds is left for a later pass rather than waited on. + // `= ANY(ARRAY(...))`, not `IN (...)`: IN can plan as a hash join over a sequential scan of + // the whole table; an array of TIDs is always a TID scan. + deleted += changes( + await database.query( + `DELETE FROM ${this.table} WHERE ctid = ANY(ARRAY( + SELECT ctid FROM ${this.table} + WHERE ctid >= CAST(? AS tid) AND ctid < CAST(? AS tid) AND ${this.predicate} + FOR UPDATE SKIP LOCKED + ))`, + [`(${page},0)`, `(${end},0)`, ...params] + ) + ) + scanned += end - page + page = end + } + this.nextPage = page + return deleted + } +} + +async function heapPages(database: RelayDatabase, table: string): Promise { + const row = ( + await database.query( + `SELECT pg_relation_size(CAST(? AS regclass)) / current_setting('block_size')::bigint + AS pages`, + [table] + ) + )[0] + return Number(row?.pages ?? 0) +} + +function changes(rows: SqlRow[]): number { + return Number(rows[0]?.changes ?? 0) +} diff --git a/cloud/apps/relay/src/released-reservation-prune.test.ts b/cloud/apps/relay/src/released-reservation-prune.test.ts new file mode 100644 index 00000000000..9dc81c32391 --- /dev/null +++ b/cloud/apps/relay/src/released-reservation-prune.test.ts @@ -0,0 +1,236 @@ +import pg from 'pg' +import { afterAll, afterEach, describe, expect, it } from 'vitest' +import { + RelayAssignmentStore, + RELEASED_CONTROL_RESERVATION_RETENTION_MS +} from './assignment-store.js' +import { + openInMemoryRelayDatabase, + openRelayDatabase, + type RelayDatabase +} from './database.js' +import { HeapWindowReaper } from './heap-window-reaper.js' + +// Released reservations were never deleted: 13.7M rows and 9 GB in production, every one of them +// still locked by each placement of its host. These pin what the prune may take and how it walks. +const databaseUrl = process.env.ORCA_RELAY_TEST_POSTGRES_URL +const describePostgres = databaseUrl ? describe : describe.skip +const schema = 'relay_released_reservation_prune_test' +const NOW = 100 * 24 * 60 * 60 * 1_000 +const STALE = NOW - RELEASED_CONTROL_RESERVATION_RETENTION_MS +const PREDICATE = `state = 'released' AND released_at <= ?` + +function scopedUrl(): string { + const url = new URL(databaseUrl!) + url.searchParams.set('options', `-c search_path=${schema}`) + return url.toString() +} + +async function onAdmin(sql: string): Promise { + const client = new pg.Client({ connectionString: databaseUrl }) + await client.connect() + try { + await client.query(sql) + } finally { + await client.end() + } +} + +async function openPostgres(): Promise { + await onAdmin(`DROP SCHEMA IF EXISTS ${schema} CASCADE`) + await onAdmin(`CREATE SCHEMA ${schema}`) + return await openRelayDatabase({ databaseUrl: scopedUrl(), dataDir: '' }) +} + +async function insertReservation( + database: RelayDatabase, + id: string, + state: string, + releasedAt: number | null +): Promise { + await database.query( + `INSERT INTO relay_control_connection_reservations + (reservation_id, idempotency_key, user_id, relay_host_id, assignment_epoch, + cell_id, state, created_at, timeout_at, released_at, updated_at) + VALUES (?, ?, 'user-1', 'host000000000001', 1, 'cell-a', ?, 1, 1, ?, 1)`, + [id, id, state, releasedAt] + ) +} + +async function remainingIds(database: RelayDatabase): Promise { + const rows = await database.query( + `SELECT reservation_id FROM relay_control_connection_reservations ORDER BY reservation_id` + ) + return rows.map((row) => String(row.reservation_id)) +} + +const dialects: [string, () => Promise][] = [ + ['sqlite', openInMemoryRelayDatabase], + ...(databaseUrl ? [['postgres', openPostgres] as [string, () => Promise]] : []) +] + +describe.each(dialects)('released reservation prune (%s)', (_dialect, open) => { + let database: RelayDatabase | undefined + + afterEach(async () => { + await database?.close() + database = undefined + }) + + afterAll(async () => { + if (databaseUrl) await onAdmin(`DROP SCHEMA IF EXISTS ${schema} CASCADE`) + }) + + it('deletes only released rows older than the retention window', async () => { + database = await open() + await insertReservation(database, 'a-old-released', 'released', STALE - 1) + await insertReservation(database, 'b-edge-released', 'released', STALE) + await insertReservation(database, 'c-fresh-released', 'released', STALE + 1) + await insertReservation(database, 'd-reserved', 'reserved', null) + await insertReservation(database, 'e-debt', 'late-arrival-debt', null) + await insertReservation(database, 'f-claimed', 'claimed', null) + // A row that once was released and came back is judged by its state, not its timestamp. + await insertReservation(database, 'g-claimed-old-release', 'claimed', STALE - 1) + const store = new RelayAssignmentStore(database, () => NOW) + + expect(await store.pruneReleasedControlReservations()).toBe(2) + expect(await remainingIds(database)).toEqual([ + 'c-fresh-released', + 'd-reserved', + 'e-debt', + 'f-claimed', + 'g-claimed-old-release' + ]) + }) +}) + +describePostgres('heap window reaper against PostgreSQL', () => { + let database: RelayDatabase + // 200 rows of this shape fill about two pages; 6,000 spread the table over ~60 pages. + const ROWS = 6_000 + + afterEach(async () => { + await database?.close() + }) + + afterAll(async () => { + await onAdmin(`DROP SCHEMA IF EXISTS ${schema} CASCADE`) + }) + + async function seed(): Promise { + database = await openPostgres() + // Every third row is still live, so each page holds rows the walk must keep. + await database.query( + `INSERT INTO relay_control_connection_reservations + (reservation_id, idempotency_key, user_id, relay_host_id, assignment_epoch, + cell_id, state, created_at, timeout_at, released_at, updated_at) + SELECT 'r-' || lpad(n::text, 6, '0'), 'r-' || n, 'user-1', 'host000000000001', 1, + 'cell-a', CASE WHEN n % 3 = 0 THEN 'claimed' ELSE 'released' END, + 1, 1, CASE WHEN n % 3 = 0 THEN NULL ELSE 1 END, 1 + FROM generate_series(1, ?) AS n`, + [ROWS] + ) + const pages = ( + await database.query( + `SELECT pg_relation_size('relay_control_connection_reservations') / 8192 AS pages` + ) + )[0]! + return Number(pages.pages) + } + + it('stops each tick at its row cap and drains the backlog over later ticks', async () => { + const pages = await seed() + expect(pages).toBeGreaterThan(20) + const reaper = new HeapWindowReaper( + 'relay_control_connection_reservations', + PREDICATE, + { pagesPerStatement: 2, maxPagesPerTick: 1_000, maxRowsPerTick: 300, budgetMs: 60_000 }, + () => 0.5 + ) + + const first = await reaper.reap(database, [NOW]) + // One statement past the cap at most: two pages of this shape hold well under 300 rows. + expect(first).toBeGreaterThanOrEqual(300) + expect(first).toBeLessThan(600) + + let ticks = 1 + while ((await reaper.reap(database, [NOW])) > 0) ticks += 1 + expect(ticks).toBeGreaterThan(5) + const left = await database.query( + `SELECT state, COUNT(*) AS rows FROM relay_control_connection_reservations GROUP BY state` + ) + expect(left).toEqual([{ state: 'claimed', rows: String(ROWS / 3) }]) + }) + + it('wraps from the end of the heap back to its first page', async () => { + const pages = await seed() + // Starts on the last page, so every other page is reached only after the wrap. + const reaper = new HeapWindowReaper( + 'relay_control_connection_reservations', + PREDICATE, + { pagesPerStatement: 4, maxPagesPerTick: pages + 4, maxRowsPerTick: 1_000_000, budgetMs: 60_000 }, + () => (pages - 1) / pages + ) + + await reaper.reap(database, [NOW]) + + expect( + await database.query( + `SELECT COUNT(*) AS rows FROM relay_control_connection_reservations + WHERE state = 'released'` + ) + ).toEqual([{ rows: '0' }]) + }) + + it('skips a row a request holds instead of waiting for it', async () => { + await seed() + const holder = new pg.Client({ connectionString: scopedUrl() }) + await holder.connect() + try { + await holder.query('BEGIN') + await holder.query( + `SELECT reservation_id FROM relay_control_connection_reservations + WHERE reservation_id = 'r-000001' FOR UPDATE` + ) + const reaper = new HeapWindowReaper( + 'relay_control_connection_reservations', + PREDICATE, + { pagesPerStatement: 1_000, maxPagesPerTick: 1_000, maxRowsPerTick: 1_000_000, budgetMs: 60_000 }, + () => 0 + ) + + // The pool's lock_timeout would fail this statement if it waited on the held row. + expect(await reaper.reap(database, [NOW])).toBe((ROWS * 2) / 3 - 1) + expect( + await database.query( + `SELECT reservation_id FROM relay_control_connection_reservations + WHERE state = 'released'` + ) + ).toEqual([{ reservation_id: 'r-000001' }]) + } finally { + await holder.query('ROLLBACK') + await holder.end() + } + }) + + it('plans each statement as a TID range scan, never a sequential scan', async () => { + await seed() + await database.query('ANALYZE relay_control_connection_reservations') + const client = new pg.Client({ connectionString: scopedUrl() }) + await client.connect() + try { + const plan = await client.query( + `EXPLAIN DELETE FROM relay_control_connection_reservations WHERE ctid = ANY(ARRAY( + SELECT ctid FROM relay_control_connection_reservations + WHERE ctid >= CAST($1 AS tid) AND ctid < CAST($2 AS tid) AND ${PREDICATE.replace('?', '$3')} + FOR UPDATE SKIP LOCKED))`, + ['(0,0)', '(16,0)', NOW] + ) + const text = plan.rows.map((row) => String(row['QUERY PLAN'])).join('\n') + expect(text).toContain('Tid Range Scan') + expect(text).not.toContain('Seq Scan') + } finally { + await client.end() + } + }) +})