From 81e98a2e9db2828b276c0bf2535bbf37ff0cbfe5 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Gerg=C5=91=20M=C3=B3ricz?= Date: Thu, 2 Oct 2025 20:46:40 +0200 Subject: [PATCH] feat(api/extract): timing data --- apps/api/src/controllers/v1/extract.ts | 5 +- apps/api/src/controllers/v2/extract.ts | 5 +- .../api/src/lib/extract/extraction-service.ts | 14 +- .../extract/fire-0/extraction-service-f0.ts | 14 +- apps/api/src/services/queue-worker.ts | 143 ------------------ 5 files changed, 23 insertions(+), 158 deletions(-) diff --git a/apps/api/src/controllers/v1/extract.ts b/apps/api/src/controllers/v1/extract.ts index 0400763a8..09218c158 100644 --- a/apps/api/src/controllers/v1/extract.ts +++ b/apps/api/src/controllers/v1/extract.ts @@ -110,6 +110,8 @@ export async function extractController( isUrlBlocked(url, req.acuc?.flags ?? null), ) ?? []; + const createdAt = Date.now(); + if (invalidURLs.length > 0 && !req.body.ignoreInvalidURLs) { if (!res.headersSent) { return res.status(403).json({ @@ -149,6 +151,7 @@ export async function extractController( extractId, agent: req.body.agent, apiKeyId: req.acuc?.api_key_id ?? null, + createdAt, }; if ( @@ -164,7 +167,7 @@ export async function extractController( await saveExtract(extractId, { id: extractId, team_id: req.auth.team_id, - createdAt: Date.now(), + createdAt, status: "processing", showSteps: req.body.__experimental_streamSteps, showLLMUsage: req.body.__experimental_llmUsage, diff --git a/apps/api/src/controllers/v2/extract.ts b/apps/api/src/controllers/v2/extract.ts index 6a614e78f..08dce3a62 100644 --- a/apps/api/src/controllers/v2/extract.ts +++ b/apps/api/src/controllers/v2/extract.ts @@ -48,7 +48,7 @@ export async function extractController( } const extractId = crypto.randomUUID(); - + const createdAt = Date.now(); _logger.info("Extract starting...", { request: req.body, originalRequest, @@ -65,12 +65,13 @@ export async function extractController( subId: req.acuc?.sub_id, extractId, agent: req.body.agent, + createdAt, }; await saveExtract(extractId, { id: extractId, team_id: req.auth.team_id, - createdAt: Date.now(), + createdAt, status: "processing", showSteps: req.body.__experimental_streamSteps, showLLMUsage: req.body.__experimental_llmUsage, diff --git a/apps/api/src/lib/extract/extraction-service.ts b/apps/api/src/lib/extract/extraction-service.ts index 4e70fc516..c4952a160 100644 --- a/apps/api/src/lib/extract/extraction-service.ts +++ b/apps/api/src/lib/extract/extraction-service.ts @@ -50,6 +50,7 @@ interface ExtractServiceOptions { cacheKey?: string; agent?: boolean; apiKeyId: number | null; + createdAt?: number; } export interface ExtractResult { @@ -78,6 +79,7 @@ export async function performExtraction( options: ExtractServiceOptions, ): Promise { const { request, teamId, subId, apiKeyId } = options; + const createdAt = options.createdAt ? new Date(options.createdAt) : new Date(); const urlTraces: URLTrace[] = []; let docsMap: Map = new Map(); let singleAnswerCompletions: completions | null = null; @@ -134,7 +136,7 @@ export async function performExtraction( message: "No search results found", num_docs: 1, docs: [], - time_taken: (new Date().getTime() - Date.now()) / 1000, + time_taken: (new Date().getTime() - createdAt.getTime()) / 1000, team_id: teamId, mode: "extract", url: request.urls?.join(", ") || "", @@ -680,7 +682,7 @@ export async function performExtraction( : "Failed to transform array to object", num_docs: 1, docs: [], - time_taken: (new Date().getTime() - Date.now()) / 1000, + time_taken: (new Date().getTime() - createdAt.getTime()) / 1000, team_id: teamId, mode: "extract", url: request.urls?.join(", ") || "", @@ -797,7 +799,7 @@ export async function performExtraction( message: error.message, num_docs: 1, docs: [], - time_taken: (new Date().getTime() - Date.now()) / 1000, + time_taken: (new Date().getTime() - createdAt.getTime()) / 1000, team_id: teamId, mode: "extract", url: request.urls?.join(", ") || "", @@ -855,7 +857,7 @@ export async function performExtraction( message: errorMessage, num_docs: 1, docs: [], - time_taken: (new Date().getTime() - Date.now()) / 1000, + time_taken: (new Date().getTime() - createdAt.getTime()) / 1000, team_id: teamId, mode: "extract", url: request.urls?.join(", ") || "", @@ -1049,7 +1051,7 @@ export async function performExtraction( message: "Extract completed", num_docs: 1, docs: finalResult ?? {}, - time_taken: (new Date().getTime() - Date.now()) / 1000, + time_taken: (new Date().getTime() - createdAt.getTime()) / 1000, team_id: teamId, mode: "extract", url: request.urls?.join(", ") || "", @@ -1127,7 +1129,7 @@ export async function performExtraction( : "An unexpected error occurred", num_docs: 1, docs: [], - time_taken: (new Date().getTime() - Date.now()) / 1000, + time_taken: (new Date().getTime() - createdAt.getTime()) / 1000, team_id: teamId, mode: "extract", url: request.urls?.join(", ") || "", diff --git a/apps/api/src/lib/extract/fire-0/extraction-service-f0.ts b/apps/api/src/lib/extract/fire-0/extraction-service-f0.ts index 951c0a536..edaf96ca9 100644 --- a/apps/api/src/lib/extract/fire-0/extraction-service-f0.ts +++ b/apps/api/src/lib/extract/fire-0/extraction-service-f0.ts @@ -47,6 +47,7 @@ interface ExtractServiceOptions { cacheMode?: "load" | "save" | "direct"; cacheKey?: string; apiKeyId: number | null; + createdAt?: number; } interface ExtractResult { @@ -75,6 +76,7 @@ export async function performExtraction_F0( options: ExtractServiceOptions, ): Promise { const { request, teamId, subId, apiKeyId } = options; + const createdAt = options.createdAt ? new Date(options.createdAt) : new Date(); const urlTraces: URLTrace[] = []; let docsMap: Map = new Map(); let singleAnswerCompletions: completions | null = null; @@ -120,7 +122,7 @@ export async function performExtraction_F0( message: "No search results found", num_docs: 1, docs: [], - time_taken: (new Date().getTime() - Date.now()) / 1000, + time_taken: (new Date().getTime() - createdAt.getTime()) / 1000, team_id: teamId, mode: "extract", url: request.urls?.join(", ") || "", @@ -228,7 +230,7 @@ export async function performExtraction_F0( message: "No valid URLs found to scrape", num_docs: 1, docs: [], - time_taken: (new Date().getTime() - Date.now()) / 1000, + time_taken: (new Date().getTime() - createdAt.getTime()) / 1000, team_id: teamId, mode: "extract", url: request.urls?.join(", ") || "", @@ -604,7 +606,7 @@ export async function performExtraction_F0( message: "Failed to transform array to object", num_docs: 1, docs: [], - time_taken: (new Date().getTime() - Date.now()) / 1000, + time_taken: (new Date().getTime() - createdAt.getTime()) / 1000, team_id: teamId, mode: "extract", url: request.urls?.join(", ") || "", @@ -702,7 +704,7 @@ export async function performExtraction_F0( message: "Failed to scrape documents", num_docs: 1, docs: [], - time_taken: (new Date().getTime() - Date.now()) / 1000, + time_taken: (new Date().getTime() - createdAt.getTime()) / 1000, team_id: teamId, mode: "extract", url: request.urls?.join(", ") || "", @@ -732,7 +734,7 @@ export async function performExtraction_F0( message: "All provided URLs are invalid", num_docs: 1, docs: [], - time_taken: (new Date().getTime() - Date.now()) / 1000, + time_taken: (new Date().getTime() - createdAt.getTime()) / 1000, team_id: teamId, mode: "extract", url: request.urls?.join(", ") || "", @@ -904,7 +906,7 @@ export async function performExtraction_F0( message: "Extract completed", num_docs: 1, docs: finalResult ?? {}, - time_taken: (new Date().getTime() - Date.now()) / 1000, + time_taken: (new Date().getTime() - createdAt.getTime()) / 1000, team_id: teamId, mode: "extract", url: request.urls?.join(", ") || "", diff --git a/apps/api/src/services/queue-worker.ts b/apps/api/src/services/queue-worker.ts index b409e5554..4a1ba3b69 100644 --- a/apps/api/src/services/queue-worker.ts +++ b/apps/api/src/services/queue-worker.ts @@ -3,7 +3,6 @@ import { shutdownOtel } from "../otel"; import "./sentry"; import * as Sentry from "@sentry/node"; import { - getExtractQueue, getDeepResearchQueue, getGenerateLlmsTxtQueue, getRedisConnection, @@ -13,24 +12,13 @@ import { logger as _logger } from "../lib/logger"; import systemMonitor from "./system-monitor"; import { v4 as uuidv4 } from "uuid"; import { configDotenv } from "dotenv"; -import { - ExtractResult, - performExtraction, -} from "../lib/extract/extraction-service"; -import { updateExtract } from "../lib/extract/extract-redis"; import { updateDeepResearch } from "../lib/deep-research/deep-research-redis"; import { performDeepResearch } from "../lib/deep-research/deep-research-service"; import { performGenerateLlmsTxt } from "../lib/generate-llmstxt/generate-llmstxt-service"; import { updateGeneratedLlmsTxt } from "../lib/generate-llmstxt/generate-llmstxt-redis"; -import { performExtraction_F0 } from "../lib/extract/fire-0/extraction-service-f0"; -import { createWebhookSender, WebhookEvent } from "./webhook"; import Express from "express"; -import http from "http"; -import https from "https"; -import { cacheableLookup } from "../scraper/scrapeURL/lib/cacheableLookup"; import { robustFetch } from "../scraper/scrapeURL/lib/fetch"; import { BullMQOtel } from "bullmq-otel"; -import { getErrorContactMessage } from "../lib/deployment"; import { initializeBlocklist } from "../scraper/WebScraper/utils/blocklist"; configDotenv(); @@ -50,137 +38,6 @@ const gotJobInterval = Number(process.env.CONNECTION_MONITOR_INTERVAL) || 20; const runningJobs: Set = new Set(); -// Install cacheable lookup for all other requests -// cacheableLookup.install(http.globalAgent); -// cacheableLookup.install(https.globalAgent); - -const processExtractJobInternal = async ( - token: string, - job: Job & { id: string }, -) => { - const logger = _logger.child({ - module: "extract-worker", - method: "processJobInternal", - jobId: job.id, - extractId: job.data.extractId, - teamId: job.data?.teamId ?? undefined, - }); - - const extendLockInterval = setInterval(async () => { - logger.info(`🔄 Worker extending lock on job ${job.id}`); - await job.extendLock(token, jobLockExtensionTime); - }, jobLockExtendInterval); - - const sender = await createWebhookSender({ - teamId: job.data.teamId, - jobId: job.data.extractId, - webhook: job.data.request.webhook, - v0: false, - }); - - try { - if (sender) { - sender.send(WebhookEvent.EXTRACT_STARTED, { - success: true, - }); - } - - let result: ExtractResult | null = null; - - const model = job.data.request.agent?.model; - if ( - job.data.request.agent && - model && - model.toLowerCase().includes("fire-1") - ) { - result = await performExtraction(job.data.extractId, { - request: job.data.request, - teamId: job.data.teamId, - subId: job.data.subId, - apiKeyId: job.data.apiKeyId, - }); - } else { - result = await performExtraction_F0(job.data.extractId, { - request: job.data.request, - teamId: job.data.teamId, - subId: job.data.subId, - apiKeyId: job.data.apiKeyId, - }); - } - // result = await performExtraction_F0(job.data.extractId, { - // request: job.data.request, - // teamId: job.data.teamId, - // subId: job.data.subId, - // }); - - if (result && result.success) { - // Move job to completed state in Redis - await job.moveToCompleted(result, token, false); - - if (sender) { - sender.send(WebhookEvent.EXTRACT_COMPLETED, { - success: true, - data: [result], - }); - } - - return result; - } else { - // throw new Error(result.error || "Unknown error during extraction"); - - await job.moveToCompleted(result, token, false); - await updateExtract(job.data.extractId, { - error: result?.error ?? getErrorContactMessage(job.data.extractId), - }); - - if (sender) { - sender.send(WebhookEvent.EXTRACT_FAILED, { - success: false, - error: result?.error ?? getErrorContactMessage(job.data.extractId), - }); - } - - return result; - } - } catch (error) { - logger.error(`🚫 Job errored ${job.id} - ${error}`, { error }); - - Sentry.captureException(error, { - data: { - job: job.id, - }, - }); - - try { - // Move job to failed state in Redis - await job.moveToFailed(error, token, false); - } catch (e) { - logger.log("Failed to move job to failed state in Redis", { error }); - } - - await updateExtract(job.data.extractId, { - status: "failed", - error: error.error ?? error ?? getErrorContactMessage(job.data.extractId), - }); - - if (sender) { - sender.send(WebhookEvent.EXTRACT_FAILED, { - success: false, - error: - (error as any)?.message ?? getErrorContactMessage(job.data.extractId), - }); - } - - return { - success: false, - error: error.error ?? error ?? getErrorContactMessage(job.data.extractId), - }; - // throw error; - } finally { - clearInterval(extendLockInterval); - } -}; - const processDeepResearchJobInternal = async ( token: string, job: Job & { id: string },