diff --git a/apps/sim/lib/table/jobs/service.ts b/apps/sim/lib/table/jobs/service.ts index d9519ce0c16..fe220783025 100644 --- a/apps/sim/lib/table/jobs/service.ts +++ b/apps/sim/lib/table/jobs/service.ts @@ -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 { @@ -360,27 +360,36 @@ export async function selectExportRowPage( limit: number ): Promise> { 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 })) diff --git a/apps/sim/lib/table/rows/row-writes.integration.ts b/apps/sim/lib/table/rows/row-writes.integration.ts index 4ebc908c01a..84c3e7bb17b 100644 --- a/apps/sim/lib/table/rows/row-writes.integration.ts +++ b/apps/sim/lib/table/rows/row-writes.integration.ts @@ -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, @@ -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 | 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) + } + }) + }) }) diff --git a/apps/sim/lib/table/rows/service.ts b/apps/sim/lib/table/rows/service.ts index a93dc4e12ec..4874906e43e 100644 --- a/apps/sim/lib/table/rows/service.ts +++ b/apps/sim/lib/table/rows/service.ts @@ -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, @@ -1492,24 +1492,45 @@ export async function fetchRowsBounded(params: BoundedFetchParams): Promise { 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 } diff --git a/apps/sim/lib/table/service-filter-threading.test.ts b/apps/sim/lib/table/service-filter-threading.test.ts index e09632de11c..7577b1e7a15 100644 --- a/apps/sim/lib/table/service-filter-threading.test.ts +++ b/apps/sim/lib/table/service-filter-threading.test.ts @@ -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' @@ -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 } @@ -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. @@ -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, @@ -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') }) @@ -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) }) @@ -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), diff --git a/packages/testing/src/mocks/database.mock.ts b/packages/testing/src/mocks/database.mock.ts index c2406d06b59..2187accd0e2 100644 --- a/packages/testing/src/mocks/database.mock.ts +++ b/packages/testing/src/mocks/database.mock.ts @@ -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 @@ -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() @@ -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 } @@ -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, () => { @@ -371,6 +377,7 @@ export const dbChainMockFns = { limit, offset, orderBy, + unionAll, returning, innerJoin, leftJoin,