Skip to content

Commit 06a6262

Browse files
committed
fix(jobs): bound the retry-history read, count member purges, match driver connection errors
- the connector run-history read runs under short statement and lock timeouts and falls back to this run alone - members-mode progress counts lifecycle purges (docs_purged in the run log) - postgres.js connection codes count only on an error the driver built, or under a query error; the driver's query signature is its four own properties, which a refused connection carries with no SQL yet - real-error PostgreSQL test for terminated transactions, refused connections, and an ending pool, wired into CI - knowledge-processing tests use the static task import
1 parent 91c0228 commit 06a6262

10 files changed

Lines changed: 311 additions & 30 deletions

File tree

‎.github/workflows/test-build.yml‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -102,6 +102,7 @@ jobs:
102102
bunx vitest run
103103
scripts/retired-columns.postgres.test.ts
104104
scripts/connector-sync-schedule-precision.postgres.test.ts
105+
scripts/database-failure-classification.postgres.test.ts
105106
106107
- name: Verify OAuth lifecycle and SCIM membership guards in PostgreSQL
107108
working-directory: apps/sim

‎apps/sim/background/knowledge-processing.test.ts‎

Lines changed: 1 addition & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -49,6 +49,7 @@ import {
4949
DOCUMENT_PROCESSING_RETRY_POLICY,
5050
DocumentProcessingDatabaseRetryError,
5151
getDocumentProcessingRetry,
52+
processDocument,
5253
resolveQuotaContinuationDelayMs,
5354
runDocumentProcessing,
5455
} from '@/background/knowledge-processing'
@@ -821,14 +822,10 @@ describe('knowledge-process-document task configuration', () => {
821822
* `attempt_count = 1`, so each was left `failed` having never been retried.
822823
*/
823824
it('escalates to a larger machine on an out-of-memory kill', async () => {
824-
const { processDocument } = await import('@/background/knowledge-processing')
825-
826825
expect(processDocument.retry?.outOfMemory?.machine).toBe('large-2x')
827826
})
828827

829828
it('declares enough attempts for database retries and routes failures through catchError', async () => {
830-
const { processDocument } = await import('@/background/knowledge-processing')
831-
832829
expect(processDocument.retry?.maxAttempts).toBe(
833830
Math.max(
834831
DOCUMENT_PROCESSING_RETRY_POLICY.maxAttempts,

‎apps/sim/lib/knowledge/connectors/member-sync-engine.test.ts‎

Lines changed: 24 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -25,8 +25,10 @@ import {
2525
buildMemberSyncDatabaseRetryUpdate,
2626
buildMemberSyncFailureUpdate,
2727
deriveMemberActive,
28+
type MemberSyncResult,
2829
memberFailureBackoffMs,
2930
memberNextAttemptAt,
31+
memberRunMadeProgress,
3032
nextMemberSyncTime,
3133
resolveMemberSyncFailureUpdate,
3234
shouldListFully,
@@ -285,6 +287,28 @@ describe('member sync engine decisions', () => {
285287
})
286288
})
287289

290+
describe('memberRunMadeProgress', () => {
291+
const idle = {
292+
membersCompleted: 0,
293+
docsAdded: 0,
294+
docsUpdated: 0,
295+
docsDeleted: 0,
296+
} as MemberSyncResult
297+
298+
it('reports no progress for a run that wrote nothing', () => {
299+
expect(memberRunMadeProgress(idle)).toBe(false)
300+
})
301+
302+
it.each([
303+
['completed a member', { membersCompleted: 1 }],
304+
['added documents', { docsAdded: 2 }],
305+
['updated documents', { docsUpdated: 1 }],
306+
['purged documents in the lifecycle pass', { docsDeleted: 3 }],
307+
])('reports progress for a run that %s', (_label, writes) => {
308+
expect(memberRunMadeProgress({ ...idle, ...writes })).toBe(true)
309+
})
310+
})
311+
288312
describe('resolveMemberSyncFailureUpdate', () => {
289313
beforeEach(() => {
290314
resetDbChainMock()

‎apps/sim/lib/knowledge/connectors/member-sync-engine.ts‎

Lines changed: 9 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -350,6 +350,14 @@ export function buildMemberSyncDatabaseRetryUpdate(
350350
}
351351
}
352352

353+
/**
354+
* Whether a members-mode run moved the sync forward: it completed a member, or wrote documents.
355+
* `docsDeleted` holds the document lifecycle's purges, which the run log records as `docs_purged`.
356+
*/
357+
export function memberRunMadeProgress(result: MemberSyncResult): boolean {
358+
return result.membersCompleted + result.docsAdded + result.docsUpdated + result.docsDeleted > 0
359+
}
360+
353361
/**
354362
* The connector row a failed members-mode run writes. A deterministic capacity rejection waits for
355363
* an operator, a transient database failure retries without touching the breaker, and anything
@@ -2502,7 +2510,7 @@ export async function executeMemberSync(
25022510
previousFailures: connector.memberSyncConsecutiveFailures,
25032511
errorMessage,
25042512
retryAfterMs,
2505-
madeProgress: result.membersCompleted + result.docsAdded + result.docsUpdated > 0,
2513+
madeProgress: memberRunMadeProgress(result),
25062514
})
25072515
const written = await db
25082516
.update(knowledgeConnector)

‎apps/sim/lib/knowledge/connectors/sync-database-retry.test.ts‎

Lines changed: 33 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,8 @@ import {
2121
DATABASE_FAILURE_ALERT_STREAK,
2222
DATABASE_RETRY_AFTER_PROGRESS_MS,
2323
databaseRetryDelayMs,
24+
RUN_HISTORY_LOCK_TIMEOUT_MS,
25+
RUN_HISTORY_STATEMENT_TIMEOUT_MS,
2426
resolveDatabaseRetryDelayMs,
2527
} from '@/lib/knowledge/connectors/sync-database-retry'
2628
import {
@@ -40,6 +42,7 @@ const memberRun = (status: string, writes: Record<string, number> = {}) => ({
4042
membersCompleted: 0,
4143
docsAdded: 0,
4244
docsUpdated: 0,
45+
docsPurged: 0,
4346
...writes,
4447
})
4548

@@ -99,6 +102,7 @@ describe('countZeroProgressFailedRuns', () => {
99102
['completed a member', { membersCompleted: 1 }],
100103
['added documents', { docsAdded: 3 }],
101104
['updated documents', { docsUpdated: 1 }],
105+
['purged documents', { docsPurged: 4 }],
102106
])('ends a members-mode streak at a failed run that %s', async (_label, writes) => {
103107
queueTableRows(schemaMock.knowledgeConnectorMemberSyncLog, [
104108
memberRun('failed'),
@@ -126,6 +130,35 @@ describe('countZeroProgressFailedRuns', () => {
126130
)
127131
})
128132

133+
it('bounds the history read with its own statement and lock timeouts', async () => {
134+
queueTableRows(schemaMock.knowledgeConnectorSyncLog, [])
135+
await countZeroProgressFailedRuns('content', 'c-1', 'run-1')
136+
const bound = dbChainMockFns.execute.mock.calls[0]?.[0] as {
137+
toSQL: () => { sql: string; params: unknown[] }
138+
}
139+
const { sql, params } = bound.toSQL()
140+
expect(sql).toContain("set_config('statement_timeout'")
141+
expect(sql).toContain("set_config('lock_timeout'")
142+
expect(params).toEqual([
143+
String(RUN_HISTORY_STATEMENT_TIMEOUT_MS),
144+
String(RUN_HISTORY_LOCK_TIMEOUT_MS),
145+
])
146+
expect(dbChainMockFns.execute.mock.invocationCallOrder[0]).toBeLessThan(
147+
dbChainMockFns.select.mock.invocationCallOrder[0]
148+
)
149+
})
150+
151+
it('falls back to this run alone when the history read times out', async () => {
152+
queueTableRows(schemaMock.knowledgeConnectorSyncLog, [
153+
contentRun('failed'),
154+
contentRun('failed'),
155+
])
156+
dbChainMockFns.limit.mockRejectedValueOnce(
157+
Object.assign(new Error('canceling statement due to statement timeout'), { code: '57014' })
158+
)
159+
expect(await countZeroProgressFailedRuns('content', 'c-1', 'run-1')).toBe(1)
160+
})
161+
129162
it('falls back to this run alone when the history cannot be read', async () => {
130163
dbChainMockFns.limit.mockRejectedValueOnce(new Error('canceling statement'))
131164
expect(await countZeroProgressFailedRuns('content', 'c-1', 'run-1')).toBe(1)

‎apps/sim/lib/knowledge/connectors/sync-database-retry.ts‎

Lines changed: 26 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -3,7 +3,8 @@ import { knowledgeConnectorMemberSyncLog, knowledgeConnectorSyncLog } from '@sim
33
import { createLogger } from '@sim/logger'
44
import { describeError } from '@sim/utils/errors'
55
import { randomInt } from '@sim/utils/random'
6-
import { and, desc, eq, ne } from 'drizzle-orm'
6+
import { and, desc, eq, ne, sql } from 'drizzle-orm'
7+
import type { DbTransaction } from '@/lib/db/types'
78
import {
89
CONNECTOR_FAILURE_BACKOFF_CAP_MINUTES,
910
CONNECTOR_FAILURE_BACKOFF_STEP_MINUTES,
@@ -33,6 +34,14 @@ const LADDER_RUNGS = Math.ceil(
3334
CONNECTOR_FAILURE_BACKOFF_CAP_MINUTES / CONNECTOR_FAILURE_BACKOFF_STEP_MINUTES
3435
)
3536

37+
/**
38+
* Limits on the run-history read. It runs on the failure path, usually while the database is the
39+
* thing failing, so it must give up fast rather than queue behind the slow window; a read that
40+
* times out falls back to counting this run alone.
41+
*/
42+
export const RUN_HISTORY_STATEMENT_TIMEOUT_MS = 2_000
43+
export const RUN_HISTORY_LOCK_TIMEOUT_MS = 500
44+
3645
/** Which run log a sync writes: the content sync's, or the members-mode run's. */
3746
export type SyncRunLogKind = 'content' | 'member'
3847

@@ -43,14 +52,15 @@ interface LoggedRun {
4352

4453
/** Earlier runs of this connector, newest first, and whether each one moved the sync forward. */
4554
async function readEarlierRuns(
55+
tx: DbTransaction,
4656
kind: SyncRunLogKind,
4757
connectorId: string,
4858
runId: string,
4959
limit: number
5060
): Promise<LoggedRun[]> {
5161
if (kind === 'content') {
5262
const log = knowledgeConnectorSyncLog
53-
const rows = await db
63+
const rows = await tx
5464
.select({
5565
status: log.status,
5666
docsAdded: log.docsAdded,
@@ -67,20 +77,21 @@ async function readEarlierRuns(
6777
}))
6878
}
6979
const log = knowledgeConnectorMemberSyncLog
70-
const rows = await db
80+
const rows = await tx
7181
.select({
7282
status: log.status,
7383
membersCompleted: log.membersCompleted,
7484
docsAdded: log.docsAdded,
7585
docsUpdated: log.docsUpdated,
86+
docsPurged: log.docsPurged,
7687
})
7788
.from(log)
7889
.where(and(eq(log.connectorId, connectorId), ne(log.id, runId)))
7990
.orderBy(desc(log.startedAt))
8091
.limit(limit)
8192
return rows.map((row) => ({
8293
status: row.status,
83-
progressed: row.membersCompleted + row.docsAdded + row.docsUpdated > 0,
94+
progressed: row.membersCompleted + row.docsAdded + row.docsUpdated + row.docsPurged > 0,
8495
}))
8596
}
8697

@@ -93,17 +104,24 @@ async function readEarlierRuns(
93104
* cannot spend the breaker that disables connectors for persistent source failures. The run log
94105
* already records every attempt and what it wrote, so it measures the streak instead: a statement
95106
* that fails every run without progress still backs off rung by rung, while a run that added,
96-
* updated, or deleted documents (or, in members mode, completed a member) before failing ends it.
97-
* The read uses the log's `(connector_id, started_at DESC)` index. If it fails too, the streak
98-
* counts only this run and the caller's own floor applies.
107+
* updated, or deleted documents (in members mode, purged by the document lifecycle and logged as
108+
* `docs_purged`), or completed a member, before failing ends it.
109+
* The read uses the log's `(connector_id, started_at DESC)` index and is bounded by
110+
* {@link RUN_HISTORY_STATEMENT_TIMEOUT_MS}. If it fails or times out, the streak counts only this
111+
* run and the caller's own floor applies.
99112
*/
100113
export async function countZeroProgressFailedRuns(
101114
kind: SyncRunLogKind,
102115
connectorId: string,
103116
runId: string
104117
): Promise<number> {
105118
try {
106-
const earlier = await readEarlierRuns(kind, connectorId, runId, LADDER_RUNGS - 1)
119+
const earlier = await db.transaction(async (tx) => {
120+
await tx.execute(
121+
sql`SELECT set_config('statement_timeout', ${String(RUN_HISTORY_STATEMENT_TIMEOUT_MS)}, true), set_config('lock_timeout', ${String(RUN_HISTORY_LOCK_TIMEOUT_MS)}, true)`
122+
)
123+
return readEarlierRuns(tx, kind, connectorId, runId, LADDER_RUNGS - 1)
124+
})
107125
const streakEnd = earlier.findIndex((run) => run.status !== 'failed' || run.progressed)
108126
return 1 + (streakEnd === -1 ? earlier.length : streakEnd)
109127
} catch (error) {

‎packages/db/script-migrations/0021_embedding_search_connector.test.ts‎

Lines changed: 12 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -20,8 +20,19 @@ const SERVER_MESSAGES: Record<string, string> = {
2020
'40P01': 'deadlock detected',
2121
}
2222

23-
/** A driver error carrying a SQLSTATE, the shape `postgres` throws: the server's message as is. */
23+
/**
24+
* The error `postgres` throws for `code`: a SQLSTATE carries the server's message as is, and a
25+
* lost connection is the driver's own connection error (`write <code> <host:port>`).
26+
*/
2427
function postgresError(code: string, message = SERVER_MESSAGES[code] ?? 'failed'): Error {
28+
if (code.startsWith('CONNECTION_')) {
29+
return Object.assign(new Error(`write ${code} localhost:5432`), {
30+
code,
31+
errno: code,
32+
address: 'localhost',
33+
port: 5432,
34+
})
35+
}
2536
return Object.assign(new Error(message), { code })
2637
}
2738

Lines changed: 80 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,80 @@
1+
import { migrationTestDatabaseUrl } from '@sim/db/scripts/migration-fixture'
2+
import { classifyDatabaseFailure } from '@sim/utils/errors'
3+
import { sql as statement } from 'drizzle-orm'
4+
import { drizzle } from 'drizzle-orm/postgres-js'
5+
import postgres from 'postgres'
6+
import { describe, expect, it } from 'vitest'
7+
8+
/** A port on the database host with nothing listening, so a connection is refused at once. */
9+
function closedPortUrl(): string {
10+
const url = new URL(migrationTestDatabaseUrl!)
11+
url.port = '1'
12+
return url.toString()
13+
}
14+
15+
async function rejectionOf(work: () => Promise<unknown>): Promise<unknown> {
16+
try {
17+
await work()
18+
} catch (error) {
19+
return error
20+
}
21+
throw new Error('Expected the work to fail')
22+
}
23+
24+
/**
25+
* Classifies the errors postgres.js and Drizzle actually raise, rather than hand-built shapes:
26+
* a transaction that loses its connection is rejected with the driver's bare connection error,
27+
* with no query attached and no Drizzle wrapper, and must still read as a connection failure.
28+
*/
29+
describe.skipIf(!migrationTestDatabaseUrl)('database failure classification', () => {
30+
it('reads a connection terminated mid-transaction as a connection failure', async () => {
31+
const admin = postgres(migrationTestDatabaseUrl!, { max: 1 })
32+
const client = postgres(migrationTestDatabaseUrl!, { max: 1 })
33+
try {
34+
const error = await rejectionOf(() =>
35+
drizzle(client).transaction(async (tx) => {
36+
const [row] = await tx.execute<{ pid: number }>(statement`SELECT pg_backend_pid() AS pid`)
37+
await admin`SELECT pg_terminate_backend(${row.pid})`
38+
await tx.execute(statement`SELECT pg_sleep(0.2)`)
39+
})
40+
)
41+
expect(classifyDatabaseFailure(error)).toBe('connection')
42+
} finally {
43+
await client.end({ timeout: 1 }).catch(() => {})
44+
await admin.end()
45+
}
46+
})
47+
48+
it('reads a refused connection as a connection failure, in a query and in a transaction', async () => {
49+
const client = postgres(closedPortUrl(), { max: 1, connect_timeout: 2 })
50+
try {
51+
const db = drizzle(client)
52+
expect(
53+
classifyDatabaseFailure(await rejectionOf(() => db.execute(statement`SELECT 1`)))
54+
).toBe('connection')
55+
expect(
56+
classifyDatabaseFailure(
57+
await rejectionOf(() => db.transaction((tx) => tx.execute(statement`SELECT 1`)))
58+
)
59+
).toBe('connection')
60+
} finally {
61+
await client.end({ timeout: 1 }).catch(() => {})
62+
}
63+
})
64+
65+
it('reads a transaction begun on an ending pool as a connection failure', async () => {
66+
const client = postgres(migrationTestDatabaseUrl!, { max: 1 })
67+
await client`SELECT 1`
68+
const ending = client.end({ timeout: 5 })
69+
const error = await rejectionOf(() =>
70+
drizzle(client).transaction((tx) => tx.execute(statement`SELECT 1`))
71+
)
72+
await ending
73+
expect(classifyDatabaseFailure(error)).toBe('connection')
74+
})
75+
76+
it('does not read a refused connection from another client as a database failure', async () => {
77+
const error = await rejectionOf(() => fetch(`http://${new URL(closedPortUrl()).host}/`))
78+
expect(classifyDatabaseFailure(error)).toBe('permanent')
79+
})
80+
})

0 commit comments

Comments
 (0)