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
53 changes: 31 additions & 22 deletions apps/sim/lib/table/jobs/service.ts
Original file line number Diff line number Diff line change
Expand Up @@ -14,7 +14,7 @@
import { db } from '@sim/db'
import { tableJobs, userTableDefinitions, userTableRows } from '@sim/db/schema'
import type { Column, SQL } from 'drizzle-orm'
import { and, asc, desc, eq, gt, inArray, ne, or, sql } from 'drizzle-orm'
import { and, asc, desc, eq, gt, inArray, isNull, ne, or, sql } from 'drizzle-orm'
import type { CsvSkippedRecord } from '@/lib/table/import'
import { pendingDeleteMask } from '@/lib/table/rows/pending-delete-mask'
import type {
Expand Down Expand Up @@ -360,27 +360,36 @@ export async function selectExportRowPage(
limit: number
): Promise<Array<{ id: string; data: RowData; orderKey: string | null }>> {
const deleteMask = await pendingDeleteMask(table)
const rows = await db
.select({ id: userTableRows.id, data: userTableRows.data, orderKey: userTableRows.orderKey })
.from(userTableRows)
.where(
and(
eq(userTableRows.tableId, table.id),
eq(userTableRows.workspaceId, table.workspaceId),
deleteMask,
// `order_key` is nullable and sorts LAST, so the page order is "keyed rows by
// key, then the unkeyed tail by id". A bare row-constructor comparison is NULL
// for unkeyed rows (dropping them) and NULL for an unkeyed anchor (dropping
// everything), which silently truncates exports. Seek per anchor kind instead.
after
? after.orderKey === null
? sql`${userTableRows.orderKey} IS NULL AND ${userTableRows.id} > ${after.id}`
: sql`(${userTableRows.orderKey} IS NULL OR (${userTableRows.orderKey}, ${userTableRows.id}) > (${after.orderKey}, ${after.id}))`
: undefined
)
)
.orderBy(asc(userTableRows.orderKey), asc(userTableRows.id))
.limit(limit)
const tableRows = and(
eq(userTableRows.tableId, table.id),
eq(userTableRows.workspaceId, table.workspaceId),
deleteMask
)
const selectPage = (seek: SQL | undefined) =>
db
.select({ id: userTableRows.id, data: userTableRows.data, orderKey: userTableRows.orderKey })
.from(userTableRows)
.where(and(tableRows, seek))
.orderBy(asc(userTableRows.orderKey), asc(userTableRows.id))
.limit(limit)
// `order_key` is nullable and sorts LAST, so the page order is "keyed rows by key, then the
// unkeyed tail by id". A bare row-constructor comparison is NULL for unkeyed rows (dropping
// them) and NULL for an unkeyed anchor (dropping everything), which silently truncates exports.
// Past a keyed anchor the next page is the keyed rows after it followed by the unkeyed tail;
// each half is its own `(table_id, order_key, id)` range, merged in order. One WHERE with
// `order_key IS NULL OR (...)` returns the same rows but scans from the table's first key.
const rows = !after
? await selectPage(undefined)
: after.orderKey === null
? await selectPage(
sql`${userTableRows.orderKey} IS NULL AND ${userTableRows.id} > ${after.id}`
)
: await selectPage(
sql`(${userTableRows.orderKey}, ${userTableRows.id}) > (${after.orderKey}, ${after.id})`
)
.unionAll(selectPage(isNull(userTableRows.orderKey)))
.orderBy(asc(userTableRows.orderKey), asc(userTableRows.id))
.limit(limit)
// drizzle types a jsonb column as `unknown`; every writer goes through the
// row-data validators, so narrowing here is a projection, not an assumption.
return rows.map((r) => ({ ...r, data: r.data as RowData }))
Expand Down
52 changes: 52 additions & 0 deletions apps/sim/lib/table/rows/row-writes.integration.ts
Original file line number Diff line number Diff line change
Expand Up @@ -24,14 +24,17 @@ vi.mock('@/lib/uploads/core/storage-service', () => storageServiceMock)
import { deleteColumn, updateColumnConstraints } from '@/lib/table/columns/service'
import { getMaxRowSizeBytes } from '@/lib/table/constants'
import { bulkInsertImportBatch, importAppendRows, importReplaceRows } from '@/lib/table/import-data'
import { selectExportRowPage } from '@/lib/table/jobs/service'
import type { DbTransaction } from '@/lib/table/planner'
import { readCurrentRowsVersion } from '@/lib/table/row-changes'
import { decodeCursor } from '@/lib/table/rows/cursor'
import { lockLiveTableSchema } from '@/lib/table/rows/live-schema'
import {
batchInsertRows,
batchUpdateRows,
deleteRowsByFilter,
insertRow,
queryRows,
replaceTableRows,
updateRow,
updateRowsByFilter,
Expand Down Expand Up @@ -1937,4 +1940,53 @@ describe('table row writes against real PostgreSQL', () => {
expect(stored.get(snapshot.key)).toContain('during')
})
})

describe('keyset paging across the unkeyed tail', () => {
/**
* `order_key` is nullable and sorts LAST, so a seek past a keyed anchor has to reach the keyed
* rows after it and then every unkeyed row, and a compound cursor has to resume inside that
* unkeyed tail. Both page walkers must return each row exactly once, in the table's order.
*/
it('pages keyed rows and then the unkeyed tail exactly once, in order', async () => {
const table = await createTable(textColumns('name'))
const id = (suffix: string) => `${table.id}-${suffix}`
await seedRows(table.id, [
{ id: id('k0'), data: { name: 'k0' }, orderKey: 'a3' },
{ id: id('u2'), data: { name: 'u2' }, orderKey: null },
{ id: id('k1'), data: { name: 'k1' }, orderKey: 'a1' },
{ id: id('k2'), data: { name: 'k2' }, orderKey: 'a2' },
{ id: id('u0'), data: { name: 'u0' }, orderKey: null },
{ id: id('k3'), data: { name: 'k3' }, orderKey: 'a2' },
{ id: id('k4'), data: { name: 'k4' }, orderKey: 'a0' },
{ id: id('u1'), data: { name: 'u1' }, orderKey: null },
])
const expected = ['k4', 'k1', 'k2', 'k3', 'k0', 'u0', 'u1', 'u2'].map(id)

for (const limit of [1, 2, 3]) {
const paged: string[] = []
let cursor: ReturnType<typeof decodeCursor> | undefined
do {
const page = await queryRows(
table,
{ ...cursor, limit, includeTotal: false, withExecutions: false },
'keyset-paging'
)
paged.push(...page.rows.map((row) => row.id))
cursor = page.nextCursor ? decodeCursor(page.nextCursor) : undefined
} while (cursor)
expect(paged).toEqual(expected)

const exported: string[] = []
let after: { orderKey: string | null; id: string } | null = null
while (true) {
const page = await selectExportRowPage(table, after, limit)
exported.push(...page.map((row) => row.id))
const last = page.at(-1)
if (!last || page.length < limit) break
after = { orderKey: last.orderKey, id: last.id }
}
expect(exported).toEqual(expected)
}
})
})
})
49 changes: 35 additions & 14 deletions apps/sim/lib/table/rows/service.ts
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,7 @@ import { userTableRows } from '@sim/db/schema'
import { createLogger } from '@sim/logger'
import { generateId } from '@sim/utils/id'
import { escapeLikePattern } from '@sim/utils/string'
import { and, asc, count, eq, inArray, type SQL, sql } from 'drizzle-orm'
import { and, asc, count, eq, inArray, isNull, type SQL, sql } from 'drizzle-orm'
import { OrchestrationError } from '@/lib/core/orchestration/types'
import {
assertRowCapacity,
Expand Down Expand Up @@ -1492,24 +1492,45 @@ export async function fetchRowsBounded(params: BoundedFetchParams): Promise<Boun
ask: number
) => {
const buildQuery = (executor: DbExecutor) => {
// `order_key` is nullable (rows predating the backfill, and forked rows that
// inherit a NULL key). A bare row-constructor comparison evaluates to NULL for
// those rows, so they are dropped by WHERE — and because NULLs sort LAST under
// `ORDER BY order_key, id`, the entire unkeyed tail becomes unreachable and the
// drain terminates early reporting `hasMore: false`. Admitting NULLs keeps the
// seek set exactly "the tail after the anchor", which is also what the compound
// `{k,i,o}` cursor's `offsetFromAnchor` accounting assumes.
const seekWhere = batchSeek
? and(
if (!batchSeek) {
const query = executor
.select()
.from(userTableRows)
.where(baseWhere)
.orderBy(orderBy)
.limit(ask)
return batchOffset > 0 ? query.offset(batchOffset) : query
}
// The seek set is "every row after the anchor in `(order_key, id)` order". `order_key` is
// nullable (rows predating the backfill, and forked rows that inherit a NULL key) and NULLs
// sort LAST, so that set is the keyed rows past the anchor followed by the whole unkeyed
// tail — the set the compound `{k,i,o}` cursor's `offsetFromAnchor` counts into. A bare
// row-constructor comparison is NULL for unkeyed rows and would drop that tail. Admitting it
// with `order_key IS NULL OR (...)` in one WHERE keeps the rows but defeats the
// `(table_id, order_key, id)` seek: the scan starts at the table's first key and filters out
// every row before the anchor. Each branch below is an index range on its own, each yields
// at most the rows the page can reach, and the planner merges the two ordered streams.
const reach = ask + batchOffset
const keyedPastAnchor = executor
.select()
.from(userTableRows)
.where(
and(
baseWhere,
sql`(${userTableRows.orderKey} IS NULL OR (${userTableRows.orderKey}, ${userTableRows.id}) > (${batchSeek.orderKey}, ${batchSeek.id}))`
sql`(${userTableRows.orderKey}, ${userTableRows.id}) > (${batchSeek.orderKey}, ${batchSeek.id})`
)
: baseWhere
const query = executor
)
.orderBy(orderBy)
.limit(reach)
const unkeyedTail = executor
.select()
.from(userTableRows)
.where(seekWhere)
.where(and(baseWhere, isNull(userTableRows.orderKey)))
.orderBy(orderBy)
.limit(reach)
const query = keyedPastAnchor
.unionAll(unkeyedTail)
.orderBy(asc(userTableRows.orderKey), asc(userTableRows.id))
.limit(ask)
return batchOffset > 0 ? query.offset(batchOffset) : query
}
Expand Down
51 changes: 13 additions & 38 deletions apps/sim/lib/table/service-filter-threading.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,7 @@
* timestamp for dates) are always available at the SQL builder layer — the
* latent bug that PR #4657 was originally fixing.
*/
import { dbChainMockFns, resetDbChainMock, setEnv } from '@sim/testing'
import { dbChainMockFns, queueTableRows, resetDbChainMock, schemaMock, setEnv } from '@sim/testing'
import { tableTriggerMock, tableTriggerMockFns } from '@sim/testing/mocks/table-trigger.mock'
import { tableWorkflowColumnsMock } from '@sim/testing/mocks/table-workflow-columns.mock'
import { sql } from 'drizzle-orm'
Expand Down Expand Up @@ -224,14 +224,18 @@ describe('queryRows byte budget', () => {
const largeRow = row(1, TABLE_LIMITS.MAX_ROW_SIZE_BYTES)
const smallRow = row(2, 0)
const state = { drainBatch: 0 }
dbChainMockFns.limit.mockResolvedValueOnce([])
dbChainMockFns.limit.mockImplementation(async (ask: number) => {
const drainBatch = async (ask: number) => {
state.drainBatch++
if (state.drainBatch > 1001) return []
const rows = Array.from({ length: ask }, () => smallRow)
if (state.drainBatch === 1) rows[0] = largeRow
return rows
})
}
dbChainMockFns.limit.mockResolvedValueOnce([])
dbChainMockFns.limit.mockImplementationOnce(drainBatch)
// Every later batch seeks past the last keyed row: its page is the outer `.limit()` of the
// keyed/unkeyed union.
dbChainMockFns.unionAll.mockImplementation(() => ({ orderBy: () => ({ limit: drainBatch }) }))
return state
}

Expand Down Expand Up @@ -300,35 +304,6 @@ describe('queryRows byte budget', () => {
Math.floor((4 * TABLE_LIMITS.MAX_QUERY_RESULT_BYTES) / TABLE_LIMITS.MAX_ROW_SIZE_BYTES)
)

/**
* The seek predicate that pages the drain. `order_key` is nullable (rows
* predating the backfill, forked rows), and a bare
* `(order_key, id) > (:k, :i)` evaluates to NULL for those rows, so WHERE drops
* them. Because NULLs sort last, the whole unkeyed tail then becomes
* unreachable and the drain reports `hasMore: false` — 52 of 121 rows returned
* with no error, reproduced against a real table.
*/
it('seeks with a NULL-admitting comparison so unkeyed rows stay reachable', async () => {
dbChainMockFns.limit.mockResolvedValueOnce([])
dbChainMockFns.limit.mockResolvedValueOnce([])

await queryRows(
TABLE,
{
after: { orderKey: 'a5', id: 'row_5' },
limit: 10,
includeTotal: false,
withExecutions: false,
},
'req-1'
)

const seekWhere = JSON.stringify(dbChainMockFns.where.mock.calls.at(-1))
expect(seekWhere).toMatch(/is null/i)
expect(seekWhere).toContain('a5')
expect(seekWhere).toContain('row_5')
})

it('resumes exactly at the last returned row when the cursor is fed back', async () => {
dbChainMockFns.limit.mockResolvedValueOnce([])
// limit 2 + witness row_3 → page ends at row_2 with a resume cursor.
Expand All @@ -347,7 +322,7 @@ describe('queryRows byte budget', () => {
vi.clearAllMocks()
resetDbChainMock()
dbChainMockFns.limit.mockResolvedValueOnce([])
dbChainMockFns.limit.mockResolvedValueOnce([row(3, 8), row(4, 8)])
queueTableRows(schemaMock.userTableRows, [row(3, 8), row(4, 8)])

const page2 = await queryRows(
TABLE,
Expand All @@ -357,7 +332,7 @@ describe('queryRows byte budget', () => {

// No gap (row_3 is the first row past the anchor) and no duplicate of row_2.
expect(page2.rows.map((r) => r.id)).toEqual(['row_3', 'row_4'])
const seekWhere = JSON.stringify(dbChainMockFns.where.mock.calls.at(-1))
const seekWhere = JSON.stringify(dbChainMockFns.where.mock.calls)
expect(seekWhere).toContain('a2')
expect(seekWhere).toContain('row_2')
})
Expand Down Expand Up @@ -392,14 +367,14 @@ describe('queryRows byte budget', () => {
const batch1 = Array.from({ length: firstAsk }, (_, i) => row(i, 8))
dbChainMockFns.limit.mockResolvedValueOnce([])
dbChainMockFns.limit.mockResolvedValueOnce(batch1)
dbChainMockFns.limit.mockResolvedValueOnce([row(firstAsk, 8)])
queueTableRows(schemaMock.userTableRows, [row(firstAsk, 8)])

const result = await queryRows(TABLE, { includeTotal: false, withExecutions: false }, 'req-1')

expect(result.rows).toHaveLength(firstAsk + 1)
expect(result.nextCursor).toBeNull()
const firstBatchAsk = (dbChainMockFns.limit.mock.calls[1] ?? [])[0] as number
const secondBatchAsk = (dbChainMockFns.limit.mock.calls[2] ?? [])[0] as number
const secondBatchAsk = (dbChainMockFns.limit.mock.calls.at(-1) ?? [])[0] as number
expect(secondBatchAsk).toBeGreaterThan(firstBatchAsk)
})

Expand Down Expand Up @@ -429,7 +404,7 @@ describe('queryRows byte budget', () => {
const unkeyedRow = (i: number, blobBytes: number) => ({ ...row(i, blobBytes), orderKey: null })
const perRow = Math.floor(TABLE_LIMITS.MAX_QUERY_RESULT_BYTES * 0.4)
dbChainMockFns.limit.mockResolvedValueOnce([])
dbChainMockFns.limit.mockResolvedValueOnce([
queueTableRows(schemaMock.userTableRows, [
unkeyedRow(6, perRow),
unkeyedRow(7, perRow),
unkeyedRow(8, perRow),
Expand Down
11 changes: 9 additions & 2 deletions packages/testing/src/mocks/database.mock.ts
Original file line number Diff line number Diff line change
Expand Up @@ -182,6 +182,9 @@ function dequeueChainRows(tables: unknown[]): unknown[] | null {
* implementation returns a sentinel that the chain replaces with the
* chain-local builder, while any `mock*` override on the spy wins verbatim.
*
* `unionAll` resolves to the LEFT chain's rows (its own queued set or the outer
* `.limit` override); the right chain is never awaited, so it consumes nothing.
*
* `for` mirrors drizzle's `.for('update')` — it returns a Promise with
* `.limit` / `.orderBy` / `.returning` / `.groupBy` attached, so both
* `await .where().for('update')` (terminal) and
Expand Down Expand Up @@ -233,6 +236,7 @@ const where = chainSpy()
const limit = chainSpy()
const offset = chainSpy()
const orderBy = chainSpy()
const unionAll = chainSpy()
const groupBy = chainSpy()
const having = chainSpy()
const asAlias = chainSpy()
Expand Down Expand Up @@ -308,12 +312,13 @@ const lazyRowsThenable = (getRows: RowsSupplier): any => ({
})

// `.limit()` returns a builder that is awaitable and also exposes `.offset()`
// for keyset/OFFSET paging (`.limit(n).offset(m)`) and `.for()` for drizzle's
// `.limit(1).for('update')` row-lock form.
// for keyset/OFFSET paging (`.limit(n).offset(m)`), `.for()` for drizzle's
// `.limit(1).for('update')` row-lock form, and `.unionAll()` for a limited branch.
const limitBuilder = (getRows: RowsSupplier, fields: SelectedFields = {}) => {
const thenable = lazyRowsThenable(getRows)
thenable.offset = spyOrDefault(offset, () => limitBuilder(getRows, fields))
thenable.for = spyOrDefault(forClause, () => limitBuilder(getRows, fields))
thenable.unionAll = spyOrDefault(unionAll, () => terminalBuilder(getRows, fields))
thenable.as = spyOrDefault(asAlias, (alias: string) => subqueryFields(fields, alias))
return thenable
}
Expand All @@ -322,6 +327,7 @@ const terminalBuilder = (getRows: RowsSupplier, fields: SelectedFields = {}): an
const thenable = lazyRowsThenable(getRows)
thenable.limit = spyOrDefault(limit, () => limitBuilder(getRows, fields))
thenable.orderBy = spyOrDefault(orderBy, () => terminalBuilder(getRows, fields))
thenable.unionAll = spyOrDefault(unionAll, () => terminalBuilder(getRows, fields))
thenable.as = spyOrDefault(asAlias, (alias: string) => subqueryFields(fields, alias))
thenable.returning = returning
thenable.groupBy = spyOrDefault(groupBy, () => {
Expand Down Expand Up @@ -371,6 +377,7 @@ export const dbChainMockFns = {
limit,
offset,
orderBy,
unionAll,
returning,
innerJoin,
leftJoin,
Expand Down
Loading