Skip to content

Commit b124b72

Browse files
committed
fix(db): re-decide the deletion guard inside each deleting transaction and reset member relist state
Each deleting transaction share-locks the search index and its connectors and re-decides the guard, so a connector resumed between pages waits for the page in flight and the next page refuses. Every document-deleting transaction also resets the stopped connectors' listing state, now including the directory checkpoint and each member's retry time, so a run stopped partway never leaves a connector that would skip deleted documents on resume.
1 parent cd837b6 commit b124b72

5 files changed

Lines changed: 183 additions & 75 deletions

File tree

‎apps/sim/scripts/dormant-org-search/README.md‎

Lines changed: 6 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -152,12 +152,14 @@ Source-connected documents are not metered as uploaded storage, so there is no s
152152

153153
The **knowledge base row and its paused connectors are kept**: their configuration, credentials, members, permission snapshots and sync history. Re-enabling is then a matter of resuming the connectors, with nothing to set up again.
154154

155-
Once no connector-owned document remains, the run clears each stopped connector's listing state. These are the same columns the app clears when a connector must list everything again (an access-mode switch or a source change):
155+
Every page that deletes documents also clears each stopped connector's listing state, in the same transaction, so a run stopped partway never leaves a connector that would skip what was already deleted. These are the columns the app clears when a connector must list everything again (an access-mode switch or a source change), plus the directory checkpoint and each member's retry time:
156156

157-
- `lastSyncAt`, `lastSyncDocCount`, `listingCheckpoint` and `memberTombstoneCursor` on `knowledge_connector`
158-
- `listingCheckpoint`, `changeCursor`, `memberSyncedThrough`, `lastCompleteListingAt` and `lastListedCount` on its `knowledge_connector_member` rows
157+
- `lastSyncAt`, `lastSyncDocCount`, `listingCheckpoint`, `directoryCheckpoint` and `memberTombstoneCursor` on `knowledge_connector`
158+
- `listingCheckpoint`, `changeCursor`, `memberSyncedThrough`, `lastCompleteListingAt`, `lastListedCount` and `nextAttemptAt` (so every member is due) on its `knowledge_connector_member` rows
159159

160-
Without this, a resumed connector would sync incrementally from its old cursor ("changed since last sync"). It would never re-list the documents deleted here, and the index would stay silently incomplete. Partition work rows (`knowledge_connector_partition`) belong to the old listing generation, and the next full listing replaces them. `--no-connector-reset` skips the reset. The reset runs only when the run reaches the end with no connector documents left: a run resumed with `--after-id` checks for documents before its cursor and says so.
160+
Without this, a resumed connector would sync incrementally from its old cursor ("changed since last sync"), skip directory reconciliation behind a `complete` checkpoint, or leave members waiting on a future retry. It would never re-list the documents deleted here, and the index would stay silently incomplete. Partition work rows (`knowledge_connector_partition`) belong to the old listing generation, and the next full listing replaces them. `--no-connector-reset` skips the reset. When a run finishes with no connector documents left, it resets once more to catch rows a page could not (for example a member added mid-run).
161+
162+
Each deleting transaction also re-decides the guard while holding the knowledge base and its connectors `FOR SHARE`. Resuming a connector or claiming a sync updates the connector row, so it waits for the page in flight to commit, and the next page refuses.
161163

162164
### Duration and monitoring
163165

‎apps/sim/scripts/dormant-org-search/dormant-org-search.integration.ts‎

Lines changed: 25 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -262,13 +262,25 @@ describe('dormant organization search runbook in PostgreSQL', () => {
262262
sleep,
263263
})
264264
).rejects.toThrow('not an organization search index')
265+
/** A deleting transaction re-decides the guard itself, so a resume between pages cannot race it. */
266+
await expect(
267+
store.deleteChunkBatch(indexId, indexDocumentIds.slice(0, 1), 10)
268+
).rejects.toBeInstanceOf(SearchIndexDeletionRefused)
269+
await expect(
270+
store.deleteDocuments(indexId, indexDocumentIds.slice(0, 1), 'fixture', true)
271+
).rejects.toBeInstanceOf(SearchIndexDeletionRefused)
265272
expect((await rowsOf(indexId)).documents).toBe(DOCUMENTS)
266273
})
267274

268275
it('deletes the search index in resumable pages, queues storage cleanup, and marks nothing', async () => {
269276
await db
270277
.update(knowledgeConnector)
271-
.set({ status: 'paused' })
278+
.set({
279+
status: 'paused',
280+
lastSyncAt: new Date(),
281+
listingCheckpoint: { cursor: 'fixture' },
282+
directoryCheckpoint: { phase: 'complete' },
283+
})
272284
.where(eq(knowledgeConnector.id, indexConnectorId))
273285
const store = drizzleSearchIndexDeletionStore(timeouts)
274286
const options = {
@@ -295,6 +307,16 @@ describe('dormant organization search runbook in PostgreSQL', () => {
295307
})
296308
expect(first).toMatchObject({ pages: 1, documentsDeleted: 3, done: false })
297309
expect((await rowsOf(indexId)).documents).toBe(DOCUMENTS - 3)
310+
/** A run stopped after one page already leaves the connector listing from scratch on resume. */
311+
const [partial] = await db
312+
.select()
313+
.from(knowledgeConnector)
314+
.where(eq(knowledgeConnector.id, indexConnectorId))
315+
expect(partial).toMatchObject({
316+
lastSyncAt: null,
317+
listingCheckpoint: null,
318+
directoryCheckpoint: null,
319+
})
298320

299321
const rest = await deleteSearchIndexDocuments(store, {
300322
...options,
@@ -304,7 +326,7 @@ describe('dormant organization search runbook in PostgreSQL', () => {
304326
expect(rest).toMatchObject({
305327
documentsDeleted: DOCUMENTS - 3,
306328
done: true,
307-
connectorsReset: { connectors: 1, members: 0 },
329+
connectorsReset: { connectors: 0, members: 0 },
308330
standaloneDocumentsRemain: false,
309331
projectionMarks: { before: 0, after: 0 },
310332
})
@@ -343,6 +365,7 @@ describe('dormant organization search runbook in PostgreSQL', () => {
343365
status: 'paused',
344366
lastSyncAt: null,
345367
listingCheckpoint: null,
368+
directoryCheckpoint: null,
346369
})
347370
const [kbRow] = await db.select().from(knowledgeBase).where(eq(knowledgeBase.id, indexId))
348371
expect(kbRow).toMatchObject({ isSearchIndex: true, deletedAt: null })

‎apps/sim/scripts/dormant-org-search/search-index-deletion-store.ts‎

Lines changed: 132 additions & 61 deletions
Original file line numberDiff line numberDiff line change
@@ -8,13 +8,15 @@ import {
88
knowledgeProjectionDirty,
99
outboxEvent,
1010
} from '@sim/db/schema'
11-
import { and, asc, count, eq, gt, inArray, isNotNull, isNull, sql } from 'drizzle-orm'
11+
import { and, asc, count, eq, gt, inArray, isNotNull, isNull, or, sql } from 'drizzle-orm'
1212
import type { DbTransaction } from '@/lib/db/types'
1313
import {
1414
enqueueKnowledgeStorageCleanup,
1515
KNOWLEDGE_STORAGE_CLEANUP_EVENT,
1616
} from '@/lib/knowledge/documents/storage-cleanup'
1717
import {
18+
evaluateDeletionGuard,
19+
SearchIndexDeletionRefused,
1820
type SearchIndexDeletionStore,
1921
STOPPED_CONNECTOR_STATUSES,
2022
} from '@/scripts/dormant-org-search/search-index-deletion'
@@ -32,6 +34,129 @@ async function enterBoundedTransaction(tx: DbTransaction, timeouts: DeletionTime
3234
)
3335
}
3436

37+
/**
38+
* Re-decides the deletion guard inside a deleting transaction, holding the base and its
39+
* connectors `FOR SHARE` until commit. A resume or a sync claim updates the connector row, so it
40+
* waits for this page to commit and the next page's guard refuses it: no page can delete
41+
* documents under a sync that started after the page-level guard ran.
42+
*/
43+
async function lockGuard(tx: DbTransaction, knowledgeBaseId: string) {
44+
const [base] = await tx
45+
.select({
46+
id: knowledgeBase.id,
47+
isSearchIndex: knowledgeBase.isSearchIndex,
48+
deletedAt: knowledgeBase.deletedAt,
49+
workspaceId: knowledgeBase.workspaceId,
50+
organizationId: knowledgeBase.organizationId,
51+
userId: knowledgeBase.userId,
52+
})
53+
.from(knowledgeBase)
54+
.where(eq(knowledgeBase.id, knowledgeBaseId))
55+
.for('share')
56+
.limit(1)
57+
const connectors = await tx
58+
.select({
59+
id: knowledgeConnector.id,
60+
status: knowledgeConnector.status,
61+
syncLockToken: knowledgeConnector.syncLockToken,
62+
memberSyncLockToken: knowledgeConnector.memberSyncLockToken,
63+
deletedAt: knowledgeConnector.deletedAt,
64+
detachedAt: knowledgeConnector.detachedAt,
65+
})
66+
.from(knowledgeConnector)
67+
.where(eq(knowledgeConnector.knowledgeBaseId, knowledgeBaseId))
68+
.orderBy(asc(knowledgeConnector.id))
69+
.for('share')
70+
const reasons = evaluateDeletionGuard(
71+
base ?? null,
72+
connectors.map((row) => ({
73+
id: row.id,
74+
status: row.status,
75+
syncLockHeld: row.syncLockToken !== null,
76+
memberSyncLockHeld: row.memberSyncLockToken !== null,
77+
deletedAt: row.deletedAt,
78+
detachedAt: row.detachedAt,
79+
}))
80+
)
81+
if (!base || reasons.length > 0) throw new SearchIndexDeletionRefused(reasons)
82+
return base
83+
}
84+
85+
/**
86+
* Makes every stopped connector of the base list its sources from scratch when it resumes: the
87+
* same columns the app clears when a connector must list everything again (an access mode
88+
* switch, a source change), plus the directory checkpoint, and every member made due. Runs in
89+
* each deleting transaction, so a run stopped partway never leaves a connector whose cursors
90+
* would skip the documents already deleted; rows already reset are left alone.
91+
*/
92+
async function resetCursors(tx: DbTransaction, knowledgeBaseId: string, now: Date) {
93+
const connectors = await tx
94+
.update(knowledgeConnector)
95+
.set({
96+
lastSyncAt: null,
97+
lastSyncDocCount: null,
98+
listingCheckpoint: null,
99+
directoryCheckpoint: null,
100+
memberTombstoneCursor: null,
101+
updatedAt: now,
102+
})
103+
.where(
104+
and(
105+
eq(knowledgeConnector.knowledgeBaseId, knowledgeBaseId),
106+
isNull(knowledgeConnector.deletedAt),
107+
isNull(knowledgeConnector.detachedAt),
108+
inArray(knowledgeConnector.status, [...STOPPED_CONNECTOR_STATUSES]),
109+
isNull(knowledgeConnector.syncLockToken),
110+
isNull(knowledgeConnector.memberSyncLockToken),
111+
or(
112+
isNotNull(knowledgeConnector.lastSyncAt),
113+
isNotNull(knowledgeConnector.lastSyncDocCount),
114+
isNotNull(knowledgeConnector.listingCheckpoint),
115+
isNotNull(knowledgeConnector.directoryCheckpoint),
116+
isNotNull(knowledgeConnector.memberTombstoneCursor)
117+
)
118+
)
119+
)
120+
.returning({ id: knowledgeConnector.id })
121+
const stopped = tx
122+
.select({ id: knowledgeConnector.id })
123+
.from(knowledgeConnector)
124+
.where(
125+
and(
126+
eq(knowledgeConnector.knowledgeBaseId, knowledgeBaseId),
127+
isNull(knowledgeConnector.deletedAt),
128+
isNull(knowledgeConnector.detachedAt),
129+
inArray(knowledgeConnector.status, [...STOPPED_CONNECTOR_STATUSES])
130+
)
131+
)
132+
const members = await tx
133+
.update(knowledgeConnectorMember)
134+
.set({
135+
listingCheckpoint: null,
136+
changeCursor: null,
137+
memberSyncedThrough: null,
138+
lastCompleteListingAt: null,
139+
lastListedCount: null,
140+
nextAttemptAt: null,
141+
updatedAt: now,
142+
})
143+
.where(
144+
and(
145+
inArray(knowledgeConnectorMember.connectorId, stopped),
146+
or(
147+
isNotNull(knowledgeConnectorMember.listingCheckpoint),
148+
isNotNull(knowledgeConnectorMember.changeCursor),
149+
isNotNull(knowledgeConnectorMember.memberSyncedThrough),
150+
isNotNull(knowledgeConnectorMember.lastCompleteListingAt),
151+
isNotNull(knowledgeConnectorMember.lastListedCount),
152+
isNotNull(knowledgeConnectorMember.nextAttemptAt)
153+
)
154+
)
155+
)
156+
.returning({ id: knowledgeConnectorMember.id })
157+
return { connectors: connectors.length, members: members.length }
158+
}
159+
35160
/**
36161
* The deletion store on the app's database client, so storage cleanup intents are queued by the
37162
* app's own `enqueueKnowledgeStorageCleanup` in the deleting transaction, and deleted by the app's
@@ -101,9 +226,10 @@ export function drizzleSearchIndexDeletionStore(
101226
return Number(row?.chunks ?? 0)
102227
},
103228

104-
async deleteChunkBatch(documentIds, limit) {
229+
async deleteChunkBatch(knowledgeBaseId, documentIds, limit) {
105230
return db.transaction(async (tx) => {
106231
await enterBoundedTransaction(tx, timeouts)
232+
await lockGuard(tx, knowledgeBaseId)
107233
const batch = tx
108234
.select({ id: embedding.id })
109235
.from(embedding)
@@ -117,23 +243,10 @@ export function drizzleSearchIndexDeletionStore(
117243
})
118244
},
119245

120-
async deleteDocuments(knowledgeBaseId, documentIds, requestId) {
246+
async deleteDocuments(knowledgeBaseId, documentIds, requestId, resetConnectors) {
121247
return db.transaction(async (tx) => {
122248
await enterBoundedTransaction(tx, timeouts)
123-
const [owner] = await tx
124-
.select({
125-
workspaceId: knowledgeBase.workspaceId,
126-
organizationId: knowledgeBase.organizationId,
127-
userId: knowledgeBase.userId,
128-
isSearchIndex: knowledgeBase.isSearchIndex,
129-
})
130-
.from(knowledgeBase)
131-
.where(eq(knowledgeBase.id, knowledgeBaseId))
132-
.for('share')
133-
.limit(1)
134-
if (!owner?.isSearchIndex) {
135-
throw new Error('The knowledge base stopped being a search index during the run')
136-
}
249+
const owner = await lockGuard(tx, knowledgeBaseId)
137250
/** Locks the documents against a late indexing commit between the chunk check and the delete. */
138251
const docs = await tx
139252
.select({ id: document.id, fileUrl: document.fileUrl })
@@ -171,6 +284,7 @@ export function drizzleSearchIndexDeletionStore(
171284
.delete(document)
172285
.where(inArray(document.id, ids))
173286
.returning({ id: document.id })
287+
if (resetConnectors) await resetCursors(tx, knowledgeBaseId, new Date())
174288
return {
175289
kind: 'deleted',
176290
deleted: deleted.length,
@@ -221,50 +335,7 @@ export function drizzleSearchIndexDeletionStore(
221335
async resetConnectorCursors(knowledgeBaseId) {
222336
return db.transaction(async (tx) => {
223337
await enterBoundedTransaction(tx, timeouts)
224-
const now = new Date()
225-
/**
226-
* The same columns the app clears when a connector must list everything again (an access
227-
* mode switch, a source change), limited to stopped connectors no sync holds.
228-
*/
229-
const connectors = await tx
230-
.update(knowledgeConnector)
231-
.set({
232-
lastSyncAt: null,
233-
lastSyncDocCount: null,
234-
listingCheckpoint: null,
235-
memberTombstoneCursor: null,
236-
updatedAt: now,
237-
})
238-
.where(
239-
and(
240-
eq(knowledgeConnector.knowledgeBaseId, knowledgeBaseId),
241-
isNull(knowledgeConnector.deletedAt),
242-
isNull(knowledgeConnector.detachedAt),
243-
inArray(knowledgeConnector.status, [...STOPPED_CONNECTOR_STATUSES]),
244-
isNull(knowledgeConnector.syncLockToken),
245-
isNull(knowledgeConnector.memberSyncLockToken)
246-
)
247-
)
248-
.returning({ id: knowledgeConnector.id })
249-
if (connectors.length === 0) return { connectors: 0, members: 0 }
250-
const members = await tx
251-
.update(knowledgeConnectorMember)
252-
.set({
253-
listingCheckpoint: null,
254-
changeCursor: null,
255-
memberSyncedThrough: null,
256-
lastCompleteListingAt: null,
257-
lastListedCount: null,
258-
updatedAt: now,
259-
})
260-
.where(
261-
inArray(
262-
knowledgeConnectorMember.connectorId,
263-
connectors.map((connector) => connector.id)
264-
)
265-
)
266-
.returning({ id: knowledgeConnectorMember.id })
267-
return { connectors: connectors.length, members: members.length }
338+
return resetCursors(tx, knowledgeBaseId, new Date())
268339
})
269340
},
270341
}

‎apps/sim/scripts/dormant-org-search/search-index-deletion.test.ts‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -65,7 +65,7 @@ class FakeStore implements SearchIndexDeletionStore {
6565
async countChunks(ids: readonly string[]) {
6666
return ids.reduce((sum, id) => sum + (this.documents.get(id)?.chunks ?? 0), 0)
6767
}
68-
async deleteChunkBatch(ids: readonly string[], limit: number) {
68+
async deleteChunkBatch(_knowledgeBaseId: string, ids: readonly string[], limit: number) {
6969
this.mutations.push('deleteChunkBatch')
7070
let deleted = 0
7171
for (const id of ids) {
@@ -78,7 +78,7 @@ class FakeStore implements SearchIndexDeletionStore {
7878
}
7979
return deleted
8080
}
81-
async deleteDocuments(_kb: string, ids: readonly string[]) {
81+
async deleteDocuments(_kb: string, ids: readonly string[], _requestId: string, _reset: boolean) {
8282
this.mutations.push('deleteDocuments')
8383
for (const id of ids) {
8484
const late = this.lateChunks.get(id)

‎apps/sim/scripts/dormant-org-search/search-index-deletion.ts‎

Lines changed: 18 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -70,15 +70,22 @@ export interface SearchIndexDeletionStore {
7070
/** Chunks of these documents; read only by a dry run. */
7171
countChunks(documentIds: readonly string[]): Promise<number>
7272
/** Deletes at most `limit` chunks of these documents in one transaction; returns how many. */
73-
deleteChunkBatch(documentIds: readonly string[], limit: number): Promise<number>
73+
deleteChunkBatch(
74+
knowledgeBaseId: string,
75+
documentIds: readonly string[],
76+
limit: number
77+
): Promise<number>
7478
/**
75-
* In one transaction: share-locks the knowledge base and re-checks it is a search index, locks
76-
* the documents, and, when none has a chunk left, queues their storage cleanup and deletes them.
79+
* In one transaction: share-locks the knowledge base and its connectors and re-decides the
80+
* guard, locks the documents, and, when none has a chunk left, queues their storage cleanup and
81+
* deletes them. With `resetConnectors`, the same transaction resets the connectors' listing
82+
* cursors, so a run stopped partway never leaves a connector that would skip deleted documents.
7783
*/
7884
deleteDocuments(
7985
knowledgeBaseId: string,
8086
documentIds: readonly string[],
81-
requestId: string
87+
requestId: string,
88+
resetConnectors: boolean
8289
): Promise<DeleteDocumentsOutcome>
8390
/** Pending storage cleanup events, counted up to `cap`. */
8491
pendingStorageCleanup(cap: number): Promise<number>
@@ -320,15 +327,20 @@ export async function deleteSearchIndexDocuments(
320327
) {
321328
for (;;) {
322329
const deleted = await retry('delete chunk batch', () =>
323-
store.deleteChunkBatch(documentIds, chunkBatchSize)
330+
store.deleteChunkBatch(knowledgeBaseId, documentIds, chunkBatchSize)
324331
)
325332
summary.chunksDeleted += deleted
326333
if (deleted === 0) break
327334
await pause()
328335
if (deleted < chunkBatchSize) break
329336
}
330337
outcome = await retry('delete documents', () =>
331-
store.deleteDocuments(knowledgeBaseId, documentIds, options.requestId)
338+
store.deleteDocuments(
339+
knowledgeBaseId,
340+
documentIds,
341+
options.requestId,
342+
options.resetConnectors !== false
343+
)
332344
)
333345
}
334346
if (outcome.kind === 'chunks-remain') {

0 commit comments

Comments
 (0)