|
| 1 | +/** |
| 2 | + * Keeps row writes on the table's live schema rather than the definition their caller resolved |
| 3 | + * before the write transaction opened. |
| 4 | + * |
| 5 | + * Schema changes hold the table's schema lock exclusively (`withLockedTable`) while they check the |
| 6 | + * stored rows against the new schema: a column made unique is scanned for duplicates, one made |
| 7 | + * required for empty cells. A write validated against an older definition skips the check the new |
| 8 | + * schema adds, and if it commits after that scan, the column ends up holding the duplicate or the |
| 9 | + * empty cell. Each write transaction therefore takes the schema lock shared and reads the schema |
| 10 | + * under it, then validates against that: a schema change waits for writes already in flight, and a |
| 11 | + * write that waited sees the change. |
| 12 | + */ |
| 13 | + |
| 14 | +import { compareStrings } from '@sim/utils/string' |
| 15 | +import { sql } from 'drizzle-orm' |
| 16 | +import { canonicalJson } from '@/lib/api/cursor-binding' |
| 17 | +import { OrchestrationError } from '@/lib/core/orchestration/types' |
| 18 | +import { getColumnId } from '@/lib/table/column-keys' |
| 19 | +import type { DbTransaction } from '@/lib/table/planner' |
| 20 | +import { type TableTxTimeouts, tableTxTimeoutSettings } from '@/lib/table/tx' |
| 21 | +import type { RowData, TableDefinition, TableSchema, ValidationResult } from '@/lib/table/types' |
| 22 | +import { |
| 23 | + coerceRowToSchema, |
| 24 | + type PatchedKeys, |
| 25 | + type UncoercibleValuePolicy, |
| 26 | + validateRowSize, |
| 27 | +} from '@/lib/table/validation' |
| 28 | + |
| 29 | +/** Compares schemas by content, whatever order their columns are listed in. */ |
| 30 | +function schemaFingerprint(schema: TableSchema): string { |
| 31 | + const columns = [...schema.columns].sort((a, b) => compareStrings(getColumnId(a), getColumnId(b))) |
| 32 | + return canonicalJson({ ...schema, columns }) |
| 33 | +} |
| 34 | + |
| 35 | +/** `table` itself when `schema` matches its schema, else a copy carrying `schema`. */ |
| 36 | +export function withLiveSchema(table: TableDefinition, schema: TableSchema): TableDefinition { |
| 37 | + return schemaFingerprint(schema) === schemaFingerprint(table.schema) |
| 38 | + ? table |
| 39 | + : { ...table, schema } |
| 40 | +} |
| 41 | + |
| 42 | +/** |
| 43 | + * Takes the table's schema lock shared and reads its live schema, and returns the definition to |
| 44 | + * validate and write against: `table` itself when its schema is still current, else a copy carrying |
| 45 | + * the live schema. Call it first in the transaction, before its other locks, passing the |
| 46 | + * transaction's `timeouts` (see `setTableTxTimeouts`) in place of a separate timeouts statement. |
| 47 | + * |
| 48 | + * One statement: `user_table_schema_for_write` (script migration 0026) takes the lock and then |
| 49 | + * reads the schema. It is VOLATILE, so under READ COMMITTED its read takes a fresh snapshot and sees |
| 50 | + * a schema change that committed while the lock waited. The timeouts are applied in a subquery the call |
| 51 | + * reads from, first. The lock waits as long as the transaction's `statement_timeout` allows, not |
| 52 | + * its shorter `lock_timeout`, which the function leaves as it found it for the locks that follow. |
| 53 | + */ |
| 54 | +export async function lockLiveTableSchema( |
| 55 | + trx: DbTransaction, |
| 56 | + table: TableDefinition, |
| 57 | + timeouts?: TableTxTimeouts |
| 58 | +): Promise<TableDefinition> { |
| 59 | + const read = sql`SELECT user_table_schema_for_write(${table.id}) AS schema` |
| 60 | + const [live] = await trx.execute<{ schema: TableSchema | null }>( |
| 61 | + timeouts ? sql`${read} FROM (SELECT ${tableTxTimeoutSettings(timeouts)}) AS settings` : read |
| 62 | + ) |
| 63 | + if (!live?.schema) throw new OrchestrationError('not_found', 'Table not found') |
| 64 | + return withLiveSchema(table, live.schema) |
| 65 | +} |
| 66 | + |
| 67 | +/** |
| 68 | + * Removes, in place, the cells of columns `snapshot` defines and `live` no longer does. A column |
| 69 | + * delete reclaims its cells in the background, so a cell written after that pass would stay behind. |
| 70 | + */ |
| 71 | +export function dropDeletedColumns( |
| 72 | + rows: readonly RowData[], |
| 73 | + snapshot: TableSchema, |
| 74 | + live: TableSchema |
| 75 | +): void { |
| 76 | + const liveIds = new Set(live.columns.map(getColumnId)) |
| 77 | + const deleted = snapshot.columns.map(getColumnId).filter((id) => !liveIds.has(id)) |
| 78 | + if (deleted.length === 0) return |
| 79 | + for (const row of rows) { |
| 80 | + for (const id of deleted) delete row[id] |
| 81 | + } |
| 82 | +} |
| 83 | + |
| 84 | +/** |
| 85 | + * Rebuilds `row`, in place, for `live` once the schema has moved since `snapshot`: from `raw`, the |
| 86 | + * row as the caller wrote it before coercing it against `snapshot`, so a value that schema would |
| 87 | + * have reshaped (`"007"` read as a number) reaches the live column as it was sent. Then drops the |
| 88 | + * cells of deleted columns, coerces and validates against `live` exactly as the write first did, |
| 89 | + * and re-checks the row's size, which a coercion to a wider type can grow. |
| 90 | + */ |
| 91 | +export function refitRowToSchema( |
| 92 | + row: RowData, |
| 93 | + raw: RowData, |
| 94 | + snapshot: TableSchema, |
| 95 | + live: TableSchema, |
| 96 | + policy?: UncoercibleValuePolicy, |
| 97 | + patchedKeys?: PatchedKeys |
| 98 | +): ValidationResult { |
| 99 | + for (const key of Object.keys(row)) delete row[key] |
| 100 | + Object.assign(row, raw) |
| 101 | + dropDeletedColumns([row], snapshot, live) |
| 102 | + const result = coerceRowToSchema(row, live, policy, patchedKeys) |
| 103 | + return result.valid ? validateRowSize(row) : result |
| 104 | +} |
0 commit comments