Skip to content

Commit dbeb732

Browse files
authored
fix(knowledge): resume the member resurrection walk from a persisted cursor (#8365)
* fix(knowledge): resume the member resurrection walk from a persisted cursor * chore(knowledge): brace the resurrection cursor branches and drop a vacuous assertion
1 parent 9070d1b commit dbeb732

10 files changed

Lines changed: 29536 additions & 28 deletions

File tree

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

Lines changed: 67 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -11,7 +11,16 @@ import {
1111
} from '@sim/db/schema'
1212
import { generateId } from '@sim/utils/id'
1313
import { and, eq, inArray, isNotNull, sql } from 'drizzle-orm'
14-
import { afterAll, afterEach, beforeEach, describe, expect, it } from 'vitest'
14+
import {
15+
afterAll,
16+
afterEach,
17+
beforeEach,
18+
describe,
19+
expect,
20+
it,
21+
type MockInstance,
22+
vi,
23+
} from 'vitest'
1524
import {
1625
type createKnowledgeAclFixtureIds,
1726
seedKnowledgeAclFixture,
@@ -142,6 +151,13 @@ describe('member document lifecycle in PostgreSQL', () => {
142151
.from(knowledgeConnector)
143152
.where(eq(knowledgeConnector.id, members.connectorId))
144153
)[0].cursor
154+
const resurrectionCursor = async () =>
155+
(
156+
await db
157+
.select({ cursor: knowledgeConnector.memberResurrectionCursor })
158+
.from(knowledgeConnector)
159+
.where(eq(knowledgeConnector.id, members.connectorId))
160+
)[0].cursor
145161

146162
it('never grants a re-owned document to observers of the connector it left', async () => {
147163
const moved = row('re-owned')
@@ -390,6 +406,56 @@ describe('member document lifecycle in PostgreSQL', () => {
390406
expect(await tombstonedIds()).toEqual(new Set([firstInEveryOrder.id]))
391407
})
392408

409+
it('resumes a resurrection walk from where the deadline stopped it, not from the first document', async () => {
410+
const observedAgain = Array.from({ length: 700 }, (_, index) => ({
411+
...row(`observed-again-${index}`),
412+
deletedAt,
413+
}))
414+
await insertRows(observedAgain)
415+
await observe(observedAgain.map(({ id }) => id))
416+
const ordered = (
417+
await db
418+
.select({ id: document.id })
419+
.from(document)
420+
.where(eq(document.connectorId, members.connectorId))
421+
.orderBy(document.id)
422+
).map(({ id }) => id)
423+
/** Stopping reads the clock past the deadline, which the walk captured when it started. */
424+
const resurrect = async (stopAfterFirstPage: boolean) => {
425+
const deadlineAt = Date.now() + 60_000
426+
let clock: MockInstance<typeof Date.now> | undefined
427+
try {
428+
return await applyMemberDocumentLifecycle({
429+
connectorId: members.connectorId,
430+
knowledgeBaseId: ids.knowledgeBaseId,
431+
runId: members.runId,
432+
allowRemoval: false,
433+
unobservedDocumentIds: [],
434+
deadlineAt,
435+
lease: { beatIfDue: async () => {} },
436+
withLease: async (fn) => {
437+
const written = await db.transaction(fn)
438+
if (stopAfterFirstPage) clock ??= vi.spyOn(Date, 'now').mockReturnValue(deadlineAt)
439+
return written
440+
},
441+
})
442+
} finally {
443+
clock?.mockRestore()
444+
}
445+
}
446+
447+
expect(await resurrect(true)).toMatchObject({ resurrected: 500, finished: false })
448+
expect(await resurrectionCursor()).toBe(ordered[499])
449+
450+
/** Resurrectable again, but behind the cursor: the resumed walk does not revisit it. */
451+
await db.update(document).set({ deletedAt }).where(eq(document.id, ordered[0]))
452+
expect(await resurrect(false)).toMatchObject({ resurrected: 200, finished: true })
453+
expect(await resurrectionCursor()).toBeNull()
454+
expect((await tombstonedIds()).has(ordered[0])).toBe(true)
455+
456+
expect(await resurrect(false)).toMatchObject({ resurrected: 1, finished: true })
457+
})
458+
393459
it('continues past a full selected batch even if its observations change before UPDATE', async () => {
394460
const rows = Array.from({ length: 501 }, (_, index) => row(String(index)))
395461
await db.insert(document).values(rows)

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

Lines changed: 20 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -890,7 +890,8 @@ async function reconcileUnobservedPages(
890890
* then, once a member has completed a listing, by a bounded slice of a
891891
* resumable pass over the whole connector, so a
892892
* run never evaluates every live document of a large connector in one
893-
* statement.
893+
* statement. The resurrection walk likewise resumes where a deadline last
894+
* stopped it, rather than from the connector's first document.
894895
*
895896
* A document whose content refresh failed this run is not resurrected: its
896897
* stored content is known-stale, and surfacing it would show pre-tombstone
@@ -912,9 +913,15 @@ export async function applyMemberDocumentLifecycle(
912913
if (!(await tombstoneUnobserved(input, unobserved, now, result))) return result
913914
if (input.allowRemoval && !(await reconcileUnobservedPages(input, now, result))) return result
914915

916+
const [connector] = await db
917+
.select({ cursor: knowledgeConnector.memberResurrectionCursor })
918+
.from(knowledgeConnector)
919+
.where(eq(knowledgeConnector.id, connectorId))
920+
const startAfterId = connector?.cursor ?? undefined
915921
const condition = resurrectableDocument(connectorId)
916-
const resurrectionFinished = await walkReconciliationWindows({
922+
const resurrection = await walkReconciliationWindows({
917923
connectorId,
924+
startAfterId,
918925
condition,
919926
pageSize: LIFECYCLE_PAGE_SIZE,
920927
deadlineAt: input.deadlineAt,
@@ -939,7 +946,17 @@ export async function applyMemberDocumentLifecycle(
939946
result.resurrected += changed.length
940947
},
941948
})
942-
if (!resurrectionFinished) return result
949+
/** One write per run, and only when the resume point moved: the connector row is hot. */
950+
const resumeAfterId = resurrection.finished ? undefined : resurrection.lastId
951+
if (resumeAfterId !== startAfterId) {
952+
await input.withLease((tx) =>
953+
tx
954+
.update(knowledgeConnector)
955+
.set({ memberResurrectionCursor: resumeAfterId ?? null })
956+
.where(eq(knowledgeConnector.id, connectorId))
957+
)
958+
}
959+
if (!resurrection.finished) return result
943960

944961
const purgeCutoff = new Date(now.getTime() - MEMBER_TOMBSTONE_PURGE_DAYS * 24 * 60 * 60 * 1000)
945962
const purgeCandidates = input.allowRemoval

‎apps/sim/lib/knowledge/connectors/reconciliation-window.ts‎

Lines changed: 28 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,8 @@ interface ReconciliationRow {
1919

2020
export interface ReconciliationWalk {
2121
connectorId: string
22+
/** Resumes after this document id, the `lastId` an earlier walk stopped at. */
23+
startAfterId?: string
2224
/**
2325
* Evaluated over the window's rows, which carry only the document's id, connector, exclusion,
2426
* archival, tombstone, seen, ACL and content-hash columns.
@@ -31,6 +33,13 @@ export interface ReconciliationWalk {
3133
onPage: (rows: ReconciliationRow[]) => Promise<void>
3234
}
3335

36+
export interface ReconciliationWalkResult {
37+
/** False when the deadline stopped the walk. */
38+
finished: boolean
39+
/** The last document id the walk covered, from which a stopped walk resumes. */
40+
lastId: string | undefined
41+
}
42+
3443
/** A type alias, not an interface, so it satisfies `db.execute`'s row-record constraint. */
3544
type WindowScan = {
3645
size: number
@@ -97,29 +106,31 @@ async function scanWindow(walk: ReconciliationWalk, afterId: string | undefined)
97106
}
98107

99108
/**
100-
* Walks a connector's owned documents in id order, one window per statement, handing the rows
101-
* matching `condition` to `onPage` in pages of at most `pageSize`. A window without matches still
102-
* advances the walk, which ends after the first window shorter than a full one. Returns false when
103-
* the deadline stopped it.
109+
* Walks a connector's owned documents in id order from `startAfterId`, one window per statement,
110+
* handing the rows matching `condition` to `onPage` in pages of at most `pageSize`. A window
111+
* without matches still advances the walk, which ends after the first window shorter than a full
112+
* one.
104113
*/
105-
export async function walkReconciliationWindows(walk: ReconciliationWalk): Promise<boolean> {
106-
let after: string | undefined
114+
export async function walkReconciliationWindows(
115+
walk: ReconciliationWalk
116+
): Promise<ReconciliationWalkResult> {
117+
let covered = walk.startAfterId
107118
for (;;) {
108-
if (Date.now() >= walk.deadlineAt) return false
119+
if (Date.now() >= walk.deadlineAt) return { finished: false, lastId: covered }
109120
await walk.beforePage()
110-
if (Date.now() >= walk.deadlineAt) return false
111-
const { size, last, ids, tombstoned } = await scanWindow(walk, after)
121+
if (Date.now() >= walk.deadlineAt) return { finished: false, lastId: covered }
122+
const { size, last, ids, tombstoned } = await scanWindow(walk, covered)
112123
for (let offset = 0; offset < ids.length; offset += walk.pageSize) {
113124
if (offset > 0) await walk.beforePage()
114125
/** Materialized ids are acted on only inside the budget, so a late page is left for the next run. */
115-
if (Date.now() >= walk.deadlineAt) return false
116-
await walk.onPage(
117-
ids
118-
.slice(offset, offset + walk.pageSize)
119-
.map((id, index) => ({ id, tombstoned: tombstoned[offset + index] }))
120-
)
126+
if (Date.now() >= walk.deadlineAt) return { finished: false, lastId: covered }
127+
const page = ids.slice(offset, offset + walk.pageSize)
128+
await walk.onPage(page.map((id, index) => ({ id, tombstoned: tombstoned[offset + index] })))
129+
covered = page.at(-1)
130+
}
131+
if (size < RECONCILIATION_WINDOW_SIZE || !last) {
132+
return { finished: true, lastId: last ?? covered }
121133
}
122-
if (size < RECONCILIATION_WINDOW_SIZE || !last) return true
123-
after = last
134+
covered = last
124135
}
125136
}

‎apps/sim/lib/knowledge/connectors/sync-content-pass.ts‎

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -433,7 +433,7 @@ async function reconcileCompletedListing(
433433
hardHeld
434434
)
435435
if (input.documentAccess === 'admin' && aclCount > 0) {
436-
const finished = await walk(aclAbsent, 500, async (rows) => {
436+
const { finished } = await walk(aclAbsent, 500, async (rows) => {
437437
await revokeDocumentAcls(
438438
withAclPage,
439439
rows.map((row) => row.id),
@@ -445,7 +445,7 @@ async function reconcileCompletedListing(
445445
}
446446
if (!allowDeletion) return { finished: true, notice }
447447
if (!checkpoint.fullSync && !softHeld && softCount > 0) {
448-
const finished = await walk(soft, 500, async (rows) => {
448+
const { finished } = await walk(soft, 500, async (rows) => {
449449
const removed = await withLease((tx) =>
450450
tx
451451
.update(document)
@@ -466,7 +466,7 @@ async function reconcileCompletedListing(
466466
if (!finished) return { finished: false, notice }
467467
}
468468
if (!hardHeld && hardCount > 0) {
469-
const finished = await walk(hard, 25, async (rows) => {
469+
const { finished } = await walk(hard, 25, async (rows) => {
470470
/** Report newly removed documents once; purging existing tombstones is storage cleanup. */
471471
for (const tombstoned of [false, true]) {
472472
const ids = rows.filter((row) => row.tombstoned === tombstoned).map((row) => row.id)

‎apps/sim/lib/knowledge/orchestration/connectors.ts‎

Lines changed: 12 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -138,10 +138,13 @@ export type KnowledgeConnectorRow = typeof knowledgeConnector.$inferSelect
138138
type ConnectorRow = KnowledgeConnectorRow
139139
/**
140140
* The connector row as it reaches every caller: never carrying the stored API
141-
* key, nor the members-mode reconcile cursor, which names a document the
142-
* caller may not be able to read.
141+
* key, nor the members-mode reconcile cursors, which name documents the caller
142+
* may not be able to read.
143143
*/
144-
export type ConnectorWithoutSecret = Omit<ConnectorRow, 'encryptedApiKey' | 'memberTombstoneCursor'>
144+
export type ConnectorWithoutSecret = Omit<
145+
ConnectorRow,
146+
'encryptedApiKey' | 'memberTombstoneCursor' | 'memberResurrectionCursor'
147+
>
145148

146149
/** A refused `sourceConfig`, with the failure class the caller wants surfaced. */
147150
export interface SourceConfigRejection {
@@ -159,7 +162,12 @@ export interface ConnectorKnowledgeBase {
159162
}
160163

161164
export function withoutSecret(row: ConnectorRow): ConnectorWithoutSecret {
162-
const { encryptedApiKey: _encryptedApiKey, memberTombstoneCursor: _cursor, ...rest } = row
165+
const {
166+
encryptedApiKey: _encryptedApiKey,
167+
memberTombstoneCursor: _tombstoneCursor,
168+
memberResurrectionCursor: _resurrectionCursor,
169+
...rest
170+
} = row
163171
return rest
164172
}
165173

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1 @@
1+
ALTER TABLE "knowledge_connector" ADD COLUMN "member_resurrection_cursor" text;

0 commit comments

Comments
 (0)