From 322eef22a22174741363c13b5fbe3c7946ca856b Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Tue, 22 Sep 2026 17:01:03 -0700 Subject: [PATCH] fix(knowledge): revoke absent documents without rewriting unchanged ACLs --- .github/workflows/test-build.yml | 3 +- .../sync-content-pass.postgres.test.ts | 290 ++++++++++++++++++ .../connectors/sync-content-pass.test.ts | 69 ++++- .../knowledge/connectors/sync-content-pass.ts | 45 +-- .../connectors/sync-persistence.test.ts | 80 ++++- .../knowledge/connectors/sync-persistence.ts | 36 ++- 6 files changed, 491 insertions(+), 32 deletions(-) create mode 100644 apps/sim/lib/knowledge/connectors/sync-content-pass.postgres.test.ts diff --git a/.github/workflows/test-build.yml b/.github/workflows/test-build.yml index 9901b6889f9..9fc6012f302 100644 --- a/.github/workflows/test-build.yml +++ b/.github/workflows/test-build.yml @@ -300,7 +300,8 @@ jobs: bunx vitest run \ lib/knowledge/access/predicate.postgres.test.ts \ lib/knowledge/connectors/external-directory.postgres.test.ts \ - lib/knowledge/connectors/sync-persistence.postgres.test.ts + lib/knowledge/connectors/sync-persistence.postgres.test.ts \ + lib/knowledge/connectors/sync-content-pass.postgres.test.ts - name: Verify the projection source and ACL trigger and backfill in PostgreSQL working-directory: packages/db diff --git a/apps/sim/lib/knowledge/connectors/sync-content-pass.postgres.test.ts b/apps/sim/lib/knowledge/connectors/sync-content-pass.postgres.test.ts new file mode 100644 index 00000000000..49b7fdfd92a --- /dev/null +++ b/apps/sim/lib/knowledge/connectors/sync-content-pass.postgres.test.ts @@ -0,0 +1,290 @@ +/** + * @vitest-environment node + */ +import { installProjectionSourceAcl } from '@sim/db/script-migrations/0021_embedding_search_connector' +import { generateId } from '@sim/utils/id' +import postgres, { type Sql } from 'postgres' +import { afterAll, beforeAll, beforeEach, describe, expect, it, vi } from 'vitest' + +const holder = vi.hoisted(() => ({ db: undefined as unknown })) +const hardDelete = vi.hoisted(() => vi.fn(async (ids: string[]) => ids.length)) + +vi.unmock('drizzle-orm') +vi.unmock('@sim/db/schema') +vi.mock('@sim/db', () => ({ + get db() { + return holder.db + }, +})) +vi.mock('@/lib/knowledge/documents/service', () => ({ + hardDeleteDocuments: hardDelete, + isTriggerAvailable: () => true, + processDocumentsWithQueue: vi.fn(), +})) +vi.mock('@/lib/uploads', () => ({ StorageService: {} })) +vi.mock('@/lib/uploads/core/storage-service', () => ({ deleteFile: vi.fn() })) +vi.mock('@/lib/uploads/server/metadata', () => ({ deleteFileMetadata: vi.fn() })) +vi.mock('@/connectors/registry.server', () => ({ CONNECTOR_REGISTRY: {} })) + +const { drizzle } = await import('drizzle-orm/postgres-js') +const { eq, inArray } = await import('drizzle-orm') +const schema = await import('@sim/db/schema') +const { beginListingCheckpoint } = await import('@/lib/knowledge/connectors/listing-checkpoint') +const { runConnectorContentPass } = await import('@/lib/knowledge/connectors/sync-content-pass') +const { revokeDocumentAcls } = await import('@/lib/knowledge/connectors/sync-persistence') +const { confluenceConnector } = await import('@/connectors/confluence/confluence') + +const databaseUrl = process.env.KNOWLEDGE_ACL_TEST_DATABASE_URL + +const ALICE = 'u:alice@corp.com' +const CONNECTOR = 'admin' +const STARTED_AT = new Date('2026-09-08T11:00:00Z') +const LISTED_AT = '2026-09-08 12:00:00' +const ABSENT_SINCE = '2026-09-01 12:00:00' + +/** + * Absence reconciliation of an admin-mode connector against the real projection trigger. The + * document carries only the columns reconciliation reads and writes; a projection row whose + * `acl` is NULL is one the backfill has not filled yet. + */ +describe.runIf(Boolean(databaseUrl))('completed listing reconciliation in PostgreSQL', () => { + let admin: Sql + let sql: Sql + const schemaName = `acl_revoke_${generateId().replaceAll('-', '')}` + + const fannedOut = async () => + ( + await sql<{ document_id: string }[]>`SELECT document_id FROM fan_out ORDER BY document_id` + ).map((row) => row.document_id) + + const insertDocuments = ( + rows: { id: string; acl: string[]; seenAt: string; verified?: boolean }[] + ) => + sql`INSERT INTO document ${sql( + rows.map((row) => ({ + id: row.id, + external_id: row.id, + connector_id: CONNECTOR, + acl: row.acl, + acl_verified_at: row.verified ? '2026-09-01 12:00:00' : null, + source_seen_at: row.seenAt, + })) + )}` + + /** Listed documents keep both deletion guards open: the listing is neither empty nor collapsed. */ + const insertListed = (count: number) => + insertDocuments( + Array.from({ length: count }, (_unused, index) => ({ + id: `listed-${String(index).padStart(3, '0')}`, + acl: [ALICE], + seenAt: LISTED_AT, + })) + ) + + async function reconcile(fullSync = false) { + const result = { + docsAdded: 0, + docsUpdated: 0, + docsDeleted: 0, + docsUnchanged: 0, + docsSkipped: 0, + docsFailed: 0, + } + const pass = await runConnectorContentPass({ + connectorId: CONNECTOR, + connector: { + knowledgeBaseId: 'kb', + connectorType: 'confluence', + listingCheckpoint: { + ...beginListingCheckpoint({ + fingerprint: 'a'.repeat(64), + generationId: 'completed-listing', + startedAt: STARTED_AT, + fullSync, + }), + complete: true, + }, + }, + connectorConfig: confluenceConnector, + sourceConfig: {}, + syncContext: {}, + kbOwner: { workspaceId: 'workspace', userId: 'owner' }, + billingAttribution: { + actorUserId: 'owner', + workspaceId: 'workspace', + organizationId: null, + billedAccountUserId: 'owner', + billingEntity: { type: 'user', id: 'owner' }, + billingPeriod: { start: '2026-09-01T00:00:00Z', end: '2026-10-01T00:00:00Z' }, + payerSubscription: null, + }, + result, + lease: { + stillHeld: () => eq(schema.knowledgeConnector.id, CONNECTOR), + beatIfDue: async () => undefined, + beatLive: async () => undefined, + }, + leaseKind: 'content', + runId: 'run', + fingerprint: 'a'.repeat(64), + documentAccess: 'admin', + getAccessToken: async () => 'token', + hydration: { getDocument: vi.fn() }, + forceRehydrate: false, + fullSync, + deadlineAt: Date.now() + 60_000, + }) + expect(pass.complete).toBe(true) + return { result, notice: pass.holdNotice } + } + + beforeAll(async () => { + const url = new URL(databaseUrl!) + if ( + !['localhost', '127.0.0.1'].includes(url.hostname) || + !url.pathname.startsWith('/sim_acl_test') + ) { + throw new Error('Reconciliation tests require a disposable local integration database') + } + admin = postgres(url.toString(), { max: 1, onnotice: () => undefined }) + await admin.unsafe(`CREATE SCHEMA "${schemaName}"`) + sql = postgres(url.toString(), { + max: 1, + onnotice: () => undefined, + connection: { search_path: schemaName }, + }) + holder.db = drizzle(sql, { schema }) + await sql`CREATE TABLE knowledge_connector (id text PRIMARY KEY)` + await sql`CREATE TABLE document ( + id text PRIMARY KEY, external_id text, connector_id text, + user_excluded boolean NOT NULL DEFAULT false, archived_at timestamp, deleted_at timestamp, + source_seen_at timestamp, + acl text[] NOT NULL DEFAULT '{ws}', acl_requirements jsonb NOT NULL DEFAULT '[]', + acl_verified_at timestamp + )` + for (const projection of ['embedding_search', 'embedding_keyword_tin']) { + await sql`CREATE TABLE ${sql(projection)} ( + id text PRIMARY KEY, document_id text NOT NULL, enabled boolean NOT NULL DEFAULT true, + connector_id text, acl text[] + )` + } + await installProjectionSourceAcl(sql) + /** + * Counts the documents whose write fired the projection fan-out: the trigger fires on any + * assignment of `acl`, changed or not, so this is the cost a no-op revocation must not pay. + */ + await sql`CREATE TABLE fan_out (document_id text NOT NULL)` + await sql.unsafe(`CREATE FUNCTION count_fan_out() RETURNS trigger LANGUAGE plpgsql AS $$ + BEGIN INSERT INTO fan_out VALUES (NEW.id); RETURN NEW; END; $$`) + await sql`CREATE TRIGGER count_fan_out AFTER UPDATE OF connector_id, acl ON document + FOR EACH ROW EXECUTE FUNCTION count_fan_out()` + await sql`ALTER TABLE embedding_search DISABLE TRIGGER embedding_search_source_acl_set` + await sql`ALTER TABLE embedding_keyword_tin DISABLE TRIGGER embedding_keyword_tin_source_acl_set` + }, 60_000) + + afterAll(async () => { + await sql?.end() + await admin?.unsafe(`DROP SCHEMA IF EXISTS "${schemaName}" CASCADE`) + await admin?.end() + }) + + beforeEach(async () => { + hardDelete.mockClear() + await sql`TRUNCATE embedding_search, embedding_keyword_tin, document, fan_out, knowledge_connector` + await sql`INSERT INTO knowledge_connector (id) VALUES (${CONNECTOR})` + }) + + it('revokes an absent document that still grants someone without rewriting one that grants nobody', async () => { + await insertListed(8) + await insertDocuments([ + { id: 'absent-empty', acl: [], seenAt: ABSENT_SINCE }, + { id: 'absent-granted', acl: [ALICE], seenAt: ABSENT_SINCE, verified: true }, + ]) + for (const projection of ['embedding_search', 'embedding_keyword_tin']) { + const prefix = projection === 'embedding_search' ? 'vec' : 'kw' + await sql`INSERT INTO ${sql(projection)} (id, document_id, connector_id, acl) VALUES + (${`${prefix}-empty-filled`}, 'absent-empty', ${CONNECTOR}, '{}'), + (${`${prefix}-granted-unfilled`}, 'absent-granted', NULL, NULL), + (${`${prefix}-granted-filled`}, 'absent-granted', ${CONNECTOR}, ARRAY[${ALICE}])` + } + + const { result, notice } = await reconcile() + + expect(await fannedOut()).toEqual(['absent-granted']) + const absent = await sql< + { id: string; acl: string[]; verified: boolean; deleted: boolean }[] + >`SELECT id, acl, acl_verified_at IS NOT NULL AS verified, deleted_at IS NOT NULL AS deleted + FROM document WHERE id LIKE 'absent-%' ORDER BY id` + expect(absent).toEqual([ + { id: 'absent-empty', acl: [], verified: false, deleted: true }, + { id: 'absent-granted', acl: [], verified: false, deleted: true }, + ]) + expect( + await sql`SELECT id, acl FROM embedding_search + UNION ALL SELECT id, acl FROM embedding_keyword_tin ORDER BY id` + ).toEqual([ + { id: 'kw-empty-filled', acl: [] }, + { id: 'kw-granted-filled', acl: [] }, + { id: 'kw-granted-unfilled', acl: null }, + { id: 'vec-empty-filled', acl: [] }, + { id: 'vec-granted-filled', acl: [] }, + { id: 'vec-granted-unfilled', acl: null }, + ]) + const listed = await sql<{ acl: string[] }[]>` + SELECT acl FROM document WHERE id LIKE 'listed-%'` + expect(listed.every((row) => row.acl.join() === ALICE)).toBe(true) + expect(result.docsDeleted).toBe(2) + expect(notice).toBeNull() + }) + + it('revokes a backlog larger than one change batch completely and reports every removal', async () => { + const ids = Array.from( + { length: 60 }, + (_unused, index) => `absent-${String(index).padStart(3, '0')}` + ) + await insertListed(60) + await insertDocuments(ids.map((id) => ({ id, acl: [ALICE], seenAt: ABSENT_SINCE }))) + await sql`INSERT INTO embedding_search ${sql( + ids.map((id) => ({ id: `vec-${id}`, document_id: id, connector_id: CONNECTOR, acl: [ALICE] })) + )}` + + const { result } = await reconcile(true) + + expect(await fannedOut()).toEqual(ids) + const documents = await sql<{ acl: string[] }[]>` + SELECT acl FROM document WHERE id LIKE 'absent-%'` + expect(documents).toHaveLength(ids.length) + expect(documents.every((row) => row.acl.length === 0)).toBe(true) + const chunks = await sql<{ acl: string[] }[]>`SELECT acl FROM embedding_search` + expect(chunks).toHaveLength(ids.length) + expect(chunks.every((row) => row.acl.length === 0)).toBe(true) + expect(hardDelete.mock.calls.flatMap(([batch]) => batch).sort()).toEqual(ids) + expect(result.docsDeleted).toBe(ids.length) + }) + + /** + * The permission-only page revokes a changed document whatever its ACL; one that already + * grants nobody still has leftover evidence cleared, without an `acl` assignment. + */ + it('clears leftover evidence on a revoked document that already grants nobody without fan-out', async () => { + await insertDocuments([ + { id: 'stale-evidence', acl: [], seenAt: LISTED_AT, verified: true }, + { id: 'granted', acl: [ALICE], seenAt: LISTED_AT, verified: true }, + ]) + await sql`UPDATE document SET acl_requirements = '[[], ["g:confluence:tenant:space"]]' + WHERE id = 'stale-evidence'` + + await revokeDocumentAcls(holder.db as never, ['stale-evidence', 'granted'], (batch) => + inArray(schema.document.id, batch) + ) + + expect(await fannedOut()).toEqual(['granted']) + expect( + await sql`SELECT id, acl, acl_requirements AS requirements, acl_verified_at AS verified + FROM document ORDER BY id` + ).toEqual([ + { id: 'granted', acl: [], requirements: [], verified: null }, + { id: 'stale-evidence', acl: [], requirements: [], verified: null }, + ]) + }) +}) diff --git a/apps/sim/lib/knowledge/connectors/sync-content-pass.test.ts b/apps/sim/lib/knowledge/connectors/sync-content-pass.test.ts index 897aa702fc7..967dbe31b71 100644 --- a/apps/sim/lib/knowledge/connectors/sync-content-pass.test.ts +++ b/apps/sim/lib/knowledge/connectors/sync-content-pass.test.ts @@ -1,6 +1,7 @@ /** @vitest-environment node */ import { dbChainMockFns, + flattenMockConditions, queueTableRows, resetDbChainMock as resetDatabaseMock, schemaMock, @@ -146,6 +147,29 @@ afterEach(() => { vi.unstubAllGlobals() }) +/** + * Every statement that assigns `acl`, with its flattened WHERE. Selects call `where` too, so + * each `set` is paired with the next `where` by call order. + */ +function aclAssignments() { + const whereOrder = dbChainMockFns.where.mock.invocationCallOrder + return dbChainMockFns.set.mock.calls + .map(([values], index) => { + const setOrder = dbChainMockFns.set.mock.invocationCallOrder[index] + const next = whereOrder.findIndex((order) => order > setOrder) + return { + values, + conditions: flattenMockConditions(dbChainMockFns.where.mock.calls[next]?.[0]), + } + }) + .filter(({ values }) => 'acl' in values) +} + +/** The guard that keeps an ACL already readable by nobody out of an `acl` assignment. */ +function grantsSomeone(node: Record): boolean { + return Array.isArray(node.strings) && node.strings.join('?') === 'cardinality(?) > 0' +} + describe('completed listing removal counts', () => { interface AbsentDocument { id: string @@ -164,6 +188,7 @@ describe('completed listing removal counts', () => { hard?: AbsentDocument[] fullSync?: boolean updated?: { id: string }[] + revoked?: AbsentDocument[] }) { resetDbChainMock() const soft = options.soft ?? [] @@ -181,7 +206,11 @@ describe('completed listing removal counts', () => { queueTableRows(schemaMock.document, [ { ownedCount: 10, listedCount: 8, softCount: soft.length, hardCount: hard.length }, ]) - queueTableRows(schemaMock.document, []) + queueTableRows(schemaMock.document, options.revoked ?? []) + if (options.revoked?.length) { + queueTableRows(schemaMock.knowledgeConnector, [{ id: 'connector' }]) + queueTableRows(schemaMock.document, []) + } if (!options.fullSync) { queueTableRows(schemaMock.document, soft) if (soft.length) { @@ -272,6 +301,29 @@ describe('completed listing removal counts', () => { expect(mocks.hardDelete.mock.calls.map(([ids]) => ids)).toEqual([['live'], ['already-hidden']]) }) + /** + * Each document in an `acl` assignment costs a rewrite of its chunks' projection rows, so a + * backlog of absent documents is revoked a small batch at a time, and only where a grant is left. + */ + it('revokes absent documents in small batches, only where they still grant someone', async () => { + const revoked = Array.from({ length: 30 }, (_unused, index) => absent(`absent-${index}`)) + const result = await reconcile({ revoked }) + + const writes = aclAssignments() + expect(writes.map(({ values }) => values)).toEqual([ + { acl: [], aclRequirements: [], aclVerifiedAt: null }, + { acl: [], aclRequirements: [], aclVerifiedAt: null }, + ]) + expect( + writes.map( + ({ conditions }) => + (conditions.find((node) => node.type === 'inArray')?.values as string[]).length + ) + ).toEqual([25, 5]) + expect(writes.every(({ conditions }) => conditions.some(grantsSomeone))).toBe(true) + expect(result.docsDeleted).toBe(0) + }) + it('does not report a full-sync removal when the guarded delete removed no live rows', async () => { mocks.hardDelete.mockResolvedValue(0) const result = await reconcile({ fullSync: true, hard: [absent('detached')] }) @@ -763,6 +815,21 @@ describe('permission refresh through the shared content pass', () => { ) }) + /** Assigning `acl` fires the projection fan-out even when the document already grants nobody. */ + it('assigns acl while revoking a changed document only where it still grants someone', async () => { + sourceBody = { value: '

Current safe body

' } + await runPass({ + access: 'admin', + permissionsOnly: true, + readCurrent: true, + existing: { ...current, contentHash: 'old-body' }, + permissionStoredAfter: current, + }) + const revocations = aclAssignments().filter(({ values }) => values.acl.length === 0) + expect(revocations).toHaveLength(1) + expect(revocations[0].conditions.some(grantsSomeone)).toBe(true) + }) + it('never renews a changed body after its hydration fails', async () => { const { hydrate, result } = await runPass({ access: 'admin', diff --git a/apps/sim/lib/knowledge/connectors/sync-content-pass.ts b/apps/sim/lib/knowledge/connectors/sync-content-pass.ts index 150c3635638..86f515bd805 100644 --- a/apps/sim/lib/knowledge/connectors/sync-content-pass.ts +++ b/apps/sim/lib/knowledge/connectors/sync-content-pass.ts @@ -30,6 +30,7 @@ import { assertSyncLeaseHeldInTx, type SyncRunLease } from '@/lib/knowledge/conn import { type KnowledgeBaseOwner, persistSourceDocumentFailures, + revokeDocumentAcls, } from '@/lib/knowledge/connectors/sync-persistence' import { buildReconciliationHoldNotice, @@ -179,23 +180,18 @@ export async function runConnectorContentPass(input: ContentPassInput) { }) /** Revoke grants without matching stored content, retaining the content crawl's observation for EOF reconciliation. */ if (changed.length) - await withLease(async (tx) => { - for (let offset = 0; offset < changed.length; offset += 500) { - await tx - .update(document) - .set({ acl: [], aclRequirements: [], aclVerifiedAt: null }) - .where( - and( - eq(document.connectorId, input.connectorId), - inArray( - document.externalId, - changed.slice(offset, offset + 500).map((item) => item.externalId) - ), - isNull(document.archivedAt) - ) + await withLease((tx) => + revokeDocumentAcls( + tx, + changed.map((item) => item.externalId), + (batch) => + and( + eq(document.connectorId, input.connectorId), + inArray(document.externalId, batch), + isNull(document.archivedAt) ) - } - }) + ) + ) } const state = createSyncRunState(input.result) const startedAt = new Date(cycle.startedAt) @@ -431,18 +427,11 @@ async function reconcileCompletedListing( const rows = await loadBatch(and(absent, sql`cardinality(${document.acl}) > 0`), 500, after) if (rows.length === 0) break await withLease((tx) => - tx - .update(document) - .set({ acl: [], aclRequirements: [], aclVerifiedAt: null }) - .where( - and( - absent, - inArray( - document.id, - rows.map((row) => row.id) - ) - ) - ) + revokeDocumentAcls( + tx, + rows.map((row) => row.id), + (batch) => and(absent, inArray(document.id, batch)) + ) ) after = rows.at(-1) } diff --git a/apps/sim/lib/knowledge/connectors/sync-persistence.test.ts b/apps/sim/lib/knowledge/connectors/sync-persistence.test.ts index ac84bae837a..86931ddbed5 100644 --- a/apps/sim/lib/knowledge/connectors/sync-persistence.test.ts +++ b/apps/sim/lib/knowledge/connectors/sync-persistence.test.ts @@ -1,7 +1,16 @@ /** * @vitest-environment node */ -import { dbChainMockFns, queueTableRows, resetDbChainMock, schemaMock } from '@sim/testing' + +import { db } from '@sim/db' +import { + dbChainMockFns, + flattenMockConditions, + queueTableRows, + resetDbChainMock, + schemaMock, +} from '@sim/testing' +import { inArray } from 'drizzle-orm' import { beforeEach, describe, expect, it, vi } from 'vitest' vi.mock('@/lib/knowledge/documents/service', () => ({ hardDeleteDocuments: vi.fn() })) @@ -45,6 +54,7 @@ import { persistDocumentAcls, persistSourceDocumentFailures, resolveTagMapping, + revokeDocumentAcls, } from '@/lib/knowledge/connectors/sync-persistence' const CONNECTOR = 'connector-1' @@ -281,6 +291,74 @@ describe('persistDocumentAcls', () => { }) }) +describe('revokeDocumentAcls', () => { + beforeEach(() => { + vi.clearAllMocks() + resetDbChainMock() + }) + + const REVOKED = { acl: [], aclRequirements: [], aclVerifiedAt: null } + const EVIDENCE_ONLY = { aclRequirements: [], aclVerifiedAt: null } + const scope = (batch: string[]) => inArray(schemaMock.document.id, batch) + + /** The SQL text of a mock `sql` node, or undefined for an operator node. */ + const sqlText = (node: Record) => + Array.isArray(node.strings) ? node.strings.join('?') : undefined + const grants = (node: Record) => + sqlText(node)?.startsWith('cardinality(') && sqlText(node)?.endsWith(') > 0') + + /** Every `where` condition, paired with the `set` of the same statement. */ + function statements() { + return dbChainMockFns.set.mock.calls.map(([values], index) => ({ + values, + conditions: flattenMockConditions(dbChainMockFns.where.mock.calls[index]?.[0]), + })) + } + + /** + * Assigning `acl` fires the projection fan-out whether or not the value changes, so a + * document that already grants nobody must never be in an `acl` assignment. + */ + it('assigns acl only to documents that still grant someone', async () => { + await revokeDocumentAcls(db, ['a', 'b'], scope) + + const writes = statements().filter(({ values }) => 'acl' in values) + expect(writes).toHaveLength(1) + expect(writes[0].values).toEqual(REVOKED) + expect(writes[0].conditions.some(grants)).toBe(true) + }) + + it('clears leftover evidence on an already-empty ACL without assigning acl', async () => { + await revokeDocumentAcls(db, ['a', 'b'], scope) + + const clears = statements().filter(({ values }) => !('acl' in values)) + expect(clears).toHaveLength(1) + expect(clears[0].values).toEqual(EVIDENCE_ONLY) + expect( + clears[0].conditions.some( + (node) => node.type === 'not' && grants(node.condition as Record) + ) + ).toBe(true) + }) + + /** Each document in an `acl` assignment costs a rewrite of every one of its chunks' projection rows. */ + it('assigns acl in batches of 25 and clears evidence in batches of 500', async () => { + const ids = Array.from({ length: 60 }, (_unused, index) => `doc-${index}`) + + await revokeDocumentAcls(db, ids, scope) + + const batchSizes = (writesAcl: boolean) => + statements() + .filter(({ values }) => 'acl' in values === writesAcl) + .map(({ conditions }) => { + const inArray = conditions.find((node) => node.type === 'inArray') + return (inArray?.values as string[]).length + }) + expect(batchSizes(true)).toEqual([25, 25, 10]) + expect(batchSizes(false)).toEqual([60]) + }) +}) + describe('persistSourceDocumentFailures', () => { beforeEach(() => { vi.clearAllMocks() diff --git a/apps/sim/lib/knowledge/connectors/sync-persistence.ts b/apps/sim/lib/knowledge/connectors/sync-persistence.ts index b675700125d..ec3799f1dd5 100644 --- a/apps/sim/lib/knowledge/connectors/sync-persistence.ts +++ b/apps/sim/lib/knowledge/connectors/sync-persistence.ts @@ -4,7 +4,7 @@ import { createLogger } from '@sim/logger' import { chunkArray } from '@sim/utils/helpers' import { generateId } from '@sim/utils/id' import { truncateAtCodePoint } from '@sim/utils/string' -import { and, eq, exists, inArray, isNull, lt, not, or, sql } from 'drizzle-orm' +import { and, eq, exists, inArray, isNull, lt, not, or, type SQL, sql } from 'drizzle-orm' import { getInternalApiBaseUrl } from '@/lib/core/utils/urls' import type { DbOrTx } from '@/lib/db/types' import { textArrayLiteral } from '@/lib/knowledge/access/predicate' @@ -196,6 +196,40 @@ export async function persistDocumentAcls( return { updated, rejected } } +/** + * Revokes every grant on the documents `target` selects from `ids`, leaving each readable by + * nobody with its permission evidence cleared. Only a document that still grants someone has + * `acl` assigned, {@link ACL_CHANGE_BATCH_SIZE} at a time: the projection trigger fires on every + * assignment of `acl`, changed or not, and each document costs a rewrite of its chunks' + * projection rows. A document already readable by nobody only has leftover evidence cleared, + * which fires no fan-out. + */ +export async function revokeDocumentAcls( + executor: DbOrTx, + ids: string[], + target: (batch: string[]) => SQL | undefined +): Promise { + const grants = sql`cardinality(${document.acl}) > 0` + for (const batch of chunkArray(ids, ACL_WRITE_BATCH_SIZE)) { + await executor + .update(document) + .set({ aclRequirements: [], aclVerifiedAt: null }) + .where( + and( + target(batch), + not(grants), + sql`(${document.aclRequirements} <> '[]'::jsonb OR ${document.aclVerifiedAt} IS NOT NULL)` + ) + ) + } + for (const batch of chunkArray(ids, ACL_CHANGE_BATCH_SIZE)) { + await executor + .update(document) + .set({ acl: [], aclRequirements: [], aclVerifiedAt: null }) + .where(and(target(batch), grants)) + } +} + const MAX_SAFE_TITLE_LENGTH = 200 /** The suffix a cut value carries, counted inside {@link MAX_DOCUMENT_INDEXED_TEXT_LENGTH}. */