Skip to content

Commit caabe9b

Browse files
committed
fix(knowledge): land member observations and ACLs only under the lease
A run that stalled past the lease TTL and resumed after a replacement took over could commit stale observations, membership rows, and document ACLs over the replacement's. Every such write now runs in a transaction that first proves the run still holds the connector's member lease, holding the connector row's lock so the scheduler cannot reclaim it mid-transaction; a run that lost it ends as superseded.
1 parent 88fbe59 commit caabe9b

2 files changed

Lines changed: 62 additions & 22 deletions

File tree

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

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -209,13 +209,14 @@ export async function rewriteConnectorAcls(
209209

210210
export async function materializeDocumentAcls(
211211
connectorId: string,
212-
documentIds: Iterable<string>
212+
documentIds: Iterable<string>,
213+
executor: DbOrTx = db
213214
): Promise<number> {
214215
const ids = [...new Set(documentIds)]
215216
let updated = 0
216217
for (let offset = 0; offset < ids.length; offset += MATERIALIZE_BATCH_SIZE) {
217218
const batch = ids.slice(offset, offset + MATERIALIZE_BATCH_SIZE)
218-
const rows = await db
219+
const rows = await executor
219220
.update(document)
220221
.set({ acl: observedAcl() })
221222
.where(

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

Lines changed: 59 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,7 @@ import {
2222
isManagedCredentialGroupBindingLive,
2323
loadCredentialGroupCredentialListContext,
2424
} from '@/lib/credential-groups/credentials'
25+
import type { DbOrTx } from '@/lib/db/types'
2526
import { isKnowledgeMemberAccessAvailable } from '@/lib/knowledge/access/availability'
2627
import { EMPTY_ACL, subjectToken } from '@/lib/knowledge/access/tokens'
2728
import {
@@ -397,6 +398,28 @@ function createMemberTokenCache(input: {
397398
}
398399
}
399400

401+
/**
402+
* Runs `fn` in a transaction that first proves this run still holds the
403+
* connector's member lease, taking the connector row's lock so the scheduler
404+
* cannot reclaim the lease mid-transaction. A run that stalled past the lease
405+
* TTL and resumed after a replacement took over therefore never lands its
406+
* observations or ACLs over the replacement's; it ends as superseded.
407+
*/
408+
async function withMemberLease<T>(
409+
run: Pick<MemberSyncRun, 'connectorId' | 'runId'>,
410+
fn: (tx: DbOrTx) => Promise<T>
411+
): Promise<T> {
412+
return db.transaction(async (tx) => {
413+
const [held] = await tx
414+
.select({ id: knowledgeConnector.id })
415+
.from(knowledgeConnector)
416+
.where(stillHoldsMemberSyncLock(run.connectorId, run.runId))
417+
.for('update')
418+
if (!held) throw new SyncLockLostException(run.connectorId)
419+
return fn(tx)
420+
})
421+
}
422+
400423
async function acquireMemberSyncLock(
401424
connectorId: string,
402425
runId: string,
@@ -525,6 +548,10 @@ async function reconcileMembership(
525548

526549
const inserts: (typeof knowledgeConnectorMember.$inferInsert)[] = []
527550
const affectedMemberIds: string[] = []
551+
const updates: Array<{
552+
id: string
553+
values: Partial<typeof knowledgeConnectorMember.$inferInsert>
554+
}> = []
528555
const deleteMemberIds: string[] = []
529556

530557
for (const snapshot of snapshots.values()) {
@@ -558,17 +585,17 @@ async function reconcileMembership(
558585
const tokenChanged = row.subjectToken !== snapshot.subjectToken
559586
const statusChanged = row.status !== status
560587
if (!tokenChanged && !statusChanged) continue
561-
await db
562-
.update(knowledgeConnectorMember)
563-
.set({
588+
updates.push({
589+
id: row.id,
590+
values: {
564591
subjectToken: snapshot.subjectToken,
565592
status,
566593
suspendedAt: snapshot.active ? null : (row.suspendedAt ?? now),
567594
/** A reactivated member is due immediately; their observations may be stale. */
568595
...(statusChanged && snapshot.active ? { nextAttemptAt: now, consecutiveFailures: 0 } : {}),
569596
updatedAt: now,
570-
})
571-
.where(eq(knowledgeConnectorMember.id, row.id))
597+
},
598+
})
572599
affectedMemberIds.push(row.id)
573600
}
574601

@@ -585,18 +612,28 @@ async function reconcileMembership(
585612
affectedDocumentIds.add(documentId)
586613
}
587614
}
588-
if (deleteMemberIds.length > 0) {
589-
await db
590-
.delete(knowledgeConnectorMember)
591-
.where(
592-
and(
593-
eq(knowledgeConnectorMember.connectorId, run.connectorId),
594-
inArray(knowledgeConnectorMember.id, deleteMemberIds)
595-
)
596-
)
597-
}
598-
if (inserts.length > 0) {
599-
await db.insert(knowledgeConnectorMember).values(inserts).onConflictDoNothing()
615+
if (updates.length > 0 || deleteMemberIds.length > 0 || inserts.length > 0) {
616+
await withMemberLease(run, async (tx) => {
617+
for (const update of updates) {
618+
await tx
619+
.update(knowledgeConnectorMember)
620+
.set(update.values)
621+
.where(eq(knowledgeConnectorMember.id, update.id))
622+
}
623+
if (deleteMemberIds.length > 0) {
624+
await tx
625+
.delete(knowledgeConnectorMember)
626+
.where(
627+
and(
628+
eq(knowledgeConnectorMember.connectorId, run.connectorId),
629+
inArray(knowledgeConnectorMember.id, deleteMemberIds)
630+
)
631+
)
632+
}
633+
if (inserts.length > 0) {
634+
await tx.insert(knowledgeConnectorMember).values(inserts).onConflictDoNothing()
635+
}
636+
})
600637
}
601638

602639
logger.info('Reconciled members-mode membership', {
@@ -895,7 +932,7 @@ async function applyMemberListing(
895932
const exhaustedFailures = (outcome.member.consecutiveFailures ?? 0) + 1
896933
const now = new Date()
897934

898-
await db.transaction(async (tx) => {
935+
await withMemberLease(run, async (tx) => {
899936
const added = await recordMemberObservations(tx, outcome.member.id, seenDocumentIds, run.runId)
900937
run.result.observationsAdded += added
901938
/**
@@ -1145,7 +1182,7 @@ async function disableMemberSync(run: MemberSyncRun, reason: string): Promise<vo
11451182
db,
11461183
suspended.map((row) => row.id)
11471184
)
1148-
await materializeDocumentAcls(run.connectorId, affected)
1185+
await withMemberLease(run, (tx) => materializeDocumentAcls(run.connectorId, affected, tx))
11491186
}
11501187
await failMemberSyncLog(run.runId, run.result, reason)
11511188
await db
@@ -1472,7 +1509,9 @@ export async function executeMemberSync(
14721509
for (const documentId of affected) affectedDocumentIds.add(documentId)
14731510
}
14741511

1475-
await materializeDocumentAcls(connectorId, affectedDocumentIds)
1512+
await withMemberLease(run, (tx) =>
1513+
materializeDocumentAcls(connectorId, affectedDocumentIds, tx)
1514+
)
14761515

14771516
/**
14781517
* Nobody has completed a listing yet — a connector that just entered

0 commit comments

Comments
 (0)