Skip to content

Commit a53830c

Browse files
committed
feat(knowledge): find a fill chain in flight by its tag before starting another
A start no longer relies on an idempotency key, which would have blocked a legitimate restart for its window and collided across options. Every run of a chain carries the chain's tag, continuations included, and a start lists in-flight runs by that tag first: a range whose chain is still running is left to it, whether the start is sliced or not, and a start after a chain ended starts anew.
1 parent 0de8533 commit a53830c

4 files changed

Lines changed: 92 additions & 31 deletions

File tree

‎apps/sim/background/projection-source-acl-backfill.ts‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -4,6 +4,7 @@ import {
44
PROJECTION_SOURCE_ACL_BACKFILL_SHARDS,
55
PROJECTION_SOURCE_ACL_BACKFILL_TASK_ID,
66
type ProjectionSourceAclBackfillPayload,
7+
projectionSourceAclChainTag,
78
runProjectionSourceAclBackfill,
89
} from '@/lib/knowledge/search/projection-source-acl-backfill'
910

@@ -36,6 +37,8 @@ export const projectionSourceAclBackfillTask = task({
3637
const continuation: ProjectionSourceAclBackfillPayload = { ...payload, cursor }
3738
await tasks.trigger(PROJECTION_SOURCE_ACL_BACKFILL_TASK_ID, continuation, {
3839
region: await resolveTriggerRegion(),
40+
/** The chain's tag rides on every continuation, so a start finds the chain wherever it is. */
41+
tags: [projectionSourceAclChainTag(payload.shard)],
3942
})
4043
},
4144
})

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

Lines changed: 49 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -3,15 +3,25 @@
33
*/
44
import { beforeEach, describe, expect, it, vi } from 'vitest'
55

6-
const { mockBackfill, mockEnd, mockPostgres, mockPrewarm, mockTasksTrigger, mockUnsafe } =
7-
vi.hoisted(() => ({
8-
mockBackfill: vi.fn(),
9-
mockEnd: vi.fn(async () => undefined),
10-
mockPostgres: vi.fn(),
11-
mockPrewarm: vi.fn(async () => []),
12-
mockTasksTrigger: vi.fn(async () => ({ id: 'run-1' })),
13-
mockUnsafe: vi.fn(async () => [{ unfilled: false }]),
14-
}))
6+
const {
7+
mockBackfill,
8+
mockEnd,
9+
mockPostgres,
10+
mockPrewarm,
11+
mockRunsList,
12+
mockTasksTrigger,
13+
mockUnsafe,
14+
} = vi.hoisted(() => ({
15+
mockBackfill: vi.fn(),
16+
mockEnd: vi.fn(async () => undefined),
17+
mockPostgres: vi.fn(),
18+
mockPrewarm: vi.fn(async () => []),
19+
mockRunsList: vi.fn(
20+
(_query: unknown): AsyncIterable<{ id: string }> => (async function* () {})()
21+
),
22+
mockTasksTrigger: vi.fn(async () => ({ id: 'run-1' })),
23+
mockUnsafe: vi.fn(async () => [{ unfilled: false }]),
24+
}))
1525

1626
vi.mock('@sim/db', () => ({ resolveDbUrl: () => 'postgres://localhost:5432/sim' }))
1727
vi.mock('@sim/db/script-migrations/0021_embedding_search_connector', () => ({
@@ -20,7 +30,10 @@ vi.mock('@sim/db/script-migrations/0021_embedding_search_connector', () => ({
2030
}))
2131
vi.mock('postgres', () => ({ default: mockPostgres }))
2232
vi.mock('@/lib/knowledge/search/prewarm', () => ({ prewarmSearchProjection: mockPrewarm }))
23-
vi.mock('@trigger.dev/sdk', () => ({ tasks: { trigger: mockTasksTrigger } }))
33+
vi.mock('@trigger.dev/sdk', () => ({
34+
runs: { list: mockRunsList },
35+
tasks: { trigger: mockTasksTrigger },
36+
}))
2437
vi.mock('@/lib/core/async-jobs/region', () => ({ resolveTriggerRegion: async () => 'us-east-1' }))
2538
vi.mock('@/lib/core/utils/background', () => ({
2639
runDetached: (_label: string, work: () => Promise<unknown>) => {
@@ -164,6 +177,8 @@ describe('projectionSourceAclShardRange', () => {
164177
describe('enqueueProjectionSourceAclBackfill', () => {
165178
beforeEach(() => {
166179
vi.clearAllMocks()
180+
/** No chain in flight unless a case says so. */
181+
mockRunsList.mockImplementation(() => (async function* () {})())
167182
mockPostgres.mockReturnValue(connection)
168183
mockBackfill.mockResolvedValue({
169184
projection: 'embedding_search',
@@ -177,28 +192,44 @@ describe('enqueueProjectionSourceAclBackfill', () => {
177192
it('hands the backfill to the Trigger.dev worker when one is configured', async () => {
178193
await expect(enqueueProjectionSourceAclBackfill({ pageSize: 25 })).resolves.toEqual({
179194
runIds: ['run-1'],
195+
inFlight: [],
180196
})
197+
expect(mockRunsList).toHaveBeenCalledWith(
198+
expect.objectContaining({ tag: 'projection-source-acl-backfill:shard:0/1' })
199+
)
181200
expect(mockTasksTrigger).toHaveBeenCalledWith(
182201
'projection-source-acl-backfill',
183202
{ pageSize: 25 },
184-
{
185-
region: 'us-east-1',
186-
idempotencyKey: 'projection-source-acl-backfill:1:0',
187-
idempotencyKeyTTL: '1h',
188-
}
203+
{ region: 'us-east-1', tags: ['projection-source-acl-backfill:shard:0/1'] }
189204
)
190205
expect(mockBackfill).not.toHaveBeenCalled()
191206
})
192207

193-
it('starts one run per shard, each on its own slice with its own idempotency key', async () => {
208+
it('leaves a range whose chain is still in flight to that chain', async () => {
209+
mockRunsList.mockImplementation((query: unknown) =>
210+
(async function* () {
211+
if ((query as { tag: string }).tag.endsWith(':shard:1/4')) yield { id: 'run-live' }
212+
})()
213+
)
214+
await expect(enqueueProjectionSourceAclBackfill({}, 4)).resolves.toEqual({
215+
runIds: ['run-1', 'run-1', 'run-1'],
216+
inFlight: ['run-live'],
217+
})
218+
expect(mockTasksTrigger.mock.calls.map(([, payload]) => payload.shard?.index)).toEqual([
219+
0, 2, 3,
220+
])
221+
})
222+
223+
it('starts one run per shard, each on its own slice under its own chain tag', async () => {
194224
await expect(enqueueProjectionSourceAclBackfill({ pageSize: 25 }, 4)).resolves.toEqual({
195225
runIds: ['run-1', 'run-1', 'run-1', 'run-1'],
226+
inFlight: [],
196227
})
197228
expect(mockTasksTrigger.mock.calls.map(([, payload]) => payload)).toEqual(
198229
[0, 1, 2, 3].map((index) => ({ pageSize: 25, shard: { index, count: 4 } }))
199230
)
200-
expect(mockTasksTrigger.mock.calls.map(([, , options]) => options.idempotencyKey)).toEqual(
201-
[0, 1, 2, 3].map((index) => `projection-source-acl-backfill:4:${index}`)
231+
expect(mockTasksTrigger.mock.calls.map(([, , options]) => options.tags)).toEqual(
232+
[0, 1, 2, 3].map((index) => [`projection-source-acl-backfill:shard:${index}/4`])
202233
)
203234
})
204235

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

Lines changed: 38 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -156,25 +156,38 @@ async function projectionsFilled(sql: postgres.Sql): Promise<boolean> {
156156
return true
157157
}
158158

159-
/** How long a start stays idempotent: long enough that a retried command finds its runs, not a second set. */
160-
const ENQUEUE_IDEMPOTENCY_TTL = '1h'
159+
/** The tag every run of one chain carries, so a chain in flight is found before another is started. */
160+
export function projectionSourceAclChainTag(shard?: ProjectionSourceAclBackfillShard): string {
161+
const { index, count } = shard ?? { index: 0, count: 1 }
162+
return `${PROJECTION_SOURCE_ACL_BACKFILL_TASK_ID}:shard:${index}/${count}`
163+
}
164+
165+
/** A run that has not ended: it, or the continuation it triggers, still owns its range. */
166+
const IN_FLIGHT_RUN_STATUSES = [
167+
'PENDING_VERSION',
168+
'QUEUED',
169+
'DEQUEUED',
170+
'EXECUTING',
171+
'WAITING',
172+
'DELAYED',
173+
] as const
161174

162175
/**
163176
* Starts the backfill on the deployment's Trigger.dev worker, where bounded runs chain until the
164177
* projections are filled: one chain over the whole id space, or one per shard, each filling its
165178
* own slice at the same time. Safe to call again at any time: a run only fills rows still unset,
166-
* and a start repeated within the hour finds the runs it already made rather than making more. A
167-
* cursor belongs to one chain, so a sliced start takes none: each slice begins at its own bound.
179+
* and a range whose chain is still in flight is left to that chain rather than given a second.
180+
* A cursor belongs to one chain, so a sliced start takes none: each slice begins at its own bound.
168181
*/
169182
export async function enqueueProjectionSourceAclBackfill(
170183
payload: ProjectionSourceAclBackfillPayload = {},
171184
shards = 1
172-
): Promise<{ runIds: string[] }> {
185+
): Promise<{ runIds: string[]; inFlight: string[] }> {
173186
if (shards !== 1) {
174187
assertProjectionSourceAclShard({ index: 0, count: shards })
175188
if (payload.cursor) throw new Error('A sliced projection backfill cannot start from a cursor')
176189
}
177-
const { tasks } = await import('@trigger.dev/sdk')
190+
const { runs, tasks } = await import('@trigger.dev/sdk')
178191
const region = await resolveTriggerRegion()
179192
const payloads: ProjectionSourceAclBackfillPayload[] =
180193
shards === 1
@@ -184,14 +197,28 @@ export async function enqueueProjectionSourceAclBackfill(
184197
shard: { index, count: shards },
185198
}))
186199
const runIds: string[] = []
187-
for (const [index, shardPayload] of payloads.entries()) {
200+
const inFlight: string[] = []
201+
for (const shardPayload of payloads) {
202+
const tag = projectionSourceAclChainTag(shardPayload.shard)
203+
let running: string | undefined
204+
for await (const run of runs.list({
205+
taskIdentifier: PROJECTION_SOURCE_ACL_BACKFILL_TASK_ID,
206+
tag,
207+
status: [...IN_FLIGHT_RUN_STATUSES],
208+
limit: 1,
209+
})) {
210+
running = run.id
211+
}
212+
if (running) {
213+
inFlight.push(running)
214+
continue
215+
}
188216
const handle = await tasks.trigger(PROJECTION_SOURCE_ACL_BACKFILL_TASK_ID, shardPayload, {
189217
region,
190-
idempotencyKey: `${PROJECTION_SOURCE_ACL_BACKFILL_TASK_ID}:${shards}:${index}`,
191-
idempotencyKeyTTL: ENQUEUE_IDEMPOTENCY_TTL,
218+
tags: [tag],
192219
})
193220
runIds.push(handle.id)
194221
}
195-
logger.info('Projection source and ACL backfill enqueued', { runIds, shards })
196-
return { runIds }
222+
logger.info('Projection source and ACL backfill enqueued', { runIds, inFlight, shards })
223+
return { runIds, inFlight }
197224
}

‎apps/sim/scripts/backfill-projection-source-acl.ts‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -41,8 +41,8 @@ function shardsFlag(): number {
4141
async function main(): Promise<void> {
4242
const shards = shardsFlag()
4343
if (isTriggerDevEnabled && env.TRIGGER_SECRET_KEY) {
44-
const handle = await enqueueProjectionSourceAclBackfill({}, shards)
45-
logger.info('Backfill enqueued on the Trigger.dev worker', handle)
44+
const started = await enqueueProjectionSourceAclBackfill({}, shards)
45+
logger.info('Backfill enqueued on the Trigger.dev worker', started)
4646
return
4747
}
4848
if (shards !== 1)

0 commit comments

Comments
 (0)