Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 4 additions & 3 deletions apps/sim/background/cleanup-table-row-ttl.integration.ts
Original file line number Diff line number Diff line change
Expand Up @@ -65,9 +65,10 @@ async function seedRows(tableId: string, count: number, value: string | null = e
async function rowCount(tableId: string): Promise<number> {
const [result] =
await control`SELECT count(*)::int AS count FROM user_table_rows WHERE table_id = ${tableId}`
const [definition] =
await control`SELECT row_count FROM user_table_definitions WHERE id = ${tableId}`
expect(definition.row_count).toBe(result.count)
const [definition] = await control`SELECT d.row_count + coalesce(sum(c.row_delta), 0)::int AS live
FROM user_table_definitions d LEFT JOIN user_table_row_changes c ON c.table_id = d.id
WHERE d.id = ${tableId} GROUP BY d.id`
expect(definition.live).toBe(result.count)
return result.count
}

Expand Down
5 changes: 2 additions & 3 deletions apps/sim/lib/table/rows/ordering.ts
Original file line number Diff line number Diff line change
Expand Up @@ -623,9 +623,8 @@ export async function guardBatch(

/**
* Deletes one page of rows for the async delete-job worker, committing each `DELETE_BATCH_SIZE`
* chunk in its own short transaction. One statement per transaction bounds how long the
* statement-level row_count trigger's lock on the definition row is held (a page-wide transaction
* held it for the entire page, starving concurrent inserts and overrunning `statement_timeout`),
* chunk in its own short transaction. One statement per transaction keeps each transaction short
* (a page-wide transaction held its row locks for the entire page and overran `statement_timeout`),
* and a mid-page failure loses at most one uncommitted batch — the keyset walker (or a task
* retry) re-walks whatever remains. Skips legacy position compaction: under fractional ordering
* it's unnecessary, and in the legacy path `position` gaps are harmless — rows still order by
Expand Down
167 changes: 165 additions & 2 deletions apps/sim/lib/table/rows/row-writes.integration.ts
Original file line number Diff line number Diff line change
Expand Up @@ -26,11 +26,13 @@ import { deleteColumn, updateColumnConstraints } from '@/lib/table/columns/servi
import { getMaxRowSizeBytes } from '@/lib/table/constants'
import { bulkInsertImportBatch, importReplaceRows } from '@/lib/table/import-data'
import type { DbTransaction } from '@/lib/table/planner'
import { readCurrentRowsVersion } from '@/lib/table/row-changes'
import { lockLiveTableSchema } from '@/lib/table/rows/live-schema'
import { acquireRowOrderLock } from '@/lib/table/rows/ordering'
import {
batchInsertRows,
batchUpdateRows,
deleteRowsByFilter,
insertRow,
replaceTableRows,
updateRow,
Expand Down Expand Up @@ -66,6 +68,12 @@ const [{ migrated }] = await control<{ migrated: boolean }[]>`SELECT EXISTS (
AND p.proname = 'bump_user_table_rows_version_at_commit'
) AS migrated`

/** Whether the row triggers log to `user_table_row_changes` (0402) rather than lock the definition. */
const [{ logsRowChanges }] = await control<{ logsRowChanges: boolean }[]>`SELECT EXISTS (
SELECT 1 FROM pg_proc
WHERE proname = 'increment_user_table_row_count_stmt' AND prosrc LIKE '%user_table_row_changes%'
) AS "logsRowChanges"`

async function createTable(columns: ColumnDefinition[]): Promise<TableDefinition> {
const id = generateId()
await db
Expand All @@ -84,8 +92,9 @@ async function seedRows(
}

async function rowsVersion(tableId: string): Promise<number> {
const [row] = await control`SELECT rows_version FROM user_table_definitions WHERE id = ${tableId}`
return Number(row.rows_version)
const version = await readCurrentRowsVersion(tableId)
if (version === null) throw new Error('Fixture table missing')
return version
}

const textColumns = (...ids: string[]): ColumnDefinition[] =>
Expand Down Expand Up @@ -1382,6 +1391,160 @@ describe('table row writes against real PostgreSQL', () => {
})
})

describe.skipIf(!migrated || !logsRowChanges)('a held definition row', () => {
it('lets every row write path commit while another session holds the definition row', async () => {
const table = await createTable([
{ id: 'key', name: 'key', type: 'string', unique: true },
{ id: 'name', name: 'name', type: 'string' },
])
const versionBefore = await rowsVersion(table.id)
const rowIdByKey = async (key: string) => {
const [row] = await control<{ id: string }[]>`SELECT id FROM user_table_rows
WHERE table_id = ${table.id} AND data->>'key' = ${key}`
return row.id
}
const writes: Array<[string, () => Promise<unknown>]> = [
[
'insert',
() =>
insertRow(
{
tableId: table.id,
workspaceId,
data: { key: 'a', name: 'a' },
secretProvenance: undefined,
capabilityGovernedUserId: null,
},
table,
'held-insert'
),
],
[
'batch insert',
() =>
batchInsertRows(
{
tableId: table.id,
workspaceId,
rows: [
{ key: 'b', name: 'b' },
{ key: 'c', name: 'c' },
],
secretProvenance: undefined,
capabilityGovernedUserId: null,
},
table,
'held-batch-insert'
),
],
[
'upsert',
() =>
upsertRow(
{
tableId: table.id,
workspaceId,
data: { key: 'd', name: 'd' },
Comment thread
TheodoreSpeaks marked this conversation as resolved.
conflictTarget: 'key',
Comment thread
TheodoreSpeaks marked this conversation as resolved.
secretProvenance: undefined,
capabilityGovernedUserId: null,
},
table,
'held-upsert'
),
],
[
'upsert of an existing key',
() =>
upsertRow(
{
tableId: table.id,
workspaceId,
data: { key: 'd', name: 'd2' },
conflictTarget: 'key',
secretProvenance: undefined,
capabilityGovernedUserId: null,
},
table,
'held-upsert-existing'
),
],
[
'update by id',
async () =>
updateRow(
{
tableId: table.id,
rowId: await rowIdByKey('a'),
workspaceId,
data: { name: 'a2' },
secretProvenance: undefined,
capabilityGovernedUserId: null,
},
table,
'held-update'
),
],
[
'update by filter',
() =>
updateRowsByFilter(
table,
{
filter: { name: 'b' },
data: { name: 'b2' },
limit: 10,
secretProvenance: undefined,
capabilityGovernedUserId: null,
},
'held-update-by-filter'
),
],
[
'delete by filter',
() => deleteRowsByFilter(table, { filter: { name: 'c' } }, 'held-delete-by-filter'),
],
[
'replace',
() =>
replaceTableRows(
{
tableId: table.id,
workspaceId,
rows: [{ key: 'x', name: 'x' }],
secretProvenance: undefined,
},
table,
'held-replace'
),
],
]

const holder = await control.reserve()
const elapsedMs: Record<string, number> = {}
try {
await holder`BEGIN`
await holder`SELECT 1 FROM user_table_definitions WHERE id = ${table.id} FOR NO KEY UPDATE`
for (const [name, write] of writes) {
const started = Date.now()
await write()
elapsedMs[name] = Date.now() - started
}
} finally {
await holder`ROLLBACK`.catch(() => {})
holder.release()
}

for (const [name] of writes) expect(elapsedMs[name], name).toBeLessThan(2_000)
const [{ count }] = await control<{ count: number }[]>`SELECT count(*)::int AS count
FROM user_table_rows WHERE table_id = ${table.id}`
expect(count).toBe(1)
expect((await getTableById(table.id))?.rowCount).toBe(count)
// One log entry per write, except replace, which logs its DELETE and its INSERT.
expect(await rowsVersion(table.id)).toBe(versionBefore + writes.length + 1)
})
})

describe.skipIf(!migrated)('rows_version', () => {
it('advances once for a transaction that edits cells across several statements', async () => {
const table = await createTable(textColumns('name'))
Expand Down
6 changes: 3 additions & 3 deletions apps/sim/lib/table/service.ts
Original file line number Diff line number Diff line change
Expand Up @@ -123,9 +123,9 @@ function readLocks(row: {
* validates and computes against the prior writer's committed columns.
*
* Uses an advisory lock (not `SELECT ... FOR UPDATE` on the definition row) so
* it adds no edges to the row-lock graph — the row-count trigger (migration
* 0198) locks the definition row from `insertRow`/`deleteRow`, and a FOR UPDATE
* here would invert that order. Mirrors `acquireRowOrderLock`. The lock and
* it adds no edges to the row-lock graph — a FOR UPDATE here would also block the
* foreign-key check (KEY SHARE) of every row write to the table. Mirrors
* `acquireRowOrderLock`. The lock and
* the read both release at COMMIT/ROLLBACK; the wait is bounded by the
* `statement_timeout` set in `setTableTxTimeouts`.
*/
Expand Down
72 changes: 72 additions & 0 deletions packages/db/migrations/0402_table_row_change_triggers.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1,72 @@
-- Row writes append to user_table_row_changes instead of updating the table's definition row.
--
-- Every row-count and rows_version trigger updated the one user_table_definitions row per table
-- and held that lock until the commit was acknowledged, so a commit stalled on synchronous
-- replication made every other writer to the table queue behind it and hit lock_timeout. An
-- INSERT into the log takes no lock another writer waits on (its foreign-key check is KEY SHARE,
-- which only a table delete conflicts with). The fold cron moves the log into the definition row,
-- and readers already add the unfolded tail (0396), so the deployed app reads the same values.
--
-- One log row per table per INSERT or DELETE statement carries +n / -n rows and one version bump,
-- which replaces both the row-count triggers (0224) and the version insert/delete triggers (0240).
-- The deferred UPDATE trigger (0390) keeps its column filter and per-transaction dedupe and logs a
-- zero-delta row instead of bumping. Function and trigger names stay, so nothing else is renamed.
-- The DELETE and UPDATE triggers skip tables whose definition is gone: deleting a table or
-- workspace cascades into its rows after the definition row is already deleted, and nothing is
-- left to count.
--
-- Earlier migrations in the same deploy may COMMIT the runner's batch transaction. Re-open one so
-- the swap is atomic: between a separately committed function swap and trigger drop, a write
-- would bump the version twice or log no row at all.
BEGIN;--> statement-breakpoint

CREATE OR REPLACE FUNCTION increment_user_table_row_count_stmt()
RETURNS TRIGGER AS $$
BEGIN
INSERT INTO user_table_row_changes (table_id, row_delta)
SELECT table_id, count(*)::integer FROM new_rows GROUP BY table_id;

RETURN NULL;
END;
$$ LANGUAGE plpgsql;
--> statement-breakpoint

CREATE OR REPLACE FUNCTION decrement_user_table_row_count_stmt()
RETURNS TRIGGER AS $$
BEGIN
INSERT INTO user_table_row_changes (table_id, row_delta)
SELECT o.table_id, -count(*)::integer FROM old_rows o
WHERE EXISTS (SELECT 1 FROM user_table_definitions d WHERE d.id = o.table_id)
GROUP BY o.table_id;

RETURN NULL;
END;
$$ LANGUAGE plpgsql;
--> statement-breakpoint

CREATE OR REPLACE FUNCTION bump_user_table_rows_version_at_commit()
RETURNS TRIGGER AS $$
DECLARE
bumped text := coalesce(current_setting('sim_rows_version.bumped', true), '');
BEGIN
IF position(',' || NEW.table_id || ',' IN ',' || bumped || ',') = 0 THEN
PERFORM set_config(
'sim_rows_version.bumped',
CASE WHEN bumped = '' THEN NEW.table_id ELSE bumped || ',' || NEW.table_id END,
true
);
INSERT INTO user_table_row_changes (table_id, row_delta)
SELECT NEW.table_id, 0
WHERE EXISTS (SELECT 1 FROM user_table_definitions d WHERE d.id = NEW.table_id);
END IF;

RETURN NULL;
END;
$$ LANGUAGE plpgsql;
--> statement-breakpoint

DROP TRIGGER IF EXISTS user_table_rows_version_insert_trigger ON user_table_rows;--> statement-breakpoint
DROP TRIGGER IF EXISTS user_table_rows_version_delete_trigger ON user_table_rows;--> statement-breakpoint
DROP FUNCTION IF EXISTS bump_user_table_rows_version();--> statement-breakpoint

COMMIT;
Loading
Loading