|
| 1 | +/** |
| 2 | + * @vitest-environment node |
| 3 | + */ |
| 4 | +import { beforeEach, describe, expect, it, vi } from 'vitest' |
| 5 | + |
| 6 | +const { mockBackfill, mockEnd, mockPostgres, mockTasksTrigger } = vi.hoisted(() => ({ |
| 7 | + mockBackfill: vi.fn(), |
| 8 | + mockEnd: vi.fn(async () => undefined), |
| 9 | + mockPostgres: vi.fn(), |
| 10 | + mockTasksTrigger: vi.fn(async () => ({ id: 'run-1' })), |
| 11 | +})) |
| 12 | + |
| 13 | +vi.mock('@sim/db', () => ({ resolveDbUrl: () => 'postgres://localhost:5432/sim' })) |
| 14 | +vi.mock('@sim/db/script-migrations/0021_embedding_search_connector', () => ({ |
| 15 | + PROJECTION_SOURCE_ACL_TABLES: ['embedding_search', 'embedding_keyword_tin'], |
| 16 | + backfillProjectionSourceAcl: mockBackfill, |
| 17 | +})) |
| 18 | +vi.mock('postgres', () => ({ default: mockPostgres })) |
| 19 | +vi.mock('@trigger.dev/sdk', () => ({ tasks: { trigger: mockTasksTrigger } })) |
| 20 | +vi.mock('@/lib/core/async-jobs/region', () => ({ resolveTriggerRegion: async () => 'us-east-1' })) |
| 21 | +vi.mock('@/lib/core/utils/background', () => ({ |
| 22 | + runDetached: (_label: string, work: () => Promise<unknown>) => { |
| 23 | + void work() |
| 24 | + }, |
| 25 | +})) |
| 26 | + |
| 27 | +import { |
| 28 | + enqueueProjectionSourceAclBackfill, |
| 29 | + runProjectionSourceAclBackfill, |
| 30 | +} from '@/lib/knowledge/search/projection-source-acl-backfill' |
| 31 | + |
| 32 | +const connection = { end: mockEnd } |
| 33 | + |
| 34 | +describe('runProjectionSourceAclBackfill', () => { |
| 35 | + beforeEach(() => { |
| 36 | + vi.clearAllMocks() |
| 37 | + mockPostgres.mockReturnValue(connection) |
| 38 | + mockBackfill.mockImplementation(async (_sql, projection, options) => ({ |
| 39 | + projection, |
| 40 | + scanned: 0, |
| 41 | + written: 0, |
| 42 | + afterId: options?.afterId ?? '', |
| 43 | + done: true, |
| 44 | + })) |
| 45 | + }) |
| 46 | + |
| 47 | + it('fills both projections in order on its own connection and closes it', async () => { |
| 48 | + await expect(runProjectionSourceAclBackfill({ pageSize: 50, pauseMs: 10 })).resolves.toBeNull() |
| 49 | + expect(mockBackfill.mock.calls.map(([, projection]) => projection)).toEqual([ |
| 50 | + 'embedding_search', |
| 51 | + 'embedding_keyword_tin', |
| 52 | + ]) |
| 53 | + for (const [sql, , options] of mockBackfill.mock.calls) { |
| 54 | + expect(sql).toBe(connection) |
| 55 | + expect(options).toMatchObject({ afterId: undefined, pageSize: 50, pauseMs: 10 }) |
| 56 | + } |
| 57 | + expect(mockEnd).toHaveBeenCalledTimes(1) |
| 58 | + }) |
| 59 | + |
| 60 | + it('resumes after the cursor in its projection and from the start of the next', async () => { |
| 61 | + await runProjectionSourceAclBackfill({ |
| 62 | + cursor: { projection: 'embedding_keyword_tin', afterId: 'chunk-9' }, |
| 63 | + }) |
| 64 | + expect(mockBackfill).toHaveBeenCalledTimes(1) |
| 65 | + expect(mockBackfill.mock.calls[0][1]).toBe('embedding_keyword_tin') |
| 66 | + expect(mockBackfill.mock.calls[0][2]).toMatchObject({ afterId: 'chunk-9' }) |
| 67 | + }) |
| 68 | + |
| 69 | + it('returns where a budgeted run stopped so the next run can carry on', async () => { |
| 70 | + mockBackfill.mockResolvedValueOnce({ |
| 71 | + projection: 'embedding_search', |
| 72 | + scanned: 100, |
| 73 | + written: 100, |
| 74 | + afterId: 'chunk-100', |
| 75 | + done: false, |
| 76 | + }) |
| 77 | + await expect(runProjectionSourceAclBackfill({}, { budgetMs: 1000 })).resolves.toEqual({ |
| 78 | + projection: 'embedding_search', |
| 79 | + afterId: 'chunk-100', |
| 80 | + }) |
| 81 | + expect(mockBackfill).toHaveBeenCalledTimes(1) |
| 82 | + expect(mockBackfill.mock.calls[0][2].budgetMs).toBeLessThanOrEqual(1000) |
| 83 | + expect(mockEnd).toHaveBeenCalledTimes(1) |
| 84 | + }) |
| 85 | + |
| 86 | + it('closes the connection when a page fails', async () => { |
| 87 | + mockBackfill.mockRejectedValueOnce(new Error('canceling statement due to statement timeout')) |
| 88 | + await expect(runProjectionSourceAclBackfill({})).rejects.toThrow('statement timeout') |
| 89 | + expect(mockEnd).toHaveBeenCalledTimes(1) |
| 90 | + }) |
| 91 | +}) |
| 92 | + |
| 93 | +describe('enqueueProjectionSourceAclBackfill', () => { |
| 94 | + beforeEach(() => { |
| 95 | + vi.clearAllMocks() |
| 96 | + mockPostgres.mockReturnValue(connection) |
| 97 | + mockBackfill.mockResolvedValue({ |
| 98 | + projection: 'embedding_search', |
| 99 | + scanned: 0, |
| 100 | + written: 0, |
| 101 | + afterId: '', |
| 102 | + done: true, |
| 103 | + }) |
| 104 | + }) |
| 105 | + |
| 106 | + it('hands the backfill to the Trigger.dev worker when one is configured', async () => { |
| 107 | + await expect(enqueueProjectionSourceAclBackfill({ pageSize: 25 }, true)).resolves.toEqual({ |
| 108 | + runId: 'run-1', |
| 109 | + }) |
| 110 | + expect(mockTasksTrigger).toHaveBeenCalledWith( |
| 111 | + 'projection-source-acl-backfill', |
| 112 | + { pageSize: 25 }, |
| 113 | + { region: 'us-east-1' } |
| 114 | + ) |
| 115 | + expect(mockBackfill).not.toHaveBeenCalled() |
| 116 | + }) |
| 117 | + |
| 118 | + it('fills the projections detached in this process without one', async () => { |
| 119 | + await expect(enqueueProjectionSourceAclBackfill({}, false)).resolves.toBeNull() |
| 120 | + expect(mockTasksTrigger).not.toHaveBeenCalled() |
| 121 | + await vi.waitFor(() => expect(mockBackfill).toHaveBeenCalledTimes(2)) |
| 122 | + }) |
| 123 | +}) |
0 commit comments