Skip to content

Commit 5dbd5ca

Browse files
committed
Merge file ownership foundation staging alignment into Project files
2 parents 34c6a46 + 903c669 commit 5dbd5ca

38 files changed

Lines changed: 36554 additions & 6005 deletions

‎apps/sim/background/cleanup-table-row-ttl.integration.ts‎

Lines changed: 4 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -65,9 +65,10 @@ async function seedRows(tableId: string, count: number, value: string | null = e
6565
async function rowCount(tableId: string): Promise<number> {
6666
const [result] =
6767
await control`SELECT count(*)::int AS count FROM user_table_rows WHERE table_id = ${tableId}`
68-
const [definition] =
69-
await control`SELECT row_count FROM user_table_definitions WHERE id = ${tableId}`
70-
expect(definition.row_count).toBe(result.count)
68+
const [definition] = await control`SELECT d.row_count + coalesce(sum(c.row_delta), 0)::int AS live
69+
FROM user_table_definitions d LEFT JOIN user_table_row_changes c ON c.table_id = d.id
70+
WHERE d.id = ${tableId} GROUP BY d.id`
71+
expect(definition.live).toBe(result.count)
7172
return result.count
7273
}
7374

‎apps/sim/lib/projects/__integration__/foundation.integration.ts‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -121,7 +121,7 @@ async function enforce() {
121121
try {
122122
const migration = await readFile(
123123
new URL(
124-
'../../../../../packages/db/migrations/0402_project_membership_enforcement.sql',
124+
'../../../../../packages/db/migrations/0403_project_membership_enforcement.sql',
125125
import.meta.url
126126
),
127127
'utf8'

‎apps/sim/lib/table/rows/ordering.ts‎

Lines changed: 2 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -623,9 +623,8 @@ export async function guardBatch(
623623

624624
/**
625625
* Deletes one page of rows for the async delete-job worker, committing each `DELETE_BATCH_SIZE`
626-
* chunk in its own short transaction. One statement per transaction bounds how long the
627-
* statement-level row_count trigger's lock on the definition row is held (a page-wide transaction
628-
* held it for the entire page, starving concurrent inserts and overrunning `statement_timeout`),
626+
* chunk in its own short transaction. One statement per transaction keeps each transaction short
627+
* (a page-wide transaction held its row locks for the entire page and overran `statement_timeout`),
629628
* and a mid-page failure loses at most one uncommitted batch — the keyset walker (or a task
630629
* retry) re-walks whatever remains. Skips legacy position compaction: under fractional ordering
631630
* it's unnecessary, and in the legacy path `position` gaps are harmless — rows still order by

‎apps/sim/lib/table/rows/row-writes.integration.ts‎

Lines changed: 165 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -26,11 +26,13 @@ import { deleteColumn, updateColumnConstraints } from '@/lib/table/columns/servi
2626
import { getMaxRowSizeBytes } from '@/lib/table/constants'
2727
import { bulkInsertImportBatch, importReplaceRows } from '@/lib/table/import-data'
2828
import type { DbTransaction } from '@/lib/table/planner'
29+
import { readCurrentRowsVersion } from '@/lib/table/row-changes'
2930
import { lockLiveTableSchema } from '@/lib/table/rows/live-schema'
3031
import { acquireRowOrderLock } from '@/lib/table/rows/ordering'
3132
import {
3233
batchInsertRows,
3334
batchUpdateRows,
35+
deleteRowsByFilter,
3436
insertRow,
3537
replaceTableRows,
3638
updateRow,
@@ -66,6 +68,12 @@ const [{ migrated }] = await control<{ migrated: boolean }[]>`SELECT EXISTS (
6668
AND p.proname = 'bump_user_table_rows_version_at_commit'
6769
) AS migrated`
6870

71+
/** Whether the row triggers log to `user_table_row_changes` (0402) rather than lock the definition. */
72+
const [{ logsRowChanges }] = await control<{ logsRowChanges: boolean }[]>`SELECT EXISTS (
73+
SELECT 1 FROM pg_proc
74+
WHERE proname = 'increment_user_table_row_count_stmt' AND prosrc LIKE '%user_table_row_changes%'
75+
) AS "logsRowChanges"`
76+
6977
async function createTable(columns: ColumnDefinition[]): Promise<TableDefinition> {
7078
const id = generateId()
7179
await db
@@ -84,8 +92,9 @@ async function seedRows(
8492
}
8593

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

91100
const textColumns = (...ids: string[]): ColumnDefinition[] =>
@@ -1389,6 +1398,160 @@ describe('table row writes against real PostgreSQL', () => {
13891398
})
13901399
})
13911400

1401+
describe.skipIf(!migrated || !logsRowChanges)('a held definition row', () => {
1402+
it('lets every row write path commit while another session holds the definition row', async () => {
1403+
const table = await createTable([
1404+
{ id: 'key', name: 'key', type: 'string', unique: true },
1405+
{ id: 'name', name: 'name', type: 'string' },
1406+
])
1407+
const versionBefore = await rowsVersion(table.id)
1408+
const rowIdByKey = async (key: string) => {
1409+
const [row] = await control<{ id: string }[]>`SELECT id FROM user_table_rows
1410+
WHERE table_id = ${table.id} AND data->>'key' = ${key}`
1411+
return row.id
1412+
}
1413+
const writes: Array<[string, () => Promise<unknown>]> = [
1414+
[
1415+
'insert',
1416+
() =>
1417+
insertRow(
1418+
{
1419+
tableId: table.id,
1420+
workspaceId,
1421+
data: { key: 'a', name: 'a' },
1422+
secretProvenance: undefined,
1423+
capabilityGovernedUserId: null,
1424+
},
1425+
table,
1426+
'held-insert'
1427+
),
1428+
],
1429+
[
1430+
'batch insert',
1431+
() =>
1432+
batchInsertRows(
1433+
{
1434+
tableId: table.id,
1435+
workspaceId,
1436+
rows: [
1437+
{ key: 'b', name: 'b' },
1438+
{ key: 'c', name: 'c' },
1439+
],
1440+
secretProvenance: undefined,
1441+
capabilityGovernedUserId: null,
1442+
},
1443+
table,
1444+
'held-batch-insert'
1445+
),
1446+
],
1447+
[
1448+
'upsert',
1449+
() =>
1450+
upsertRow(
1451+
{
1452+
tableId: table.id,
1453+
workspaceId,
1454+
data: { key: 'd', name: 'd' },
1455+
conflictTarget: 'key',
1456+
secretProvenance: undefined,
1457+
capabilityGovernedUserId: null,
1458+
},
1459+
table,
1460+
'held-upsert'
1461+
),
1462+
],
1463+
[
1464+
'upsert of an existing key',
1465+
() =>
1466+
upsertRow(
1467+
{
1468+
tableId: table.id,
1469+
workspaceId,
1470+
data: { key: 'd', name: 'd2' },
1471+
conflictTarget: 'key',
1472+
secretProvenance: undefined,
1473+
capabilityGovernedUserId: null,
1474+
},
1475+
table,
1476+
'held-upsert-existing'
1477+
),
1478+
],
1479+
[
1480+
'update by id',
1481+
async () =>
1482+
updateRow(
1483+
{
1484+
tableId: table.id,
1485+
rowId: await rowIdByKey('a'),
1486+
workspaceId,
1487+
data: { name: 'a2' },
1488+
secretProvenance: undefined,
1489+
capabilityGovernedUserId: null,
1490+
},
1491+
table,
1492+
'held-update'
1493+
),
1494+
],
1495+
[
1496+
'update by filter',
1497+
() =>
1498+
updateRowsByFilter(
1499+
table,
1500+
{
1501+
filter: { name: 'b' },
1502+
data: { name: 'b2' },
1503+
limit: 10,
1504+
secretProvenance: undefined,
1505+
capabilityGovernedUserId: null,
1506+
},
1507+
'held-update-by-filter'
1508+
),
1509+
],
1510+
[
1511+
'delete by filter',
1512+
() => deleteRowsByFilter(table, { filter: { name: 'c' } }, 'held-delete-by-filter'),
1513+
],
1514+
[
1515+
'replace',
1516+
() =>
1517+
replaceTableRows(
1518+
{
1519+
tableId: table.id,
1520+
workspaceId,
1521+
rows: [{ key: 'x', name: 'x' }],
1522+
secretProvenance: undefined,
1523+
},
1524+
table,
1525+
'held-replace'
1526+
),
1527+
],
1528+
]
1529+
1530+
const holder = await control.reserve()
1531+
const elapsedMs: Record<string, number> = {}
1532+
try {
1533+
await holder`BEGIN`
1534+
await holder`SELECT 1 FROM user_table_definitions WHERE id = ${table.id} FOR NO KEY UPDATE`
1535+
for (const [name, write] of writes) {
1536+
const started = Date.now()
1537+
await write()
1538+
elapsedMs[name] = Date.now() - started
1539+
}
1540+
} finally {
1541+
await holder`ROLLBACK`.catch(() => {})
1542+
holder.release()
1543+
}
1544+
1545+
for (const [name] of writes) expect(elapsedMs[name], name).toBeLessThan(2_000)
1546+
const [{ count }] = await control<{ count: number }[]>`SELECT count(*)::int AS count
1547+
FROM user_table_rows WHERE table_id = ${table.id}`
1548+
expect(count).toBe(1)
1549+
expect((await getTableById(table.id))?.rowCount).toBe(count)
1550+
// One log entry per write, except replace, which logs its DELETE and its INSERT.
1551+
expect(await rowsVersion(table.id)).toBe(versionBefore + writes.length + 1)
1552+
})
1553+
})
1554+
13921555
describe.skipIf(!migrated)('rows_version', () => {
13931556
it('advances once for a transaction that edits cells across several statements', async () => {
13941557
const table = await createTable(textColumns('name'))

‎apps/sim/lib/table/service.ts‎

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -123,9 +123,9 @@ function readLocks(row: {
123123
* validates and computes against the prior writer's committed columns.
124124
*
125125
* Uses an advisory lock (not `SELECT ... FOR UPDATE` on the definition row) so
126-
* it adds no edges to the row-lock graph — the row-count trigger (migration
127-
* 0198) locks the definition row from `insertRow`/`deleteRow`, and a FOR UPDATE
128-
* here would invert that order. Mirrors `acquireRowOrderLock`. The lock and
126+
* it adds no edges to the row-lock graph — a FOR UPDATE here would also block the
127+
* foreign-key check (KEY SHARE) of every row write to the table. Mirrors
128+
* `acquireRowOrderLock`. The lock and
129129
* the read both release at COMMIT/ROLLBACK; the wait is bounded by the
130130
* `statement_timeout` set in `setTableTxTimeouts`.
131131
*/

‎apps/sim/lib/workspace-files/search/README.md‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -69,7 +69,7 @@ The design uses PostgreSQL's documented [TOAST behavior](https://www.postgresql.
6969

7070
### File owner expansion
7171

72-
Migration `0408_file_search_owner_scope.sql` adds canonical workspace/Project ownership to the existing build, revision, and chunk pipeline. Legacy workspace inserts are translated by triggers, and the existing workspace trigram index remains usable. Project rows use a separate owner-scoped trigram index on the same chunk table. Organizations and personal users are not enabled search owners.
72+
Migration `0409_file_search_owner_scope.sql` adds canonical workspace/Project ownership to the existing build, revision, and chunk pipeline. Legacy workspace inserts are translated by triggers, and the existing workspace trigram index remains usable. Project rows use a separate owner-scoped trigram index on the same chunk table. Organizations and personal users are not enabled search owners.
7373

7474
The new `file_search_dispatch_queue` replaces the workspace-keyed scheduler because older workers cannot decode a Project without a workspace. During expansion, workspace mutations enqueue both queues, while the new dispatcher reads the owner queue. Both deployments share the dispatcher lock, revision claim token, build lease, and publication fence. Drain old app dispatchers and Trigger workers/retries before enabling Project dispatch. Retire the legacy queue and its dual enqueue in a later contract deploy; they are not a second indexing pipeline. The `file-search-owners-v3` cursor performs bounded owner/file keyset backfill and hourly reconciliation without rewriting extracted text in the migration.
7575

‎apps/sim/lib/workspace-files/search/chunks.integration.ts‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -173,7 +173,7 @@ describe('chunked workspace file search on PostgreSQL', () => {
173173
'0359_workspace_file_search_chunks.sql',
174174
ginWriteMigration,
175175
'0382_workspace_file_search_dispatch_handoff.sql',
176-
'0408_file_search_owner_scope.sql',
176+
'0409_file_search_owner_scope.sql',
177177
]) {
178178
await applyMigration(migration)
179179
}

‎apps/sim/lib/workspace-files/search/dispatcher.integration.ts‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -70,7 +70,7 @@ describe('workspace file search dispatch PostgreSQL deadlines', () => {
7070
})
7171
)
7272
const migration = readFileSync(
73-
resolve(process.cwd(), '../../packages/db/migrations/0408_file_search_owner_scope.sql'),
73+
resolve(process.cwd(), '../../packages/db/migrations/0409_file_search_owner_scope.sql'),
7474
'utf8'
7575
).replaceAll('"public".', `"${schemaName}".`)
7676
const ownerIndexes = [

‎packages/db/file-creator-lifetime.integration.ts‎

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -9,9 +9,9 @@ import { afterAll, beforeAll, beforeEach, describe, expect, it } from 'vitest'
99

1010
const databaseUrl = readTestDatabaseUrl()
1111
const migrations = [
12-
'0403_file_entity_ownership.sql',
13-
'0404_file_folder_version_ownership.sql',
14-
'0405_file_creator_lifetime.sql',
12+
'0404_file_entity_ownership.sql',
13+
'0405_file_folder_version_ownership.sql',
14+
'0406_file_creator_lifetime.sql',
1515
].map((name) => readFileSync(new URL(`./migrations/${name}`, import.meta.url), 'utf8'))
1616
const checks: { name: string; status: 'passed' | 'failed'; durationMs: number; error?: string }[] =
1717
[]

‎packages/db/file-entity-ownership.integration.ts‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -10,12 +10,12 @@ import { afterAll, beforeAll, beforeEach, describe, expect, it } from 'vitest'
1010

1111
const databaseUrl = readTestDatabaseUrl()
1212
const migration = readFileSync(
13-
new URL('./migrations/0403_file_entity_ownership.sql', import.meta.url),
13+
new URL('./migrations/0404_file_entity_ownership.sql', import.meta.url),
1414
'utf8'
1515
)
1616

1717
const workspaceBindingMigration = readFileSync(
18-
new URL('./migrations/0409_file_workspace_binding.sql', import.meta.url),
18+
new URL('./migrations/0410_file_workspace_binding.sql', import.meta.url),
1919
'utf8'
2020
)
2121

0 commit comments

Comments
 (0)