diff --git a/apps/sim/.gitignore b/apps/sim/.gitignore index 4609b7bec2..55770d96da 100644 --- a/apps/sim/.gitignore +++ b/apps/sim/.gitignore @@ -47,3 +47,6 @@ next-env.d.ts # Sentry Config File .env.sentry-build-plugin + +# Uploads +/uploads diff --git a/apps/sim/app/api/knowledge/[id]/documents/[documentId]/chunks/[chunkId]/route.ts b/apps/sim/app/api/knowledge/[id]/documents/[documentId]/chunks/[chunkId]/route.ts new file mode 100644 index 0000000000..234836f27c --- /dev/null +++ b/apps/sim/app/api/knowledge/[id]/documents/[documentId]/chunks/[chunkId]/route.ts @@ -0,0 +1,267 @@ +import { and, eq, isNull } from 'drizzle-orm' +import { type NextRequest, NextResponse } from 'next/server' +import { z } from 'zod' +import { getSession } from '@/lib/auth' +import { createLogger } from '@/lib/logs/console-logger' +import { db } from '@/db' +import { document, embedding, knowledgeBase } from '@/db/schema' + +const logger = createLogger('ChunkByIdAPI') + +// Schema for chunk updates +const UpdateChunkSchema = z.object({ + content: z.string().min(1, 'Content is required').optional(), + enabled: z.boolean().optional(), + searchRank: z.number().min(0).optional(), + qualityScore: z.number().min(0).max(1).optional(), +}) + +async function checkChunkAccess( + knowledgeBaseId: string, + documentId: string, + chunkId: string, + userId: string +) { + // First check knowledge base access + const kb = await db + .select({ + id: knowledgeBase.id, + userId: knowledgeBase.userId, + }) + .from(knowledgeBase) + .where(and(eq(knowledgeBase.id, knowledgeBaseId), isNull(knowledgeBase.deletedAt))) + .limit(1) + + if (kb.length === 0) { + return { hasAccess: false, notFound: true, reason: 'Knowledge base not found' } + } + + const kbData = kb[0] + + // Check if user owns the knowledge base + if (kbData.userId !== userId) { + return { hasAccess: false, reason: 'Unauthorized knowledge base access' } + } + + // Check if document exists and belongs to the knowledge base + const doc = await db + .select() + .from(document) + .where( + and( + eq(document.id, documentId), + eq(document.knowledgeBaseId, knowledgeBaseId), + isNull(document.deletedAt) + ) + ) + .limit(1) + + if (doc.length === 0) { + return { hasAccess: false, notFound: true, reason: 'Document not found' } + } + + // Check if chunk exists and belongs to the document + const chunk = await db + .select() + .from(embedding) + .where(and(eq(embedding.id, chunkId), eq(embedding.documentId, documentId))) + .limit(1) + + if (chunk.length === 0) { + return { hasAccess: false, notFound: true, reason: 'Chunk not found' } + } + + return { hasAccess: true, chunk: chunk[0], document: doc[0], knowledgeBase: kbData } +} + +export async function GET( + req: NextRequest, + { params }: { params: Promise<{ id: string; documentId: string; chunkId: string }> } +) { + const requestId = crypto.randomUUID().slice(0, 8) + const { id: knowledgeBaseId, documentId, chunkId } = await params + + try { + const session = await getSession() + if (!session?.user?.id) { + logger.warn(`[${requestId}] Unauthorized chunk access attempt`) + return NextResponse.json({ error: 'Unauthorized' }, { status: 401 }) + } + + const accessCheck = await checkChunkAccess( + knowledgeBaseId, + documentId, + chunkId, + session.user.id + ) + + if (accessCheck.notFound) { + logger.warn( + `[${requestId}] ${accessCheck.reason}: KB=${knowledgeBaseId}, Doc=${documentId}, Chunk=${chunkId}` + ) + return NextResponse.json({ error: accessCheck.reason }, { status: 404 }) + } + + if (!accessCheck.hasAccess) { + logger.warn( + `[${requestId}] User ${session.user.id} attempted unauthorized chunk access: ${accessCheck.reason}` + ) + return NextResponse.json({ error: 'Unauthorized' }, { status: 401 }) + } + + logger.info( + `[${requestId}] Retrieved chunk: ${chunkId} from document ${documentId} in knowledge base ${knowledgeBaseId}` + ) + + return NextResponse.json({ + success: true, + data: accessCheck.chunk, + }) + } catch (error) { + logger.error(`[${requestId}] Error fetching chunk`, error) + return NextResponse.json({ error: 'Failed to fetch chunk' }, { status: 500 }) + } +} + +export async function PUT( + req: NextRequest, + { params }: { params: Promise<{ id: string; documentId: string; chunkId: string }> } +) { + const requestId = crypto.randomUUID().slice(0, 8) + const { id: knowledgeBaseId, documentId, chunkId } = await params + + try { + const session = await getSession() + if (!session?.user?.id) { + logger.warn(`[${requestId}] Unauthorized chunk update attempt`) + return NextResponse.json({ error: 'Unauthorized' }, { status: 401 }) + } + + const accessCheck = await checkChunkAccess( + knowledgeBaseId, + documentId, + chunkId, + session.user.id + ) + + if (accessCheck.notFound) { + logger.warn( + `[${requestId}] ${accessCheck.reason}: KB=${knowledgeBaseId}, Doc=${documentId}, Chunk=${chunkId}` + ) + return NextResponse.json({ error: accessCheck.reason }, { status: 404 }) + } + + if (!accessCheck.hasAccess) { + logger.warn( + `[${requestId}] User ${session.user.id} attempted unauthorized chunk update: ${accessCheck.reason}` + ) + return NextResponse.json({ error: 'Unauthorized' }, { status: 401 }) + } + + const body = await req.json() + + try { + const validatedData = UpdateChunkSchema.parse(body) + + const updateData: any = { + updatedAt: new Date(), + } + + if (validatedData.content !== undefined) { + updateData.content = validatedData.content + updateData.contentLength = validatedData.content.length + // Update token count estimation (rough approximation: 4 chars per token) + updateData.tokenCount = Math.ceil(validatedData.content.length / 4) + } + if (validatedData.enabled !== undefined) updateData.enabled = validatedData.enabled + if (validatedData.searchRank !== undefined) + updateData.searchRank = validatedData.searchRank.toString() + if (validatedData.qualityScore !== undefined) + updateData.qualityScore = validatedData.qualityScore.toString() + + await db.update(embedding).set(updateData).where(eq(embedding.id, chunkId)) + + // Fetch the updated chunk + const updatedChunk = await db + .select() + .from(embedding) + .where(eq(embedding.id, chunkId)) + .limit(1) + + logger.info( + `[${requestId}] Chunk updated: ${chunkId} in document ${documentId} in knowledge base ${knowledgeBaseId}` + ) + + return NextResponse.json({ + success: true, + data: updatedChunk[0], + }) + } catch (validationError) { + if (validationError instanceof z.ZodError) { + logger.warn(`[${requestId}] Invalid chunk update data`, { + errors: validationError.errors, + }) + return NextResponse.json( + { error: 'Invalid request data', details: validationError.errors }, + { status: 400 } + ) + } + throw validationError + } + } catch (error) { + logger.error(`[${requestId}] Error updating chunk`, error) + return NextResponse.json({ error: 'Failed to update chunk' }, { status: 500 }) + } +} + +export async function DELETE( + req: NextRequest, + { params }: { params: Promise<{ id: string; documentId: string; chunkId: string }> } +) { + const requestId = crypto.randomUUID().slice(0, 8) + const { id: knowledgeBaseId, documentId, chunkId } = await params + + try { + const session = await getSession() + if (!session?.user?.id) { + logger.warn(`[${requestId}] Unauthorized chunk delete attempt`) + return NextResponse.json({ error: 'Unauthorized' }, { status: 401 }) + } + + const accessCheck = await checkChunkAccess( + knowledgeBaseId, + documentId, + chunkId, + session.user.id + ) + + if (accessCheck.notFound) { + logger.warn( + `[${requestId}] ${accessCheck.reason}: KB=${knowledgeBaseId}, Doc=${documentId}, Chunk=${chunkId}` + ) + return NextResponse.json({ error: accessCheck.reason }, { status: 404 }) + } + + if (!accessCheck.hasAccess) { + logger.warn( + `[${requestId}] User ${session.user.id} attempted unauthorized chunk deletion: ${accessCheck.reason}` + ) + return NextResponse.json({ error: 'Unauthorized' }, { status: 401 }) + } + + // Delete the chunk + await db.delete(embedding).where(eq(embedding.id, chunkId)) + + logger.info( + `[${requestId}] Chunk deleted: ${chunkId} from document ${documentId} in knowledge base ${knowledgeBaseId}` + ) + + return NextResponse.json({ + success: true, + data: { message: 'Chunk deleted successfully' }, + }) + } catch (error) { + logger.error(`[${requestId}] Error deleting chunk`, error) + return NextResponse.json({ error: 'Failed to delete chunk' }, { status: 500 }) + } +} diff --git a/apps/sim/app/api/knowledge/[id]/documents/[documentId]/chunks/route.ts b/apps/sim/app/api/knowledge/[id]/documents/[documentId]/chunks/route.ts new file mode 100644 index 0000000000..0bf7ebc40c --- /dev/null +++ b/apps/sim/app/api/knowledge/[id]/documents/[documentId]/chunks/route.ts @@ -0,0 +1,161 @@ +import { and, asc, eq, ilike, isNull, sql } from 'drizzle-orm' +import { type NextRequest, NextResponse } from 'next/server' +import { z } from 'zod' +import { getSession } from '@/lib/auth' +import { createLogger } from '@/lib/logs/console-logger' +import { db } from '@/db' +import { document, embedding, knowledgeBase } from '@/db/schema' + +const logger = createLogger('DocumentChunksAPI') + +// Schema for query parameters +const GetChunksQuerySchema = z.object({ + search: z.string().optional(), + enabled: z.enum(['true', 'false', 'all']).optional().default('all'), + limit: z.coerce.number().min(1).max(100).optional().default(50), + offset: z.coerce.number().min(0).optional().default(0), +}) + +async function checkDocumentAccess(knowledgeBaseId: string, documentId: string, userId: string) { + // First check knowledge base access + const kb = await db + .select({ + id: knowledgeBase.id, + userId: knowledgeBase.userId, + }) + .from(knowledgeBase) + .where(and(eq(knowledgeBase.id, knowledgeBaseId), isNull(knowledgeBase.deletedAt))) + .limit(1) + + if (kb.length === 0) { + return { hasAccess: false, notFound: true, reason: 'Knowledge base not found' } + } + + const kbData = kb[0] + + // Check if user owns the knowledge base + if (kbData.userId !== userId) { + return { hasAccess: false, reason: 'Unauthorized knowledge base access' } + } + + // Now check if document exists and belongs to the knowledge base + const doc = await db + .select() + .from(document) + .where( + and( + eq(document.id, documentId), + eq(document.knowledgeBaseId, knowledgeBaseId), + isNull(document.deletedAt) + ) + ) + .limit(1) + + if (doc.length === 0) { + return { hasAccess: false, notFound: true, reason: 'Document not found' } + } + + return { hasAccess: true, document: doc[0], knowledgeBase: kbData } +} + +export async function GET( + req: NextRequest, + { params }: { params: Promise<{ id: string; documentId: string }> } +) { + const requestId = crypto.randomUUID().slice(0, 8) + const { id: knowledgeBaseId, documentId } = await params + + try { + const session = await getSession() + if (!session?.user?.id) { + logger.warn(`[${requestId}] Unauthorized chunks access attempt`) + return NextResponse.json({ error: 'Unauthorized' }, { status: 401 }) + } + + const accessCheck = await checkDocumentAccess(knowledgeBaseId, documentId, session.user.id) + + if (accessCheck.notFound) { + logger.warn(`[${requestId}] ${accessCheck.reason}: KB=${knowledgeBaseId}, Doc=${documentId}`) + return NextResponse.json({ error: accessCheck.reason }, { status: 404 }) + } + + if (!accessCheck.hasAccess) { + logger.warn( + `[${requestId}] User ${session.user.id} attempted unauthorized chunks access: ${accessCheck.reason}` + ) + return NextResponse.json({ error: 'Unauthorized' }, { status: 401 }) + } + + // Parse query parameters + const { searchParams } = new URL(req.url) + const queryParams = GetChunksQuerySchema.parse({ + search: searchParams.get('search') || undefined, + enabled: searchParams.get('enabled') || undefined, + limit: searchParams.get('limit') || undefined, + offset: searchParams.get('offset') || undefined, + }) + + // Build query conditions + const conditions = [eq(embedding.documentId, documentId)] + + // Add enabled filter + if (queryParams.enabled === 'true') { + conditions.push(eq(embedding.enabled, true)) + } else if (queryParams.enabled === 'false') { + conditions.push(eq(embedding.enabled, false)) + } + + // Add search filter + if (queryParams.search) { + conditions.push(ilike(embedding.content, `%${queryParams.search}%`)) + } + + // Fetch chunks + const chunks = await db + .select({ + id: embedding.id, + chunkIndex: embedding.chunkIndex, + content: embedding.content, + contentLength: embedding.contentLength, + tokenCount: embedding.tokenCount, + enabled: embedding.enabled, + startOffset: embedding.startOffset, + endOffset: embedding.endOffset, + overlapTokens: embedding.overlapTokens, + metadata: embedding.metadata, + searchRank: embedding.searchRank, + qualityScore: embedding.qualityScore, + createdAt: embedding.createdAt, + updatedAt: embedding.updatedAt, + }) + .from(embedding) + .where(and(...conditions)) + .orderBy(asc(embedding.chunkIndex)) + .limit(queryParams.limit) + .offset(queryParams.offset) + + // Get total count for pagination + const totalCount = await db + .select({ count: sql`count(*)` }) + .from(embedding) + .where(and(...conditions)) + + logger.info( + `[${requestId}] Retrieved ${chunks.length} chunks for document ${documentId} in knowledge base ${knowledgeBaseId}` + ) + + return NextResponse.json({ + success: true, + data: chunks, + pagination: { + total: Number(totalCount[0]?.count || 0), + limit: queryParams.limit, + offset: queryParams.offset, + hasMore: chunks.length === queryParams.limit, + }, + }) + } catch (error) { + logger.error(`[${requestId}] Error fetching chunks`, error) + return NextResponse.json({ error: 'Failed to fetch chunks' }, { status: 500 }) + } +} diff --git a/apps/sim/app/api/knowledge/[id]/documents/[documentId]/route.ts b/apps/sim/app/api/knowledge/[id]/documents/[documentId]/route.ts new file mode 100644 index 0000000000..7367e891e0 --- /dev/null +++ b/apps/sim/app/api/knowledge/[id]/documents/[documentId]/route.ts @@ -0,0 +1,229 @@ +import { and, eq, isNull } from 'drizzle-orm' +import { type NextRequest, NextResponse } from 'next/server' +import { z } from 'zod' +import { getSession } from '@/lib/auth' +import { createLogger } from '@/lib/logs/console-logger' +import { db } from '@/db' +import { document, knowledgeBase } from '@/db/schema' + +const logger = createLogger('DocumentByIdAPI') + +// Schema for document updates +const UpdateDocumentSchema = z.object({ + filename: z.string().min(1, 'Filename is required').optional(), + enabled: z.boolean().optional(), + chunkCount: z.number().min(0).optional(), + tokenCount: z.number().min(0).optional(), + characterCount: z.number().min(0).optional(), +}) + +async function checkDocumentAccess(knowledgeBaseId: string, documentId: string, userId: string) { + // First check knowledge base access + const kb = await db + .select({ + id: knowledgeBase.id, + userId: knowledgeBase.userId, + }) + .from(knowledgeBase) + .where(and(eq(knowledgeBase.id, knowledgeBaseId), isNull(knowledgeBase.deletedAt))) + .limit(1) + + if (kb.length === 0) { + return { hasAccess: false, notFound: true, reason: 'Knowledge base not found' } + } + + const kbData = kb[0] + + // Check if user owns the knowledge base + if (kbData.userId !== userId) { + return { hasAccess: false, reason: 'Unauthorized knowledge base access' } + } + + // Now check if document exists and belongs to the knowledge base + const doc = await db + .select() + .from(document) + .where( + and( + eq(document.id, documentId), + eq(document.knowledgeBaseId, knowledgeBaseId), + isNull(document.deletedAt) + ) + ) + .limit(1) + + if (doc.length === 0) { + return { hasAccess: false, notFound: true, reason: 'Document not found' } + } + + return { hasAccess: true, document: doc[0], knowledgeBase: kbData } +} + +export async function GET( + req: NextRequest, + { params }: { params: Promise<{ id: string; documentId: string }> } +) { + const requestId = crypto.randomUUID().slice(0, 8) + const { id: knowledgeBaseId, documentId } = await params + + try { + const session = await getSession() + if (!session?.user?.id) { + logger.warn(`[${requestId}] Unauthorized document access attempt`) + return NextResponse.json({ error: 'Unauthorized' }, { status: 401 }) + } + + const accessCheck = await checkDocumentAccess(knowledgeBaseId, documentId, session.user.id) + + if (accessCheck.notFound) { + logger.warn(`[${requestId}] ${accessCheck.reason}: KB=${knowledgeBaseId}, Doc=${documentId}`) + return NextResponse.json({ error: accessCheck.reason }, { status: 404 }) + } + + if (!accessCheck.hasAccess) { + logger.warn( + `[${requestId}] User ${session.user.id} attempted unauthorized document access: ${accessCheck.reason}` + ) + return NextResponse.json({ error: 'Unauthorized' }, { status: 401 }) + } + + logger.info( + `[${requestId}] Retrieved document: ${documentId} from knowledge base ${knowledgeBaseId}` + ) + + return NextResponse.json({ + success: true, + data: accessCheck.document, + }) + } catch (error) { + logger.error(`[${requestId}] Error fetching document`, error) + return NextResponse.json({ error: 'Failed to fetch document' }, { status: 500 }) + } +} + +export async function PUT( + req: NextRequest, + { params }: { params: Promise<{ id: string; documentId: string }> } +) { + const requestId = crypto.randomUUID().slice(0, 8) + const { id: knowledgeBaseId, documentId } = await params + + try { + const session = await getSession() + if (!session?.user?.id) { + logger.warn(`[${requestId}] Unauthorized document update attempt`) + return NextResponse.json({ error: 'Unauthorized' }, { status: 401 }) + } + + const accessCheck = await checkDocumentAccess(knowledgeBaseId, documentId, session.user.id) + + if (accessCheck.notFound) { + logger.warn(`[${requestId}] ${accessCheck.reason}: KB=${knowledgeBaseId}, Doc=${documentId}`) + return NextResponse.json({ error: accessCheck.reason }, { status: 404 }) + } + + if (!accessCheck.hasAccess) { + logger.warn( + `[${requestId}] User ${session.user.id} attempted unauthorized document update: ${accessCheck.reason}` + ) + return NextResponse.json({ error: 'Unauthorized' }, { status: 401 }) + } + + const body = await req.json() + + try { + const validatedData = UpdateDocumentSchema.parse(body) + + const updateData: any = {} + + if (validatedData.filename !== undefined) updateData.filename = validatedData.filename + if (validatedData.enabled !== undefined) updateData.enabled = validatedData.enabled + if (validatedData.chunkCount !== undefined) updateData.chunkCount = validatedData.chunkCount + if (validatedData.tokenCount !== undefined) updateData.tokenCount = validatedData.tokenCount + if (validatedData.characterCount !== undefined) + updateData.characterCount = validatedData.characterCount + + await db.update(document).set(updateData).where(eq(document.id, documentId)) + + // Fetch the updated document + const updatedDocument = await db + .select() + .from(document) + .where(eq(document.id, documentId)) + .limit(1) + + logger.info( + `[${requestId}] Document updated: ${documentId} in knowledge base ${knowledgeBaseId}` + ) + + return NextResponse.json({ + success: true, + data: updatedDocument[0], + }) + } catch (validationError) { + if (validationError instanceof z.ZodError) { + logger.warn(`[${requestId}] Invalid document update data`, { + errors: validationError.errors, + }) + return NextResponse.json( + { error: 'Invalid request data', details: validationError.errors }, + { status: 400 } + ) + } + throw validationError + } + } catch (error) { + logger.error(`[${requestId}] Error updating document`, error) + return NextResponse.json({ error: 'Failed to update document' }, { status: 500 }) + } +} + +export async function DELETE( + req: NextRequest, + { params }: { params: Promise<{ id: string; documentId: string }> } +) { + const requestId = crypto.randomUUID().slice(0, 8) + const { id: knowledgeBaseId, documentId } = await params + + try { + const session = await getSession() + if (!session?.user?.id) { + logger.warn(`[${requestId}] Unauthorized document delete attempt`) + return NextResponse.json({ error: 'Unauthorized' }, { status: 401 }) + } + + const accessCheck = await checkDocumentAccess(knowledgeBaseId, documentId, session.user.id) + + if (accessCheck.notFound) { + logger.warn(`[${requestId}] ${accessCheck.reason}: KB=${knowledgeBaseId}, Doc=${documentId}`) + return NextResponse.json({ error: accessCheck.reason }, { status: 404 }) + } + + if (!accessCheck.hasAccess) { + logger.warn( + `[${requestId}] User ${session.user.id} attempted unauthorized document deletion: ${accessCheck.reason}` + ) + return NextResponse.json({ error: 'Unauthorized' }, { status: 401 }) + } + + // Soft delete by setting deletedAt timestamp + await db + .update(document) + .set({ + deletedAt: new Date(), + }) + .where(eq(document.id, documentId)) + + logger.info( + `[${requestId}] Document deleted: ${documentId} from knowledge base ${knowledgeBaseId}` + ) + + return NextResponse.json({ + success: true, + data: { message: 'Document deleted successfully' }, + }) + } catch (error) { + logger.error(`[${requestId}] Error deleting document`, error) + return NextResponse.json({ error: 'Failed to delete document' }, { status: 500 }) + } +} diff --git a/apps/sim/app/api/knowledge/[id]/documents/route.ts b/apps/sim/app/api/knowledge/[id]/documents/route.ts new file mode 100644 index 0000000000..d1453cd7b1 --- /dev/null +++ b/apps/sim/app/api/knowledge/[id]/documents/route.ts @@ -0,0 +1,213 @@ +import { and, eq, isNull } from 'drizzle-orm' +import { type NextRequest, NextResponse } from 'next/server' +import { z } from 'zod' +import { getSession } from '@/lib/auth' +import { createLogger } from '@/lib/logs/console-logger' +import { db } from '@/db' +import { document, knowledgeBase } from '@/db/schema' + +const logger = createLogger('DocumentsAPI') + +// Schema for document creation +const CreateDocumentSchema = z.object({ + filename: z.string().min(1, 'Filename is required'), + fileUrl: z.string().url('File URL must be valid'), + fileSize: z.number().min(1, 'File size must be greater than 0'), + mimeType: z.string().min(1, 'MIME type is required'), + fileHash: z.string().optional(), +}) + +async function checkKnowledgeBaseAccess(knowledgeBaseId: string, userId: string) { + const kb = await db + .select({ + id: knowledgeBase.id, + userId: knowledgeBase.userId, + }) + .from(knowledgeBase) + .where(and(eq(knowledgeBase.id, knowledgeBaseId), isNull(knowledgeBase.deletedAt))) + .limit(1) + + if (kb.length === 0) { + return { hasAccess: false, notFound: true } + } + + const kbData = kb[0] + + // Check if user owns the knowledge base + if (kbData.userId === userId) { + return { hasAccess: true, knowledgeBase: kbData } + } + + return { hasAccess: false, knowledgeBase: kbData } +} + +export async function GET(req: NextRequest, { params }: { params: Promise<{ id: string }> }) { + const requestId = crypto.randomUUID().slice(0, 8) + const { id: knowledgeBaseId } = await params + + try { + const session = await getSession() + if (!session?.user?.id) { + logger.warn(`[${requestId}] Unauthorized documents access attempt`) + return NextResponse.json({ error: 'Unauthorized' }, { status: 401 }) + } + + const accessCheck = await checkKnowledgeBaseAccess(knowledgeBaseId, session.user.id) + + if (accessCheck.notFound) { + logger.warn(`[${requestId}] Knowledge base not found: ${knowledgeBaseId}`) + return NextResponse.json({ error: 'Knowledge base not found' }, { status: 404 }) + } + + if (!accessCheck.hasAccess) { + logger.warn( + `[${requestId}] User ${session.user.id} attempted to access unauthorized knowledge base documents ${knowledgeBaseId}` + ) + return NextResponse.json({ error: 'Unauthorized' }, { status: 401 }) + } + + const url = new URL(req.url) + const includeDisabled = url.searchParams.get('includeDisabled') === 'true' + + // Build where conditions + const whereConditions = [ + eq(document.knowledgeBaseId, knowledgeBaseId), + isNull(document.deletedAt), + ] + + // Filter out disabled documents unless specifically requested + if (!includeDisabled) { + whereConditions.push(eq(document.enabled, true)) + } + + const documents = await db + .select({ + id: document.id, + knowledgeBaseId: document.knowledgeBaseId, + filename: document.filename, + fileUrl: document.fileUrl, + fileSize: document.fileSize, + mimeType: document.mimeType, + fileHash: document.fileHash, + chunkCount: document.chunkCount, + tokenCount: document.tokenCount, + characterCount: document.characterCount, + enabled: document.enabled, + uploadedAt: document.uploadedAt, + }) + .from(document) + .where(and(...whereConditions)) + .orderBy(document.uploadedAt) + + logger.info( + `[${requestId}] Retrieved ${documents.length} documents for knowledge base ${knowledgeBaseId}` + ) + + return NextResponse.json({ + success: true, + data: documents, + }) + } catch (error) { + logger.error(`[${requestId}] Error fetching documents`, error) + return NextResponse.json({ error: 'Failed to fetch documents' }, { status: 500 }) + } +} + +export async function POST(req: NextRequest, { params }: { params: Promise<{ id: string }> }) { + const requestId = crypto.randomUUID().slice(0, 8) + const { id: knowledgeBaseId } = await params + + try { + const session = await getSession() + if (!session?.user?.id) { + logger.warn(`[${requestId}] Unauthorized document creation attempt`) + return NextResponse.json({ error: 'Unauthorized' }, { status: 401 }) + } + + const accessCheck = await checkKnowledgeBaseAccess(knowledgeBaseId, session.user.id) + + if (accessCheck.notFound) { + logger.warn(`[${requestId}] Knowledge base not found: ${knowledgeBaseId}`) + return NextResponse.json({ error: 'Knowledge base not found' }, { status: 404 }) + } + + if (!accessCheck.hasAccess) { + logger.warn( + `[${requestId}] User ${session.user.id} attempted to create document in unauthorized knowledge base ${knowledgeBaseId}` + ) + return NextResponse.json({ error: 'Unauthorized' }, { status: 401 }) + } + + const body = await req.json() + + try { + const validatedData = CreateDocumentSchema.parse(body) + + // Check for duplicate file hash if provided + if (validatedData.fileHash) { + const existingDocument = await db + .select({ id: document.id }) + .from(document) + .where( + and( + eq(document.knowledgeBaseId, knowledgeBaseId), + eq(document.fileHash, validatedData.fileHash), + isNull(document.deletedAt) + ) + ) + .limit(1) + + if (existingDocument.length > 0) { + logger.warn(`[${requestId}] Duplicate file hash detected: ${validatedData.fileHash}`) + return NextResponse.json( + { error: 'Document with this file hash already exists' }, + { status: 409 } + ) + } + } + + const documentId = crypto.randomUUID() + const now = new Date() + + const newDocument = { + id: documentId, + knowledgeBaseId, + filename: validatedData.filename, + fileUrl: validatedData.fileUrl, + fileSize: validatedData.fileSize, + mimeType: validatedData.mimeType, + fileHash: validatedData.fileHash || null, + chunkCount: 0, + tokenCount: 0, + characterCount: 0, + enabled: true, + uploadedAt: now, + } + + await db.insert(document).values(newDocument) + + logger.info( + `[${requestId}] Document created: ${documentId} in knowledge base ${knowledgeBaseId}` + ) + + return NextResponse.json({ + success: true, + data: newDocument, + }) + } catch (validationError) { + if (validationError instanceof z.ZodError) { + logger.warn(`[${requestId}] Invalid document data`, { + errors: validationError.errors, + }) + return NextResponse.json( + { error: 'Invalid request data', details: validationError.errors }, + { status: 400 } + ) + } + throw validationError + } + } catch (error) { + logger.error(`[${requestId}] Error creating document`, error) + return NextResponse.json({ error: 'Failed to create document' }, { status: 500 }) + } +} diff --git a/apps/sim/app/api/knowledge/[id]/process-documents/route.ts b/apps/sim/app/api/knowledge/[id]/process-documents/route.ts new file mode 100644 index 0000000000..73521a7d47 --- /dev/null +++ b/apps/sim/app/api/knowledge/[id]/process-documents/route.ts @@ -0,0 +1,392 @@ +import { and, eq, isNull } from 'drizzle-orm' +import { type NextRequest, NextResponse } from 'next/server' +import { z } from 'zod' +import { getSession } from '@/lib/auth' +import { type ProcessedDocument, processDocuments } from '@/lib/document-processor' +import { env } from '@/lib/env' +import { createLogger } from '@/lib/logs/console-logger' +import { db } from '@/db' +import { document, embedding, knowledgeBase } from '@/db/schema' + +const logger = createLogger('ProcessDocumentsAPI') + +// Schema for document processing request +const ProcessDocumentsSchema = z.object({ + documents: z + .array( + z.object({ + filename: z.string().min(1, 'Filename is required'), + fileUrl: z.string().url('File URL must be valid'), + fileSize: z.number().min(1, 'File size must be greater than 0'), + mimeType: z.string().min(1, 'MIME type is required'), + fileHash: z.string().optional(), + }) + ) + .min(1, 'At least one document is required'), + processingOptions: z + .object({ + chunkSize: z.number().min(100).max(2048).default(512), + minCharactersPerChunk: z.number().min(10).max(1000).default(24), + recipe: z.string().default('default'), + lang: z.string().default('en'), + }) + .optional(), +}) + +async function checkKnowledgeBaseAccess(knowledgeBaseId: string, userId: string) { + const kb = await db + .select({ + id: knowledgeBase.id, + userId: knowledgeBase.userId, + chunkingConfig: knowledgeBase.chunkingConfig, + }) + .from(knowledgeBase) + .where(and(eq(knowledgeBase.id, knowledgeBaseId), isNull(knowledgeBase.deletedAt))) + .limit(1) + + if (kb.length === 0) { + return { hasAccess: false, notFound: true } + } + + const kbData = kb[0] + + // Check if user owns the knowledge base + if (kbData.userId === userId) { + return { hasAccess: true, knowledgeBase: kbData } + } + + return { hasAccess: false, knowledgeBase: kbData } +} + +async function generateEmbeddings( + texts: string[], + embeddingModel = 'text-embedding-3-small' +): Promise { + const openaiApiKey = env.OPENAI_API_KEY + if (!openaiApiKey) { + throw new Error('OPENAI_API_KEY not configured') + } + + try { + // Batch process embeddings for efficiency + const batchSize = 100 // OpenAI allows up to 2048 inputs per request + const allEmbeddings: number[][] = [] + + for (let i = 0; i < texts.length; i += batchSize) { + const batch = texts.slice(i, i + batchSize) + + logger.info( + `Generating embeddings for batch ${Math.floor(i / batchSize) + 1}/${Math.ceil(texts.length / batchSize)} (${batch.length} texts)` + ) + + // Make direct API call to OpenAI embeddings + const response = await fetch('https://api.openai.com/v1/embeddings', { + method: 'POST', + headers: { + Authorization: `Bearer ${openaiApiKey}`, + 'Content-Type': 'application/json', + }, + body: JSON.stringify({ + input: batch, + model: embeddingModel, + encoding_format: 'float', + }), + }) + + if (!response.ok) { + const errorText = await response.text() + throw new Error( + `OpenAI API error: ${response.status} ${response.statusText} - ${errorText}` + ) + } + + const data = await response.json() + + if (!data.data || !Array.isArray(data.data)) { + throw new Error('Invalid response format from OpenAI embeddings API') + } + + // Extract embeddings from response + const batchEmbeddings = data.data.map((item: any) => item.embedding) + allEmbeddings.push(...batchEmbeddings) + } + + logger.info(`Successfully generated ${allEmbeddings.length} embeddings`) + return allEmbeddings + } catch (error) { + logger.error('Failed to generate embeddings:', error) + throw new Error( + `Embedding generation failed: ${error instanceof Error ? error.message : 'Unknown error'}` + ) + } +} + +async function saveProcessedDocuments( + knowledgeBaseId: string, + processedDocuments: ProcessedDocument[], + requestedDocuments: Array<{ + filename: string + fileUrl: string + fileSize: number + mimeType: string + fileHash?: string + }> +) { + const now = new Date() + const results: Array<{ + documentId: string + chunkCount: number + success: boolean + error?: string + }> = [] + + // Collect all chunk texts for batch embedding generation + const allChunkTexts: string[] = [] + const chunkMapping: Array<{ docIndex: number; chunkIndex: number }> = [] + + processedDocuments.forEach((processed, docIndex) => { + processed.chunks.forEach((chunk, chunkIndex) => { + allChunkTexts.push(chunk.text) + chunkMapping.push({ docIndex, chunkIndex }) + }) + }) + + // Generate embeddings for all chunks at once + let allEmbeddings: number[][] = [] + if (allChunkTexts.length > 0) { + try { + logger.info( + `Generating embeddings for ${allChunkTexts.length} chunks across ${processedDocuments.length} documents` + ) + allEmbeddings = await generateEmbeddings(allChunkTexts, 'text-embedding-3-small') + logger.info(`Successfully generated ${allEmbeddings.length} embeddings`) + } catch (error) { + logger.error('Failed to generate embeddings for chunks:', error) + // Continue without embeddings rather than failing completely + allEmbeddings = [] + } + } + + for (let i = 0; i < processedDocuments.length; i++) { + const processed = processedDocuments[i] + const original = requestedDocuments.find((doc) => doc.filename === processed.metadata.filename) + + if (!original) { + results.push({ + documentId: '', + chunkCount: 0, + success: false, + error: `Original document data not found for ${processed.metadata.filename}`, + }) + continue + } + + try { + // Check for duplicate file hash if provided + if (original.fileHash) { + const existingDocument = await db + .select({ id: document.id }) + .from(document) + .where( + and( + eq(document.knowledgeBaseId, knowledgeBaseId), + eq(document.fileHash, original.fileHash), + isNull(document.deletedAt) + ) + ) + .limit(1) + + if (existingDocument.length > 0) { + results.push({ + documentId: existingDocument[0].id, + chunkCount: 0, + success: false, + error: 'Document with this file hash already exists', + }) + continue + } + } + + // Insert document record + const documentId = crypto.randomUUID() + const newDocument = { + id: documentId, + knowledgeBaseId, + filename: original.filename, + fileUrl: processed.metadata.s3Url || original.fileUrl, + fileSize: original.fileSize, + mimeType: original.mimeType, + fileHash: original.fileHash || null, + chunkCount: processed.metadata.chunkCount, + tokenCount: processed.metadata.tokenCount, + characterCount: processed.metadata.characterCount, + enabled: true, + uploadedAt: now, + } + + await db.insert(document).values(newDocument) + + // Insert embedding records for chunks with generated embeddings + const embeddingRecords = processed.chunks.map((chunk, chunkIndex) => { + // Find the corresponding embedding for this chunk + const globalChunkIndex = chunkMapping.findIndex( + (mapping) => mapping.docIndex === i && mapping.chunkIndex === chunkIndex + ) + const embedding = + globalChunkIndex >= 0 && globalChunkIndex < allEmbeddings.length + ? allEmbeddings[globalChunkIndex] + : null + + return { + id: crypto.randomUUID(), + knowledgeBaseId, + documentId, + chunkIndex: chunkIndex, + chunkHash: crypto.randomUUID(), // Generate a hash for the chunk + content: chunk.text, + contentLength: chunk.text.length, + tokenCount: Math.ceil(chunk.text.length / 4), // Rough token estimation + embedding: embedding, // Store the generated OpenAI embedding + embeddingModel: 'text-embedding-3-small', + startOffset: chunk.startIndex || 0, + endOffset: chunk.endIndex || chunk.text.length, + overlapTokens: 0, + metadata: {}, + searchRank: '1.0', + accessCount: 0, + lastAccessedAt: null, + qualityScore: null, + createdAt: now, + updatedAt: now, + } + }) + + if (embeddingRecords.length > 0) { + await db.insert(embedding).values(embeddingRecords) + } + + results.push({ + documentId, + chunkCount: processed.metadata.chunkCount, + success: true, + }) + + logger.info( + `Document processed and saved: ${documentId} with ${processed.metadata.chunkCount} chunks and ${embeddingRecords.filter((r) => r.embedding).length} embeddings` + ) + } catch (error) { + logger.error(`Failed to save processed document ${processed.metadata.filename}:`, error) + results.push({ + documentId: '', + chunkCount: 0, + success: false, + error: error instanceof Error ? error.message : 'Unknown error during save', + }) + } + } + + return results +} + +export async function POST(req: NextRequest, { params }: { params: Promise<{ id: string }> }) { + const requestId = crypto.randomUUID().slice(0, 8) + const { id: knowledgeBaseId } = await params + + try { + const session = await getSession() + if (!session?.user?.id) { + logger.warn(`[${requestId}] Unauthorized document processing attempt`) + return NextResponse.json({ error: 'Unauthorized' }, { status: 401 }) + } + + const accessCheck = await checkKnowledgeBaseAccess(knowledgeBaseId, session.user.id) + + if (accessCheck.notFound) { + logger.warn(`[${requestId}] Knowledge base not found: ${knowledgeBaseId}`) + return NextResponse.json({ error: 'Knowledge base not found' }, { status: 404 }) + } + + if (!accessCheck.hasAccess) { + logger.warn( + `[${requestId}] User ${session.user.id} attempted to process documents in unauthorized knowledge base ${knowledgeBaseId}` + ) + return NextResponse.json({ error: 'Unauthorized' }, { status: 401 }) + } + + const body = await req.json() + + try { + const validatedData = ProcessDocumentsSchema.parse(body) + + logger.info( + `[${requestId}] Starting processing of ${validatedData.documents.length} documents` + ) + + // Get chunking config from knowledge base or use defaults + const kbChunkingConfig = accessCheck.knowledgeBase?.chunkingConfig as any + const processingOptions = { + knowledgeBaseId, + chunkSize: validatedData.processingOptions?.chunkSize || kbChunkingConfig?.maxSize || 512, + minCharactersPerChunk: + validatedData.processingOptions?.minCharactersPerChunk || kbChunkingConfig?.minSize || 24, + recipe: validatedData.processingOptions?.recipe || 'default', + lang: validatedData.processingOptions?.lang || 'en', + } + + // Process documents (parsing + chunking) + const processedDocuments = await processDocuments( + validatedData.documents.map((doc) => ({ + fileUrl: doc.fileUrl, + filename: doc.filename, + mimeType: doc.mimeType, + fileSize: doc.fileSize, + })), + processingOptions + ) + + // Save processed documents and chunks to database + const saveResults = await saveProcessedDocuments( + knowledgeBaseId, + processedDocuments, + validatedData.documents + ) + + const successfulCount = saveResults.filter((r) => r.success).length + const totalChunks = saveResults.reduce((sum, r) => sum + r.chunkCount, 0) + + logger.info( + `[${requestId}] Document processing completed: ${successfulCount}/${validatedData.documents.length} documents, ${totalChunks} total chunks` + ) + + return NextResponse.json({ + success: true, + data: { + processed: successfulCount, + total: validatedData.documents.length, + totalChunks, + results: saveResults, + }, + }) + } catch (validationError) { + if (validationError instanceof z.ZodError) { + logger.warn(`[${requestId}] Invalid document processing data`, { + errors: validationError.errors, + }) + return NextResponse.json( + { error: 'Invalid request data', details: validationError.errors }, + { status: 400 } + ) + } + throw validationError + } + } catch (error) { + logger.error(`[${requestId}] Error processing documents`, error) + return NextResponse.json( + { + error: 'Failed to process documents', + details: error instanceof Error ? error.message : 'Unknown error', + }, + { status: 500 } + ) + } +} diff --git a/apps/sim/app/api/knowledge/[id]/route.ts b/apps/sim/app/api/knowledge/[id]/route.ts new file mode 100644 index 0000000000..51ca1e2091 --- /dev/null +++ b/apps/sim/app/api/knowledge/[id]/route.ts @@ -0,0 +1,224 @@ +import { and, eq, isNull } from 'drizzle-orm' +import { type NextRequest, NextResponse } from 'next/server' +import { z } from 'zod' +import { getSession } from '@/lib/auth' +import { createLogger } from '@/lib/logs/console-logger' +import { db } from '@/db' +import { knowledgeBase } from '@/db/schema' + +const logger = createLogger('KnowledgeBaseByIdAPI') + +// Schema for knowledge base updates +const UpdateKnowledgeBaseSchema = z.object({ + name: z.string().min(1, 'Name is required').optional(), + description: z.string().optional(), + embeddingModel: z.literal('text-embedding-3-small').optional(), + embeddingDimension: z.literal(1536).optional(), + chunkingConfig: z + .object({ + maxSize: z.number(), + minSize: z.number(), + overlap: z.number(), + }) + .optional(), +}) + +async function checkKnowledgeBaseAccess(knowledgeBaseId: string, userId: string) { + const kb = await db + .select({ + id: knowledgeBase.id, + userId: knowledgeBase.userId, + }) + .from(knowledgeBase) + .where(and(eq(knowledgeBase.id, knowledgeBaseId), isNull(knowledgeBase.deletedAt))) + .limit(1) + + if (kb.length === 0) { + return { hasAccess: false, notFound: true } + } + + const kbData = kb[0] + + // Check if user owns the knowledge base + if (kbData.userId === userId) { + return { hasAccess: true, knowledgeBase: kbData } + } + + return { hasAccess: false, knowledgeBase: kbData } +} + +export async function GET(req: NextRequest, { params }: { params: Promise<{ id: string }> }) { + const requestId = crypto.randomUUID().slice(0, 8) + const { id } = await params + + try { + const session = await getSession() + if (!session?.user?.id) { + logger.warn(`[${requestId}] Unauthorized knowledge base access attempt`) + return NextResponse.json({ error: 'Unauthorized' }, { status: 401 }) + } + + const accessCheck = await checkKnowledgeBaseAccess(id, session.user.id) + + if (accessCheck.notFound) { + logger.warn(`[${requestId}] Knowledge base not found: ${id}`) + return NextResponse.json({ error: 'Knowledge base not found' }, { status: 404 }) + } + + if (!accessCheck.hasAccess) { + logger.warn( + `[${requestId}] User ${session.user.id} attempted to access unauthorized knowledge base ${id}` + ) + return NextResponse.json({ error: 'Unauthorized' }, { status: 401 }) + } + + const knowledgeBases = await db + .select() + .from(knowledgeBase) + .where(and(eq(knowledgeBase.id, id), isNull(knowledgeBase.deletedAt))) + .limit(1) + + if (knowledgeBases.length === 0) { + return NextResponse.json({ error: 'Knowledge base not found' }, { status: 404 }) + } + + logger.info(`[${requestId}] Retrieved knowledge base: ${id} for user ${session.user.id}`) + + return NextResponse.json({ + success: true, + data: knowledgeBases[0], + }) + } catch (error) { + logger.error(`[${requestId}] Error fetching knowledge base`, error) + return NextResponse.json({ error: 'Failed to fetch knowledge base' }, { status: 500 }) + } +} + +export async function PUT(req: NextRequest, { params }: { params: Promise<{ id: string }> }) { + const requestId = crypto.randomUUID().slice(0, 8) + const { id } = await params + + try { + const session = await getSession() + if (!session?.user?.id) { + logger.warn(`[${requestId}] Unauthorized knowledge base update attempt`) + return NextResponse.json({ error: 'Unauthorized' }, { status: 401 }) + } + + const accessCheck = await checkKnowledgeBaseAccess(id, session.user.id) + + if (accessCheck.notFound) { + logger.warn(`[${requestId}] Knowledge base not found: ${id}`) + return NextResponse.json({ error: 'Knowledge base not found' }, { status: 404 }) + } + + if (!accessCheck.hasAccess) { + logger.warn( + `[${requestId}] User ${session.user.id} attempted to update unauthorized knowledge base ${id}` + ) + return NextResponse.json({ error: 'Unauthorized' }, { status: 401 }) + } + + const body = await req.json() + + try { + const validatedData = UpdateKnowledgeBaseSchema.parse(body) + + const updateData: any = { + updatedAt: new Date(), + } + + if (validatedData.name !== undefined) updateData.name = validatedData.name + if (validatedData.description !== undefined) + updateData.description = validatedData.description + + // Handle embedding model and dimension together to ensure consistency + if ( + validatedData.embeddingModel !== undefined || + validatedData.embeddingDimension !== undefined + ) { + updateData.embeddingModel = 'text-embedding-3-small' + updateData.embeddingDimension = 1536 + } + + if (validatedData.chunkingConfig !== undefined) + updateData.chunkingConfig = validatedData.chunkingConfig + + await db.update(knowledgeBase).set(updateData).where(eq(knowledgeBase.id, id)) + + // Fetch the updated knowledge base + const updatedKnowledgeBase = await db + .select() + .from(knowledgeBase) + .where(eq(knowledgeBase.id, id)) + .limit(1) + + logger.info(`[${requestId}] Knowledge base updated: ${id} for user ${session.user.id}`) + + return NextResponse.json({ + success: true, + data: updatedKnowledgeBase[0], + }) + } catch (validationError) { + if (validationError instanceof z.ZodError) { + logger.warn(`[${requestId}] Invalid knowledge base update data`, { + errors: validationError.errors, + }) + return NextResponse.json( + { error: 'Invalid request data', details: validationError.errors }, + { status: 400 } + ) + } + throw validationError + } + } catch (error) { + logger.error(`[${requestId}] Error updating knowledge base`, error) + return NextResponse.json({ error: 'Failed to update knowledge base' }, { status: 500 }) + } +} + +export async function DELETE(req: NextRequest, { params }: { params: Promise<{ id: string }> }) { + const requestId = crypto.randomUUID().slice(0, 8) + const { id } = await params + + try { + const session = await getSession() + if (!session?.user?.id) { + logger.warn(`[${requestId}] Unauthorized knowledge base delete attempt`) + return NextResponse.json({ error: 'Unauthorized' }, { status: 401 }) + } + + const accessCheck = await checkKnowledgeBaseAccess(id, session.user.id) + + if (accessCheck.notFound) { + logger.warn(`[${requestId}] Knowledge base not found: ${id}`) + return NextResponse.json({ error: 'Knowledge base not found' }, { status: 404 }) + } + + if (!accessCheck.hasAccess) { + logger.warn( + `[${requestId}] User ${session.user.id} attempted to delete unauthorized knowledge base ${id}` + ) + return NextResponse.json({ error: 'Unauthorized' }, { status: 401 }) + } + + // Soft delete by setting deletedAt timestamp + await db + .update(knowledgeBase) + .set({ + deletedAt: new Date(), + updatedAt: new Date(), + }) + .where(eq(knowledgeBase.id, id)) + + logger.info(`[${requestId}] Knowledge base deleted: ${id} for user ${session.user.id}`) + + return NextResponse.json({ + success: true, + data: { message: 'Knowledge base deleted successfully' }, + }) + } catch (error) { + logger.error(`[${requestId}] Error deleting knowledge base`, error) + return NextResponse.json({ error: 'Failed to delete knowledge base' }, { status: 500 }) + } +} diff --git a/apps/sim/app/api/knowledge/route.ts b/apps/sim/app/api/knowledge/route.ts new file mode 100644 index 0000000000..3094040cd5 --- /dev/null +++ b/apps/sim/app/api/knowledge/route.ts @@ -0,0 +1,137 @@ +import { and, count, eq, isNull } from 'drizzle-orm' +import { type NextRequest, NextResponse } from 'next/server' +import { z } from 'zod' +import { getSession } from '@/lib/auth' +import { createLogger } from '@/lib/logs/console-logger' +import { db } from '@/db' +import { document, knowledgeBase } from '@/db/schema' + +const logger = createLogger('KnowledgeBaseAPI') + +// Schema for knowledge base creation +const CreateKnowledgeBaseSchema = z.object({ + name: z.string().min(1, 'Name is required'), + description: z.string().optional(), + workspaceId: z.string().optional(), + embeddingModel: z.literal('text-embedding-3-small').default('text-embedding-3-small'), + embeddingDimension: z.literal(1536).default(1536), + chunkingConfig: z + .object({ + maxSize: z.number().default(1024), + minSize: z.number().default(100), + overlap: z.number().default(200), + }) + .default({}), +}) + +export async function GET(req: NextRequest) { + const requestId = crypto.randomUUID().slice(0, 8) + + try { + const session = await getSession() + if (!session?.user?.id) { + logger.warn(`[${requestId}] Unauthorized knowledge base access attempt`) + return NextResponse.json({ error: 'Unauthorized' }, { status: 401 }) + } + + // Build where conditions + const whereConditions = [ + eq(knowledgeBase.userId, session.user.id), + isNull(knowledgeBase.deletedAt), + ] + + // Get knowledge bases with document counts + const knowledgeBasesWithCounts = await db + .select({ + id: knowledgeBase.id, + name: knowledgeBase.name, + description: knowledgeBase.description, + tokenCount: knowledgeBase.tokenCount, + embeddingModel: knowledgeBase.embeddingModel, + embeddingDimension: knowledgeBase.embeddingDimension, + chunkingConfig: knowledgeBase.chunkingConfig, + createdAt: knowledgeBase.createdAt, + updatedAt: knowledgeBase.updatedAt, + workspaceId: knowledgeBase.workspaceId, + docCount: count(document.id), + }) + .from(knowledgeBase) + .leftJoin( + document, + and(eq(document.knowledgeBaseId, knowledgeBase.id), isNull(document.deletedAt)) + ) + .where(and(...whereConditions)) + .groupBy(knowledgeBase.id) + .orderBy(knowledgeBase.createdAt) + + logger.info( + `[${requestId}] Retrieved ${knowledgeBasesWithCounts.length} knowledge bases for user ${session.user.id}` + ) + + return NextResponse.json({ + success: true, + data: knowledgeBasesWithCounts, + }) + } catch (error) { + logger.error(`[${requestId}] Error fetching knowledge bases`, error) + return NextResponse.json({ error: 'Failed to fetch knowledge bases' }, { status: 500 }) + } +} + +export async function POST(req: NextRequest) { + const requestId = crypto.randomUUID().slice(0, 8) + + try { + const session = await getSession() + if (!session?.user?.id) { + logger.warn(`[${requestId}] Unauthorized knowledge base creation attempt`) + return NextResponse.json({ error: 'Unauthorized' }, { status: 401 }) + } + + const body = await req.json() + + try { + const validatedData = CreateKnowledgeBaseSchema.parse(body) + + const id = crypto.randomUUID() + const now = new Date() + + const newKnowledgeBase = { + id, + userId: session.user.id, + workspaceId: validatedData.workspaceId || null, + name: validatedData.name, + description: validatedData.description || null, + tokenCount: 0, + embeddingModel: validatedData.embeddingModel, + embeddingDimension: validatedData.embeddingDimension, + chunkingConfig: validatedData.chunkingConfig, + createdAt: now, + updatedAt: now, + } + + await db.insert(knowledgeBase).values(newKnowledgeBase) + + logger.info(`[${requestId}] Knowledge base created: ${id} for user ${session.user.id}`) + + return NextResponse.json({ + success: true, + data: newKnowledgeBase, + }) + } catch (validationError) { + if (validationError instanceof z.ZodError) { + logger.warn(`[${requestId}] Invalid knowledge base data`, { + errors: validationError.errors, + }) + return NextResponse.json( + { error: 'Invalid request data', details: validationError.errors }, + { status: 400 } + ) + } + throw validationError + } + } catch (error) { + logger.error(`[${requestId}] Error creating knowledge base`, error) + return NextResponse.json({ error: 'Failed to create knowledge base' }, { status: 500 }) + } +} diff --git a/apps/sim/app/api/knowledge/search/route.ts b/apps/sim/app/api/knowledge/search/route.ts new file mode 100644 index 0000000000..4176d85980 --- /dev/null +++ b/apps/sim/app/api/knowledge/search/route.ts @@ -0,0 +1,166 @@ +import { and, eq, isNull, sql } from 'drizzle-orm' +import { type NextRequest, NextResponse } from 'next/server' +import { z } from 'zod' +import { getSession } from '@/lib/auth' +import { env } from '@/lib/env' +import { createLogger } from '@/lib/logs/console-logger' +import { db } from '@/db' +import { embedding, knowledgeBase } from '@/db/schema' + +const logger = createLogger('VectorSearchAPI') + +// Schema for vector search request +const VectorSearchSchema = z.object({ + knowledgeBaseId: z.string().min(1, 'Knowledge base ID is required'), + query: z.string().min(1, 'Search query is required'), + topK: z.number().min(1).max(100).default(10), +}) + +async function generateSearchEmbedding(query: string): Promise { + const openaiApiKey = env.OPENAI_API_KEY + if (!openaiApiKey) { + throw new Error('OPENAI_API_KEY not configured') + } + + try { + const response = await fetch('https://api.openai.com/v1/embeddings', { + method: 'POST', + headers: { + Authorization: `Bearer ${openaiApiKey}`, + 'Content-Type': 'application/json', + }, + body: JSON.stringify({ + input: query, + model: 'text-embedding-3-small', + encoding_format: 'float', + }), + }) + + if (!response.ok) { + const errorText = await response.text() + throw new Error(`OpenAI API error: ${response.status} ${response.statusText} - ${errorText}`) + } + + const data = await response.json() + + if (!data.data || !Array.isArray(data.data) || data.data.length === 0) { + throw new Error('Invalid response format from OpenAI embeddings API') + } + + return data.data[0].embedding + } catch (error) { + logger.error('Failed to generate search embedding:', error) + throw new Error( + `Embedding generation failed: ${error instanceof Error ? error.message : 'Unknown error'}` + ) + } +} + +export async function POST(request: NextRequest) { + const requestId = crypto.randomUUID().slice(0, 8) + + try { + logger.info(`[${requestId}] Processing vector search request`) + + const session = await getSession() + if (!session?.user?.id) { + logger.warn(`[${requestId}] Unauthorized vector search attempt`) + return NextResponse.json({ error: 'Unauthorized' }, { status: 401 }) + } + + const body = await request.json() + + try { + const validatedData = VectorSearchSchema.parse(body) + + // Verify the knowledge base exists and user has access + const kb = await db + .select() + .from(knowledgeBase) + .where( + and( + eq(knowledgeBase.id, validatedData.knowledgeBaseId), + eq(knowledgeBase.userId, session.user.id), + isNull(knowledgeBase.deletedAt) + ) + ) + .limit(1) + + if (kb.length === 0) { + logger.warn( + `[${requestId}] Knowledge base not found or access denied: ${validatedData.knowledgeBaseId}` + ) + return NextResponse.json( + { error: 'Knowledge base not found or access denied' }, + { status: 404 } + ) + } + + // Generate embedding for the search query + logger.info(`[${requestId}] Generating embedding for search query`) + const queryEmbedding = await generateSearchEmbedding(validatedData.query) + + // Perform vector similarity search using pgvector cosine similarity + logger.info(`[${requestId}] Performing vector search with topK=${validatedData.topK}`) + + const results = await db + .select({ + id: embedding.id, + content: embedding.content, + documentId: embedding.documentId, + chunkIndex: embedding.chunkIndex, + metadata: embedding.metadata, + similarity: sql`1 - (${embedding.embedding} <=> ${JSON.stringify(queryEmbedding)}::vector)`, + }) + .from(embedding) + .where( + and( + eq(embedding.knowledgeBaseId, validatedData.knowledgeBaseId), + eq(embedding.enabled, true) + ) + ) + .orderBy(sql`${embedding.embedding} <=> ${JSON.stringify(queryEmbedding)}::vector`) + .limit(validatedData.topK) + + logger.info(`[${requestId}] Vector search completed. Found ${results.length} results`) + + return NextResponse.json({ + success: true, + data: { + results: results.map((result) => ({ + id: result.id, + content: result.content, + documentId: result.documentId, + chunkIndex: result.chunkIndex, + metadata: result.metadata, + similarity: result.similarity, + })), + query: validatedData.query, + knowledgeBaseId: validatedData.knowledgeBaseId, + topK: validatedData.topK, + totalResults: results.length, + }, + }) + } catch (validationError) { + if (validationError instanceof z.ZodError) { + logger.warn(`[${requestId}] Invalid vector search data`, { + errors: validationError.errors, + }) + return NextResponse.json( + { error: 'Invalid request data', details: validationError.errors }, + { status: 400 } + ) + } + throw validationError + } + } catch (error) { + logger.error(`[${requestId}] Error performing vector search`, error) + return NextResponse.json( + { + error: 'Failed to perform vector search', + message: error instanceof Error ? error.message : 'Unknown error', + }, + { status: 500 } + ) + } +} diff --git a/apps/sim/app/w/[id]/components/control-bar/control-bar.tsx b/apps/sim/app/w/[id]/components/control-bar/control-bar.tsx index 1d2ccfd868..9ef1a4bf0f 100644 --- a/apps/sim/app/w/[id]/components/control-bar/control-bar.tsx +++ b/apps/sim/app/w/[id]/components/control-bar/control-bar.tsx @@ -236,7 +236,9 @@ export function ControlBar() { // Use cache if available and not expired if (!forceRefresh && usageDataCache.data && cacheAge < usageDataCache.expirationMs) { - logger.info('Using cached usage data', { cacheAge: `${Math.round(cacheAge / 1000)}s` }) + logger.info('Using cached usage data', { + cacheAge: `${Math.round(cacheAge / 1000)}s`, + }) return usageDataCache.data } @@ -658,7 +660,7 @@ export function ControlBar() { // size="icon" // onClick={handlePublishWorkflow} // disabled={isPublishing} - // className={cn('hover:text-[#802FFF]', isPublished && 'text-[#802FFF]')} + // className={cn('hover:text-[#701FFC]', isPublished && 'text-[#701FFC]')} // > // {isPublishing ? ( // @@ -869,13 +871,13 @@ export function ControlBar() { + + + +
+
+ {/* Scrollable Content */} +
+
+
+ +