Skip to content

Commit 9ea67ba

Browse files
committed
fix(knowledge): disable a member sync only under the lease
A run whose binding was removed suspended members and rewrote ACLs before proving it still held the lease; a run reclaimed meanwhile could suspend the replacement's members and then fail. Suspension, the ACLs it changes, and the disable now land in one lease-guarded transaction, and a reclaimed run ends as superseded.
1 parent caabe9b commit 9ea67ba

1 file changed

Lines changed: 30 additions & 27 deletions

File tree

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

Lines changed: 30 additions & 27 deletions
Original file line numberDiff line numberDiff line change
@@ -1167,35 +1167,38 @@ async function deferMemberSync(run: MemberSyncRun, syncIntervalMinutes: number):
11671167
*/
11681168
async function disableMemberSync(run: MemberSyncRun, reason: string): Promise<void> {
11691169
const now = new Date()
1170-
const suspended = await db
1171-
.update(knowledgeConnectorMember)
1172-
.set({ status: 'suspended', suspendedAt: now, updatedAt: now })
1173-
.where(
1174-
and(
1175-
eq(knowledgeConnectorMember.connectorId, run.connectorId),
1176-
eq(knowledgeConnectorMember.status, 'active')
1170+
/** Suspension, the ACLs it changes, and the disable itself land together, and only under the lease. */
1171+
await withMemberLease(run, async (tx) => {
1172+
const suspended = await tx
1173+
.update(knowledgeConnectorMember)
1174+
.set({ status: 'suspended', suspendedAt: now, updatedAt: now })
1175+
.where(
1176+
and(
1177+
eq(knowledgeConnectorMember.connectorId, run.connectorId),
1178+
eq(knowledgeConnectorMember.status, 'active')
1179+
)
11771180
)
1178-
)
1179-
.returning({ id: knowledgeConnectorMember.id })
1180-
if (suspended.length > 0) {
1181-
const affected = await listObservedDocumentIds(
1182-
db,
1183-
suspended.map((row) => row.id)
1184-
)
1185-
await withMemberLease(run, (tx) => materializeDocumentAcls(run.connectorId, affected, tx))
1186-
}
1181+
.returning({ id: knowledgeConnectorMember.id })
1182+
if (suspended.length > 0) {
1183+
const affected = await listObservedDocumentIds(
1184+
tx,
1185+
suspended.map((row) => row.id)
1186+
)
1187+
await materializeDocumentAcls(run.connectorId, affected, tx)
1188+
}
1189+
await tx
1190+
.update(knowledgeConnector)
1191+
.set({
1192+
memberSyncStatus: 'disabled',
1193+
lastMemberSyncError: reason,
1194+
nextMemberSyncAt: null,
1195+
memberSyncLockToken: null,
1196+
memberSyncLockLeaseAt: null,
1197+
updatedAt: now,
1198+
})
1199+
.where(holdsMemberSyncLockToken(run.connectorId, run.runId))
1200+
})
11871201
await failMemberSyncLog(run.runId, run.result, reason)
1188-
await db
1189-
.update(knowledgeConnector)
1190-
.set({
1191-
memberSyncStatus: 'disabled',
1192-
lastMemberSyncError: reason,
1193-
nextMemberSyncAt: null,
1194-
memberSyncLockToken: null,
1195-
memberSyncLockLeaseAt: null,
1196-
updatedAt: now,
1197-
})
1198-
.where(holdsMemberSyncLockToken(run.connectorId, run.runId))
11991202
logger.warn('Member sync disabled', { connectorId: run.connectorId, reason })
12001203
}
12011204

0 commit comments

Comments
 (0)