Skip to content

Commit fcd27da

Browse files
committed
fix(knowledge): walk the members-mode reconcile by external id and tombstone explicit removals without a completed listing
The absence reconcile ordered its resumable walk by source_seen_at, which every member listing rewrites for what it observed. On a connector larger than one run's page budget the walk chased re-stamped documents and a pass never ended, so a document that lost its observers behind the cursor was never revisited. The walk now keys on external id through doc_connector_external_id_idx, which a document never changes, so a pass completes within ceil(documents / budget) runs. Documents a run explicitly unobserved (a complete listing, a change-feed withdrawal, or a member removal) are now tombstoned even when no member with a completed listing remains, as the stale-member sweep already does. Removing the only listed member previously left every document only it observed live and unobserved indefinitely. The reconcile and the purge stay gated on a completed listing.
1 parent eb6bd27 commit fcd27da

5 files changed

Lines changed: 176 additions & 62 deletions

File tree

‎apps/sim/lib/knowledge/__integration__/member-document-lifecycle.integration.ts‎

Lines changed: 102 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,7 @@ import { db } from '@sim/db'
33
import {
44
document,
55
knowledgeConnector,
6+
knowledgeConnectorMember,
67
knowledgeDocumentObservation,
78
organization,
89
user,
@@ -21,9 +22,11 @@ import {
2122
recordMemberObservations,
2223
removeMemberObservationsForDocuments,
2324
} from '@/lib/knowledge/connectors/member-observations'
25+
import { resumeMembershipRewrites } from '@/lib/knowledge/connectors/member-sync-engine'
2426
import { MEMBER_TOMBSTONE_RECONCILE_PAGES_PER_RUN } from '@/lib/knowledge/connectors/sync-limits'
2527
import {
2628
assertSyncLeaseHeldInTx,
29+
createMemberSyncLease,
2730
SyncLockLostException,
2831
stillHoldsMemberSyncLock,
2932
} from '@/lib/knowledge/connectors/sync-lock'
@@ -143,10 +146,10 @@ describe('member document lifecycle in PostgreSQL', () => {
143146
const unobserved = Array.from({ length: pageBudget + 20 }, (_, index) =>
144147
row(`unobserved-${index}`)
145148
)
146-
const lastSeen = sql`'2026-06-01 00:00:00'::timestamp`
149+
/** Both sort after every unobserved document, beyond what this run's reconcile reaches. */
147150
const [lostByThisRun, stillObservedByBob] = [
148-
{ ...row('lost-by-this-run'), sourceSeenAt: lastSeen },
149-
{ ...row('still-observed-by-bob'), sourceSeenAt: lastSeen },
151+
row('zz-lost-by-this-run'),
152+
row('zz-still-observed-by-bob'),
150153
]
151154
await insertRows([...unobserved, lostByThisRun, stillObservedByBob])
152155
await observe([lostByThisRun.id, stillObservedByBob.id])
@@ -172,7 +175,7 @@ describe('member document lifecycle in PostgreSQL', () => {
172175
expect(afterFirst.has(lostByThisRun.id)).toBe(true)
173176
expect(afterFirst.has(stillObservedByBob.id)).toBe(false)
174177
expect(unobserved.filter(({ id }) => !afterFirst.has(id))).toHaveLength(20)
175-
expect(await savedCursor()).toEqual({ seenAt: expect.any(String), id: expect.any(String) })
178+
expect(await savedCursor()).toEqual({ externalId: expect.any(String) })
176179

177180
expect(await run()).toEqual({ tombstoned: 20, resurrected: 0, purged: 0, finished: true })
178181
const afterSecond = await tombstonedIds()
@@ -181,6 +184,101 @@ describe('member document lifecycle in PostgreSQL', () => {
181184
expect(await savedCursor()).toBeNull()
182185
})
183186

187+
it('tombstones what removing the only listed member unobserved, though no completed listing remains', async () => {
188+
const [removedMember, otherMember] = members.members
189+
await db
190+
.update(knowledgeConnectorMember)
191+
.set({
192+
lastCompleteListingAt: new Date(),
193+
listingCheckpoint: { kind: 'membership', cursor: null, removeMember: true },
194+
})
195+
.where(eq(knowledgeConnectorMember.id, removedMember.id))
196+
const onlyRemoved = row('only-the-removed-member')
197+
const sharedWithOther = row('shared-with-other-member')
198+
const neverObserved = row('never-observed')
199+
await insertRows([onlyRemoved, sharedWithOther, neverObserved])
200+
await observe([onlyRemoved.id, sharedWithOther.id])
201+
await recordMemberObservations(db, otherMember.id, [sharedWithOther.id], members.runId)
202+
203+
const unobservedDocumentIds = new Set<string>()
204+
expect(
205+
await resumeMembershipRewrites({
206+
connectorId: members.connectorId,
207+
runId: members.runId,
208+
deadlineAt: Date.now() + 60_000,
209+
lease: createMemberSyncLease(members.connectorId, members.runId),
210+
unobservedDocumentIds,
211+
})
212+
).toBe(true)
213+
const [completed] = await db
214+
.select({ count: sql<number>`count(*)::int` })
215+
.from(knowledgeConnectorMember)
216+
.where(
217+
and(
218+
eq(knowledgeConnectorMember.connectorId, members.connectorId),
219+
isNotNull(knowledgeConnectorMember.lastCompleteListingAt)
220+
)
221+
)
222+
expect(completed.count).toBe(0)
223+
224+
expect(
225+
await run({
226+
allowRemoval: completed.count > 0,
227+
unobservedDocumentIds: [...unobservedDocumentIds],
228+
})
229+
).toEqual({ tombstoned: 1, resurrected: 0, purged: 0, finished: true })
230+
const tombstoned = await tombstonedIds()
231+
expect(tombstoned.has(onlyRemoved.id)).toBe(true)
232+
expect(tombstoned.has(sharedWithOther.id)).toBe(false)
233+
expect(tombstoned.has(neverObserved.id)).toBe(false)
234+
})
235+
236+
it('finishes a pass within its page budget while listings re-stamp every observed document', async () => {
237+
const pageBudget = MEMBER_TOMBSTONE_RECONCILE_PAGES_PER_RUN * 500
238+
const total = pageBudget + 500
239+
const runsPerPass = Math.ceil(total / pageBudget)
240+
const firstInEveryOrder = {
241+
...row('walk-0000000'),
242+
id: '00000000-0000-4000-8000-000000000000',
243+
}
244+
const rest = Array.from({ length: total - 1 }, (_, index) =>
245+
row(`walk-${String(index + 1).padStart(7, '0')}`)
246+
)
247+
await insertRows([firstInEveryOrder, ...rest])
248+
const all = [firstInEveryOrder, ...rest].map(({ id }) => id)
249+
for (let offset = 0; offset < all.length; offset += 5000)
250+
await observe(all.slice(offset, offset + 5000))
251+
/** A listing stamps what it saw with its start; the document nobody observes keeps its old stamp. */
252+
const relist = () =>
253+
db
254+
.update(document)
255+
.set({ sourceSeenAt: new Date() })
256+
.where(
257+
and(
258+
eq(document.connectorId, members.connectorId),
259+
sql`EXISTS (SELECT 1 FROM knowledge_document_observation o WHERE o.document_id = ${document.id})`
260+
)
261+
)
262+
263+
expect(await run()).toMatchObject({ tombstoned: 0 })
264+
expect(await savedCursor()).not.toBeNull()
265+
await db
266+
.delete(knowledgeDocumentObservation)
267+
.where(eq(knowledgeDocumentObservation.documentId, firstInEveryOrder.id))
268+
for (let pass = 1; pass < runsPerPass; pass++) {
269+
await relist()
270+
await run()
271+
}
272+
expect(await savedCursor()).toBeNull()
273+
expect((await tombstonedIds()).has(firstInEveryOrder.id)).toBe(false)
274+
275+
for (let next = 0; next < runsPerPass; next++) {
276+
await relist()
277+
await run()
278+
}
279+
expect(await tombstonedIds()).toEqual(new Set([firstInEveryOrder.id]))
280+
})
281+
184282
it('continues past a full selected batch even if its observations change before UPDATE', async () => {
185283
const rows = Array.from({ length: 501 }, (_, index) => row(String(index)))
186284
await db.insert(document).values(rows)

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

Lines changed: 33 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -261,7 +261,7 @@ describe('applyMemberDocumentLifecycle', () => {
261261
withLease: (fn) => fn(db),
262262
...overrides,
263263
})
264-
const pageRow = (id: string, live = true) => ({ id, seenAt: '2026-01-01 00:00:00', live })
264+
const pageRow = (id: string) => ({ id, externalId: `ext-${id}` })
265265
/** The WHERE of every statement that targets documents by id, in call order. */
266266
const documentUpdateConditions = () =>
267267
dbChainMockFns.where.mock.calls
@@ -288,11 +288,7 @@ describe('applyMemberDocumentLifecycle', () => {
288288
})
289289

290290
it('checks observations only for the live documents of one bounded page', async () => {
291-
queueTableRows(schemaMock.document, [
292-
pageRow('observed'),
293-
pageRow('unobserved'),
294-
pageRow('already-tombstoned', false),
295-
])
291+
queueTableRows(schemaMock.document, [pageRow('observed'), pageRow('unobserved')])
296292
dbChainMockFns.returning.mockResolvedValueOnce([{ id: 'unobserved' }])
297293

298294
await expect(applyMemberDocumentLifecycle(lifecycleInput())).resolves.toEqual({
@@ -314,6 +310,17 @@ describe('applyMemberDocumentLifecycle', () => {
314310
expect(
315311
hasMockCondition(pageCondition, (node) => node.type === 'notExists' || node.type === 'exists')
316312
).toBe(false)
313+
expect(
314+
hasMockCondition(
315+
pageCondition,
316+
(node) => node.type === 'isNull' && node.column === schemaMock.document.deletedAt
317+
)
318+
).toBe(true)
319+
/** An immutable key: `source_seen_at` moves on every listing, so a walk ordered by it never ends. */
320+
expect(dbChainMockFns.orderBy).toHaveBeenCalledWith({
321+
type: 'asc',
322+
column: schemaMock.document.externalId,
323+
})
317324
expect(dbChainMockFns.limit).toHaveBeenCalledWith(500)
318325
const [tombstone] = documentUpdateConditions()
319326
expect(updatedIds(tombstone)).toEqual(['observed', 'unobserved'])
@@ -322,9 +329,7 @@ describe('applyMemberDocumentLifecycle', () => {
322329
})
323330

324331
it('stops after its page budget, saving where the next run resumes', async () => {
325-
queueTableRows(schemaMock.knowledgeConnector, [
326-
{ cursor: { seenAt: '2025-12-31 00:00:00', id: 'previous-run' } },
327-
])
332+
queueTableRows(schemaMock.knowledgeConnector, [{ cursor: { externalId: 'previous-run' } }])
328333
for (let page = 0; page < MEMBER_TOMBSTONE_RECONCILE_PAGES_PER_RUN; page++) {
329334
queueTableRows(
330335
schemaMock.document,
@@ -344,15 +349,18 @@ describe('applyMemberDocumentLifecycle', () => {
344349
const firstPage = dbChainMockFns.where.mock.calls
345350
.map(([condition]) => flattenMockConditions(condition))
346351
.find((nodes) => nodes.some((node) => node.left === schemaMock.document.connectorId))
347-
expect(JSON.stringify(firstPage)).toContain('previous-run')
352+
expect(firstPage).toContainEqual({
353+
type: 'gt',
354+
left: schemaMock.document.externalId,
355+
right: 'previous-run',
356+
})
348357
const cursors = dbChainMockFns.set.mock.calls
349358
.map(([value]) => value)
350359
.filter((value) => 'memberTombstoneCursor' in value)
351360
expect(cursors).toHaveLength(MEMBER_TOMBSTONE_RECONCILE_PAGES_PER_RUN)
352361
expect(cursors.at(-1)).toEqual({
353362
memberTombstoneCursor: {
354-
seenAt: '2026-01-01 00:00:00',
355-
id: `p${MEMBER_TOMBSTONE_RECONCILE_PAGES_PER_RUN - 1}-499`,
363+
externalId: `ext-p${MEMBER_TOMBSTONE_RECONCILE_PAGES_PER_RUN - 1}-499`,
356364
},
357365
})
358366
expect(documentUpdateConditions()).toHaveLength(MEMBER_TOMBSTONE_RECONCILE_PAGES_PER_RUN)
@@ -380,12 +388,20 @@ describe('applyMemberDocumentLifecycle', () => {
380388
})
381389
})
382390

383-
it('removes nothing until a member has completed a listing', async () => {
391+
it('before any completed listing, tombstones only what this run explicitly unobserved', async () => {
384392
queueTableRows(schemaMock.document, [])
385-
await applyMemberDocumentLifecycle(
386-
lifecycleInput({ allowRemoval: false, unobservedDocumentIds: ['d-1'] })
387-
)
388-
expect(dbChainMockFns.update).not.toHaveBeenCalled()
393+
dbChainMockFns.returning.mockResolvedValueOnce([{ id: 'd-removed-member' }])
394+
await expect(
395+
applyMemberDocumentLifecycle(
396+
lifecycleInput({ allowRemoval: false, unobservedDocumentIds: ['d-removed-member'] })
397+
)
398+
).resolves.toEqual({ tombstoned: 1, resurrected: 0, purged: 0, finished: true })
399+
const [targeted] = documentUpdateConditions()
400+
expect(updatedIds(targeted)).toEqual(['d-removed-member'])
401+
expect(hasMockCondition(targeted, (node) => node.type === 'notExists')).toBe(true)
402+
/** Neither the absence reconcile nor the purge runs: absence alone still says nothing. */
403+
expect(dbChainMockFns.update).not.toHaveBeenCalledWith(schemaMock.knowledgeConnector)
404+
expect(documentUpdateConditions()).toHaveLength(1)
389405
expect(hardDeleteDocuments).not.toHaveBeenCalled()
390406
})
391407

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

Lines changed: 34 additions & 33 deletions
Original file line numberDiff line numberDiff line change
@@ -371,15 +371,17 @@ interface MemberDocumentLifecycleInput {
371371
withLease: <T>(fn: (tx: DbOrTx) => Promise<T>) => Promise<T>
372372
deadlineAt: number
373373
/**
374-
* Whether absence of observers may hide or purge a document. False until at
375-
* least one member has completed a listing: before that, nothing has been
376-
* observed yet, so absence says nothing.
374+
* Whether absence of observers may hide or purge a document the run has no
375+
* explicit word on. False until at least one member has completed a listing:
376+
* before that, nothing has been observed yet, so absence says nothing.
377377
*/
378378
allowRemoval: boolean
379379
/**
380380
* Documents whose observations this run removed. Each is tombstoned if it has
381-
* no observer left; absence arising any other way is found by the resumable
382-
* reconcile.
381+
* no observer left, whether or not `allowRemoval` holds: a removal is the
382+
* source's or the directory's explicit word, not an absence, just as the
383+
* stale-member sweep tombstones what it removes. Absence arising any other
384+
* way is found by the resumable reconcile.
383385
*/
384386
unobservedDocumentIds: Iterable<string>
385387
}
@@ -395,7 +397,7 @@ function unobservedLiveDocument(connectorId: string) {
395397
)
396398
}
397399

398-
/** The order of `doc_connector_reconciliation_idx`, which every lifecycle page walks. */
400+
/** The order of `doc_connector_reconciliation_idx`, which the resurrection pages walk. */
399401
const seenOrder = sql`COALESCE(${document.sourceSeenAt}, '-infinity'::timestamp)`
400402

401403
type TombstoneCursor = NonNullable<
@@ -432,13 +434,16 @@ async function tombstoneUnobserved(
432434
* The backstop for absence the run did not cause itself — a member row
433435
* deleted by an earlier run, a document restored from exclusion, a connector
434436
* whose first listing just completed, or a run that stopped between removing
435-
* observations and tombstoning. Walks the connector's documents in index
436-
* order, one page per statement, and checks observations only in the UPDATE
437-
* over that page's live ids: filtering the walk itself by observation lets
438-
* the LIMIT stop bounding it, and a select-list `EXISTS` can be planned as a
439-
* hash over every observation. Resumes from the cursor the previous run saved,
440-
* so each run's cost is bounded by the page budget rather than the
441-
* connector's size. Returns false when the deadline stopped it.
437+
* observations and tombstoning. Walks the connector's live documents by
438+
* external id through `doc_connector_external_id_idx`, one page per
439+
* statement, and checks observations only in the UPDATE over that page's ids:
440+
* filtering the walk itself by observation lets the LIMIT stop bounding it,
441+
* and a select-list `EXISTS` can be planned as a hash over every observation.
442+
* The key never changes for a document, unlike `source_seen_at`, which every
443+
* listing rewrites: a walk ordered by it would chase the documents each run
444+
* re-stamps and never reach the end of a pass. Resumes from the cursor the
445+
* previous run saved, so each run's cost is bounded by the page budget rather
446+
* than the connector's size. Returns false when the deadline stopped it.
442447
*/
443448
async function reconcileUnobservedPages(
444449
input: MemberDocumentLifecycleInput,
@@ -454,29 +459,26 @@ async function reconcileUnobservedPages(
454459
for (let page = 0; page < MEMBER_TOMBSTONE_RECONCILE_PAGES_PER_RUN; page++) {
455460
if (Date.now() >= input.deadlineAt) return false
456461
await input.lease.beatIfDue()
462+
/** Exclusion and archival are left to the UPDATE so every page is exactly one LIMIT of index entries. */
457463
const rows = await db
458-
.select({
459-
id: document.id,
460-
seenAt: sql<string>`${seenOrder}::text`,
461-
live: sql<boolean>`(${document.deletedAt} IS NULL)`,
462-
})
464+
.select({ id: document.id, externalId: document.externalId })
463465
.from(document)
464466
.where(
465467
and(
466468
eq(document.connectorId, connectorId),
467-
eq(document.userExcluded, false),
468-
isNull(document.archivedAt),
469-
after
470-
? sql`(${seenOrder}, ${document.id}) > (${after.seenAt}::timestamp, ${after.id})`
471-
: undefined
469+
isNull(document.deletedAt),
470+
isNotNull(document.externalId),
471+
after ? gt(document.externalId, after.externalId) : undefined
472472
)
473473
)
474-
.orderBy(seenOrder, asc(document.id))
474+
.orderBy(asc(document.externalId))
475475
.limit(MATERIALIZE_BATCH_SIZE)
476-
const candidates = rows.filter((row) => row.live).map((row) => row.id)
477-
const last = rows.at(-1)
476+
const candidates = rows.map((row) => row.id)
477+
const lastExternalId = rows.at(-1)?.externalId
478478
const next: TombstoneCursor | null =
479-
rows.length < MATERIALIZE_BATCH_SIZE || !last ? null : { seenAt: last.seenAt, id: last.id }
479+
rows.length < MATERIALIZE_BATCH_SIZE || !lastExternalId
480+
? null
481+
: { externalId: lastExternalId }
480482
if (Date.now() >= input.deadlineAt) return false
481483
const changed = await input.withLease(async (tx) => {
482484
const tombstoned =
@@ -508,7 +510,8 @@ async function reconcileUnobservedPages(
508510
* observation graph; visibility follows the active observers through the ACL.
509511
*
510512
* Tombstoning is driven by the documents whose observations this run removed,
511-
* then by a bounded slice of a resumable pass over the whole connector, so a
513+
* then, once a member has completed a listing, by a bounded slice of a
514+
* resumable pass over the whole connector, so a
512515
* run never evaluates every live document of a large connector in one
513516
* statement.
514517
*
@@ -528,11 +531,9 @@ export async function applyMemberDocumentLifecycle(
528531
purged: 0,
529532
finished: false,
530533
}
531-
if (input.allowRemoval) {
532-
const unobserved = [...new Set(input.unobservedDocumentIds)]
533-
if (!(await tombstoneUnobserved(input, unobserved, now, result))) return result
534-
if (!(await reconcileUnobservedPages(input, now, result))) return result
535-
}
534+
const unobserved = [...new Set(input.unobservedDocumentIds)]
535+
if (!(await tombstoneUnobserved(input, unobserved, now, result))) return result
536+
if (input.allowRemoval && !(await reconcileUnobservedPages(input, now, result))) return result
536537

537538
let after: { id: string; seenAt: string } | undefined
538539
for (;;) {

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

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -2267,7 +2267,9 @@ export async function executeMemberSync(
22672267
* Nobody has completed a listing yet — a connector that just entered
22682268
* members mode, waiting for its first member to connect — so an
22692269
* unobserved document says nothing about access and must not be
2270-
* tombstoned, let alone purged a week later.
2270+
* tombstoned, let alone purged a week later. What this run explicitly
2271+
* unobserved is tombstoned regardless, including after removing the
2272+
* last member that had completed a listing.
22712273
*/
22722274
const [listed] = await db
22732275
.select({ count: sql<number>`count(*)::int` })

‎packages/db/schema.ts‎

Lines changed: 4 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -5914,14 +5914,11 @@ export const knowledgeConnector = pgTable(
59145914
*/
59155915
accessRewritePending: boolean('access_rewrite_pending').notNull().default(false),
59165916
/**
5917-
* Where the members-mode absence reconcile resumes: the last document it
5918-
* checked, in `doc_connector_reconciliation_idx` order. NULL starts a new
5919-
* pass from the beginning.
5917+
* Where the members-mode absence reconcile resumes: the external id of the
5918+
* last live document it checked, in `doc_connector_external_id_idx` order.
5919+
* NULL starts a new pass from the beginning.
59205920
*/
5921-
memberTombstoneCursor: jsonb('member_tombstone_cursor').$type<{
5922-
seenAt: string
5923-
id: string
5924-
}>(),
5921+
memberTombstoneCursor: jsonb('member_tombstone_cursor').$type<{ externalId: string }>(),
59255922
/**
59265923
* One of `active`, `pending`, `syncing`, `error`, `paused`, `disabled`.
59275924
*

0 commit comments

Comments
 (0)