From e6562c980bbf4d6319991809e2fac1e8d4121818 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Gerg=C5=91=20M=C3=B3ricz?= Date: Sat, 27 Sep 2025 14:17:29 +0300 Subject: [PATCH] feat(api): otel sidecar --- apps/api/package.json | 22 +++++----- apps/api/src/harness.ts | 8 ++-- apps/api/src/index.ts | 13 +----- apps/api/src/lib/bigquery-jobs.ts | 42 +++++++++---------- apps/api/src/lib/job-transform.ts | 34 ++++++++------- apps/api/src/otel.ts | 14 +++++++ apps/api/src/services/queue-worker.ts | 13 +----- apps/api/src/services/worker/scrape-worker.ts | 14 +------ 8 files changed, 75 insertions(+), 85 deletions(-) create mode 100644 apps/api/src/otel.ts diff --git a/apps/api/package.json b/apps/api/package.json index c29baf861..69245b839 100644 --- a/apps/api/package.json +++ b/apps/api/package.json @@ -6,12 +6,12 @@ "scripts": { "start": "tsc && node dist/src/harness.js --start-built", "dev": "tsx src/harness.ts --start", - "server": "tsc-watch --onSuccess \"node dist/src/index.js\"", - "server:production": "tsc && node dist/src/index.js", - "server:production:nobuild": "node dist/src/index.js", + "server": "tsc-watch --onSuccess \"node --import ./dist/src/otel.js dist/src/index.js\"", + "server:production": "tsc && node --import ./dist/src/otel.js dist/src/index.js", + "server:production:nobuild": "node --import ./dist/src/otel.js dist/src/index.js", "format": "prettier --write \"src/**/*.(js|ts)\"", - "flyio": "node dist/src/index.js", - "start:dev": "tsc-watch --onSuccess \"node dist/src/index.js\"", + "flyio": "node --import ./dist/src/otel.js dist/src/index.js", + "start:dev": "tsc-watch --onSuccess \"node --import ./dist/src/otel.js dist/src/index.js\"", "build": "tsc", "build:nosentry": "tsc", "test": "jest --testPathIgnorePatterns=\"src/__tests__/e2e_noAuth/*\"", @@ -20,12 +20,12 @@ "test:prod": "jest --testPathIgnorePatterns=\"(src/__tests__/e2e_noAuth|src/__tests__/e2e_full_withAuth|src/scraper/scrapeURL)\"", "test:snips": "jest \"src/__tests__/snips/v[12]/.+\\.test\\.ts\"", "harness": "tsx src/harness.ts", - "workers": "tsc-watch --onSuccess \"node dist/src/services/queue-worker.js\"", - "worker:production": "node dist/src/services/queue-worker.js", - "nuq-worker": "tsc-watch --onSuccess \"node dist/src/services/worker/nuq-worker.js\"", - "nuq-worker:production": "node dist/src/services/worker/nuq-worker.js", - "index-worker": "tsc-watch --onSuccess \"node dist/src/services/indexing/index-worker.js\"", - "index-worker:production": "node dist/src/services/indexing/index-worker.js", + "workers": "tsc-watch --onSuccess \"node --import ./dist/src/otel.js dist/src/services/queue-worker.js\"", + "worker:production": "node --import ./dist/src/otel.js dist/src/services/queue-worker.js", + "nuq-worker": "tsc-watch --onSuccess \"node --import ./dist/src/otel.js dist/src/services/worker/nuq-worker.js\"", + "nuq-worker:production": "node --import ./dist/src/otel.js dist/src/services/worker/nuq-worker.js", + "index-worker": "tsc-watch --onSuccess \"node --import ./dist/src/otel.js dist/src/services/indexing/index-worker.js\"", + "index-worker:production": "node --import ./dist/src/otel.js dist/src/services/indexing/index-worker.js", "mongo-docker": "docker run -d -p 2717:27017 -v ./mongo-data:/data/db --name mongodb mongo:latest", "mongo-docker-console": "docker exec -it mongodb mongosh", "run-example": "npx ts-node src/example.ts", diff --git a/apps/api/src/harness.ts b/apps/api/src/harness.ts index a740f444a..77e4d0b3b 100644 --- a/apps/api/src/harness.ts +++ b/apps/api/src/harness.ts @@ -320,7 +320,7 @@ function startServices(command?: string[]): Services { const api = execForward( "api", process.argv[2] === "--start-docker" - ? "node dist/src/index.js" + ? "node --import ./dist/src/otel.js dist/src/index.js" : "pnpm server:production:nobuild", { NUQ_REDUCE_NOISE: "true", @@ -330,7 +330,7 @@ function startServices(command?: string[]): Services { const worker = execForward( "worker", process.argv[2] === "--start-docker" - ? "node dist/src/services/queue-worker.js" + ? "node --import ./dist/src/otel.js dist/src/services/queue-worker.js" : "pnpm worker:production", { NUQ_REDUCE_NOISE: "true", @@ -341,7 +341,7 @@ function startServices(command?: string[]): Services { execForward( `nuq-worker-${i}`, process.argv[2] === "--start-docker" - ? "node dist/src/services/worker/nuq-worker.js" + ? "node --import ./dist/src/otel.js dist/src/services/worker/nuq-worker.js" : "pnpm nuq-worker:production", { NUQ_WORKER_PORT: String(3006 + i), @@ -355,7 +355,7 @@ function startServices(command?: string[]): Services { ? execForward( "index-worker", process.argv[2] === "--start-docker" - ? "node dist/src/services/indexing/index-worker.js" + ? "node --import ./dist/src/otel.js dist/src/services/indexing/index-worker.js" : "pnpm index-worker:production", { NUQ_REDUCE_NOISE: "true", diff --git a/apps/api/src/index.ts b/apps/api/src/index.ts index 598565536..e6b4854da 100644 --- a/apps/api/src/index.ts +++ b/apps/api/src/index.ts @@ -1,4 +1,5 @@ import "dotenv/config"; +import { shutdownOtel } from "./otel"; import "./services/sentry"; import * as Sentry from "@sentry/node"; import express, { NextFunction, Request, Response } from "express"; @@ -30,9 +31,6 @@ import { attachWsProxy } from "./services/agentLivecastWS"; import { cacheableLookup } from "./scraper/scrapeURL/lib/cacheableLookup"; import { v2Router } from "./routes/v2"; import domainFrequencyRouter from "./routes/domain-frequency"; -import { NodeSDK } from "@opentelemetry/sdk-node"; -import { getNodeAutoInstrumentations } from "@opentelemetry/auto-instrumentations-node"; -import { OTLPTraceExporter } from "@opentelemetry/exporter-trace-otlp-http"; import { nuqShutdown } from "./services/worker/nuq"; import { getErrorContactMessage } from "./lib/deployment"; import { initializeBlocklist } from "./scraper/WebScraper/utils/blocklist"; @@ -53,13 +51,6 @@ logger.info("Network info dump", { cacheableLookup.install(http.globalAgent); cacheableLookup.install(https.globalAgent); -const otelSdk = new NodeSDK({ - traceExporter: new OTLPTraceExporter(), - instrumentations: [getNodeAutoInstrumentations()], -}); - -otelSdk.start(); - // Initialize Express with WebSocket support const expressApp = express(); const ws = expressWs(expressApp); @@ -142,7 +133,7 @@ async function startServer(port = DEFAULT_PORT) { logger.info("Server closed."); nuqShutdown().finally(() => { logger.info("NUQ shutdown complete"); - otelSdk.shutdown().finally(() => { + shutdownOtel().finally(() => { logger.info("OTEL shutdown"); process.exit(0); }); diff --git a/apps/api/src/lib/bigquery-jobs.ts b/apps/api/src/lib/bigquery-jobs.ts index 0fb7b7c7e..41517e30b 100644 --- a/apps/api/src/lib/bigquery-jobs.ts +++ b/apps/api/src/lib/bigquery-jobs.ts @@ -345,26 +345,26 @@ export async function saveJobToBigQuery( /** * Queries jobs from BigQuery */ -export async function queryJobsFromBigQuery( - query: string, - params?: any[] -): Promise { - if (!bigquery || !process.env.BIGQUERY_DATASET_ID) { - logger.warn("BigQuery not configured"); - return []; - } +// export async function queryJobsFromBigQuery( +// query: string, +// params?: any[] +// ): Promise { +// if (!bigquery || !process.env.BIGQUERY_DATASET_ID) { +// logger.warn("BigQuery not configured"); +// return []; +// } - try { - const options = { - query, - params, - location: process.env.BIGQUERY_LOCATION || "US", - }; +// try { +// const options = { +// query, +// params, +// location: process.env.BIGQUERY_LOCATION || "US", +// }; - const [rows] = await bigquery.query(options); - return rows; - } catch (error) { - logger.error("Error querying jobs from BigQuery", { error, query }); - throw error; - } -} +// const [rows] = await bigquery.query(options); +// return rows; +// } catch (error) { +// logger.error("Error querying jobs from BigQuery", { error, query }); +// throw error; +// } +// } diff --git a/apps/api/src/lib/job-transform.ts b/apps/api/src/lib/job-transform.ts index 38746ef32..fac5a8f02 100644 --- a/apps/api/src/lib/job-transform.ts +++ b/apps/api/src/lib/job-transform.ts @@ -14,7 +14,7 @@ function cleanOfNull(x: T): T { } } -export interface TransformOptions { +interface TransformOptions { /** Whether to include timestamp field (for BigQuery) */ includeTimestamp?: boolean; /** Whether to serialize objects to JSON strings (for BigQuery) */ @@ -28,7 +28,7 @@ export interface TransformOptions { */ export function transformJobForLogging( job: FirecrawlJob, - options: TransformOptions = {} + options: TransformOptions = {}, ) { const { includeTimestamp = false, @@ -39,16 +39,22 @@ export function transformJobForLogging( const zeroDataRetention = job.zeroDataRetention ?? false; // Determine if docs should be included based on zero data retention and GCS usage - const shouldIncludeDocs = !zeroDataRetention && - !((job.mode === "single_urls" || job.mode === "scrape") && process.env.GCS_BUCKET_NAME); + const shouldIncludeDocs = + !zeroDataRetention && + !( + (job.mode === "single_urls" || job.mode === "scrape") && + process.env.GCS_BUCKET_NAME + ); const baseTransform = { job_id: job.job_id ? job.job_id : null, success: job.success, message: zeroDataRetention ? null : job.message, num_docs: job.num_docs, - docs: shouldIncludeDocs - ? (cleanNullValues ? cleanOfNull(job.docs) : job.docs) + docs: shouldIncludeDocs + ? cleanNullValues + ? cleanOfNull(job.docs) + : job.docs : null, time_taken: job.time_taken, team_id: @@ -56,9 +62,7 @@ export function transformJobForLogging( ? null : job.team_id, mode: job.mode, - url: zeroDataRetention - ? "" - : job.url, + url: zeroDataRetention ? "" : job.url, crawler_options: zeroDataRetention ? null : job.crawlerOptions, page_options: zeroDataRetention ? null : job.scrapeOptions, origin: zeroDataRetention ? null : job.origin, @@ -90,14 +94,14 @@ export function transformJobForLogging( return { ...withTimestamp, docs: withTimestamp.docs ? JSON.stringify(withTimestamp.docs) : null, - crawler_options: withTimestamp.crawler_options - ? JSON.stringify(withTimestamp.crawler_options) + crawler_options: withTimestamp.crawler_options + ? JSON.stringify(withTimestamp.crawler_options) : null, - page_options: withTimestamp.page_options - ? JSON.stringify(withTimestamp.page_options) + page_options: withTimestamp.page_options + ? JSON.stringify(withTimestamp.page_options) : null, - cost_tracking: withTimestamp.cost_tracking - ? JSON.stringify(withTimestamp.cost_tracking) + cost_tracking: withTimestamp.cost_tracking + ? JSON.stringify(withTimestamp.cost_tracking) : null, }; } diff --git a/apps/api/src/otel.ts b/apps/api/src/otel.ts new file mode 100644 index 000000000..ed3e608af --- /dev/null +++ b/apps/api/src/otel.ts @@ -0,0 +1,14 @@ +import { getNodeAutoInstrumentations } from "@opentelemetry/auto-instrumentations-node"; +import { NodeSDK } from "@opentelemetry/sdk-node"; +import { OTLPTraceExporter } from "@opentelemetry/exporter-trace-otlp-http"; + +const otelSdk = new NodeSDK({ + traceExporter: new OTLPTraceExporter(), + instrumentations: [getNodeAutoInstrumentations()], +}); + +otelSdk.start(); + +export const shutdownOtel = async () => { + await otelSdk.shutdown(); +}; diff --git a/apps/api/src/services/queue-worker.ts b/apps/api/src/services/queue-worker.ts index e128d7ed1..7550f419b 100644 --- a/apps/api/src/services/queue-worker.ts +++ b/apps/api/src/services/queue-worker.ts @@ -1,4 +1,5 @@ import "dotenv/config"; +import { shutdownOtel } from "../otel"; import "./sentry"; import * as Sentry from "@sentry/node"; import { @@ -28,9 +29,6 @@ import http from "http"; import https from "https"; import { cacheableLookup } from "../scraper/scrapeURL/lib/cacheableLookup"; import { robustFetch } from "../scraper/scrapeURL/lib/fetch"; -import { NodeSDK } from "@opentelemetry/sdk-node"; -import { getNodeAutoInstrumentations } from "@opentelemetry/auto-instrumentations-node"; -import { OTLPTraceExporter } from "@opentelemetry/exporter-trace-otlp-http"; import { BullMQOtel } from "bullmq-otel"; import { getErrorContactMessage } from "../lib/deployment"; import { initializeBlocklist } from "../scraper/WebScraper/utils/blocklist"; @@ -56,13 +54,6 @@ const runningJobs: Set = new Set(); cacheableLookup.install(http.globalAgent); cacheableLookup.install(https.globalAgent); -const otelSdk = new NodeSDK({ - traceExporter: new OTLPTraceExporter(), - instrumentations: [getNodeAutoInstrumentations()], -}); - -otelSdk.start(); - const processExtractJobInternal = async ( token: string, job: Job & { id: string }, @@ -488,6 +479,6 @@ app.listen(workerPort, () => { } console.log("All jobs finished. Worker out!"); - await otelSdk.shutdown(); + await shutdownOtel(); process.exit(0); })(); diff --git a/apps/api/src/services/worker/scrape-worker.ts b/apps/api/src/services/worker/scrape-worker.ts index 536a760c1..17bcd4bea 100644 --- a/apps/api/src/services/worker/scrape-worker.ts +++ b/apps/api/src/services/worker/scrape-worker.ts @@ -50,8 +50,6 @@ import { calculateCreditsToBeBilled } from "../../lib/scrape-billing"; import { getBillingQueue } from "../queue-service"; import type { Logger } from "winston"; import { finishCrawlIfNeeded } from "./crawl-logic"; -import { getNodeAutoInstrumentations } from "@opentelemetry/auto-instrumentations-node"; -import { NodeSDK } from "@opentelemetry/sdk-node"; import { RacedRedirectError, ScrapeJobTimeoutError, @@ -59,7 +57,6 @@ import { UnknownError, } from "../../lib/error"; import { serializeTransportableError } from "../../lib/error-serde"; -import { OTLPTraceExporter } from "@opentelemetry/exporter-trace-otlp-http"; import type { NuQJob } from "./nuq"; import { ScrapeJobData, @@ -68,7 +65,7 @@ import { ScrapeJobSingleUrls, } from "../../types"; import { scrapeSitemap } from "../../scraper/crawler/sitemap"; -import { filterUrl } from "@mendable/firecrawl-rs"; +import { shutdownOtel } from "../../otel"; const sleep = (ms: number) => new Promise(resolve => setTimeout(resolve, ms)); @@ -1140,15 +1137,8 @@ export const processJobInternal = async (job: NuQJob) => { } }; -const otelSdk = new NodeSDK({ - traceExporter: new OTLPTraceExporter(), - instrumentations: [getNodeAutoInstrumentations()], -}); - -otelSdk.start(); - const exitHandler = () => { - otelSdk.shutdown().finally(() => { + shutdownOtel().finally(() => { _logger.debug("OTEL shutdown"); process.exit(0); });