From 947f33d55be6727252ed89f7e139c9ac63ae1659 Mon Sep 17 00:00:00 2001 From: Eric Allam Date: Wed, 25 Mar 2026 16:08:30 +0000 Subject: [PATCH] fix: downgrade queue size limit errors to warnings (#3243) Queue limit ServiceValidationErrors were being logged at error level. These are expected validation rejections, not bugs. - Add logLevel property to ServiceValidationError (webapp + run-engine) - Set logLevel: warn on all queue limit throws - Schedule engine: detect queue limit failures and log as warn - Redis-worker: respect logLevel on thrown errors --- .../runEngine/services/triggerTask.server.ts | 4 +- apps/webapp/app/v3/scheduleEngine.server.ts | 15 +++- .../app/v3/services/batchTriggerV3.server.ts | 4 +- apps/webapp/app/v3/services/common.server.ts | 8 +- .../app/v3/services/triggerTaskV1.server.ts | 4 +- .../run-engine/src/engine/errors.ts | 5 +- .../run-engine/src/run-queue/index.ts | 10 +-- .../schedule-engine/src/engine/index.ts | 25 ++++-- .../schedule-engine/src/engine/types.ts | 8 +- .../schedule-engine/src/index.ts | 1 + packages/redis-worker/src/worker.ts | 82 +++++++++++-------- 11 files changed, 113 insertions(+), 53 deletions(-) diff --git a/apps/webapp/app/runEngine/services/triggerTask.server.ts b/apps/webapp/app/runEngine/services/triggerTask.server.ts index bbd1331e6..1fe3b839e 100644 --- a/apps/webapp/app/runEngine/services/triggerTask.server.ts +++ b/apps/webapp/app/runEngine/services/triggerTask.server.ts @@ -264,7 +264,9 @@ export class RunEngineTriggerTaskService { if (!queueSizeGuard.ok) { throw new ServiceValidationError( - `Cannot trigger ${taskId} as the queue size limit for this environment has been reached. The maximum size is ${queueSizeGuard.maximumSize}` + `Cannot trigger ${taskId} as the queue size limit for this environment has been reached. The maximum size is ${queueSizeGuard.maximumSize}`, + undefined, + "warn" ); } } diff --git a/apps/webapp/app/v3/scheduleEngine.server.ts b/apps/webapp/app/v3/scheduleEngine.server.ts index cbec21a3b..4a12046a8 100644 --- a/apps/webapp/app/v3/scheduleEngine.server.ts +++ b/apps/webapp/app/v3/scheduleEngine.server.ts @@ -1,4 +1,5 @@ import { ScheduleEngine } from "@internal/schedule-engine"; +import type { TriggerScheduledTaskErrorType } from "@internal/schedule-engine"; import { stringifyIO } from "@trigger.dev/core/v3"; import { prisma } from "~/db.server"; import { env } from "~/env.server"; @@ -8,6 +9,7 @@ import { singleton } from "~/utils/singleton"; import { TriggerTaskService } from "./services/triggerTask.server"; import { meter, tracer } from "./tracer.server"; import { workerQueue } from "~/services/worker.server"; +import { ServiceValidationError } from "./services/common.server"; export const scheduleEngine = singleton("ScheduleEngine", createScheduleEngine); @@ -113,9 +115,20 @@ function createScheduleEngine() { return { success: !!result }; } catch (error) { + const errorMessage = error instanceof Error ? error.message : String(error); + let errorType: TriggerScheduledTaskErrorType = "SYSTEM_ERROR"; + + if ( + error instanceof ServiceValidationError && + errorMessage.includes("queue size limit for this environment has been reached") + ) { + errorType = "QUEUE_LIMIT"; + } + return { success: false, - error: error instanceof Error ? error.message : String(error), + error: errorMessage, + errorType, }; } }, diff --git a/apps/webapp/app/v3/services/batchTriggerV3.server.ts b/apps/webapp/app/v3/services/batchTriggerV3.server.ts index 2ed8a38b3..f0531f844 100644 --- a/apps/webapp/app/v3/services/batchTriggerV3.server.ts +++ b/apps/webapp/app/v3/services/batchTriggerV3.server.ts @@ -251,7 +251,9 @@ export class BatchTriggerV3Service extends BaseService { if (!queueSizeGuard.isWithinLimits) { throw new ServiceValidationError( - `Cannot trigger ${newRunCount} tasks as the queue size limit for this environment has been reached. The maximum size is ${queueSizeGuard.maximumSize}` + `Cannot trigger ${newRunCount} tasks as the queue size limit for this environment has been reached. The maximum size is ${queueSizeGuard.maximumSize}`, + undefined, + "warn" ); } diff --git a/apps/webapp/app/v3/services/common.server.ts b/apps/webapp/app/v3/services/common.server.ts index 31f401e00..015e8e23d 100644 --- a/apps/webapp/app/v3/services/common.server.ts +++ b/apps/webapp/app/v3/services/common.server.ts @@ -1,5 +1,11 @@ +export type ServiceValidationErrorLevel = "error" | "warn" | "info"; + export class ServiceValidationError extends Error { - constructor(message: string, public status?: number) { + constructor( + message: string, + public status?: number, + public logLevel?: ServiceValidationErrorLevel + ) { super(message); this.name = "ServiceValidationError"; } diff --git a/apps/webapp/app/v3/services/triggerTaskV1.server.ts b/apps/webapp/app/v3/services/triggerTaskV1.server.ts index efc6510ef..d36130109 100644 --- a/apps/webapp/app/v3/services/triggerTaskV1.server.ts +++ b/apps/webapp/app/v3/services/triggerTaskV1.server.ts @@ -134,7 +134,9 @@ export class TriggerTaskServiceV1 extends BaseService { if (!queueSizeGuard.isWithinLimits) { throw new ServiceValidationError( - `Cannot trigger ${taskId} as the queue size limit for this environment has been reached. The maximum size is ${queueSizeGuard.maximumSize}` + `Cannot trigger ${taskId} as the queue size limit for this environment has been reached. The maximum size is ${queueSizeGuard.maximumSize}`, + undefined, + "warn" ); } } diff --git a/internal-packages/run-engine/src/engine/errors.ts b/internal-packages/run-engine/src/engine/errors.ts index 772282deb..820f0ec4c 100644 --- a/internal-packages/run-engine/src/engine/errors.ts +++ b/internal-packages/run-engine/src/engine/errors.ts @@ -69,11 +69,14 @@ export function runStatusFromError( } } +export type ServiceValidationErrorLevel = "error" | "warn" | "info"; + export class ServiceValidationError extends Error { constructor( message: string, public status?: number, - public metadata?: Record + public metadata?: Record, + public logLevel?: ServiceValidationErrorLevel ) { super(message); this.name = "ServiceValidationError"; diff --git a/internal-packages/run-engine/src/run-queue/index.ts b/internal-packages/run-engine/src/run-queue/index.ts index 8432aaaaf..9088099ef 100644 --- a/internal-packages/run-engine/src/run-queue/index.ts +++ b/internal-packages/run-engine/src/run-queue/index.ts @@ -2545,7 +2545,7 @@ export class RunQueue { return; } - this.logger.info("Processing concurrency keys from stream", { + this.logger.debug("Processing concurrency keys from stream", { keys: uniqueKeys, }); @@ -2615,9 +2615,9 @@ export class RunQueue { } private async processCurrentConcurrencyRunIds(concurrencyKey: string, runIds: string[]) { - this.logger.info("Processing concurrency set with runs", { + this.logger.debug("Processing concurrency set with runs", { concurrencyKey, - runIds: runIds.slice(0, 5), // Log first 5 for debugging, + runIds: runIds.slice(0, 5), runIdsLength: runIds.length, }); @@ -2625,12 +2625,12 @@ export class RunQueue { const completedRuns = await this.options.concurrencySweeper?.callback(runIds); if (!completedRuns) { - this.logger.info("No completed runs found in concurrency set", { concurrencyKey }); + this.logger.debug("No completed runs found in concurrency set", { concurrencyKey }); return; } if (completedRuns.length === 0) { - this.logger.info("No completed runs found in concurrency set", { concurrencyKey }); + this.logger.debug("No completed runs found in concurrency set", { concurrencyKey }); return; } diff --git a/internal-packages/schedule-engine/src/engine/index.ts b/internal-packages/schedule-engine/src/engine/index.ts index 47dc1e31e..4eb641176 100644 --- a/internal-packages/schedule-engine/src/engine/index.ts +++ b/internal-packages/schedule-engine/src/engine/index.ts @@ -497,17 +497,28 @@ export class ScheduleEngine { span.setAttribute("trigger_success", true); } else { - this.logger.error("Failed to trigger scheduled task", { - instanceId: params.instanceId, - taskIdentifier: instance.taskSchedule.taskIdentifier, - durationMs: triggerDuration, - error: result.error, - }); + const isQueueLimit = result.errorType === "QUEUE_LIMIT"; + + if (isQueueLimit) { + this.logger.warn("Scheduled task trigger skipped due to queue limit", { + instanceId: params.instanceId, + taskIdentifier: instance.taskSchedule.taskIdentifier, + durationMs: triggerDuration, + error: result.error, + }); + } else { + this.logger.error("Failed to trigger scheduled task", { + instanceId: params.instanceId, + taskIdentifier: instance.taskSchedule.taskIdentifier, + durationMs: triggerDuration, + error: result.error, + }); + } this.scheduleExecutionFailureCounter.add(1, { environment_type: environmentType, schedule_type: scheduleType, - error_type: "task_failure", + error_type: isQueueLimit ? "queue_limit" : "task_failure", }); span.setAttribute("trigger_success", false); diff --git a/internal-packages/schedule-engine/src/engine/types.ts b/internal-packages/schedule-engine/src/engine/types.ts index 8cad80c30..f4888c447 100644 --- a/internal-packages/schedule-engine/src/engine/types.ts +++ b/internal-packages/schedule-engine/src/engine/types.ts @@ -24,8 +24,14 @@ export type TriggerScheduledTaskParams = { exactScheduleTime?: Date; }; +export type TriggerScheduledTaskErrorType = "QUEUE_LIMIT" | "SYSTEM_ERROR"; + export interface TriggerScheduledTaskCallback { - (params: TriggerScheduledTaskParams): Promise<{ success: boolean; error?: string }>; + (params: TriggerScheduledTaskParams): Promise<{ + success: boolean; + error?: string; + errorType?: TriggerScheduledTaskErrorType; + }>; } export interface ScheduleEngineOptions { diff --git a/internal-packages/schedule-engine/src/index.ts b/internal-packages/schedule-engine/src/index.ts index 22f6c07b3..6c96f2cd5 100644 --- a/internal-packages/schedule-engine/src/index.ts +++ b/internal-packages/schedule-engine/src/index.ts @@ -3,4 +3,5 @@ export type { ScheduleEngineOptions, TriggerScheduleParams, TriggerScheduledTaskCallback, + TriggerScheduledTaskErrorType, } from "./engine/types.js"; diff --git a/packages/redis-worker/src/worker.ts b/packages/redis-worker/src/worker.ts index 2435cd9fb..51812e575 100644 --- a/packages/redis-worker/src/worker.ts +++ b/packages/redis-worker/src/worker.ts @@ -745,23 +745,25 @@ class Worker { ).catch(async (error) => { const errorMessage = error instanceof Error ? error.message : String(error); const shouldLogError = catalogItem.logErrors ?? true; + const errorLogLevel = + error && typeof error === "object" && "logLevel" in error ? error.logLevel : undefined; - if (shouldLogError) { - this.logger.error(`Worker error processing batch`, { - name: this.options.name, - jobType, - batchSize: items.length, - error, - errorMessage, - }); + const logAttributes = { + name: this.options.name, + jobType, + batchSize: items.length, + error, + errorMessage, + }; + + if (!shouldLogError) { + this.logger.info(`Worker failed to process batch`, logAttributes); + } else if (errorLogLevel === "warn") { + this.logger.warn(`Worker error processing batch`, logAttributes); + } else if (errorLogLevel === "info") { + this.logger.info(`Worker error processing batch`, logAttributes); } else { - this.logger.info(`Worker failed to process batch`, { - name: this.options.name, - jobType, - batchSize: items.length, - error, - errorMessage, - }); + this.logger.error(`Worker error processing batch`, logAttributes); } // Re-enqueue each item individually with retry logic @@ -775,20 +777,21 @@ class Worker { const retryDelay = calculateNextRetryDelay(retrySettings, newAttempt); if (!retryDelay) { - if (shouldLogError) { - this.logger.error(`Worker batch item reached max attempts. Moving to DLQ.`, { - name: this.options.name, - id: item.id, - jobType, - attempt: newAttempt, - }); + const dlqLogAttributes = { + name: this.options.name, + id: item.id, + jobType, + attempt: newAttempt, + }; + + if (!shouldLogError) { + this.logger.info(`Worker batch item reached max attempts. Moving to DLQ.`, dlqLogAttributes); + } else if (errorLogLevel === "warn") { + this.logger.warn(`Worker batch item reached max attempts. Moving to DLQ.`, dlqLogAttributes); + } else if (errorLogLevel === "info") { + this.logger.info(`Worker batch item reached max attempts. Moving to DLQ.`, dlqLogAttributes); } else { - this.logger.info(`Worker batch item reached max attempts. Moving to DLQ.`, { - name: this.options.name, - id: item.id, - jobType, - attempt: newAttempt, - }); + this.logger.error(`Worker batch item reached max attempts. Moving to DLQ.`, dlqLogAttributes); } await this.queue.moveToDeadLetterQueue(item.id, errorMessage); @@ -895,6 +898,8 @@ class Worker { const errorMessage = error instanceof Error ? error.message : String(error); const shouldLogError = catalogItem.logErrors ?? true; + const errorLogLevel = + error && typeof error === "object" && "logLevel" in error ? error.logLevel : undefined; const logAttributes = { name: this.options.name, @@ -906,10 +911,14 @@ class Worker { errorMessage, }; - if (shouldLogError) { - this.logger.error(`Worker error processing item`, logAttributes); - } else { + if (!shouldLogError) { this.logger.info(`Worker failed to process item`, logAttributes); + } else if (errorLogLevel === "warn") { + this.logger.warn(`Worker error processing item`, logAttributes); + } else if (errorLogLevel === "info") { + this.logger.info(`Worker error processing item`, logAttributes); + } else { + this.logger.error(`Worker error processing item`, logAttributes); } // Attempt requeue logic. @@ -922,13 +931,18 @@ class Worker { const retryDelay = calculateNextRetryDelay(retrySettings, newAttempt); if (!retryDelay) { - if (shouldLogError) { - this.logger.error(`Worker item reached max attempts. Moving to DLQ.`, { + if (!shouldLogError || errorLogLevel === "info") { + this.logger.info(`Worker item reached max attempts. Moving to DLQ.`, { + ...logAttributes, + attempt: newAttempt, + }); + } else if (errorLogLevel === "warn") { + this.logger.warn(`Worker item reached max attempts. Moving to DLQ.`, { ...logAttributes, attempt: newAttempt, }); } else { - this.logger.info(`Worker item reached max attempts. Moving to DLQ.`, { + this.logger.error(`Worker item reached max attempts. Moving to DLQ.`, { ...logAttributes, attempt: newAttempt, });