Skip to content

Commit 624793f

Browse files
authored
fix(knowledge): bound every connector-lease ACL transaction to one short page (#8195)
* fix(knowledge): bound every connector-lease ACL transaction to one short page * fix(knowledge): isolate stale sweep lock failures, close failed disables, update ACL test callers * fix(knowledge): prove connector leases last, page ACL writes by projection rows, and skip unchanged work * fix(knowledge): row-bound every observation page, size ACL pages from locked chunk counts, and check the sweep budget per page * fix(knowledge): stop every ACL page walker at the run budget and keep pending rewrites until they finish * test(knowledge): seed filled projection rows so the fan-out test holds on any provisioned database * fix(knowledge): check the page deadline after the lease heartbeat * test(knowledge): expire the restore budget at the heartbeat before the next window
1 parent 7bbe3f5 commit 624793f

21 files changed

Lines changed: 3258 additions & 493 deletions

.github/workflows/test-build.yml

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -277,6 +277,7 @@ jobs:
277277
lib/knowledge/__integration__/listing-continuation.integration.ts
278278
lib/knowledge/__integration__/member-scope-renewal.integration.ts
279279
lib/knowledge/__integration__/member-document-lifecycle.integration.ts
280+
lib/knowledge/__integration__/connector-lease-pages.integration.ts
280281
lib/knowledge/__integration__/slack-empty-threads.integration.ts
281282
lib/knowledge/__integration__/kb-block-search.integration.ts
282283
lib/knowledge/__integration__/gitlab-workspace.integration.ts

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

0 commit comments

Comments
 (0)