From f1cbba477bb6b1517af848b8236494f027bb077f Mon Sep 17 00:00:00 2001 From: Nicolas Date: Sun, 12 Oct 2025 02:16:31 -0300 Subject: [PATCH] Nick: --- apps/api/src/controllers/v2/f-search.ts | 50 +-- apps/api/src/lib/search-index-client.ts | 363 ++++++++++++++++++ .../transformers/sendToSearchIndex.ts | 25 +- .../api/src/services/indexing/index-worker.ts | 53 ++- 4 files changed, 413 insertions(+), 78 deletions(-) create mode 100644 apps/api/src/lib/search-index-client.ts diff --git a/apps/api/src/controllers/v2/f-search.ts b/apps/api/src/controllers/v2/f-search.ts index b791c02a0..063825d1e 100644 --- a/apps/api/src/controllers/v2/f-search.ts +++ b/apps/api/src/controllers/v2/f-search.ts @@ -1,18 +1,12 @@ import { Request, Response } from "express"; import { z } from "zod"; -import { - search_index_supabase_service, - isSearchIndexEnabled, -} from "../../services/search-index-db"; -import { - search, - searchChunks, - getSearchStats, - type SearchQuery, - type SearchFilters, -} from "../../lib/search-index"; import { logger as _logger } from "../../lib/logger"; import { CostTracking } from "../../lib/cost-tracking"; +import { + getSearchIndexClient, + SearchIndexClient, + type SearchRequest, +} from "../../lib/search-index-client"; // Validation schemas const searchRequestSchema = z.object({ @@ -75,29 +69,27 @@ export async function realtimeSearchController( filters, }); - // Check if search index is enabled - if (!isSearchIndexEnabled()) { + // Get search index client + const client = getSearchIndexClient(); + + if (!client) { res.status(503).json({ success: false, - error: "Search index is not configured", + error: "Search index service is not configured", }); return; } - // Perform search - const searchQuery: SearchQuery = { + // Perform search via HTTP client + const searchRequest: SearchRequest = { query, limit, offset, mode, - filters: filters as SearchFilters, + filters, }; - const result = await search( - search_index_supabase_service, - searchQuery, - logger, - ); + const result = await client.search(searchRequest, logger); // Track cost (if applicable) const costTracking = new CostTracking(); @@ -105,19 +97,7 @@ export async function realtimeSearchController( res.status(200).json({ success: true, - data: { - results: result.results, - total: result.total, - query: result.query, - mode: result.mode, - took: result.took, - pagination: { - limit, - offset, - // hasMore is true if we got a full page of results, suggesting there may be more - hasMore: result.results.length >= limit, - }, - }, + data: result, costTracking: costTracking.toJSON(), }); } catch (error) { diff --git a/apps/api/src/lib/search-index-client.ts b/apps/api/src/lib/search-index-client.ts new file mode 100644 index 000000000..2fd562e5d --- /dev/null +++ b/apps/api/src/lib/search-index-client.ts @@ -0,0 +1,363 @@ +/** + * HTTP Client for Search Index Service + * + * This client communicates with the standalone search index service + * (firecrawl-search backend) via HTTP. + */ + +import { logger as _logger } from "./logger"; +import type { Logger } from "winston"; + +export interface SearchIndexClientConfig { + baseUrl: string; + apiSecret?: string; + timeout?: number; +} + +export interface IndexDocumentRequest { + 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 SearchRequest { + query: string; + limit?: number; + offset?: number; + mode?: "hybrid" | "keyword" | "semantic"; + filters?: { + 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; + pagination: { + limit: number; + offset: number; + hasMore: boolean; + }; +} + +export class SearchIndexClient { + private baseUrl: string; + private apiSecret?: string; + private timeout: number; + + constructor(config: SearchIndexClientConfig) { + this.baseUrl = config.baseUrl.replace(/\/$/, ""); // Remove trailing slash + this.apiSecret = config.apiSecret; + this.timeout = config.timeout || 30000; + } + + /** + * Check if search index service is enabled + */ + static isEnabled(): boolean { + return !!( + process.env.SEARCH_SERVICE_URL && + process.env.ENABLE_SEARCH_INDEX === "true" + ); + } + + /** + * Make HTTP request to search service + */ + private async request( + method: string, + path: string, + body?: any, + logger?: Logger, + ): Promise { + const log = logger ?? _logger.child({ module: "search-index-client" }); + const url = `${this.baseUrl}${path}`; + + const headers: Record = { + "Content-Type": "application/json", + }; + + if (this.apiSecret) { + headers["X-API-Secret"] = this.apiSecret; + } + + const controller = new AbortController(); + const timeoutId = setTimeout(() => controller.abort(), this.timeout); + + try { + log.debug("Making request to search service", { + method, + url, + hasBody: !!body, + }); + + const response = await fetch(url, { + method, + headers, + body: body ? JSON.stringify(body) : undefined, + signal: controller.signal, + }); + + clearTimeout(timeoutId); + + const data = await response.json(); + + if (!response.ok) { + log.error("Search service request failed", { + status: response.status, + error: data.error || "Unknown error", + }); + throw new Error( + data.error || `Search service returned ${response.status}`, + ); + } + + return data; + } catch (error) { + clearTimeout(timeoutId); + + if (error.name === "AbortError") { + log.error("Search service request timed out", { url, timeout: this.timeout }); + throw new Error("Search service request timed out"); + } + + log.error("Search service request failed", { + error: (error as Error).message, + url, + }); + throw error; + } + } + + /** + * Index a document (async - queues for processing) + */ + async indexDocument( + request: IndexDocumentRequest, + logger?: Logger, + ): Promise<{ success: boolean; message: string }> { + const log = logger ?? _logger.child({ module: "search-index-client" }); + + try { + const response = await this.request( + "POST", + "/api/index", + request, + log, + ); + + log.info("Document queued for indexing", { + url: request.url, + success: response.success, + }); + + return { + success: response.success, + message: response.message || "Document queued for indexing", + }; + } catch (error) { + log.error("Failed to queue document for indexing", { + error: (error as Error).message, + url: request.url, + }); + + // Don't throw - indexing failures shouldn't break scraping + return { + success: false, + message: (error as Error).message, + }; + } + } + + /** + * Search indexed documents + */ + async search( + request: SearchRequest, + logger?: Logger, + ): Promise { + const log = logger ?? _logger.child({ module: "search-index-client" }); + + try { + const response = await this.request<{ success: boolean; data: SearchResponse }>( + "POST", + "/api/search", + request, + log, + ); + + if (!response.success || !response.data) { + throw new Error("Invalid response from search service"); + } + + log.info("Search completed", { + query: request.query, + results: response.data.results.length, + took: response.data.took, + }); + + return response.data; + } catch (error) { + log.error("Search request failed", { + error: (error as Error).message, + query: request.query, + }); + + // Return empty results on error + return { + results: [], + total: 0, + query: request.query, + mode: request.mode || "hybrid", + took: 0, + pagination: { + limit: request.limit || 50, + offset: request.offset || 0, + hasMore: false, + }, + }; + } + } + + /** + * Get search index statistics + */ + async getStats(logger?: Logger): Promise<{ + index: any; + queue: any; + }> { + const log = logger ?? _logger.child({ module: "search-index-client" }); + + try { + const response = await this.request<{ success: boolean; data: any }>( + "GET", + "/api/stats", + undefined, + log, + ); + + if (!response.success || !response.data) { + throw new Error("Invalid response from search service"); + } + + return response.data; + } catch (error) { + log.error("Failed to get search stats", { + error: (error as Error).message, + }); + + return { + index: {}, + queue: {}, + }; + } + } + + /** + * Health check + */ + async health(logger?: Logger): Promise { + const log = logger ?? _logger.child({ module: "search-index-client" }); + + try { + const response = await this.request<{ success: boolean }>( + "GET", + "/health", + undefined, + log, + ); + + return response.success; + } catch (error) { + log.warn("Search service health check failed", { + error: (error as Error).message, + }); + return false; + } + } +} + +// Singleton instance +let searchIndexClient: SearchIndexClient | null = null; + +/** + * Get search index client instance + */ +export function getSearchIndexClient(): SearchIndexClient | null { + if (!SearchIndexClient.isEnabled()) { + return null; + } + + if (!searchIndexClient) { + const baseUrl = process.env.SEARCH_SERVICE_URL!; + const apiSecret = process.env.SEARCH_SERVICE_API_SECRET; + + searchIndexClient = new SearchIndexClient({ + baseUrl, + apiSecret, + timeout: 30000, + }); + + _logger.info("Search index client initialized", { + baseUrl: baseUrl.substring(0, 30) + "...", + }); + } + + return searchIndexClient; +} + +/** + * Helper: Index document if search service is enabled + */ +export async function indexDocumentIfEnabled( + request: IndexDocumentRequest, + logger?: Logger, +): Promise { + const client = getSearchIndexClient(); + + if (!client) { + return; + } + + try { + await client.indexDocument(request, logger); + } catch (error) { + // Silently fail - indexing is optional + (logger ?? _logger).warn("Failed to index document", { + error: (error as Error).message, + url: request.url, + }); + } +} + diff --git a/apps/api/src/scraper/scrapeURL/transformers/sendToSearchIndex.ts b/apps/api/src/scraper/scrapeURL/transformers/sendToSearchIndex.ts index e71c7e868..07ac02613 100644 --- a/apps/api/src/scraper/scrapeURL/transformers/sendToSearchIndex.ts +++ b/apps/api/src/scraper/scrapeURL/transformers/sendToSearchIndex.ts @@ -10,7 +10,7 @@ */ import { Document } from "../../../controllers/v1/types"; -import { addSearchIndexJob } from "../../../lib/search-index/queue"; +import { indexDocumentIfEnabled } from "../../../lib/search-index-client"; import { logger as _logger } from "../../../lib/logger"; import { Meta } from ".."; @@ -80,15 +80,16 @@ export async function sendDocumentToSearchIndex( meta: Meta, document: Document, ): Promise { - // Check if search indexing is enabled + // Check if search indexing is enabled via the SEARCH_SERVICE_URL const searchIndexEnabled = process.env.ENABLE_SEARCH_INDEX === "true" && - process.env.SEARCH_INDEX_SUPABASE_URL && - process.env.SEARCH_INDEX_SUPABASE_SERVICE_TOKEN; + process.env.SEARCH_SERVICE_URL; + + meta.logger.debug("Sending document to search index", { + url: meta.url, + searchIndexEnabled, + }); - meta.logger.debug("Sending document to search index", { - url: meta.url, - }); if (!searchIndexEnabled) { return document; } @@ -119,10 +120,10 @@ export async function sendDocumentToSearchIndex( // Remove indexId from metadata after extracting it (internal field, shouldn't be exposed to user) delete document.metadata.indexId; - // Queue job for async processing (don't block scraper) + // Send to search service via HTTP (async, don't block scraper) (async () => { try { - await addSearchIndexJob({ + await indexDocumentIfEnabled({ url: meta.url, resolvedUrl: document.metadata.url ?? @@ -143,13 +144,13 @@ export async function sendDocumentToSearchIndex( language: document.metadata.language ?? "en", country: meta.options.location?.country ?? undefined, isMobile: meta.options.mobile ?? false, - }); + }, meta.logger); - meta.logger.debug("Queued document for search indexing", { + meta.logger.debug("Sent document to search service", { url: meta.url, }); } catch (error) { - meta.logger.error("Failed to queue document for search indexing", { + meta.logger.error("Failed to send document to search service", { error: (error as Error).message, url: meta.url, }); diff --git a/apps/api/src/services/indexing/index-worker.ts b/apps/api/src/services/indexing/index-worker.ts index af0b7615a..d0257eb72 100644 --- a/apps/api/src/services/indexing/index-worker.ts +++ b/apps/api/src/services/indexing/index-worker.ts @@ -23,7 +23,9 @@ import { processOMCEJobs, processDomainFrequencyJobs, } from ".."; -import { processSearchIndexJobs } from "../../lib/search-index/queue"; +import { getSearchIndexClient } from "../../lib/search-index-client"; +// Search indexing is now handled by the separate search service +// import { processSearchIndexJobs } from "../../lib/search-index/queue"; import { processWebhookInsertJobs } from "../webhook"; import { scrapeOptions as scrapeOptionsSchema, @@ -332,7 +334,8 @@ 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 +// Search indexing is now handled by separate search service, not this worker +// const SEARCH_INDEX_INTERVAL = 10000; // Start the workers (async () => { @@ -414,34 +417,23 @@ const SEARCH_INDEX_INTERVAL = 10000; // Process search index queue every 10 seco ); }, 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); + // Search indexing is now handled by separate search service + // The search service has its own worker that processes the queue + // This worker no longer needs to process search index jobs + + // Health check for search service (optional) + const searchClient = getSearchIndexClient(); + if (searchClient) { + searchClient.health().then(healthy => { + if (healthy) { + logger.info("Search service is healthy"); + } else { + logger.warn("Search service health check failed"); + } + }).catch(error => { + logger.error("Search service health check error", { error }); + }); + } // Wait for all workers to complete (which should only happen on shutdown) await Promise.all([billingWorkerPromise, precrawlWorkerPromise]); @@ -451,5 +443,4 @@ const SEARCH_INDEX_INTERVAL = 10000; // Process search index queue every 10 seco clearInterval(indexRFInserterInterval); clearInterval(omceInserterInterval); clearInterval(domainFrequencyInterval); - clearInterval(searchIndexInterval); })();