mirror of
https://github.com/stablyai/orca.git
synced 2026-10-02 08:02:02 +00:00
fix(relay): pre-check constraint swaps so a warm boot sends no DDL at all
The two ALTER TABLE constraint statements were the last lock-taking statements without a pre-check, so every boot still took ACCESS EXCLUSIVE on relay_region_rehome_attempts twice. A lock target now carries the catalog answer that means there is nothing left to do. ADD CONSTRAINT skips when pg_constraint already names it; DROP CONSTRAINT IF EXISTS is the inverse and skips when it does not, because nothing to drop is nothing to do. The match is by name only: the CHECK body is generated from RELAY_REGIONS, so comparing it would re-run the swap on every region change. Changing a definition under the same name is an operator migration, and the rule comment beside SCHEMA says so. A bare DROP CONSTRAINT gets no target and throws at boot, because skipping it would swallow the undefined_object the server is supposed to raise. The census invariant is now that every lock-taking statement has a pre-check, with no exceptions, and the warm-boot Postgres test asserts zero statements sent rather than two.
This commit is contained in:
@@ -76,6 +76,9 @@ export interface RelayDatabase {
|
||||
// bounds how long that build waits for its lock, not how long it holds it. Build the index out of
|
||||
// band with CREATE INDEX CONCURRENTLY first, then add it here, where the pre-check skips it forever
|
||||
// after. relay-schema-lock-targets.test.ts pins the current list, so an addition fails CI.
|
||||
// Constraint swaps are matched by NAME in pg_constraint, never by body, because the CHECK list is
|
||||
// generated from REGION_LIST. Changing a constraint's definition under the same name therefore does
|
||||
// nothing on boot: an operator drops it, and the next boot adds the current definition back.
|
||||
const SCHEMA = `
|
||||
CREATE TABLE IF NOT EXISTS relay_invites (
|
||||
user_id TEXT NOT NULL,
|
||||
|
||||
@@ -83,12 +83,6 @@ describePostgres('relay boot-time schema against PostgreSQL', () => {
|
||||
return sent.filter(takesRelationLock)
|
||||
}
|
||||
|
||||
// The two constraint swaps look up pg_constraint, which the pre-check has no query for, so they
|
||||
// are the only lock-taking statements a warm boot is still allowed to send.
|
||||
function preCheckable(): string[] {
|
||||
return lockTaking().filter((statement) => schemaLockTarget(statement) !== undefined)
|
||||
}
|
||||
|
||||
it('creates the schema cold, then issues no lock-taking statement on the next boot', async () => {
|
||||
const cold = await applyRecording()
|
||||
expect(lockTaking().length).toBeGreaterThan(0)
|
||||
@@ -96,16 +90,15 @@ describePostgres('relay boot-time schema against PostgreSQL', () => {
|
||||
|
||||
sent = []
|
||||
const warm = await applyRecording()
|
||||
expect(preCheckable()).toEqual([])
|
||||
expect(lockTaking()).toHaveLength(2)
|
||||
// Every pre-checkable statement skipped, plus the ADD CONSTRAINT the server answers 42710 to.
|
||||
// The CREATE TABLEs still run: they resolve a name and take no lock on an existing table.
|
||||
const preCheckableCount = relayPostgresSchemaStatements().filter(
|
||||
// Zero, with no exceptions: every statement that takes a relation lock has a pre-check.
|
||||
expect(lockTaking()).toEqual([])
|
||||
const preCheckedCount = relayPostgresSchemaStatements().filter(
|
||||
(statement) => schemaLockTarget(statement) !== undefined
|
||||
).length
|
||||
expect(preCheckableCount).toBe(28)
|
||||
expect(warm.skipped).toBe(preCheckableCount + 1)
|
||||
expect(warm.ran).toBe(relayPostgresSchemaStatements().length - warm.skipped)
|
||||
expect(warm.skipped).toBe(preCheckedCount)
|
||||
expect(warm.ran).toBe(relayPostgresSchemaStatements().length - preCheckedCount)
|
||||
// The CREATE TABLEs still run: they resolve a name and take no lock on an existing table.
|
||||
expect(warm.ran).toBeGreaterThan(0)
|
||||
await pool.end()
|
||||
})
|
||||
|
||||
@@ -119,7 +112,38 @@ describePostgres('relay boot-time schema against PostgreSQL', () => {
|
||||
)
|
||||
sent = []
|
||||
await applyRecording()
|
||||
expect(preCheckable()).toEqual([])
|
||||
expect(lockTaking()).toEqual([])
|
||||
await pool.end()
|
||||
})
|
||||
|
||||
it('re-adds a constraint an operator dropped, then skips it again', async () => {
|
||||
// The pre-check matches the constraint by name only, so this is the one shape where a changed
|
||||
// body needs an operator: drop it, and the next boot puts the current definition back.
|
||||
await applyRecording()
|
||||
const named = async (): Promise<number> =>
|
||||
Number(
|
||||
(
|
||||
await pool.query(
|
||||
`SELECT count(*) AS present FROM pg_catalog.pg_constraint
|
||||
WHERE conrelid = to_regclass('relay_region_rehome_attempts')
|
||||
AND conname = 'relay_region_rehome_attempts_preferred_region_valid'`
|
||||
)
|
||||
).rows[0]?.present
|
||||
)
|
||||
expect(await named()).toBe(1)
|
||||
|
||||
await pool.query(
|
||||
`ALTER TABLE relay_region_rehome_attempts
|
||||
DROP CONSTRAINT relay_region_rehome_attempts_preferred_region_valid`
|
||||
)
|
||||
sent = []
|
||||
await applyRecording()
|
||||
expect(lockTaking()).toHaveLength(1)
|
||||
expect(await named()).toBe(1)
|
||||
|
||||
sent = []
|
||||
await applyRecording()
|
||||
expect(lockTaking()).toEqual([])
|
||||
await pool.end()
|
||||
})
|
||||
|
||||
@@ -132,7 +156,7 @@ describePostgres('relay boot-time schema against PostgreSQL', () => {
|
||||
(
|
||||
await catalogObjectPresence(
|
||||
async (sql, params) => (await pool.query(sql, params)).rows,
|
||||
{ kind: 'index', table, name: 'relay_audit_events_at' }
|
||||
{ kind: 'index', table, name: 'relay_audit_events_at', skipWhen: 'present' }
|
||||
)
|
||||
).present
|
||||
|
||||
|
||||
@@ -14,121 +14,120 @@ import { relayPostgresSchemaStatements } from './database.js'
|
||||
// brand-new index reports missing on every director at once and each one runs a non-concurrent
|
||||
// build over the whole table, which is how a boot takes the site down. Build it out of band with
|
||||
// CREATE INDEX CONCURRENTLY first, then add it to SCHEMA and update this list.
|
||||
const GOLDEN_LOCK_TAKING: (SchemaLockTarget | { unchecked: string })[] = [
|
||||
{ kind: 'index', table: 'relay_invites', name: 'relay_invites_device' },
|
||||
{ kind: 'index', table: 'relay_devices', name: 'relay_devices_current_hash' },
|
||||
{ kind: 'index', table: 'relay_devices', name: 'relay_devices_grace_hash' },
|
||||
{ kind: 'index', table: 'relay_connection_bases', name: 'relay_connection_bases_active_deadline' },
|
||||
const GOLDEN_LOCK_TAKING: SchemaLockTarget[] = [
|
||||
{ kind: 'index', table: 'relay_invites', name: 'relay_invites_device', skipWhen: 'present' },
|
||||
{ kind: 'index', table: 'relay_devices', name: 'relay_devices_current_hash', skipWhen: 'present' },
|
||||
{ kind: 'index', table: 'relay_devices', name: 'relay_devices_grace_hash', skipWhen: 'present' },
|
||||
{ kind: 'index', table: 'relay_connection_bases', name: 'relay_connection_bases_active_deadline', skipWhen: 'present' },
|
||||
{
|
||||
kind: 'index',
|
||||
table: 'relay_assignment_region_preferences',
|
||||
name: 'relay_assignment_region_preferences_observed'
|
||||
name: 'relay_assignment_region_preferences_observed',
|
||||
skipWhen: 'present'
|
||||
},
|
||||
{
|
||||
kind: 'index',
|
||||
table: 'relay_region_rehome_attempts',
|
||||
name: 'relay_region_rehome_attempts_pending'
|
||||
name: 'relay_region_rehome_attempts_pending',
|
||||
skipWhen: 'present'
|
||||
},
|
||||
{
|
||||
kind: 'index',
|
||||
table: 'relay_region_rehome_attempts',
|
||||
name: 'relay_region_rehome_attempts_host_recency'
|
||||
name: 'relay_region_rehome_attempts_host_recency',
|
||||
skipWhen: 'present'
|
||||
},
|
||||
{ kind: 'index', table: 'relay_cell_runtime', name: 'relay_cell_runtime_heartbeat' },
|
||||
{ kind: 'index', table: 'relay_cell_runtime', name: 'relay_cell_runtime_heartbeat', skipWhen: 'present' },
|
||||
{
|
||||
kind: 'index',
|
||||
table: 'relay_cell_connection_runtime',
|
||||
name: 'relay_cell_connection_runtime_heartbeat'
|
||||
name: 'relay_cell_connection_runtime_heartbeat',
|
||||
skipWhen: 'present'
|
||||
},
|
||||
{
|
||||
kind: 'index',
|
||||
table: 'relay_cell_connection_snapshots',
|
||||
name: 'relay_cell_connection_snapshot_freshness'
|
||||
name: 'relay_cell_connection_snapshot_freshness',
|
||||
skipWhen: 'present'
|
||||
},
|
||||
{ kind: 'index', table: 'relay_cell_fences', name: 'relay_cell_fences_expiry' },
|
||||
{ kind: 'index', table: 'relay_cell_committed_fences', name: 'relay_cell_committed_fences_expiry' },
|
||||
{ kind: 'index', table: 'relay_cell_fences', name: 'relay_cell_fences_expiry', skipWhen: 'present' },
|
||||
{ kind: 'index', table: 'relay_cell_committed_fences', name: 'relay_cell_committed_fences_expiry', skipWhen: 'present' },
|
||||
{
|
||||
kind: 'index',
|
||||
table: 'relay_cell_legacy_fence_adoptions',
|
||||
name: 'relay_cell_legacy_fence_adoptions_expiry'
|
||||
name: 'relay_cell_legacy_fence_adoptions_expiry',
|
||||
skipWhen: 'present'
|
||||
},
|
||||
{ kind: 'index', table: 'relay_cell_fence_attempts', name: 'relay_cell_fence_attempts_expiry' },
|
||||
{ kind: 'index', table: 'relay_cell_fence_attempts', name: 'relay_cell_fence_attempts_cell' },
|
||||
{ kind: 'index', table: 'relay_cell_fence_attempts', name: 'relay_cell_fence_attempts_expiry', skipWhen: 'present' },
|
||||
{ kind: 'index', table: 'relay_cell_fence_attempts', name: 'relay_cell_fence_attempts_cell', skipWhen: 'present' },
|
||||
{
|
||||
kind: 'index',
|
||||
table: 'relay_cell_fence_apply_invocations',
|
||||
name: 'relay_cell_fence_apply_invocations_attempt'
|
||||
name: 'relay_cell_fence_apply_invocations_attempt',
|
||||
skipWhen: 'present'
|
||||
},
|
||||
{
|
||||
kind: 'index',
|
||||
table: 'relay_cell_drain_attempt_states',
|
||||
name: 'relay_cell_drain_attempt_states_cell'
|
||||
name: 'relay_cell_drain_attempt_states_cell',
|
||||
skipWhen: 'present'
|
||||
},
|
||||
{
|
||||
kind: 'index',
|
||||
table: 'relay_assignment_activity_leases',
|
||||
name: 'relay_assignment_activity_expiry'
|
||||
name: 'relay_assignment_activity_expiry',
|
||||
skipWhen: 'present'
|
||||
},
|
||||
{
|
||||
kind: 'index',
|
||||
table: 'relay_control_connection_reservations',
|
||||
name: 'relay_control_connection_reservation_headroom'
|
||||
name: 'relay_control_connection_reservation_headroom',
|
||||
skipWhen: 'present'
|
||||
},
|
||||
{
|
||||
kind: 'index',
|
||||
table: 'relay_control_connection_reservations',
|
||||
name: 'relay_control_connection_reservation_assignment'
|
||||
name: 'relay_control_connection_reservation_assignment',
|
||||
skipWhen: 'present'
|
||||
},
|
||||
{ kind: 'index', table: 'relay_assignment_migrations', name: 'relay_assignment_migrations_active' },
|
||||
{ kind: 'index', table: 'relay_assignment_migrations', name: 'relay_assignment_migrations_active', skipWhen: 'present' },
|
||||
{
|
||||
kind: 'index',
|
||||
table: 'relay_post_drain_migration_pins',
|
||||
name: 'relay_post_drain_migration_pins_attempt'
|
||||
name: 'relay_post_drain_migration_pins_attempt',
|
||||
skipWhen: 'present'
|
||||
},
|
||||
{ kind: 'index', table: 'relay_audit_events', name: 'relay_audit_events_at' },
|
||||
{ kind: 'column', table: 'relay_region_decisions', name: 'last_considered_at' },
|
||||
{ kind: 'column', table: 'relay_region_decisions', name: 'cohort_bucket' },
|
||||
// Constraint swaps look up pg_constraint, not pg_class or pg_attribute, so the pre-check has no
|
||||
// answer for them and they still take ACCESS EXCLUSIVE on every boot. Both are cheap on
|
||||
// relay_region_rehome_attempts today and both are pinned here so a third one cannot slip in.
|
||||
{ kind: 'index', table: 'relay_audit_events', name: 'relay_audit_events_at', skipWhen: 'present' },
|
||||
{ kind: 'column', table: 'relay_region_decisions', name: 'last_considered_at', skipWhen: 'present' },
|
||||
{ kind: 'column', table: 'relay_region_decisions', name: 'cohort_bucket', skipWhen: 'present' },
|
||||
// Constraint swaps are matched by name in pg_constraint, with opposite polarities: nothing to
|
||||
// drop is nothing to do, and a name already there is nothing to add.
|
||||
{
|
||||
unchecked:
|
||||
'DROP CONSTRAINT relay_region_rehome_attempts.relay_region_rehome_attempts_preferred_region_check'
|
||||
kind: 'constraint',
|
||||
table: 'relay_region_rehome_attempts',
|
||||
name: 'relay_region_rehome_attempts_preferred_region_check',
|
||||
skipWhen: 'absent'
|
||||
},
|
||||
{
|
||||
unchecked:
|
||||
'ADD CONSTRAINT relay_region_rehome_attempts.relay_region_rehome_attempts_preferred_region_valid'
|
||||
kind: 'constraint',
|
||||
table: 'relay_region_rehome_attempts',
|
||||
name: 'relay_region_rehome_attempts_preferred_region_valid',
|
||||
skipWhen: 'present'
|
||||
},
|
||||
{ kind: 'column', table: 'relay_region_rehome_control', name: 'host_cooldown_ms' },
|
||||
{ kind: 'column', table: 'relay_control_capabilities', name: 'idle_regional_rehome' },
|
||||
{ kind: 'column', table: 'relay_region_rehome_attempts', name: 'source_generation' }
|
||||
{ kind: 'column', table: 'relay_region_rehome_control', name: 'host_cooldown_ms', skipWhen: 'present' },
|
||||
{ kind: 'column', table: 'relay_control_capabilities', name: 'idle_regional_rehome', skipWhen: 'present' },
|
||||
{ kind: 'column', table: 'relay_region_rehome_attempts', name: 'source_generation', skipWhen: 'present' }
|
||||
]
|
||||
|
||||
const INDEX_OR_ADD_COLUMN = /^(?:CREATE\s+(?:UNIQUE\s+)?INDEX|ALTER\s+TABLE\s+[^\s]+\s+ADD\s+COLUMN)/i
|
||||
|
||||
const CONSTRAINT_SWAP =
|
||||
/^ALTER\s+TABLE\s+(\S+)\s+(ADD|DROP)\s+CONSTRAINT\s+(?:IF\s+EXISTS\s+)?(\S+)/i
|
||||
|
||||
// A constraint swap is pinned by the constraint it names, not by its body: the CHECK list is
|
||||
// generated from RELAY_REGIONS, and adding a region must not have to touch this golden. Anything
|
||||
// else unchecked falls back to its whole text, so a new shape fails here loudly.
|
||||
function uncheckedIdentity(statement: string): string {
|
||||
const collapsed = sqlWithoutLeadingComments(statement).replace(/\s+/g, ' ').trim()
|
||||
const swap = CONSTRAINT_SWAP.exec(collapsed)
|
||||
return swap ? `${swap[2]!.toUpperCase()} CONSTRAINT ${swap[1]}.${swap[3]}` : collapsed
|
||||
}
|
||||
|
||||
function lockTakingStatements(): string[] {
|
||||
return relayPostgresSchemaStatements().filter(takesRelationLock)
|
||||
}
|
||||
|
||||
describe('relay boot-time lock targets', () => {
|
||||
it('matches the pinned list of lock-taking statements', () => {
|
||||
expect(
|
||||
lockTakingStatements().map(
|
||||
(statement) => schemaLockTarget(statement) ?? { unchecked: uncheckedIdentity(statement) }
|
||||
)
|
||||
).toEqual(GOLDEN_LOCK_TAKING)
|
||||
expect(lockTakingStatements().map(schemaLockTarget)).toEqual(GOLDEN_LOCK_TAKING)
|
||||
})
|
||||
|
||||
it('derives a target for every CREATE INDEX and every ALTER TABLE ADD COLUMN', () => {
|
||||
@@ -155,11 +154,14 @@ describe('relay boot-time lock targets', () => {
|
||||
}
|
||||
})
|
||||
|
||||
it('pre-checks every lock-taking statement except the two pinned constraint swaps', () => {
|
||||
it('pre-checks every lock-taking statement, with no exceptions', () => {
|
||||
// The invariant the rule comment beside SCHEMA depends on: nothing that takes a relation lock
|
||||
// reaches the server on a warm boot. A statement with no target breaks it.
|
||||
const unchecked = lockTakingStatements().filter(
|
||||
(statement) => schemaLockTarget(statement) === undefined
|
||||
)
|
||||
expect(unchecked).toHaveLength(2)
|
||||
expect(unchecked).toEqual([])
|
||||
expect(lockTakingStatements()).toHaveLength(GOLDEN_LOCK_TAKING.length)
|
||||
})
|
||||
|
||||
it('derives a target through the comment block a split schema glues on', () => {
|
||||
@@ -176,7 +178,8 @@ describe('relay boot-time lock targets', () => {
|
||||
expect(commented.map(schemaLockTarget)).toContainEqual({
|
||||
kind: 'index',
|
||||
table: 'relay_connection_bases',
|
||||
name: 'relay_connection_bases_active_deadline'
|
||||
name: 'relay_connection_bases_active_deadline',
|
||||
skipWhen: 'present'
|
||||
})
|
||||
})
|
||||
|
||||
|
||||
@@ -88,6 +88,8 @@ describe('applyPostgresSchema classification', () => {
|
||||
})
|
||||
|
||||
it('treats an already-applied constraint as skipped rather than an error', async () => {
|
||||
// Still the answer for a caller with no pre-check, and for a constraint another director
|
||||
// committed between this boot's pre-check and its ALTER TABLE.
|
||||
const query = vi.fn(async () => {
|
||||
throw postgresError('42710')
|
||||
})
|
||||
@@ -223,11 +225,72 @@ describe('applyPostgresSchema catalog pre-check', () => {
|
||||
it('never probes the catalog for a statement that takes no relation lock', async () => {
|
||||
const query = vi.fn(async (_statement: string) => undefined)
|
||||
const { catalogQuery, asked } = catalogAnswers([{ indisvalid: true }])
|
||||
await applyPostgresSchema([COMMENTED_TABLE, 'ALTER TABLE t DROP CONSTRAINT IF EXISTS c'], query, {
|
||||
catalogQuery
|
||||
})
|
||||
await applyPostgresSchema([COMMENTED_TABLE], query, { catalogQuery })
|
||||
expect(asked).toEqual([])
|
||||
expect(query).toHaveBeenCalledTimes(2)
|
||||
expect(query).toHaveBeenCalledTimes(1)
|
||||
})
|
||||
|
||||
it('skips an ADD CONSTRAINT the catalog already names', async () => {
|
||||
vi.spyOn(console, 'log').mockImplementation(() => undefined)
|
||||
const query = vi.fn(async (_statement: string) => undefined)
|
||||
const { catalogQuery, asked } = catalogAnswers([{}])
|
||||
const summary = await applyPostgresSchema(
|
||||
['ALTER TABLE relay_region_rehome_attempts ADD CONSTRAINT region_valid CHECK (r IN (1))'],
|
||||
query,
|
||||
{ catalogQuery }
|
||||
)
|
||||
expect(asked).toEqual([
|
||||
[
|
||||
expect.stringContaining('pg_catalog.pg_constraint'),
|
||||
'relay_region_rehome_attempts',
|
||||
'region_valid'
|
||||
]
|
||||
])
|
||||
expect(query).not.toHaveBeenCalled()
|
||||
expect(summary).toEqual({ ran: 0, skipped: 1 })
|
||||
})
|
||||
|
||||
it('skips a DROP CONSTRAINT IF EXISTS when the constraint is already gone', async () => {
|
||||
// Inverse polarity: an absent constraint is what means there is nothing to drop. Sending it
|
||||
// anyway takes ACCESS EXCLUSIVE to discover the same thing.
|
||||
const logged: { event?: string }[] = []
|
||||
vi.spyOn(console, 'log').mockImplementation((line: string) => {
|
||||
logged.push(JSON.parse(line))
|
||||
})
|
||||
const query = vi.fn(async (_statement: string) => undefined)
|
||||
const { catalogQuery } = catalogAnswers([])
|
||||
const summary = await applyPostgresSchema(
|
||||
['ALTER TABLE relay_region_rehome_attempts DROP CONSTRAINT IF EXISTS region_check'],
|
||||
query,
|
||||
{ catalogQuery }
|
||||
)
|
||||
expect(query).not.toHaveBeenCalled()
|
||||
expect(summary).toEqual({ ran: 0, skipped: 1 })
|
||||
expect(logged).toContainEqual({
|
||||
event: 'orca_relay_postgres_schema_object_absent',
|
||||
kind: 'constraint',
|
||||
table: 'relay_region_rehome_attempts',
|
||||
name: 'region_check',
|
||||
indisvalid: undefined
|
||||
})
|
||||
})
|
||||
|
||||
it('sends a DROP CONSTRAINT IF EXISTS when the constraint is still there', async () => {
|
||||
const query = vi.fn(async (_statement: string) => undefined)
|
||||
const { catalogQuery } = catalogAnswers([{}])
|
||||
const statement = 'ALTER TABLE t DROP CONSTRAINT IF EXISTS region_check'
|
||||
const summary = await applyPostgresSchema([statement], query, { catalogQuery })
|
||||
expect(query.mock.calls.map(([sql]) => sql)).toEqual([statement])
|
||||
expect(summary).toEqual({ ran: 1, skipped: 0 })
|
||||
})
|
||||
|
||||
it('sends an ADD CONSTRAINT the catalog does not name yet', async () => {
|
||||
const query = vi.fn(async (_statement: string) => undefined)
|
||||
const { catalogQuery } = catalogAnswers([])
|
||||
const statement = 'ALTER TABLE t ADD CONSTRAINT region_valid CHECK (r IN (1))'
|
||||
const summary = await applyPostgresSchema([statement], query, { catalogQuery })
|
||||
expect(query.mock.calls.map(([sql]) => sql)).toEqual([statement])
|
||||
expect(summary).toEqual({ ran: 1, skipped: 0 })
|
||||
})
|
||||
|
||||
it('sends every statement when no catalog query is supplied', async () => {
|
||||
|
||||
@@ -75,9 +75,9 @@ function retryableSchemaError(error: unknown, sql: string): boolean {
|
||||
return RETRYABLE_SCHEMA_CODES.has(String(value.code)) || concurrentCreateCollision(value, sql)
|
||||
}
|
||||
|
||||
// Evaluated immediately before each statement, so a pre-check still sees the tables the statements
|
||||
// Evaluated immediately before each statement, so a pre-check still sees the objects the statements
|
||||
// ahead of it created in this same boot.
|
||||
async function alreadyPresent(
|
||||
async function nothingToDo(
|
||||
target: SchemaLockTarget | undefined,
|
||||
options: SchemaStartupOptions,
|
||||
eventPrefix: string
|
||||
@@ -85,10 +85,10 @@ async function alreadyPresent(
|
||||
const catalogQuery = options.catalogQuery
|
||||
if (!catalogQuery || !target) return false
|
||||
const presence = await catalogObjectPresence(catalogQuery, target)
|
||||
if (!presence.present) return false
|
||||
if (presence.present !== (target.skipWhen === 'present')) return false
|
||||
console.log(
|
||||
JSON.stringify({
|
||||
event: `${eventPrefix}_object_present`,
|
||||
event: `${eventPrefix}_object_${target.skipWhen}`,
|
||||
kind: target.kind,
|
||||
table: target.table,
|
||||
name: target.name,
|
||||
@@ -114,7 +114,7 @@ export async function applyPostgresSchema(
|
||||
// Throws when an index or column statement's target cannot be read, rather than sending it
|
||||
// unchecked into the lock queue.
|
||||
const target = requireSchemaLockTarget(statement)
|
||||
if (await alreadyPresent(target, options, eventPrefix)) {
|
||||
if (await nothingToDo(target, options, eventPrefix)) {
|
||||
summary.skipped += 1
|
||||
continue
|
||||
}
|
||||
@@ -150,7 +150,7 @@ export async function applyPostgresSchema(
|
||||
// an object that is already there.
|
||||
if (
|
||||
concurrentCreateCollision((error as { code?: unknown; constraint?: unknown }) ?? {}, sql) &&
|
||||
(await alreadyPresent(target, options, eventPrefix))
|
||||
(await nothingToDo(target, options, eventPrefix))
|
||||
) {
|
||||
summary.skipped += 1
|
||||
break
|
||||
|
||||
@@ -21,6 +21,17 @@ WHERE t.oid = to_regclass($1)`
|
||||
const COLUMN_PRESENT = `SELECT 1 FROM pg_catalog.pg_attribute
|
||||
WHERE attrelid = to_regclass($1) AND attname = $2 AND attnum > 0 AND NOT attisdropped`
|
||||
|
||||
// Name only. The CHECK body is generated from RELAY_REGIONS, so comparing it would re-run the swap
|
||||
// on every region change, and an ADD CONSTRAINT is the one statement here that scans the table.
|
||||
const CONSTRAINT_PRESENT = `SELECT 1 FROM pg_catalog.pg_constraint
|
||||
WHERE conrelid = to_regclass($1) AND conname = $2`
|
||||
|
||||
const PRESENCE_SQL = {
|
||||
index: INDEX_PRESENT,
|
||||
column: COLUMN_PRESENT,
|
||||
constraint: CONSTRAINT_PRESENT
|
||||
} as const
|
||||
|
||||
export type SchemaCatalogPresence = { present: boolean; indisvalid: unknown }
|
||||
|
||||
// Row presence is the answer, whatever the row says. An index left invalid by a cancelled
|
||||
@@ -30,7 +41,7 @@ export async function catalogObjectPresence(
|
||||
query: SchemaCatalogQuery,
|
||||
target: SchemaLockTarget
|
||||
): Promise<SchemaCatalogPresence> {
|
||||
const sql = target.kind === 'index' ? INDEX_PRESENT : COLUMN_PRESENT
|
||||
const sql = PRESENCE_SQL[target.kind]
|
||||
const rows = await query(sql, [target.table, target.name])
|
||||
const row = rows[0]
|
||||
return row ? { present: true, indisvalid: row.indisvalid } : { present: false, indisvalid: undefined }
|
||||
|
||||
@@ -54,7 +54,8 @@ describe('schemaLockTarget', () => {
|
||||
expect(schemaLockTarget(COMMENTED_INDEX)).toEqual({
|
||||
kind: 'index',
|
||||
table: 'relay_connection_bases',
|
||||
name: 'relay_connection_bases_active_deadline'
|
||||
name: 'relay_connection_bases_active_deadline',
|
||||
skipWhen: 'present'
|
||||
})
|
||||
})
|
||||
|
||||
@@ -67,7 +68,8 @@ describe('schemaLockTarget', () => {
|
||||
).toEqual({
|
||||
kind: 'index',
|
||||
table: 'relay_control_connection_reservations',
|
||||
name: 'relay_reservation_assignment'
|
||||
name: 'relay_reservation_assignment',
|
||||
skipWhen: 'present'
|
||||
})
|
||||
})
|
||||
|
||||
@@ -75,7 +77,8 @@ describe('schemaLockTarget', () => {
|
||||
expect(schemaLockTarget('CREATE UNIQUE INDEX CONCURRENTLY i ON t(c)')).toEqual({
|
||||
kind: 'index',
|
||||
table: 't',
|
||||
name: 'i'
|
||||
name: 'i',
|
||||
skipWhen: 'present'
|
||||
})
|
||||
})
|
||||
|
||||
@@ -85,7 +88,8 @@ describe('schemaLockTarget', () => {
|
||||
expect(schemaLockTarget('CREATE INDEX IF NOT EXISTS app.i ON app.t(c)')).toEqual({
|
||||
kind: 'index',
|
||||
table: 'app.t',
|
||||
name: 'i'
|
||||
name: 'i',
|
||||
skipWhen: 'present'
|
||||
})
|
||||
})
|
||||
|
||||
@@ -93,7 +97,8 @@ describe('schemaLockTarget', () => {
|
||||
expect(schemaLockTarget('CREATE INDEX IF NOT EXISTS "od""d" ON "My Table"(c)')).toEqual({
|
||||
kind: 'index',
|
||||
table: '"My Table"',
|
||||
name: 'od"d'
|
||||
name: 'od"d',
|
||||
skipWhen: 'present'
|
||||
})
|
||||
})
|
||||
|
||||
@@ -102,21 +107,53 @@ describe('schemaLockTarget', () => {
|
||||
schemaLockTarget(`ALTER TABLE relay_region_rehome_control
|
||||
ADD COLUMN IF NOT EXISTS host_cooldown_ms BIGINT NOT NULL
|
||||
DEFAULT 604800000`)
|
||||
).toEqual({ kind: 'column', table: 'relay_region_rehome_control', name: 'host_cooldown_ms' })
|
||||
).toEqual({
|
||||
kind: 'column',
|
||||
table: 'relay_region_rehome_control',
|
||||
name: 'host_cooldown_ms',
|
||||
skipWhen: 'present'
|
||||
})
|
||||
})
|
||||
|
||||
it('derives a column target without IF NOT EXISTS', () => {
|
||||
expect(schemaLockTarget('ALTER TABLE ONLY t ADD COLUMN c TEXT')).toEqual({
|
||||
kind: 'column',
|
||||
table: 't',
|
||||
name: 'c'
|
||||
name: 'c',
|
||||
skipWhen: 'present'
|
||||
})
|
||||
})
|
||||
|
||||
it('gives a constraint swap no target', () => {
|
||||
// pg_constraint is a different lookup; these statements stay unchecked and the census pins them.
|
||||
expect(schemaLockTarget('ALTER TABLE t ADD CONSTRAINT c CHECK (x > 0)')).toBeUndefined()
|
||||
expect(schemaLockTarget('ALTER TABLE t DROP CONSTRAINT IF EXISTS c')).toBeUndefined()
|
||||
it('skips an ADD CONSTRAINT once the constraint name is there', () => {
|
||||
expect(schemaLockTarget('ALTER TABLE t ADD CONSTRAINT c CHECK (x > 0)')).toEqual({
|
||||
kind: 'constraint',
|
||||
table: 't',
|
||||
name: 'c',
|
||||
skipWhen: 'present'
|
||||
})
|
||||
})
|
||||
|
||||
it('skips a DROP CONSTRAINT IF EXISTS when the constraint is already gone', () => {
|
||||
// The inverse polarity: nothing to drop is nothing to do.
|
||||
expect(schemaLockTarget('ALTER TABLE t DROP CONSTRAINT IF EXISTS c')).toEqual({
|
||||
kind: 'constraint',
|
||||
table: 't',
|
||||
name: 'c',
|
||||
skipWhen: 'absent'
|
||||
})
|
||||
})
|
||||
|
||||
it('derives a constraint target across a line break', () => {
|
||||
expect(
|
||||
schemaLockTarget(`ALTER TABLE relay_region_rehome_attempts
|
||||
ADD CONSTRAINT relay_region_rehome_attempts_preferred_region_valid
|
||||
CHECK (preferred_region IN ('us-central1'))`)
|
||||
).toEqual({
|
||||
kind: 'constraint',
|
||||
table: 'relay_region_rehome_attempts',
|
||||
name: 'relay_region_rehome_attempts_preferred_region_valid',
|
||||
skipWhen: 'present'
|
||||
})
|
||||
})
|
||||
|
||||
it('gives CREATE TABLE IF NOT EXISTS no target', () => {
|
||||
@@ -151,14 +188,13 @@ describe('requireSchemaLockTarget', () => {
|
||||
expect(requireSchemaLockTarget(COMMENTED_INDEX)).toEqual({
|
||||
kind: 'index',
|
||||
table: 'relay_connection_bases',
|
||||
name: 'relay_connection_bases_active_deadline'
|
||||
name: 'relay_connection_bases_active_deadline',
|
||||
skipWhen: 'present'
|
||||
})
|
||||
})
|
||||
|
||||
it.each([
|
||||
['CREATE TABLE IF NOT EXISTS t (id TEXT)'],
|
||||
['ALTER TABLE t ADD CONSTRAINT c CHECK (x > 0)'],
|
||||
['ALTER TABLE t DROP CONSTRAINT IF EXISTS c'],
|
||||
['ALTER TABLE t ALTER COLUMN c SET DEFAULT 0']
|
||||
])('leaves %s alone, because no target is expected of it', (statement) => {
|
||||
expect(requireSchemaLockTarget(statement)).toBeUndefined()
|
||||
|
||||
@@ -3,9 +3,12 @@
|
||||
// statement wrote it, schema qualification and quoting included, because it is fed to
|
||||
// `to_regclass`; `name` is the bare identifier the catalog stores in `relname`/`attname`.
|
||||
export type SchemaLockTarget = {
|
||||
kind: 'index' | 'column'
|
||||
kind: 'index' | 'column' | 'constraint'
|
||||
table: string
|
||||
name: string
|
||||
// The catalog answer that means this statement has nothing left to do. Creating statements skip
|
||||
// on present; `DROP CONSTRAINT IF EXISTS` is the inverse, because nothing to drop is done.
|
||||
skipWhen: 'present' | 'absent'
|
||||
}
|
||||
|
||||
// Keywords that sit in an identifier position when the optional clause before them is absent.
|
||||
@@ -36,6 +39,19 @@ const ADD_COLUMN = new RegExp(
|
||||
`ADD\\s+COLUMN\\s+(?:IF\\s+NOT\\s+EXISTS\\s+)?${QUALIFIED}`,
|
||||
'i'
|
||||
)
|
||||
const ADD_CONSTRAINT = new RegExp(
|
||||
`^ALTER\\s+TABLE\\s+(?:IF\\s+EXISTS\\s+)?(?:ONLY\\s+)?${QUALIFIED}\\s+` +
|
||||
`ADD\\s+CONSTRAINT\\s+${QUALIFIED}`,
|
||||
'i'
|
||||
)
|
||||
// `IF EXISTS` is required, not optional. A bare `DROP CONSTRAINT` on a missing constraint is an
|
||||
// error the server is supposed to raise, and skipping it would swallow that. Without a target the
|
||||
// statement throws at boot instead, which tells the author to write `IF EXISTS`.
|
||||
const DROP_CONSTRAINT = new RegExp(
|
||||
`^ALTER\\s+TABLE\\s+(?:IF\\s+EXISTS\\s+)?(?:ONLY\\s+)?${QUALIFIED}\\s+` +
|
||||
`DROP\\s+CONSTRAINT\\s+IF\\s+EXISTS\\s+${QUALIFIED}`,
|
||||
'i'
|
||||
)
|
||||
|
||||
// Every statement shape that takes a relation lock before Postgres evaluates its existence test.
|
||||
// `CREATE TABLE IF NOT EXISTS` is absent on purpose: it resolves a name against the schema and
|
||||
@@ -57,7 +73,9 @@ function bareIdentifier(written: string): string {
|
||||
// rather than falling through to the lock path.
|
||||
const MUST_PARSE = [
|
||||
/^CREATE\s+(?:UNIQUE\s+)?INDEX\b/i,
|
||||
/^ALTER\s+TABLE\b[\s\S]*\bADD\s+COLUMN\b/i
|
||||
/^ALTER\s+TABLE\b[\s\S]*\bADD\s+COLUMN\b/i,
|
||||
/^ALTER\s+TABLE\b[\s\S]*\bADD\s+CONSTRAINT\b/i,
|
||||
/^ALTER\s+TABLE\b[\s\S]*\bDROP\s+CONSTRAINT\b/i
|
||||
]
|
||||
|
||||
// Derived from the statement itself so a renamed index cannot drift away from its pre-check.
|
||||
@@ -65,11 +83,29 @@ export function schemaLockTarget(statement: string): SchemaLockTarget | undefine
|
||||
const sql = sqlWithoutLeadingComments(statement)
|
||||
const index = CREATE_INDEX.exec(sql)
|
||||
if (index?.[1] && index[2]) {
|
||||
return { kind: 'index', table: index[2], name: bareIdentifier(index[1]) }
|
||||
return { kind: 'index', table: index[2], name: bareIdentifier(index[1]), skipWhen: 'present' }
|
||||
}
|
||||
const column = ADD_COLUMN.exec(sql)
|
||||
if (column?.[1] && column[2]) {
|
||||
return { kind: 'column', table: column[1], name: bareIdentifier(column[2]) }
|
||||
return { kind: 'column', table: column[1], name: bareIdentifier(column[2]), skipWhen: 'present' }
|
||||
}
|
||||
const added = ADD_CONSTRAINT.exec(sql)
|
||||
if (added?.[1] && added[2]) {
|
||||
return {
|
||||
kind: 'constraint',
|
||||
table: added[1],
|
||||
name: bareIdentifier(added[2]),
|
||||
skipWhen: 'present'
|
||||
}
|
||||
}
|
||||
const dropped = DROP_CONSTRAINT.exec(sql)
|
||||
if (dropped?.[1] && dropped[2]) {
|
||||
return {
|
||||
kind: 'constraint',
|
||||
table: dropped[1],
|
||||
name: bareIdentifier(dropped[2]),
|
||||
skipWhen: 'absent'
|
||||
}
|
||||
}
|
||||
return undefined
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user