From ce6e164a535a28bf46705973bc4d64e7bb1f0ef2 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Gerg=C5=91=20M=C3=B3ricz?= Date: Thu, 14 Aug 2025 01:21:39 +0200 Subject: [PATCH 1/6] fix(concurrency-limit): give getNextConcurrentJob some room to breathe --- apps/api/src/lib/concurrency-limit.ts | 13 ++++++++++++- 1 file changed, 12 insertions(+), 1 deletion(-) diff --git a/apps/api/src/lib/concurrency-limit.ts b/apps/api/src/lib/concurrency-limit.ts index 7fbb7837b..f245a5413 100644 --- a/apps/api/src/lib/concurrency-limit.ts +++ b/apps/api/src/lib/concurrency-limit.ts @@ -216,8 +216,19 @@ async function getNextConcurrentJob(teamId: string, i = 0): Promise<{ zeroDataRetention: finalJob.job.data?.zeroDataRetention, i }); + } else if (i > 100) { + logger.error("Failed to remove job from concurrency limit queue, hard bailing", { + teamId, + jobId: finalJob.job.id, + zeroDataRetention: finalJob.job.data?.zeroDataRetention, + i + }); + return null; } - return await getNextConcurrentJob(teamId, i + 1); + + return await new Promise((resolve, reject) => setTimeout(() => { + getNextConcurrentJob(teamId, i + 1).then(resolve).catch(reject); + }, 10)); } } From 2e2967cede34bce82e223782d22c44e5211b7c51 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Gerg=C5=91=20M=C3=B3ricz?= Date: Thu, 14 Aug 2025 01:33:55 +0200 Subject: [PATCH 2/6] fix(concurency-limit): bump timeout and log if removed --- apps/api/src/lib/concurrency-limit.ts | 9 ++++++++- 1 file changed, 8 insertions(+), 1 deletion(-) diff --git a/apps/api/src/lib/concurrency-limit.ts b/apps/api/src/lib/concurrency-limit.ts index f245a5413..132c7356c 100644 --- a/apps/api/src/lib/concurrency-limit.ts +++ b/apps/api/src/lib/concurrency-limit.ts @@ -228,7 +228,14 @@ async function getNextConcurrentJob(teamId: string, i = 0): Promise<{ return await new Promise((resolve, reject) => setTimeout(() => { getNextConcurrentJob(teamId, i + 1).then(resolve).catch(reject); - }, 10)); + }, 250)); + } else { + logger.debug("Removed job from concurrency limit queue", { + teamId, + jobId: finalJob.job.id, + zeroDataRetention: finalJob.job.data?.zeroDataRetention, + i + }); } } From a9c3e7ac35c68831d30a0fcb57320a85befe7d8d Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Gerg=C5=91=20M=C3=B3ricz?= Date: Thu, 14 Aug 2025 01:48:13 +0200 Subject: [PATCH 3/6] fix(concurrency-limit): add staggering to the setTimeout --- apps/api/src/lib/concurrency-limit.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/apps/api/src/lib/concurrency-limit.ts b/apps/api/src/lib/concurrency-limit.ts index 132c7356c..52fc2fbdb 100644 --- a/apps/api/src/lib/concurrency-limit.ts +++ b/apps/api/src/lib/concurrency-limit.ts @@ -228,7 +228,7 @@ async function getNextConcurrentJob(teamId: string, i = 0): Promise<{ return await new Promise((resolve, reject) => setTimeout(() => { getNextConcurrentJob(teamId, i + 1).then(resolve).catch(reject); - }, 250)); + }, Math.floor(Math.random() * 300))); // Stagger the workers off to break up the clump that causes the race condition } else { logger.debug("Removed job from concurrency limit queue", { teamId, From 308e4f4368032ba8ec29740404125a6928ab628c Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Gerg=C5=91=20M=C3=B3ricz?= Date: Thu, 14 Aug 2025 21:48:00 +0200 Subject: [PATCH 4/6] fix(queue-service): reduce redis connections to BullMQ (#1966) --- .../unit/deep-research-redis.test.ts | 26 +++++++------- .../queue-concurrency-integration.test.ts | 32 ++++++++--------- .../src/lib/__tests__/job-priority.test.ts | 34 +++++++++--------- .../api/src/services/billing/batch_billing.ts | 14 ++++---- .../api/src/services/indexing/index-worker.ts | 4 +-- apps/api/src/services/queue-service.ts | 35 ++++++++++--------- apps/api/src/services/queue-worker.ts | 9 +++-- 7 files changed, 77 insertions(+), 77 deletions(-) diff --git a/apps/api/src/__tests__/deep-research/unit/deep-research-redis.test.ts b/apps/api/src/__tests__/deep-research/unit/deep-research-redis.test.ts index 2f2ff7359..0ffee0d44 100644 --- a/apps/api/src/__tests__/deep-research/unit/deep-research-redis.test.ts +++ b/apps/api/src/__tests__/deep-research/unit/deep-research-redis.test.ts @@ -1,4 +1,4 @@ -import { redisConnection } from "../../../services/queue-service"; +import { redisEvictConnection } from "../../../services/redis"; import { saveDeepResearch, getDeepResearch, @@ -40,11 +40,11 @@ describe("Deep Research Redis Operations", () => { it("should save research data to Redis with TTL", async () => { await saveDeepResearch("test-id", mockResearch); - expect(redisConnection.set).toHaveBeenCalledWith( + expect(redisEvictConnection.set).toHaveBeenCalledWith( "deep-research:test-id", JSON.stringify(mockResearch) ); - expect(redisConnection.expire).toHaveBeenCalledWith( + expect(redisEvictConnection.expire).toHaveBeenCalledWith( "deep-research:test-id", 6 * 60 * 60 ); @@ -53,17 +53,17 @@ describe("Deep Research Redis Operations", () => { describe("getDeepResearch", () => { it("should retrieve research data from Redis", async () => { - (redisConnection.get as jest.Mock).mockResolvedValue( + (redisEvictConnection.get as jest.Mock).mockResolvedValue( JSON.stringify(mockResearch) ); const result = await getDeepResearch("test-id"); expect(result).toEqual(mockResearch); - expect(redisConnection.get).toHaveBeenCalledWith("deep-research:test-id"); + expect(redisEvictConnection.get).toHaveBeenCalledWith("deep-research:test-id"); }); it("should return null when research not found", async () => { - (redisConnection.get as jest.Mock).mockResolvedValue(null); + (redisEvictConnection.get as jest.Mock).mockResolvedValue(null); const result = await getDeepResearch("non-existent-id"); expect(result).toBeNull(); @@ -72,7 +72,7 @@ describe("Deep Research Redis Operations", () => { describe("updateDeepResearch", () => { it("should update existing research with new data", async () => { - (redisConnection.get as jest.Mock).mockResolvedValue( + (redisEvictConnection.get as jest.Mock).mockResolvedValue( JSON.stringify(mockResearch) ); @@ -98,30 +98,30 @@ describe("Deep Research Redis Operations", () => { activities: [...mockResearch.activities, ...update.activities], }; - expect(redisConnection.set).toHaveBeenCalledWith( + expect(redisEvictConnection.set).toHaveBeenCalledWith( "deep-research:test-id", JSON.stringify(expectedUpdate) ); - expect(redisConnection.expire).toHaveBeenCalledWith( + expect(redisEvictConnection.expire).toHaveBeenCalledWith( "deep-research:test-id", 6 * 60 * 60 ); }); it("should do nothing if research not found", async () => { - (redisConnection.get as jest.Mock).mockResolvedValue(null); + (redisEvictConnection.get as jest.Mock).mockResolvedValue(null); await updateDeepResearch("test-id", { status: "completed" }); - expect(redisConnection.set).not.toHaveBeenCalled(); - expect(redisConnection.expire).not.toHaveBeenCalled(); + expect(redisEvictConnection.set).not.toHaveBeenCalled(); + expect(redisEvictConnection.expire).not.toHaveBeenCalled(); }); }); describe("getDeepResearchExpiry", () => { it("should return correct expiry date", async () => { const mockTTL = 3600000; // 1 hour in milliseconds - (redisConnection.pttl as jest.Mock).mockResolvedValue(mockTTL); + (redisEvictConnection.pttl as jest.Mock).mockResolvedValue(mockTTL); const result = await getDeepResearchExpiry("test-id"); diff --git a/apps/api/src/__tests__/queue-concurrency-integration.test.ts b/apps/api/src/__tests__/queue-concurrency-integration.test.ts index c2b75befd..c60fe7276 100644 --- a/apps/api/src/__tests__/queue-concurrency-integration.test.ts +++ b/apps/api/src/__tests__/queue-concurrency-integration.test.ts @@ -1,4 +1,4 @@ -import { redisConnection } from "../services/queue-service"; +import { redisEvictConnection } from "../services/redis"; import { addScrapeJob, addScrapeJobs } from "../services/queue-jobs"; import { cleanOldConcurrencyLimitEntries, @@ -85,18 +85,18 @@ describe("Queue Concurrency Integration", () => { it("should add job directly to BullMQ when under concurrency limit", async () => { // Mock current active jobs to be under limit - (redisConnection.zrangebyscore as jest.Mock).mockResolvedValue([]); + (redisEvictConnection.zrangebyscore as jest.Mock).mockResolvedValue([]); await addScrapeJob(mockWebScraperOptions); // Should have checked concurrency - expect(redisConnection.zrangebyscore).toHaveBeenCalled(); + expect(redisEvictConnection.zrangebyscore).toHaveBeenCalled(); // Should have added to BullMQ expect(mockAdd).toHaveBeenCalled(); // Should have added to active jobs - expect(redisConnection.zadd).toHaveBeenCalledWith( + expect(redisEvictConnection.zadd).toHaveBeenCalledWith( expect.stringContaining("concurrency-limiter"), expect.any(Number), expect.any(String), @@ -109,20 +109,20 @@ describe("Queue Concurrency Integration", () => { concurrency: 15, } as any); const activeJobs = Array(15).fill("active-job"); - (redisConnection.zrangebyscore as jest.Mock).mockResolvedValue( + (redisEvictConnection.zrangebyscore as jest.Mock).mockResolvedValue( activeJobs, ); await addScrapeJob(mockWebScraperOptions); // Should have checked concurrency - expect(redisConnection.zrangebyscore).toHaveBeenCalled(); + expect(redisEvictConnection.zrangebyscore).toHaveBeenCalled(); // Should NOT have added to BullMQ expect(mockAdd).not.toHaveBeenCalled(); // Should have added to concurrency queue - expect(redisConnection.zadd).toHaveBeenCalledWith( + expect(redisEvictConnection.zadd).toHaveBeenCalledWith( expect.stringContaining("concurrency-limit-queue"), expect.any(Number), expect.stringContaining("mock-uuid"), @@ -157,7 +157,7 @@ describe("Queue Concurrency Integration", () => { const mockJobs = createMockJobs(totalJobs); // Mock current active jobs to be empty - (redisConnection.zrangebyscore as jest.Mock).mockResolvedValue([]); + (redisEvictConnection.zrangebyscore as jest.Mock).mockResolvedValue([]); await addScrapeJobs(mockJobs); @@ -165,7 +165,7 @@ describe("Queue Concurrency Integration", () => { expect(mockAdd).toHaveBeenCalledTimes(maxConcurrency); // Should have added remaining jobs to concurrency queue - expect(redisConnection.zadd).toHaveBeenCalledWith( + expect(redisEvictConnection.zadd).toHaveBeenCalledWith( expect.stringContaining("concurrency-limit-queue"), expect.any(Number), expect.any(String), @@ -176,7 +176,7 @@ describe("Queue Concurrency Integration", () => { const result = await addScrapeJobs([]); expect(result).toBe(true); expect(mockAdd).not.toHaveBeenCalled(); - expect(redisConnection.zadd).not.toHaveBeenCalled(); + expect(redisEvictConnection.zadd).not.toHaveBeenCalled(); }); }); @@ -196,7 +196,7 @@ describe("Queue Concurrency Integration", () => { data: { test: "data" }, opts: {}, }; - (redisConnection.zmpop as jest.Mock).mockResolvedValueOnce([ + (redisEvictConnection.zmpop as jest.Mock).mockResolvedValueOnce([ "key", [[JSON.stringify(queuedJob)]], ]); @@ -212,7 +212,7 @@ describe("Queue Concurrency Integration", () => { // Should have added new job to active jobs await pushConcurrencyLimitActiveJob(mockTeamId, nextJob!.id, 2 * 60 * 1000); - expect(redisConnection.zadd).toHaveBeenCalledWith( + expect(redisEvictConnection.zadd).toHaveBeenCalledWith( expect.stringContaining("concurrency-limiter"), expect.any(Number), nextJob!.id, @@ -235,7 +235,7 @@ describe("Queue Concurrency Integration", () => { await cleanOldConcurrencyLimitEntries(mockTeamId); // Verify job was removed from active jobs - expect(redisConnection.zrem).toHaveBeenCalledWith( + expect(redisEvictConnection.zrem).toHaveBeenCalledWith( expect.stringContaining("concurrency-limiter"), mockJob.id, ); @@ -247,14 +247,14 @@ describe("Queue Concurrency Integration", () => { const stalledTime = mockNow - 3 * 60 * 1000; // 3 minutes ago // Mock stalled jobs in Redis - (redisConnection.zrangebyscore as jest.Mock).mockResolvedValueOnce([ + (redisEvictConnection.zrangebyscore as jest.Mock).mockResolvedValueOnce([ "stalled-job", ]); await cleanOldConcurrencyLimitEntries(mockTeamId, mockNow); // Should have cleaned up stalled jobs - expect(redisConnection.zremrangebyscore).toHaveBeenCalledWith( + expect(redisEvictConnection.zremrangebyscore).toHaveBeenCalledWith( expect.stringContaining("concurrency-limiter"), -Infinity, mockNow, @@ -263,7 +263,7 @@ describe("Queue Concurrency Integration", () => { it("should handle race conditions in job queue processing", async () => { // Mock a race condition where job is taken by another worker - (redisConnection.zmpop as jest.Mock).mockResolvedValueOnce(null); + (redisEvictConnection.zmpop as jest.Mock).mockResolvedValueOnce(null); const nextJob = await takeConcurrencyLimitedJob(mockTeamId); diff --git a/apps/api/src/lib/__tests__/job-priority.test.ts b/apps/api/src/lib/__tests__/job-priority.test.ts index 118f3fe21..5e1605813 100644 --- a/apps/api/src/lib/__tests__/job-priority.test.ts +++ b/apps/api/src/lib/__tests__/job-priority.test.ts @@ -3,7 +3,7 @@ import { addJobPriority, deleteJobPriority, } from "../job-priority"; -import { redisConnection } from "../../services/queue-service"; +import { redisEvictConnection } from "../../services/redis"; import { } from "../../types"; jest.mock("../../services/queue-service", () => ({ @@ -24,11 +24,11 @@ describe("Job Priority Tests", () => { const team_id = "team1"; const job_id = "job1"; await addJobPriority(team_id, job_id); - expect(redisConnection.sadd).toHaveBeenCalledWith( + expect(redisEvictConnection.sadd).toHaveBeenCalledWith( `limit_team_id:${team_id}`, job_id, ); - expect(redisConnection.expire).toHaveBeenCalledWith( + expect(redisEvictConnection.expire).toHaveBeenCalledWith( `limit_team_id:${team_id}`, 60, ); @@ -38,7 +38,7 @@ describe("Job Priority Tests", () => { const team_id = "team1"; const job_id = "job1"; await deleteJobPriority(team_id, job_id); - expect(redisConnection.srem).toHaveBeenCalledWith( + expect(redisEvictConnection.srem).toHaveBeenCalledWith( `limit_team_id:${team_id}`, job_id, ); @@ -47,12 +47,12 @@ describe("Job Priority Tests", () => { test("getJobPriority should return correct priority based on plan and set length", async () => { const team_id = "team1"; const plan = "standard"; - (redisConnection.scard as jest.Mock).mockResolvedValue(150); + (redisEvictConnection.scard as jest.Mock).mockResolvedValue(150); const priority = await getJobPriority({ team_id }); expect(priority).toBe(10); - (redisConnection.scard as jest.Mock).mockResolvedValue(250); + (redisEvictConnection.scard as jest.Mock).mockResolvedValue(250); const priorityExceeded = await getJobPriority({ team_id }); expect(priorityExceeded).toBe(20); // basePriority + Math.ceil((250 - 200) * 0.4) }); @@ -60,22 +60,22 @@ describe("Job Priority Tests", () => { test("getJobPriority should handle different plans correctly", async () => { const team_id = "team1"; - (redisConnection.scard as jest.Mock).mockResolvedValue(50); + (redisEvictConnection.scard as jest.Mock).mockResolvedValue(50); let plan = "hobby"; let priority = await getJobPriority({ team_id }); expect(priority).toBe(10); - (redisConnection.scard as jest.Mock).mockResolvedValue(150); + (redisEvictConnection.scard as jest.Mock).mockResolvedValue(150); plan = "hobby"; priority = await getJobPriority({ team_id }); expect(priority).toBe(25); // basePriority + Math.ceil((150 - 50) * 0.3) - (redisConnection.scard as jest.Mock).mockResolvedValue(25); + (redisEvictConnection.scard as jest.Mock).mockResolvedValue(25); plan = "free"; priority = await getJobPriority({ team_id }); expect(priority).toBe(10); - (redisConnection.scard as jest.Mock).mockResolvedValue(60); + (redisEvictConnection.scard as jest.Mock).mockResolvedValue(60); plan = "free"; priority = await getJobPriority({ team_id }); expect(priority).toBe(28); // basePriority + Math.ceil((60 - 25) * 0.5) @@ -87,17 +87,17 @@ describe("Job Priority Tests", () => { const job_id2 = "job2"; await addJobPriority(team_id, job_id1); - expect(redisConnection.expire).toHaveBeenCalledWith( + expect(redisEvictConnection.expire).toHaveBeenCalledWith( `limit_team_id:${team_id}`, 60, ); // Clear the mock calls - (redisConnection.expire as jest.Mock).mockClear(); + (redisEvictConnection.expire as jest.Mock).mockClear(); // Add another job await addJobPriority(team_id, job_id2); - expect(redisConnection.expire).toHaveBeenCalledWith( + expect(redisEvictConnection.expire).toHaveBeenCalledWith( `limit_team_id:${team_id}`, 60, ); @@ -110,7 +110,7 @@ describe("Job Priority Tests", () => { jest.useFakeTimers(); await addJobPriority(team_id, job_id); - expect(redisConnection.expire).toHaveBeenCalledWith( + expect(redisEvictConnection.expire).toHaveBeenCalledWith( `limit_team_id:${team_id}`, 60, ); @@ -119,14 +119,14 @@ describe("Job Priority Tests", () => { jest.advanceTimersByTime(59000); // The set should still exist - expect(redisConnection.scard).not.toHaveBeenCalled(); + expect(redisEvictConnection.scard).not.toHaveBeenCalled(); // Fast-forward time by 2 more seconds (total 61 seconds) jest.advanceTimersByTime(2000); // Check if the set has been removed (scard should return 0) - (redisConnection.scard as jest.Mock).mockResolvedValue(0); - const setSize = await redisConnection.scard(`limit_team_id:${team_id}`); + (redisEvictConnection.scard as jest.Mock).mockResolvedValue(0); + const setSize = await redisEvictConnection.scard(`limit_team_id:${team_id}`); expect(setSize).toBe(0); jest.useRealTimers(); diff --git a/apps/api/src/services/billing/batch_billing.ts b/apps/api/src/services/billing/batch_billing.ts index 1a13e8e0a..679bf844e 100644 --- a/apps/api/src/services/billing/batch_billing.ts +++ b/apps/api/src/services/billing/batch_billing.ts @@ -1,5 +1,5 @@ import { logger } from "../../lib/logger"; -import { redisConnection } from "../queue-service"; +import { getRedisConnection } from "../queue-service"; import { supabase_service } from "../supabase"; import * as Sentry from "@sentry/node"; import { Queue } from "bullmq"; @@ -33,7 +33,7 @@ interface GroupedBillingOperation { // Function to acquire a lock for batch processing async function acquireLock(): Promise { - const redis = redisConnection; + const redis = getRedisConnection(); // Set lock with NX (only if it doesn't exist) and PX (millisecond expiry) const result = await redis.set(BATCH_LOCK_KEY, "1", "PX", LOCK_TIMEOUT, "NX"); const acquired = result === "OK"; @@ -45,14 +45,14 @@ async function acquireLock(): Promise { // Function to release the lock async function releaseLock() { - const redis = redisConnection; + const redis = getRedisConnection(); await redis.del(BATCH_LOCK_KEY); logger.info("🔓 Released billing batch processing lock"); } // Main function to process the billing batch export async function processBillingBatch() { - const redis = redisConnection; + const redis = getRedisConnection(); // Try to acquire lock if (!(await acquireLock())) { @@ -156,7 +156,7 @@ export function startBillingBatchProcessing() { logger.info("🔄 Starting periodic billing batch processing"); batchInterval = setInterval(async () => { - const queueLength = await redisConnection.llen(BATCH_KEY); + const queueLength = await getRedisConnection().llen(BATCH_KEY); logger.info(`Checking billing batch queue (${queueLength} items pending)`); await processBillingBatch(); }, BATCH_TIMEOUT); @@ -195,9 +195,9 @@ export async function queueBillingOperation( }; // Add operation to Redis list - const redis = redisConnection; + const redis = getRedisConnection(); await redis.rpush(BATCH_KEY, JSON.stringify(operation)); - const queueLength = await redis.llen(BATCH_KEY); + const queueLength = await getRedisConnection().llen(BATCH_KEY); logger.info(`📥 Added billing operation to queue (${queueLength} total pending)`, { team_id, credits diff --git a/apps/api/src/services/indexing/index-worker.ts b/apps/api/src/services/indexing/index-worker.ts index f39fef2a7..2cf3dd41e 100644 --- a/apps/api/src/services/indexing/index-worker.ts +++ b/apps/api/src/services/indexing/index-worker.ts @@ -4,7 +4,7 @@ import * as Sentry from "@sentry/node"; import { Job, Queue, Worker } from "bullmq"; import { logger as _logger, logger } from "../../lib/logger"; import { - redisConnection, + getRedisConnection, getBillingQueue, getPrecrawlQueue, precrawlQueueName, @@ -218,7 +218,7 @@ const workerFun = async (queue: Queue, jobProcessor: (token: string, job: Job) = const logger = _logger.child({ module: "index-worker", method: "workerFun" }); const worker = new Worker(queue.name, null, { - connection: redisConnection, + connection: getRedisConnection(), lockDuration: workerLockDuration, stalledInterval: workerStalledCheckInterval, maxStalledCount: queue.name === precrawlQueueName ? 0 : 10, diff --git a/apps/api/src/services/queue-service.ts b/apps/api/src/services/queue-service.ts index 4ba840a91..76c5e4a66 100644 --- a/apps/api/src/services/queue-service.ts +++ b/apps/api/src/services/queue-service.ts @@ -13,20 +13,21 @@ let deepResearchQueue: Queue; let generateLlmsTxtQueue: Queue; let billingQueue: Queue; let precrawlQueue: Queue; +let redisConnection: IORedis; -export function createRedisConnection() { - const connection = new IORedis(process.env.REDIS_URL!, { - maxRetriesPerRequest: null, - }); +export function getRedisConnection(): IORedis { + if (!redisConnection) { + redisConnection = new IORedis(process.env.REDIS_URL!, { + maxRetriesPerRequest: null, + }); + redisConnection.on("connect", () => logger.info("Redis connected")); + redisConnection.on("reconnecting", () => logger.warn("Redis reconnecting")); + redisConnection.on("error", (err) => logger.warn("Redis error", { err })); - connection.on("reconnecting", () => logger.warn("Redis reconnecting")); - connection.on("error", (err) => logger.warn("Redis error", { err })); - - return connection; + } + return redisConnection; } -export const redisConnection = createRedisConnection(); - export const scrapeQueueName = "{scrapeQueue}"; export const extractQueueName = "{extractQueue}"; export const loggingQueueName = "{loggingQueue}"; @@ -39,7 +40,7 @@ export const precrawlQueueName = "{precrawlQueue}"; export function getScrapeQueue() { if (!scrapeQueue) { scrapeQueue = new Queue(scrapeQueueName, { - connection: createRedisConnection(), + connection: getRedisConnection(), defaultJobOptions: { removeOnComplete: { age: 3600, // 1 hour @@ -56,7 +57,7 @@ export function getScrapeQueue() { export function getScrapeQueueEvents() { if (!scrapeQueueEvents) { scrapeQueueEvents = new QueueEvents(scrapeQueueName, { - connection: createRedisConnection(), + connection: getRedisConnection(), }); } @@ -66,7 +67,7 @@ export function getScrapeQueueEvents() { export function getExtractQueue() { if (!extractQueue) { extractQueue = new Queue(extractQueueName, { - connection: redisConnection, + connection: getRedisConnection(), defaultJobOptions: { removeOnComplete: { age: 90000, // 25 hours @@ -83,7 +84,7 @@ export function getExtractQueue() { export function getGenerateLlmsTxtQueue() { if (!generateLlmsTxtQueue) { generateLlmsTxtQueue = new Queue(generateLlmsTxtQueueName, { - connection: redisConnection, + connection: getRedisConnection(), defaultJobOptions: { removeOnComplete: { age: 90000, // 25 hours @@ -100,7 +101,7 @@ export function getGenerateLlmsTxtQueue() { export function getDeepResearchQueue() { if (!deepResearchQueue) { deepResearchQueue = new Queue(deepResearchQueueName, { - connection: redisConnection, + connection: getRedisConnection(), defaultJobOptions: { removeOnComplete: { age: 90000, // 25 hours @@ -117,7 +118,7 @@ export function getDeepResearchQueue() { export function getBillingQueue() { if (!billingQueue) { billingQueue = new Queue(billingQueueName, { - connection: redisConnection, + connection: getRedisConnection(), defaultJobOptions: { removeOnComplete: { age: 60, // 1 minute @@ -134,7 +135,7 @@ export function getBillingQueue() { export function getPrecrawlQueue() { if (!precrawlQueue) { precrawlQueue = new Queue(precrawlQueueName, { - connection: redisConnection, + connection: getRedisConnection(), defaultJobOptions: { removeOnComplete: { age: 24 * 60 * 60, // 1 day diff --git a/apps/api/src/services/queue-worker.ts b/apps/api/src/services/queue-worker.ts index 28fd8a25b..fdb4b111a 100644 --- a/apps/api/src/services/queue-worker.ts +++ b/apps/api/src/services/queue-worker.ts @@ -5,10 +5,9 @@ import { getScrapeQueue, getExtractQueue, getDeepResearchQueue, - redisConnection, getGenerateLlmsTxtQueue, scrapeQueueName, - createRedisConnection, + getRedisConnection, } from "./queue-service"; import { Job, Queue, QueueEvents } from "bullmq"; import { logger as _logger } from "../lib/logger"; @@ -357,7 +356,7 @@ const separateWorkerFun = ( .filter(arg => !arg.startsWith('--max-old-space-size')); const worker = new Worker(queue.name, path, { - connection: createRedisConnection(), + connection: getRedisConnection(), lockDuration: 60 * 1000, // 60 seconds stalledInterval: 60 * 1000, // 60 seconds maxStalledCount: 10, // 10 times @@ -386,7 +385,7 @@ const workerFun = async ( const logger = _logger.child({ module: "queue-worker", method: "workerFun" }); const worker = new Worker(queue.name, null, { - connection: redisConnection, + connection: getRedisConnection(), lockDuration: 60 * 1000, // 60 seconds stalledInterval: 60 * 1000, // 60 seconds maxStalledCount: 10, // 10 times @@ -534,7 +533,7 @@ app.listen(workerPort, () => { } } - const scrapeQueueEvents = new QueueEvents(scrapeQueueName, { connection: redisConnection }); + const scrapeQueueEvents = new QueueEvents(scrapeQueueName, { connection: getRedisConnection() }); scrapeQueueEvents.on("failed", failedListener); const results = await Promise.all([ From c0c4c7b890d591bf9f266bfa12ca5c4e9371c0da Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Gerg=C5=91=20M=C3=B3ricz?= Date: Fri, 15 Aug 2025 11:48:44 +0200 Subject: [PATCH 5/6] fix(crawl-status, extract-status): reduce the use of BullMQ getState (#1968) --- apps/api/src/controllers/v0/crawl-status.ts | 2 +- apps/api/src/controllers/v1/crawl-status.ts | 5 ++--- apps/api/src/controllers/v1/extract-status.ts | 2 +- 3 files changed, 4 insertions(+), 5 deletions(-) diff --git a/apps/api/src/controllers/v0/crawl-status.ts b/apps/api/src/controllers/v0/crawl-status.ts index 78c2ae0b9..843eb56d1 100644 --- a/apps/api/src/controllers/v0/crawl-status.ts +++ b/apps/api/src/controllers/v0/crawl-status.ts @@ -55,7 +55,7 @@ export async function getJobs(crawlId: string, ids: string[]): Promise = { id, - getState: bullJob ? (() => bullJob.getState()) : (() => dbJob!.success ? "completed" : "failed"), + getState: dbJob ? (() => dbJob.success ? "completed" : "failed") : (() => bullJob!.getState()), returnvalue: Array.isArray(data) ? data[0] : data, diff --git a/apps/api/src/controllers/v1/crawl-status.ts b/apps/api/src/controllers/v1/crawl-status.ts index 4da7112cf..34e16e567 100644 --- a/apps/api/src/controllers/v1/crawl-status.ts +++ b/apps/api/src/controllers/v1/crawl-status.ts @@ -58,7 +58,7 @@ export async function getJob(id: string): Promise | null> { const job: PseudoJob = { id, - getState: bullJob ? bullJob.getState : (() => dbJob!.success ? "completed" : "failed"), + getState: dbJob ? (() => dbJob.success ? "completed" : "failed") : bullJob!.getState, returnvalue: Array.isArray(data) ? data[0] : data, @@ -111,10 +111,9 @@ export async function getJobs(ids: string[]): Promise[]> { }); } - const state = await bullJob?.getState(); const job: PseudoJob = { id, - getState: bullJob ? (() => state!) : (() => dbJob!.success ? "completed" : "failed"), + getState: dbJob ? (() => dbJob.success ? "completed" : "failed") : (() => bullJob!.getState()), returnvalue: Array.isArray(data) ? data[0] : data, diff --git a/apps/api/src/controllers/v1/extract-status.ts b/apps/api/src/controllers/v1/extract-status.ts index 09d5e42ed..763482926 100644 --- a/apps/api/src/controllers/v1/extract-status.ts +++ b/apps/api/src/controllers/v1/extract-status.ts @@ -18,7 +18,7 @@ export async function getExtractJob(id: string): Promise = { id, - getState: bullJob ? bullJob.getState.bind(bullJob) : (() => dbJob!.success ? "completed" : "failed"), + getState: dbJob ? (() => dbJob.success ? "completed" : "failed") : (() => bullJob!.getState()), returnvalue: data, data: { scrapeOptions: bullJob ? bullJob.data.scrapeOptions : dbJob!.page_options, From 06a0198a73cb727479259607a065b354a9d0c9fe Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Gerg=C5=91=20M=C3=B3ricz?= Date: Fri, 15 Aug 2025 13:15:14 +0200 Subject: [PATCH 6/6] feat(api): add OTEL everywhere (#1969) --- apps/api/package.json | 6 +++ apps/api/pnpm-lock.yaml | 36 ++++++++++++++++-- apps/api/src/index.ts | 29 +++++++++++---- .../api/src/services/indexing/index-worker.ts | 2 + apps/api/src/services/queue-service.ts | 7 ++++ apps/api/src/services/queue-worker.ts | 37 +++++++++++++------ apps/api/src/services/worker/scrape-worker.ts | 27 ++++++++++---- 7 files changed, 114 insertions(+), 30 deletions(-) diff --git a/apps/api/package.json b/apps/api/package.json index 25673ba51..7440caafe 100644 --- a/apps/api/package.json +++ b/apps/api/package.json @@ -73,7 +73,12 @@ "@nangohq/node": "^0.40.8", "@openrouter/ai-sdk-provider": "^0.4.5", "@opentelemetry/auto-instrumentations-node": "^0.62.0", + "@opentelemetry/core": "^2.0.1", + "@opentelemetry/exporter-trace-otlp-grpc": "^0.203.0", + "@opentelemetry/resources": "^2.0.1", "@opentelemetry/sdk-node": "^0.203.0", + "@opentelemetry/sdk-trace-node": "^2.0.1", + "@opentelemetry/semantic-conventions": "^1.36.0", "@sentry/cli": "^2.33.1", "@sentry/node": "^9.40.0", "@sentry/profiling-node": "^9.40.0", @@ -88,6 +93,7 @@ "body-parser": "^1.20.3", "bottleneck": "^2.19.5", "bullmq": "^5.56.7", + "bullmq-otel": "^1.0.1", "cacheable-lookup": "^6.1.0", "cheerio": "^1.0.0-rc.12", "cohere": "^1.1.1", diff --git a/apps/api/pnpm-lock.yaml b/apps/api/pnpm-lock.yaml index 1c562dc4f..b1ba5e628 100644 --- a/apps/api/pnpm-lock.yaml +++ b/apps/api/pnpm-lock.yaml @@ -64,10 +64,25 @@ importers: version: 0.4.5(zod@3.24.2) '@opentelemetry/auto-instrumentations-node': specifier: ^0.62.0 - version: 0.62.0(@opentelemetry/api@1.9.0)(@opentelemetry/core@1.30.1(@opentelemetry/api@1.9.0))(encoding@0.1.13) + version: 0.62.0(@opentelemetry/api@1.9.0)(@opentelemetry/core@2.0.1(@opentelemetry/api@1.9.0))(encoding@0.1.13) + '@opentelemetry/core': + specifier: ^2.0.1 + version: 2.0.1(@opentelemetry/api@1.9.0) + '@opentelemetry/exporter-trace-otlp-grpc': + specifier: ^0.203.0 + version: 0.203.0(@opentelemetry/api@1.9.0) + '@opentelemetry/resources': + specifier: ^2.0.1 + version: 2.0.1(@opentelemetry/api@1.9.0) '@opentelemetry/sdk-node': specifier: ^0.203.0 version: 0.203.0(@opentelemetry/api@1.9.0) + '@opentelemetry/sdk-trace-node': + specifier: ^2.0.1 + version: 2.0.1(@opentelemetry/api@1.9.0) + '@opentelemetry/semantic-conventions': + specifier: ^1.36.0 + version: 1.36.0 '@sentry/cli': specifier: ^2.33.1 version: 2.33.1(encoding@0.1.13) @@ -110,6 +125,9 @@ importers: bullmq: specifier: ^5.56.7 version: 5.56.7 + bullmq-otel: + specifier: ^1.0.1 + version: 1.0.1 cacheable-lookup: specifier: ^6.1.0 version: 6.1.0 @@ -3210,6 +3228,9 @@ packages: resolution: {integrity: sha512-WDtdLmJvAuNNPzByAYpRo2rF1Mmradw6gvWsQKf63476DDXmomT9zUiGypLcG4ibIM67vhAj8jJRdbmEws2Aqw==} engines: {node: '>=6.14.2'} + bullmq-otel@1.0.1: + resolution: {integrity: sha512-9SSBA/iq8bFUJ9I5bDwv1UmtbJAmIGHt796g5lsYZWasqouTqkFA3+Z/n17mCc8VzfxxygCCCBsXdC1mICbxdw==} + bullmq@5.56.7: resolution: {integrity: sha512-Aa4Y7rmkuZOuyHvitGd2rSETgQ0SMawPD0XVsU9gAUg9Q4bwNpSD9ZDVcWGhfAQAUZFJH82JUmLsm0B59Prmtg==} @@ -7947,10 +7968,10 @@ snapshots: '@opentelemetry/api@1.9.0': {} - '@opentelemetry/auto-instrumentations-node@0.62.0(@opentelemetry/api@1.9.0)(@opentelemetry/core@1.30.1(@opentelemetry/api@1.9.0))(encoding@0.1.13)': + '@opentelemetry/auto-instrumentations-node@0.62.0(@opentelemetry/api@1.9.0)(@opentelemetry/core@2.0.1(@opentelemetry/api@1.9.0))(encoding@0.1.13)': dependencies: '@opentelemetry/api': 1.9.0 - '@opentelemetry/core': 1.30.1(@opentelemetry/api@1.9.0) + '@opentelemetry/core': 2.0.1(@opentelemetry/api@1.9.0) '@opentelemetry/instrumentation': 0.203.0(@opentelemetry/api@1.9.0) '@opentelemetry/instrumentation-amqplib': 0.50.0(@opentelemetry/api@1.9.0) '@opentelemetry/instrumentation-aws-lambda': 0.54.0(@opentelemetry/api@1.9.0) @@ -10816,7 +10837,7 @@ snapshots: async-mutex@0.2.6: dependencies: - tslib: 2.6.3 + tslib: 2.8.1 async-mutex@0.5.0: dependencies: @@ -11014,6 +11035,13 @@ snapshots: dependencies: node-gyp-build: 4.8.4 + bullmq-otel@1.0.1: + dependencies: + '@opentelemetry/api': 1.9.0 + bullmq: 5.56.7 + transitivePeerDependencies: + - supports-color + bullmq@5.56.7: dependencies: cron-parser: 4.9.0 diff --git a/apps/api/src/index.ts b/apps/api/src/index.ts index 9144d0c8a..a099ff188 100644 --- a/apps/api/src/index.ts +++ b/apps/api/src/index.ts @@ -30,6 +30,10 @@ import domainFrequencyRouter from "./routes/domain-frequency"; import { NodeSDK } from "@opentelemetry/sdk-node"; import { LangfuseExporter } from "langfuse-vercel"; import { getNodeAutoInstrumentations } from "@opentelemetry/auto-instrumentations-node"; +import { ATTR_SERVICE_NAME } from "@opentelemetry/semantic-conventions"; +import { resourceFromAttributes } from "@opentelemetry/resources"; +import { BatchSpanProcessor } from "@opentelemetry/sdk-trace-node"; +import { OTLPTraceExporter } from "@opentelemetry/exporter-trace-otlp-grpc"; const { createBullBoard } = require("@bull-board/api"); const { BullMQAdapter } = require('@bull-board/api/bullMQAdapter'); @@ -46,12 +50,22 @@ logger.info("Network info dump", { cacheableLookup.install(http.globalAgent); cacheableLookup.install(https.globalAgent); -const langfuseOtel = process.env.LANGFUSE_PUBLIC_KEY ? new NodeSDK({ - traceExporter: new LangfuseExporter(), - instrumentations: [getNodeAutoInstrumentations()] +const shouldOtel = process.env.LANGFUSE_PUBLIC_KEY || process.env.OTEL_EXPORTER_OTLP_ENDPOINT; +const otelSdk = shouldOtel ? new NodeSDK({ + resource: resourceFromAttributes({ + [ATTR_SERVICE_NAME]: "firecrawl-app", + }), + spanProcessors: [ + ...(process.env.LANGFUSE_PUBLIC_KEY ? [new BatchSpanProcessor(new LangfuseExporter())] : []), + ...(process.env.OTEL_EXPORTER_OTLP_ENDPOINT ? [new BatchSpanProcessor(new OTLPTraceExporter({ + url: process.env.OTEL_EXPORTER_OTLP_ENDPOINT, + }))] : []), + ], + instrumentations: [getNodeAutoInstrumentations()], }) : null; -if (langfuseOtel) { - langfuseOtel.start(); + +if (otelSdk) { + otelSdk.start(); } // Initialize Express with WebSocket support @@ -121,8 +135,9 @@ function startServer(port = DEFAULT_PORT) { } server.close(() => { logger.info("Server closed."); - if (langfuseOtel) { - langfuseOtel.shutdown().then(() => { + if (otelSdk) { + otelSdk.shutdown().then(() => { + logger.info("OTEL shutdown"); process.exit(0); }); } else { diff --git a/apps/api/src/services/indexing/index-worker.ts b/apps/api/src/services/indexing/index-worker.ts index 2cf3dd41e..758169c25 100644 --- a/apps/api/src/services/indexing/index-worker.ts +++ b/apps/api/src/services/indexing/index-worker.ts @@ -17,6 +17,7 @@ import { processWebhookInsertJobs } from "../webhook"; import { scrapeOptions as scrapeOptionsSchema, crawlRequestSchema, toLegacyCrawlerOptions } from "../../controllers/v1/types"; import { StoredCrawl, crawlToCrawler, saveCrawl } from "../../lib/crawl-redis"; import { _addScrapeJobToBullMQ } from "../queue-jobs"; +import { BullMQOtel } from "bullmq-otel"; const workerLockDuration = Number(process.env.WORKER_LOCK_DURATION) || 60000; const workerStalledCheckInterval = @@ -222,6 +223,7 @@ const workerFun = async (queue: Queue, jobProcessor: (token: string, job: Job) = lockDuration: workerLockDuration, stalledInterval: workerStalledCheckInterval, maxStalledCount: queue.name === precrawlQueueName ? 0 : 10, + telemetry: new BullMQOtel("firecrawl-bullmq"), }); worker.startStalledCheckTimer(); diff --git a/apps/api/src/services/queue-service.ts b/apps/api/src/services/queue-service.ts index 76c5e4a66..6041b9ec8 100644 --- a/apps/api/src/services/queue-service.ts +++ b/apps/api/src/services/queue-service.ts @@ -1,6 +1,7 @@ import { Queue, QueueEvents } from "bullmq"; import { logger } from "../lib/logger"; import IORedis from "ioredis"; +import { BullMQOtel } from "bullmq-otel"; export type QueueFunction = () => Queue; @@ -49,6 +50,7 @@ export function getScrapeQueue() { age: 3600, // 1 hour }, }, + telemetry: new BullMQOtel("firecrawl-bullmq"), }); } return scrapeQueue; @@ -76,6 +78,7 @@ export function getExtractQueue() { age: 90000, // 25 hours }, }, + telemetry: new BullMQOtel("firecrawl-bullmq"), }); } return extractQueue; @@ -93,6 +96,7 @@ export function getGenerateLlmsTxtQueue() { age: 90000, // 25 hours }, }, + telemetry: new BullMQOtel("firecrawl-bullmq"), }); } return generateLlmsTxtQueue; @@ -110,6 +114,7 @@ export function getDeepResearchQueue() { age: 90000, // 25 hours }, }, + telemetry: new BullMQOtel("firecrawl-bullmq"), }); } return deepResearchQueue; @@ -127,6 +132,7 @@ export function getBillingQueue() { age: 3600, // 1 hour }, }, + telemetry: new BullMQOtel("firecrawl-bullmq"), }); } return billingQueue; @@ -144,6 +150,7 @@ export function getPrecrawlQueue() { age: 24 * 60 * 60, // 1 day }, }, + telemetry: new BullMQOtel("firecrawl-bullmq"), }); } return precrawlQueue; diff --git a/apps/api/src/services/queue-worker.ts b/apps/api/src/services/queue-worker.ts index fdb4b111a..03c8998b1 100644 --- a/apps/api/src/services/queue-worker.ts +++ b/apps/api/src/services/queue-worker.ts @@ -46,6 +46,11 @@ import { finishCrawlIfNeeded } from "./worker/crawl-logic"; import { NodeSDK } from "@opentelemetry/sdk-node"; import { getNodeAutoInstrumentations } from "@opentelemetry/auto-instrumentations-node"; import { LangfuseExporter } from "langfuse-vercel"; +import { BatchSpanProcessor } from "@opentelemetry/sdk-trace-node"; +import { ATTR_SERVICE_NAME } from "@opentelemetry/semantic-conventions"; +import { resourceFromAttributes } from "@opentelemetry/resources"; +import { OTLPTraceExporter } from "@opentelemetry/exporter-trace-otlp-grpc"; +import { BullMQOtel } from "bullmq-otel"; configDotenv(); @@ -68,16 +73,22 @@ const runningJobs: Set = new Set(); cacheableLookup.install(http.globalAgent); cacheableLookup.install(https.globalAgent); -const langfuseOtel = process.env.LANGFUSE_PUBLIC_KEY ? new NodeSDK({ - traceExporter: new LangfuseExporter(), - instrumentations: [getNodeAutoInstrumentations({ - '@opentelemetry/instrumentation-undici': { enabled: false }, - '@opentelemetry/instrumentation-http': { enabled: false }, - })], +const shouldOtel = process.env.LANGFUSE_PUBLIC_KEY || process.env.OTEL_EXPORTER_OTLP_ENDPOINT; +const otelSdk = shouldOtel ? new NodeSDK({ + resource: resourceFromAttributes({ + [ATTR_SERVICE_NAME]: "firecrawl-worker", + }), + spanProcessors: [ + ...(process.env.LANGFUSE_PUBLIC_KEY ? [new BatchSpanProcessor(new LangfuseExporter())] : []), + ...(process.env.OTEL_EXPORTER_OTLP_ENDPOINT ? [new BatchSpanProcessor(new OTLPTraceExporter({ + url: process.env.OTEL_EXPORTER_OTLP_ENDPOINT, + }))] : []), + ], + instrumentations: [getNodeAutoInstrumentations()], }) : null; - -if (langfuseOtel) { - langfuseOtel.start(); + +if (otelSdk) { + otelSdk.start(); } const processExtractJobInternal = async ( @@ -372,7 +383,8 @@ const separateWorkerFun = ( resourceLimits: maxOldSpaceSize ? { maxOldGenerationSizeMb: parseInt(maxOldSpaceSize) } : undefined - } + }, + telemetry: new BullMQOtel("firecrawl-bullmq"), }); return worker; @@ -389,6 +401,7 @@ const workerFun = async ( lockDuration: 60 * 1000, // 60 seconds stalledInterval: 60 * 1000, // 60 seconds maxStalledCount: 10, // 10 times + telemetry: new BullMQOtel("firecrawl-bullmq"), }); worker.startStalledCheckTimer(); @@ -562,8 +575,8 @@ app.listen(workerPort, () => { await scrapeQueueEvents.close(); console.log("All jobs finished. Worker out!"); - if (langfuseOtel) { - await langfuseOtel.shutdown(); + if (otelSdk) { + await otelSdk.shutdown(); } process.exit(0); })(); \ No newline at end of file diff --git a/apps/api/src/services/worker/scrape-worker.ts b/apps/api/src/services/worker/scrape-worker.ts index 84f80c9be..8cd7541d0 100644 --- a/apps/api/src/services/worker/scrape-worker.ts +++ b/apps/api/src/services/worker/scrape-worker.ts @@ -48,6 +48,10 @@ import { finishCrawlIfNeeded } from "./crawl-logic"; import { LangfuseExporter } from "langfuse-vercel"; import { getNodeAutoInstrumentations } from "@opentelemetry/auto-instrumentations-node"; import { NodeSDK } from "@opentelemetry/sdk-node"; +import { BatchSpanProcessor } from "@opentelemetry/sdk-trace-node"; +import { OTLPTraceExporter } from "@opentelemetry/exporter-trace-otlp-grpc"; +import { resourceFromAttributes } from "@opentelemetry/resources"; +import { ATTR_SERVICE_NAME } from "@opentelemetry/semantic-conventions"; class RacedRedirectError extends Error { constructor() { @@ -849,21 +853,30 @@ export const processJobInternal = async (job: Job & { id: string }) => { } }; -const langfuseOtel = process.env.LANGFUSE_PUBLIC_KEY ? new NodeSDK({ - traceExporter: new LangfuseExporter(), +const shouldOtel = process.env.LANGFUSE_PUBLIC_KEY || process.env.OTEL_EXPORTER_OTLP_ENDPOINT; +const otelSdk = shouldOtel ? new NodeSDK({ + resource: resourceFromAttributes({ + [ATTR_SERVICE_NAME]: "firecrawl-worker-scrape", + }), + spanProcessors: [ + ...(process.env.LANGFUSE_PUBLIC_KEY ? [new BatchSpanProcessor(new LangfuseExporter())] : []), + ...(process.env.OTEL_EXPORTER_OTLP_ENDPOINT ? [new BatchSpanProcessor(new OTLPTraceExporter({ + url: process.env.OTEL_EXPORTER_OTLP_ENDPOINT, + }))] : []), + ], instrumentations: [getNodeAutoInstrumentations()], }) : null; -if (langfuseOtel) { - langfuseOtel.start(); +if (otelSdk) { + otelSdk.start(); } module.exports = processJobInternal; const exitHandler = () => { - if (langfuseOtel) { - langfuseOtel.shutdown().then(() => { - _logger.debug("Langfuse OTEL shutdown"); + if (otelSdk) { + otelSdk.shutdown().then(() => { + _logger.debug("OTEL shutdown"); process.exit(0); }); }