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 @@ -17,6 +17,7 @@ import {
workspace,
workspaceFiles,
} from '@sim/db/schema'
import { createDeferred } from '@sim/testing/helpers/deferred'
import { generateId } from '@sim/utils/id'
import { and, eq, inArray, sql } from 'drizzle-orm'
import { afterAll, beforeAll, describe, expect, it, vi } from 'vitest'
Expand Down Expand Up @@ -50,7 +51,11 @@ import { listKnowledgeChunks } from '@/lib/knowledge/application/chunks'
import { searchKnowledge } from '@/lib/knowledge/application/search'
import { KNOWLEDGE_DOCUMENT_PROCESSING_OUTBOX_EVENT } from '@/lib/knowledge/documents/processing-events'
import { knowledgeDocumentProcessingOutboxHandlers } from '@/lib/knowledge/documents/processing-outbox-handler'
import { createDocumentRecords, createSingleDocument } from '@/lib/knowledge/documents/service'
import {
createDocumentRecords,
createSingleDocument,
hardDeleteDocuments,
} from '@/lib/knowledge/documents/service'
import { KNOWLEDGE_STORAGE_CLEANUP_EVENT } from '@/lib/knowledge/documents/storage-cleanup-event'
import { uploadKnowledgeArtifact } from '@/lib/knowledge/documents/storage-upload'
import {
Expand Down Expand Up @@ -88,6 +93,42 @@ async function sourceFile(
)
}

/**
* Runs `work` while another transaction holds `FOR KEY SHARE` on the knowledge base row, as an
* in-flight embedding insert does through its foreign key, and requires it to settle well before
* that holder's idle timeout would release the lock.
*/
async function settlesWhileKeyShareHeld<T>(knowledgeBaseId: string, work: () => Promise<T>) {
const release = createDeferred<void>()
const acquired = createDeferred<void>()
const holding = db.transaction(async (tx) => {
await tx.execute(sql`SET LOCAL idle_in_transaction_session_timeout = '15s'`)
await tx
.select({ id: knowledgeBase.id })
.from(knowledgeBase)
.where(eq(knowledgeBase.id, knowledgeBaseId))
.for('key share')
acquired.resolve()
await release.promise
})
await acquired.promise
const running = Promise.allSettled([work()])
let settled = false
void running.then(() => {
settled = true
})
try {
await expect.poll(() => settled, { interval: 5, timeout: 5000 }).toBe(true)
const [outcome] = await running
if (outcome.status === 'rejected') throw outcome.reason
return outcome.value
} finally {
release.resolve()
await holding
await running
}
}

/** Drains real audit writes before asserting persisted rows, including duplicates from one operation. */
function observeAbsenceAudits(workspaceId: string) {
let submitted = 0
Expand Down Expand Up @@ -229,6 +270,33 @@ describe('durable workspace file import', () => {
}
})

it.each(['single', 'bulk'] as const)(
'%s admission and hard delete commit while an embedding holds the KB key-share lock',
async (mode) => {
const ids = await seed()
const source = await sourceFile(ids, 'ordinary source bytes')
const input = {
filename: source.name,
fileUrl: `/api/files/serve/${encodeURIComponent(source.key)}?context=workspace`,
fileSize: source.size,
mimeType: source.type,
}
const documentIds = await settlesWhileKeyShareHeld(ids.knowledgeBaseId, async () =>
mode === 'single'
? [(await createSingleDocument(input, ids.knowledgeBaseId, generateId(), ids.aliceId)).id]
: (
await createDocumentRecords([input], ids.knowledgeBaseId, generateId(), ids.aliceId)
).map((created) => created.documentId)
)
expect(documentIds).toHaveLength(1)
expect(
await settlesWhileKeyShareHeld(ids.knowledgeBaseId, () =>
hardDeleteDocuments(documentIds, generateId())
)
).toBe(1)
}
)

it('indexes the admitted snapshot after the source was deleted and temporary access would have expired', async () => {
const ids = await seed()
const content =
Expand Down
19 changes: 16 additions & 3 deletions apps/sim/lib/knowledge/documents/service.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2546,7 +2546,15 @@ export async function createDocumentRecords(
const { returnData, storageNotification, unrecordedCount } = await db.transaction(async (tx) => {
let storageNotification: DocumentStorageNotification | null = null

await tx.execute(sql`SELECT 1 FROM knowledge_base WHERE id = ${knowledgeBaseId} FOR UPDATE`)
/**
* `FOR NO KEY UPDATE` still serializes this insert with other document writes,
* KB moves, and archive, but not with the `FOR KEY SHARE` that each embedding
* insert takes through its foreign key, so uploads never queue behind a
* long-running indexing transaction in the same knowledge base.
*/
await tx.execute(
sql`SELECT 1 FROM knowledge_base WHERE id = ${knowledgeBaseId} FOR NO KEY UPDATE`
)

const kb = await tx
.select({
Expand Down Expand Up @@ -3229,7 +3237,9 @@ export async function createSingleDocument(
await claimKnowledgeUploadForAttachment(tx, options.uploadedArtifact.cleanupEventId)
}

await tx.execute(sql`SELECT 1 FROM knowledge_base WHERE id = ${knowledgeBaseId} FOR UPDATE`)
await tx.execute(
sql`SELECT 1 FROM knowledge_base WHERE id = ${knowledgeBaseId} FOR NO KEY UPDATE`
)

const kb = await tx
.select({
Expand Down Expand Up @@ -4332,6 +4342,9 @@ async function hardDeleteDocumentBatch(
* Lock every parent KB in stable ID order before deleting document rows.
* Normal inserts and KB moves take the same parent lock first, so the
* workspace snapshots used for accounting cannot change mid-delete.
* `no key update` lets embedding inserts keep taking their foreign-key
* `KEY SHARE` lock, so a delete neither waits on indexing elsewhere in the
* KB nor deadlocks against an indexing pass on one of its own documents.
*/
const knowledgeBaseIds = [
...new Set(documentsToDelete.map((doc) => doc.knowledgeBaseId)),
Expand All @@ -4351,7 +4364,7 @@ async function hardDeleteDocumentBatch(
)
)
.orderBy(asc(knowledgeBase.id))
.for('update')
.for('no key update')
const lockedKnowledgeBaseById = new Map(lockedKnowledgeBases.map((kb) => [kb.id, kb]))
for (const doc of documentsToDelete) {
const lockedKb = lockedKnowledgeBaseById.get(doc.knowledgeBaseId)
Expand Down
2 changes: 1 addition & 1 deletion apps/sim/lib/knowledge/orchestration/connectors.ts
Original file line number Diff line number Diff line change
Expand Up @@ -471,7 +471,7 @@ export async function performCreateKnowledgeConnector(
try {
created = await db.transaction(async (tx) => {
if (owner.organizationId) await lockOrganizationSearchApproval(tx, owner.organizationId)
await tx.execute(sql`SELECT 1 FROM knowledge_base WHERE id = ${kb.id} FOR UPDATE`)
await tx.execute(sql`SELECT 1 FROM knowledge_base WHERE id = ${kb.id} FOR NO KEY UPDATE`)

const activeKb = await tx
.select({ id: knowledgeBase.id })
Expand Down
4 changes: 3 additions & 1 deletion apps/sim/lib/knowledge/service.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1136,7 +1136,9 @@ export async function restoreKnowledgeBase(
try {
await db.transaction(async (tx) => {
if (kb.organizationId) await lockOrganizationSearchApproval(tx, kb.organizationId)
await tx.execute(sql`SELECT 1 FROM knowledge_base WHERE id = ${knowledgeBaseId} FOR UPDATE`)
await tx.execute(
sql`SELECT 1 FROM knowledge_base WHERE id = ${knowledgeBaseId} FOR NO KEY UPDATE`
)

attemptedRestoreName = await generateRestoreName(kb.name, async (candidate) => {
if (!kb.workspaceId) return false
Expand Down
Loading