|
17 | 17 | */ |
18 | 18 |
|
19 | 19 | import { db } from '@sim/db' |
20 | | -import { userTableDefinitions, userTableRowChanges } from '@sim/db/schema' |
| 20 | +import { userTableDefinitions } from '@sim/db/schema' |
21 | 21 | import { createLogger } from '@sim/logger' |
22 | | -import { and, eq, notInArray, sql } from 'drizzle-orm' |
| 22 | +import { and, eq, sql } from 'drizzle-orm' |
23 | 23 | import { setTableTxTimeouts } from '@/lib/table/tx' |
24 | 24 |
|
25 | 25 | const logger = createLogger('TableRowChanges') |
26 | 26 |
|
27 | | -/** Tables folded per sweep query; the sweep loops until the log is empty or its budget runs out. */ |
| 27 | +/** Tables per pending-table page; the sweep pages until the log is empty or its budget runs out. */ |
28 | 28 | const FOLD_SWEEP_PAGE_SIZE = 500 |
29 | 29 |
|
30 | 30 | /** |
@@ -100,36 +100,55 @@ export async function foldTableRowChanges(tableId: string): Promise<boolean> { |
100 | 100 | }) |
101 | 101 | } |
102 | 102 |
|
| 103 | +/** Outcome of one fold sweep. */ |
103 | 104 | export interface FoldSweepResult { |
| 105 | + /** Tables whose log rows were folded into their definition row. */ |
104 | 106 | folded: number |
| 107 | + /** Tables left for the next sweep because another transaction held their definition row. */ |
105 | 108 | skipped: number |
| 109 | + /** Whether the sweep stopped at its budget with tables still unvisited. */ |
106 | 110 | budgetExhausted: boolean |
107 | 111 | } |
108 | 112 |
|
109 | 113 | /** |
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. |
| 114 | + * The next page of distinct table ids in the log after `afterTableId`, by loose index scan: each |
| 115 | + * step seeks the next id through the `table_id` index, so the cost follows the number of tables, |
| 116 | + * not the number of log rows a busy table has piled up. |
| 117 | + */ |
| 118 | +async function nextPendingTableIds(afterTableId: string): Promise<string[]> { |
| 119 | + const rows = await db.execute<{ table_id: string }>(sql` |
| 120 | + WITH RECURSIVE pending AS ( |
| 121 | + (SELECT table_id FROM user_table_row_changes |
| 122 | + WHERE table_id > ${afterTableId} ORDER BY table_id LIMIT 1) |
| 123 | + UNION ALL |
| 124 | + SELECT (SELECT c.table_id FROM user_table_row_changes c |
| 125 | + WHERE c.table_id > p.table_id ORDER BY c.table_id LIMIT 1) |
| 126 | + FROM pending p WHERE p.table_id IS NOT NULL |
| 127 | + ) |
| 128 | + SELECT table_id FROM pending WHERE table_id IS NOT NULL LIMIT ${FOLD_SWEEP_PAGE_SIZE}`) |
| 129 | + return rows.map((row) => row.table_id) |
| 130 | +} |
| 131 | + |
| 132 | +/** |
| 133 | + * Folds each table with unfolded log rows once, in table-id order, until none is left or |
| 134 | + * `budgetMs` elapses. Rows a table logs after its fold wait for the next sweep. |
112 | 135 | */ |
113 | 136 | export async function foldPendingTableRowChanges(budgetMs: number): Promise<FoldSweepResult> { |
114 | 137 | const deadline = Date.now() + budgetMs |
115 | | - const visited = new Set<string>() |
| 138 | + let afterTableId = '' |
116 | 139 | let folded = 0 |
117 | 140 | let skipped = 0 |
118 | 141 |
|
119 | 142 | 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 } |
| 143 | + const tableIds = await nextPendingTableIds(afterTableId) |
| 144 | + if (tableIds.length === 0) return { folded, skipped, budgetExhausted: false } |
126 | 145 |
|
127 | | - for (const { tableId } of pending) { |
| 146 | + for (const tableId of tableIds) { |
128 | 147 | if (Date.now() >= deadline) { |
129 | 148 | logger.warn('Table row-change fold sweep ran out of budget', { folded, skipped, budgetMs }) |
130 | 149 | return { folded, skipped, budgetExhausted: true } |
131 | 150 | } |
132 | | - visited.add(tableId) |
| 151 | + afterTableId = tableId |
133 | 152 | if (await foldTableRowChanges(tableId)) folded++ |
134 | 153 | else skipped++ |
135 | 154 | } |
|
0 commit comments