Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,16 @@ import {
} from '@sim/db/schema'
import { generateId } from '@sim/utils/id'
import { and, eq, inArray, isNotNull, sql } from 'drizzle-orm'
import { afterAll, afterEach, beforeEach, describe, expect, it } from 'vitest'
import {
afterAll,
afterEach,
beforeEach,
describe,
expect,
it,
type MockInstance,
vi,
} from 'vitest'
import {
type createKnowledgeAclFixtureIds,
seedKnowledgeAclFixture,
Expand Down Expand Up @@ -142,6 +151,13 @@ describe('member document lifecycle in PostgreSQL', () => {
.from(knowledgeConnector)
.where(eq(knowledgeConnector.id, members.connectorId))
)[0].cursor
const resurrectionCursor = async () =>
(
await db
.select({ cursor: knowledgeConnector.memberResurrectionCursor })
.from(knowledgeConnector)
.where(eq(knowledgeConnector.id, members.connectorId))
)[0].cursor

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

it('resumes a resurrection walk from where the deadline stopped it, not from the first document', async () => {
const observedAgain = Array.from({ length: 700 }, (_, index) => ({
...row(`observed-again-${index}`),
deletedAt,
}))
await insertRows(observedAgain)
await observe(observedAgain.map(({ id }) => id))
const ordered = (
await db
.select({ id: document.id })
.from(document)
.where(eq(document.connectorId, members.connectorId))
.orderBy(document.id)
).map(({ id }) => id)
/** Stopping reads the clock past the deadline, which the walk captured when it started. */
const resurrect = async (stopAfterFirstPage: boolean) => {
const deadlineAt = Date.now() + 60_000
let clock: MockInstance<typeof Date.now> | undefined
try {
return await applyMemberDocumentLifecycle({
connectorId: members.connectorId,
knowledgeBaseId: ids.knowledgeBaseId,
runId: members.runId,
allowRemoval: false,
unobservedDocumentIds: [],
deadlineAt,
lease: { beatIfDue: async () => {} },
withLease: async (fn) => {
const written = await db.transaction(fn)
if (stopAfterFirstPage) clock ??= vi.spyOn(Date, 'now').mockReturnValue(deadlineAt)
return written
},
})
} finally {
clock?.mockRestore()
}
}

expect(await resurrect(true)).toMatchObject({ resurrected: 500, finished: false })
expect(await resurrectionCursor()).toBe(ordered[499])

/** Resurrectable again, but behind the cursor: the resumed walk does not revisit it. */
await db.update(document).set({ deletedAt }).where(eq(document.id, ordered[0]))
expect(await resurrect(false)).toMatchObject({ resurrected: 200, finished: true })
expect(await resurrectionCursor()).toBeNull()
expect((await tombstonedIds()).has(ordered[0])).toBe(true)

expect(await resurrect(false)).toMatchObject({ resurrected: 1, finished: true })
})

it('continues past a full selected batch even if its observations change before UPDATE', async () => {
const rows = Array.from({ length: 501 }, (_, index) => row(String(index)))
await db.insert(document).values(rows)
Expand Down
23 changes: 20 additions & 3 deletions apps/sim/lib/knowledge/connectors/member-observations.ts
Original file line number Diff line number Diff line change
Expand Up @@ -890,7 +890,8 @@ async function reconcileUnobservedPages(
* then, once a member has completed a listing, by a bounded slice of a
* resumable pass over the whole connector, so a
* run never evaluates every live document of a large connector in one
* statement.
* statement. The resurrection walk likewise resumes where a deadline last
* stopped it, rather than from the connector's first document.
*
* A document whose content refresh failed this run is not resurrected: its
* stored content is known-stale, and surfacing it would show pre-tombstone
Expand All @@ -912,9 +913,15 @@ export async function applyMemberDocumentLifecycle(
if (!(await tombstoneUnobserved(input, unobserved, now, result))) return result
if (input.allowRemoval && !(await reconcileUnobservedPages(input, now, result))) return result

const [connector] = await db
.select({ cursor: knowledgeConnector.memberResurrectionCursor })
.from(knowledgeConnector)
.where(eq(knowledgeConnector.id, connectorId))
const startAfterId = connector?.cursor ?? undefined
const condition = resurrectableDocument(connectorId)
const resurrectionFinished = await walkReconciliationWindows({
const resurrection = await walkReconciliationWindows({
connectorId,
startAfterId,
condition,
pageSize: LIFECYCLE_PAGE_SIZE,
deadlineAt: input.deadlineAt,
Expand All @@ -939,7 +946,17 @@ export async function applyMemberDocumentLifecycle(
result.resurrected += changed.length
},
})
if (!resurrectionFinished) return result
/** One write per run, and only when the resume point moved: the connector row is hot. */
const resumeAfterId = resurrection.finished ? undefined : resurrection.lastId
if (resumeAfterId !== startAfterId) {
await input.withLease((tx) =>
tx
.update(knowledgeConnector)
.set({ memberResurrectionCursor: resumeAfterId ?? null })
.where(eq(knowledgeConnector.id, connectorId))
)
}
if (!resurrection.finished) return result

const purgeCutoff = new Date(now.getTime() - MEMBER_TOMBSTONE_PURGE_DAYS * 24 * 60 * 60 * 1000)
const purgeCandidates = input.allowRemoval
Expand Down
45 changes: 28 additions & 17 deletions apps/sim/lib/knowledge/connectors/reconciliation-window.ts
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,8 @@ interface ReconciliationRow {

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

export interface ReconciliationWalkResult {
/** False when the deadline stopped the walk. */
finished: boolean
/** The last document id the walk covered, from which a stopped walk resumes. */
lastId: string | undefined
}

/** A type alias, not an interface, so it satisfies `db.execute`'s row-record constraint. */
type WindowScan = {
size: number
Expand Down Expand Up @@ -97,29 +106,31 @@ async function scanWindow(walk: ReconciliationWalk, afterId: string | undefined)
}

/**
* Walks a connector's owned documents in id order, one window per statement, handing the rows
* matching `condition` to `onPage` in pages of at most `pageSize`. A window without matches still
* advances the walk, which ends after the first window shorter than a full one. Returns false when
* the deadline stopped it.
* Walks a connector's owned documents in id order from `startAfterId`, one window per statement,
* handing the rows matching `condition` to `onPage` in pages of at most `pageSize`. A window
* without matches still advances the walk, which ends after the first window shorter than a full
* one.
*/
export async function walkReconciliationWindows(walk: ReconciliationWalk): Promise<boolean> {
let after: string | undefined
export async function walkReconciliationWindows(
walk: ReconciliationWalk
): Promise<ReconciliationWalkResult> {
let covered = walk.startAfterId
for (;;) {
if (Date.now() >= walk.deadlineAt) return false
if (Date.now() >= walk.deadlineAt) return { finished: false, lastId: covered }
await walk.beforePage()
if (Date.now() >= walk.deadlineAt) return false
const { size, last, ids, tombstoned } = await scanWindow(walk, after)
if (Date.now() >= walk.deadlineAt) return { finished: false, lastId: covered }
const { size, last, ids, tombstoned } = await scanWindow(walk, covered)
for (let offset = 0; offset < ids.length; offset += walk.pageSize) {
if (offset > 0) await walk.beforePage()
/** Materialized ids are acted on only inside the budget, so a late page is left for the next run. */
if (Date.now() >= walk.deadlineAt) return false
await walk.onPage(
ids
.slice(offset, offset + walk.pageSize)
.map((id, index) => ({ id, tombstoned: tombstoned[offset + index] }))
)
if (Date.now() >= walk.deadlineAt) return { finished: false, lastId: covered }
const page = ids.slice(offset, offset + walk.pageSize)
await walk.onPage(page.map((id, index) => ({ id, tombstoned: tombstoned[offset + index] })))
covered = page.at(-1)
}
if (size < RECONCILIATION_WINDOW_SIZE || !last) {
return { finished: true, lastId: last ?? covered }
}
if (size < RECONCILIATION_WINDOW_SIZE || !last) return true
after = last
covered = last
}
}
6 changes: 3 additions & 3 deletions apps/sim/lib/knowledge/connectors/sync-content-pass.ts
Original file line number Diff line number Diff line change
Expand Up @@ -433,7 +433,7 @@ async function reconcileCompletedListing(
hardHeld
)
if (input.documentAccess === 'admin' && aclCount > 0) {
const finished = await walk(aclAbsent, 500, async (rows) => {
const { finished } = await walk(aclAbsent, 500, async (rows) => {
await revokeDocumentAcls(
withAclPage,
rows.map((row) => row.id),
Expand All @@ -445,7 +445,7 @@ async function reconcileCompletedListing(
}
if (!allowDeletion) return { finished: true, notice }
if (!checkpoint.fullSync && !softHeld && softCount > 0) {
const finished = await walk(soft, 500, async (rows) => {
const { finished } = await walk(soft, 500, async (rows) => {
const removed = await withLease((tx) =>
tx
.update(document)
Expand All @@ -466,7 +466,7 @@ async function reconcileCompletedListing(
if (!finished) return { finished: false, notice }
}
if (!hardHeld && hardCount > 0) {
const finished = await walk(hard, 25, async (rows) => {
const { finished } = await walk(hard, 25, async (rows) => {
/** Report newly removed documents once; purging existing tombstones is storage cleanup. */
for (const tombstoned of [false, true]) {
const ids = rows.filter((row) => row.tombstoned === tombstoned).map((row) => row.id)
Expand Down
16 changes: 12 additions & 4 deletions apps/sim/lib/knowledge/orchestration/connectors.ts
Original file line number Diff line number Diff line change
Expand Up @@ -138,10 +138,13 @@ export type KnowledgeConnectorRow = typeof knowledgeConnector.$inferSelect
type ConnectorRow = KnowledgeConnectorRow
/**
* The connector row as it reaches every caller: never carrying the stored API
* key, nor the members-mode reconcile cursor, which names a document the
* caller may not be able to read.
* key, nor the members-mode reconcile cursors, which name documents the caller
* may not be able to read.
*/
export type ConnectorWithoutSecret = Omit<ConnectorRow, 'encryptedApiKey' | 'memberTombstoneCursor'>
export type ConnectorWithoutSecret = Omit<
ConnectorRow,
'encryptedApiKey' | 'memberTombstoneCursor' | 'memberResurrectionCursor'
>

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

export function withoutSecret(row: ConnectorRow): ConnectorWithoutSecret {
const { encryptedApiKey: _encryptedApiKey, memberTombstoneCursor: _cursor, ...rest } = row
const {
encryptedApiKey: _encryptedApiKey,
memberTombstoneCursor: _tombstoneCursor,
memberResurrectionCursor: _resurrectionCursor,
...rest
} = row
return rest
}

Expand Down
1 change: 1 addition & 0 deletions packages/db/migrations/0389_member_resurrection_cursor.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
ALTER TABLE "knowledge_connector" ADD COLUMN "member_resurrection_cursor" text;
Loading
Loading