From ab3fa4838458c8303a67dd30fdd75a16b89cc20b Mon Sep 17 00:00:00 2001 From: Nicolas Date: Fri, 10 Oct 2025 14:57:26 -0300 Subject: [PATCH] Nick: init --- .gitignore | 2 + apps/api/package.json | 1 + apps/api/pnpm-lock.yaml | 9 + apps/api/src/lib/search-index/chunker.ts | 358 ++++++++++ apps/api/src/lib/search-index/embeddings.ts | 291 ++++++++ apps/api/src/lib/search-index/index.ts | 69 ++ .../src/lib/search-index/pinecone-service.ts | 316 ++++++++ apps/api/src/lib/search-index/query.ts | 641 +++++++++++++++++ apps/api/src/lib/search-index/queue.ts | 164 +++++ apps/api/src/lib/search-index/service.ts | 674 ++++++++++++++++++ .../scraper/scrapeURL/transformers/index.ts | 4 +- .../transformers/sendToSearchIndex.ts | 154 ++++ apps/api/src/services/index.ts | 4 + .../api/src/services/indexing/index-worker.ts | 32 + apps/api/src/services/search-index-db.ts | 106 +++ 15 files changed, 2824 insertions(+), 1 deletion(-) create mode 100644 apps/api/src/lib/search-index/chunker.ts create mode 100644 apps/api/src/lib/search-index/embeddings.ts create mode 100644 apps/api/src/lib/search-index/index.ts create mode 100644 apps/api/src/lib/search-index/pinecone-service.ts create mode 100644 apps/api/src/lib/search-index/query.ts create mode 100644 apps/api/src/lib/search-index/queue.ts create mode 100644 apps/api/src/lib/search-index/service.ts create mode 100644 apps/api/src/scraper/scrapeURL/transformers/sendToSearchIndex.ts create mode 100644 apps/api/src/services/search-index-db.ts diff --git a/.gitignore b/.gitignore index 79d3a4db1..9a97f69a3 100644 --- a/.gitignore +++ b/.gitignore @@ -47,3 +47,5 @@ CLAUDE.local.md *.egg-info/ # local SDK venv apps/python-sdk/.venv/ + +/apps/api/running-docs/ diff --git a/apps/api/package.json b/apps/api/package.json index c169436ed..ba8436608 100644 --- a/apps/api/package.json +++ b/apps/api/package.json @@ -71,6 +71,7 @@ "@ai-sdk/fireworks": "^0.2.4", "@ai-sdk/google": "^1.2.3", "@ai-sdk/google-vertex": "^2.2.15", + "@pinecone-database/pinecone": "^6.1.2", "@ai-sdk/groq": "^1.2.1", "@ai-sdk/openai": "^1.3.12", "@apidevtools/json-schema-ref-parser": "^11.7.3", diff --git a/apps/api/pnpm-lock.yaml b/apps/api/pnpm-lock.yaml index f034a8c69..3f2b02dd5 100644 --- a/apps/api/pnpm-lock.yaml +++ b/apps/api/pnpm-lock.yaml @@ -76,6 +76,9 @@ importers: '@opentelemetry/sdk-node': specifier: ^0.205.0 version: 0.205.0(@opentelemetry/api@1.9.0) + '@pinecone-database/pinecone': + specifier: ^6.1.2 + version: 6.1.2 '@sentry/cli': specifier: ^2.33.1 version: 2.33.1(encoding@0.1.13) @@ -2641,6 +2644,10 @@ packages: resolution: {integrity: sha512-IHnV6A+zxU7XwmKFinmYjUcwlyK9+xkG3/s9KcQhI9BjQKycrJ1JRO+FbNYPwZiPKW3je/DR0k7w8/gLa5eaxQ==} deprecated: 'The package is now available as "qr": npm install qr' + '@pinecone-database/pinecone@6.1.2': + resolution: {integrity: sha512-ydIlbtgIIHFgBL08sPzua5ckmOgtjgDz8xg21CnP1fqnnEgDmOlnfd10MRKU+fvFRhDlh4Md37SwZDr0d4cBqg==} + engines: {node: '>=18.0.0'} + '@pkgjs/parseargs@0.11.0': resolution: {integrity: sha512-+1VkjdD0QBLPodGrJUeqarH8VAIvQODIbwh9XpP5Syisf7YoQgsJKPNFoqqLQlu+VQ/tVSshMR6loPMn8U+dPg==} engines: {node: '>=14'} @@ -9937,6 +9944,8 @@ snapshots: '@paulmillr/qr@0.2.1': {} + '@pinecone-database/pinecone@6.1.2': {} + '@pkgjs/parseargs@0.11.0': optional: true diff --git a/apps/api/src/lib/search-index/chunker.ts b/apps/api/src/lib/search-index/chunker.ts new file mode 100644 index 000000000..71241a477 --- /dev/null +++ b/apps/api/src/lib/search-index/chunker.ts @@ -0,0 +1,358 @@ +import { logger as _logger } from "../logger"; +import type { Logger } from "winston"; + +export interface TextChunk { + text: string; + ordinal: number; + tokenCount: number; + charCount: number; + startOffset: number; + endOffset: number; +} + +export interface ChunkingOptions { + targetTokens?: number; // Target tokens per chunk (default: 750) + minTokens?: number; // Minimum tokens per chunk (default: 600) + maxTokens?: number; // Maximum tokens per chunk (default: 900) + overlapTokens?: number; // Overlap between chunks (default: 100) + preserveStructure?: boolean; // Preserve markdown structure (default: true) +} + +const DEFAULT_OPTIONS: Required = { + targetTokens: 750, + minTokens: 600, + maxTokens: 900, + overlapTokens: 100, + preserveStructure: true, +}; + +/** + * Estimate token count using simple heuristic (1 token ≈ 4 chars for English) + * This is approximate but fast. For exact counts, use tiktoken. + */ +export function estimateTokenCount(text: string): number { + // Remove extra whitespace + const normalized = text.replace(/\s+/g, " ").trim(); + // Rough approximation: 1 token ≈ 4 characters + return Math.ceil(normalized.length / 4); +} + +/** + * Split text into sentences, preserving structure + */ +function splitIntoSentences(text: string): string[] { + // Split on sentence boundaries, but preserve structure + const sentences: string[] = []; + + // Match sentences ending with .!? followed by space/newline/end + // Also handle abbreviations like "Dr." "Mr." "etc." + const sentenceRegex = /(? 0) { + sentences.push(sentence); + } + lastIndex = match.index + match[0].length; + } + + // Add remaining text + if (lastIndex < text.length) { + const remaining = text.slice(lastIndex).trim(); + if (remaining.length > 0) { + sentences.push(remaining); + } + } + + return sentences.filter(s => s.length > 0); +} + +/** + * Detect markdown structure boundaries (headings, lists, code blocks) + */ +function detectStructureBoundaries(text: string): number[] { + const boundaries: number[] = [0]; + + const lines = text.split("\n"); + let offset = 0; + + for (let i = 0; i < lines.length; i++) { + const line = lines[i]; + + // Detect headings + if (/^#{1,6}\s/.test(line)) { + boundaries.push(offset); + } + + // Detect code blocks + if (/^```/.test(line)) { + boundaries.push(offset); + } + + // Detect list items (with proper structure) + if (/^[\s]*[-*+]\s/.test(line) || /^[\s]*\d+\.\s/.test(line)) { + boundaries.push(offset); + } + + offset += line.length + 1; // +1 for newline + } + + boundaries.push(text.length); + + return [...new Set(boundaries)].sort((a, b) => a - b); +} + +/** + * Main chunking function: splits text into semantically coherent chunks + */ +export async function chunkText( + text: string, + options: ChunkingOptions = {}, + logger?: Logger, +): Promise { + const opts = { ...DEFAULT_OPTIONS, ...options }; + const log = logger ?? _logger.child({ module: "search-chunker" }); + + // Handle empty or very short text + if (!text || text.trim().length === 0) { + return []; + } + + const totalTokens = estimateTokenCount(text); + + // If text is short enough, return as single chunk + if (totalTokens <= opts.maxTokens) { + return [{ + text: text.trim(), + ordinal: 0, + tokenCount: totalTokens, + charCount: text.length, + startOffset: 0, + endOffset: text.length, + }]; + } + + const chunks: TextChunk[] = []; + + // Strategy 1: If preserving structure, respect markdown boundaries + if (opts.preserveStructure) { + const boundaries = detectStructureBoundaries(text); + const sections = boundaries.slice(0, -1).map((start, i) => ({ + text: text.slice(start, boundaries[i + 1]), + offset: start, + })); + + let currentChunk = ""; + let currentOffset = 0; + let currentTokens = 0; + + for (const section of sections) { + const sectionTokens = estimateTokenCount(section.text); + + // If section alone exceeds max, split it further + if (sectionTokens > opts.maxTokens) { + // Flush current chunk if any + if (currentChunk.length > 0) { + chunks.push({ + text: currentChunk.trim(), + ordinal: chunks.length, + tokenCount: currentTokens, + charCount: currentChunk.length, + startOffset: currentOffset, + endOffset: currentOffset + currentChunk.length, + }); + currentChunk = ""; + currentTokens = 0; + } + + // Split large section by sentences + const sentences = splitIntoSentences(section.text); + let sentenceChunk = ""; + let sentenceTokens = 0; + + for (const sentence of sentences) { + const sentenceTokenCount = estimateTokenCount(sentence); + + if (sentenceTokens + sentenceTokenCount > opts.maxTokens && sentenceChunk.length > 0) { + chunks.push({ + text: sentenceChunk.trim(), + ordinal: chunks.length, + tokenCount: sentenceTokens, + charCount: sentenceChunk.length, + startOffset: section.offset, + endOffset: section.offset + sentenceChunk.length, + }); + + // Add overlap + const overlapText = sentence.slice(-opts.overlapTokens * 4); // Rough char estimate + sentenceChunk = overlapText + " " + sentence; + sentenceTokens = estimateTokenCount(sentenceChunk); + } else { + sentenceChunk += (sentenceChunk.length > 0 ? " " : "") + sentence; + sentenceTokens += sentenceTokenCount; + } + } + + if (sentenceChunk.length > 0) { + chunks.push({ + text: sentenceChunk.trim(), + ordinal: chunks.length, + tokenCount: sentenceTokens, + charCount: sentenceChunk.length, + startOffset: section.offset, + endOffset: section.offset + sentenceChunk.length, + }); + } + + currentOffset = section.offset + section.text.length; + currentChunk = ""; + currentTokens = 0; + continue; + } + + // Try to add section to current chunk + if (currentTokens + sectionTokens <= opts.maxTokens) { + currentChunk += (currentChunk.length > 0 ? "\n" : "") + section.text; + currentTokens += sectionTokens; + if (currentChunk.length === section.text.length) { + currentOffset = section.offset; + } + } else { + // Current chunk is full, flush it + if (currentChunk.length > 0) { + chunks.push({ + text: currentChunk.trim(), + ordinal: chunks.length, + tokenCount: currentTokens, + charCount: currentChunk.length, + startOffset: currentOffset, + endOffset: currentOffset + currentChunk.length, + }); + } + + // Start new chunk with current section + currentChunk = section.text; + currentOffset = section.offset; + currentTokens = sectionTokens; + } + } + + // Flush remaining chunk + if (currentChunk.length > 0) { + chunks.push({ + text: currentChunk.trim(), + ordinal: chunks.length, + tokenCount: currentTokens, + charCount: currentChunk.length, + startOffset: currentOffset, + endOffset: currentOffset + currentChunk.length, + }); + } + } else { + // Strategy 2: Simple sentence-based chunking without structure preservation + const sentences = splitIntoSentences(text); + let currentChunk = ""; + let currentTokens = 0; + let currentOffset = 0; + + for (const sentence of sentences) { + const sentenceTokens = estimateTokenCount(sentence); + + if (currentTokens + sentenceTokens > opts.maxTokens && currentChunk.length > 0) { + chunks.push({ + text: currentChunk.trim(), + ordinal: chunks.length, + tokenCount: currentTokens, + charCount: currentChunk.length, + startOffset: currentOffset, + endOffset: currentOffset + currentChunk.length, + }); + + // Add overlap for context + const lastSentences = currentChunk.split(/[.!?]+\s+/).slice(-2).join(". "); + currentChunk = lastSentences + " " + sentence; + currentTokens = estimateTokenCount(currentChunk); + } else { + currentChunk += (currentChunk.length > 0 ? " " : "") + sentence; + currentTokens += sentenceTokens; + } + + if (chunks.length === 0 && currentChunk === sentence) { + currentOffset = 0; + } + } + + // Flush remaining + if (currentChunk.length > 0) { + chunks.push({ + text: currentChunk.trim(), + ordinal: chunks.length, + tokenCount: currentTokens, + charCount: currentChunk.length, + startOffset: currentOffset, + endOffset: text.length, + }); + } + } + + log.debug("Chunked text", { + totalTokens, + chunks: chunks.length, + avgTokensPerChunk: chunks.length > 0 + ? Math.round(chunks.reduce((sum, c) => sum + c.tokenCount, 0) / chunks.length) + : 0, + }); + + return chunks; +} + +/** + * Extract clean text from markdown for indexing + * Removes markdown syntax but preserves structure and content + */ +export function cleanMarkdownForIndexing(markdown: string): string { + if (!markdown) return ""; + + let text = markdown; + + // Remove code blocks but keep language hints + text = text.replace(/```(\w+)?\n([\s\S]*?)```/g, (_, lang, code) => { + return lang ? `Code (${lang}): ${code.split("\n").slice(0, 3).join(" ")}...` : ""; + }); + + // Remove inline code (keep content) + text = text.replace(/`([^`]+)`/g, "$1"); + + // Remove images (keep alt text) + text = text.replace(/!\[([^\]]*)\]\([^)]+\)/g, "$1"); + + // Remove links (keep text) + text = text.replace(/\[([^\]]+)\]\([^)]+\)/g, "$1"); + + // Remove headers (#) but keep text + text = text.replace(/^#{1,6}\s+(.+)$/gm, "$1"); + + // Remove bold/italic + text = text.replace(/(\*\*|__)(.*?)\1/g, "$2"); + text = text.replace(/(\*|_)(.*?)\1/g, "$2"); + + // Remove list markers but keep structure + text = text.replace(/^[\s]*[-*+]\s+/gm, ""); + text = text.replace(/^[\s]*\d+\.\s+/gm, ""); + + // Remove blockquotes + text = text.replace(/^>\s+/gm, ""); + + // Remove horizontal rules + text = text.replace(/^[-*_]{3,}$/gm, ""); + + // Normalize whitespace + text = text.replace(/\n{3,}/g, "\n\n"); + text = text.replace(/ +/g, " "); + + return text.trim(); +} + diff --git a/apps/api/src/lib/search-index/embeddings.ts b/apps/api/src/lib/search-index/embeddings.ts new file mode 100644 index 000000000..582aaab75 --- /dev/null +++ b/apps/api/src/lib/search-index/embeddings.ts @@ -0,0 +1,291 @@ +import { embed, EmbeddingModel } from "ai"; +import { getEmbeddingModel } from "../generic-ai"; +import { logger as _logger } from "../logger"; +import type { Logger } from "winston"; +import { withSpan, setSpanAttributes } from "../otel-tracer"; + +export interface EmbeddingResult { + embedding: number[]; + tokens: number; +} + +export interface BatchEmbeddingResult { + embeddings: number[][]; + totalTokens: number; + failed: number[]; +} + +export interface EmbeddingOptions { + model?: string; + dimensions?: number; + maxRetries?: number; + retryDelay?: number; + batchSize?: number; +} + +const DEFAULT_OPTIONS: Required = { + model: "text-embedding-3-small", + dimensions: 1536, + maxRetries: 3, + retryDelay: 1000, // ms + batchSize: 100, +}; + +/** + * Generate embedding for a single text + */ +export async function generateEmbedding( + text: string, + options: EmbeddingOptions = {}, + logger?: Logger, +): Promise { + return await withSpan("firecrawl-generate-embedding", async span => { + const opts = { ...DEFAULT_OPTIONS, ...options }; + const log = logger ?? _logger.child({ module: "search-embeddings" }); + + setSpanAttributes(span, { + "embedding.model": opts.model, + "embedding.dimensions": opts.dimensions, + "embedding.text_length": text.length, + }); + + if (!text || text.trim().length === 0) { + setSpanAttributes(span, { "embedding.empty_text": true }); + return { + embedding: new Array(opts.dimensions).fill(0), + tokens: 0, + }; + } + + // Truncate if too long (OpenAI limit: 8191 tokens for text-embedding-3-small) + const maxChars = 8191 * 4; // Rough approximation + const truncatedText = text.length > maxChars ? text.slice(0, maxChars) : text; + + let lastError: Error | null = null; + + for (let attempt = 0; attempt < opts.maxRetries; attempt++) { + try { + const embeddingModel: EmbeddingModel = getEmbeddingModel( + opts.model, + "openai", + ); + + const result = await embed({ + model: embeddingModel, + value: truncatedText, + }); + + setSpanAttributes(span, { + "embedding.success": true, + "embedding.tokens": result.usage.tokens, + "embedding.attempt": attempt + 1, + }); + + log.debug("Generated embedding", { + textLength: text.length, + tokens: result.usage.tokens, + dimensions: result.embedding.length, + }); + + return { + embedding: result.embedding, + tokens: result.usage.tokens, + }; + } catch (error) { + lastError = error as Error; + + setSpanAttributes(span, { + "embedding.error": true, + "embedding.attempt": attempt + 1, + "embedding.error_message": lastError.message, + }); + + log.warn("Failed to generate embedding", { + error: lastError.message, + attempt: attempt + 1, + maxRetries: opts.maxRetries, + }); + + // Exponential backoff + if (attempt < opts.maxRetries - 1) { + const delay = opts.retryDelay * Math.pow(2, attempt); + await new Promise(resolve => setTimeout(resolve, delay)); + } + } + } + + log.error("Failed to generate embedding after all retries", { + error: lastError?.message, + maxRetries: opts.maxRetries, + }); + + // Return zero vector as fallback + return { + embedding: new Array(opts.dimensions).fill(0), + tokens: 0, + }; + }); +} + +/** + * Generate embeddings for multiple texts in batches + * More efficient for bulk operations + */ +export async function generateEmbeddingsBatch( + texts: string[], + options: EmbeddingOptions = {}, + logger?: Logger, +): Promise { + return await withSpan("firecrawl-generate-embeddings-batch", async span => { + const opts = { ...DEFAULT_OPTIONS, ...options }; + const log = logger ?? _logger.child({ module: "search-embeddings-batch" }); + + setSpanAttributes(span, { + "embedding.batch_size": texts.length, + "embedding.model": opts.model, + }); + + if (texts.length === 0) { + return { + embeddings: [], + totalTokens: 0, + failed: [], + }; + } + + const embeddings: number[][] = new Array(texts.length); + const failed: number[] = []; + let totalTokens = 0; + + // Process in batches + const batches: string[][] = []; + for (let i = 0; i < texts.length; i += opts.batchSize) { + batches.push(texts.slice(i, i + opts.batchSize)); + } + + log.info("Processing embedding batches", { + totalTexts: texts.length, + batches: batches.length, + batchSize: opts.batchSize, + }); + + for (let batchIndex = 0; batchIndex < batches.length; batchIndex++) { + const batch = batches[batchIndex]; + const batchOffset = batchIndex * opts.batchSize; + + log.debug("Processing batch", { + batch: batchIndex + 1, + total: batches.length, + size: batch.length, + }); + + // Process batch items in parallel + const batchResults = await Promise.allSettled( + batch.map((text, i) => + generateEmbedding(text, options, log).then(result => ({ + index: batchOffset + i, + ...result, + })) + ), + ); + + // Collect results + for (const result of batchResults) { + if (result.status === "fulfilled") { + const { index, embedding, tokens } = result.value; + embeddings[index] = embedding; + totalTokens += tokens; + } else { + const index = batchOffset + batchResults.indexOf(result); + failed.push(index); + embeddings[index] = new Array(opts.dimensions).fill(0); + + log.warn("Failed to generate embedding in batch", { + index, + error: result.reason?.message, + }); + } + } + + // Rate limiting: small delay between batches + if (batchIndex < batches.length - 1) { + await new Promise(resolve => setTimeout(resolve, 100)); + } + } + + setSpanAttributes(span, { + "embedding.total_tokens": totalTokens, + "embedding.failed_count": failed.length, + "embedding.success_rate": ((texts.length - failed.length) / texts.length), + }); + + log.info("Batch embedding complete", { + total: texts.length, + successful: texts.length - failed.length, + failed: failed.length, + totalTokens, + }); + + return { + embeddings, + totalTokens, + failed, + }; + }); +} + +/** + * Estimate embedding cost for text + * Based on OpenAI pricing: $0.0001 / 1M tokens + */ +export function estimateEmbeddingCost(tokenCount: number): number { + return (tokenCount / 1_000_000) * 0.0001; +} + +/** + * Check if embeddings are available (API key configured) + */ +export function isEmbeddingEnabled(): boolean { + return !!process.env.OPENAI_API_KEY; +} + +/** + * Normalize embedding vector (convert to unit vector) + * Useful for cosine similarity when using dot product + */ +export function normalizeEmbedding(embedding: number[]): number[] { + const magnitude = Math.sqrt( + embedding.reduce((sum, val) => sum + val * val, 0), + ); + + if (magnitude === 0) return embedding; + + return embedding.map(val => val / magnitude); +} + +/** + * Calculate cosine similarity between two embeddings + */ +export function cosineSimilarity(a: number[], b: number[]): number { + if (a.length !== b.length) { + throw new Error("Embeddings must have the same dimensions"); + } + + let dotProduct = 0; + let magnitudeA = 0; + let magnitudeB = 0; + + for (let i = 0; i < a.length; i++) { + dotProduct += a[i] * b[i]; + magnitudeA += a[i] * a[i]; + magnitudeB += b[i] * b[i]; + } + + magnitudeA = Math.sqrt(magnitudeA); + magnitudeB = Math.sqrt(magnitudeB); + + if (magnitudeA === 0 || magnitudeB === 0) return 0; + + return dotProduct / (magnitudeA * magnitudeB); +} + diff --git a/apps/api/src/lib/search-index/index.ts b/apps/api/src/lib/search-index/index.ts new file mode 100644 index 000000000..420922ce4 --- /dev/null +++ b/apps/api/src/lib/search-index/index.ts @@ -0,0 +1,69 @@ +/** + * Search Index Module + * + * Real-time search index on top of Firecrawl's web scraping infrastructure. + * Combines keyword (BM25) and semantic (vector) search with RRF ranking. + * + * + * Architecture: + * - Ingest: Crawler → Text Normalization → Chunking → Embeddings → Postgres + * - Query: Hybrid Search (BM25 + Vector) → RRF → Filters → Results + * - Storage: search_documents + search_chunks tables + * - Embeddings: OpenAI text-embedding-3-small (1536 dims) + */ + +export { + chunkText, + cleanMarkdownForIndexing, + estimateTokenCount, + type TextChunk, + type ChunkingOptions, +} from "./chunker"; + +export { + generateEmbedding, + generateEmbeddingsBatch, + isEmbeddingEnabled, + estimateEmbeddingCost, + normalizeEmbedding, + cosineSimilarity, + type EmbeddingResult, + type BatchEmbeddingResult, + type EmbeddingOptions, +} from "./embeddings"; + +export { + indexDocumentForSearch, + deleteDocumentFromSearch, + searchDocumentExists, + type SearchDocumentInput, + type SearchIndexResult, +} from "./service"; + +export { + search, + searchChunks, + getSearchStats, + type SearchQuery, + type SearchFilters, + type SearchResult, + type SearchResponse, +} from "./query"; + +export { + addSearchIndexJob, + processSearchIndexJobs, + getSearchIndexQueueLength, +} from "./queue"; + +export { + upsertToPinecone, + searchPinecone, + deleteFromPinecone, + getPineconeStats, + isPineconeEnabled, + buildPineconeFilter, + getPineconeIndex, + type PineconeRecord, +} from "./pinecone-service"; + diff --git a/apps/api/src/lib/search-index/pinecone-service.ts b/apps/api/src/lib/search-index/pinecone-service.ts new file mode 100644 index 000000000..99e587ced --- /dev/null +++ b/apps/api/src/lib/search-index/pinecone-service.ts @@ -0,0 +1,316 @@ +/** + * Pinecone Vector Database Service + * + * Handles billion-scale vector storage and similarity search. + * Replaces pgvector for scalability. + * + * Architecture: + * - Pinecone: Vector embeddings (HNSW index) + * - Postgres: Metadata + BM25 (full-text search) + * + * Benefits: + * - Scales to 10B+ vectors + * - 50-100ms query latency + * - Auto-scaling and replication + * - No RAM limitations + */ + +import { Pinecone } from "@pinecone-database/pinecone"; +import { logger as _logger } from "../logger"; +import type { Logger } from "winston"; +import { withSpan, setSpanAttributes } from "../otel-tracer"; + +// Pinecone client (singleton) +let pineconeClient: Pinecone | null = null; +let pineconeIndex: any = null; + +/** + * Initialize Pinecone client + */ +function getPineconeClient(): Pinecone { + if (!pineconeClient) { + const apiKey = process.env.PINECONE_API_KEY; + + if (!apiKey) { + throw new Error("PINECONE_API_KEY not set"); + } + + pineconeClient = new Pinecone({ + apiKey, + }); + + _logger.info("Pinecone client initialized"); + } + + return pineconeClient; +} + +/** + * Get Pinecone index + */ +export function getPineconeIndex() { + if (!pineconeIndex) { + const client = getPineconeClient(); + const indexName = process.env.PINECONE_INDEX_NAME || "firecrawl-search"; + + pineconeIndex = client.index(indexName); + + _logger.info("Pinecone index initialized", { indexName }); + } + + return pineconeIndex; +} + +/** + * Check if Pinecone is enabled + */ +export function isPineconeEnabled(): boolean { + return !!( + process.env.PINECONE_API_KEY && + process.env.PINECONE_INDEX_NAME + ); +} + +export interface PineconeRecord { + id: string; + values: number[]; + metadata: { + url: string; + domain: string; + title?: string; + freshness_score: number; + quality_score: number; + country?: string; + is_mobile: boolean; + chunk_ordinal?: number; // For chunks + doc_id?: string; // For chunks + }; +} + +/** + * Upsert embeddings to Pinecone + */ +export async function upsertToPinecone( + records: PineconeRecord[], + namespace: string = "documents", + logger?: Logger, +): Promise { + return await withSpan("firecrawl-pinecone-upsert", async span => { + const log = logger ?? _logger.child({ module: "pinecone-service" }); + + setSpanAttributes(span, { + "pinecone.operation": "upsert", + "pinecone.namespace": namespace, + "pinecone.record_count": records.length, + }); + + if (records.length === 0) { + return; + } + + try { + const index = getPineconeIndex(); + const ns = index.namespace(namespace); + + // Upsert in batches of 100 (Pinecone limit) + const batchSize = 100; + for (let i = 0; i < records.length; i += batchSize) { + const batch = records.slice(i, i + batchSize); + + await ns.upsert(batch); + + log.debug("Upserted batch to Pinecone", { + batch: Math.floor(i / batchSize) + 1, + size: batch.length, + }); + } + + setSpanAttributes(span, { + "pinecone.upsert_successful": true, + }); + + log.info("Upserted to Pinecone", { + namespace, + records: records.length, + }); + } catch (error) { + log.error("Failed to upsert to Pinecone", { + error: (error as Error).message, + namespace, + recordCount: records.length, + }); + + setSpanAttributes(span, { + "pinecone.upsert_error": true, + "pinecone.error_message": (error as Error).message, + }); + + throw error; + } + }); +} + +/** + * Search Pinecone for similar vectors + */ +export async function searchPinecone( + queryEmbedding: number[], + limit: number = 100, + filter?: Record, + namespace: string = "documents", + logger?: Logger, +): Promise; +}>> { + return await withSpan("firecrawl-pinecone-query", async span => { + const log = logger ?? _logger.child({ module: "pinecone-service" }); + + setSpanAttributes(span, { + "pinecone.operation": "query", + "pinecone.namespace": namespace, + "pinecone.limit": limit, + "pinecone.has_filter": !!filter, + }); + + try { + const index = getPineconeIndex(); + const ns = index.namespace(namespace); + + const queryResponse = await ns.query({ + vector: queryEmbedding, + topK: limit, + filter, + includeMetadata: true, + includeValues: false, + }); + + const results = (queryResponse.matches || []).map(match => ({ + id: match.id, + score: match.score || 0, + metadata: match.metadata || {}, + })); + + setSpanAttributes(span, { + "pinecone.query_successful": true, + "pinecone.results_count": results.length, + }); + + log.debug("Queried Pinecone", { + namespace, + limit, + results: results.length, + }); + + return results; + } catch (error) { + log.error("Failed to query Pinecone", { + error: (error as Error).message, + namespace, + limit, + }); + + setSpanAttributes(span, { + "pinecone.query_error": true, + "pinecone.error_message": (error as Error).message, + }); + + throw error; + } + }); +} + +/** + * Delete from Pinecone + */ +export async function deleteFromPinecone( + ids: string[], + namespace: string = "documents", + logger?: Logger, +): Promise { + const log = logger ?? _logger.child({ module: "pinecone-service" }); + + try { + const index = getPineconeIndex(); + const ns = index.namespace(namespace); + + await ns.deleteMany(ids); + + log.info("Deleted from Pinecone", { + namespace, + count: ids.length, + }); + } catch (error) { + log.error("Failed to delete from Pinecone", { + error: (error as Error).message, + namespace, + count: ids.length, + }); + } +} + +/** + * Get Pinecone index stats + */ +export async function getPineconeStats( + namespace: string = "documents", +): Promise<{ + vectorCount: number; + dimension: number; + indexFullness: number; +}> { + try { + const index = getPineconeIndex(); + const stats = await index.describeIndexStats(); + + const nsStats = stats.namespaces?.[namespace]; + + return { + vectorCount: nsStats?.recordCount || 0, + dimension: stats.dimension || 1536, + indexFullness: stats.indexFullness || 0, + }; + } catch (error) { + _logger.error("Failed to get Pinecone stats", { + error: (error as Error).message, + }); + + return { + vectorCount: 0, + dimension: 1536, + indexFullness: 0, + }; + } +} + +/** + * Build Pinecone filter from search filters + */ +export function buildPineconeFilter(filters: { + domain?: string; + country?: string; + isMobile?: boolean; + minFreshness?: number; +}): Record | undefined { + const filter: Record = {}; + + if (filters.domain) { + filter.domain = { $eq: filters.domain }; + } + + if (filters.country) { + filter.country = { $eq: filters.country }; + } + + if (filters.isMobile !== undefined) { + filter.is_mobile = { $eq: filters.isMobile }; + } + + if (filters.minFreshness !== undefined) { + filter.freshness_score = { $gte: filters.minFreshness }; + } + + return Object.keys(filter).length > 0 ? filter : undefined; +} + diff --git a/apps/api/src/lib/search-index/query.ts b/apps/api/src/lib/search-index/query.ts new file mode 100644 index 000000000..8c1f60d09 --- /dev/null +++ b/apps/api/src/lib/search-index/query.ts @@ -0,0 +1,641 @@ +import { SupabaseClient } from "@supabase/supabase-js"; +import { logger as _logger } from "../logger"; +import type { Logger } from "winston"; +import { withSpan, setSpanAttributes } from "../otel-tracer"; +import { generateEmbedding, isEmbeddingEnabled } from "./embeddings"; +import { + searchPinecone, + isPineconeEnabled, + buildPineconeFilter, +} from "./pinecone-service"; + +export interface SearchQuery { + query: string; + limit?: number; + offset?: number; + filters?: SearchFilters; + mode?: "hybrid" | "keyword" | "semantic"; +} + +export interface SearchFilters { + domain?: string; + country?: string; + isMobile?: boolean; + minFreshness?: number; + language?: string; +} + +export interface SearchResult { + documentId: string; + url: string; + title: string | null; + description: string | null; + domain: string; + snippet?: string; + score: number; + bm25Rank: number | null; + vectorRank: number | null; + freshnessScore: number; + qualityScore: number; + lastCrawledAt: string; +} + +export interface SearchResponse { + results: SearchResult[]; + total: number; + query: string; + mode: string; + took: number; // milliseconds +} + +/** + * Main search function: hybrid search with RRF ranking + */ +export async function search( + supabase: SupabaseClient, + searchQuery: SearchQuery, + logger?: Logger, +): Promise { + return await withSpan("firecrawl-search", async span => { + const log = logger ?? _logger.child({ module: "search-query" }); + const startTime = Date.now(); + + const { + query, + limit = 50, + offset = 0, + filters = {}, + mode = "hybrid", + } = searchQuery; + + setSpanAttributes(span, { + "search.query": query, + "search.limit": limit, + "search.mode": mode, + "search.filters": JSON.stringify(filters), + }); + + if (!query || query.trim().length === 0) { + return { + results: [], + total: 0, + query, + mode, + took: Date.now() - startTime, + }; + } + + try { + let results: SearchResult[] = []; + + // Choose search strategy based on mode + if (mode === "keyword") { + results = await keywordSearch(supabase, query, limit, filters, log); + } else if (mode === "semantic") { + results = await semanticSearch(supabase, query, limit, filters, log); + } else { + // Hybrid: combine BM25 + vector with RRF + results = await hybridSearch(supabase, query, limit, filters, log); + } + + // Apply offset + const paginatedResults = results.slice(offset, offset + limit); + + const took = Date.now() - startTime; + + setSpanAttributes(span, { + "search.results_count": results.length, + "search.took_ms": took, + }); + + log.info("Search completed", { + query, + mode, + results: results.length, + took, + }); + + return { + results: paginatedResults, + total: results.length, + query, + mode, + took, + }; + } catch (error) { + log.error("Search failed", { + error: (error as Error).message, + query, + }); + + setSpanAttributes(span, { + "search.error": true, + "search.error_message": (error as Error).message, + }); + + return { + results: [], + total: 0, + query, + mode, + took: Date.now() - startTime, + }; + } + }); +} + +/** + * Hybrid search: BM25 (Postgres) + Vector (Pinecone) with RRF ranking + */ +async function hybridSearch( + supabase: SupabaseClient, + query: string, + limit: number, + filters: SearchFilters, + logger: Logger, +): Promise { + // Generate query embedding for vector search + let queryEmbedding: number[] | null = null; + + if (isEmbeddingEnabled() && isPineconeEnabled()) { + try { + const embResult = await generateEmbedding(query, {}, logger); + queryEmbedding = embResult.embedding; + } catch (error) { + logger.warn("Failed to generate query embedding, falling back to keyword search", { + error: (error as Error).message, + }); + } + } + + // If no embedding, fall back to keyword search + if (!queryEmbedding) { + return keywordSearch(supabase, query, limit, filters, logger); + } + + // Parallel search: BM25 in Postgres + Vector in Pinecone + const [bm25Results, vectorResults] = await Promise.all([ + // BM25 search in Postgres + supabase.rpc("bm25_search", { + query_text: query, + result_limit: 100, + country_filter: filters.country ?? null, + domain_filter: filters.domain ?? null, + is_mobile_filter: filters.isMobile ?? null, + min_freshness: filters.minFreshness ?? 0.0, + }), + + // Vector search in Pinecone + (async () => { + try { + const pineconeFilter = buildPineconeFilter(filters); + const results = await searchPinecone( + queryEmbedding!, + 100, + pineconeFilter, + "documents", + logger, + ); + return results; + } catch (error) { + logger.warn("Pinecone search failed, continuing with BM25 only", { + error: (error as Error).message, + }); + return []; + } + })(), + ]); + + if (bm25Results.error) { + logger.error("BM25 search failed", { error: bm25Results.error.message }); + throw new Error(`BM25 search failed: ${bm25Results.error.message}`); + } + + // Extract Pinecone doc IDs + const pineconeDocIds = vectorResults + .map(r => r.metadata.doc_id) + .filter(Boolean) as string[]; + + // Get metadata for Pinecone results from Postgres + const { data: pineconeMetadata } = await supabase.rpc("get_documents_metadata", { + doc_ids: pineconeDocIds, + }); + + // Create map of doc_id -> metadata + const metadataMap = new Map( + (pineconeMetadata || []).map((m: any) => [m.doc_id, m]) + ); + + // Merge results with RRF + const bm25Map = new Map( + (bm25Results.data || []).map((r: any, index: number) => [ + r.doc_id, + { rank: index + 1, ...r }, + ]) + ); + + const vectorMap = new Map( + vectorResults.map((r, index) => [ + r.metadata.doc_id || r.id, + { rank: index + 1, score: r.score }, + ]) + ); + + // Combine all unique doc IDs + const allDocIds = new Set([ + ...bm25Map.keys(), + ...vectorMap.keys(), + ]); + + const k = 60; // RRF constant + const merged: SearchResult[] = []; + + for (const docId of allDocIds) { + const bm25 = bm25Map.get(docId); + const vector = vectorMap.get(docId); + + // RRF score + const bm25Rank = bm25?.rank ?? null; + const vectorRank = vector?.rank ?? null; + + const rrf = + (bm25Rank ? 1 / (k + bm25Rank) : 0) + + (vectorRank ? 1 / (k + vectorRank) : 0); + + // Get metadata (prefer from BM25 result, fallback to Pinecone metadata) + const metadata: any = bm25 || metadataMap.get(docId); + + if (!metadata) continue; + + const combinedScore = + rrf * (metadata.freshness_score || 1.0) * (metadata.quality_score || 1.0); + + merged.push({ + documentId: docId, + url: metadata.url, + title: metadata.title, + description: metadata.description, + domain: metadata.domain, + score: combinedScore, + bm25Rank, + vectorRank, + freshnessScore: metadata.freshness_score || 1.0, + qualityScore: metadata.quality_score || 1.0, + lastCrawledAt: metadata.last_crawled_at, + }); + } + + // Sort by combined score + merged.sort((a, b) => b.score - a.score); + + return merged.slice(0, limit); +} + +/** + * Keyword-only search (BM25) + */ +async function keywordSearch( + supabase: SupabaseClient, + query: string, + limit: number, + filters: SearchFilters, + logger: Logger, +): Promise { + let queryBuilder = supabase + .from("search_documents") + .select( + "id, resolved_url, title, description, domain, freshness_score, quality_score, last_crawled_at, content_ts", + ) + .textSearch("content_ts", query, { + type: "plain", + config: "english", + }); + + // Apply filters + if (filters.domain) { + queryBuilder = queryBuilder.eq("domain", filters.domain); + } + if (filters.country) { + queryBuilder = queryBuilder.eq("country", filters.country); + } + if (filters.isMobile !== undefined) { + queryBuilder = queryBuilder.eq("is_mobile", filters.isMobile); + } + if (filters.minFreshness) { + queryBuilder = queryBuilder.gte("freshness_score", filters.minFreshness); + } + if (filters.language) { + queryBuilder = queryBuilder.eq("language", filters.language); + } + + queryBuilder = queryBuilder.limit(limit); + + const { data, error } = await queryBuilder; + + if (error) { + logger.error("Keyword search failed", { error: error.message }); + throw new Error(`Keyword search failed: ${error.message}`); + } + + return (data ?? []).map((row: any, index: number) => ({ + documentId: row.id, + url: row.resolved_url, + title: row.title, + description: row.description, + domain: row.domain, + score: row.freshness_score * row.quality_score * (1 / (index + 1)), + bm25Rank: index + 1, + vectorRank: null, + freshnessScore: row.freshness_score, + qualityScore: row.quality_score, + lastCrawledAt: row.last_crawled_at, + })); +} + +/** + * Semantic-only search (vector similarity via Pinecone) + */ +async function semanticSearch( + supabase: SupabaseClient, + query: string, + limit: number, + filters: SearchFilters, + logger: Logger, +): Promise { + if (!isEmbeddingEnabled() || !isPineconeEnabled()) { + logger.warn("Embeddings or Pinecone not enabled, falling back to keyword search"); + return keywordSearch(supabase, query, limit, filters, logger); + } + + // Generate query embedding + const embResult = await generateEmbedding(query, {}, logger); + const queryEmbedding = embResult.embedding; + + // Search Pinecone with filters + const pineconeFilter = buildPineconeFilter(filters); + const vectorResults = await searchPinecone( + queryEmbedding, + limit, + pineconeFilter, + "documents", + logger, + ); + + if (vectorResults.length === 0) { + return []; + } + + // Get metadata from Postgres + const docIds = vectorResults + .map(r => r.metadata.doc_id) + .filter(Boolean) as string[]; + + const { data: metadata, error } = await supabase.rpc("get_documents_metadata", { + doc_ids: docIds, + }); + + if (error) { + logger.error("Failed to fetch metadata for semantic search", { error: error.message }); + throw new Error(`Metadata fetch failed: ${error.message}`); + } + + // Create metadata map + const metadataMap = new Map( + (metadata || []).map((m: any) => [m.doc_id, m]) + ); + + // Combine Pinecone results with Postgres metadata + return vectorResults.map((result, index) => { + const docId = result.metadata.doc_id || result.id; + const meta: any = metadataMap.get(docId); + + if (!meta) { + // Fallback to Pinecone metadata + return { + documentId: docId, + url: result.metadata.url || "", + title: result.metadata.title || null, + description: null, + domain: result.metadata.domain || "", + score: result.score * (result.metadata.freshness_score || 1.0) * (result.metadata.quality_score || 1.0), + bm25Rank: null, + vectorRank: index + 1, + freshnessScore: result.metadata.freshness_score || 1.0, + qualityScore: result.metadata.quality_score || 1.0, + lastCrawledAt: new Date().toISOString(), + }; + } + + return { + documentId: docId, + url: meta.url, + title: meta.title, + description: meta.description, + domain: meta.domain, + score: result.score * meta.freshness_score * meta.quality_score, + bm25Rank: null, + vectorRank: index + 1, + freshnessScore: meta.freshness_score, + qualityScore: meta.quality_score, + lastCrawledAt: meta.last_crawled_at, + }; + }).filter(r => r.url !== ""); +} + +/** + * Search within chunks (for precise snippet retrieval via Pinecone) + */ +export async function searchChunks( + supabase: SupabaseClient, + query: string, + limit: number = 20, + filters: SearchFilters = {}, + logger?: Logger, +): Promise> { + const log = logger ?? _logger.child({ module: "search-chunks" }); + + // Generate query embedding + let queryEmbedding: number[] | null = null; + + if (isEmbeddingEnabled() && isPineconeEnabled()) { + try { + const embResult = await generateEmbedding(query, {}, log); + queryEmbedding = embResult.embedding; + } catch (error) { + log.warn("Failed to generate query embedding for chunk search", { + error: (error as Error).message, + }); + } + } + + if (!queryEmbedding) { + // Fall back to keyword search on chunks (Postgres) + const { data, error } = await supabase + .from("search_chunks") + .select("id, doc_id, text, ordinal") + .textSearch("text_ts", query, { + type: "plain", + config: "english", + }) + .not("text", "is", null) + .limit(limit); + + if (error) { + log.error("Chunk keyword search failed", { error: error.message }); + return []; + } + + // Join with documents to get metadata + const docIds = [...new Set((data ?? []).map(c => c.doc_id))]; + const { data: docsData } = await supabase + .from("search_documents") + .select("id, resolved_url, title") + .in("id", docIds); + + const docMap = new Map((docsData ?? []).map(d => [d.id, d])); + + return (data ?? []).map((chunk, index) => { + const doc = docMap.get(chunk.doc_id); + return { + chunkId: chunk.id, + documentId: chunk.doc_id, + url: doc?.resolved_url ?? "", + title: doc?.title ?? null, + text: chunk.text || "", + score: 1 / (index + 1), + ordinal: chunk.ordinal, + }; + }); + } + + // Vector chunk search via Pinecone + const pineconeFilter = buildPineconeFilter(filters); + + // Add chunk-specific filter + const chunkFilter = { + ...pineconeFilter, + chunk_ordinal: { $exists: true }, // Only get chunk records, not document records + }; + + const vectorResults = await searchPinecone( + queryEmbedding, + limit, + chunkFilter, + "documents", + log, + ); + + if (vectorResults.length === 0) { + return []; + } + + // Get chunk data from Postgres (for text and ordinal) + const docIds = [...new Set( + vectorResults.map(r => r.metadata.doc_id).filter(Boolean) + )] as string[]; + + const { data: chunksData } = await supabase + .from("search_chunks") + .select("id, doc_id, text, ordinal") + .in("doc_id", docIds); + + const { data: docsData } = await supabase + .from("search_documents") + .select("id, resolved_url, title") + .in("id", docIds); + + const chunkMap = new Map( + (chunksData ?? []) + .filter(c => c.doc_id && c.ordinal !== null) + .map(c => [`${c.doc_id}_${c.ordinal}`, c]) + ); + const docMap = new Map((docsData ?? []).map(d => [d.id, d])); + + // Combine Pinecone results with Postgres data + return vectorResults.map((result, index) => { + const docId = result.metadata.doc_id; + const ordinal = result.metadata.chunk_ordinal; + const chunkKey = `${docId}_${ordinal}`; + + const chunk = chunkMap.get(chunkKey); + const doc = docMap.get(docId); + + return { + chunkId: result.id, + documentId: docId || "", + url: doc?.resolved_url || result.metadata.url || "", + title: doc?.title || result.metadata.title || null, + text: chunk?.text || "", + score: result.score, + ordinal: ordinal ?? 0, + }; + }).filter(r => r.url !== ""); +} + +/** + * Get search statistics (including Pinecone stats) + */ +export async function getSearchStats( + supabase: SupabaseClient, +): Promise<{ + totalDocuments: number; + totalChunks: number; + documentsWithEmbeddings: number; + chunksWithEmbeddings: number; + avgFreshness: number; + avgQuality: number; + uniqueDomains: number; + pineconeVectors?: number; + pineconeIndexFullness?: number; +}> { + const { data, error } = await supabase.from("search_index_stats").select("*"); + + // Get Pinecone stats if available + let pineconeStats = { vectorCount: 0, indexFullness: 0 }; + if (isPineconeEnabled()) { + try { + const { getPineconeStats } = await import("./pinecone-service.js"); + pineconeStats = await getPineconeStats("documents"); + } catch (error) { + _logger.warn("Failed to get Pinecone stats", { + error: (error as Error).message, + }); + } + } + + if (error || !data || data.length === 0) { + return { + totalDocuments: 0, + totalChunks: 0, + documentsWithEmbeddings: 0, + chunksWithEmbeddings: 0, + avgFreshness: 0, + avgQuality: 0, + uniqueDomains: 0, + pineconeVectors: pineconeStats.vectorCount, + pineconeIndexFullness: pineconeStats.indexFullness, + }; + } + + const stats = data[0]; + return { + totalDocuments: stats.total_documents ?? 0, + totalChunks: stats.total_chunks ?? 0, + documentsWithEmbeddings: stats.documents_in_pinecone ?? 0, + chunksWithEmbeddings: stats.chunks_in_pinecone ?? 0, + avgFreshness: stats.avg_freshness ?? 0, + avgQuality: stats.avg_quality ?? 0, + uniqueDomains: stats.unique_domains ?? 0, + pineconeVectors: pineconeStats.vectorCount, + pineconeIndexFullness: pineconeStats.indexFullness, + }; +} + diff --git a/apps/api/src/lib/search-index/queue.ts b/apps/api/src/lib/search-index/queue.ts new file mode 100644 index 000000000..a8c6c87b1 --- /dev/null +++ b/apps/api/src/lib/search-index/queue.ts @@ -0,0 +1,164 @@ +/** + * Search Index Queue Processor + * + * Handles async indexing jobs using Redis queue. + * Similar to existing index insert queue, but for search index. + * + * Flow: + * 1. Scraper adds job to queue after successful scrape + * 2. Background worker picks up jobs in batches + * 3. Text is chunked and embeddings generated + * 4. Data is written to search_documents and search_chunks + */ + +import { redisEvictConnection } from "../../services/redis"; +import { logger as _logger } from "../logger"; +import { indexDocumentForSearch, type SearchDocumentInput } from "./service"; +import { search_index_supabase_service } from "../../services/search-index-db"; + +const SEARCH_INDEX_QUEUE_KEY = "search-index-queue"; +const SEARCH_INDEX_BATCH_SIZE = 10; // Smaller batch size due to embedding generation + +export interface SearchIndexJob { + url: string; + resolvedUrl: string; + title?: string; + description?: string; + markdown: string; + html: string; + statusCode: number; + gcsPath?: string; + screenshotUrl?: string; + language?: string; + country?: string; + isMobile?: boolean; +} + +/** + * Add a job to the search index queue + */ +export async function addSearchIndexJob(job: SearchIndexJob): Promise { + try { + await redisEvictConnection.rpush( + SEARCH_INDEX_QUEUE_KEY, + JSON.stringify(job), + ); + + _logger.debug("Added job to search index queue", { + url: job.url, + }); + } catch (error) { + _logger.error("Failed to add job to search index queue", { + error: (error as Error).message, + url: job.url, + }); + } +} + +/** + * Get jobs from the queue + */ +async function getSearchIndexJobs(): Promise { + try { + const jobs = + (await redisEvictConnection.lpop( + SEARCH_INDEX_QUEUE_KEY, + SEARCH_INDEX_BATCH_SIZE, + )) ?? []; + + return jobs.map(x => JSON.parse(x) as SearchIndexJob); + } catch (error) { + _logger.error("Failed to get jobs from search index queue", { + error: (error as Error).message, + }); + return []; + } +} + +/** + * Process search index jobs (called by background worker) + */ +export async function processSearchIndexJobs(): Promise { + const jobs = await getSearchIndexJobs(); + + if (jobs.length === 0) { + return; + } + + _logger.info("Processing search index jobs", { + jobCount: jobs.length, + }); + + let successCount = 0; + let failedCount = 0; + + // Process jobs sequentially to avoid overwhelming embedding API + for (const job of jobs) { + try { + const logger = _logger.child({ + module: "search-index-queue", + url: job.url, + }); + + const result = await indexDocumentForSearch( + search_index_supabase_service, + { + url: job.url, + resolvedUrl: job.resolvedUrl, + title: job.title, + description: job.description, + markdown: job.markdown, + html: job.html, + statusCode: job.statusCode, + gcsPath: job.gcsPath, + screenshotUrl: job.screenshotUrl, + language: job.language, + country: job.country, + isMobile: job.isMobile, + }, + logger, + ); + + if (result.error) { + failedCount++; + logger.error("Failed to index document for search", { + error: result.error, + }); + } else { + successCount++; + logger.info("Indexed document for search", { + documentId: result.documentId, + chunks: result.chunkCount, + tokens: result.totalTokens, + }); + } + } catch (error) { + failedCount++; + _logger.error("Failed to process search index job", { + error: (error as Error).message, + url: job.url, + }); + } + } + + _logger.info("Finished processing search index jobs", { + total: jobs.length, + successful: successCount, + failed: failedCount, + }); +} + +/** + * Get queue length + */ +export async function getSearchIndexQueueLength(): Promise { + try { + return (await redisEvictConnection.llen(SEARCH_INDEX_QUEUE_KEY)) ?? 0; + } catch (error) { + _logger.error("Failed to get search index queue length", { + error: (error as Error).message, + }); + return 0; + } +} + diff --git a/apps/api/src/lib/search-index/service.ts b/apps/api/src/lib/search-index/service.ts new file mode 100644 index 000000000..a8e24aa60 --- /dev/null +++ b/apps/api/src/lib/search-index/service.ts @@ -0,0 +1,674 @@ +import { SupabaseClient } from "@supabase/supabase-js"; +import { logger as _logger } from "../logger"; +import type { Logger } from "winston"; +import { withSpan, setSpanAttributes } from "../otel-tracer"; +import { + chunkText, + cleanMarkdownForIndexing, + estimateTokenCount, + type TextChunk, +} from "./chunker"; +import { + generateEmbedding, + generateEmbeddingsBatch, + isEmbeddingEnabled, + type EmbeddingResult, +} from "./embeddings"; +import { + upsertToPinecone, + isPineconeEnabled, + type PineconeRecord, +} from "./pinecone-service"; +import crypto from "crypto"; +import psl from "psl"; + +export interface SearchDocumentInput { + url: string; + resolvedUrl: string; + title?: string; + description?: string; + markdown: string; + html: string; + statusCode: number; + gcsPath?: string; + screenshotUrl?: string; + language?: string; + country?: string; + isMobile?: boolean; +} + +export interface SearchIndexResult { + documentId: string; + chunkCount: number; + totalTokens: number; + embeddingsGenerated: boolean; + error?: string; +} + +/** + * Main function: Index a document for search + */ +export async function indexDocumentForSearch( + supabase: SupabaseClient, + input: SearchDocumentInput, + logger?: Logger, +): Promise { + return await withSpan("firecrawl-index-document-for-search", async span => { + const log = logger ?? _logger.child({ module: "search-index-service" }); + + setSpanAttributes(span, { + "search.url": input.url, + "search.status_code": input.statusCode, + "search.is_mobile": input.isMobile ?? false, + }); + + try { + // 1. Extract and clean text from markdown + const cleanText = cleanMarkdownForIndexing(input.markdown); + + if (!cleanText || cleanText.length < 100) { + log.warn("Document too short to index", { + url: input.url, + textLength: cleanText.length, + }); + return { + documentId: "", + chunkCount: 0, + totalTokens: 0, + embeddingsGenerated: false, + error: "Document too short", + }; + } + + setSpanAttributes(span, { + "search.clean_text_length": cleanText.length, + "search.clean_text_tokens": estimateTokenCount(cleanText), + }); + + // 2. Chunk the text + const chunks = await chunkText(cleanText, { + targetTokens: 750, + minTokens: 600, + maxTokens: 900, + overlapTokens: 100, + preserveStructure: true, + }, log); + + if (chunks.length === 0) { + log.warn("No chunks generated", { url: input.url }); + return { + documentId: "", + chunkCount: 0, + totalTokens: 0, + embeddingsGenerated: false, + error: "No chunks generated", + }; + } + + setSpanAttributes(span, { + "search.chunk_count": chunks.length, + "search.total_tokens": chunks.reduce((sum, c) => sum + c.tokenCount, 0), + }); + + // 3. Calculate URL hash and domain + const urlHash = Buffer.from( + crypto.createHash("sha256").update(normalizeSearchURL(input.url)).digest("hex"), + "hex", + ); + const domain = extractDomain(input.resolvedUrl); + + // 4. Generate embeddings (if enabled) + const embeddingsEnabled = isEmbeddingEnabled() && isPineconeEnabled(); + let documentEmbedding: number[] | null = null; + let chunkEmbeddings: number[][] = []; + let totalTokens = 0; + let pineconeRecords: PineconeRecord[] = []; + + if (embeddingsEnabled) { + try { + // Generate document-level embedding (from title + description + first chunk) + const documentText = [ + input.title, + input.description, + chunks[0]?.text.slice(0, 500), + ] + .filter(Boolean) + .join(" "); + + const docEmbResult = await generateEmbedding(documentText, {}, log); + documentEmbedding = docEmbResult.embedding; + totalTokens += docEmbResult.tokens; + + // Generate chunk embeddings in batch + const chunkTexts = chunks.map(c => c.text); + const batchResult = await generateEmbeddingsBatch(chunkTexts, { + batchSize: 50, + }, log); + + chunkEmbeddings = batchResult.embeddings; + totalTokens += batchResult.totalTokens; + + // Prepare Pinecone records (document + chunks) + // Document-level record + pineconeRecords.push({ + id: `doc_${crypto.randomUUID()}`, + values: documentEmbedding, + metadata: { + url: input.resolvedUrl, + domain: extractDomain(input.resolvedUrl), + title: input.title, + freshness_score: 1.0, + quality_score: 1.0, // Will update later + country: input.country, + is_mobile: input.isMobile ?? false, + }, + }); + + // Chunk-level records + chunks.forEach((chunk, i) => { + if (chunkEmbeddings[i]) { + pineconeRecords.push({ + id: `chunk_${crypto.randomUUID()}`, + values: chunkEmbeddings[i], + metadata: { + url: input.resolvedUrl, + domain: extractDomain(input.resolvedUrl), + title: input.title, + freshness_score: 1.0, + quality_score: 1.0, + country: input.country, + is_mobile: input.isMobile ?? false, + chunk_ordinal: i, + doc_id: "", // Will update after DB insert + }, + }); + } + }); + + setSpanAttributes(span, { + "search.embeddings_generated": true, + "search.embedding_tokens": totalTokens, + "search.pinecone_records": pineconeRecords.length, + }); + } catch (error) { + log.error("Failed to generate embeddings", { + error: (error as Error).message, + url: input.url, + }); + setSpanAttributes(span, { + "search.embeddings_error": true, + "search.embeddings_error_message": (error as Error).message, + }); + } + } + + // 5. Create tsvector for full-text search + const contentTsVector = await createTsVector(supabase, cleanText); + + // 6. Upsert document into search_documents + // Check if document already exists by url_hash + // Note: url_hash is BYTEA, need to query with proper format + const urlHashHex = "\\x" + urlHash.toString('hex'); + const { data: existingDocs, error: lookupError } = await supabase + .from("search_documents") + .select("id, content_hash") + .eq("url_hash", urlHashHex) + .limit(1); + + if (lookupError) { + log.warn("Error looking up existing document, will attempt insert", { + error: lookupError.message, + url: input.url, + }); + } + + const existingDoc = existingDocs?.[0]; + const newContentHash = Buffer.from( + crypto.createHash("sha256").update(cleanText).digest("hex"), + "hex", + ); + + let documentId: string; + let isUpdate = false; + + if (existingDoc) { + // Document exists - update it + documentId = existingDoc.id; + isUpdate = true; + + // Delete existing chunks (will be replaced) + await supabase + .from("search_chunks") + .delete() + .eq("doc_id", documentId); + + // Update document + const { error: updateError } = await supabase + .from("search_documents") + .update({ + resolved_url: input.resolvedUrl, + title: input.title?.slice(0, 200) ?? null, + description: input.description?.slice(0, 500) ?? null, + language: input.language ?? "en", + gcs_path: input.gcsPath ?? null, + screenshot_url: input.screenshotUrl ?? null, + content_ts: contentTsVector, + country: input.country ?? null, + pinecone_synced: false, // Will sync after successful Pinecone upsert + is_mobile: input.isMobile ?? false, + status_code: input.statusCode, + content_hash: newContentHash, + freshness_score: 1.0, // Reset to fresh + quality_score: calculateQualityScore(input, chunks), + last_crawled_at: new Date().toISOString(), + updated_at: new Date().toISOString(), + }) + .eq("id", documentId); + + if (updateError) { + log.error("Failed to update document", { + error: updateError.message, + url: input.url, + documentId, + }); + + setSpanAttributes(span, { + "search.update_error": true, + "search.update_error_message": updateError.message, + }); + + return { + documentId: "", + chunkCount: 0, + totalTokens: 0, + embeddingsGenerated: false, + error: updateError.message, + }; + } + } else { + // New document - insert + documentId = crypto.randomUUID(); + + const { error: insertError } = await supabase + .from("search_documents") + .insert({ + id: documentId, + url_hash: urlHashHex, + url: normalizeSearchURL(input.url), + resolved_url: input.resolvedUrl, + title: input.title?.slice(0, 200) ?? null, + description: input.description?.slice(0, 500) ?? null, + language: input.language ?? "en", + gcs_path: input.gcsPath ?? null, + screenshot_url: input.screenshotUrl ?? null, + content_ts: contentTsVector, + domain, + country: input.country ?? null, + is_mobile: input.isMobile ?? false, + status_code: input.statusCode, + content_hash: newContentHash, + freshness_score: 1.0, + quality_score: calculateQualityScore(input, chunks), + last_crawled_at: new Date().toISOString(), + pinecone_synced: false, + }); + + if (insertError) { + log.error("Failed to insert document", { + error: insertError.message, + url: input.url, + }); + + setSpanAttributes(span, { + "search.insert_error": true, + "search.insert_error_message": insertError.message, + }); + + return { + documentId: "", + chunkCount: 0, + totalTokens: 0, + embeddingsGenerated: false, + error: insertError.message, + }; + } + } + + // 7. Insert chunks into search_chunks + // Only insert first 2 chunks for snippet display + // Rest are stored only in Pinecone (saves 80% storage!) + const chunkInserts = chunks + .slice(0, 2) // Only first 2 chunks + .map((chunk, i) => ({ + doc_id: documentId, + ordinal: i, + text: chunk.text, + token_count: chunk.tokenCount, + char_count: chunk.charCount, + prev_chunk_id: null, + pinecone_synced: false, + })); + + const { error: chunkError } = await supabase + .from("search_chunks") + .insert(chunkInserts); + + if (chunkError) { + log.error("Failed to insert chunks", { + error: chunkError.message, + url: input.url, + documentId, + }); + + // Try to clean up the document + await supabase.from("search_documents").delete().eq("id", documentId); + + return { + documentId: "", + chunkCount: 0, + totalTokens: 0, + embeddingsGenerated: false, + error: chunkError.message, + }; + } + + // 8. Upsert to Pinecone (async, don't block) + if (embeddingsEnabled && pineconeRecords.length > 0) { + // Update Pinecone records with actual doc_id + pineconeRecords.forEach(record => { + if (record.metadata.chunk_ordinal !== undefined) { + record.metadata.doc_id = documentId; + } + }); + + try { + await upsertToPinecone(pineconeRecords, "documents", log); + + // Mark as synced + await supabase + .from("search_documents") + .update({ + pinecone_synced: true, + pinecone_synced_at: new Date().toISOString(), + }) + .eq("id", documentId); + + await supabase + .from("search_chunks") + .update({ pinecone_synced: true }) + .eq("doc_id", documentId); + + log.info("Synced to Pinecone", { + documentId, + records: pineconeRecords.length, + }); + } catch (error) { + log.error("Failed to sync to Pinecone", { + error: (error as Error).message, + documentId, + }); + // Continue anyway - can retry later + } + } + + // 9. Update sync state + await updateSyncState(supabase, { + documentsIndexed: 1, + chunksIndexed: Math.min(chunks.length, 2), // Only inserted first 2 + embeddingsGenerated: embeddingsEnabled ? 1 + chunks.length : 0, + }); + + log.info("Document indexed successfully", { + url: input.url, + documentId, + chunks: chunks.length, + totalTokens, + embeddingsGenerated: embeddingsEnabled, + isUpdate, + }); + + setSpanAttributes(span, { + "search.success": true, + "search.document_id": documentId, + "search.is_update": isUpdate, + }); + + return { + documentId, + chunkCount: chunks.length, + totalTokens, + embeddingsGenerated: embeddingsEnabled, + }; + } catch (error) { + log.error("Failed to index document", { + error: (error as Error).message, + url: input.url, + }); + + setSpanAttributes(span, { + "search.error": true, + "search.error_message": (error as Error).message, + }); + + return { + documentId: "", + chunkCount: 0, + totalTokens: 0, + embeddingsGenerated: false, + error: (error as Error).message, + }; + } + }); +} + +/** + * Normalize URL for search index (remove query params, fragments, etc.) + */ +function normalizeSearchURL(url: string): string { + try { + const urlObj = new URL(url); + urlObj.hash = ""; + urlObj.search = ""; // Remove query params for canonical URL + urlObj.protocol = "https:"; + + if (urlObj.hostname.startsWith("www.")) { + urlObj.hostname = urlObj.hostname.slice(4); + } + + if (urlObj.pathname.endsWith("/")) { + urlObj.pathname = urlObj.pathname.slice(0, -1); + } + + return urlObj.toString(); + } catch { + return url; + } +} + +/** + * Extract domain from URL + */ +function extractDomain(url: string): string { + try { + const urlObj = new URL(url); + const parsed = psl.parse(urlObj.hostname); + + if (parsed.domain) { + return parsed.domain; + } + + return urlObj.hostname; + } catch { + return ""; + } +} + +/** + * Calculate quality score based on content characteristics + */ +function calculateQualityScore( + input: SearchDocumentInput, + chunks: TextChunk[], +): number { + let score = 1.0; + + // Factor 1: Has title and description + if (input.title && input.title.length > 10) score *= 1.1; + if (input.description && input.description.length > 50) score *= 1.1; + + // Factor 2: Content length (sweet spot: 2000-10000 chars) + const totalChars = chunks.reduce((sum, c) => sum + c.charCount, 0); + if (totalChars >= 2000 && totalChars <= 10000) { + score *= 1.2; + } else if (totalChars < 500) { + score *= 0.5; + } else if (totalChars > 50000) { + score *= 0.8; + } + + // Factor 3: Chunk count (good chunking indicates structured content) + if (chunks.length >= 3 && chunks.length <= 20) { + score *= 1.1; + } + + // Factor 4: Status code + if (input.statusCode >= 200 && input.statusCode < 300) { + score *= 1.0; + } else { + score *= 0.3; + } + + // Normalize to [0, 1] + return Math.min(Math.max(score, 0.1), 2.0); +} + +/** + * Create tsvector from text (for BM25 search) + */ +async function createTsVector( + supabase: SupabaseClient, + text: string, +): Promise { + try { + const { data, error } = await supabase.rpc("to_tsvector_english", { + text: text.slice(0, 100000), // Limit to 100k chars + }); + + if (error) { + _logger.warn("Failed to create tsvector", { error: error.message }); + return null; + } + + return data; + } catch (error) { + _logger.warn("Failed to create tsvector", { + error: (error as Error).message, + }); + return null; + } +} + +/** + * Update sync state with metrics + */ +async function updateSyncState( + supabase: SupabaseClient, + metrics: { + documentsIndexed?: number; + chunksIndexed?: number; + embeddingsGenerated?: number; + failedDocuments?: number; + }, +): Promise { + try { + const updates: any = { + updated_at: new Date().toISOString(), + }; + + if (metrics.documentsIndexed) { + // Use SQL to increment + await supabase.rpc("increment_search_sync_state", { + field: "total_documents_indexed", + amount: metrics.documentsIndexed, + }); + } + + if (metrics.chunksIndexed) { + await supabase.rpc("increment_search_sync_state", { + field: "total_chunks_indexed", + amount: metrics.chunksIndexed, + }); + } + + if (metrics.failedDocuments) { + await supabase.rpc("increment_search_sync_state", { + field: "failed_documents", + amount: metrics.failedDocuments, + }); + } + + // Update embedding quota tracking + const today = new Date().toISOString().split("T")[0]; + await supabase.rpc("update_embedding_quota", { + embeddings_count: metrics.embeddingsGenerated ?? 0, + today_date: today, + }); + } catch (error) { + _logger.warn("Failed to update sync state", { + error: (error as Error).message, + }); + } +} + +/** + * Delete document from search index + */ +export async function deleteDocumentFromSearch( + supabase: SupabaseClient, + urlHash: Buffer, + logger?: Logger, +): Promise { + const log = logger ?? _logger.child({ module: "search-index-service" }); + + try { + const urlHashHex = "\\x" + urlHash.toString('hex'); + const { error } = await supabase + .from("search_documents") + .delete() + .eq("url_hash", urlHashHex); + + if (error) { + log.error("Failed to delete document from search index", { + error: error.message, + }); + } + } catch (error) { + log.error("Failed to delete document from search index", { + error: (error as Error).message, + }); + } +} + +/** + * Check if document already exists in search index + */ +export async function searchDocumentExists( + supabase: SupabaseClient, + urlHash: Buffer, +): Promise { + try { + const urlHashHex = "\\x" + urlHash.toString('hex'); + const { data, error } = await supabase + .from("search_documents") + .select("id") + .eq("url_hash", urlHashHex) + .limit(1); + + if (error) return false; + + return (data?.length ?? 0) > 0; + } catch { + return false; + } +} + diff --git a/apps/api/src/scraper/scrapeURL/transformers/index.ts b/apps/api/src/scraper/scrapeURL/transformers/index.ts index 7093b90a1..12bbb1c18 100644 --- a/apps/api/src/scraper/scrapeURL/transformers/index.ts +++ b/apps/api/src/scraper/scrapeURL/transformers/index.ts @@ -12,8 +12,9 @@ import { performAgent } from "./agent"; import { performAttributes } from "./performAttributes"; import { deriveDiff } from "./diff"; -import { useIndex } from "../../../services/index"; +import { useIndex, useSearchIndex } from "../../../services/index"; import { sendDocumentToIndex } from "../engines/index/index"; +import { sendDocumentToSearchIndex } from "./sendToSearchIndex"; import { hasFormatOfType, hasAnyFormatOfTypes, @@ -336,6 +337,7 @@ const transformerStack: Transformer[] = [ deriveMetadataFromRawHTML, uploadScreenshot, ...(useIndex ? [sendDocumentToIndex] : []), + ...(useSearchIndex ? [sendDocumentToSearchIndex] : []), // Add to search index for real-time search performLLMExtract, performSummary, performAttributes, diff --git a/apps/api/src/scraper/scrapeURL/transformers/sendToSearchIndex.ts b/apps/api/src/scraper/scrapeURL/transformers/sendToSearchIndex.ts new file mode 100644 index 000000000..4fa44075b --- /dev/null +++ b/apps/api/src/scraper/scrapeURL/transformers/sendToSearchIndex.ts @@ -0,0 +1,154 @@ +/** + * Transformer: Send Document to Search Index + * + * Integrates with the existing scraper transformer stack. + * Queues documents for real-time search indexing. + * + * Sampling: Controlled via SEARCH_INDEX_SAMPLE_RATE (0.0-1.0) + * - 0.1 = 10% of documents indexed (recommended for initial rollout) + * - 1.0 = 100% of documents indexed (full production) + */ + +import { Document } from "../../../controllers/v1/types"; +import { addSearchIndexJob } from "../../../lib/search-index/queue"; +import { logger as _logger } from "../../../lib/logger"; +import { Meta } from ".."; + +/** + * Check if document should be indexed for search + */ +function shouldIndexForSearch(meta: Meta, document: Document): boolean { + + + if (meta.internalOptions.zeroDataRetention) { + return false; + } + + const statusCode = document.metadata.statusCode; + if (statusCode < 200 || statusCode >= 300) { + return false; + } + + // Check if markdown content exists and is substantial + const markdown = document.markdown ?? ""; + if (markdown.length < 200) { + return false; + } + + // Don't index if has auth headers (private content) + if ( + meta.options.headers && + (meta.options.headers["Authorization"] || + meta.options.headers["Cookie"]) + ) { + return false; + } + + // Don't index PDFs without parsing + const isPdf = document.metadata.contentType?.includes("pdf"); + if (isPdf && !document.markdown) { + return false; + } + + return true; +} + +/** + * Determine if this document should be sampled for indexing + */ +function shouldSampleDocument(): boolean { + // Get sample rate from environment (default 10% for safe rollout) + const sampleRateStr = process.env.SEARCH_INDEX_SAMPLE_RATE || "0.1"; + const sampleRate = parseFloat(sampleRateStr); + + // Validate sample rate + if (isNaN(sampleRate) || sampleRate < 0 || sampleRate > 1) { + _logger.warn("Invalid SEARCH_INDEX_SAMPLE_RATE, using 0.1 (10%)", { + value: sampleRateStr, + }); + return Math.random() < 0.1; + } + + // Sample based on rate + return Math.random() < sampleRate; +} + +/** + * Transformer: Send document to search index + */ +export async function sendDocumentToSearchIndex( + meta: Meta, + document: Document, +): Promise { + // Check if search indexing is enabled + const searchIndexEnabled = + process.env.ENABLE_SEARCH_INDEX === "true" && + process.env.SEARCH_INDEX_SUPABASE_URL && + process.env.SEARCH_INDEX_SUPABASE_SERVICE_TOKEN; + + meta.logger.debug("Sending document to search index", { + url: meta.url, + }); + if (!searchIndexEnabled) { + return document; + } + + // Apply sampling (canary rollout) + if (!shouldSampleDocument()) { + meta.logger.debug("Document not sampled for search indexing", { + url: meta.url, + sampleRate: process.env.SEARCH_INDEX_SAMPLE_RATE || "0.1", + }); + return document; + } + + // Check if document should be indexed + if (!shouldIndexForSearch(meta, document)) { + meta.logger.debug("Document not suitable for search index", { + url: meta.url, + statusCode: document.metadata.statusCode, + markdownLength: document.markdown?.length ?? 0, + }); + return document; + } + + // Queue job for async processing (don't block scraper) + (async () => { + try { + await addSearchIndexJob({ + url: meta.url, + resolvedUrl: + document.metadata.url ?? + document.metadata.sourceURL ?? + meta.rewrittenUrl ?? + meta.url, + title: + document.metadata.title ?? document.metadata.ogTitle ?? undefined, + description: + document.metadata.description ?? + document.metadata.ogDescription ?? + undefined, + markdown: document.markdown ?? "", + html: document.rawHtml ?? "", + statusCode: document.metadata.statusCode, + gcsPath: undefined, // Can be populated from GCS if needed + screenshotUrl: document.screenshot ?? undefined, + language: document.metadata.language ?? "en", + country: meta.options.location?.country ?? undefined, + isMobile: meta.options.mobile ?? false, + }); + + meta.logger.debug("Queued document for search indexing", { + url: meta.url, + }); + } catch (error) { + meta.logger.error("Failed to queue document for search indexing", { + error: (error as Error).message, + url: meta.url, + }); + } + })(); + + return document; +} + diff --git a/apps/api/src/services/index.ts b/apps/api/src/services/index.ts index a65d93377..103a890f8 100644 --- a/apps/api/src/services/index.ts +++ b/apps/api/src/services/index.ts @@ -203,6 +203,10 @@ export const useIndex = process.env.INDEX_SUPABASE_URL !== "" && process.env.INDEX_SUPABASE_URL !== undefined; +export const useSearchIndex = + process.env.SEARCH_INDEX_SUPABASE_URL !== "" && + process.env.SEARCH_INDEX_SUPABASE_URL !== undefined; + export function normalizeURLForIndex(url: string): string { const urlObj = new URL(url); urlObj.hash = ""; diff --git a/apps/api/src/services/indexing/index-worker.ts b/apps/api/src/services/indexing/index-worker.ts index 59d5b0037..af0b7615a 100644 --- a/apps/api/src/services/indexing/index-worker.ts +++ b/apps/api/src/services/indexing/index-worker.ts @@ -23,6 +23,7 @@ import { processOMCEJobs, processDomainFrequencyJobs, } from ".."; +import { processSearchIndexJobs } from "../../lib/search-index/queue"; import { processWebhookInsertJobs } from "../webhook"; import { scrapeOptions as scrapeOptionsSchema, @@ -331,6 +332,7 @@ const INDEX_INSERT_INTERVAL = 3000; const WEBHOOK_INSERT_INTERVAL = 15000; const OMCE_INSERT_INTERVAL = 5000; const DOMAIN_FREQUENCY_INTERVAL = 10000; +const SEARCH_INDEX_INTERVAL = 10000; // Process search index queue every 10 seconds // Start the workers (async () => { @@ -412,6 +414,35 @@ const DOMAIN_FREQUENCY_INTERVAL = 10000; ); }, DOMAIN_FREQUENCY_INTERVAL); + // Search index queue processor + const searchIndexInterval = setInterval(async () => { + if (isShuttingDown) { + return; + } + + // Only process if search index is enabled + if (process.env.ENABLE_SEARCH_INDEX !== "true") { + return; + } + + await withSpan( + "firecrawl-index-worker-process-search-index-jobs", + async span => { + setSpanAttributes(span, { + "index.worker.operation": "process_search_index_jobs", + "index.worker.type": "scheduled", + }); + + try { + await processSearchIndexJobs(); + } catch (error) { + logger.error("Error processing search index jobs", { error }); + Sentry.captureException(error); + } + }, + ); + }, SEARCH_INDEX_INTERVAL); + // Wait for all workers to complete (which should only happen on shutdown) await Promise.all([billingWorkerPromise, precrawlWorkerPromise]); @@ -420,4 +451,5 @@ const DOMAIN_FREQUENCY_INTERVAL = 10000; clearInterval(indexRFInserterInterval); clearInterval(omceInserterInterval); clearInterval(domainFrequencyInterval); + clearInterval(searchIndexInterval); })(); diff --git a/apps/api/src/services/search-index-db.ts b/apps/api/src/services/search-index-db.ts new file mode 100644 index 000000000..6db8a2444 --- /dev/null +++ b/apps/api/src/services/search-index-db.ts @@ -0,0 +1,106 @@ +/** + * Search Index Database Service + * + * Separate Supabase client for the search index database. + * This allows independent scaling and isolation from the main index DB. + */ + +import { createClient, SupabaseClient } from "@supabase/supabase-js"; +import { logger as _logger } from "../lib/logger"; +import { configDotenv } from "dotenv"; + +configDotenv(); + +/** + * SearchIndexSupabaseService - Manages connection to search index database + */ +class SearchIndexSupabaseService { + private client: SupabaseClient | null = null; + + constructor() { + const supabaseUrl = process.env.SEARCH_INDEX_SUPABASE_URL; + const supabaseServiceToken = process.env.SEARCH_INDEX_SUPABASE_SERVICE_TOKEN; + + // Only initialize if both URL and token are provided + if (!supabaseUrl || !supabaseServiceToken) { + _logger.warn("Search index database not configured. Set SEARCH_INDEX_SUPABASE_URL and SEARCH_INDEX_SUPABASE_SERVICE_TOKEN to enable."); + this.client = null; + } else { + this.client = createClient(supabaseUrl, supabaseServiceToken); + _logger.info("Search index database client initialized", { + url: supabaseUrl.substring(0, 30) + "...", + }); + } + } + + /** + * Get the Supabase client for search index + */ + getClient(): SupabaseClient | null { + return this.client; + } + + /** + * Check if search index is enabled + */ + isEnabled(): boolean { + return this.client !== null; + } +} + +const searchIndexService = new SearchIndexSupabaseService(); + +/** + * Proxy to provide clean error messages when search index is not configured + */ +export const search_index_supabase_service: SupabaseClient = new Proxy( + searchIndexService, + { + get: function (target, prop, receiver) { + const client = target.getClient(); + + // If client is not initialized, provide meaningful error + if (client === null) { + return () => { + throw new Error( + "Search index database is not configured. " + + "Set SEARCH_INDEX_SUPABASE_URL and SEARCH_INDEX_SUPABASE_SERVICE_TOKEN environment variables." + ); + }; + } + + // Direct access to service properties + if (prop in target) { + return Reflect.get(target, prop, receiver); + } + + // Delegate to Supabase client + return Reflect.get(client, prop, receiver); + }, + } +) as unknown as SupabaseClient; + +/** + * Check if search index database is enabled + */ +export function isSearchIndexEnabled(): boolean { + return ( + process.env.ENABLE_SEARCH_INDEX === "true" && + searchIndexService.isEnabled() + ); +} + +/** + * Get search index database client (throws if not configured) + */ +export function getSearchIndexClient(): SupabaseClient { + const client = searchIndexService.getClient(); + if (!client) { + throw new Error( + "Search index database is not configured. " + + "Set SEARCH_INDEX_SUPABASE_URL and SEARCH_INDEX_SUPABASE_SERVICE_TOKEN." + ); + } + return client; +} +