From f02249a802b332f8a710f725642d8caee0bbb3d9 Mon Sep 17 00:00:00 2001 From: Vikhyath Mondreti Date: Tue, 15 Sep 2026 11:20:23 -0700 Subject: [PATCH] fix(knowledge): remove connections with durable background cleanup --- .../[connectorId]/source-detail.test.tsx | 14 +- .../sources/[connectorId]/source-detail.tsx | 9 +- .../connectors-section.test.tsx | 21 +- .../use-connector-actions.ts | 27 +-- .../queries/kb/connectors-cache.test.tsx | 40 ++++ apps/sim/hooks/queries/kb/connectors.ts | 24 +- .../storage-accounting.integration.ts | 142 +++++++++++- .../lib/knowledge/access/predicate.test.ts | 8 +- apps/sim/lib/knowledge/access/predicate.ts | 5 +- apps/sim/lib/knowledge/access/types.ts | 5 +- apps/sim/lib/knowledge/connectors/deletion.md | 13 ++ .../lib/knowledge/connectors/deletion.test.ts | 189 ++++++++++++++++ apps/sim/lib/knowledge/connectors/deletion.ts | 205 ++++++++++++++++++ .../documents/connector-lifecycle.ts | 13 ++ .../documents/processing-outbox-handler.ts | 5 + apps/sim/lib/knowledge/documents/service.ts | 15 +- .../orchestration/connectors.test.ts | 91 +++++++- .../lib/knowledge/orchestration/connectors.ts | 202 +++++++++-------- apps/sim/lib/knowledge/tags/service.test.ts | 46 +++- apps/sim/lib/knowledge/tags/service.ts | 29 ++- 20 files changed, 943 insertions(+), 160 deletions(-) create mode 100644 apps/sim/lib/knowledge/connectors/deletion.md create mode 100644 apps/sim/lib/knowledge/connectors/deletion.test.ts create mode 100644 apps/sim/lib/knowledge/connectors/deletion.ts create mode 100644 apps/sim/lib/knowledge/documents/connector-lifecycle.ts diff --git a/apps/sim/app/o/[organizationId]/settings/integrations/sources/[connectorId]/source-detail.test.tsx b/apps/sim/app/o/[organizationId]/settings/integrations/sources/[connectorId]/source-detail.test.tsx index 0dbbe4341de..d277c6a025d 100644 --- a/apps/sim/app/o/[organizationId]/settings/integrations/sources/[connectorId]/source-detail.test.tsx +++ b/apps/sim/app/o/[organizationId]/settings/integrations/sources/[connectorId]/source-detail.test.tsx @@ -13,6 +13,7 @@ const mocks = vi.hoisted(() => ({ detail: vi.fn(), integrations: vi.fn(), push: vi.fn(), + replace: vi.fn(), documents: vi.fn(), actions: vi.fn(), recovery: vi.fn(), @@ -23,7 +24,7 @@ const mocks = vi.hoisted(() => ({ save: vi.fn(), })) vi.mock('next/navigation', () => ({ - useRouter: () => ({ push: mocks.push }), + useRouter: () => ({ push: mocks.push, replace: mocks.replace }), usePathname: () => '/o/org-one/settings/integrations/sources/source-one', })) vi.mock('@/app/o/[organizationId]/providers/organization-provider', () => ({ @@ -185,6 +186,17 @@ describe('organization source detail navigation', () => { expect(button, `Missing ${text}`).toBeTruthy() await act(async () => button!.click()) } + it.each(['documents', 'settings', 'history'])( + 'replaces the removed connection with Sources from the %s view', + async (view) => { + await render(`?view=${view}`) + const options: ConnectorActionsOptions = mocks.actions.mock.lastCall![0] + act(() => options.onRemoved?.()) + expect(mocks.replace).toHaveBeenCalledWith('/o/org-one/settings/integrations') + expect(mocks.push).not.toHaveBeenCalled() + } + ) + it('opens documents by default and uses the exact canonical search index', async () => { await render() expect(mocks.detail).toHaveBeenLastCalledWith('index-one', 'source-one') diff --git a/apps/sim/app/o/[organizationId]/settings/integrations/sources/[connectorId]/source-detail.tsx b/apps/sim/app/o/[organizationId]/settings/integrations/sources/[connectorId]/source-detail.tsx index 939dfd52563..df92230fa62 100644 --- a/apps/sim/app/o/[organizationId]/settings/integrations/sources/[connectorId]/source-detail.tsx +++ b/apps/sim/app/o/[organizationId]/settings/integrations/sources/[connectorId]/source-detail.tsx @@ -197,6 +197,8 @@ function SourceDetailContent({ const description = [title === meta?.name ? undefined : meta?.name, status].filter(Boolean).join(' · ') || undefined const onBack = () => router.push(backHref) + const onRemoved = () => + router.replace(organizationRoutes(organization.id).settingsSection('integrations')) const onViewChange = (value: string) => { const next = sourceViewParam.parser.parse(value) if (next) void setView(next) @@ -254,6 +256,7 @@ function SourceDetailContent({ queryError={integrationFeedback} backText={backText} onBack={onBack} + onRemoved={onRemoved} onViewChange={onViewChange} /> ) @@ -264,7 +267,7 @@ function SourceDetailContent({ title={title} description={description} docsLink={meta?.searchDocsUrl} - onRemoved={onBack} + onRemoved={onRemoved} > {integrationFeedback} @@ -367,6 +370,7 @@ interface SourceSettingsEditorProps { queryError?: ReactNode backText: string onBack: () => void + onRemoved: () => void onViewChange: (view: string) => void } @@ -400,6 +404,7 @@ function SourceSettingsForm({ queryError, backText, onBack, + onRemoved, onViewChange, onSaved, onDiscard, @@ -420,7 +425,7 @@ function SourceSettingsForm({ description={description} docsLink={form.docsUrl} lifecycleDisabled={form.dirty || form.saving} - onRemoved={onBack} + onRemoved={onRemoved} actions={saveDiscardActions({ dirty: form.dirty, saving: form.saving, diff --git a/apps/sim/app/workspace/[workspaceId]/knowledge/[id]/components/connectors-section/connectors-section.test.tsx b/apps/sim/app/workspace/[workspaceId]/knowledge/[id]/components/connectors-section/connectors-section.test.tsx index b8559b55027..c8574915cbe 100644 --- a/apps/sim/app/workspace/[workspaceId]/knowledge/[id]/components/connectors-section/connectors-section.test.tsx +++ b/apps/sim/app/workspace/[workspaceId]/knowledge/[id]/components/connectors-section/connectors-section.test.tsx @@ -36,6 +36,7 @@ const { isFetching: false, }, lifecycle: { + removeOptions: { onSuccess: undefined as (() => void) | undefined }, sync: { mutate: vi.fn(), reset: vi.fn(), error: null as Error | null, isPending: false }, update: { mutate: vi.fn(), reset: vi.fn(), error: null as Error | null, isPending: false }, remove: { mutate: vi.fn(), reset: vi.fn(), error: null as Error | null, isPending: false }, @@ -235,7 +236,10 @@ vi.mock('@/hooks/queries/kb/connectors', () => ({ isPlaceholderData: lifecycle.detail.isPlaceholderData, refetch: lifecycle.detail.refetch, })), - useDeleteConnector: () => lifecycle.remove, + useDeleteConnector: (options: { onSuccess: () => void }) => { + lifecycle.removeOptions = options + return lifecycle.remove + }, useTriggerSync: () => lifecycle.sync, useUpdateConnector: () => lifecycle.update, })) @@ -943,15 +947,12 @@ describe('shared connector lifecycle actions', () => { expect(dialog.textContent).not.toContain('remain unless') } act(() => findButton(dialog, 'Remove').click()) - expect(lifecycle.remove.mutate).toHaveBeenCalledWith( - { - knowledgeBaseId: 'knowledge-1', - connectorId: 'connector-1', - deleteDocuments: accessMode !== 'workspace', - }, - expect.any(Object) - ) - act(() => lifecycle.remove.mutate.mock.calls[0][1].onSuccess()) + expect(lifecycle.remove.mutate).toHaveBeenCalledWith({ + knowledgeBaseId: 'knowledge-1', + connectorId: 'connector-1', + deleteDocuments: accessMode !== 'workspace', + }) + act(() => lifecycle.removeOptions.onSuccess?.()) expect(onRemoved).toHaveBeenCalledOnce() expect(container.querySelector('[role="dialog"]')).toBeNull() } diff --git a/apps/sim/app/workspace/[workspaceId]/knowledge/[id]/components/connectors-section/use-connector-actions.ts b/apps/sim/app/workspace/[workspaceId]/knowledge/[id]/components/connectors-section/use-connector-actions.ts index 9fbbcda5690..c413b5733b4 100644 --- a/apps/sim/app/workspace/[workspaceId]/knowledge/[id]/components/connectors-section/use-connector-actions.ts +++ b/apps/sim/app/workspace/[workspaceId]/knowledge/[id]/components/connectors-section/use-connector-actions.ts @@ -31,9 +31,15 @@ export function useConnectorActions({ }: ConnectorActionsOptions) { const sync = useTriggerSync() const update = useUpdateConnector() - const remove = useDeleteConnector() const [confirmRemove, setConfirmRemove] = useState(false) const [deleteDocuments, setDeleteDocuments] = useState(false) + const remove = useDeleteConnector({ + onSuccess: () => { + setConfirmRemove(false) + setDeleteDocuments(false) + onRemoved?.() + }, + }) const requiresDocumentDeletion = connector.accessMode !== 'workspace' const state = getConnectorSyncState(connector) const actionsDisabled = disabled || sync.isPending || update.isPending || remove.isPending @@ -117,20 +123,11 @@ export function useConnectorActions({ error: remove.error, onConfirm: () => { if (!canEdit || actionsDisabled) return - remove.mutate( - { - knowledgeBaseId, - connectorId: connector.id, - deleteDocuments: requiresDocumentDeletion || deleteDocuments, - }, - { - onSuccess: () => { - setConfirmRemove(false) - setDeleteDocuments(false) - onRemoved?.() - }, - } - ) + remove.mutate({ + knowledgeBaseId, + connectorId: connector.id, + deleteDocuments: requiresDocumentDeletion || deleteDocuments, + }) }, }, } diff --git a/apps/sim/hooks/queries/kb/connectors-cache.test.tsx b/apps/sim/hooks/queries/kb/connectors-cache.test.tsx index 89ee868601c..72af149fd60 100644 --- a/apps/sim/hooks/queries/kb/connectors-cache.test.tsx +++ b/apps/sim/hooks/queries/kb/connectors-cache.test.tsx @@ -388,6 +388,46 @@ describe('connector Search result cache reconciliation', () => { }) describe('Search source list reconciliation', () => { + it('runs removal navigation before refetches and retains it after the caller unmounts', async () => { + const client = createQueryClient() + const request = Promise.withResolvers() + mocks.requestJson.mockReturnValueOnce(request.promise) + const invalidated = vi.spyOn(client, 'invalidateQueries') + const onSuccess = vi.fn(() => expect(invalidated).not.toHaveBeenCalled()) + const mutation = renderMutation(client, () => useDeleteConnector({ onSuccess })) + let done!: Promise + await act(async () => { + done = mutation().mutateAsync({ + knowledgeBaseId: KNOWLEDGE_BASE_ID, + connectorId: CONNECTOR_ID, + deleteDocuments: true, + }) + }) + act(() => mountedRoots.pop()!.unmount()) + request.resolve({ success: true }) + await act(async () => { + await done + }) + expect(onSuccess).toHaveBeenCalledOnce() + expect(invalidated).toHaveBeenCalledWith({ + queryKey: connectorKeys.detail(KNOWLEDGE_BASE_ID, CONNECTOR_ID), + refetchType: 'none', + }) + }) + + it('does not navigate when removal fails', async () => { + const client = createQueryClient() + const onSuccess = vi.fn() + mocks.requestJson.mockRejectedValueOnce(new Error('Removal failed')) + const mutation = renderMutation(client, () => useDeleteConnector({ onSuccess })) + await act(async () => { + await expect( + mutation().mutateAsync({ knowledgeBaseId: KNOWLEDGE_BASE_ID, connectorId: CONNECTOR_ID }) + ).rejects.toThrow('Removal failed') + }) + expect(onSuccess).not.toHaveBeenCalled() + }) + it('refreshes summaries after editing source configuration or pausing sync', async () => { const queryClient = createQueryClient() const mutation = renderMutation(queryClient, useUpdateConnector) diff --git a/apps/sim/hooks/queries/kb/connectors.ts b/apps/sim/hooks/queries/kb/connectors.ts index 66d5fee637c..fa507b4d24f 100644 --- a/apps/sim/hooks/queries/kb/connectors.ts +++ b/apps/sim/hooks/queries/kb/connectors.ts @@ -657,19 +657,37 @@ async function deleteConnector({ }) } -export function useDeleteConnector() { +interface UseDeleteConnectorOptions { + onSuccess?: () => void +} + +export function useDeleteConnector(options?: UseDeleteConnectorOptions) { const queryClient = useQueryClient() return useMutation({ mutationFn: deleteConnector, + /** Run before invalidation can unmount the source page on a 404 response. */ + onSuccess: () => options?.onSuccess?.(), /** * Removing a connector can take its documents with it, so the document * lists and the base's own totals move — but nothing below them does. * Invalidating `knowledgeKeys.detail` as a prefix would also refetch every * cached document detail, chunk page, and chunk search in the base. */ - onSettled: (_data, _error, { knowledgeBaseId, deleteDocuments }) => { - queryClient.invalidateQueries({ queryKey: connectorKeys.all(knowledgeBaseId) }) + onSettled: (_data, error, { knowledgeBaseId, connectorId, deleteDocuments }) => { + if (error) { + queryClient.invalidateQueries({ queryKey: connectorKeys.all(knowledgeBaseId) }) + } else { + queryClient.invalidateQueries({ queryKey: connectorKeys.lists(knowledgeBaseId) }) + /** Retire stale detail pages without fetching the just-deleted resource during navigation. */ + void queryClient.cancelQueries({ + queryKey: connectorKeys.detail(knowledgeBaseId, connectorId), + }) + queryClient.invalidateQueries({ + queryKey: connectorKeys.detail(knowledgeBaseId, connectorId), + refetchType: 'none', + }) + } queryClient.invalidateQueries({ queryKey: searchSourceKeys.lists() }) queryClient.invalidateQueries({ queryKey: searchIntegrationKeys.lists() }) queryClient.invalidateQueries({ queryKey: knowledgeKeys.documentLists(knowledgeBaseId) }) diff --git a/apps/sim/lib/knowledge/__integration__/storage-accounting.integration.ts b/apps/sim/lib/knowledge/__integration__/storage-accounting.integration.ts index 8103bfb464a..1d678bd8dc8 100644 --- a/apps/sim/lib/knowledge/__integration__/storage-accounting.integration.ts +++ b/apps/sim/lib/knowledge/__integration__/storage-accounting.integration.ts @@ -5,20 +5,33 @@ import { promisify } from 'node:util' import { db } from '@sim/db' import { document, + embedding, knowledgeBase, knowledgeConnector, organization, + outboxEvent, user, workspace, } from '@sim/db/schema' import { generateId } from '@sim/utils/id' import { and, eq, inArray, isNull, sql } from 'drizzle-orm' import { afterAll, describe, expect, it, vi } from 'vitest' +import { processOutboxEventById } from '@/lib/core/outbox/service' import { createKnowledgeAclFixtureIds, seedKnowledgeAclFixture, } from '@/lib/knowledge/__integration__/seed-source-access-fixture' -import { createSingleDocument, hardDeleteDocuments } from '@/lib/knowledge/documents/service' +import { WORKSPACE_ACCESS_SCOPE } from '@/lib/knowledge/access/scope' +import { SYSTEM_ACCESS_SCOPE } from '@/lib/knowledge/access/types' +import { KNOWLEDGE_CONNECTOR_CLEANUP_EVENT } from '@/lib/knowledge/connectors/deletion' +import { createContentSyncLease, SyncLockLostException } from '@/lib/knowledge/connectors/sync-lock' +import { persistSkippedDocuments } from '@/lib/knowledge/connectors/sync-persistence' +import { knowledgeDocumentProcessingOutboxHandlers } from '@/lib/knowledge/documents/processing-outbox-handler' +import { + createSingleDocument, + getKnowledgeDocument, + hardDeleteDocuments, +} from '@/lib/knowledge/documents/service' import { performDeleteKnowledgeConnector } from '@/lib/knowledge/orchestration/connectors' type Fixture = ReturnType @@ -93,6 +106,14 @@ function disconnect(ids: Fixture, deleteDocuments = false) { afterAll(async () => { for (const ids of fixtures) { + await db + .delete(outboxEvent) + .where( + and( + eq(outboxEvent.eventType, KNOWLEDGE_CONNECTOR_CLEANUP_EVENT), + sql`${outboxEvent.payload}->>'knowledgeBaseId' = ${ids.knowledgeBaseId}` + ) + ) await db.delete(knowledgeBase).where(eq(knowledgeBase.id, ids.knowledgeBaseId)) await db.delete(workspace).where(eq(workspace.id, ids.workspaceId)) await db.delete(organization).where(eq(organization.id, ids.organizationId)) @@ -176,19 +197,100 @@ describe('knowledge document storage ledgers', () => { expect(await ledger(ids)).toEqual({ workspaceBytes: 0, payerBytes: 0 }) }) - it('deletes a paginated source including archived documents without debiting manual storage', async () => { + it('hides a source immediately and cleans bounded batches without debiting manual storage', async () => { const ids = await seed() await manualDocument(ids, 31) const rows = Array.from({ length: 501 }, (_, index) => sourceDocument(ids, index + 1, index % 2 ? { archivedAt: new Date() } : {}) ) await db.insert(document).values(rows) + await db.insert(embedding).values( + Array.from({ length: 1_001 }, (_, chunkIndex) => ({ + id: generateId(), + knowledgeBaseId: ids.knowledgeBaseId, + documentId: rows[0].id, + chunkIndex, + chunkHash: `hash-${chunkIndex}`, + content: 'test chunk', + contentLength: 10, + tokenCount: 2, + startOffset: 0, + endOffset: 10, + embedding384: Array(384).fill(0.1), + })) + ) expect(await disconnect(ids, true)).toEqual({ success: true, documentsKept: 0, documentsDeleted: 501, }) + const [tombstone] = await db + .select() + .from(knowledgeConnector) + .where(eq(knowledgeConnector.id, ids.connectorId)) + expect(tombstone.deletedAt).toBeInstanceOf(Date) + expect(tombstone.status).toBe('disabled') + expect(tombstone.syncLockToken).toBeNull() + const [retained] = await db + .select({ count: sql`COUNT(*)::integer` }) + .from(document) + .where(eq(document.connectorId, ids.connectorId)) + expect(retained.count).toBe(501) + for (const access of [SYSTEM_ACCESS_SCOPE, WORKSPACE_ACCESS_SCOPE]) { + expect(await getKnowledgeDocument(ids.knowledgeBaseId, rows[0].id, access)).toBeNull() + } + const [job] = await db + .select() + .from(outboxEvent) + .where( + and( + eq(outboxEvent.eventType, KNOWLEDGE_CONNECTOR_CLEANUP_EVENT), + sql`${outboxEvent.payload}->>'connectorId' = ${ids.connectorId}` + ) + ) + .limit(1) + expect(job).toBeDefined() + await expect( + persistSkippedDocuments( + ids.knowledgeBaseId, + ids.connectorId, + 'confluence', + [ + { + type: 'skip', + extDoc: { + externalId: generateId(), + title: 'Late source write', + content: '', + mimeType: 'text/plain', + contentHash: 'late-source', + skippedReason: 'Too large', + }, + }, + ], + undefined, + 'workspace', + createContentSyncLease(ids.connectorId, ids.lockId) + ) + ).rejects.toBeInstanceOf(SyncLockLostException) + const handlers = knowledgeDocumentProcessingOutboxHandlers + let status = await processOutboxEventById(job.id, handlers) + expect(status).toBe('pending') + for (let attempt = 0; status === 'pending' && attempt < 5; attempt++) { + await db + .update(outboxEvent) + .set({ availableAt: new Date() }) + .where(eq(outboxEvent.id, job.id)) + status = await processOutboxEventById(job.id, handlers) + } + expect(status).toBe('completed') + expect( + await db + .select({ id: knowledgeConnector.id }) + .from(knowledgeConnector) + .where(eq(knowledgeConnector.id, ids.connectorId)) + ).toHaveLength(0) const [remaining] = await db .select({ count: sql`COUNT(*)::integer` }) .from(document) @@ -197,6 +299,42 @@ describe('knowledge document storage ledgers', () => { expect(await ledger(ids)).toEqual({ workspaceBytes: 31, payerBytes: 31 }) }) + it('rolls back both the source tombstone and cleanup intent on a failed commit', async () => { + const ids = await seed() + const row = sourceDocument(ids, 10) + await db.insert(document).values(row) + const transaction = db.transaction.bind(db) + const failure = vi.spyOn(db, 'transaction').mockImplementationOnce((callback, config) => + transaction(async (tx) => { + await callback(tx) + throw new Error('Removal transaction failed') + }, config) + ) + try { + expect(await disconnect(ids, true)).toMatchObject({ success: false }) + } finally { + failure.mockRestore() + } + const [source] = await db + .select({ deletedAt: knowledgeConnector.deletedAt }) + .from(knowledgeConnector) + .where(eq(knowledgeConnector.id, ids.connectorId)) + expect(source.deletedAt).toBeNull() + expect( + await getKnowledgeDocument(ids.knowledgeBaseId, row.id, WORKSPACE_ACCESS_SCOPE) + ).not.toBeNull() + const events = await db + .select({ id: outboxEvent.id }) + .from(outboxEvent) + .where( + and( + eq(outboxEvent.eventType, KNOWLEDGE_CONNECTOR_CLEANUP_EVENT), + sql`${outboxEvent.payload}->>'connectorId' = ${ids.connectorId}` + ) + ) + expect(events).toHaveLength(0) + }) + it('keeps the source and every document attached when detachment exceeds the quota', async () => { const ids = await seed() const rows = [sourceDocument(ids, 800_000_000), sourceDocument(ids, 800_000_000)] diff --git a/apps/sim/lib/knowledge/access/predicate.test.ts b/apps/sim/lib/knowledge/access/predicate.test.ts index 2ff74b38a51..11582cc360e 100644 --- a/apps/sim/lib/knowledge/access/predicate.test.ts +++ b/apps/sim/lib/knowledge/access/predicate.test.ts @@ -77,7 +77,11 @@ describe('knowledgeAccessCondition', () => { ) }) - it('exempts only the branded system scope', () => { - expect(render(knowledgeAccessCondition(SYSTEM_ACCESS_SCOPE)).sql).toBe('true') + it('exempts system jobs from ACL checks while refusing removed sources', () => { + const { sql } = render(knowledgeAccessCondition(SYSTEM_ACCESS_SCOPE)) + expect(sql).toContain('"document"."connector_id" IS NULL OR EXISTS') + expect(sql).toContain('"knowledge_connector"."deleted_at" IS NULL') + expect(sql).toContain('"knowledge_connector"."archived_at" IS NULL') + expect(sql).not.toContain('"document"."acl"') }) }) diff --git a/apps/sim/lib/knowledge/access/predicate.ts b/apps/sim/lib/knowledge/access/predicate.ts index 21a653271fb..bbce2808f25 100644 --- a/apps/sim/lib/knowledge/access/predicate.ts +++ b/apps/sim/lib/knowledge/access/predicate.ts @@ -16,6 +16,7 @@ import { type SQL, sql } from 'drizzle-orm' import { EXTERNAL_GROUP_STALE_AFTER_MS } from '@/lib/knowledge/access/external-groups' import { SOURCE_ACL_MAX_AGE_MS } from '@/lib/knowledge/access/freshness' import type { KnowledgeAccessScope, SystemAccessScope } from '@/lib/knowledge/access/types' +import { documentConnectorIsActive } from '@/lib/knowledge/documents/connector-lifecycle' import { searchIntegrationAccessCondition } from '@/lib/knowledge/search/integration-policy' import { GITHUB_INSTALLATION_PROVIDER_ID } from '@/lib/oauth/github-installation-types' import { ATLASSIAN_SERVICE_ACCOUNT_PROVIDER_ID } from '@/lib/oauth/types' @@ -202,7 +203,7 @@ function storedKnowledgeAccessCondition( scope: KnowledgeAccessScope | SystemAccessScope, liveSourceAccess: SQL ): SQL { - if (scope.kind === 'system') return sql`true` + if (scope.kind === 'system') return documentConnectorIsActive() if (scope.tokens.length === 0) return sql`false` const tokens = textArrayLiteral(scope.tokens) const cutoff = sql`statement_timestamp() - (${SOURCE_ACL_MAX_AGE_MS} * interval '1 millisecond')` @@ -217,6 +218,8 @@ function storedKnowledgeAccessCondition( OR EXISTS ( SELECT 1 FROM ${knowledgeConnector} WHERE ${knowledgeConnector.id} = ${document.connectorId} + AND ${knowledgeConnector.deletedAt} IS NULL + AND ${knowledgeConnector.archivedAt} IS NULL AND ${knowledgeConnector.accessRewritePending} = false AND ${searchIntegrationAccessCondition()} AND ${liveSourceAccess} diff --git a/apps/sim/lib/knowledge/access/types.ts b/apps/sim/lib/knowledge/access/types.ts index d4db3abd081..1ad91607e41 100644 --- a/apps/sim/lib/knowledge/access/types.ts +++ b/apps/sim/lib/knowledge/access/types.ts @@ -110,8 +110,9 @@ export const MAX_KNOWLEDGE_ACCESS_CANDIDATES = 400 declare const systemAccessScopeBrand: unique symbol /** - * The one exemption from access filtering: a background job acting on rows it - * owns (document processing, connector sync). It is a branded type so it cannot + * The exemption from ACL filtering for a background job acting on rows it + * owns (document processing, connector sync). Removed sources remain inaccessible. + * It is a branded type so it cannot * be assembled from a literal, and this module is its only source, so every * caller is one grep away. Never construct it on a request path. */ diff --git a/apps/sim/lib/knowledge/connectors/deletion.md b/apps/sim/lib/knowledge/connectors/deletion.md new file mode 100644 index 00000000000..826ea5dbecf --- /dev/null +++ b/apps/sim/lib/knowledge/connectors/deletion.md @@ -0,0 +1,13 @@ +# Connection removal + +All existing UI, API, and Copilot callers still enter `knowledge.connectors.delete` through its authorized application use case. Roles, scope checks, response shapes, audit attribution, and the workspace-only option to retain documents are unchanged. + +When documents are removed, a transaction locks the canonical knowledge base and connector, counts the attached documents, marks the connector deleted, invalidates both sync leases, disables scheduling, and inserts one `knowledge.connector.cleanup` outbox event. Failure rolls back both the deletion and the event. Returned deletion counts describe documents logically removed from Sim; physical deletion follows asynchronously. The keep-documents path still performs its storage quota check and detachment atomically. + +The shared document access predicate excludes deleted connectors, including public and workspace access and internal indexing reads. Ingestion also checks connector liveness before claiming or committing document processing. Existing source-write leases reject late sync writes. Documents keep their connector reference until physical deletion, so they cannot become standalone readable or billable documents during cleanup. Restoring a knowledge base does not restore a directly deleted connector. + +The existing outbox worker runs cleanup, with 48 failure attempts and bounded continuations that do not consume that retry budget. Each transaction removes at most 1,000 chunks, 250 documents, or 1,000 sync-history/member rows. A run does at most four batches and yields after its time budget. Transactions use lock and statement timeouts. The worker verifies the connector's deletion timestamp, locks documents against late indexing commits, and commits storage cleanup intents before deleting those documents. It resolves storage ownership from the currently locked knowledge base. Already committed batches survive worker restarts. Failures remain retryable; exhausted events remain as dead letters while the source stays inaccessible. + +After document and connector cleanup, credential grant revocation and unused-tag cleanup are retried as needed. Tag cleanup checks for existence rather than counting the entire remaining corpus. No provider credentials or document contents enter the connector cleanup payload. + +The mutation's success handler navigates before cache invalidation and survives the source component unmounting. Source detail replaces its history entry with the Sources page from Documents, Settings, or Sync history; failed removal stays on the current page with the error. diff --git a/apps/sim/lib/knowledge/connectors/deletion.test.ts b/apps/sim/lib/knowledge/connectors/deletion.test.ts new file mode 100644 index 00000000000..d885a707a0b --- /dev/null +++ b/apps/sim/lib/knowledge/connectors/deletion.test.ts @@ -0,0 +1,189 @@ +/** @vitest-environment node */ +import { db } from '@sim/db' +import { + document, + embedding, + knowledgeBase, + knowledgeConnector, + knowledgeConnectorMember, + knowledgeConnectorMemberSyncLog, + knowledgeConnectorSyncLog, +} from '@sim/db/schema' +import { dbChainMockFns, queueTableRows, resetDbChainMock } from '@sim/testing' +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' +import type { OutboxEventContext } from '@/lib/core/outbox/service' + +const mocks = vi.hoisted(() => ({ storage: vi.fn(), tags: vi.fn(), revoke: vi.fn() })) +vi.mock('@/lib/knowledge/documents/storage-cleanup', () => ({ + enqueueKnowledgeStorageCleanup: mocks.storage, +})) +vi.mock('@/lib/knowledge/tags/service', () => ({ cleanupUnusedTagDefinitions: mocks.tags })) +vi.mock('@/lib/knowledge/connectors/member-access', () => ({ + revokeKnowledgeConnectorCredentialAccess: mocks.revoke, +})) + +import { + cleanupKnowledgeConnector, + enqueueConnectorDeletion, + KNOWLEDGE_CONNECTOR_CLEANUP_EVENT, +} from '@/lib/knowledge/connectors/deletion' + +const payload = { + version: 1 as const, + knowledgeBaseId: 'kb-1', + connectorId: 'connector-1', + deletedAt: '2026-09-15T12:00:00.000Z', + credentialAccess: { workspaceId: 'ws-1', credentialGroupId: 'group-1', actorUserId: 'user-1' }, +} +const owner = { workspaceId: 'ws-1', organizationId: null, userId: 'user-1' } + +function context(): OutboxEventContext { + return { + eventId: 'event-1', + eventType: KNOWLEDGE_CONNECTOR_CLEANUP_EVENT, + attempts: 0, + maxAttempts: 48, + signal: new AbortController().signal, + checkpointPayload: vi.fn(), + } +} + +function queueBatch(docs: { id: string; fileUrl: string }[], chunks: { id: string }[] = []) { + queueTableRows(knowledgeBase, [owner]) + queueTableRows(knowledgeConnector, [{ deletedAt: new Date(payload.deletedAt) }]) + queueTableRows(document, docs) + if (docs.length) queueTableRows(embedding, chunks) +} + +describe('durable connector cleanup', () => { + beforeEach(() => { + vi.clearAllMocks() + resetDbChainMock() + mocks.storage.mockResolvedValue([]) + mocks.tags.mockResolvedValue(0) + mocks.revoke.mockResolvedValue(undefined) + }) + afterEach(resetDbChainMock) + + it('enqueues a bounded immutable identity with a retry budget', async () => { + await enqueueConnectorDeletion(db, payload) + expect(dbChainMockFns.values).toHaveBeenCalledWith( + expect.objectContaining({ + eventType: KNOWLEDGE_CONNECTOR_CLEANUP_EVENT, + payload, + maxAttempts: 48, + }) + ) + }) + + it('preserves storage cleanup intent before removing documents and then the connector', async () => { + const docs = [{ id: 'doc-1', fileUrl: '/file.txt' }] + queueBatch(docs) + queueBatch([]) + await cleanupKnowledgeConnector(payload, context()) + expect(mocks.storage).toHaveBeenCalledWith( + expect.anything(), + [{ ...docs[0], ...owner }], + 'event-1' + ) + expect(mocks.storage.mock.invocationCallOrder[0]).toBeLessThan( + dbChainMockFns.delete.mock.invocationCallOrder[0] + ) + expect(dbChainMockFns.delete.mock.calls.map(([table]) => table)).toEqual([ + document, + knowledgeConnector, + ]) + expect(mocks.revoke).toHaveBeenCalledWith( + { + workspaceId: 'ws-1', + credentialGroupId: 'group-1', + connectorId: 'connector-1', + }, + 'user-1' + ) + expect(mocks.tags).toHaveBeenCalledOnce() + }) + + it('yields after four bounded chunk batches without spending the failure retry budget', async () => { + const docs = Array.from({ length: 250 }, (_, index) => ({ id: `doc-${index}`, fileUrl: '' })) + const chunks = Array.from({ length: 1000 }, (_, index) => ({ id: `chunk-${index}` })) + for (let batch = 0; batch < 4; batch++) queueBatch(docs, chunks) + expect(await cleanupKnowledgeConnector(payload, context())).toMatchObject({ + outcome: 'deferred', + consumeAttempt: false, + }) + expect(dbChainMockFns.transaction).toHaveBeenCalledTimes(4) + expect(dbChainMockFns.limit).toHaveBeenCalledWith(250) + expect(dbChainMockFns.limit).toHaveBeenCalledWith(1000) + expect(dbChainMockFns.delete.mock.calls.map(([table]) => table)).toEqual( + Array(4).fill(embedding) + ) + expect(mocks.storage).not.toHaveBeenCalled() + expect(mocks.revoke).not.toHaveBeenCalled() + }) + + it('keeps the documents when storage cleanup intent cannot be persisted', async () => { + queueBatch([{ id: 'doc-1', fileUrl: '/file.txt' }]) + mocks.storage.mockRejectedValueOnce(new Error('Outbox unavailable')) + await expect(cleanupKnowledgeConnector(payload, context())).rejects.toThrow( + 'Outbox unavailable' + ) + expect(dbChainMockFns.delete).not.toHaveBeenCalled() + expect(mocks.tags).not.toHaveBeenCalled() + }) + + it.each([knowledgeConnectorSyncLog, knowledgeConnectorMemberSyncLog, knowledgeConnectorMember])( + 'drains related rows before deleting their connector', + async (table) => { + const rows = Array.from({ length: 1000 }, (_, index) => ({ id: `row-${index}` })) + for (let batch = 0; batch < 4; batch++) { + queueBatch([]) + queueTableRows(table, rows) + } + expect(await cleanupKnowledgeConnector(payload, context())).toMatchObject({ + outcome: 'deferred', + consumeAttempt: false, + }) + expect(dbChainMockFns.delete.mock.calls.map(([target]) => target)).toEqual( + Array(4).fill(table) + ) + expect(mocks.revoke).not.toHaveBeenCalled() + expect(mocks.tags).not.toHaveBeenCalled() + } + ) + + it.each([null, new Date('2026-09-14T12:00:00.000Z')])( + 'leaves a connector with a different deletion generation untouched: %s', + async (deletedAt) => { + queueTableRows(knowledgeBase, [owner]) + queueTableRows(knowledgeConnector, [{ deletedAt }]) + await cleanupKnowledgeConnector(payload, context()) + expect(dbChainMockFns.delete).not.toHaveBeenCalled() + expect(mocks.storage).not.toHaveBeenCalled() + expect(mocks.revoke).not.toHaveBeenCalled() + } + ) + + it('retries final effects after the connector deletion already committed', async () => { + queueBatch([]) + mocks.tags.mockRejectedValueOnce(new Error('Temporary tag failure')) + await expect(cleanupKnowledgeConnector(payload, context())).rejects.toThrow( + 'Temporary tag failure' + ) + resetDbChainMock() + queueTableRows(knowledgeBase, [owner]) + queueTableRows(knowledgeConnector, []) + await cleanupKnowledgeConnector(payload, context()) + expect(mocks.tags).toHaveBeenCalledTimes(2) + expect(dbChainMockFns.delete).not.toHaveBeenCalled() + }) + + it('does no work after cancellation', async () => { + const controller = new AbortController() + controller.abort() + await expect( + cleanupKnowledgeConnector(payload, { ...context(), signal: controller.signal }) + ).rejects.toThrow() + expect(dbChainMockFns.transaction).not.toHaveBeenCalled() + }) +}) diff --git a/apps/sim/lib/knowledge/connectors/deletion.ts b/apps/sim/lib/knowledge/connectors/deletion.ts new file mode 100644 index 00000000000..b95d27c6fc7 --- /dev/null +++ b/apps/sim/lib/knowledge/connectors/deletion.ts @@ -0,0 +1,205 @@ +import { db } from '@sim/db' +import { + document, + embedding, + knowledgeBase, + knowledgeConnector, + knowledgeConnectorMember, + knowledgeConnectorMemberSyncLog, + knowledgeConnectorSyncLog, +} from '@sim/db/schema' +import { and, eq, inArray, sql } from 'drizzle-orm' +import { z } from 'zod' +import { + continueOutboxHandler, + enqueueOutboxEvent, + type OutboxHandler, +} from '@/lib/core/outbox/service' +import type { DbOrTx } from '@/lib/db/types' +import { revokeKnowledgeConnectorCredentialAccess } from '@/lib/knowledge/connectors/member-access' +import { enqueueKnowledgeStorageCleanup } from '@/lib/knowledge/documents/storage-cleanup' +import { cleanupUnusedTagDefinitions } from '@/lib/knowledge/tags/service' + +export const KNOWLEDGE_CONNECTOR_CLEANUP_EVENT = 'knowledge.connector.cleanup' +const DOCUMENT_BATCH_SIZE = 250 +const EMBEDDING_BATCH_SIZE = 1_000 +const RELATED_ROW_BATCH_SIZE = 1_000 +const MAX_BATCHES_PER_RUN = 4 +const RUN_BUDGET_MS = 30_000 + +const deletionPayloadSchema = z + .object({ + version: z.literal(1), + knowledgeBaseId: z.string().min(1).max(256), + connectorId: z.string().min(1).max(256), + deletedAt: z.iso.datetime(), + credentialAccess: z + .object({ + workspaceId: z.string().min(1).max(256), + credentialGroupId: z.string().min(1).max(256), + actorUserId: z.string().min(1).max(256), + }) + .optional(), + }) + .strict() + +type ConnectorDeletionPayload = z.infer + +/** The tombstone and its durable cleanup intent must commit together. */ +export async function enqueueConnectorDeletion( + tx: DbOrTx, + payload: Omit +): Promise { + await enqueueOutboxEvent( + tx, + KNOWLEDGE_CONNECTOR_CLEANUP_EVENT, + deletionPayloadSchema.parse({ version: 1, ...payload }), + { maxAttempts: 48 } + ) +} + +/** + * Each transaction releases at most 250 documents or 1,000 chunks. Locks prevent a late + * indexing commit from inserting chunks between the final chunk scan and document deletion. + * The connector stays attached until every document is gone, preserving storage accounting + * and preventing removed content from becoming a standalone workspace document. + */ +export const cleanupKnowledgeConnector: OutboxHandler = async (rawPayload, context) => { + const payload = deletionPayloadSchema.parse(rawPayload) + const deadline = Math.min( + Date.now() + RUN_BUDGET_MS, + context.deadlineAt ?? Number.POSITIVE_INFINITY + ) + for (let batch = 0; batch < MAX_BATCHES_PER_RUN; batch++) { + context.signal.throwIfAborted() + const outcome = await db.transaction(async (tx) => { + await tx.execute(sql`SET LOCAL lock_timeout = '5s'`) + await tx.execute(sql`SET LOCAL statement_timeout = '10s'`) + const [owner] = await tx + .select({ + workspaceId: knowledgeBase.workspaceId, + organizationId: knowledgeBase.organizationId, + userId: knowledgeBase.userId, + }) + .from(knowledgeBase) + .where(eq(knowledgeBase.id, payload.knowledgeBaseId)) + .for('share') + .limit(1) + if (!owner) return 'complete' + const [connector] = await tx + .select({ deletedAt: knowledgeConnector.deletedAt }) + .from(knowledgeConnector) + .where( + and( + eq(knowledgeConnector.id, payload.connectorId), + eq(knowledgeConnector.knowledgeBaseId, payload.knowledgeBaseId) + ) + ) + .for('update') + .limit(1) + if (!connector) return 'complete' + if (connector.deletedAt?.toISOString() !== payload.deletedAt) return 'obsolete' + + /** Drain the indexed connector bucket; sorting the whole remaining corpus on every batch is unnecessary. */ + const docs = await tx + .select({ id: document.id, fileUrl: document.fileUrl }) + .from(document) + .where( + and( + eq(document.connectorId, payload.connectorId), + eq(document.knowledgeBaseId, payload.knowledgeBaseId) + ) + ) + .limit(DOCUMENT_BATCH_SIZE) + .for('update') + context.signal.throwIfAborted() + if (docs.length === 0) { + for (const table of [ + knowledgeConnectorSyncLog, + knowledgeConnectorMemberSyncLog, + knowledgeConnectorMember, + ]) { + const rows = await tx + .select({ id: table.id }) + .from(table) + .where(eq(table.connectorId, payload.connectorId)) + .limit(RELATED_ROW_BATCH_SIZE) + if (rows.length === 0) continue + await tx.delete(table).where( + inArray( + table.id, + rows.map(({ id }) => id) + ) + ) + context.signal.throwIfAborted() + return 'progress' + } + await tx + .delete(knowledgeConnector) + .where( + and( + eq(knowledgeConnector.id, payload.connectorId), + eq(knowledgeConnector.knowledgeBaseId, payload.knowledgeBaseId), + eq(knowledgeConnector.deletedAt, new Date(payload.deletedAt)) + ) + ) + return 'complete' + } + const documentIds = docs.map(({ id }) => id) + const chunks = await tx + .select({ id: embedding.id }) + .from(embedding) + .where(inArray(embedding.documentId, documentIds)) + .limit(EMBEDDING_BATCH_SIZE) + if (chunks.length > 0) { + await tx.delete(embedding).where( + inArray( + embedding.id, + chunks.map(({ id }) => id) + ) + ) + } else { + await enqueueKnowledgeStorageCleanup( + tx, + docs.map((doc) => ({ ...doc, ...owner })), + context.eventId + ) + await tx.delete(document).where(inArray(document.id, documentIds)) + } + context.signal.throwIfAborted() + return 'progress' + }) + if (outcome === 'obsolete') return + if (outcome === 'complete') { + context.signal.throwIfAborted() + if (payload.credentialAccess) { + await revokeKnowledgeConnectorCredentialAccess( + { + workspaceId: payload.credentialAccess.workspaceId, + credentialGroupId: payload.credentialAccess.credentialGroupId, + connectorId: payload.connectorId, + }, + payload.credentialAccess.actorUserId + ) + } + context.signal.throwIfAborted() + await db.transaction(async (tx) => { + await tx.execute(sql`SET LOCAL lock_timeout = '5s'`) + await tx.execute(sql`SET LOCAL statement_timeout = '10s'`) + await tx + .select({ id: knowledgeBase.id }) + .from(knowledgeBase) + .where(eq(knowledgeBase.id, payload.knowledgeBaseId)) + .for('no key update') + await cleanupUnusedTagDefinitions(payload.knowledgeBaseId, context.eventId, { + executor: tx, + signal: context.signal, + }) + context.signal.throwIfAborted() + }) + return + } + if (Date.now() >= deadline) break + } + return continueOutboxHandler('Connector cleanup committed a bounded batch', 1_000) +} diff --git a/apps/sim/lib/knowledge/documents/connector-lifecycle.ts b/apps/sim/lib/knowledge/documents/connector-lifecycle.ts new file mode 100644 index 00000000000..747b8ff9ee7 --- /dev/null +++ b/apps/sim/lib/knowledge/documents/connector-lifecycle.ts @@ -0,0 +1,13 @@ +import { document, knowledgeConnector } from '@sim/db/schema' +import { sql } from 'drizzle-orm' + +/** A removed source hides its documents immediately, before permanent cleanup catches up. */ +export function documentConnectorIsActive() { + return sql`(${document.connectorId} IS NULL OR EXISTS ( + SELECT 1 FROM ${knowledgeConnector} + WHERE ${knowledgeConnector.id} = ${document.connectorId} + AND ${knowledgeConnector.knowledgeBaseId} = ${document.knowledgeBaseId} + AND ${knowledgeConnector.deletedAt} IS NULL + AND ${knowledgeConnector.archivedAt} IS NULL + ))` +} diff --git a/apps/sim/lib/knowledge/documents/processing-outbox-handler.ts b/apps/sim/lib/knowledge/documents/processing-outbox-handler.ts index e1ee59af466..eef0ceb2a25 100644 --- a/apps/sim/lib/knowledge/documents/processing-outbox-handler.ts +++ b/apps/sim/lib/knowledge/documents/processing-outbox-handler.ts @@ -7,6 +7,10 @@ import { } from '@/lib/core/outbox/service' import { isBYOKEmbeddingCredentialRejection, isEmbeddingQuotaExhaustion } from '@/lib/embeddings' import { SYSTEM_ACCESS_SCOPE } from '@/lib/knowledge/access/types' +import { + cleanupKnowledgeConnector, + KNOWLEDGE_CONNECTOR_CLEANUP_EVENT, +} from '@/lib/knowledge/connectors/deletion' import { getOcrRequestRejection, isPermanentDocumentProcessingError, @@ -228,6 +232,7 @@ const KNOWLEDGE_HANDLER_TIMEOUT_MS = Math.min( ) export const knowledgeDocumentProcessingOutboxHandlers = { + [KNOWLEDGE_CONNECTOR_CLEANUP_EVENT]: cleanupKnowledgeConnector, [KNOWLEDGE_STORAGE_CLEANUP_EVENT]: cleanupKnowledgeStorage, [OCR_CHECKPOINT_CLEANUP_OUTBOX_EVENT]: cleanupOcrCheckpoint, [EMBEDDING_CHECKPOINT_CLEANUP_EVENT]: cleanupEmbeddingCheckpoint, diff --git a/apps/sim/lib/knowledge/documents/service.ts b/apps/sim/lib/knowledge/documents/service.ts index 188a5c0308d..4b5d9cf2761 100644 --- a/apps/sim/lib/knowledge/documents/service.ts +++ b/apps/sim/lib/knowledge/documents/service.ts @@ -81,6 +81,7 @@ import { SYSTEM_ACCESS_SCOPE, } from '@/lib/knowledge/access/types' import { assertSyncLeaseHeldInTx, type SyncWriteLease } from '@/lib/knowledge/connectors/sync-lock' +import { documentConnectorIsActive } from '@/lib/knowledge/documents/connector-lifecycle' import { assertDocumentChunkCountWithinLimit, getOcrRequestRejection, @@ -1506,6 +1507,7 @@ export async function processDocumentAsync( eq(document.userExcluded, false), isNull(document.archivedAt), isNull(document.deletedAt), + documentConnectorIsActive(), isNull(knowledgeBase.deletedAt) ) ) @@ -1533,7 +1535,8 @@ export async function processDocumentAsync( ...queueGenerationConditions(attemptContext), eq(document.userExcluded, false), isNull(document.archivedAt), - isNull(document.deletedAt) + isNull(document.deletedAt), + documentConnectorIsActive() ) ) return @@ -1601,7 +1604,8 @@ export async function processDocumentAsync( : queueGenerationConditions(attemptContext)), eq(document.userExcluded, false), isNull(document.archivedAt), - isNull(document.deletedAt) + isNull(document.deletedAt), + documentConnectorIsActive() ) ) .returning({ id: document.id }) @@ -1861,6 +1865,7 @@ export async function processDocumentAsync( eq(document.userExcluded, false), isNull(document.archivedAt), isNull(document.deletedAt), + documentConnectorIsActive(), isNull(knowledgeBase.deletedAt) ) ) @@ -1930,7 +1935,8 @@ export async function processDocumentAsync( ...queueGenerationConditions(attemptContext), eq(document.userExcluded, false), isNull(document.archivedAt), - isNull(document.deletedAt) + isNull(document.deletedAt), + documentConnectorIsActive() ) ) signal.throwIfAborted() @@ -2124,7 +2130,8 @@ export async function processDocumentAsync( ...queueGenerationConditions(attemptContext), eq(document.userExcluded, false), isNull(document.archivedAt), - isNull(document.deletedAt) + isNull(document.deletedAt), + documentConnectorIsActive() ) ) diff --git a/apps/sim/lib/knowledge/orchestration/connectors.test.ts b/apps/sim/lib/knowledge/orchestration/connectors.test.ts index 0f4783471ac..2a08a2bc246 100644 --- a/apps/sim/lib/knowledge/orchestration/connectors.test.ts +++ b/apps/sim/lib/knowledge/orchestration/connectors.test.ts @@ -25,6 +25,7 @@ const { mockResolveStorageBillingContext, mockIncrementStorage, mockNotifyStorage, + mockEnqueueConnectorDeletion, } = vi.hoisted(() => ({ mockCaptureServerEvent: vi.fn(), mockDispatchSync: vi.fn(), @@ -38,6 +39,7 @@ const { mockResolveStorageBillingContext: vi.fn(), mockIncrementStorage: vi.fn(), mockNotifyStorage: vi.fn(), + mockEnqueueConnectorDeletion: vi.fn(), })) vi.mock('@sim/audit', () => ({ @@ -63,6 +65,9 @@ vi.mock('@/lib/billing/storage', () => ({ vi.mock('@/lib/knowledge/documents/storage-cleanup', () => ({ enqueueKnowledgeStorageCleanup: vi.fn().mockResolvedValue(undefined), })) +vi.mock('@/lib/knowledge/connectors/deletion', () => ({ + enqueueConnectorDeletion: mockEnqueueConnectorDeletion, +})) vi.mock('@/lib/knowledge/connectors/queue', () => ({ dispatchSync: mockDispatchSync })) vi.mock('@/lib/knowledge/connectors/member-queue', () => ({ dispatchMemberSync: mockDispatchMemberSync, @@ -296,11 +301,11 @@ const STORAGE_CONTEXT = { customStorageLimitGB: null, } -function queueConnectorDeletionOwnerAndLock(accessMode = 'workspace') { +function queueConnectorDeletionOwnerAndLock(accessMode = 'workspace', credentialGroupId?: string) { const owner = { id: 'kb-1', workspaceId: 'ws-1', organizationId: null, userId: 'user-1' } queueTableRows(schemaMock.knowledgeBase, [owner]) queueTableRows(schemaMock.knowledgeBase, [owner]) - queueTableRows(schemaMock.knowledgeConnector, [{ accessMode }]) + queueTableRows(schemaMock.knowledgeConnector, [{ accessMode, credentialGroupId }]) } describe('performDeleteKnowledgeConnector', () => { @@ -340,11 +345,11 @@ describe('performDeleteKnowledgeConnector', () => { ) }) - it('reports the documents it deleted when asked to delete them', async () => { + it('hides the connector and queues cleanup without deleting documents in the request', async () => { dbChainMockFns.limit.mockResolvedValueOnce([ { id: 'conn-1', connectorType: 'notion', accessMode: 'workspace' }, ]) - queueTableRows(document, [{ id: 'doc-1', fileUrl: '/a.txt' }]) + queueTableRows(document, [{ count: 501 }]) dbChainMockFns.returning.mockResolvedValueOnce([{ id: 'conn-1' }]) const outcome = await performDeleteKnowledgeConnector({ @@ -354,9 +359,67 @@ describe('performDeleteKnowledgeConnector', () => { deleteDocuments: true, }) - expect(outcome).toMatchObject({ success: true, documentsDeleted: 1, documentsKept: 0 }) + expect(outcome).toMatchObject({ success: true, documentsDeleted: 501, documentsKept: 0 }) + expect(dbChainMockFns.delete).not.toHaveBeenCalled() + expect(dbChainMockFns.update).not.toHaveBeenCalledWith(document) + expect(dbChainMockFns.set).toHaveBeenCalledWith( + expect.objectContaining({ + deletedAt: expect.any(Date), + status: 'disabled', + memberSyncStatus: 'disabled', + syncLockToken: null, + memberSyncLockToken: null, + }) + ) + expect(mockEnqueueConnectorDeletion).toHaveBeenCalledWith(expect.anything(), { + knowledgeBaseId: KB.id, + connectorId: 'conn-1', + deletedAt: expect.any(String), + }) }) + it('fails the transaction without auditing success when cleanup cannot be queued', async () => { + dbChainMockFns.limit.mockResolvedValueOnce([ + { id: 'conn-1', connectorType: 'notion', accessMode: 'workspace' }, + ]) + queueTableRows(document, [{ count: 2 }]) + mockEnqueueConnectorDeletion.mockRejectedValueOnce(new Error('Queue unavailable')) + const outcome = await performDeleteKnowledgeConnector({ + ...ACTOR, + knowledgeBase: KB, + connectorId: 'conn-1', + deleteDocuments: true, + }) + expect(outcome).toMatchObject({ success: false }) + expect(mockRecordAudit).not.toHaveBeenCalled() + expect(mockCaptureServerEvent).not.toHaveBeenCalled() + }) + + it.each(['55P03', '57014', '40P01'])( + 'returns a retryable message for transaction contention %s', + async (code) => { + dbChainMockFns.limit.mockResolvedValueOnce([ + { id: 'conn-1', connectorType: 'notion', accessMode: 'workspace' }, + ]) + dbChainMockFns.transaction.mockRejectedValueOnce( + Object.assign(new Error('Database detail'), { code }) + ) + expect( + await performDeleteKnowledgeConnector({ + ...ACTOR, + knowledgeBase: KB, + connectorId: 'conn-1', + deleteDocuments: true, + }) + ).toMatchObject({ + success: false, + errorCode: 'conflict', + error: 'Connection is busy. Try removing it again in a moment.', + }) + expect(mockRecordAudit).not.toHaveBeenCalled() + } + ) + it('reports a missing connector as not found', async () => { dbChainMockFns.limit.mockResolvedValueOnce([]) @@ -1256,9 +1319,9 @@ describe('members-mode connectors', () => { expect(dbChainMockFns.delete).not.toHaveBeenCalled() }) - it('revokes the credential grant once the connector and its documents are gone', async () => { + it('defers credential grant cleanup with the connector deletion', async () => { queueTableRows(schemaMock.knowledgeConnector, [MEMBERS_CONNECTOR]) - queueConnectorDeletionOwnerAndLock('members') + queueConnectorDeletionOwnerAndLock('members', 'group-1') queueTableRows(schemaMock.document, []) dbChainMockFns.returning.mockResolvedValueOnce([{ id: 'c-1' }]) @@ -1270,9 +1333,17 @@ describe('members-mode connectors', () => { }) expect(outcome).toMatchObject({ success: true }) - expect(mockRevoke).toHaveBeenCalledWith( - { workspaceId: 'ws-1', credentialGroupId: 'group-1', connectorId: 'c-1' }, - 'user-1' + expect(mockRevoke).not.toHaveBeenCalled() + expect(mockEnqueueConnectorDeletion).toHaveBeenCalledWith( + expect.anything(), + expect.objectContaining({ + connectorId: 'c-1', + credentialAccess: { + workspaceId: 'ws-1', + credentialGroupId: 'group-1', + actorUserId: ACTOR.userId, + }, + }) ) }) diff --git a/apps/sim/lib/knowledge/orchestration/connectors.ts b/apps/sim/lib/knowledge/orchestration/connectors.ts index e7985025438..f3ee6da4858 100644 --- a/apps/sim/lib/knowledge/orchestration/connectors.ts +++ b/apps/sim/lib/knowledge/orchestration/connectors.ts @@ -3,15 +3,15 @@ import { db } from '@sim/db' import { credentialGroup, document, - embedding, knowledgeBase, knowledgeBaseTagDefinitions, knowledgeConnector, knowledgeConnectorMember, } from '@sim/db/schema' import { createLogger } from '@sim/logger' +import { getPostgresErrorCode } from '@sim/utils/errors' import { generateId } from '@sim/utils/id' -import { and, asc, eq, gt, inArray, isNull, sql } from 'drizzle-orm' +import { and, eq, isNull, sql } from 'drizzle-orm' import { encryptApiKey } from '@/lib/api-key/crypto' import type { BillingAttributionSnapshot } from '@/lib/billing/core/billing-attribution' import { @@ -44,6 +44,7 @@ import { type ConnectorAccessToken, syncContextForToken, } from '@/lib/knowledge/connectors/access-token' +import { enqueueConnectorDeletion } from '@/lib/knowledge/connectors/deletion' import { findListingCapViolation, grantKnowledgeConnectorCredentialAccess, @@ -52,7 +53,6 @@ import { } from '@/lib/knowledge/connectors/member-access' import type { PreparedConnectorPermissions } from '@/lib/knowledge/connectors/permission-config' import { allocateTagSlots } from '@/lib/knowledge/constants' -import { enqueueKnowledgeStorageCleanup } from '@/lib/knowledge/documents/storage-cleanup' import { auditActorFields, classifyKnowledgeFailure, @@ -60,7 +60,7 @@ import { type KnowledgeOperationContext, type KnowledgeOrchestrationResult, } from '@/lib/knowledge/orchestration/shared' -import { cleanupUnusedTagDefinitions, createTagDefinition } from '@/lib/knowledge/tags/service' +import { createTagDefinition } from '@/lib/knowledge/tags/service' import { captureServerEvent } from '@/lib/posthog/server' import { searchSourceIdentity } from '@/lib/sim-search/source-identity' import { getConnectorApiKeyConfig } from '@/connectors/auth' @@ -1077,8 +1077,8 @@ export interface PerformDeleteKnowledgeConnectorParams extends KnowledgeOperatio knowledgeBase: ConnectorKnowledgeBase connectorId: string /** - * Also hard-delete the documents the connector produced. Defaults to keeping - * them, which turns them into ordinary standalone knowledge base entries. + * Immediately hide the connector's documents and durably queue permanent cleanup. + * Defaults to keeping them as ordinary standalone knowledge base entries. */ deleteDocuments?: boolean /** False only when an authorized application use case projects the semantic audit. */ @@ -1094,8 +1094,8 @@ export type PerformDeleteKnowledgeConnectorResult = KnowledgeOrchestrationResult }> /** - * Hard-deletes a connector, either removing the documents it produced or - * releasing them as standalone entries. + * Removes a connector, either tombstoning it for bounded background cleanup or + * releasing its documents as standalone entries before deleting the connector. * * Returns the counts so callers state what happened rather than assert it. The * copilot tool used to reach this through an internal HTTP self-call that sent @@ -1153,6 +1153,8 @@ export async function performDeleteKnowledgeConnector( : undefined docCount = await db.transaction(async (tx) => { + await tx.execute(sql`SET LOCAL lock_timeout = '5s'`) + await tx.execute(sql`SET LOCAL statement_timeout = '10s'`) /** Match source writes and document deletion: parent KB, connector, then storage ledgers. */ const [lockedOwner] = await tx .select({ @@ -1162,7 +1164,7 @@ export async function performDeleteKnowledgeConnector( }) .from(knowledgeBase) .where(and(eq(knowledgeBase.id, kb.id), isNull(knowledgeBase.deletedAt))) - .for('update') + .for(deleteDocuments ? 'share' : 'update') .limit(1) if ( !lockedOwner || @@ -1173,7 +1175,10 @@ export async function performDeleteKnowledgeConnector( throw new OrchestrationError('conflict', 'Knowledge base ownership changed; retry deletion') } const [lockedConnector] = await tx - .select({ accessMode: knowledgeConnector.accessMode }) + .select({ + accessMode: knowledgeConnector.accessMode, + credentialGroupId: knowledgeConnector.credentialGroupId, + }) .from(knowledgeConnector) .where( and( @@ -1193,93 +1198,99 @@ export async function performDeleteKnowledgeConnector( ) } - let count = 0 if (deleteDocuments) { - let afterId: string | undefined - for (;;) { - /** Archived rows also lose their connector FK and must not escape deletion or cleanup. */ - const docs = await tx - .select({ id: document.id, fileUrl: document.fileUrl }) - .from(document) - .where( - and( - eq(document.connectorId, connectorId), - eq(document.knowledgeBaseId, kb.id), - afterId ? gt(document.id, afterId) : undefined - ) - ) - .orderBy(asc(document.id)) - .limit(250) - if (docs.length === 0) break - const documentIds = docs.map((doc) => doc.id) - await tx.delete(embedding).where(inArray(embedding.documentId, documentIds)) - await tx.delete(document).where(inArray(document.id, documentIds)) - await enqueueKnowledgeStorageCleanup( - tx, - docs.map((doc) => ({ - ...doc, - workspaceId: owner.workspaceId, - organizationId: owner.organizationId, - userId: owner.userId, - })), - requestId - ) - count += docs.length - afterId = docs.at(-1)?.id - } - } else { - /** Legacy skipped rows used remote size despite retaining no artifact. */ - await tx - .update(document) - .set({ fileSize: 0 }) - .where( - and( - eq(document.connectorId, connectorId), - eq(document.knowledgeBaseId, kb.id), - isNull(document.storageKey), - eq(document.fileUrl, '') - ) - ) - /** - * Connector bytes are unmetered until detachment. Count retained archived files too; - * live tombstones are resurrected below, while archived tombstones remain nonbillable. - */ const [totals] = await tx - .select({ - count: sql`COUNT(*)::integer`, - bytes: sql`COALESCE(SUM(${document.fileSize}::bigint) FILTER ( - WHERE ${document.archivedAt} IS NULL OR ${document.deletedAt} IS NULL - ), 0)::text`, - }) + .select({ count: sql`COUNT(*)::integer` }) .from(document) .where(and(eq(document.connectorId, connectorId), eq(document.knowledgeBaseId, kb.id))) - count = totals?.count ?? 0 - const retainedBytes = Number(totals?.bytes ?? 0) - if (!Number.isSafeInteger(retainedBytes) || retainedBytes < 0) { - throw new Error('Invalid retained connector storage size') - } - if (retainedBytes > 0) { - if (storageContext) { - const updatedUsage = await incrementStorageUsageForBillingContextInTx( - tx, - storageContext, - retainedBytes - ) - if (updatedUsage !== undefined) - storageNotification = { context: storageContext, updatedUsage } - } - } + const deletedAt = new Date() await tx - .update(document) - .set({ deletedAt: null }) + .update(knowledgeConnector) + .set({ + deletedAt, + updatedAt: deletedAt, + status: 'disabled', + memberSyncStatus: 'disabled', + syncLockToken: null, + syncLockLeaseAt: null, + memberSyncLockToken: null, + memberSyncLockLeaseAt: null, + nextSyncAt: null, + nextMemberSyncAt: null, + }) .where( and( - eq(document.connectorId, connectorId), - eq(document.knowledgeBaseId, kb.id), - isNull(document.archivedAt) + eq(knowledgeConnector.id, connectorId), + eq(knowledgeConnector.knowledgeBaseId, kb.id) ) ) + await enqueueConnectorDeletion(tx, { + knowledgeBaseId: kb.id, + connectorId, + deletedAt: deletedAt.toISOString(), + ...(lockedConnector.credentialGroupId && owner.workspaceId + ? { + credentialAccess: { + workspaceId: owner.workspaceId, + credentialGroupId: lockedConnector.credentialGroupId, + actorUserId: params.userId, + }, + } + : {}), + }) + return totals?.count ?? 0 + } + /** Legacy skipped rows used remote size despite retaining no artifact. */ + await tx + .update(document) + .set({ fileSize: 0 }) + .where( + and( + eq(document.connectorId, connectorId), + eq(document.knowledgeBaseId, kb.id), + isNull(document.storageKey), + eq(document.fileUrl, '') + ) + ) + /** + * Connector bytes are unmetered until detachment. Count retained archived files too; + * live tombstones are resurrected below, while archived tombstones remain nonbillable. + */ + const [totals] = await tx + .select({ + count: sql`COUNT(*)::integer`, + bytes: sql`COALESCE(SUM(${document.fileSize}::bigint) FILTER ( + WHERE ${document.archivedAt} IS NULL OR ${document.deletedAt} IS NULL + ), 0)::text`, + }) + .from(document) + .where(and(eq(document.connectorId, connectorId), eq(document.knowledgeBaseId, kb.id))) + const count = totals?.count ?? 0 + const retainedBytes = Number(totals?.bytes ?? 0) + if (!Number.isSafeInteger(retainedBytes) || retainedBytes < 0) { + throw new Error('Invalid retained connector storage size') } + if (retainedBytes > 0) { + if (storageContext) { + const updatedUsage = await incrementStorageUsageForBillingContextInTx( + tx, + storageContext, + retainedBytes + ) + if (updatedUsage !== undefined) + storageNotification = { context: storageContext, updatedUsage } + } + } + await tx + .update(document) + .set({ deletedAt: null }) + .where( + and( + eq(document.connectorId, connectorId), + eq(document.knowledgeBaseId, kb.id), + isNull(document.archivedAt) + ) + ) const deletedConnectors = await tx .delete(knowledgeConnector) @@ -1298,6 +1309,13 @@ export async function performDeleteKnowledgeConnector( return count }) } catch (error) { + if (['55P03', '57014', '40P01'].includes(getPostgresErrorCode(error) ?? '')) { + logger.warn(`[${requestId}] Connector removal could not acquire or finish its transaction`, { + connectorId, + error, + }) + return fail('Connection is busy. Try removing it again in a moment.', 'conflict') + } return classifyKnowledgeFailure(error, requestId, `Delete connector ${connectorId}`) } @@ -1308,15 +1326,7 @@ export async function performDeleteKnowledgeConnector( ) } - if (deleteDocuments) { - await Promise.all([ - cleanupUnusedTagDefinitions(kb.id, requestId).catch((error) => { - logger.warn(`[${requestId}] Failed to cleanup tag definitions`, error) - }), - ]) - } - - if (existing.credentialGroupId && kb.workspaceId) { + if (!deleteDocuments && existing.credentialGroupId && kb.workspaceId) { await revokeKnowledgeConnectorCredentialAccess( { workspaceId: kb.workspaceId, diff --git a/apps/sim/lib/knowledge/tags/service.test.ts b/apps/sim/lib/knowledge/tags/service.test.ts index 92e8c31f515..43f36ddcef4 100644 --- a/apps/sim/lib/knowledge/tags/service.test.ts +++ b/apps/sim/lib/knowledge/tags/service.test.ts @@ -2,7 +2,8 @@ * @vitest-environment node */ -import { knowledgeBaseTagDefinitions } from '@sim/db/schema' +import { db } from '@sim/db' +import { document, embedding, knowledgeBaseTagDefinitions } from '@sim/db/schema' import { dbChainMockFns, hasMockCondition, queueTableRows, resetDbChainMock } from '@sim/testing' import { beforeEach, describe, expect, it, vi } from 'vitest' @@ -12,6 +13,7 @@ vi.mock('@sim/utils/id', () => ({ })) import { + cleanupUnusedTagDefinitions, createOrUpdateTagDefinitionsBulk, createTagDefinition, getDocumentTagDefinitions, @@ -34,6 +36,48 @@ function existingDefinition(overrides: Record) { } } +describe('cleanupUnusedTagDefinitions', () => { + beforeEach(() => { + vi.clearAllMocks() + resetDbChainMock() + }) + + it('keeps tags used by either documents or chunks and removes only unused definitions', async () => { + queueTableRows(knowledgeBaseTagDefinitions, [ + existingDefinition({ id: 'document-tag', tagSlot: 'tag1' }), + existingDefinition({ id: 'chunk-tag', tagSlot: 'tag2' }), + existingDefinition({ id: 'unused-tag', tagSlot: 'tag3' }), + ]) + queueTableRows(document, [{ id: 'doc-1' }]) + queueTableRows(document, []) + queueTableRows(embedding, [{ id: 'chunk-1' }]) + queueTableRows(document, []) + queueTableRows(embedding, []) + + expect(await cleanupUnusedTagDefinitions('kb-1', 'request-1')).toBe(1) + expect(dbChainMockFns.delete).toHaveBeenCalledOnce() + expect( + hasMockCondition( + dbChainMockFns.where.mock.calls.at(-1)?.[0], + (node) => node.type === 'eq' && node.right === 'unused-tag' + ) + ).toBe(true) + }) + + it('stops cleanup before deleting tags when its worker is cancelled', async () => { + queueTableRows(knowledgeBaseTagDefinitions, [existingDefinition({})]) + const controller = new AbortController() + controller.abort() + await expect( + cleanupUnusedTagDefinitions('kb-1', 'request-1', { + executor: db, + signal: controller.signal, + }) + ).rejects.toThrow() + expect(dbChainMockFns.delete).not.toHaveBeenCalled() + }) +}) + describe('getDocumentTagDefinitionsByKnowledgeBaseIds', () => { beforeEach(() => { vi.clearAllMocks() diff --git a/apps/sim/lib/knowledge/tags/service.ts b/apps/sim/lib/knowledge/tags/service.ts index 7a599fea365..b414c969f1c 100644 --- a/apps/sim/lib/knowledge/tags/service.ts +++ b/apps/sim/lib/knowledge/tags/service.ts @@ -563,17 +563,20 @@ export async function getTagDefinitionById( */ export async function cleanupUnusedTagDefinitions( knowledgeBaseId: string, - requestId: string + requestId: string, + options?: { executor: DbOrTx; signal: AbortSignal } ): Promise { - const definitions = await getDocumentTagDefinitions(knowledgeBaseId) + const executor = options?.executor ?? db + const definitions = await getDocumentTagDefinitions(knowledgeBaseId, executor) let cleanedUp = 0 for (const def of definitions) { + options?.signal.throwIfAborted() const tagSlot = def.tagSlot validateTagSlot(tagSlot) - const docCountResult = await db - .select({ count: sql`count(*)` }) + const [taggedDocument] = await executor + .select({ id: document.id }) .from(document) .where( and( @@ -583,9 +586,12 @@ export async function cleanupUnusedTagDefinitions( sql`${sql.raw(tagSlot)} IS NOT NULL` ) ) + .limit(1) + if (taggedDocument) continue - const chunkCountResult = await db - .select({ count: sql`count(*)` }) + options?.signal.throwIfAborted() + const [taggedChunk] = await executor + .select({ id: embedding.id }) .from(embedding) .innerJoin(document, eq(embedding.documentId, document.id)) .where( @@ -596,12 +602,13 @@ export async function cleanupUnusedTagDefinitions( sql`${sql.raw(`embedding.${tagSlot}`)} IS NOT NULL` ) ) + .limit(1) - const docCount = Number(docCountResult[0]?.count || 0) - const chunkCount = Number(chunkCountResult[0]?.count || 0) - - if (docCount === 0 && chunkCount === 0) { - await db.delete(knowledgeBaseTagDefinitions).where(eq(knowledgeBaseTagDefinitions.id, def.id)) + if (!taggedChunk) { + options?.signal.throwIfAborted() + await executor + .delete(knowledgeBaseTagDefinitions) + .where(eq(knowledgeBaseTagDefinitions.id, def.id)) cleanedUp++ logger.info(