Skip to content

Commit 8b89622

Browse files
improvement(tables): read row counts and versions through an append-only change log
1 parent d4f26a7 commit 8b89622

16 files changed

Lines changed: 30435 additions & 34 deletions

File tree

Lines changed: 26 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,26 @@
1+
import { createLogger } from '@sim/logger'
2+
import { type NextRequest, NextResponse } from 'next/server'
3+
import { verifyCronAuth } from '@/lib/auth/internal'
4+
import { withRouteHandler } from '@/lib/core/utils/with-route-handler'
5+
import { foldPendingTableRowChanges } from '@/lib/table/row-changes'
6+
7+
const logger = createLogger('FoldTableRowChangesApi')
8+
9+
export const dynamic = 'force-dynamic'
10+
export const maxDuration = 60
11+
12+
/** Leaves the route's `maxDuration` room for a fold started at the deadline to finish. */
13+
const FOLD_SWEEP_BUDGET_MS = 45_000
14+
15+
/**
16+
* Folds the table row-change log into each table's definition row. Each fold is one short
17+
* statement, so the sweep runs in the request; overlapping sweeps skip each other's tables.
18+
*/
19+
export const GET = withRouteHandler(async (request: NextRequest) => {
20+
const authError = verifyCronAuth(request, 'Table row-change fold')
21+
if (authError) return authError
22+
23+
const result = await foldPendingTableRowChanges(FOLD_SWEEP_BUDGET_MS)
24+
logger.info('Table row-change fold sweep completed', { ...result })
25+
return NextResponse.json({ success: true, ...result })
26+
})

‎apps/sim/ee/workspace-forking/lib/copy/copy-resources.ts‎

Lines changed: 2 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -706,8 +706,8 @@ export async function copyForkResourceContainers(
706706
schema: remappedSchema,
707707
createdBy: userId,
708708
rowsVersion: 0,
709-
// Start at 0 - the post-commit content copy raises it to the rows actually
710-
// copied, so a failed/partial copy never advertises the source's count.
709+
// Start at 0 - the row-count trigger counts the rows the content copy actually
710+
// inserts, so a failed/partial copy never advertises the source's count.
711711
rowCount: 0,
712712
// Locks are workspace-local governance and never transit a fork edge —
713713
// mirrors `copy-workflows.ts` writing `locked: false`. Inheriting them
@@ -1343,10 +1343,6 @@ export async function copyForkResourceContent(params: {
13431343
}
13441344
if (rows.length < PROVENANCE_CONTENT_PAGE) break
13451345
}
1346-
await db
1347-
.update(userTableDefinitions)
1348-
.set({ rowCount: copied })
1349-
.where(eq(userTableDefinitions.id, table.childId))
13501346
await completeForkCopyResource(control, `table:${table.childId}`)
13511347
copiedResources += 1
13521348
} catch (error) {
Lines changed: 184 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,184 @@
1+
/**
2+
* The table row-change log against real PostgreSQL: reads add the unfolded tail, and the fold
3+
* moves it into the definition row without changing what a reader sees, under concurrent appends.
4+
*/
5+
import { db } from '@sim/db'
6+
import { userTableDefinitions } from '@sim/db/schema'
7+
import { readTestDatabaseUrl } from '@sim/db/testing/test-infrastructure'
8+
import { generateId } from '@sim/utils/id'
9+
import postgres from 'postgres'
10+
import { afterAll, beforeAll, describe, expect, it } from 'vitest'
11+
import {
12+
foldPendingTableRowChanges,
13+
foldTableRowChanges,
14+
readCurrentRowsVersion,
15+
} from '@/lib/table/row-changes'
16+
import { getTableById, listTables } from '@/lib/table/service'
17+
18+
const url = readTestDatabaseUrl()
19+
if (process.env.DATABASE_URL !== url) {
20+
throw new Error('This suite requires only the disposable local test database')
21+
}
22+
const control = postgres(url, { max: 8, onnotice: () => {} })
23+
const workspaceId = generateId()
24+
const userId = generateId()
25+
26+
async function createTable(): Promise<string> {
27+
const id = generateId()
28+
await db
29+
.insert(userTableDefinitions)
30+
.values({ id, workspaceId, name: id, schema: { columns: [] }, createdBy: userId })
31+
return id
32+
}
33+
34+
async function logChanges(tableId: string, ...deltas: number[]) {
35+
for (const delta of deltas) {
36+
await control`INSERT INTO user_table_row_changes (table_id, row_delta) VALUES (${tableId}, ${delta})`
37+
}
38+
}
39+
40+
async function stored(tableId: string) {
41+
const [row] = await control<{ row_count: number; rows_version: string; updated_at: Date }[]>`
42+
SELECT row_count, rows_version, updated_at FROM user_table_definitions WHERE id = ${tableId}`
43+
return {
44+
rowCount: row.row_count,
45+
rowsVersion: Number(row.rows_version),
46+
updatedAt: row.updated_at,
47+
}
48+
}
49+
50+
async function tailLength(tableId: string): Promise<number> {
51+
const [{ n }] = await control<{ n: number }[]>`
52+
SELECT count(*)::int AS n FROM user_table_row_changes WHERE table_id = ${tableId}`
53+
return n
54+
}
55+
56+
async function liveRowCount(tableId: string): Promise<number> {
57+
const table = await getTableById(tableId)
58+
if (!table) throw new Error('Fixture table missing')
59+
return table.rowCount
60+
}
61+
62+
describe('table row-change log against real PostgreSQL', () => {
63+
beforeAll(async () => {
64+
await control`INSERT INTO "user" (id, name, email, email_verified, created_at, updated_at)
65+
VALUES (${userId}, 'Row change fixture', ${`${userId}@example.test`}, true, now(), now())`
66+
await control`INSERT INTO workspace (id, name, owner_id, billed_account_user_id)
67+
VALUES (${workspaceId}, 'Row change fixtures', ${userId}, ${userId})`
68+
})
69+
70+
afterAll(async () => {
71+
await control`DELETE FROM workspace WHERE id = ${workspaceId}`
72+
await control`DELETE FROM "user" WHERE id = ${userId}`
73+
await control.end()
74+
})
75+
76+
it('reads the stored count and version plus the unfolded tail', async () => {
77+
const tableId = await createTable()
78+
await logChanges(tableId, 5, -2, 0)
79+
80+
expect(await liveRowCount(tableId)).toBe(3)
81+
expect(await readCurrentRowsVersion(tableId)).toBe(3)
82+
expect(await readCurrentRowsVersion(tableId, workspaceId)).toBe(3)
83+
expect(await readCurrentRowsVersion(tableId, generateId())).toBeNull()
84+
const listed = (await listTables(workspaceId)).find((table) => table.id === tableId)
85+
expect(listed?.rowCount).toBe(3)
86+
})
87+
88+
it('folds the tail into the definition row without changing what readers see', async () => {
89+
const tableId = await createTable()
90+
const before = await stored(tableId)
91+
await logChanges(tableId, 4, 0, -1)
92+
93+
expect(await foldTableRowChanges(tableId)).toBe(true)
94+
95+
expect(await tailLength(tableId)).toBe(0)
96+
const after = await stored(tableId)
97+
expect(after.rowCount).toBe(3)
98+
expect(after.rowsVersion).toBe(3)
99+
expect(after.updatedAt.getTime()).toBeGreaterThanOrEqual(before.updatedAt.getTime())
100+
expect(await liveRowCount(tableId)).toBe(3)
101+
expect(await readCurrentRowsVersion(tableId)).toBe(3)
102+
})
103+
104+
it('leaves updated_at alone when the tail holds only updates', async () => {
105+
const tableId = await createTable()
106+
const before = await stored(tableId)
107+
await logChanges(tableId, 0, 0)
108+
109+
await foldTableRowChanges(tableId)
110+
111+
const after = await stored(tableId)
112+
expect(after.rowsVersion).toBe(2)
113+
expect(after.updatedAt.getTime()).toBe(before.updatedAt.getTime())
114+
})
115+
116+
it('skips a table whose definition row is held instead of waiting on it', async () => {
117+
const tableId = await createTable()
118+
await logChanges(tableId, 2)
119+
120+
const held = control.begin(async (holder) => {
121+
await holder`SELECT 1 FROM user_table_definitions WHERE id = ${tableId} FOR NO KEY UPDATE`
122+
const started = Date.now()
123+
const folded = await foldTableRowChanges(tableId)
124+
return { folded, elapsedMs: Date.now() - started }
125+
})
126+
const { folded, elapsedMs } = await held
127+
128+
expect(folded).toBe(false)
129+
expect(elapsedMs).toBeLessThan(1_000)
130+
expect(await tailLength(tableId)).toBe(1)
131+
expect(await liveRowCount(tableId)).toBe(2)
132+
})
133+
134+
it('keeps the live version exact and non-decreasing while writers append during folds', async () => {
135+
const tableId = await createTable()
136+
const writers = 8
137+
const appendsPerWriter = 40
138+
let appending = true
139+
140+
const appenders = Array.from({ length: writers }, async () => {
141+
for (let i = 0; i < appendsPerWriter; i++) await logChanges(tableId, 1)
142+
})
143+
const folder = (async () => {
144+
while (appending) await foldTableRowChanges(tableId)
145+
})()
146+
const versionsSeen: number[] = []
147+
const reader = (async () => {
148+
while (appending) {
149+
const version = await readCurrentRowsVersion(tableId)
150+
if (version === null) throw new Error('Fixture table missing')
151+
versionsSeen.push(version)
152+
}
153+
})()
154+
155+
await Promise.all(appenders)
156+
appending = false
157+
await Promise.all([folder, reader])
158+
159+
const total = writers * appendsPerWriter
160+
expect(await readCurrentRowsVersion(tableId)).toBe(total)
161+
expect(await liveRowCount(tableId)).toBe(total)
162+
for (let i = 1; i < versionsSeen.length; i++) {
163+
expect(versionsSeen[i]).toBeGreaterThanOrEqual(versionsSeen[i - 1])
164+
}
165+
166+
await foldTableRowChanges(tableId)
167+
expect(await tailLength(tableId)).toBe(0)
168+
expect(await stored(tableId)).toMatchObject({ rowCount: total, rowsVersion: total })
169+
})
170+
171+
it('sweeps every table with a tail', async () => {
172+
const tableIds = [await createTable(), await createTable(), await createTable()]
173+
for (const tableId of tableIds) await logChanges(tableId, 1, 1)
174+
175+
const result = await foldPendingTableRowChanges(30_000)
176+
177+
expect(result.budgetExhausted).toBe(false)
178+
expect(result.folded).toBeGreaterThanOrEqual(tableIds.length)
179+
for (const tableId of tableIds) {
180+
expect(await tailLength(tableId)).toBe(0)
181+
expect(await stored(tableId)).toMatchObject({ rowCount: 2, rowsVersion: 2 })
182+
}
183+
})
184+
})

‎apps/sim/lib/table/row-changes.ts‎

Lines changed: 137 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,137 @@
1+
/**
2+
* A table's live `row_count` and `rows_version`: the values stored on `user_table_definitions`
3+
* plus the tail of `user_table_row_changes` not yet folded into them.
4+
*
5+
* Row writes only append to the change log, so no writer waits on another for the shared
6+
* definition row, which otherwise stays locked until a standby acknowledges the commit. The fold
7+
* moves a table's log rows into the definition row and deletes them in one transaction, so stored
8+
* value + tail is the same immediately before and after it. `updated_at` moves only at the fold, so
9+
* it trails the last insert or delete by up to one sweep.
10+
*
11+
* Three rules keep `rows_version` exact for the snapshot cache and its mount-safety check:
12+
* - Read stored value and tail in ONE statement on the primary, so both come from one snapshot.
13+
* Two statements could straddle a fold and count its rows twice or not at all.
14+
* - Only the fold deletes log rows.
15+
* - The version is a count of log rows, never `max(id)`: ids are allocated
16+
* before commit, so a row that commits late can carry a lower id than one already seen.
17+
*/
18+
19+
import { db } from '@sim/db'
20+
import { userTableDefinitions, userTableRowChanges } from '@sim/db/schema'
21+
import { createLogger } from '@sim/logger'
22+
import { and, eq, notInArray, sql } from 'drizzle-orm'
23+
import { setTableTxTimeouts } from '@/lib/table/tx'
24+
25+
const logger = createLogger('TableRowChanges')
26+
27+
/** Tables folded per sweep query; the sweep loops until the log is empty or its budget runs out. */
28+
const FOLD_SWEEP_PAGE_SIZE = 500
29+
30+
/**
31+
* Live row count, for a select over `user_table_definitions`. Spelled with explicit table names:
32+
* drizzle prints bare column names in a subquery, where `id` would bind to the log's own column.
33+
*/
34+
export const currentRowCountSql = sql<number>`(user_table_definitions.row_count + coalesce((
35+
select sum(c.row_delta) from user_table_row_changes c
36+
where c.table_id = user_table_definitions.id
37+
), 0))::integer`.mapWith(Number)
38+
39+
/** Live rows version, for a select over `user_table_definitions`. */
40+
const currentRowsVersionSql = sql<number>`(user_table_definitions.rows_version + (
41+
select count(*) from user_table_row_changes c
42+
where c.table_id = user_table_definitions.id
43+
))::bigint`.mapWith(Number)
44+
45+
/**
46+
* The table's live `rows_version`, or `null` when no such table exists (in `workspaceId`, when
47+
* given). Always reads the primary.
48+
*/
49+
export async function readCurrentRowsVersion(
50+
tableId: string,
51+
workspaceId?: string
52+
): Promise<number | null> {
53+
const [row] = await db
54+
.select({ rowsVersion: currentRowsVersionSql })
55+
.from(userTableDefinitions)
56+
.where(
57+
workspaceId
58+
? and(
59+
eq(userTableDefinitions.id, tableId),
60+
eq(userTableDefinitions.workspaceId, workspaceId)
61+
)
62+
: eq(userTableDefinitions.id, tableId)
63+
)
64+
.limit(1)
65+
return row?.rowsVersion ?? null
66+
}
67+
68+
/**
69+
* Folds one table's change log into its definition row. Skips the table, returning `false`, when
70+
* another transaction holds the definition row (a schema change, or a concurrent fold), so the fold
71+
* never waits on it; the next sweep retries.
72+
*/
73+
export async function foldTableRowChanges(tableId: string): Promise<boolean> {
74+
return db.transaction(async (tx) => {
75+
await setTableTxTimeouts(tx)
76+
const rows = await tx.execute<{ locked: boolean }>(sql`
77+
WITH target AS (
78+
SELECT id FROM user_table_definitions WHERE id = ${tableId}
79+
FOR NO KEY UPDATE SKIP LOCKED
80+
), folded AS (
81+
DELETE FROM user_table_row_changes c USING target
82+
WHERE c.table_id = target.id
83+
RETURNING c.row_delta
84+
), totals AS (
85+
SELECT count(*) AS n,
86+
coalesce(sum(row_delta), 0) AS delta,
87+
bool_or(row_delta <> 0) AS rows_added_or_removed
88+
FROM folded
89+
), applied AS (
90+
UPDATE user_table_definitions d SET
91+
rows_version = d.rows_version + totals.n,
92+
row_count = greatest(d.row_count + totals.delta, 0),
93+
updated_at = CASE WHEN totals.rows_added_or_removed
94+
THEN timezone('UTC', now()) ELSE d.updated_at END
95+
FROM target, totals
96+
WHERE d.id = target.id AND totals.n > 0
97+
)
98+
SELECT EXISTS (SELECT 1 FROM target) AS locked`)
99+
return rows[0]?.locked === true
100+
})
101+
}
102+
103+
export interface FoldSweepResult {
104+
folded: number
105+
skipped: number
106+
budgetExhausted: boolean
107+
}
108+
109+
/**
110+
* Folds each table with unfolded log rows once, until none is left or `budgetMs` elapses. Rows a
111+
* table logs after its fold wait for the next sweep.
112+
*/
113+
export async function foldPendingTableRowChanges(budgetMs: number): Promise<FoldSweepResult> {
114+
const deadline = Date.now() + budgetMs
115+
const visited = new Set<string>()
116+
let folded = 0
117+
let skipped = 0
118+
119+
for (;;) {
120+
const pending = await db
121+
.selectDistinct({ tableId: userTableRowChanges.tableId })
122+
.from(userTableRowChanges)
123+
.where(visited.size > 0 ? notInArray(userTableRowChanges.tableId, [...visited]) : undefined)
124+
.limit(FOLD_SWEEP_PAGE_SIZE)
125+
if (pending.length === 0) return { folded, skipped, budgetExhausted: false }
126+
127+
for (const { tableId } of pending) {
128+
if (Date.now() >= deadline) {
129+
logger.warn('Table row-change fold sweep ran out of budget', { folded, skipped, budgetMs })
130+
return { folded, skipped, budgetExhausted: true }
131+
}
132+
visited.add(tableId)
133+
if (await foldTableRowChanges(tableId)) folded++
134+
else skipped++
135+
}
136+
}
137+
}

‎apps/sim/lib/table/rows/secret-provenance.integration.ts‎

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -189,6 +189,10 @@ describe('table provenance in PostgreSQL', () => {
189189
database.current = drizzle(connection, { schema })
190190
await connection.unsafe(`
191191
CREATE TABLE user_table_definitions (id text PRIMARY KEY, workspace_id text NOT NULL, rows_version integer NOT NULL, schema jsonb NOT NULL);
192+
CREATE TABLE user_table_row_changes (
193+
id bigint GENERATED ALWAYS AS IDENTITY PRIMARY KEY, table_id text NOT NULL,
194+
row_delta integer NOT NULL
195+
);
192196
CREATE TABLE user_table_rows (
193197
id text PRIMARY KEY, table_id text NOT NULL, workspace_id text NOT NULL,
194198
data jsonb NOT NULL, updated_at timestamp NOT NULL, secret_provenance_version integer,
@@ -226,7 +230,7 @@ describe('table provenance in PostgreSQL', () => {
226230

227231
beforeEach(async () => {
228232
await connection.unsafe(
229-
'TRUNCATE table_row_executions, user_table_rows, user_table_row_secret_provenance, user_table_definitions'
233+
'TRUNCATE table_row_executions, user_table_rows, user_table_row_secret_provenance, user_table_row_changes, user_table_definitions'
230234
)
231235
await connection`INSERT INTO user_table_definitions VALUES ('table-1', 'workspace-1', 7, ${JSON.stringify(table.schema)}::jsonb)`
232236
})

0 commit comments

Comments
 (0)