Skip to content

Commit 6309546

Browse files
committed
fix(knowledge): keep projection guards out of historical migrations
Restores 0016, 0019 and 0021 to their staging bodies; 0024 alone re-creates the projection triggers with the mode guard, in the same transaction that installs the marks. The fill starts only inside the pass budget and reports what it marked as remaining. Passes dispatch to Trigger.dev by the rule document processing uses, and the inline coalescer clears its running flag in the same step it reads the owed flag. Bumps the Helm chart for the projection cron job.
1 parent 70f3dfb commit 6309546

14 files changed

Lines changed: 739 additions & 132 deletions

‎apps/sim/lib/knowledge/__integration__/knowledge-projection.integration.ts‎

Lines changed: 61 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -548,6 +548,67 @@ describe('the projector', () => {
548548
expect((await rowAcl(embeddingSearch))?.acl).toEqual(aclOf('alice'))
549549
})
550550

551+
it('keeps the mark when a synchronous ACL change commits under a page that then writes stale values', async () => {
552+
/** A pending revocation of Bob: the rows still name both until a pass. */
553+
await write('async', (tx) =>
554+
tx
555+
.update(document)
556+
.set({ acl: aclOf('alice') })
557+
.where(eq(document.id, documentId))
558+
)
559+
const before = await markOf()
560+
const [{ pid }] = await projector<Array<{ pid: number }>>`SELECT pg_backend_pid() AS pid`
561+
/**
562+
* A synchronous writer revokes Alice too and holds its transaction open: its fan-out has
563+
* locked the projection rows and its mark bump is not yet visible.
564+
*/
565+
const writer = postgres(process.env.DATABASE_URL!, { max: 1, onnotice: () => undefined })
566+
let release: () => void = () => {}
567+
const held = new Promise<void>((resolve) => {
568+
release = resolve
569+
})
570+
let locked: () => void = () => {}
571+
const fannedOut = new Promise<void>((resolve) => {
572+
locked = resolve
573+
})
574+
try {
575+
const writing = writer.begin(async (tx) => {
576+
await tx`UPDATE document SET acl = ${aclOf('bob')} WHERE id = ${documentId}`
577+
locked()
578+
await held
579+
})
580+
await fannedOut
581+
/**
582+
* The pass reads the committed document, Alice alone, and its page blocks on the rows the
583+
* writer holds. Once the writer commits, the page rewrites those rows from what it read.
584+
*/
585+
const pass = project()
586+
await vi.waitFor(async () => {
587+
const [activity] = await db.execute<{ waiting: boolean }>(
588+
sql`SELECT wait_event_type = 'Lock' AS waiting FROM pg_stat_activity WHERE pid = ${pid}`
589+
)
590+
expect(activity?.waiting).toBe(true)
591+
})
592+
release()
593+
await writing
594+
await pass
595+
} finally {
596+
release()
597+
await writer.end()
598+
}
599+
/** The page wrote the value it read over the writer's newer one. */
600+
expect((await rowAcl(embeddingSearch))?.acl).toEqual(aclOf('alice'))
601+
/** The writer's bump is what the pass could not settle over: the rows stay decided on the document. */
602+
expect((await markOf())?.generation).toBe((before?.generation ?? 0) + 1)
603+
expect(await admitted()).toEqual({ vector: [], keyword: [] })
604+
expect(await vectorIds()).toEqual([])
605+
await project()
606+
expect(await markOf()).toBeUndefined()
607+
expect((await rowAcl(embeddingSearch))?.acl).toEqual(aclOf('bob'))
608+
expect((await rowAcl(embeddingKeywordTin))?.acl).toEqual(aclOf('bob'))
609+
expect(await admitted()).toEqual({ vector: [], keyword: [] })
610+
})
611+
551612
it('passes over a document deleted while it is being projected', async () => {
552613
const doomed = generateId()
553614
await db.insert(document).values({

‎apps/sim/lib/knowledge/projection/enqueue-inline.test.ts‎

Lines changed: 24 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
/**
22
* @vitest-environment node
33
*/
4+
import { sleep } from '@sim/utils/helpers'
45
import { describe, expect, it, vi } from 'vitest'
56

67
const mocks = vi.hoisted(() => ({ runPass: vi.fn(), trigger: vi.fn() }))
@@ -10,6 +11,7 @@ vi.mock('@/lib/core/config/feature-flags', () => ({ isFeatureEnabled: async () =
1011

1112
vi.mock('@trigger.dev/sdk', () => ({ tasks: { trigger: mocks.trigger } }))
1213
vi.mock('@/lib/core/config/env-flags', () => ({ isTriggerDevEnabled: false }))
14+
vi.mock('@/lib/core/config/trigger-runtime', () => ({ isInsideTriggerRun: () => false }))
1315
vi.mock('@/lib/knowledge/projection/run', () => ({ runKnowledgeProjectionPass: mocks.runPass }))
1416

1517
import {
@@ -48,6 +50,28 @@ describe('knowledge projection without a Trigger.dev worker', () => {
4850
expect(mocks.trigger).not.toHaveBeenCalled()
4951
})
5052

53+
it('runs a pass for a request that arrives at any point while the last pass is finishing', async () => {
54+
for (let hops = 0; hops < 8; hops++) {
55+
mocks.runPass.mockReset()
56+
/** The pass asks for another once it has settled, `hops` microtasks later. */
57+
mocks.runPass
58+
.mockImplementationOnce(() => {
59+
void (async () => {
60+
for (let hop = 0; hop < hops; hop++) await Promise.resolve()
61+
await requestKnowledgeProjection()
62+
})()
63+
return Promise.resolve()
64+
})
65+
.mockResolvedValue(undefined)
66+
await requestKnowledgeProjection()
67+
await vi.waitFor(() => expect(mocks.runPass.mock.calls.length).toBeGreaterThanOrEqual(2), {
68+
timeout: 200,
69+
})
70+
/** Lets the loop finish before the next interleaving starts. */
71+
await sleep(1)
72+
}
73+
})
74+
5175
it('logs a failed pass and still runs the pass owed after it', async () => {
5276
mocks.runPass.mockReset()
5377
const failFirst = heldPass()

‎apps/sim/lib/knowledge/projection/enqueue.test.ts‎

Lines changed: 26 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -9,6 +9,8 @@ const mocks = vi.hoisted(() => ({
99
trigger: vi.fn(),
1010
execute: vi.fn(),
1111
isFeatureEnabled: vi.fn(),
12+
env: { TRIGGER_SECRET_KEY: 'fixture-key' as string | undefined },
13+
insideRun: vi.fn(),
1214
}))
1315

1416
vi.mock('@sim/db', () => ({ db: { execute: mocks.execute } }))
@@ -17,6 +19,8 @@ vi.mock('@/lib/core/config/feature-flags', () => ({ isFeatureEnabled: mocks.isFe
1719
vi.mock('@trigger.dev/sdk', () => ({ tasks: { trigger: mocks.trigger } }))
1820
vi.mock('@/lib/core/async-jobs/region', () => ({ resolveTriggerRegion: mocks.resolveRegion }))
1921
vi.mock('@/lib/core/config/env-flags', () => ({ isTriggerDevEnabled: true }))
22+
vi.mock('@/lib/core/config/env', () => ({ env: mocks.env }))
23+
vi.mock('@/lib/core/config/trigger-runtime', () => ({ isInsideTriggerRun: mocks.insideRun }))
2024
vi.mock('@/lib/knowledge/projection/run', () => ({ runKnowledgeProjectionPass: mocks.runPass }))
2125

2226
import {
@@ -33,6 +37,8 @@ describe('knowledge projection enqueue', () => {
3337
mocks.trigger.mockResolvedValue({ id: 'run-1' })
3438
mocks.execute.mockResolvedValue([{ pending: true }])
3539
mocks.isFeatureEnabled.mockResolvedValue(false)
40+
mocks.env.TRIGGER_SECRET_KEY = 'fixture-key'
41+
mocks.insideRun.mockReturnValue(false)
3642
})
3743

3844
afterEach(() => {
@@ -90,6 +96,26 @@ describe('knowledge projection enqueue', () => {
9096
expect(mocks.trigger).toHaveBeenCalledTimes(2)
9197
})
9298

99+
it('runs the pass in this process where Trigger.dev is enabled without its secret key', async () => {
100+
mocks.env.TRIGGER_SECRET_KEY = undefined
101+
mocks.runPass.mockResolvedValue(undefined)
102+
await expect(enqueueKnowledgeProjectionSweep()).resolves.toEqual({
103+
triggered: true,
104+
backend: 'inline',
105+
jobId: null,
106+
})
107+
expect(mocks.trigger).not.toHaveBeenCalled()
108+
await vi.waitFor(() => expect(mocks.runPass).toHaveBeenCalledTimes(1))
109+
})
110+
111+
it('enqueues from inside a Trigger.dev run whatever the environment holds', async () => {
112+
mocks.env.TRIGGER_SECRET_KEY = undefined
113+
mocks.insideRun.mockReturnValue(true)
114+
await expect(enqueueKnowledgeProjectionSweep()).resolves.toMatchObject({
115+
backend: 'trigger-dev',
116+
})
117+
})
118+
93119
it('never fails the write that asked when the request is refused', async () => {
94120
vi.advanceTimersByTime(60_000)
95121
mocks.trigger.mockRejectedValueOnce(new Error('trigger unavailable'))

‎apps/sim/lib/knowledge/projection/enqueue.ts‎

Lines changed: 29 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -3,8 +3,10 @@ import { SOURCE_ACL_PROJECTIONS } from '@sim/db/knowledge-projection'
33
import { createLogger } from '@sim/logger'
44
import { getErrorMessage } from '@sim/utils/errors'
55
import { sql } from 'drizzle-orm'
6+
import { env } from '@/lib/core/config/env'
67
import { isTriggerDevEnabled } from '@/lib/core/config/env-flags'
78
import { isFeatureEnabled } from '@/lib/core/config/feature-flags'
9+
import { isInsideTriggerRun } from '@/lib/core/config/trigger-runtime'
810

911
const logger = createLogger('KnowledgeProjectionEnqueue')
1012

@@ -31,35 +33,48 @@ const PROMPT_REQUEST_INTERVAL_MS = 5_000
3133

3234
let lastPromptAt = 0
3335

34-
/** The inline pass in flight when no Trigger.dev worker is configured, and whether another is owed. */
35-
let inlinePass: Promise<void> | undefined
36+
/** Whether an inline pass is running when no Trigger.dev worker is configured, and whether another is owed. */
37+
let inlineRunning = false
3638
let inlinePassOwed = false
3739

3840
/**
3941
* Starts a pass in this process without waiting for it: at most one runs at a time, a request while
4042
* one runs is folded into a single pass after it, and a failed pass is logged without dropping one
41-
* owed after it. The
42-
* pass module loads on first use. For deployments without a Trigger.dev worker, whose sweep and
43-
* writes run passes here, as their document processing does.
43+
* owed after it. The loop reads the owed flag and clears the running flag in the same synchronous
44+
* step, so a request can never land between the two and be dropped. The pass module loads on first
45+
* use. For deployments without a Trigger.dev worker, whose sweep and writes run passes here, as
46+
* their document processing does.
4447
*/
4548
function runInline(): void {
46-
if (inlinePass) {
49+
if (inlineRunning) {
4750
inlinePassOwed = true
4851
return
4952
}
50-
inlinePass = (async () => {
51-
do {
53+
inlineRunning = true
54+
void (async () => {
55+
for (;;) {
5256
inlinePassOwed = false
5357
try {
5458
const { runKnowledgeProjectionPass } = await import('@/lib/knowledge/projection/run')
5559
await runKnowledgeProjectionPass({ budgetMs: KNOWLEDGE_PROJECTION_PASS_BUDGET_MS })
5660
} catch (error) {
5761
logger.error('Inline knowledge projection pass failed', { error: getErrorMessage(error) })
5862
}
59-
} while (inlinePassOwed)
60-
})().finally(() => {
61-
inlinePass = undefined
62-
})
63+
if (!inlinePassOwed) {
64+
inlineRunning = false
65+
return
66+
}
67+
}
68+
})()
69+
}
70+
71+
/**
72+
* Whether passes run on Trigger.dev, by the rule document processing dispatches with: inside a
73+
* Trigger.dev run always, and otherwise only where Trigger.dev is enabled and the secret key the
74+
* SDK authenticates with is set.
75+
*/
76+
function projectsOnTrigger(): boolean {
77+
return isInsideTriggerRun() || Boolean(isTriggerDevEnabled && env.TRIGGER_SECRET_KEY)
6378
}
6479

6580
/**
@@ -70,7 +85,7 @@ function runInline(): void {
7085
* leaves the marks to the sweep, which keeps enqueueing a pass every minute until one runs.
7186
*/
7287
export async function requestKnowledgeProjection(): Promise<void> {
73-
if (!isTriggerDevEnabled) {
88+
if (!projectsOnTrigger()) {
7489
runInline()
7590
return
7691
}
@@ -125,7 +140,7 @@ export async function enqueueKnowledgeProjectionSweep(): Promise<KnowledgeProjec
125140
if (!(await hasKnowledgeProjectionWork())) {
126141
return { triggered: false, backend: null, jobId: null }
127142
}
128-
if (!isTriggerDevEnabled) {
143+
if (!projectsOnTrigger()) {
129144
runInline()
130145
return { triggered: true, backend: 'inline', jobId: null }
131146
}

‎apps/sim/lib/knowledge/projection/run.test.ts‎

Lines changed: 35 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,7 @@
11
/**
22
* @vitest-environment node
33
*/
4-
import { beforeEach, describe, expect, it, vi } from 'vitest'
4+
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
55

66
const mocks = vi.hoisted(() => ({
77
runProjection: vi.fn(),
@@ -106,6 +106,40 @@ describe('runKnowledgeProjectionPass', () => {
106106
})
107107
})
108108

109+
describe('at the budget', () => {
110+
beforeEach(() => {
111+
vi.useFakeTimers()
112+
mocks.isFeatureEnabled.mockResolvedValue(true)
113+
})
114+
115+
afterEach(() => {
116+
vi.useRealTimers()
117+
})
118+
119+
it('does not start the fill once the budget is spent', async () => {
120+
mocks.runProjection.mockImplementation(async () => {
121+
vi.advanceTimersByTime(60_001)
122+
return drained
123+
})
124+
await expect(runKnowledgeProjectionPass({ budgetMs: 60_000 })).resolves.toMatchObject({
125+
remaining: false,
126+
filled: 0,
127+
})
128+
expect(mocks.markUnfilled).not.toHaveBeenCalled()
129+
})
130+
131+
it('reports documents the fill marked but no round settled as remaining', async () => {
132+
mocks.markUnfilled.mockImplementation(async () => {
133+
vi.advanceTimersByTime(60_001)
134+
return { marked: 2, cursor: { projection: 0, afterId: 'row-2' } }
135+
})
136+
await expect(runKnowledgeProjectionPass({ budgetMs: 60_000 })).resolves.toMatchObject({
137+
remaining: true,
138+
filled: 2,
139+
})
140+
})
141+
})
142+
109143
it('does not fill while marks remain, so writers are converged first', async () => {
110144
mocks.isFeatureEnabled.mockResolvedValue(true)
111145
mocks.runProjection.mockResolvedValue({ ...drained, settled: 0, remaining: true })

‎apps/sim/lib/knowledge/projection/run.ts‎

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -106,10 +106,12 @@ export async function runKnowledgeProjectionPass(options: {
106106
result.deferred = round.reduce((most, progress) => Math.max(most, progress.deferred), 0)
107107
break
108108
}
109-
if (fillCursor === null) break
109+
/** The fill starts only inside the budget, and what it marks is left for a round to settle. */
110+
if (fillCursor === null || Date.now() >= deadline) break
110111
const fill = await markUnfilledProjectionDocuments(sessions[0], fillCursor)
111112
result.filled += fill.marked
112113
fillCursor = fill.cursor
114+
if (fill.marked > 0) result.remaining = true
113115
if (fill.marked === 0 && fillCursor === null) break
114116
}
115117
logger.info('Knowledge projection pass finished', result)

‎helm/sim/Chart.yaml‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -2,7 +2,7 @@ apiVersion: v2
22
name: sim
33
description: A Helm chart for Sim - the open-source AI workspace where teams build, deploy, and manage AI agents
44
type: application
5-
version: 1.11.3
5+
version: 1.11.4
66
appVersion: "v0.8.26"
77
kubeVersion: ">=1.25.0-0"
88
home: https://sim.ai

‎packages/db/knowledge-projection.ts‎

Lines changed: 7 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -60,18 +60,18 @@ const SOURCE_VECTOR_WIDTHS = [1536, 384, 768, 1024, 3072] as const
6060
const widthColumn = (name: string, width: number) => (width === 1536 ? name : `${name}_${width}`)
6161

6262
/** The halfvec columns of `embedding_search`, in width order. */
63-
export const SEARCH_VECTOR_COLUMNS = SEARCH_VECTOR_WIDTHS.map((width) =>
64-
widthColumn('vector', width)
65-
)
63+
const SEARCH_VECTOR_COLUMNS = SEARCH_VECTOR_WIDTHS.map((width) => widthColumn('vector', width))
6664

6765
/** The bit columns of `embedding_search`, whose width check requires exactly one. */
6866
const SEARCH_BINARY_COLUMNS = ['"binary"', 'binary_384', 'binary_768', 'binary_1024', 'binary_3072']
6967

7068
/**
7169
* The halfvec projections of an embedding row, in {@link SEARCH_VECTOR_COLUMNS} order. Shortening
72-
* is valid only for the two OpenAI models trained for prefix retrieval.
70+
* is valid only for the two OpenAI models trained for prefix retrieval. These and the bit
71+
* projections match what `sync_embedding_search()` from `0016_backfill_search_vectors` writes, so
72+
* a pass over rows a synchronous writer wrote finds nothing to rewrite.
7373
*/
74-
export function searchVectorProjections(prefix: string, shortened: string): string {
74+
function searchVectorProjections(prefix: string, shortened: string): string {
7575
return SEARCH_VECTOR_WIDTHS.map((width) =>
7676
width === 512
7777
? `CASE WHEN ${shortened} THEN subvector(coalesce(${SOURCE_VECTOR_WIDTHS.map((size) => `${prefix}.${widthColumn('embedding', size)}`).join(', ')}), 1, 512)::halfvec(512) END`
@@ -80,14 +80,14 @@ export function searchVectorProjections(prefix: string, shortened: string): stri
8080
}
8181

8282
/** The bit projections of an embedding row, in {@link SEARCH_BINARY_COLUMNS} order. */
83-
export function searchBinaryProjections(prefix: string): string {
83+
function searchBinaryProjections(prefix: string): string {
8484
return `binary_quantize(${prefix}.embedding)::bit(1536), binary_quantize(${prefix}.embedding_384)::bit(384),
8585
binary_quantize(${prefix}.embedding_768)::bit(768), binary_quantize(${prefix}.embedding_1024)::bit(1024),
8686
binary_quantize(${prefix}.embedding_3072)::bit(3072)`
8787
}
8888

8989
/** Whether a knowledge base's model is shortened for the 512 projection, for an embedding row. */
90-
export function searchVectorShortened(model: string, prefix: string): string {
90+
function searchVectorShortened(model: string, prefix: string): string {
9191
return `${model} IN ${SHORTENED_EMBEDDING_MODELS} AND ${prefix}.embedding_384 IS NULL`
9292
}
9393

0 commit comments

Comments
 (0)