Skip to content

Commit 526f659

Browse files
committed
fix(knowledge): start the projection backfill the way the table backfill is started
The outbox event and its slice-and-defer handler are gone; nothing else in the repo starts background work that way. The fill is started from the app side as the table backfill is: tasks.trigger on the Trigger.dev worker when one is configured, runDetached in-process otherwise, with the manual script as the operator path after a deploy. The registered migration is now 0022_projection_source_acl_backfill, which supersedes 0021 so a database that already recorded the synchronous shape still gains the unfilled indexes. The page statement timeout is exported so a caller bounding a run can leave it as headroom. Claude-Session: https://claude.ai/code/session_01XU6c7pKRpa5CMoMHKDdqxX
1 parent ecbda39 commit 526f659

11 files changed

Lines changed: 84 additions & 188 deletions

‎apps/sim/lib/core/outbox/processor.test.ts‎

Lines changed: 0 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -33,9 +33,6 @@ vi.mock('@/lib/knowledge/application/slack-search/outbox', () => ({
3333
vi.mock('@/lib/knowledge/documents/processing-outbox-handler', () => ({
3434
knowledgeDocumentProcessingOutboxHandlers: {},
3535
}))
36-
vi.mock('@/lib/knowledge/search/projection-source-acl-backfill', () => ({
37-
projectionSourceAclBackfillOutboxHandlers: {},
38-
}))
3936
vi.mock('@/lib/mothership/inbox/cleanup-outbox', () => ({ inboxCleanupOutboxHandlers: {} }))
4037
vi.mock('@/lib/organizations/resource-cleanup', () => ({
4138
organizationResourceCleanupOutboxHandlers: {},

‎apps/sim/lib/core/outbox/processor.ts‎

Lines changed: 0 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -18,7 +18,6 @@ import { slackSearchOutboxHandlers } from '@/lib/knowledge/application/slack-sea
1818
import { getConnectorFailureDiagnostic } from '@/lib/knowledge/connectors/connector-error'
1919
import { knowledgeDocumentProcessingOutboxHandlers } from '@/lib/knowledge/documents/processing-outbox-handler'
2020
import { recoverKnowledgeDocumentProcessing } from '@/lib/knowledge/documents/processing-recovery'
21-
import { projectionSourceAclBackfillOutboxHandlers } from '@/lib/knowledge/search/projection-source-acl-backfill'
2221
import { inboxCleanupOutboxHandlers } from '@/lib/mothership/inbox/cleanup-outbox'
2322
import { organizationResourceCleanupOutboxHandlers } from '@/lib/organizations/resource-cleanup'
2423
import { workspaceFileLiveDocOutboxHandlers } from '@/lib/uploads/contexts/workspace/workspace-file-live-doc-outbox'
@@ -43,7 +42,6 @@ const handlers = {
4342
...invitationMigrationOutboxHandlers,
4443
...directGrantOutboxHandlers,
4544
...knowledgeDocumentProcessingOutboxHandlers,
46-
...projectionSourceAclBackfillOutboxHandlers,
4745
...organizationResourceCleanupOutboxHandlers,
4846
...inboxCleanupOutboxHandlers,
4947
...permissionAccessRequestOutboxHandlers,

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

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

6-
const { mockBackfill, mockEnd, mockPostgres, mockTasksTrigger, envState } = vi.hoisted(() => ({
6+
const { mockBackfill, mockEnd, mockPostgres, mockTasksTrigger } = vi.hoisted(() => ({
77
mockBackfill: vi.fn(),
88
mockEnd: vi.fn(async () => undefined),
99
mockPostgres: vi.fn(),
1010
mockTasksTrigger: vi.fn(async () => ({ id: 'run-1' })),
11-
envState: { triggerEnabled: false, secret: undefined as string | undefined },
1211
}))
1312

1413
vi.mock('@sim/db', () => ({ resolveDbUrl: () => 'postgres://localhost:5432/sim' }))
1514
vi.mock('@sim/db/script-migrations/0021_embedding_search_connector', () => ({
1615
PROJECTION_SOURCE_ACL_TABLES: ['embedding_search', 'embedding_keyword_tin'],
17-
PROJECTION_SOURCE_ACL_BACKFILL_OUTBOX_EVENT: 'knowledge.projection.source_acl.backfill',
1816
backfillProjectionSourceAcl: mockBackfill,
1917
}))
2018
vi.mock('postgres', () => ({ default: mockPostgres }))
2119
vi.mock('@trigger.dev/sdk', () => ({ tasks: { trigger: mockTasksTrigger } }))
2220
vi.mock('@/lib/core/async-jobs/region', () => ({ resolveTriggerRegion: async () => 'us-east-1' }))
23-
vi.mock('@/lib/core/config/env', () => ({
24-
env: {
25-
get TRIGGER_SECRET_KEY() {
26-
return envState.secret
27-
},
21+
vi.mock('@/lib/core/utils/background', () => ({
22+
runDetached: (_label: string, work: () => Promise<unknown>) => {
23+
void work()
2824
},
2925
}))
30-
vi.mock('@/lib/core/config/env-flags', () => ({
31-
get isTriggerDevEnabled() {
32-
return envState.triggerEnabled
33-
},
34-
}))
35-
vi.mock('@/lib/core/outbox/service', () => ({
36-
continueOutboxHandler: (reason: string) => ({
37-
outcome: 'deferred',
38-
reason,
39-
consumeAttempt: false,
40-
}),
41-
withOutboxHandlerTimeout: (handler: unknown, timeoutMs: number) =>
42-
Object.assign(handler as object, { timeoutMs }),
43-
}))
4426

4527
import {
4628
enqueueProjectionSourceAclBackfill,
47-
projectionSourceAclBackfillOutboxHandlers,
4829
runProjectionSourceAclBackfill,
4930
} from '@/lib/knowledge/search/projection-source-acl-backfill'
5031

5132
const connection = { end: mockEnd }
52-
const handler =
53-
projectionSourceAclBackfillOutboxHandlers['knowledge.projection.source_acl.backfill']
54-
const context = { eventId: 'e', eventType: 'knowledge.projection.source_acl.backfill' }
5533

5634
describe('runProjectionSourceAclBackfill', () => {
5735
beforeEach(() => {
@@ -112,7 +90,7 @@ describe('runProjectionSourceAclBackfill', () => {
11290
})
11391
})
11492

115-
describe('the outbox event the migration leaves behind', () => {
93+
describe('enqueueProjectionSourceAclBackfill', () => {
11694
beforeEach(() => {
11795
vi.clearAllMocks()
11896
mockPostgres.mockReturnValue(connection)
@@ -125,40 +103,21 @@ describe('the outbox event the migration leaves behind', () => {
125103
})
126104
})
127105

128-
it('hands the backfill to the Trigger.dev worker when there is one', async () => {
129-
envState.triggerEnabled = true
130-
envState.secret = 'tr_secret'
131-
await expect(enqueueProjectionSourceAclBackfill({ pageSize: 25 })).resolves.toEqual({
106+
it('hands the backfill to the Trigger.dev worker when one is configured', async () => {
107+
await expect(enqueueProjectionSourceAclBackfill({ pageSize: 25 }, true)).resolves.toEqual({
132108
runId: 'run-1',
133109
})
134110
expect(mockTasksTrigger).toHaveBeenCalledWith(
135111
'projection-source-acl-backfill',
136112
{ pageSize: 25 },
137113
{ region: 'us-east-1' }
138114
)
139-
await expect(handler({}, context as never)).resolves.toBeUndefined()
140-
expect(mockTasksTrigger).toHaveBeenCalledTimes(2)
141115
expect(mockBackfill).not.toHaveBeenCalled()
142116
})
143117

144-
it('runs a bounded slice per outbox run without one, and stays pending until it is filled', async () => {
145-
envState.triggerEnabled = false
146-
envState.secret = undefined
147-
mockBackfill.mockResolvedValueOnce({
148-
projection: 'embedding_search',
149-
scanned: 100,
150-
written: 100,
151-
afterId: 'chunk-100',
152-
done: false,
153-
})
154-
await expect(handler({}, context as never)).resolves.toMatchObject({
155-
outcome: 'deferred',
156-
consumeAttempt: false,
157-
})
118+
it('fills the projections detached in this process without one', async () => {
119+
await expect(enqueueProjectionSourceAclBackfill({}, false)).resolves.toBeNull()
158120
expect(mockTasksTrigger).not.toHaveBeenCalled()
159-
expect(mockBackfill).toHaveBeenCalledTimes(1)
160-
expect(mockBackfill.mock.calls[0][2].budgetMs).toBeLessThanOrEqual(handler.timeoutMs!)
161-
await expect(handler({}, context as never)).resolves.toBeUndefined()
162-
expect(mockBackfill).toHaveBeenCalledTimes(3)
121+
await vi.waitFor(() => expect(mockBackfill).toHaveBeenCalledTimes(2))
163122
})
164123
})

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

Lines changed: 17 additions & 56 deletions
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,6 @@
11
import { resolveDbUrl } from '@sim/db'
22
import {
33
backfillProjectionSourceAcl,
4-
PROJECTION_SOURCE_ACL_BACKFILL_OUTBOX_EVENT,
54
PROJECTION_SOURCE_ACL_TABLES,
65
type ProjectionSourceAclTable,
76
} from '@sim/db/script-migrations/0021_embedding_search_connector'
@@ -10,13 +9,7 @@ import postgres from 'postgres'
109
import { resolveTriggerRegion } from '@/lib/core/async-jobs/region'
1110
import { env } from '@/lib/core/config/env'
1211
import { isTriggerDevEnabled } from '@/lib/core/config/env-flags'
13-
import {
14-
continueOutboxHandler,
15-
type DeferredOutboxHandlerResult,
16-
type OutboxEventContext,
17-
type OutboxHandlerRegistry,
18-
withOutboxHandlerTimeout,
19-
} from '@/lib/core/outbox/service'
12+
import { runDetached } from '@/lib/core/utils/background'
2013

2114
const logger = createLogger('ProjectionSourceAclBackfill')
2215

@@ -82,56 +75,24 @@ export async function runProjectionSourceAclBackfill(
8275
}
8376
}
8477

85-
/** Whether the deployment has a Trigger.dev worker to hand the backfill to. */
86-
export function projectionSourceAclBackfillUsesTrigger(): boolean {
87-
return Boolean(isTriggerDevEnabled && env.TRIGGER_SECRET_KEY)
88-
}
89-
9078
/**
9179
* Starts the backfill the way the table backfill is started: on the deployment's Trigger.dev
92-
* worker, where bounded runs chain until both projections are filled. Safe to call again at any
93-
* time — a run only fills rows still unset.
80+
* worker when one is configured, where bounded runs chain until both projections are filled, and
81+
* detached in this process otherwise. Safe to call again at any time — a run only fills rows still
82+
* unset.
9483
*/
9584
export async function enqueueProjectionSourceAclBackfill(
96-
payload: ProjectionSourceAclBackfillPayload = {}
97-
): Promise<{ runId: string }> {
98-
const { tasks } = await import('@trigger.dev/sdk')
99-
const handle = await tasks.trigger(PROJECTION_SOURCE_ACL_BACKFILL_TASK_ID, payload, {
100-
region: await resolveTriggerRegion(),
101-
})
102-
logger.info('Projection source and ACL backfill enqueued', { runId: handle.id })
103-
return { runId: handle.id }
104-
}
105-
106-
/** One outbox run's share of the backfill on a deployment without a worker, inside its window. */
107-
const OUTBOX_RUN_BUDGET_MS = 4 * 60 * 1000
108-
const OUTBOX_HANDLER_TIMEOUT_MS = 5 * 60 * 1000
109-
110-
/**
111-
* Handles the event script migration `0021_embedding_search_connector` leaves behind, so the
112-
* backfill starts once the app that ships this handler is up, on every deployment. With a
113-
* Trigger.dev worker the event is done once the task is enqueued: the task owns its own retries
114-
* and continuations. Without one the backfill runs here, one bounded slice per outbox run, and the
115-
* event stays pending until both projections are filled — every slice's writes are durable, and a
116-
* slice that starts over after a restart skips the filled rows through the unfilled index, so an
117-
* interrupted deployment loses nothing but the time of one slice.
118-
*/
119-
export const projectionSourceAclBackfillOutboxHandlers: OutboxHandlerRegistry = {
120-
[PROJECTION_SOURCE_ACL_BACKFILL_OUTBOX_EVENT]: withOutboxHandlerTimeout(
121-
async (
122-
_payload: unknown,
123-
_context: OutboxEventContext
124-
): Promise<undefined | DeferredOutboxHandlerResult> => {
125-
if (projectionSourceAclBackfillUsesTrigger()) {
126-
await enqueueProjectionSourceAclBackfill()
127-
return undefined
128-
}
129-
const cursor = await runProjectionSourceAclBackfill({}, { budgetMs: OUTBOX_RUN_BUDGET_MS })
130-
if (!cursor) return undefined
131-
return continueOutboxHandler(
132-
`projection source and ACL backfill paused at ${cursor.projection} after ${cursor.afterId}`
133-
)
134-
},
135-
OUTBOX_HANDLER_TIMEOUT_MS
136-
),
85+
payload: ProjectionSourceAclBackfillPayload = {},
86+
useTrigger = Boolean(isTriggerDevEnabled && env.TRIGGER_SECRET_KEY)
87+
): Promise<{ runId: string } | null> {
88+
if (useTrigger) {
89+
const { tasks } = await import('@trigger.dev/sdk')
90+
const handle = await tasks.trigger(PROJECTION_SOURCE_ACL_BACKFILL_TASK_ID, payload, {
91+
region: await resolveTriggerRegion(),
92+
})
93+
logger.info('Projection source and ACL backfill enqueued', { runId: handle.id })
94+
return { runId: handle.id }
95+
}
96+
runDetached(PROJECTION_SOURCE_ACL_BACKFILL_TASK_ID, () => runProjectionSourceAclBackfill(payload))
97+
return null
13798
}

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

Lines changed: 10 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -2,30 +2,31 @@
22

33
/**
44
* Starts the projection source and ACL backfill that script migration
5-
* `0021_embedding_search_connector` leaves to the background, by hand. The migration's outbox
6-
* event already starts it after a deploy; this is for starting it again — on a deployment with
7-
* Trigger.dev it enqueues the `projection-source-acl-backfill` task, which chains bounded runs
8-
* until both ranking projections are filled; without one it fills them here, paced the same way.
9-
* Safe to run at any time — a run only fills rows still unset.
5+
* `0022_projection_source_acl_backfill` leaves to the background. On a deployment with Trigger.dev
6+
* it enqueues the `projection-source-acl-backfill` task, which chains bounded runs until both
7+
* ranking projections are filled; without one it fills them here, paced the same way. Safe to run
8+
* again at any time — a run only fills rows still unset.
109
*
1110
* Usage:
1211
* bun apps/sim/scripts/backfill-projection-source-acl.ts
1312
*/
1413

1514
import { createLogger } from '@sim/logger'
1615
import { toError } from '@sim/utils/errors'
16+
import { env } from '@/lib/core/config/env'
17+
import { isTriggerDevEnabled } from '@/lib/core/config/env-flags'
1718
import {
1819
enqueueProjectionSourceAclBackfill,
19-
projectionSourceAclBackfillUsesTrigger,
2020
runProjectionSourceAclBackfill,
2121
} from '@/lib/knowledge/search/projection-source-acl-backfill'
2222

2323
const logger = createLogger('BackfillProjectionSourceAcl')
2424

25+
/** A script has no long-lived process to detach into, so without a worker it fills inline. */
2526
async function main(): Promise<void> {
26-
if (projectionSourceAclBackfillUsesTrigger()) {
27-
const handle = await enqueueProjectionSourceAclBackfill()
28-
logger.info('Backfill enqueued on the Trigger.dev worker', handle)
27+
if (isTriggerDevEnabled && env.TRIGGER_SECRET_KEY) {
28+
const handle = await enqueueProjectionSourceAclBackfill({}, true)
29+
logger.info('Backfill enqueued on the Trigger.dev worker', handle ?? {})
2930
return
3031
}
3132
await runProjectionSourceAclBackfill({})

‎packages/db/script-migrations-paused-billing-attribution.test.ts‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -453,7 +453,7 @@ describe('script migration registry', () => {
453453
'0017_index_search_documents',
454454
'0018_repair_workspace_file_content_revision',
455455
'0019_tin_keyword_projection',
456-
'0021_embedding_search_connector',
456+
'0022_projection_source_acl_backfill',
457457
])
458458
})
459459
})

‎packages/db/script-migrations/0016_backfill_search_vectors.postgres.test.ts‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -392,6 +392,7 @@ describe.runIf(Boolean(databaseUrl))('search projection upgrade in PostgreSQL',
392392
{ name: '0018_repair_workspace_file_content_revision' },
393393
{ name: '0019_tin_keyword_projection' },
394394
{ name: '0021_embedding_search_connector' },
395+
{ name: '0022_projection_source_acl_backfill' },
395396
])
396397
const [{ complete }] = await sql`SELECT count(*)::int AS complete FROM embedding e
397398
JOIN embedding_search s ON s.id = e.id JOIN embedding_keyword_search k ON k.id = e.id

‎packages/db/script-migrations/0021_embedding_search_connector.postgres.test.ts‎

Lines changed: 3 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,5 @@
1-
import {
2-
backfillProjectionSourceAcl,
3-
embeddingSearchConnectorMigration,
4-
} from '@sim/db/script-migrations/0021_embedding_search_connector'
1+
import { backfillProjectionSourceAcl } from '@sim/db/script-migrations/0021_embedding_search_connector'
2+
import { projectionSourceAclBackfillMigration as embeddingSearchConnectorMigration } from '@sim/db/script-migrations/0022_projection_source_acl_backfill'
53
import { generateId } from '@sim/utils/id'
64
import postgres, { type Sql } from 'postgres'
75
import { afterAll, beforeAll, beforeEach, describe, expect, it } from 'vitest'
@@ -40,10 +38,6 @@ describe.runIf(Boolean(databaseUrl))('projection source and ACL backfill in Post
4038
await sql`CREATE TABLE document (
4139
id text PRIMARY KEY, connector_id text, acl text[] NOT NULL DEFAULT '{ws}'
4240
)`
43-
await sql`CREATE TABLE outbox_event (
44-
id text PRIMARY KEY, event_type text NOT NULL, payload json NOT NULL,
45-
status text NOT NULL DEFAULT 'pending'
46-
)`
4741
for (const projection of ['embedding_search', 'embedding_keyword_tin']) {
4842
await sql`CREATE TABLE ${sql(projection)} (
4943
id text PRIMARY KEY, document_id text NOT NULL, enabled boolean NOT NULL DEFAULT true,
@@ -65,7 +59,7 @@ describe.runIf(Boolean(databaseUrl))('projection source and ACL backfill in Post
6559
await sql`ALTER TABLE embedding_keyword_tin DISABLE TRIGGER embedding_keyword_tin_source_acl_set`
6660
})
6761

68-
it('installs its triggers, indexes and start event again without failing, so a cut-short deploy completes', async () => {
62+
it('installs its triggers and indexes again without failing, so a cut-short deploy completes', async () => {
6963
await expect(embeddingSearchConnectorMigration.up(sql)).resolves.toBeUndefined()
7064
const indexes = await sql<{ indexname: string }[]>`
7165
SELECT indexname FROM pg_indexes WHERE schemaname = ${schemaName} ORDER BY indexname`
@@ -78,14 +72,6 @@ describe.runIf(Boolean(databaseUrl))('projection source and ACL backfill in Post
7872
'embedding_search_source_idx',
7973
])
8074
)
81-
/** Two runs of the migration leave the app one event, so the backfill is started once. */
82-
expect(await sql`SELECT id, event_type, status FROM outbox_event`).toEqual([
83-
{
84-
id: 'projection-source-acl-backfill:0021',
85-
event_type: 'knowledge.projection.source_acl.backfill',
86-
status: 'pending',
87-
},
88-
])
8975
})
9076

9177
it('fills only the rows still unset, in pages, and leaves a chunk that changed documents to its trigger', async () => {

0 commit comments

Comments
 (0)