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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -120,7 +120,7 @@ describe('connector lease ACL pages in PostgreSQL', () => {
RETURN NEW;
END $$`)
)
for (const projection of ['embedding_search', 'embedding_keyword_tin']) {
for (const projection of ['embedding_search']) {
await db.execute(
sql.raw(`DROP TRIGGER IF EXISTS log_lease_page_projection_write ON ${projection}`)
)
Expand Down Expand Up @@ -151,7 +151,7 @@ describe('connector lease ACL pages in PostgreSQL', () => {
await db.execute(sql`DROP FUNCTION IF EXISTS log_lease_page_acl_write()`)
await db.execute(sql`DROP FUNCTION IF EXISTS fail_after_acl_writes()`)
await db.execute(sql`DROP TABLE IF EXISTS lease_page_acl_writes`)
for (const projection of ['embedding_search', 'embedding_keyword_tin'])
for (const projection of ['embedding_search'])
await db.execute(
sql.raw(`DROP TRIGGER IF EXISTS log_lease_page_projection_write ON ${projection}`)
)
Expand Down Expand Up @@ -300,25 +300,14 @@ describe('connector lease ACL pages in PostgreSQL', () => {
)
for (let offset = 0; offset < rows.length; offset += 200)
await db.insert(embedding).values(rows.slice(offset, offset + 200))
/**
* Each chunk's projection rows, filled from its document as the backfill leaves them. The
* embedding insert writes the vector projection's row; the keyword projection's own sync
* trigger ships with the Tin migration, which a database without Tin skips. The fixture sets
* the filled state itself, whatever triggers an earlier suite left behind; what is under
* test is the document trigger that rewrites these rows.
*/
/** Seeds the filled vector state whose ACL updates must stay bounded. */
const chunkIds = sql.join(
documents.map((entry) => sql`${entry.id}`),
sql`, `
)
await db.execute(sql`
UPDATE embedding_search p SET enabled = true, connector_id = d.connector_id, acl = d.acl
FROM document d WHERE d.id = p.document_id AND d.id IN (${chunkIds})`)
await db.execute(sql`
INSERT INTO embedding_keyword_tin (id, knowledge_base_id, document_id, enabled, content, connector_id, acl)
SELECT e.id, e.knowledge_base_id, e.document_id, true, e.content, d.connector_id, d.acl
FROM embedding e JOIN document d ON d.id = e.document_id WHERE e.document_id IN (${chunkIds})
ON CONFLICT (id) DO UPDATE SET enabled = true, connector_id = EXCLUDED.connector_id, acl = EXCLUDED.acl`)
/** Only writes made by the code under test are counted. */
await db.execute(sql`DELETE FROM lease_page_projection_writes`)
await db
Expand All @@ -337,10 +326,7 @@ describe('connector lease ACL pages in PostgreSQL', () => {
db.execute<{ projection: string; acl: string | null; expected: string }>(sql`
SELECT 'embedding_search' AS projection, array_to_string(p.acl, ',') AS acl,
array_to_string(d.acl, ',') AS expected
FROM embedding_search p JOIN document d ON d.id = p.document_id WHERE d.connector_id = ${connectorId}
UNION ALL
SELECT 'embedding_keyword_tin', array_to_string(p.acl, ','), array_to_string(d.acl, ',')
FROM embedding_keyword_tin p JOIN document d ON d.id = p.document_id WHERE d.connector_id = ${connectorId}`)
FROM embedding_search p JOIN document d ON d.id = p.document_id WHERE d.connector_id = ${connectorId}`)

/** Projection rows each transaction rewrote, per table. */
const projectionRowsPerTransaction = async () =>
Expand All @@ -354,7 +340,7 @@ describe('connector lease ACL pages in PostgreSQL', () => {
const seeded = await seedDocuments(ids.connectorId, [alice()], 30)
await seedChunks(seeded)
const before = [...(await projectionAcls(ids.connectorId))]
expect(before).toHaveLength(2 * 30 * CHUNKS)
expect(before).toHaveLength(30 * CHUNKS)

await expect(
persistDocumentAcls(
Expand All @@ -367,7 +353,7 @@ describe('connector lease ACL pages in PostgreSQL', () => {
const after = [...(await projectionAcls(ids.connectorId))]
expect(after.filter((row) => row.acl !== bob() || row.expected !== bob())).toEqual([])
const perTransaction = await projectionRowsPerTransaction()
expect(perTransaction.reduce((total, rows) => total + rows, 0)).toBe(2 * 30 * CHUNKS)
expect(perTransaction.reduce((total, rows) => total + rows, 0)).toBe(30 * CHUNKS)
expect(Math.max(...perTransaction)).toBeLessThanOrEqual(PROJECTION_ROW_BATCH_SIZE)
})

Expand All @@ -384,7 +370,7 @@ describe('connector lease ACL pages in PostgreSQL', () => {
).resolves.toBe(true)

const after = [...(await projectionAcls(members.connectorId))]
expect(after).toHaveLength(2 * 30 * CHUNKS)
expect(after).toHaveLength(30 * CHUNKS)
expect(after.filter((row) => row.acl !== '' || row.expected !== '')).toEqual([])
expect(Math.max(...(await projectionRowsPerTransaction()))).toBeLessThanOrEqual(
PROJECTION_ROW_BATCH_SIZE
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,6 @@ import {
import {
document,
embedding,
embeddingKeywordSearch,
embeddingKeywordTin,
embeddingSearch,
knowledgeBase,
Expand Down Expand Up @@ -123,6 +122,25 @@ afterAll(async () => {
})

describe('ordinary KB vector repair after indexed Search retirement', () => {
it('keeps canonical chunk writes after the retired keyword projection is removed', async () => {
const [relations] =
await projector`SELECT to_regclass('public.embedding_keyword_search') AS keyword`
expect(relations.keyword).toBeNull()
const chunk = chunkRow(generateId(), 0)
await db.insert(embedding).values(chunk)
await db
.update(embedding)
.set({ content: 'updated searchable text' })
.where(eq(embedding.id, chunk.id))
const [canonical] =
await projector`SELECT content_tsv @@ plainto_tsquery('english', 'searchable') AS matches
FROM embedding WHERE id = ${chunk.id}`
expect(canonical.matches).toBe(true)
expect((await vectorsOf()).map((row) => row.id)).toEqual([chunk.id])
await db.delete(embedding).where(eq(embedding.id, chunk.id))
expect(await vectorsOf()).toEqual([])
})

it('repairs deferred vectors in bounded pages without creating copied ACL or keyword rows', async () => {
const chunks = Array.from({ length: 5 }, (_, index) => chunkRow(generateId(), index))
await defer((tx) => tx.insert(embedding).values(chunks))
Expand All @@ -143,12 +161,6 @@ describe('ordinary KB vector repair after indexed Search retirement', () => {
(row) => row.vector512?.[511] === 1 && row.acl === null && row.connectorId === null
)
).toBe(true)
expect(
await db
.select()
.from(embeddingKeywordSearch)
.where(eq(embeddingKeywordSearch.documentId, documentId))
).toEqual([])
expect(
await db
.select()
Expand All @@ -167,12 +179,6 @@ describe('ordinary KB vector repair after indexed Search retirement', () => {
expect(await markOf(searchDocumentId)).toMatchObject({ content: true })
await runKnowledgeProjection(projector, {})
expect(await vectorsOf(searchDocumentId)).toEqual([])
expect(
await db
.select()
.from(embeddingKeywordSearch)
.where(eq(embeddingKeywordSearch.documentId, searchDocumentId))
).toEqual([])
expect(
await db
.select()
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -685,7 +685,6 @@ describe.skipIf(!enabled)('Knowledge search latency on a realistic indexed corpu
await db.execute(sql`ANALYZE document`)
await db.execute(sql`ANALYZE embedding`)
await db.execute(sql`ANALYZE embedding_search`)
await db.execute(sql`ANALYZE embedding_keyword_search`)
if (evictSharedBuffers) await db.execute(sql`CREATE EXTENSION IF NOT EXISTS pg_buffercache`)
report.server = (
await db.execute(sql`SELECT version(), current_setting('work_mem') AS work_mem,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,6 @@ import {
credentialGroupEnrollment,
document,
embedding,
embeddingKeywordSearch,
embeddingSearch,
knowledgeBase,
knowledgeConnector,
Expand Down Expand Up @@ -135,9 +134,6 @@ describe.each([
.update(embeddingSearch)
.set({ acl: ['s:nobody:-:stale'], connectorId: null })
.where(eq(embeddingSearch.knowledgeBaseId, baseId))
await db
.delete(embeddingKeywordSearch)
.where(eq(embeddingKeywordSearch.knowledgeBaseId, baseId))
})

afterAll(async () => {
Expand Down
25 changes: 6 additions & 19 deletions apps/sim/lib/knowledge/connectors/detachment.test.ts
Original file line number Diff line number Diff line change
@@ -1,10 +1,4 @@
import {
document,
embeddingKeywordTin,
embeddingSearch,
knowledgeBase,
knowledgeConnector,
} from '@sim/db/schema'
import { document, embeddingSearch, knowledgeBase, knowledgeConnector } from '@sim/db/schema'
import { dbChainMockFns, queueTableRows, resetDbChainMock } from '@sim/testing'
import { billingStorageMock, billingStorageMockFns } from '@sim/testing/mocks/billing-storage.mock'
import {
Expand Down Expand Up @@ -68,12 +62,10 @@ function queueBatch(documentIds: string[], reservedBytes = 0) {
}

/** Row ids the projection release returns, in the order the handler updates the tables. */
function releaseProjectionRows(searchRows: number, keywordRows: number) {
function releaseProjectionRows(searchRows: number) {
const rows = (count: number) =>
Array.from({ length: count }, (_, index) => ({ id: `c-${index}` }))
dbChainMockFns.returning
.mockResolvedValueOnce(rows(searchRows))
.mockResolvedValueOnce(rows(keywordRows))
dbChainMockFns.returning.mockResolvedValueOnce(rows(searchRows))
}

function updatedTables() {
Expand All @@ -92,7 +84,7 @@ describe('connector detachment', () => {

it('releases documents against the reservation, then settles what remains with the connector', async () => {
queueBatch(['doc-1', 'doc-2'], 25)
releaseProjectionRows(3, 3)
releaseProjectionRows(3)
dbChainMockFns.returning.mockResolvedValueOnce([
{ fileSize: 10, deletedAt: null },
{ fileSize: 5, deletedAt: new Date('2026-09-01T00:00:00.000Z') },
Expand All @@ -108,12 +100,7 @@ describe('connector detachment', () => {
context()
)

expect(updatedTables()).toEqual([
embeddingSearch,
embeddingKeywordTin,
document,
knowledgeConnector,
])
expect(updatedTables()).toEqual([embeddingSearch, document, knowledgeConnector])
expect(dbChainMockFns.set).toHaveBeenCalledWith({ connectorId: null })
expect(dbChainMockFns.set).toHaveBeenCalledWith(
expect.objectContaining({ connectorId: null, deletedAt: expect.anything() })
Expand All @@ -131,7 +118,7 @@ describe('connector detachment', () => {
it('keeps documents attached until their search rows fit one page, without spending retries', async () => {
for (let batch = 0; batch < 4; batch++) {
queueBatch(['doc-1'])
releaseProjectionRows(250, 0)
releaseProjectionRows(250)
}

expect(await detachKnowledgeConnector(payload, context())).toMatchObject({
Expand Down
50 changes: 19 additions & 31 deletions apps/sim/lib/knowledge/connectors/detachment.ts
Original file line number Diff line number Diff line change
@@ -1,12 +1,6 @@
import { db } from '@sim/db'
import { PROJECTION_ROW_BATCH_SIZE } from '@sim/db/knowledge-projection'
import {
document,
embeddingKeywordTin,
embeddingSearch,
knowledgeBase,
knowledgeConnector,
} from '@sim/db/schema'
import { document, embeddingSearch, knowledgeBase, knowledgeConnector } from '@sim/db/schema'
import { and, asc, eq, inArray, isNotNull, lt, ne, sql } from 'drizzle-orm'
import { z } from 'zod'
import {
Expand Down Expand Up @@ -78,8 +72,6 @@ export function keptDocumentBytes() {
END`
}

const SEARCH_PROJECTIONS = [embeddingSearch, embeddingKeywordTin] as const

/**
* Settles what is left of a detached connector's reservation: an unreleased remainder is refunded,
* and an overdraft, released bytes beyond what removal charged, is charged as already admitted.
Expand Down Expand Up @@ -196,8 +188,8 @@ export async function settleDetachedConnectorReservations(
* Releases a detached connector's documents as standalone entries, then deletes the connector.
*
* Nulling a document's `connector_id` fires the projection trigger, which rewrites every enabled
* chunk of it in both search projections, each a fresh index entry. So each transaction first
* releases at most 250 projection rows per table of the next 100 documents, and flips those
* chunk of it in the vector projection, each a fresh index entry. So each transaction first
* releases at most 250 vector projection rows of the next 100 documents, and flips those
* documents only once none of their rows still name the connector; the trigger then finds nothing
* to rewrite. A document larger than one page spans several transactions, and its released rows
* read as an upload in the meantime, which grants the same access: only workspace-access
Expand Down Expand Up @@ -299,27 +291,23 @@ export const detachKnowledgeConnector: OutboxHandler = async (rawPayload, contex
return drained
}

let projectionPageFull = false
for (const projection of SEARCH_PROJECTIONS) {
const page = tx
.select({ id: projection.id })
.from(projection)
.where(
and(
inArray(projection.documentId, documentIds),
eq(projection.enabled, true),
isNotNull(projection.connectorId)
)
const page = tx
.select({ id: embeddingSearch.id })
.from(embeddingSearch)
.where(
and(
inArray(embeddingSearch.documentId, documentIds),
eq(embeddingSearch.enabled, true),
isNotNull(embeddingSearch.connectorId)
)
.limit(PROJECTION_ROW_BATCH_SIZE)
const released = await tx
.update(projection)
.set({ connectorId: null })
.where(inArray(projection.id, page))
.returning({ id: projection.id })
if (released.length === PROJECTION_ROW_BATCH_SIZE) projectionPageFull = true
}
if (projectionPageFull) return 'progress'
)
.limit(PROJECTION_ROW_BATCH_SIZE)
const released = await tx
.update(embeddingSearch)
.set({ connectorId: null })
.where(inArray(embeddingSearch.id, page))
.returning({ id: embeddingSearch.id })
if (released.length === PROJECTION_ROW_BATCH_SIZE) return 'progress'

const releasedDocuments = await tx
.update(document)
Expand Down
15 changes: 4 additions & 11 deletions apps/sim/lib/knowledge/connectors/sync-content-pass.integration.ts
Original file line number Diff line number Diff line change
Expand Up @@ -153,7 +153,7 @@ describe('completed listing reconciliation in PostgreSQL', () => {
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']) {
for (const projection of ['embedding_search']) {
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[]
Expand All @@ -170,7 +170,6 @@ describe('completed listing reconciliation in PostgreSQL', () => {
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 () => {
Expand All @@ -181,7 +180,7 @@ describe('completed listing reconciliation in PostgreSQL', () => {

beforeEach(async () => {
hardDelete.mockClear()
await sql`TRUNCATE embedding_search, embedding_keyword_tin, document, fan_out, knowledge_connector`
await sql`TRUNCATE embedding_search, document, fan_out, knowledge_connector`
await sql`INSERT INTO knowledge_connector (id) VALUES (${CONNECTOR})`
})

Expand All @@ -191,7 +190,7 @@ describe('completed listing reconciliation in PostgreSQL', () => {
{ 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']) {
for (const projection of ['embedding_search']) {
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}, '{}'),
Expand All @@ -210,13 +209,7 @@ describe('completed listing reconciliation in PostgreSQL', () => {
{ 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 },
expect(await sql`SELECT id, acl FROM embedding_search ORDER BY id`).toEqual([
{ id: 'vec-empty-filled', acl: [] },
{ id: 'vec-granted-filled', acl: [] },
{ id: 'vec-granted-unfilled', acl: null },
Expand Down
Loading
Loading