Skip to content

Commit ba7975d

Browse files
committed
improvement(tables): fence only containment filters and trim rows version trigger cost
1 parent e4b7c74 commit ba7975d

4 files changed

Lines changed: 163 additions & 33 deletions

File tree

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

Lines changed: 76 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -12,6 +12,7 @@ import { tableTriggerMock } from '@sim/testing/mocks/table-trigger.mock'
1212
import { tableWorkflowColumnsMock } from '@sim/testing/mocks/table-workflow-columns.mock'
1313
import { sleep } from '@sim/utils/helpers'
1414
import { generateId } from '@sim/utils/id'
15+
import { sql } from 'drizzle-orm'
1516
import postgres from 'postgres'
1617
import { afterAll, beforeAll, beforeEach, describe, expect, it, vi } from 'vitest'
1718

@@ -20,16 +21,20 @@ vi.mock('@/lib/table/trigger', () => tableTriggerMock)
2021
vi.mock('@/lib/table/workflow-columns', () => tableWorkflowColumnsMock)
2122
vi.mock('@/lib/uploads/core/storage-service', () => storageServiceMock)
2223

24+
import { USER_TABLE_ROWS_SQL_NAME } from '@/lib/table/constants'
25+
import { withSeqscanOff } from '@/lib/table/planner'
2326
import {
2427
batchInsertRows,
2528
deleteRowsByFilter,
29+
firstMatchingRowsQuery,
2630
insertRow,
2731
updateRow,
2832
updateRowsByFilter,
2933
} from '@/lib/table/rows/service'
3034
import { getTableById } from '@/lib/table/service'
3135
import { getOrCreateTableSnapshot } from '@/lib/table/snapshot-cache'
32-
import type { ColumnDefinition, TableDefinition } from '@/lib/table/types'
36+
import { buildFilterClause } from '@/lib/table/sql'
37+
import type { ColumnDefinition, Filter, TableDefinition } from '@/lib/table/types'
3338

3439
const url = readTestDatabaseUrl()
3540
if (process.env.DATABASE_URL !== url) {
@@ -67,6 +72,16 @@ async function rowsVersion(tableId: string): Promise<number> {
6772
return Number(row.rows_version)
6873
}
6974

75+
interface QueryPlan {
76+
'Node Type': string
77+
'Index Name'?: string
78+
Plans?: QueryPlan[]
79+
}
80+
81+
function planNodes(plan: QueryPlan): QueryPlan[] {
82+
return [plan, ...(plan.Plans ?? []).flatMap(planNodes)]
83+
}
84+
7085
const textColumns = (...ids: string[]): ColumnDefinition[] =>
7186
ids.map((id) => ({ id, name: id, type: 'string' }))
7287

@@ -342,6 +357,66 @@ describe.skipIf(!migrated)('table row writes against real PostgreSQL', () => {
342357
WHERE table_id = ${table.id} AND id IN ${control(expected)}`
343358
expect(remaining).toHaveLength(0)
344359
})
360+
361+
it('updates exactly the first matching rows for a filter the GIN index cannot serve', async () => {
362+
const table = await createTable(textColumns('kind', 'n', 'flag'))
363+
await seedInterleaved(table)
364+
const expected = await firstHits(table.id, 9)
365+
expect(expected).toHaveLength(9)
366+
367+
const result = await updateRowsByFilter(
368+
table,
369+
{
370+
filter: { kind: { $ne: 'miss' } },
371+
data: { flag: 'set' },
372+
limit: 9,
373+
secretProvenance: undefined,
374+
capabilityGovernedUserId: null,
375+
},
376+
'limited-update-ne'
377+
)
378+
379+
expect([...result.affectedRowIds].sort()).toEqual([...expected].sort())
380+
})
381+
382+
it('fences only a containment filter into a GIN probe and leaves any other filter unfenced', async () => {
383+
const table = await createTable(textColumns('kind', 'n'))
384+
await control`INSERT INTO user_table_rows (id, table_id, workspace_id, data, order_key)
385+
SELECT ${table.id} || '-' || n, ${table.id}, ${workspaceId},
386+
jsonb_build_object('kind', CASE WHEN n % 500 = 0 THEN 'hit' ELSE 'miss' END, 'n', n::text),
387+
'a' || lpad(n::text, 6, '0')
388+
FROM generate_series(1, 20000) AS n`
389+
await control`ANALYZE user_table_rows`
390+
391+
async function planFor(filter: Filter): Promise<QueryPlan[]> {
392+
const filterClause = buildFilterClause(
393+
filter,
394+
USER_TABLE_ROWS_SQL_NAME,
395+
table.schema.columns
396+
)
397+
if (!filterClause) throw new Error('Fixture filter compiled to no clause')
398+
const plans = await withSeqscanOff((trx) =>
399+
trx.execute<{ 'QUERY PLAN': { Plan: QueryPlan }[] }>(
400+
sql`EXPLAIN (FORMAT JSON) ${firstMatchingRowsQuery(table, filter, filterClause, 5)}`
401+
)
402+
)
403+
return planNodes(plans[0]['QUERY PLAN'][0].Plan)
404+
}
405+
406+
const fenced = await planFor({ kind: 'hit' })
407+
expect(fenced.some((node) => node['Node Type'] === 'CTE Scan')).toBe(true)
408+
expect(
409+
fenced.some(
410+
(node) =>
411+
node['Node Type'] === 'Bitmap Index Scan' &&
412+
node['Index Name'] === 'user_table_rows_tenant_data_gin_idx'
413+
)
414+
).toBe(true)
415+
416+
const walked = await planFor({ kind: { $ne: 'miss' } })
417+
expect(walked.some((node) => node['Node Type'] === 'CTE Scan')).toBe(false)
418+
expect(walked[0]['Node Type']).toBe('Limit')
419+
})
345420
})
346421

347422
describe('unique columns under concurrent inserts', () => {

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

Lines changed: 48 additions & 23 deletions
Original file line numberDiff line numberDiff line change
@@ -73,6 +73,7 @@ import {
7373
buildSortClause,
7474
escapeLikePattern,
7575
fieldPredicate,
76+
isContainmentOnlyFilter,
7677
} from '@/lib/table/sql'
7778
import { fireTableTrigger } from '@/lib/table/trigger'
7879
import { scaledStatementTimeoutMs, setTableTxTimeouts } from '@/lib/table/tx'
@@ -981,35 +982,59 @@ function buildRowOrderBySql(
981982
}
982983

983984
/**
984-
* The first `limit` rows matching `filterClause` in the default `(order_key, id)` order.
985+
* The first `limit` rows matching `filter` in the default `(order_key, id)` order.
985986
*
986-
* As one ordered query, the planner walks the `(table_id, order_key, id)` index and tests every
987-
* row against the filter until `limit` match — the whole table when matches are sparse. The
988-
* MATERIALIZED CTE fences the match into its own tenant-GIN probe; only the matches' sort keys
989-
* are sorted, and only the kept rows' `data` is read.
987+
* As one ordered query, the planner walks the `(table_id, order_key, id)` index and tests each row
988+
* against the filter until `limit` match: cheap when matches are dense, the whole table when they
989+
* are sparse. A containment-only filter is fenced into a MATERIALIZED CTE instead, so the match
990+
* runs as its own tenant-GIN probe and only the matches' sort keys are sorted. Any other filter
991+
* cannot use the GIN index, so fencing it would evaluate and sort every match before the limit;
992+
* it keeps the ordered walk, which stops early.
990993
*/
994+
export function firstMatchingRowsQuery(
995+
table: TableDefinition,
996+
filter: Filter,
997+
filterClause: SQL,
998+
limit: number
999+
): SQL {
1000+
const tenantMatch = sql`${userTableRows.tableId} = ${table.id}
1001+
AND ${userTableRows.workspaceId} = ${table.workspaceId}
1002+
AND ${filterClause}`
1003+
if (!isContainmentOnlyFilter(filter, table.schema.columns)) {
1004+
return sql`
1005+
SELECT ${userTableRows.id}, ${userTableRows.data}
1006+
FROM ${userTableRows}
1007+
WHERE ${tenantMatch}
1008+
ORDER BY ${userTableRows.orderKey}, ${userTableRows.id}
1009+
LIMIT ${limit}
1010+
`
1011+
}
1012+
return sql`
1013+
WITH matched AS MATERIALIZED (
1014+
SELECT ${userTableRows.id}, ${userTableRows.orderKey}
1015+
FROM ${userTableRows}
1016+
WHERE ${tenantMatch}
1017+
),
1018+
kept AS (
1019+
SELECT id, order_key FROM matched ORDER BY order_key, id LIMIT ${limit}
1020+
)
1021+
SELECT ${userTableRows.id}, ${userTableRows.data}
1022+
FROM kept
1023+
INNER JOIN ${userTableRows} ON ${userTableRows.id} = kept.id
1024+
ORDER BY kept.order_key, kept.id
1025+
`
1026+
}
1027+
9911028
async function selectFirstMatchingRows(
9921029
table: TableDefinition,
1030+
filter: Filter,
9931031
filterClause: SQL,
9941032
limit: number
9951033
): Promise<Array<{ id: string; data: RowData }>> {
9961034
const rows = await withSeqscanOff(async (trx) =>
997-
trx.execute<{ id: string; data: RowData }>(sql`
998-
WITH matched AS MATERIALIZED (
999-
SELECT ${userTableRows.id}, ${userTableRows.orderKey}
1000-
FROM ${userTableRows}
1001-
WHERE ${userTableRows.tableId} = ${table.id}
1002-
AND ${userTableRows.workspaceId} = ${table.workspaceId}
1003-
AND ${filterClause}
1004-
),
1005-
kept AS (
1006-
SELECT id, order_key FROM matched ORDER BY order_key, id LIMIT ${limit}
1007-
)
1008-
SELECT ${userTableRows.id}, ${userTableRows.data}
1009-
FROM kept
1010-
INNER JOIN ${userTableRows} ON ${userTableRows.id} = kept.id
1011-
ORDER BY kept.order_key, kept.id
1012-
`)
1035+
trx.execute<{ id: string; data: RowData }>(
1036+
firstMatchingRowsQuery(table, filter, filterClause, limit)
1037+
)
10131038
)
10141039
return Array.from(rows)
10151040
}
@@ -2345,7 +2370,7 @@ export async function updateRowsByFilter(
23452370
return { affectedCount: affectedRowIds.length, affectedRowIds }
23462371
}
23472372

2348-
const matchingRows = await selectFirstMatchingRows(table, filterClause, limit)
2373+
const matchingRows = await selectFirstMatchingRows(table, data.filter, filterClause, limit)
23492374
if (matchingRows.length === 0) {
23502375
return { affectedCount: 0, affectedRowIds: [] }
23512376
}
@@ -2720,7 +2745,7 @@ export async function deleteRowsByFilter(
27202745
if (page.length < TABLE_LIMITS.DELETE_PAGE_SIZE) break
27212746
}
27222747
} else {
2723-
const matchingRows = await selectFirstMatchingRows(table, filterClause, limit)
2748+
const matchingRows = await selectFirstMatchingRows(table, data.filter, filterClause, limit)
27242749
const rowIds = matchingRows.map((row) => row.id)
27252750
if (rowIds.length > 0) {
27262751
const deletedIds = await deleteOrderedRowsByIds({

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

Lines changed: 32 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -138,6 +138,38 @@ export function buildFilterClause(
138138
return buildFilterClauseInternal(filter, tableName, columnMap)
139139
}
140140

141+
/**
142+
* Whether every condition in `filter` compiles to JSONB containment, the only leaf the tenant
143+
* `(table_id, data jsonb_path_ops)` GIN index serves. `$or`/`$and` groups qualify when all their
144+
* members do. Mirrors the leaf choices in {@link fieldPredicate}: `eq`/`in` (and a multi-select's
145+
* `contains`) compile to containment unless the column type compares through an equality cast;
146+
* negations, ranges, patterns, emptiness checks, and system columns do not.
147+
*/
148+
export function isContainmentOnlyFilter(filter: Filter, columns: ColumnDefinition[]): boolean {
149+
return filterIsContainmentOnly(filter, buildColumnMap(columns))
150+
}
151+
152+
function filterIsContainmentOnly(filter: Filter, columnMap: ColumnMap): boolean {
153+
return Object.entries(filter).every(([field, condition]) => {
154+
if (condition === undefined) return true
155+
if (field === '$or' || field === '$and') {
156+
return (
157+
Array.isArray(condition) &&
158+
(condition as Filter[]).every((member) => filterIsContainmentOnly(member, columnMap))
159+
)
160+
}
161+
if (Array.isArray(condition) || isSystemColumn(field)) return false
162+
const column = columnMap.get(field)
163+
const definition = column && columnTypeOf(column)
164+
if (definition?.valueForEquality && definition.jsonbCast) return false
165+
const isMultiSelect = column?.type === 'select' && column.multiple === true
166+
const containmentOps = isMultiSelect ? ['$contains'] : ['$eq', '$in']
167+
return isRecordLike(condition)
168+
? Object.keys(condition).every((op) => containmentOps.includes(op))
169+
: true
170+
})
171+
}
172+
141173
function buildFilterClauseInternal(
142174
filter: Filter,
143175
tableName: string,

‎packages/db/migrations/0385_table_rows_version_at_commit.sql‎

Lines changed: 7 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -13,7 +13,8 @@
1313
-- (cells) or `order_key` (row order). `id`, `table_id`, and `workspace_id` are
1414
-- never updated, and the snapshot never reads `position`. Constraint triggers
1515
-- are row-level only, so one transaction-local setting lists the tables already
16-
-- bumped (a comma-delimited list of md5(table_id) tokens, matched exactly) and
16+
-- bumped (a comma-delimited list of table ids, matched exactly; ids are
17+
-- generated `tbl_<hex>` or UUIDs, so they never contain the delimiter) and
1718
-- dedupes the bump to one per table per transaction. It resets at transaction
1819
-- end, so each backend holds a single placeholder setting however many tables
1920
-- it writes. The row write and its bump still commit atomically, which the
@@ -26,12 +27,11 @@ CREATE OR REPLACE FUNCTION bump_user_table_rows_version_at_commit()
2627
RETURNS TRIGGER AS $$
2728
DECLARE
2829
bumped text := coalesce(current_setting('sim_rows_version.bumped', true), '');
29-
token text := md5(NEW.table_id);
3030
BEGIN
31-
IF position(',' || token || ',' IN ',' || bumped || ',') = 0 THEN
31+
IF position(',' || NEW.table_id || ',' IN ',' || bumped || ',') = 0 THEN
3232
PERFORM set_config(
3333
'sim_rows_version.bumped',
34-
CASE WHEN bumped = '' THEN token ELSE bumped || ',' || token END,
34+
CASE WHEN bumped = '' THEN NEW.table_id ELSE bumped || ',' || NEW.table_id END,
3535
true
3636
);
3737
UPDATE user_table_definitions
@@ -64,16 +64,14 @@ CREATE CONSTRAINT TRIGGER user_table_rows_version_update_trigger
6464
-- determines workspace_id. Without a functional-dependency statistic the
6565
-- planner multiplies the two selectivities and underestimates a table's row
6666
-- count by orders of magnitude, which steers `data @>` filters to the btree
67-
-- instead of the tenant GIN. A larger table_id sample keeps its most-common
68-
-- values representative across many tables.
67+
-- instead of the tenant GIN. The per-column sample target stays at its
68+
-- default: the dependency alone restores the GIN plan, and a larger target
69+
-- would make every future autoanalyze of this table read more rows.
6970

7071
CREATE STATISTICS IF NOT EXISTS user_table_rows_workspace_table_stats (dependencies)
7172
ON workspace_id, table_id FROM user_table_rows;
7273
--> statement-breakpoint
7374

74-
ALTER TABLE user_table_rows ALTER COLUMN table_id SET STATISTICS 1000;
75-
--> statement-breakpoint
76-
7775
-- Release the DDL locks before the sample so writers never wait on it. The
7876
-- statements above are idempotent, so replaying this file after a partial run
7977
-- is safe.

0 commit comments

Comments
 (0)