mirror of
https://github.com/stablyai/orca.git
synced 2026-10-07 08:02:21 +00:00
fix(relay): prune released control-connection reservations in bounded batches (#25352)
relay_control_connection_reservations was never deleted: 13.8M rows and 9.1 GB in production, 99.999% of them in state 'released' and ~400k more a day. Nothing reads a released row (every reader filters it out by state), but each placement still locks all of its host's rows, ~160 on average. A new director sweep step deletes released rows older than a day. It walks the heap in TID ranges of 16 pages, because without an index on released_at a `LIMIT n` delete plans as a sequential scan from page 0 that gets slower as the head of the heap empties (production EXPLAIN). Each statement selects its rows FOR UPDATE SKIP LOCKED, so a row a request holds is skipped rather than waited on, and deletes them by `ctid = ANY(ARRAY(...))` so the delete is always a TID scan. A tick stops at 400 rows, 128 pages or 250 ms, which drains the backlog over about three days at five directors.
This commit is contained in:
@@ -16,6 +16,7 @@ function stubStore(overrides: Partial<AssignmentCleanupStore> = {}) {
|
||||
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'
|
||||
|
||||
@@ -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<unknown>
|
||||
reapRegionalRehomeAttempts(): Promise<unknown>
|
||||
releaseExpiredActivityLeases(): Promise<unknown>
|
||||
pruneReleasedControlReservations(): Promise<unknown>
|
||||
releaseExpiredActivity(): Promise<unknown>
|
||||
releaseExpiredRegionPreferences(): Promise<unknown>
|
||||
evacuateDeadCells(): Promise<unknown>
|
||||
@@ -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()]
|
||||
|
||||
@@ -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<void> = Promise.resolve()
|
||||
|
||||
constructor(
|
||||
@@ -7200,6 +7219,12 @@ export class RelayAssignmentStore {
|
||||
return aborted
|
||||
}
|
||||
|
||||
async pruneReleasedControlReservations(): Promise<number> {
|
||||
return await this.releasedReservationReaper.reap(this.database, [
|
||||
this.now() - RELEASED_CONTROL_RESERVATION_RETENTION_MS
|
||||
])
|
||||
}
|
||||
|
||||
async releaseExpiredActivityLeases(): Promise<number> {
|
||||
const now = this.now()
|
||||
await this.database.query(
|
||||
|
||||
@@ -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 <retention> 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<number> {
|
||||
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<number> {
|
||||
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)
|
||||
}
|
||||
@@ -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<void> {
|
||||
const client = new pg.Client({ connectionString: databaseUrl })
|
||||
await client.connect()
|
||||
try {
|
||||
await client.query(sql)
|
||||
} finally {
|
||||
await client.end()
|
||||
}
|
||||
}
|
||||
|
||||
async function openPostgres(): Promise<RelayDatabase> {
|
||||
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<void> {
|
||||
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<string[]> {
|
||||
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<RelayDatabase>][] = [
|
||||
['sqlite', openInMemoryRelayDatabase],
|
||||
...(databaseUrl ? [['postgres', openPostgres] as [string, () => Promise<RelayDatabase>]] : [])
|
||||
]
|
||||
|
||||
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<number> {
|
||||
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()
|
||||
}
|
||||
})
|
||||
})
|
||||
Reference in New Issue
Block a user