Skip to content

Commit 7c8a5ea

Browse files
committed
fix(tables): answer a write that loses a lock race with a retryable 503
1 parent c723046 commit 7c8a5ea

4 files changed

Lines changed: 267 additions & 4 deletions

File tree

‎apps/sim/app/api/table/utils.ts‎

Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,11 @@ import {
1717
import { capabilityRefusalResponse } from '@/lib/permission-groups/capability-response'
1818
import type { ColumnDefinition, Filter, TableDefinition, TablePredicate } from '@/lib/table'
1919
import { buildFilterClause, getTableById, TableQueryValidationError } from '@/lib/table'
20+
import {
21+
isTableWriteContention,
22+
TABLE_WRITE_CONTENTION_MESSAGE,
23+
TABLE_WRITE_CONTENTION_RETRY_AFTER_SECONDS,
24+
} from '@/lib/table/api/write-contention'
2025
import { USER_TABLE_ROWS_SQL_NAME } from '@/lib/table/constants'
2126
import { TableLockedError } from '@/lib/table/mutation-locks'
2227
import {
@@ -137,6 +142,16 @@ export function orchestrationErrorResponse(error: unknown): NextResponse | null
137142
const lockResponse = tableLockErrorResponse(error)
138143
if (lockResponse) return lockResponse
139144

145+
if (isTableWriteContention(error)) {
146+
return NextResponse.json(
147+
{ error: TABLE_WRITE_CONTENTION_MESSAGE },
148+
{
149+
status: 503,
150+
headers: { 'Retry-After': String(TABLE_WRITE_CONTENTION_RETRY_AFTER_SECONDS) },
151+
}
152+
)
153+
}
154+
140155
const classified = asOrchestrationError(error)
141156
if (!classified) return null
142157

‎apps/sim/lib/table/api/route-policies.ts‎

Lines changed: 28 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,11 @@ import {
88
type V2ErrorPolicy,
99
} from '@/lib/api/server/routes'
1010
import { asOrchestrationError } from '@/lib/core/orchestration/types'
11+
import {
12+
isTableWriteContention,
13+
TABLE_WRITE_CONTENTION_MESSAGE,
14+
TABLE_WRITE_CONTENTION_RETRY_AFTER_SECONDS,
15+
} from '@/lib/table/api/write-contention'
1116
import { TABLE_DELEGATION_AUDIENCE } from '@/lib/table/application/authorization'
1217
import { TableOperationError } from '@/lib/table/application/errors'
1318
import { TableRowTtlDisabledError } from '@/lib/table/errors'
@@ -27,6 +32,9 @@ export const internalTableSessionOrExecutorAuth = createInternalSessionOrExecuto
2732
})
2833

2934
function renderTableError(error: unknown) {
35+
if (isTableWriteContention(error)) {
36+
return v2Error('SERVICE_UNAVAILABLE', TABLE_WRITE_CONTENTION_MESSAGE)
37+
}
3038
const classified = asOrchestrationError(error)
3139
if (classified instanceof TableRowTtlDisabledError) {
3240
return v2Error('BAD_REQUEST', classified.message, {
@@ -76,8 +84,24 @@ export const v2TableErrorPolicies = {
7684
} satisfies V2ErrorPolicy,
7785
} as const
7886

79-
const internalTableGroupErrorPolicy = extendInternalErrorPolicy(
87+
/**
88+
* The base of every internal table policy: a transaction that lost a lock race answers a
89+
* retryable 503 rather than falling through to the route's generic 500.
90+
*/
91+
const internalTableBaseErrorPolicy = extendInternalErrorPolicy(
8092
internalOrchestrationErrorPolicy,
93+
(error) =>
94+
isTableWriteContention(error)
95+
? internalErrorResponse(
96+
503,
97+
{ error: TABLE_WRITE_CONTENTION_MESSAGE },
98+
{ 'Retry-After': String(TABLE_WRITE_CONTENTION_RETRY_AFTER_SECONDS) }
99+
)
100+
: null
101+
)
102+
103+
const internalTableGroupErrorPolicy = extendInternalErrorPolicy(
104+
internalTableBaseErrorPolicy,
81105
(error) =>
82106
error instanceof TableLockedError
83107
? internalErrorResponse(423, { error: error.message, lock: error.lock })
@@ -99,19 +123,19 @@ export const internalTableErrorPolicies = {
99123
*/
100124
bulk: internalTableGroupErrorPolicy,
101125
concealTableAuthorization: createInternalResourceConcealmentPolicy({
102-
base: internalOrchestrationErrorPolicy,
126+
base: internalTableBaseErrorPolicy,
103127
notFoundMessage: 'Table not found',
104128
}),
105129
concealTableGroupAuthorization: createInternalResourceConcealmentPolicy({
106130
base: internalTableGroupErrorPolicy,
107131
notFoundMessage: 'Table not found',
108132
}),
109133
concealImportAuthorization: createInternalResourceConcealmentPolicy({
110-
base: internalOrchestrationErrorPolicy,
134+
base: internalTableBaseErrorPolicy,
111135
notFoundMessage: 'Table import not found',
112136
}),
113137
concealExportAuthorization: createInternalResourceConcealmentPolicy({
114-
base: internalOrchestrationErrorPolicy,
138+
base: internalTableBaseErrorPolicy,
115139
notFoundMessage: 'Table export not found',
116140
}),
117141
} as const
Lines changed: 187 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,187 @@
1+
/**
2+
* A row write that loses a lock race must answer every table surface with a retryable 503, not the
3+
* generic 500. The races are real: a second session holds the lock the write needs until the
4+
* write's 3 s `lock_timeout` fires, and the error the service throws is projected as thrown.
5+
*/
6+
import { db } from '@sim/db'
7+
import { userTableDefinitions, userTableRows } from '@sim/db/schema'
8+
import { readTestDatabaseUrl } from '@sim/db/testing/test-infrastructure'
9+
import { storageServiceMock } from '@sim/testing/mocks/storage-service.mock'
10+
import { tableBillingMock, tableBillingMockFns } from '@sim/testing/mocks/table-billing.mock'
11+
import { tableTriggerMock } from '@sim/testing/mocks/table-trigger.mock'
12+
import { tableWorkflowColumnsMock } from '@sim/testing/mocks/table-workflow-columns.mock'
13+
import { generateId } from '@sim/utils/id'
14+
import postgres from 'postgres'
15+
import { afterAll, beforeAll, beforeEach, describe, expect, it, vi } from 'vitest'
16+
17+
vi.mock('@/lib/table/billing', () => tableBillingMock)
18+
vi.mock('@/lib/table/trigger', () => tableTriggerMock)
19+
vi.mock('@/lib/table/workflow-columns', () => tableWorkflowColumnsMock)
20+
vi.mock('@/lib/uploads/core/storage-service', () => storageServiceMock)
21+
22+
import {
23+
internalTableRowsErrorPolicy,
24+
v2TableRowsErrorPolicy,
25+
} from '@/lib/table/api/row-route-policies'
26+
import { insertRow, updateRowsByFilter } from '@/lib/table/rows/service'
27+
import { getTableById } from '@/lib/table/service'
28+
import type { TableDefinition } from '@/lib/table/types'
29+
import { orchestrationErrorResponse } from '@/app/api/table/utils'
30+
31+
const url = readTestDatabaseUrl()
32+
if (process.env.DATABASE_URL !== url) {
33+
throw new Error('This suite requires only the disposable local test database')
34+
}
35+
const control = postgres(url, { max: 2, onnotice: () => {} })
36+
const workspaceId = generateId()
37+
const userId = generateId()
38+
39+
/** Only the migrated deferred trigger bumps `rows_version` at COMMIT; `db:push` installs none. */
40+
const [{ migrated }] = await control<{ migrated: boolean }[]>`SELECT EXISTS (
41+
SELECT 1 FROM pg_trigger t
42+
JOIN pg_proc p ON p.oid = t.tgfoid
43+
WHERE t.tgrelid = 'user_table_rows'::regclass
44+
AND t.tgname = 'user_table_rows_version_update_trigger'
45+
AND t.tgconstraint <> 0
46+
AND t.tginitdeferred
47+
AND p.proname = 'bump_user_table_rows_version_at_commit'
48+
) AS migrated`
49+
50+
async function createTable(): Promise<TableDefinition> {
51+
const id = generateId()
52+
await db.insert(userTableDefinitions).values({
53+
id,
54+
workspaceId,
55+
name: id,
56+
schema: { columns: [{ id: 'note', name: 'note', type: 'string' }] },
57+
createdBy: userId,
58+
})
59+
const table = await getTableById(id)
60+
if (!table) throw new Error('Fixture table was not created')
61+
return table
62+
}
63+
64+
/** Runs `write` while a second session holds whatever `hold` locks, and returns what it threw. */
65+
async function failedUnderLock(
66+
hold: (holder: postgres.ReservedSql) => Promise<unknown>,
67+
write: () => Promise<unknown>
68+
): Promise<unknown> {
69+
const holder = await control.reserve()
70+
try {
71+
await holder`BEGIN`
72+
await hold(holder)
73+
return await write().then(
74+
() => {
75+
throw new Error('The write succeeded while its lock was held')
76+
},
77+
(error: unknown) => error
78+
)
79+
} finally {
80+
await holder`ROLLBACK`
81+
holder.release()
82+
}
83+
}
84+
85+
/** The PL/pgSQL context Postgres attached to the error, from the driver error under Drizzle's wrapper. */
86+
function postgresErrorContext(error: unknown): string | undefined {
87+
let current = error
88+
while (current instanceof Error) {
89+
if ('where' in current && typeof current.where === 'string') return current.where
90+
current = current.cause
91+
}
92+
return undefined
93+
}
94+
95+
async function expectRetryableOnEverySurface(error: unknown) {
96+
const internal = internalTableRowsErrorPolicy.project(error)
97+
expect(internal?.status).toBe(503)
98+
expect(new Headers(internal?.headers).get('Retry-After')).toBe('5')
99+
100+
const v2 = v2TableRowsErrorPolicy.render(error)
101+
expect(v2?.status).toBe(503)
102+
expect(v2?.headers.get('Retry-After')).toBe('5')
103+
expect(await v2?.json()).toMatchObject({ error: { code: 'SERVICE_UNAVAILABLE' } })
104+
105+
const legacy = orchestrationErrorResponse(error)
106+
expect(legacy?.status).toBe(503)
107+
expect(legacy?.headers.get('Retry-After')).toBe('5')
108+
}
109+
110+
describe('table writes that lose a lock race', () => {
111+
beforeAll(async () => {
112+
await control`INSERT INTO "user" (id, name, email, email_verified, created_at, updated_at)
113+
VALUES (${userId}, 'Write contention fixture', ${`${userId}@example.test`}, true, now(), now())`
114+
await control`INSERT INTO workspace (id, name, owner_id, billed_account_user_id)
115+
VALUES (${workspaceId}, 'Write contention fixtures', ${userId}, ${userId})`
116+
})
117+
118+
beforeEach(() => {
119+
tableBillingMockFns.mockAssertRowCapacity.mockResolvedValue(10_000)
120+
})
121+
122+
afterAll(async () => {
123+
await control`DELETE FROM workspace WHERE id = ${workspaceId}`
124+
await control`DELETE FROM "user" WHERE id = ${userId}`
125+
await control.end()
126+
})
127+
128+
it('answers an insert blocked on the row-order lock with a retryable 503', async () => {
129+
const table = await createTable()
130+
131+
const error = await failedUnderLock(
132+
(holder) =>
133+
holder`SELECT pg_advisory_xact_lock(hashtextextended(${`user_table_rows_pos:${table.id}`}, 0))`,
134+
() =>
135+
insertRow(
136+
{
137+
tableId: table.id,
138+
workspaceId,
139+
data: { note: 'blocked' },
140+
secretProvenance: undefined,
141+
capabilityGovernedUserId: null,
142+
},
143+
table,
144+
'contention-insert'
145+
)
146+
)
147+
148+
await expectRetryableOnEverySurface(error)
149+
})
150+
151+
it.skipIf(!migrated)(
152+
'answers an update whose commit cannot take the definition row with a retryable 503',
153+
async () => {
154+
const table = await createTable()
155+
await db.insert(userTableRows).values({
156+
id: `${table.id}-row`,
157+
tableId: table.id,
158+
workspaceId,
159+
data: { note: 'n' },
160+
orderKey: 'a0',
161+
})
162+
163+
const error = await failedUnderLock(
164+
(holder) =>
165+
holder`SELECT 1 FROM user_table_definitions WHERE id = ${table.id} FOR NO KEY UPDATE`,
166+
() =>
167+
updateRowsByFilter(
168+
table,
169+
{
170+
filter: { note: 'n' },
171+
data: { note: 'patched' },
172+
limit: 1,
173+
secretProvenance: undefined,
174+
capabilityGovernedUserId: null,
175+
},
176+
'contention-update'
177+
)
178+
)
179+
180+
expect(postgresErrorContext(error)).toContain('bump_user_table_rows_version_at_commit')
181+
await expectRetryableOnEverySurface(error)
182+
const [row] = await control`SELECT data FROM user_table_rows
183+
WHERE id = ${`${table.id}-row`} AND table_id = ${table.id} AND workspace_id = ${workspaceId}`
184+
expect(row.data).toEqual({ note: 'n' })
185+
}
186+
)
187+
})
Lines changed: 37 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,37 @@
1+
import { createLogger } from '@sim/logger'
2+
import { getPostgresCancellationReason, getPostgresErrorCode } from '@sim/utils/errors'
3+
import { ADMISSION_RETRY_AFTER_SECONDS } from '@/lib/core/admission/transient-failure'
4+
5+
const logger = createLogger('TableWriteContention')
6+
7+
/** `lock_not_available` (a `lock_timeout` fired) and `deadlock_detected`. */
8+
const WRITE_CONTENTION_SQLSTATES = new Set(['55P03', '40P01'])
9+
10+
export const TABLE_WRITE_CONTENTION_MESSAGE =
11+
'Table is busy with other writes. Try again in a few seconds.'
12+
13+
export const TABLE_WRITE_CONTENTION_RETRY_AFTER_SECONDS = ADMISSION_RETRY_AFTER_SECONDS
14+
15+
/**
16+
* 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.
18+
*
19+
* Every row write locks its table's `user_table_definitions` row (the `row_count` and
20+
* `rows_version` triggers, the latter at COMMIT), so a transaction that holds it longer than the
21+
* 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.
25+
*
26+
* Logs the failure, because route builders do not log errors their policy projects. Only the
27+
* SQLSTATE is logged: the driver's message carries the query's bound parameters.
28+
*/
29+
export function isTableWriteContention(error: unknown): boolean {
30+
const code = getPostgresErrorCode(error)
31+
if (!code || !WRITE_CONTENTION_SQLSTATES.has(code)) return false
32+
logger.warn('Table operation lost a lock race', {
33+
sqlstate: code,
34+
reason: getPostgresCancellationReason(error) ?? null,
35+
})
36+
return true
37+
}

0 commit comments

Comments
 (0)