Skip to content

Commit 41db9c3

Browse files
committed
feat(knowledge): key a fill start on the chain it follows, and count only rows the fill can finish
A start's idempotency key is the chain's latest run, so two starts that saw the same state collapse into one while a start after another chain ended is its own; a start validates any shard it is handed; and the completion probe counts only unfilled rows whose document exists, since a row whose document is gone is not the fill's to finish.
1 parent 1cf2483 commit 41db9c3

2 files changed

Lines changed: 53 additions & 16 deletions

File tree

‎apps/sim/lib/knowledge/search/projection-source-acl-backfill.test.ts‎

Lines changed: 32 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -17,7 +17,7 @@ const {
1717
mockPostgres: vi.fn(),
1818
mockPrewarm: vi.fn(async () => []),
1919
mockRunsList: vi.fn(
20-
(_query: unknown): AsyncIterable<{ id: string }> => (async function* () {})()
20+
(_query: unknown): AsyncIterable<{ id: string; status: string }> => (async function* () {})()
2121
),
2222
mockTasksTrigger: vi.fn(async () => ({ id: 'run-1' })),
2323
mockUnsafe: vi.fn(async () => [{ unfilled: false }]),
@@ -78,6 +78,12 @@ describe('runProjectionSourceAclBackfill', () => {
7878

7979
it('analyzes and warms the projections on the same connection once both are filled, before closing it', async () => {
8080
await runProjectionSourceAclBackfill({})
81+
/** A row whose document is gone is not the fill's to finish; the probe joins the document. */
82+
expect(
83+
mockUnsafe.mock.calls.some(([query]) =>
84+
String(query).includes('JOIN document d ON d.id = s.document_id WHERE s.acl IS NULL')
85+
)
86+
).toBe(true)
8187
expect(mockUnsafe.mock.calls.map(([query]) => query)).toEqual(
8288
expect.arrayContaining(['ANALYZE embedding_search', 'ANALYZE embedding_keyword_tin'])
8389
)
@@ -203,17 +209,40 @@ describe('enqueueProjectionSourceAclBackfill', () => {
203209
{
204210
region: 'us-east-1',
205211
tags: ['projection-source-acl-backfill:shard:0/1'],
206-
idempotencyKey: 'projection-source-acl-backfill:shard:0/1',
212+
idempotencyKey: 'projection-source-acl-backfill:shard:0/1:after:none',
207213
idempotencyKeyTTL: '2m',
208214
}
209215
)
210216
expect(mockBackfill).not.toHaveBeenCalled()
211217
})
212218

219+
it('keys a start after a chain that ended on that chain, so a restart is its own start', async () => {
220+
mockRunsList.mockImplementation(() =>
221+
(async function* () {
222+
yield { id: 'run-done', status: 'COMPLETED' }
223+
})()
224+
)
225+
await expect(enqueueProjectionSourceAclBackfill({})).resolves.toEqual({
226+
runIds: ['run-1'],
227+
inFlight: [],
228+
})
229+
expect(mockTasksTrigger.mock.calls[0][2].idempotencyKey).toBe(
230+
'projection-source-acl-backfill:shard:0/1:after:run-done'
231+
)
232+
})
233+
234+
it('refuses a shard the id space cannot be sliced into before starting anything', async () => {
235+
await expect(
236+
enqueueProjectionSourceAclBackfill({ shard: { index: 5, count: 4 } })
237+
).rejects.toThrow('shard index must be within 0..3')
238+
expect(mockTasksTrigger).not.toHaveBeenCalled()
239+
})
240+
213241
it('leaves a range whose chain is still in flight to that chain', async () => {
214242
mockRunsList.mockImplementation((query: unknown) =>
215243
(async function* () {
216-
if ((query as { tag: string }).tag.endsWith(':shard:1/4')) yield { id: 'run-live' }
244+
if ((query as { tag: string }).tag.endsWith(':shard:1/4'))
245+
yield { id: 'run-live', status: 'EXECUTING' }
217246
})()
218247
)
219248
await expect(enqueueProjectionSourceAclBackfill({}, 4)).resolves.toEqual({

‎apps/sim/lib/knowledge/search/projection-source-acl-backfill.ts‎

Lines changed: 21 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -145,11 +145,17 @@ export async function runProjectionSourceAclBackfill(
145145
}
146146
}
147147

148-
/** Whether no projection still holds a row without its source and ACL; each read is one index probe. */
148+
/**
149+
* Whether no projection still holds a row the fill could give its source and ACL: a row without
150+
* them whose document exists. A row whose document is gone is not the fill's to finish and never
151+
* counts as left. Each read is one index probe while any such row remains.
152+
*/
149153
async function projectionsFilled(sql: postgres.Sql): Promise<boolean> {
150154
for (const projection of PROJECTION_SOURCE_ACL_TABLES) {
151155
const [row] = await sql.unsafe<Array<{ unfilled: boolean }>>(
152-
`SELECT EXISTS (SELECT 1 FROM ${projection} WHERE acl IS NULL) AS unfilled`
156+
`SELECT EXISTS (
157+
SELECT 1 FROM ${projection} s JOIN document d ON d.id = s.document_id WHERE s.acl IS NULL
158+
) AS unfilled`
153159
)
154160
if (row?.unfilled) return false
155161
}
@@ -163,21 +169,22 @@ export function projectionSourceAclChainTag(shard?: ProjectionSourceAclBackfillS
163169
}
164170

165171
/**
166-
* How long a start's trigger stays idempotent. The in-flight lookup and the trigger are two
167-
* calls, so two starts in the same instant could both find no chain; a key that lives just past
168-
* that instant closes the gap without holding a later, legitimate restart.
172+
* How long a start's trigger stays idempotent. The lookup and the trigger are two calls, so two
173+
* starts in the same instant could both find no chain in flight; the key is what they both saw,
174+
* the chain's latest run, so they collapse into one start, while a start after another chain has
175+
* ended sees a different latest run and is a new key.
169176
*/
170177
const START_IDEMPOTENCY_TTL = '2m'
171178

172179
/** A run that has not ended: it, or the continuation it triggers, still owns its range. */
173-
const IN_FLIGHT_RUN_STATUSES = [
180+
const IN_FLIGHT_RUN_STATUSES: ReadonlySet<string> = new Set([
174181
'PENDING_VERSION',
175182
'QUEUED',
176183
'DEQUEUED',
177184
'EXECUTING',
178185
'WAITING',
179186
'DELAYED',
180-
] as const
187+
])
181188

182189
/**
183190
* Starts the backfill on the deployment's Trigger.dev worker, where bounded runs chain until the
@@ -190,6 +197,7 @@ export async function enqueueProjectionSourceAclBackfill(
190197
payload: ProjectionSourceAclBackfillPayload = {},
191198
shards = 1
192199
): Promise<{ runIds: string[]; inFlight: string[] }> {
200+
if (payload.shard) assertProjectionSourceAclShard(payload.shard)
193201
if (shards !== 1) {
194202
assertProjectionSourceAclShard({ index: 0, count: shards })
195203
if (payload.cursor) throw new Error('A sliced projection backfill cannot start from a cursor')
@@ -207,23 +215,23 @@ export async function enqueueProjectionSourceAclBackfill(
207215
const inFlight: string[] = []
208216
for (const shardPayload of payloads) {
209217
const tag = projectionSourceAclChainTag(shardPayload.shard)
210-
let running: string | undefined
218+
/** The chain's latest run, newest first, whatever its state. */
219+
let latest: { id: string; status: string } | undefined
211220
for await (const run of runs.list({
212221
taskIdentifier: PROJECTION_SOURCE_ACL_BACKFILL_TASK_ID,
213222
tag,
214-
status: [...IN_FLIGHT_RUN_STATUSES],
215223
limit: 1,
216224
})) {
217-
running = run.id
225+
latest = run
218226
}
219-
if (running) {
220-
inFlight.push(running)
227+
if (latest && IN_FLIGHT_RUN_STATUSES.has(latest.status)) {
228+
inFlight.push(latest.id)
221229
continue
222230
}
223231
const handle = await tasks.trigger(PROJECTION_SOURCE_ACL_BACKFILL_TASK_ID, shardPayload, {
224232
region,
225233
tags: [tag],
226-
idempotencyKey: tag,
234+
idempotencyKey: `${tag}:after:${latest?.id ?? 'none'}`,
227235
idempotencyKeyTTL: START_IDEMPOTENCY_TTL,
228236
})
229237
runIds.push(handle.id)

0 commit comments

Comments
 (0)