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
1 change: 1 addition & 0 deletions .github/workflows/test-build.yml
Original file line number Diff line number Diff line change
Expand Up @@ -274,6 +274,7 @@ jobs:
lib/knowledge/__integration__/user-document-visibility.integration.ts
lib/knowledge/__integration__/listing-continuation.integration.ts
lib/knowledge/__integration__/member-scope-renewal.integration.ts
lib/knowledge/__integration__/member-document-lifecycle.integration.ts
lib/knowledge/__integration__/slack-empty-threads.integration.ts
lib/knowledge/__integration__/kb-block-search.integration.ts
lib/knowledge/__integration__/unfilled-projection-source.integration.ts
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -3,13 +3,14 @@ import { db } from '@sim/db'
import {
document,
knowledgeConnector,
knowledgeConnectorMember,
knowledgeDocumentObservation,
organization,
user,
workspace,
} from '@sim/db/schema'
import { generateId } from '@sim/utils/id'
import { eq, inArray, sql } from 'drizzle-orm'
import { and, eq, inArray, isNotNull, sql } from 'drizzle-orm'
import { afterAll, afterEach, beforeEach, describe, expect, it } from 'vitest'
import {
type createKnowledgeAclFixtureIds,
Expand All @@ -19,9 +20,13 @@ import {
import {
applyMemberDocumentLifecycle,
recordMemberObservations,
removeMemberObservationsForDocuments,
} from '@/lib/knowledge/connectors/member-observations'
import { resumeMembershipRewrites } from '@/lib/knowledge/connectors/member-sync-engine'
import { MEMBER_TOMBSTONE_RECONCILE_PAGES_PER_RUN } from '@/lib/knowledge/connectors/sync-limits'
import {
assertSyncLeaseHeldInTx,
createMemberSyncLease,
SyncLockLostException,
stillHoldsMemberSyncLock,
} from '@/lib/knowledge/connectors/sync-lock'
Expand Down Expand Up @@ -62,12 +67,19 @@ describe('member document lifecycle in PostgreSQL', () => {
const observe = (documentIds: string[]) =>
recordMemberObservations(db, members.members[0].id, documentIds, members.runId)

const run = (options: { beforeWrite?: () => Promise<void>; allowRemoval?: boolean } = {}) =>
const run = (
options: {
beforeWrite?: () => Promise<void>
allowRemoval?: boolean
unobservedDocumentIds?: string[]
} = {}
) =>
applyMemberDocumentLifecycle({
connectorId: members.connectorId,
knowledgeBaseId: ids.knowledgeBaseId,
runId: members.runId,
allowRemoval: options.allowRemoval ?? true,
unobservedDocumentIds: options.unobservedDocumentIds ?? [],
deadlineAt: Date.now() + 60_000,
lease: { beatIfDue: async () => {} },
withLease: async (fn) => {
Expand Down Expand Up @@ -108,6 +120,256 @@ describe('member document lifecycle in PostgreSQL', () => {
expect(byId.get(noContent.id)).toEqual(deletedAt)
})

const insertRows = async (rows: ReturnType<typeof row>[]) => {
for (let offset = 0; offset < rows.length; offset += 500)
await db.insert(document).values(rows.slice(offset, offset + 500))
}
const tombstonedIds = async () =>
new Set(
(
await db
.select({ id: document.id })
.from(document)
.where(and(eq(document.connectorId, members.connectorId), isNotNull(document.deletedAt)))
).map(({ id }) => id)
)
const savedCursor = async () =>
(
await db
.select({ cursor: knowledgeConnector.memberTombstoneCursor })
.from(knowledgeConnector)
.where(eq(knowledgeConnector.id, members.connectorId))
)[0].cursor

it('tombstones what this run unobserved right away and leaves the rest of a large connector to later runs', async () => {
const pageBudget = MEMBER_TOMBSTONE_RECONCILE_PAGES_PER_RUN * 500
const unobserved = Array.from({ length: pageBudget + 20 }, (_, index) =>
row(`unobserved-${index}`)
)
/** Both sort after every unobserved document, beyond what this run's reconcile reaches. */
const [lostByThisRun, stillObservedByBob] = [
row('zz-lost-by-this-run'),
row('zz-still-observed-by-bob'),
]
await insertRows([...unobserved, lostByThisRun, stillObservedByBob])
await observe([lostByThisRun.id, stillObservedByBob.id])
await recordMemberObservations(
db,
members.members[1].id,
[stillObservedByBob.id],
members.runId
)
const removed = await removeMemberObservationsForDocuments(db, members.members[0].id, [
lostByThisRun.id,
stillObservedByBob.id,
])
expect(removed.sort()).toEqual([lostByThisRun.id, stillObservedByBob.id].sort())

expect(await run({ unobservedDocumentIds: removed })).toEqual({
tombstoned: pageBudget + 1,
resurrected: 0,
purged: 0,
finished: true,
})
const afterFirst = await tombstonedIds()
expect(afterFirst.has(lostByThisRun.id)).toBe(true)
expect(afterFirst.has(stillObservedByBob.id)).toBe(false)
expect(unobserved.filter(({ id }) => !afterFirst.has(id))).toHaveLength(20)
expect(await savedCursor()).toEqual({ externalId: expect.any(String) })

expect(await run()).toEqual({ tombstoned: 20, resurrected: 0, purged: 0, finished: true })
const afterSecond = await tombstonedIds()
expect(unobserved.every(({ id }) => afterSecond.has(id))).toBe(true)
expect(afterSecond.has(stillObservedByBob.id)).toBe(false)
expect(await savedCursor()).toBeNull()
})

it('tombstones what removing the only listed member unobserved, though the run stops right after the removal', async () => {
const [removedMember, otherMember] = members.members
await db
.update(knowledgeConnectorMember)
.set({
lastCompleteListingAt: new Date(),
listingCheckpoint: { kind: 'membership', cursor: null, removeMember: true },
})
.where(eq(knowledgeConnectorMember.id, removedMember.id))
const onlyRemoved = Array.from({ length: 1_200 }, (_, index) => row(`only-removed-${index}`))
const sharedWithOther = row('shared-with-other-member')
const neverObserved = row('never-observed')
await insertRows([...onlyRemoved, sharedWithOther, neverObserved])
await observe([...onlyRemoved.map(({ id }) => id), sharedWithOther.id])
await recordMemberObservations(db, otherMember.id, [sharedWithOther.id], members.runId)
const removal = (stopAfterFirstPage: boolean) => {
const lease = createMemberSyncLease(members.connectorId, members.runId)
const input: Parameters<typeof resumeMembershipRewrites>[0] = {
connectorId: members.connectorId,
runId: members.runId,
deadlineAt: Date.now() + 60_000,
tombstonesUnobserved: true,
lease: {
...lease,
beatIfDue: async () => {
await lease.beatIfDue()
if (stopAfterFirstPage) input.deadlineAt = Date.now() - 1
},
},
}
return resumeMembershipRewrites(input)
}

/** The first run walks one page before its deadline; the second finishes and deletes the member. */
expect(await removal(true)).toBe(false)
const [paused] = await db
.select({ checkpoint: knowledgeConnectorMember.listingCheckpoint })
.from(knowledgeConnectorMember)
.where(eq(knowledgeConnectorMember.id, removedMember.id))
expect(paused.checkpoint).toMatchObject({ removeMember: true, cursor: expect.any(String) })
expect(await removal(false)).toBe(true)
const [completed] = await db
.select({ count: sql<number>`count(*)::int` })
.from(knowledgeConnectorMember)
.where(
and(
eq(knowledgeConnectorMember.connectorId, members.connectorId),
isNotNull(knowledgeConnectorMember.lastCompleteListingAt)
)
)
expect(completed.count).toBe(0)
/** No lifecycle ran in either run: the tombstones landed with the removal itself. */
const afterRemoval = await tombstonedIds()
expect(onlyRemoved.every(({ id }) => afterRemoval.has(id))).toBe(true)
expect(afterRemoval.has(sharedWithOther.id)).toBe(false)
expect(afterRemoval.has(neverObserved.id)).toBe(false)

expect(await run({ allowRemoval: false })).toEqual({
tombstoned: 0,
resurrected: 0,
purged: 0,
finished: true,
})
expect(await tombstonedIds()).toEqual(afterRemoval)
})

it('brings back what a withdrawn removal tombstoned within the same run, and nothing else', async () => {
const [member] = members.members
await db
.update(knowledgeConnectorMember)
.set({ listingCheckpoint: { kind: 'membership', cursor: null, removeMember: true } })
.where(eq(knowledgeConnectorMember.id, member.id))
const walked = Array.from({ length: 600 }, (_, index) => row(`walked-${index}`))
const excluded = { ...row('excluded'), userExcluded: true, deletedAt }
const archived = { ...row('archived'), archivedAt: new Date(), deletedAt }
const noContent = { ...row('no-content'), contentHash: null, deletedAt }
await insertRows([...walked, excluded, archived])
await db.insert(document).values(noContent)
await observe([...walked, excluded, archived, noContent].map(({ id }) => id))
const walk = (stopAfterFirstPage: boolean) => {
const lease = createMemberSyncLease(members.connectorId, members.runId)
const input: Parameters<typeof resumeMembershipRewrites>[0] = {
connectorId: members.connectorId,
runId: members.runId,
deadlineAt: Date.now() + 60_000,
tombstonesUnobserved: true,
lease: {
...lease,
beatIfDue: async () => {
await lease.beatIfDue()
if (stopAfterFirstPage) input.deadlineAt = Date.now() - 1
},
},
}
return resumeMembershipRewrites(input)
}

expect(await walk(true)).toBe(false)
const tombstonedByRemoval = await tombstonedIds()
const walkedIds = new Set(walked.map(({ id }) => id))
const removedPage = [...tombstonedByRemoval].filter((id) => walkedIds.has(id))
expect(removedPage.length).toBeGreaterThan(0)
expect(removedPage.length).toBeLessThan(walked.length)

/** Directory re-listing withdraws the removal: the member is active again and its walk restarts. */
await db
.update(knowledgeConnectorMember)
.set({
status: 'active',
listingCheckpoint: { kind: 'membership', cursor: null, removeMember: false },
})
.where(eq(knowledgeConnectorMember.id, member.id))
expect(await walk(false)).toBe(true)

const after = await tombstonedIds()
expect(walked.every(({ id }) => !after.has(id))).toBe(true)
for (const kept of [excluded, archived, noContent]) expect(after.has(kept.id)).toBe(true)
})

it('leaves a service-owned corpus alone when a member is removed', async () => {
const [removedMember] = members.members
await db
.update(knowledgeConnectorMember)
.set({ listingCheckpoint: { kind: 'membership', cursor: null, removeMember: true } })
.where(eq(knowledgeConnectorMember.id, removedMember.id))
const onlyRemoved = row('service-owned')
await insertRows([onlyRemoved])
await observe([onlyRemoved.id])
expect(
await resumeMembershipRewrites({
connectorId: members.connectorId,
runId: members.runId,
deadlineAt: Date.now() + 60_000,
lease: createMemberSyncLease(members.connectorId, members.runId),
tombstonesUnobserved: false,
})
).toBe(true)
expect((await tombstonedIds()).size).toBe(0)
})

it('finishes a pass within its page budget while listings re-stamp every observed document', async () => {
const pageBudget = MEMBER_TOMBSTONE_RECONCILE_PAGES_PER_RUN * 500
const total = pageBudget + 500
const runsPerPass = Math.ceil(total / pageBudget)
const firstInEveryOrder = {
...row('walk-0000000'),
id: '00000000-0000-4000-8000-000000000000',
}
const rest = Array.from({ length: total - 1 }, (_, index) =>
row(`walk-${String(index + 1).padStart(7, '0')}`)
)
await insertRows([firstInEveryOrder, ...rest])
const all = [firstInEveryOrder, ...rest].map(({ id }) => id)
for (let offset = 0; offset < all.length; offset += 5000)
await observe(all.slice(offset, offset + 5000))
/** A listing stamps what it saw with its start; the document nobody observes keeps its old stamp. */
const relist = () =>
db
.update(document)
.set({ sourceSeenAt: new Date() })
.where(
and(
eq(document.connectorId, members.connectorId),
sql`EXISTS (SELECT 1 FROM knowledge_document_observation o WHERE o.document_id = ${document.id})`
)
)

expect(await run()).toMatchObject({ tombstoned: 0 })
expect(await savedCursor()).not.toBeNull()
await db
.delete(knowledgeDocumentObservation)
.where(eq(knowledgeDocumentObservation.documentId, firstInEveryOrder.id))
for (let pass = 1; pass < runsPerPass; pass++) {
await relist()
await run()
}
expect(await savedCursor()).toBeNull()
expect((await tombstonedIds()).has(firstInEveryOrder.id)).toBe(false)

for (let next = 0; next < runsPerPass; next++) {
await relist()
await run()
}
expect(await tombstonedIds()).toEqual(new Set([firstInEveryOrder.id]))
})

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
16 changes: 10 additions & 6 deletions apps/sim/lib/knowledge/application/connectors.ts
Original file line number Diff line number Diff line change
Expand Up @@ -98,6 +98,7 @@ import {
performSyncKnowledgeConnector,
performUpdateKnowledgeConnector,
type SourceConfigRejection,
withoutSecret,
} from '@/lib/knowledge/orchestration/connectors'
import type {
KnowledgeOperationSource,
Expand Down Expand Up @@ -548,11 +549,14 @@ export const listKnowledgeConnectors = defineAuthorizedKnowledgeUseCase({
})
: new Map<string, ViewerConnectorMembership>()
return {
connectors: page.map(({ encryptedApiKey: _encryptedApiKey, ...rest }) => ({
...rest,
permissionConfig: permissionSummaries.get(rest.id),
viewerMembership: memberships.get(rest.id) ?? null,
})),
connectors: page.map((row) => {
const rest = withoutSecret(row)
return {
...rest,
permissionConfig: permissionSummaries.get(rest.id),
viewerMembership: memberships.get(rest.id) ?? null,
}
}),
hasMore,
offset,
limit: input.limit ?? page.length,
Expand Down Expand Up @@ -709,7 +713,7 @@ export const readKnowledgeConnector = defineAuthorizedKnowledgeUseCase({
? summarizeConnectorMembers(context.connectorId, connector.syncIntervalMinutes)
: { active: 0, suspended: 0, stale: 0 },
])
const { encryptedApiKey: _encryptedApiKey, ...connectorData } = connector
const connectorData = withoutSecret(connector)
const viewerUserId = principal.kind === 'session' ? principal.userId : null
const memberships =
viewerUserId && (context.workspaceId || context.organizationId)
Expand Down
Loading
Loading