Skip to content

Commit ef09857

Browse files
committed
improvement(knowledge): skip redundant document resurrection writes
1 parent 305da02 commit ef09857

2 files changed

Lines changed: 89 additions & 0 deletions

File tree

‎apps/sim/lib/knowledge/__integration__/listing-continuation.integration.ts‎

Lines changed: 88 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -166,6 +166,94 @@ describe('durable source and member cycles in PostgreSQL', () => {
166166
})
167167
}
168168

169+
it('restores rediscovered tombstones without rewriting already-live documents during resurrection', async () => {
170+
const connectorId = generateId()
171+
const runId = generateId()
172+
const [connector] = await db
173+
.insert(knowledgeConnector)
174+
.values({
175+
id: connectorId,
176+
knowledgeBaseId: ids.knowledgeBaseId,
177+
connectorType: 'google_drive',
178+
status: 'syncing',
179+
syncLockToken: runId,
180+
accessMode: 'admin',
181+
sourceConfig: {},
182+
})
183+
.returning()
184+
const deletedAt = new Date(Date.now() - 60_000)
185+
const documents = ['live', 'tombstoned'].map(sourceDoc)
186+
await db.insert(document).values(
187+
documents.map((item) => ({
188+
id: generateId(),
189+
knowledgeBaseId: ids.knowledgeBaseId,
190+
connectorId,
191+
externalId: item.externalId,
192+
filename: item.title,
193+
fileUrl: '',
194+
storageKey: `kb/${item.externalId}.txt`,
195+
fileSize: 10,
196+
mimeType: item.mimeType,
197+
contentHash: item.contentHash,
198+
processingStatus: 'completed',
199+
deletedAt: item.externalId === 'tombstoned' ? deletedAt : null,
200+
}))
201+
)
202+
let liveVersionBeforeResurrection: string | undefined
203+
const syncResult = result()
204+
const pass = await runConnectorContentPass({
205+
connectorId,
206+
connector,
207+
connectorConfig: {
208+
...CONNECTOR_REGISTRY.google_drive,
209+
listDocuments: async () => ({ documents, hasMore: false }),
210+
},
211+
sourceConfig: {},
212+
syncContext: {},
213+
kbOwner: { userId: ids.aliceId, workspaceId: ids.workspaceId },
214+
billingAttribution: billing,
215+
result: syncResult,
216+
lease: createContentSyncLease(connectorId, runId),
217+
leaseKind: 'content',
218+
runId,
219+
fingerprint: listingFingerprint({ source: 'resurrection' }),
220+
documentAccess: 'admin',
221+
getAccessToken: async () => 'fixture',
222+
hydration: { getDocument: async () => null },
223+
forceRehydrate: false,
224+
deadlineAt: Date.now() + 60_000,
225+
onPage: async () => {
226+
const rows = await db
227+
.select({
228+
externalId: document.externalId,
229+
deletedAt: document.deletedAt,
230+
version: sql<string>`xmin::text`,
231+
})
232+
.from(document)
233+
.where(eq(document.connectorId, connectorId))
234+
liveVersionBeforeResurrection = rows.find((row) => row.externalId === 'live')!.version
235+
expect(rows.find((row) => row.externalId === 'tombstoned')!.deletedAt).toEqual(deletedAt)
236+
return undefined
237+
},
238+
})
239+
const rows = await db
240+
.select({
241+
externalId: document.externalId,
242+
deletedAt: document.deletedAt,
243+
version: sql<string>`xmin::text`,
244+
})
245+
.from(document)
246+
.where(eq(document.connectorId, connectorId))
247+
expect(pass.complete).toBe(true)
248+
expect(syncResult).toEqual({ ...result(), docsUnchanged: 2 })
249+
expect(rows).toHaveLength(2)
250+
expect(rows.every((row) => row.deletedAt === null)).toBe(true)
251+
expect(liveVersionBeforeResurrection).toBeDefined()
252+
expect(rows.find((row) => row.externalId === 'live')!.version).toBe(
253+
liveVersionBeforeResurrection
254+
)
255+
})
256+
169257
it('refreshes verified existing permissions while repairing changed bodies, without indexing discoveries or reconciling absence', async () => {
170258
const connectorId = generateId()
171259
const runId = generateId()

‎apps/sim/lib/knowledge/connectors/sync-content-pass.ts‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -290,6 +290,7 @@ export async function runConnectorContentPass(input: ContentPassInput) {
290290
and(
291291
eq(document.connectorId, input.connectorId),
292292
inArray(document.externalId, verified.slice(offset, offset + 500)),
293+
isNotNull(document.deletedAt),
293294
isNotNull(document.contentHash),
294295
isNull(document.archivedAt)
295296
)

0 commit comments

Comments
 (0)