Skip to content

Commit 884cdd9

Browse files
committed
improvement(search): consolidate knowledge search around live Sim Search and document-decided workspace retrieval
- Move indexed organization search under lib/sim-search/indexed, dormant behind the single isIndexedOrgSearchEnabled() gate, with a check:indexed-org-search-boundary audit keeping callers on its public entry - Decide workspace knowledge base search on the document for every principal, so workspace retrieval never waits on the projector - Remove the knowledge-async-projection, knowledge-projection-fill and knowledge-tin-keyword flags and their code paths - Trim the projector to a single periodic sweep that releases workspace marks before deciding whether a pass is owed - Remove per-source vector index builds - Scope member sync and processing recovery to the indexed search path - Remove dead code left behind by the consolidation - Add an ops runbook and scripts for recovering and maintaining a dormant search index
1 parent 861b2cb commit 884cdd9

193 files changed

Lines changed: 7925 additions & 8851 deletions

File tree

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

‎apps/sim/app/api/cron/knowledge-projection/route.ts‎

Lines changed: 2 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -11,9 +11,8 @@ export const dynamic = 'force-dynamic'
1111
export const maxDuration = 60
1212

1313
/**
14-
* The knowledge projector's periodic sweep: enqueues one pass per window while there is work, and
15-
* returns once Trigger.dev accepts it. Writers ask for passes as they commit; this converges
16-
* whatever those requests missed.
14+
* The knowledge projector's periodic sweep: enqueues one pass per window while documents are
15+
* marked, and returns once Trigger.dev accepts it. It is the only thing that starts a pass.
1716
*/
1817
export const GET = withRouteHandler(async (request: NextRequest) => {
1918
const authError = verifyCronAuth(request, 'Knowledge projection sweep')
Lines changed: 79 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,79 @@
1+
/**
2+
* The member sync scheduler's reclaim against real PostgreSQL: a members-mode connector whose
3+
* member lease went stale is put back on the failure ladder, and a connector in any other access
4+
* mode is left exactly as it is, whatever its member columns say.
5+
*/
6+
7+
import { db } from '@sim/db'
8+
import { knowledgeConnector, organization, user, workspace } from '@sim/db/schema'
9+
import { createMockRequest } from '@sim/testing'
10+
import { generateId } from '@sim/utils/id'
11+
import { eq, inArray } from 'drizzle-orm'
12+
import { afterAll, beforeAll, describe, expect, it, vi } from 'vitest'
13+
14+
vi.mock('@/lib/auth/internal', () => ({ verifyCronAuth: () => null }))
15+
vi.mock('@/lib/knowledge/connectors/member-queue', async (importOriginal) => ({
16+
...(await importOriginal<typeof import('@/lib/knowledge/connectors/member-queue')>()),
17+
dispatchMemberSync: vi.fn(async () => undefined),
18+
}))
19+
20+
import {
21+
createKnowledgeAclFixtureIds,
22+
seedKnowledgeAclFixture,
23+
} from '@/lib/knowledge/__integration__/seed-source-access-fixture'
24+
import { GET } from '@/app/api/knowledge/connectors/member-sync/route'
25+
26+
describe('member sync reclaim in PostgreSQL', () => {
27+
const ids = createKnowledgeAclFixtureIds()
28+
const members = generateId()
29+
const admin = generateId()
30+
const staleLease = new Date(Date.now() - 24 * 60 * 60 * 1000)
31+
32+
beforeAll(async () => {
33+
await seedKnowledgeAclFixture(ids)
34+
await db.insert(knowledgeConnector).values(
35+
[
36+
{ id: members, accessMode: 'members' },
37+
{ id: admin, accessMode: 'admin' },
38+
].map(({ id, accessMode }) => ({
39+
id,
40+
knowledgeBaseId: ids.knowledgeBaseId,
41+
connectorType: 'google_drive',
42+
sourceConfig: {},
43+
accessMode,
44+
status: 'active',
45+
credentialId: ids.credentialId,
46+
memberSyncStatus: 'running',
47+
memberSyncLockToken: generateId(),
48+
memberSyncLockLeaseAt: staleLease,
49+
}))
50+
)
51+
})
52+
53+
afterAll(async () => {
54+
await db.delete(workspace).where(eq(workspace.id, ids.workspaceId))
55+
await db.delete(organization).where(eq(organization.id, ids.organizationId))
56+
await db.delete(user).where(inArray(user.id, [ids.aliceId, ids.bobId]))
57+
await db.$client.end()
58+
})
59+
60+
it('reclaims a stale members-mode lease and leaves every other access mode alone', async () => {
61+
const response = await GET(createMockRequest('GET'), undefined)
62+
expect(response.status).toBe(200)
63+
const rows = await db
64+
.select({
65+
id: knowledgeConnector.id,
66+
memberSyncStatus: knowledgeConnector.memberSyncStatus,
67+
memberSyncLockToken: knowledgeConnector.memberSyncLockToken,
68+
})
69+
.from(knowledgeConnector)
70+
.where(inArray(knowledgeConnector.id, [members, admin]))
71+
const byId = new Map(rows.map((row) => [row.id, row]))
72+
expect(byId.get(members)).toMatchObject({
73+
memberSyncStatus: 'error',
74+
memberSyncLockToken: null,
75+
})
76+
expect(byId.get(admin)?.memberSyncStatus).toBe('running')
77+
expect(byId.get(admin)?.memberSyncLockToken).not.toBeNull()
78+
})
79+
})

‎apps/sim/app/api/knowledge/connectors/member-sync/route.ts‎

Lines changed: 14 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -58,6 +58,18 @@ function reclaimedNextMemberSyncAt(): SQL {
5858
return sql`CASE WHEN ${reclaimedFailureCount()} >= ${MAX_CONSECUTIVE_FAILURES} THEN NULL ELSE now() + LEAST(${reclaimedFailureCount()} * ${CONNECTOR_FAILURE_BACKOFF_STEP_MINUTES}, ${CONNECTOR_FAILURE_BACKOFF_CAP_MINUTES}) * INTERVAL '1 minute' END`
5959
}
6060

61+
/**
62+
* Only the member engine takes the member lease, and only on a members-mode connector, which a
63+
* mode switch cannot leave while the lease is held; so both reclaims match `access_mode` too, the
64+
* predicate `kc_member_sync_due_idx` is partial on, and read that index instead of the table.
65+
*/
66+
function reclaimableMemberSync(status: 'running' | 'pending'): SQL | undefined {
67+
return and(
68+
eq(knowledgeConnector.accessMode, 'members'),
69+
eq(knowledgeConnector.memberSyncStatus, status)
70+
)
71+
}
72+
6173
/**
6274
* The write shared by both reclaims: a run that stopped making progress
6375
* re-enters the member failure ladder, which is the content engine's ladder
@@ -108,7 +120,7 @@ export const GET = withRouteHandler(async (request: NextRequest) => {
108120
.set(reclaimPayload(STALE_LOCK_ERROR_MESSAGE))
109121
.where(
110122
and(
111-
eq(knowledgeConnector.memberSyncStatus, 'running'),
123+
reclaimableMemberSync('running'),
112124
sql`${memberSyncLockLease()} <= ${sql.param(staleCutoff, knowledgeConnector.memberSyncLockLeaseAt)}`,
113125
isNull(knowledgeConnector.archivedAt),
114126
isNull(knowledgeConnector.deletedAt)
@@ -120,7 +132,7 @@ export const GET = withRouteHandler(async (request: NextRequest) => {
120132
.set(reclaimPayload(LOST_DISPATCH_ERROR_MESSAGE))
121133
.where(
122134
and(
123-
eq(knowledgeConnector.memberSyncStatus, 'pending'),
135+
reclaimableMemberSync('pending'),
124136
sql`${memberSyncLockLease()} <= ${sql.param(staleCutoff, knowledgeConnector.memberSyncLockLeaseAt)}`,
125137
isNull(knowledgeConnector.archivedAt),
126138
isNull(knowledgeConnector.deletedAt)

‎apps/sim/app/api/knowledge/github/installations/route.test.ts‎

Lines changed: 0 additions & 168 deletions
This file was deleted.

‎apps/sim/app/api/knowledge/github/installations/route.ts‎

Lines changed: 0 additions & 53 deletions
This file was deleted.

‎apps/sim/app/api/knowledge/member-connectors/route.ts‎

Lines changed: 0 additions & 22 deletions
This file was deleted.

‎apps/sim/app/api/knowledge/search/route.ts‎

Lines changed: 4 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -4,12 +4,12 @@ import {
44
internalRateLimits,
55
internalSessionAuth,
66
} from '@/lib/api/server/routes'
7-
import { isLiveEnterpriseSearchEnabled } from '@/lib/core/config/env-flags'
87
import { internalKnowledgeErrorPolicies } from '@/lib/knowledge/api/route-policies'
98
import { knowledgeOperations } from '@/lib/knowledge/application/operations'
10-
import { searchScopedKnowledge } from '@/lib/knowledge/application/workspace-search'
119
import { DEFAULT_RERANKER_MODEL } from '@/lib/knowledge/reranker-models'
1210
import { sourceAuthor } from '@/lib/knowledge/search/author'
11+
import { searchScopedKnowledge } from '@/lib/sim-search/indexed'
12+
import { isIndexedOrgSearchEnabled } from '@/lib/sim-search/indexed/gate'
1313
import { searchLiveKnowledge } from '@/lib/sim-search/live/application'
1414

1515
const DIRECT_SEARCH_VECTOR_BUDGET_MS = 3000
@@ -81,4 +81,5 @@ const liveSearchRoute = defineInternalJsonRoute({
8181
present: (data) => ({ success: true as const, data }),
8282
})
8383

84-
export const POST = isLiveEnterpriseSearchEnabled ? liveSearchRoute : indexedSearchRoute
84+
/** Indexed organization search is dormant unless its gate is on; Live Search serves otherwise. */
85+
export const POST = isIndexedOrgSearchEnabled() ? indexedSearchRoute : liveSearchRoute

0 commit comments

Comments
 (0)