f1bd11a7ef
## Summary Adds a single env flag, `DEPRECATE_V3_ENABLED` (default off), that gracefully winds down the v3 engine (`RunEngineVersion.V1`). While it's off nothing changes, so self-hosted instances still on v3 keep working. When it's on: - Triggers that resolve to v3 are rejected with a clear, actionable error pointing at the [v4 migration guide](https://trigger.dev/docs/migrating-from-v3), instead of silently creating runs that never execute. This covers single triggers, batches, scheduled fires, replays, and `triggerAndWait`, which all funnel through one place. - The legacy `trigger dev` websocket used by v3 CLIs is closed with an upgrade message (v4 CLIs use a different dev transport). - The v3 shared-queue consumer refuses to start, so no deployed v3 runs are dequeued. - The v3 run-lifecycle background jobs (heartbeat timeout, TTL expiry, retry, resume batch/dependency, delayed-run enqueue, and scheduled fires) become no-ops, so abandoned v3 runs stop generating database load. This builds on the existing deploy deprecation flag, which already rejects v3 CLI deploys. ## Design Enforcement is read through one helper, `isV3Disabled()`. Every gate combines it with a per-run or per-project engine check (`isV3Disabled() && engine === "V1"`), so a v4 run that happens to reach a shared service behaves exactly as before. v4 (V2) is never affected. The flag is a hard switch, not a drain: when it's on, in-flight v3 runs are abandoned in place rather than failed or expired, which is the intended behaviour for the final shutdown.
170 lines
5.3 KiB
TypeScript
170 lines
5.3 KiB
TypeScript
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";
|
|
import { devPresence } from "~/presenters/v3/DevPresence.server";
|
|
import { logger } from "~/services/logger.server";
|
|
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";
|
|
import { isV3Disabled } from "./engineDeprecation.server";
|
|
|
|
export const scheduleEngine = singleton("ScheduleEngine", createScheduleEngine);
|
|
|
|
export type { ScheduleEngine };
|
|
|
|
async function isDevEnvironmentConnectedHandler(environmentId: string) {
|
|
const environment = await prisma.runtimeEnvironment.findFirst({
|
|
where: {
|
|
id: environmentId,
|
|
},
|
|
select: {
|
|
currentSession: {
|
|
select: {
|
|
disconnectedAt: true,
|
|
},
|
|
},
|
|
project: {
|
|
select: {
|
|
engine: true,
|
|
},
|
|
},
|
|
},
|
|
});
|
|
|
|
if (!environment) {
|
|
return false;
|
|
}
|
|
|
|
if (environment.project.engine === "V1") {
|
|
const v3Disconnected = !environment.currentSession || environment.currentSession.disconnectedAt;
|
|
|
|
return !v3Disconnected;
|
|
}
|
|
|
|
const v4Connected = await devPresence.isConnected(environmentId);
|
|
|
|
return v4Connected;
|
|
}
|
|
|
|
function createScheduleEngine() {
|
|
const engine = new ScheduleEngine({
|
|
prisma,
|
|
logLevel: env.SCHEDULE_ENGINE_LOG_LEVEL,
|
|
redis: {
|
|
host: env.SCHEDULE_WORKER_REDIS_HOST ?? "localhost",
|
|
port: env.SCHEDULE_WORKER_REDIS_PORT ?? 6379,
|
|
username: env.SCHEDULE_WORKER_REDIS_USERNAME,
|
|
password: env.SCHEDULE_WORKER_REDIS_PASSWORD,
|
|
keyPrefix: "schedule:",
|
|
enableAutoPipelining: true,
|
|
...(env.SCHEDULE_WORKER_REDIS_TLS_DISABLED === "true" ? {} : { tls: {} }),
|
|
},
|
|
worker: {
|
|
concurrency: env.SCHEDULE_WORKER_CONCURRENCY_LIMIT,
|
|
workers: env.SCHEDULE_WORKER_CONCURRENCY_WORKERS,
|
|
tasksPerWorker: env.SCHEDULE_WORKER_CONCURRENCY_TASKS_PER_WORKER,
|
|
pollIntervalMs: env.SCHEDULE_WORKER_POLL_INTERVAL,
|
|
shutdownTimeoutMs: env.SCHEDULE_WORKER_SHUTDOWN_TIMEOUT_MS,
|
|
disabled: env.SCHEDULE_WORKER_ENABLED === "0",
|
|
},
|
|
distributionWindow: {
|
|
seconds: env.SCHEDULE_WORKER_DISTRIBUTION_WINDOW_SECONDS,
|
|
},
|
|
tracer,
|
|
meter,
|
|
onTriggerScheduledTask: async ({
|
|
taskIdentifier,
|
|
environment,
|
|
payload,
|
|
scheduleInstanceId,
|
|
scheduleId,
|
|
exactScheduleTime,
|
|
}) => {
|
|
try {
|
|
// v3 (engine V1) shutdown: skip firing schedules for V1 projects so the
|
|
// cron doesn't keep doing trigger work just to be rejected. Return success
|
|
// so the schedule engine treats it as handled and doesn't retry. v4 is
|
|
// unaffected.
|
|
if (isV3Disabled() && environment.project.engine === "V1") {
|
|
logger.debug("[ScheduleEngine] Skipping scheduled fire for shut-down v3 project", {
|
|
taskIdentifier,
|
|
scheduleId,
|
|
});
|
|
return { success: true };
|
|
}
|
|
|
|
// This will trigger either v1 or v2 depending on the engine of the project
|
|
const triggerService = new TriggerTaskService();
|
|
|
|
const payloadPacket = await stringifyIO(payload);
|
|
|
|
logger.debug("Triggering scheduled task", {
|
|
taskIdentifier,
|
|
environment,
|
|
payload,
|
|
scheduleInstanceId,
|
|
scheduleId,
|
|
exactScheduleTime,
|
|
});
|
|
|
|
const result = await triggerService.call(
|
|
taskIdentifier,
|
|
environment,
|
|
{ payload: payloadPacket.data, options: { payloadType: payloadPacket.dataType } },
|
|
{
|
|
customIcon: "scheduled",
|
|
scheduleId,
|
|
scheduleInstanceId,
|
|
queueTimestamp: exactScheduleTime,
|
|
overrideCreatedAt: exactScheduleTime,
|
|
triggerSource: "schedule",
|
|
triggerAction: "trigger",
|
|
}
|
|
);
|
|
|
|
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: errorMessage,
|
|
errorType,
|
|
};
|
|
}
|
|
},
|
|
isDevEnvironmentConnectedHandler: isDevEnvironmentConnectedHandler,
|
|
onRegisterScheduleInstance: removeDeprecatedWorkerQueueItem,
|
|
});
|
|
|
|
return engine;
|
|
}
|
|
|
|
async function removeDeprecatedWorkerQueueItem(instanceId: string) {
|
|
// We need to dequeue the instance from the existing workerQueue
|
|
try {
|
|
await workerQueue.dequeue(`scheduled-task-instance:${instanceId}`);
|
|
|
|
logger.debug("Removed deprecated worker queue item", {
|
|
instanceId,
|
|
});
|
|
} catch (error) {
|
|
logger.error("Error dequeuing scheduled task instance from deprecated queue", {
|
|
instanceId,
|
|
error: error instanceof Error ? error.message : String(error),
|
|
});
|
|
}
|
|
}
|