Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 2 additions & 3 deletions apps/sim/app/api/cron/knowledge-projection/route.ts
Original file line number Diff line number Diff line change
Expand Up @@ -11,9 +11,8 @@ export const dynamic = 'force-dynamic'
export const maxDuration = 60

/**
* The knowledge projector's periodic sweep: enqueues one pass per window while there is work, and
* returns once Trigger.dev accepts it. Writers ask for passes as they commit; this converges
* whatever those requests missed.
* The knowledge projector's periodic sweep: enqueues one pass per window while documents are
* marked, and returns once Trigger.dev accepts it. It is the only thing that starts a pass.
*/
export const GET = withRouteHandler(async (request: NextRequest) => {
const authError = verifyCronAuth(request, 'Knowledge projection sweep')
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,79 @@
/**
* The member sync scheduler's reclaim against real PostgreSQL: a members-mode connector whose
* member lease went stale is put back on the failure ladder, and a connector in any other access
* mode is left exactly as it is, whatever its member columns say.
*/

import { db } from '@sim/db'
import { knowledgeConnector, organization, user, workspace } from '@sim/db/schema'
import { createMockRequest } from '@sim/testing'
import { generateId } from '@sim/utils/id'
import { eq, inArray } from 'drizzle-orm'
import { afterAll, beforeAll, describe, expect, it, vi } from 'vitest'

vi.mock('@/lib/auth/internal', () => ({ verifyCronAuth: () => null }))
vi.mock('@/lib/knowledge/connectors/member-queue', async (importOriginal) => ({
...(await importOriginal<typeof import('@/lib/knowledge/connectors/member-queue')>()),
dispatchMemberSync: vi.fn(async () => undefined),
}))

import {
createKnowledgeAclFixtureIds,
seedKnowledgeAclFixture,
} from '@/lib/knowledge/__integration__/seed-source-access-fixture'
import { GET } from '@/app/api/knowledge/connectors/member-sync/route'

describe('member sync reclaim in PostgreSQL', () => {
const ids = createKnowledgeAclFixtureIds()
const members = generateId()
const admin = generateId()
const staleLease = new Date(Date.now() - 24 * 60 * 60 * 1000)

beforeAll(async () => {
await seedKnowledgeAclFixture(ids)
await db.insert(knowledgeConnector).values(
[
{ id: members, accessMode: 'members' },
{ id: admin, accessMode: 'admin' },
].map(({ id, accessMode }) => ({
id,
knowledgeBaseId: ids.knowledgeBaseId,
connectorType: 'google_drive',
sourceConfig: {},
accessMode,
status: 'active',
credentialId: ids.credentialId,
memberSyncStatus: 'running',
memberSyncLockToken: generateId(),
memberSyncLockLeaseAt: staleLease,
}))
)
})

afterAll(async () => {
await db.delete(workspace).where(eq(workspace.id, ids.workspaceId))
await db.delete(organization).where(eq(organization.id, ids.organizationId))
await db.delete(user).where(inArray(user.id, [ids.aliceId, ids.bobId]))
await db.$client.end()
})
Comment thread
waleedlatif1 marked this conversation as resolved.

it('reclaims a stale members-mode lease and leaves every other access mode alone', async () => {
const response = await GET(createMockRequest('GET'), undefined)
expect(response.status).toBe(200)
const rows = await db
.select({
id: knowledgeConnector.id,
memberSyncStatus: knowledgeConnector.memberSyncStatus,
memberSyncLockToken: knowledgeConnector.memberSyncLockToken,
})
.from(knowledgeConnector)
.where(inArray(knowledgeConnector.id, [members, admin]))
const byId = new Map(rows.map((row) => [row.id, row]))
expect(byId.get(members)).toMatchObject({
memberSyncStatus: 'error',
memberSyncLockToken: null,
})
expect(byId.get(admin)?.memberSyncStatus).toBe('running')
expect(byId.get(admin)?.memberSyncLockToken).not.toBeNull()
})
})
16 changes: 14 additions & 2 deletions apps/sim/app/api/knowledge/connectors/member-sync/route.ts
Original file line number Diff line number Diff line change
Expand Up @@ -58,6 +58,18 @@ function reclaimedNextMemberSyncAt(): SQL {
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`
}

/**
* Only the member engine takes the member lease, and only on a members-mode connector, which a
* mode switch cannot leave while the lease is held; so both reclaims match `access_mode` too, the
* predicate `kc_member_sync_due_idx` is partial on, and read that index instead of the table.
*/
function reclaimableMemberSync(status: 'running' | 'pending'): SQL | undefined {
return and(
eq(knowledgeConnector.accessMode, 'members'),
eq(knowledgeConnector.memberSyncStatus, status)
)
}

/**
* The write shared by both reclaims: a run that stopped making progress
* re-enters the member failure ladder, which is the content engine's ladder
Expand Down Expand Up @@ -108,7 +120,7 @@ export const GET = withRouteHandler(async (request: NextRequest) => {
.set(reclaimPayload(STALE_LOCK_ERROR_MESSAGE))
.where(
and(
eq(knowledgeConnector.memberSyncStatus, 'running'),
reclaimableMemberSync('running'),
sql`${memberSyncLockLease()} <= ${sql.param(staleCutoff, knowledgeConnector.memberSyncLockLeaseAt)}`,
isNull(knowledgeConnector.archivedAt),
isNull(knowledgeConnector.deletedAt)
Expand All @@ -120,7 +132,7 @@ export const GET = withRouteHandler(async (request: NextRequest) => {
.set(reclaimPayload(LOST_DISPATCH_ERROR_MESSAGE))
.where(
and(
eq(knowledgeConnector.memberSyncStatus, 'pending'),
reclaimableMemberSync('pending'),
sql`${memberSyncLockLease()} <= ${sql.param(staleCutoff, knowledgeConnector.memberSyncLockLeaseAt)}`,
isNull(knowledgeConnector.archivedAt),
isNull(knowledgeConnector.deletedAt)
Expand Down
168 changes: 0 additions & 168 deletions apps/sim/app/api/knowledge/github/installations/route.test.ts

This file was deleted.

53 changes: 0 additions & 53 deletions apps/sim/app/api/knowledge/github/installations/route.ts

This file was deleted.

22 changes: 0 additions & 22 deletions apps/sim/app/api/knowledge/member-connectors/route.ts

This file was deleted.

7 changes: 4 additions & 3 deletions apps/sim/app/api/knowledge/search/route.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,12 +4,12 @@ import {
internalRateLimits,
internalSessionAuth,
} from '@/lib/api/server/routes'
import { isLiveEnterpriseSearchEnabled } from '@/lib/core/config/env-flags'
import { internalKnowledgeErrorPolicies } from '@/lib/knowledge/api/route-policies'
import { knowledgeOperations } from '@/lib/knowledge/application/operations'
import { searchScopedKnowledge } from '@/lib/knowledge/application/workspace-search'
import { DEFAULT_RERANKER_MODEL } from '@/lib/knowledge/reranker-models'
import { sourceAuthor } from '@/lib/knowledge/search/author'
import { searchScopedKnowledge } from '@/lib/sim-search/indexed'
import { isIndexedOrgSearchEnabled } from '@/lib/sim-search/indexed/gate'
import { searchLiveKnowledge } from '@/lib/sim-search/live/application'

const DIRECT_SEARCH_VECTOR_BUDGET_MS = 3000
Expand Down Expand Up @@ -81,4 +81,5 @@ const liveSearchRoute = defineInternalJsonRoute({
present: (data) => ({ success: true as const, data }),
})

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