Skip to content

Commit 0a7dd0c

Browse files
committed
fix(knowledge): isolate stale sweep lock failures, close failed disables, update ACL test callers
1 parent bb4623a commit 0a7dd0c

9 files changed

Lines changed: 266 additions & 44 deletions

File tree

‎apps/sim/app/api/knowledge/connectors/member-sync/route.test.ts‎

Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -101,6 +101,21 @@ describe('member sync scheduler owner routing', () => {
101101
).toBe(true)
102102
})
103103

104+
it('still dispatches due connectors when the stale observation sweep fails', async () => {
105+
mocks.sweep.mockRejectedValue(
106+
Object.assign(new Error('canceling statement due to lock timeout'), { code: '55P03' })
107+
)
108+
queueTableRows(schemaMock.knowledgeConnector, [
109+
{ id: 'workspace-source', workspaceId: 'workspace-a', organizationId: null },
110+
])
111+
const response = await GET(createMockRequest('GET'))
112+
expect(response.status).toBe(200)
113+
expect(mocks.dispatch).toHaveBeenCalledExactlyOnceWith(
114+
'workspace-source',
115+
expect.objectContaining({ requireRunnable: true })
116+
)
117+
})
118+
104119
it('preserves workspace dispatch and refuses absent or ambiguous ownership', async () => {
105120
queueTableRows(schemaMock.knowledgeConnector, [
106121
{ id: 'missing', workspaceId: null, organizationId: null },

‎apps/sim/app/api/knowledge/connectors/member-sync/route.ts‎

Lines changed: 11 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
import { db } from '@sim/db'
22
import { knowledgeBase, knowledgeConnector, knowledgeConnectorMemberSyncLog } from '@sim/db/schema'
33
import { createLogger } from '@sim/logger'
4+
import { getErrorMessage } from '@sim/utils/errors'
45
import { and, asc, eq, inArray, isNull, lte, type SQL, sql } from 'drizzle-orm'
56
import { type NextRequest, NextResponse } from 'next/server'
67
import { verifyCronAuth } from '@/lib/auth/internal'
@@ -157,9 +158,16 @@ export const GET = withRouteHandler(async (request: NextRequest) => {
157158
logger.warn(`[${requestId}] Closed ${closedLogs.length} orphaned member sync log(s)`)
158159
}
159160

160-
const sweep = await sweepStaleMemberObservations(now)
161-
if (sweep.members > 0) {
162-
logger.warn(`[${requestId}] Swept observations of ${sweep.members} stale member(s)`, sweep)
161+
/** Observation hygiene never holds back dispatch; an unfinished sweep resumes next tick. */
162+
try {
163+
const sweep = await sweepStaleMemberObservations(now)
164+
if (sweep.members > 0) {
165+
logger.warn(`[${requestId}] Swept observations of ${sweep.members} stale member(s)`, sweep)
166+
}
167+
} catch (error) {
168+
logger.error(`[${requestId}] Stale member observation sweep failed`, {
169+
error: getErrorMessage(error),
170+
})
163171
}
164172

165173
const dueConnectors = await db

‎apps/sim/lib/knowledge/__integration__/connector-lease-pages.integration.ts‎

Lines changed: 120 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -9,6 +9,8 @@ import {
99
document,
1010
knowledgeConnector,
1111
knowledgeConnectorMember,
12+
knowledgeConnectorMemberSyncLog,
13+
knowledgeDocumentObservation,
1214
organization,
1315
resourcePolicy,
1416
user,
@@ -47,6 +49,7 @@ import {
4749
materializeDocumentAcls,
4850
recordMemberObservations,
4951
rewriteConnectorAcls,
52+
sweepStaleMemberObservations,
5053
} from '@/lib/knowledge/connectors/member-observations'
5154
import {
5255
executeMemberSync,
@@ -395,6 +398,83 @@ describe('connector lease ACL pages in PostgreSQL', () => {
395398
})
396399
})
397400

401+
describe('stale member sweep', () => {
402+
it('defers a connector whose row a member run holds and still sweeps the others', async () => {
403+
const busy = members
404+
const idle = await seedKnowledgeMemberFixture(ids)
405+
for (const fixture of [busy, idle]) {
406+
await db
407+
.update(knowledgeConnector)
408+
.set({
409+
status: 'active',
410+
memberSyncStatus: 'idle',
411+
memberSyncLockToken: null,
412+
syncIntervalMinutes: 60,
413+
lastMemberSyncAt: new Date(),
414+
})
415+
.where(eq(knowledgeConnector.id, fixture.connectorId))
416+
/** Enrolled long enough ago that never having listed makes them stale. */
417+
await db
418+
.update(knowledgeConnectorMember)
419+
.set({ createdAt: new Date(Date.now() - 3 * 24 * 60 * 60 * 1000) })
420+
.where(eq(knowledgeConnectorMember.connectorId, fixture.connectorId))
421+
const seeded = await seedDocuments(
422+
fixture.connectorId,
423+
fixture.members.map((member) => member.subjectToken).sort(),
424+
3
425+
)
426+
for (const member of fixture.members)
427+
await recordMemberObservations(
428+
db,
429+
member.id,
430+
seeded.map((row) => row.id),
431+
fixture.runId
432+
)
433+
}
434+
const observed = async (connectorId: string) =>
435+
(
436+
await db
437+
.select({ id: knowledgeDocumentObservation.memberId })
438+
.from(knowledgeDocumentObservation)
439+
.innerJoin(
440+
knowledgeConnectorMember,
441+
eq(knowledgeConnectorMember.id, knowledgeDocumentObservation.memberId)
442+
)
443+
.where(eq(knowledgeConnectorMember.connectorId, connectorId))
444+
).length
445+
446+
/** A member page of a running sync holds the busy connector's row for longer than a sweep page waits. */
447+
let release!: () => void
448+
let locked!: () => void
449+
const held = new Promise<void>((resolve) => {
450+
locked = resolve
451+
})
452+
const holder = db.transaction(async (tx) => {
453+
await tx
454+
.select({ id: knowledgeConnector.id })
455+
.from(knowledgeConnector)
456+
.where(eq(knowledgeConnector.id, busy.connectorId))
457+
.for('update')
458+
locked()
459+
await new Promise<void>((resolve) => {
460+
release = resolve
461+
})
462+
})
463+
await held
464+
try {
465+
const result = await sweepStaleMemberObservations(new Date(Date.now() + 1_000))
466+
expect(result.members).toBe(idle.members.length)
467+
} finally {
468+
release()
469+
await holder
470+
}
471+
472+
expect(await observed(busy.connectorId)).toBe(busy.members.length * 3)
473+
expect(await observed(idle.connectorId)).toBe(0)
474+
expect((await storedAcls(idle.connectorId)).every((acl) => acl.length === 0)).toBe(true)
475+
}, 30_000)
476+
})
477+
398478
describe('resumeMembershipRewrites', () => {
399479
it('rematerialises a changed member one page per lease transaction', async () => {
400480
const seeded = await seedDocuments(members.connectorId, [])
@@ -649,5 +729,45 @@ describe('connector lease ACL pages in PostgreSQL', () => {
649729
)
650730
expect(granted).toBe(0)
651731
})
732+
733+
it('closes a run whose disable of a removed option fails as an ordinary failure', async () => {
734+
/** The option id is set but no longer exists, which membership reconciliation reports. */
735+
await db
736+
.update(knowledgeConnector)
737+
.set({ credentialGroupOptionId: generateId() })
738+
.where(eq(knowledgeConnector.id, members.connectorId))
739+
await seedDocuments(
740+
members.connectorId,
741+
members.members.map((member) => member.subjectToken).sort()
742+
)
743+
await db.execute(
744+
sql.raw(`CREATE OR REPLACE FUNCTION fail_after_acl_writes() RETURNS trigger
745+
LANGUAGE plpgsql AS $$ BEGIN RAISE EXCEPTION 'fixture statement failure'; END $$`)
746+
)
747+
await db.execute(
748+
sql`CREATE TRIGGER fail_after_acl_writes BEFORE UPDATE OF acl ON document FOR EACH ROW
749+
WHEN (NEW.connector_id = ${sql.raw(`'${members.connectorId}'`)})
750+
EXECUTE FUNCTION fail_after_acl_writes()`
751+
)
752+
753+
const failed = await executeMemberSync(members.connectorId, {
754+
billingAttribution: await billing(),
755+
})
756+
757+
expect(failed.error).toBeTruthy()
758+
const state = await connectorState()
759+
expect(state?.memberSyncStatus).toBe('error')
760+
expect(state?.memberSyncLockToken).toBeNull()
761+
const [log] = await db
762+
.select({ status: knowledgeConnectorMemberSyncLog.status })
763+
.from(knowledgeConnectorMemberSyncLog)
764+
.where(eq(knowledgeConnectorMemberSyncLog.connectorId, members.connectorId))
765+
expect(log?.status).toBe('failed')
766+
767+
await db.execute(sql`DROP TRIGGER fail_after_acl_writes ON document`)
768+
await executeMemberSync(members.connectorId, { billingAttribution: await billing() })
769+
expect((await connectorState())?.memberSyncStatus).toBe('disabled')
770+
expect((await storedAcls(members.connectorId)).every((acl) => acl.length === 0)).toBe(true)
771+
})
652772
})
653773
})

‎apps/sim/lib/knowledge/access/predicate.postgres.test.ts‎

Lines changed: 23 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,7 @@ vi.mock('@/connectors/registry.server', () => ({ CONNECTOR_REGISTRY: {} }))
2222
const { drizzle } = await import('drizzle-orm/postgres-js')
2323
const schema = await import('@sim/db/schema')
2424
const { persistDocumentAcls } = await import('@/lib/knowledge/connectors/sync-persistence')
25+
const { leaseTransaction } = await import('@/lib/knowledge/connectors/sync-lock')
2526
const { mergeMirroredAcls, hideUnlistedDocuments } = await import(
2627
'@/lib/knowledge/connectors/mirrored-acls'
2728
)
@@ -750,7 +751,9 @@ describe.runIf(Boolean(databaseUrl))('knowledge ACLs in PostgreSQL', () => {
750751
const space = 'g:confluence:tenant:space'
751752
const page = 'g:confluence:tenant:page'
752753
const input = new Map([['page', { acl: [space], requirements: [[page]] }]])
753-
expect(await persistDocumentAcls('admin', input, executor)).toEqual({ updated: 1, rejected: 0 })
754+
expect(
755+
await persistDocumentAcls('admin', input, leaseTransaction('admin', undefined, executor))
756+
).toEqual({ updated: 1, rejected: 0 })
754757
expect(await readable([space, page], 'persisted')).toBe(true)
755758
expect(await readable([page], 'persisted')).toBe(false)
756759
await connection.unsafe("UPDATE document SET acl = string_to_array($1, E'\\n')", [page])
@@ -760,7 +763,7 @@ describe.runIf(Boolean(databaseUrl))('knowledge ACLs in PostgreSQL', () => {
760763
"UPDATE document SET acl_verified_at = statement_timestamp() - interval '25 hours'"
761764
)
762765
expect(await readable([space, page], 'persisted')).toBe(false)
763-
await persistDocumentAcls('admin', input, executor)
766+
await persistDocumentAcls('admin', input, leaseTransaction('admin', undefined, executor))
764767
expect(await readable([space, page], 'persisted')).toBe(true)
765768
const [stored] = await connection.unsafe(
766769
"SELECT jsonb_typeof(acl_requirements) AS shape, acl_requirements FROM document WHERE id = 'persisted'"
@@ -844,10 +847,15 @@ describe.runIf(Boolean(databaseUrl))('knowledge ACLs in PostgreSQL', () => {
844847
{}
845848
)
846849
if (step === 'unlisted') hideUnlistedDocuments(merged.acls, ['shared-file'])
847-
await persistDocumentAcls('admin', merged.acls, executor, {
848-
unresolvedExternalIds: merged.unresolvedExternalIds,
849-
generationStartedAt,
850-
})
850+
await persistDocumentAcls(
851+
'admin',
852+
merged.acls,
853+
leaseTransaction('admin', undefined, executor),
854+
{
855+
unresolvedExternalIds: merged.unresolvedExternalIds,
856+
generationStartedAt,
857+
}
858+
)
851859
const [stored] = await connection.unsafe(
852860
"SELECT to_jsonb(acl) AS acl, acl_verified_at FROM document WHERE id = 'shared'"
853861
)
@@ -882,10 +890,15 @@ describe.runIf(Boolean(databaseUrl))('knowledge ACLs in PostgreSQL', () => {
882890
[verifiedAt]
883891
)
884892
const executor = drizzle(connection, { schema })
885-
const result = await persistDocumentAcls('admin', new Map([['file', []]]), executor, {
886-
unresolvedExternalIds: new Set(['file']),
887-
generationStartedAt,
888-
})
893+
const result = await persistDocumentAcls(
894+
'admin',
895+
new Map([['file', []]]),
896+
leaseTransaction('admin', undefined, executor),
897+
{
898+
unresolvedExternalIds: new Set(['file']),
899+
generationStartedAt,
900+
}
901+
)
889902
expect(result.updated).toBe(preserved ? 0 : 1)
890903
const [stored] = await connection.unsafe(
891904
"SELECT to_jsonb(acl) AS acl, acl_verified_at FROM document WHERE id = 'boundary'"

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

Lines changed: 36 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -203,12 +203,44 @@ describe('sweepStaleMemberObservations', () => {
203203

204204
expect(dbChainMockFns.transaction).toHaveBeenCalledTimes(2)
205205
expect(dbChainMockFns.limit).toHaveBeenCalledWith(25)
206-
const bounds = dbChainMockFns.execute.mock.calls.filter(([query]) =>
207-
JSON.stringify(query).includes('lock_timeout')
206+
const bounds = dbChainMockFns.execute.mock.calls.filter((call: unknown[]) =>
207+
JSON.stringify(call).includes('lock_timeout')
208208
)
209209
expect(bounds).toHaveLength(2)
210210
})
211211

212+
/** A connector whose running member page holds its row must not stall the other connectors. */
213+
it('defers a member whose page hits a lock timeout and still sweeps the next one', async () => {
214+
queueTableRows(schemaMock.knowledgeConnectorMember, [
215+
STALE_MEMBER,
216+
{ ...STALE_MEMBER, id: 'm-2', connectorId: 'c-2' },
217+
])
218+
dbChainMockFns.transaction.mockRejectedValueOnce(
219+
Object.assign(new Error('canceling statement due to lock timeout'), { code: '55P03' })
220+
)
221+
queueTableRows(schemaMock.knowledgeConnector, [{ id: 'c-2' }])
222+
queueTableRows(schemaMock.knowledgeConnectorMember, [{ id: 'm-2' }])
223+
dbChainMockFns.returning
224+
.mockResolvedValueOnce([{ documentId: 'd-1' }])
225+
.mockResolvedValueOnce([{ id: 'd-1' }])
226+
.mockResolvedValueOnce([])
227+
228+
await expect(sweepStaleMemberObservations(NOW)).resolves.toEqual({
229+
members: 1,
230+
observationsRemoved: 1,
231+
documentsRematerialized: 1,
232+
docsTombstoned: 0,
233+
})
234+
expect(dbChainMockFns.transaction).toHaveBeenCalledTimes(2)
235+
})
236+
237+
it('fails the sweep on an error that is not a lock or capacity failure', async () => {
238+
queueTableRows(schemaMock.knowledgeConnectorMember, [STALE_MEMBER])
239+
dbChainMockFns.transaction.mockRejectedValueOnce(new TypeError('broken'))
240+
241+
await expect(sweepStaleMemberObservations(NOW)).rejects.toThrow('broken')
242+
})
243+
212244
/**
213245
* A run that claimed the member between the selection and the lock moved
214246
* `lastStartedAt` forward, so the re-check under `FOR UPDATE` finds nothing
@@ -336,8 +368,8 @@ describe('rewriteConnectorAcls', () => {
336368
expect(assignedPages()).toHaveLength(3)
337369
expect(pageSizes()).toEqual([25, 25, 10])
338370
expect(dbChainMockFns.transaction).toHaveBeenCalledTimes(3)
339-
const bounds = dbChainMockFns.execute.mock.calls.filter(([query]) =>
340-
JSON.stringify(query).includes('lock_timeout')
371+
const bounds = dbChainMockFns.execute.mock.calls.filter((call: unknown[]) =>
372+
JSON.stringify(call).includes('lock_timeout')
341373
)
342374
expect(bounds).toHaveLength(3)
343375
})

‎apps/sim/lib/knowledge/connectors/member-observations.ts‎

Lines changed: 29 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,8 @@ import {
55
knowledgeConnectorMember,
66
knowledgeDocumentObservation,
77
} from '@sim/db/schema'
8+
import { createLogger } from '@sim/logger'
9+
import { getPostgresErrorCode } from '@sim/utils/errors'
810
import { chunkArray } from '@sim/utils/helpers'
911
import {
1012
and,
@@ -49,6 +51,8 @@ import {
4951
hardDeleteDocuments,
5052
} from '@/lib/knowledge/documents/service'
5153

54+
const logger = createLogger('MemberObservations')
55+
5256
/**
5357
* Documents per tombstone or resurrection page. Those writes set `deleted_at` alone and fire no
5458
* projection fan-out, so they page wider than an ACL write.
@@ -870,6 +874,9 @@ function memberStillStale(memberId: string, cutoff: Date) {
870874
)
871875
}
872876

877+
/** Lock, statement-timeout and serialization failures a sweep page defers instead of failing the tick. */
878+
const SWEEP_DEFERRABLE_CODES = new Set(['55P03', '57014', '40P01', '40001'])
879+
873880
/**
874881
* Stale-member observations one sweep tick removes per member: the same
875882
* {@link OBSERVATION_BATCH_SIZE} as before, now in pages of
@@ -1044,14 +1051,28 @@ export async function sweepStaleMemberObservations(now: Date): Promise<StaleMemb
10441051
for (const member of staleMembers) {
10451052
const memberCutoff = new Date(now.getTime() - staleMemberWindowMs(member.syncIntervalMinutes))
10461053
let sweptAny = false
1047-
for (let page = 0; page < STALE_MEMBER_PAGES_PER_TICK; page++) {
1048-
const swept = await sweepStaleMemberPage(member, memberCutoff, now)
1049-
if (!swept) break
1050-
sweptAny = true
1051-
result.observationsRemoved += swept.observationsRemoved
1052-
result.documentsRematerialized += swept.documentsRematerialized
1053-
result.docsTombstoned += swept.docsTombstoned
1054-
if (swept.observationsRemoved < ACL_CHANGE_BATCH_SIZE) break
1054+
try {
1055+
for (let page = 0; page < STALE_MEMBER_PAGES_PER_TICK; page++) {
1056+
const swept = await sweepStaleMemberPage(member, memberCutoff, now)
1057+
if (!swept) break
1058+
sweptAny = true
1059+
result.observationsRemoved += swept.observationsRemoved
1060+
result.documentsRematerialized += swept.documentsRematerialized
1061+
result.docsTombstoned += swept.docsTombstoned
1062+
if (swept.observationsRemoved < ACL_CHANGE_BATCH_SIZE) break
1063+
}
1064+
} catch (error) {
1065+
/**
1066+
* A connector whose member run holds its row, or whose page outruns the bounds, is left
1067+
* for the next tick: its committed pages stand, and the other members are still swept.
1068+
*/
1069+
const code = getPostgresErrorCode(error)
1070+
if (!code || !SWEEP_DEFERRABLE_CODES.has(code)) throw error
1071+
logger.warn('Deferred a stale member sweep to the next tick', {
1072+
connectorId: member.connectorId,
1073+
memberId: member.id,
1074+
code,
1075+
})
10551076
}
10561077
if (sweptAny) result.members += 1
10571078
}

0 commit comments

Comments
 (0)