Skip to content

Commit db4bf8a

Browse files
committed
fix(tables): keep a batched write that already committed part of its work out of the retryable 503
1 parent 7c8a5ea commit db4bf8a

4 files changed

Lines changed: 222 additions & 88 deletions

File tree

‎apps/sim/lib/table/api/write-contention.integration.ts‎

Lines changed: 83 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -23,7 +23,9 @@ import {
2323
internalTableRowsErrorPolicy,
2424
v2TableRowsErrorPolicy,
2525
} from '@/lib/table/api/row-route-policies'
26-
import { insertRow, updateRowsByFilter } from '@/lib/table/rows/service'
26+
import { getDeleteSnapshotBatchSize, TABLE_LIMITS } from '@/lib/table/constants'
27+
import { TablePartialWriteError } from '@/lib/table/errors'
28+
import { deleteRowsByFilter, insertRow, updateRowsByFilter } from '@/lib/table/rows/service'
2729
import { getTableById } from '@/lib/table/service'
2830
import type { TableDefinition } from '@/lib/table/types'
2931
import { orchestrationErrorResponse } from '@/app/api/table/utils'
@@ -61,6 +63,26 @@ async function createTable(): Promise<TableDefinition> {
6163
return table
6264
}
6365

66+
/** Seeds `count` rows whose id order and row order agree, created before any write's cutoff. */
67+
async function seedNotes(tableId: string, count: number): Promise<string[]> {
68+
const ids = Array.from(
69+
{ length: count },
70+
(_, index) => `${tableId}-${String(index).padStart(4, '0')}`
71+
)
72+
const createdAt = new Date(Date.now() - 60_000)
73+
await db.insert(userTableRows).values(
74+
ids.map((id, index) => ({
75+
id,
76+
tableId,
77+
workspaceId,
78+
data: { note: 'n' },
79+
orderKey: `a${String(index).padStart(4, '0')}`,
80+
createdAt,
81+
}))
82+
)
83+
return ids
84+
}
85+
6486
/** Runs `write` while a second session holds whatever `hold` locks, and returns what it threw. */
6587
async function failedUnderLock(
6688
hold: (holder: postgres.ReservedSql) => Promise<unknown>,
@@ -92,6 +114,13 @@ function postgresErrorContext(error: unknown): string | undefined {
92114
return undefined
93115
}
94116

117+
/** Every surface leaves the failure to the route's generic 500 rather than inviting a retry. */
118+
async function expectNotRetryableOnAnySurface(error: unknown) {
119+
expect(internalTableRowsErrorPolicy.project(error)).toBeNull()
120+
expect(v2TableRowsErrorPolicy.render(error)).toBeNull()
121+
expect(orchestrationErrorResponse(error)).toBeNull()
122+
}
123+
95124
async function expectRetryableOnEverySurface(error: unknown) {
96125
const internal = internalTableRowsErrorPolicy.project(error)
97126
expect(internal?.status).toBe(503)
@@ -184,4 +213,57 @@ describe('table writes that lose a lock race', () => {
184213
expect(row.data).toEqual({ note: 'n' })
185214
}
186215
)
216+
217+
it('leaves a filtered update that loses a lock race after a committed page non-retryable', async () => {
218+
const table = await createTable()
219+
const ids = await seedNotes(table.id, TABLE_LIMITS.UPDATE_BATCH_SIZE + 1)
220+
const lastId = ids[ids.length - 1]
221+
222+
const error = await failedUnderLock(
223+
(holder) => holder`SELECT 1 FROM user_table_rows
224+
WHERE id = ${lastId} AND table_id = ${table.id} AND workspace_id = ${workspaceId}
225+
FOR UPDATE`,
226+
() =>
227+
updateRowsByFilter(
228+
table,
229+
{
230+
filter: { note: 'n' },
231+
data: { note: 'patched' },
232+
secretProvenance: undefined,
233+
capabilityGovernedUserId: null,
234+
},
235+
'contention-partial-update'
236+
)
237+
)
238+
239+
expect(error).toBeInstanceOf(TablePartialWriteError)
240+
expect(error).toMatchObject({ committedCount: TABLE_LIMITS.UPDATE_BATCH_SIZE })
241+
await expectNotRetryableOnAnySurface(error)
242+
})
243+
244+
it('leaves a limited filtered delete that loses a lock race after a committed batch non-retryable', async () => {
245+
const table = await createTable()
246+
const batchSize = getDeleteSnapshotBatchSize()
247+
const ids = await seedNotes(table.id, batchSize + 1)
248+
const lastId = ids[ids.length - 1]
249+
250+
const error = await failedUnderLock(
251+
(holder) => holder`SELECT 1 FROM user_table_rows
252+
WHERE id = ${lastId} AND table_id = ${table.id} AND workspace_id = ${workspaceId}
253+
FOR UPDATE`,
254+
() =>
255+
deleteRowsByFilter(
256+
table,
257+
{ filter: { note: 'n' }, limit: batchSize + 1 },
258+
'contention-partial-delete'
259+
)
260+
)
261+
262+
expect(error).toBeInstanceOf(TablePartialWriteError)
263+
expect(error).toMatchObject({ committedCount: batchSize })
264+
await expectNotRetryableOnAnySurface(error)
265+
const [{ remaining }] = await control<{ remaining: number }[]>`SELECT count(*)::int AS remaining
266+
FROM user_table_rows WHERE table_id = ${table.id} AND workspace_id = ${workspaceId}`
267+
expect(remaining).toBe(1)
268+
})
187269
})

‎apps/sim/lib/table/api/write-contention.ts‎

Lines changed: 7 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -1,36 +1,33 @@
11
import { createLogger } from '@sim/logger'
22
import { getPostgresCancellationReason, getPostgresErrorCode } from '@sim/utils/errors'
33
import { ADMISSION_RETRY_AFTER_SECONDS } from '@/lib/core/admission/transient-failure'
4+
import { isTableLockRace, TablePartialWriteError } from '@/lib/table/errors'
45

56
const logger = createLogger('TableWriteContention')
67

7-
/** `lock_not_available` (a `lock_timeout` fired) and `deadlock_detected`. */
8-
const WRITE_CONTENTION_SQLSTATES = new Set(['55P03', '40P01'])
9-
108
export const TABLE_WRITE_CONTENTION_MESSAGE =
119
'Table is busy with other writes. Try again in a few seconds.'
1210

1311
export const TABLE_WRITE_CONTENTION_RETRY_AFTER_SECONDS = ADMISSION_RETRY_AFTER_SECONDS
1412

1513
/**
1614
* Whether a table operation failed because its transaction lost a lock race, which callers should
17-
* answer with a retryable 503 instead of a generic 500.
15+
* answer with a retryable 503 rather than a generic 500.
1816
*
1917
* Every row write locks its table's `user_table_definitions` row (the `row_count` and
2018
* `rows_version` triggers, the latter at COMMIT), so a transaction that holds it longer than the
2119
* table's 3 s `lock_timeout` — typically one stuck waiting on a synchronous standby, which keeps its
22-
* locks — fails the writers queued behind it. The failed transaction rolled back. A bulk write that
23-
* commits page by page may have applied earlier pages, but every such write re-applies the same
24-
* patch or filter, so retrying it is safe.
20+
* locks — fails the writers queued behind it. Those writers rolled back, so retrying them is safe.
21+
* A batched write that had already committed part of its work is a {@link TablePartialWriteError}
22+
* instead, and is not retryable as a whole.
2523
*
2624
* Logs the failure, because route builders do not log errors their policy projects. Only the
2725
* SQLSTATE is logged: the driver's message carries the query's bound parameters.
2826
*/
2927
export function isTableWriteContention(error: unknown): boolean {
30-
const code = getPostgresErrorCode(error)
31-
if (!code || !WRITE_CONTENTION_SQLSTATES.has(code)) return false
28+
if (error instanceof TablePartialWriteError || !isTableLockRace(error)) return false
3229
logger.warn('Table operation lost a lock race', {
33-
sqlstate: code,
30+
sqlstate: getPostgresErrorCode(error),
3431
reason: getPostgresCancellationReason(error) ?? null,
3532
})
3633
return true

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

Lines changed: 25 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,30 @@
1+
import { getPostgresErrorCode } from '@sim/utils/errors'
12
import { OrchestrationError } from '@/lib/core/orchestration/types'
23

4+
/** `lock_not_available` (a `lock_timeout` fired) and `deadlock_detected`. */
5+
const LOCK_RACE_SQLSTATES = new Set(['55P03', '40P01'])
6+
7+
/** Whether a table write's transaction failed because it lost a lock race, and so rolled back. */
8+
export function isTableLockRace(error: unknown): boolean {
9+
const code = getPostgresErrorCode(error)
10+
return code !== undefined && LOCK_RACE_SQLSTATES.has(code)
11+
}
12+
13+
/**
14+
* A write that commits batch by batch lost a lock race after at least one batch committed. Unlike
15+
* a single rolled-back transaction it is not retryable as a whole: a re-run selects its rows again,
16+
* so a limited filtered delete would remove more than its limit.
17+
*/
18+
export class TablePartialWriteError extends Error {
19+
constructor(
20+
readonly committedCount: number,
21+
cause: unknown
22+
) {
23+
super(`Table write failed after ${committedCount} rows were committed`, { cause })
24+
this.name = 'TablePartialWriteError'
25+
}
26+
}
27+
328
/** A disabled TTL feature, distinct from malformed column input. */
429
export class TableRowTtlDisabledError extends OrchestrationError {
530
readonly detailCode = 'TABLE_ROW_TTL_DISABLED'

0 commit comments

Comments
 (0)