re2: env based queue selection algo (#1775)

* re2: fix @trigger.dev/core exports

* re2: WIP env based queue selection algo

* more wip

* WIP

* Get run engine tests to pass

* Adding tests for the fair dequeueing strat in the run engine

* Configure the new queue selection strategy in the webapp and get it all building and typechecks passing

* webapp now uses built packages, building redis-worker, run-engine, database, using better tsconfig setups for tests, moving isomorphic code into core/v3/isomorphic

* Fixed webapp typechecks

* dev now depends on build, fixed supervisor typecheck

* Fixed run engine tests

* Fixed e2e tests
This commit is contained in:
Eric Allam
2025-03-07 14:30:19 +00:00
committed by GitHub
parent 38ddd830d7
commit d855d55ea0
149 changed files with 4466 additions and 2822 deletions
-14
View File
@@ -1,14 +0,0 @@
module.exports = {
root: true,
// This tells ESLint to load the config from the package `eslint-config-custom`
extends: ["custom"],
settings: {
next: {
rootDir: ["apps/*/"],
},
},
parserOptions: {
sourceType: "module",
ecmaVersion: 2020,
},
};
+8
View File
@@ -141,6 +141,14 @@
"command": "pnpm run test --filter @internal/run-engine",
"cwd": "${workspaceFolder}",
"sourceMaps": true
},
{
"type": "node-terminal",
"request": "launch",
"name": "Debug RunQueue tests",
"command": "pnpm run test ./src/engine/tests/waitpoints.test.ts",
"cwd": "${workspaceFolder}/internal-packages/run-engine",
"sourceMaps": true
}
]
}
-4
View File
@@ -4,9 +4,5 @@
"compilerOptions": {
"rootDir": "src",
"outDir": "dist"
},
"paths": {
"@trigger.dev/core/v3": ["../../packages/core/src/v3"],
"@trigger.dev/core/v3/*": ["../../packages/core/src/v3/*"]
}
}
+1 -6
View File
@@ -1,10 +1,5 @@
{
"plugins": [
"@trigger.dev/eslint-plugin",
"react-hooks",
"@typescript-eslint/eslint-plugin",
"import"
],
"plugins": ["react-hooks", "@typescript-eslint/eslint-plugin", "import"],
"parser": "@typescript-eslint/parser",
"overrides": [
{
@@ -11,7 +11,7 @@ import {
RunPanelIconSection,
RunPanelProperties,
} from "./RunCard";
import { DisplayProperty } from "@trigger.dev/core";
import type { DisplayProperty } from "@trigger.dev/core";
export function TriggerDetail({
trigger,
+1 -1
View File
@@ -14,7 +14,7 @@ import {
OperatingSystemContextProvider,
OperatingSystemPlatform,
} from "./components/primitives/OperatingSystemProvider";
import { getSharedSqsEventConsumer } from "./services/events/sqsEventConsumer";
import { getSharedSqsEventConsumer } from "./services/events/sqsEventConsumer.server";
import { singleton } from "./utils/singleton";
const ABORT_DELAY = 30000;
+6
View File
@@ -419,6 +419,12 @@ const EnvironmentSchema = z.object({
RUN_ENGINE_TIMEOUT_EXECUTING: z.coerce.number().int().default(60_000),
RUN_ENGINE_TIMEOUT_EXECUTING_WITH_WAITPOINTS: z.coerce.number().int().default(60_000),
RUN_ENGINE_DEBUG_WORKER_NOTIFICATIONS: z.coerce.boolean().default(false),
RUN_ENGINE_PARENT_QUEUE_LIMIT: z.coerce.number().int().default(1000),
RUN_ENGINE_CONCURRENCY_LIMIT_BIAS: z.coerce.number().default(0.75),
RUN_ENGINE_AVAILABLE_CAPACITY_BIAS: z.coerce.number().default(0.3),
RUN_ENGINE_QUEUE_AGE_RANDOMIZATION_BIAS: z.coerce.number().default(0.25),
RUN_ENGINE_REUSE_SNAPSHOT_COUNT: z.coerce.number().int().default(0),
RUN_ENGINE_MAXIMUM_ENV_COUNT: z.coerce.number().int().optional(),
RUN_ENGINE_WORKER_REDIS_HOST: z
.string()
+1 -1
View File
@@ -1,4 +1,4 @@
import { Prettify } from "@trigger.dev/core";
import type { Prettify } from "@trigger.dev/core";
import { TaskRun } from "@trigger.dev/database";
import { SyncedShapeData, useSyncedShape } from "./useSyncedShape";
@@ -1,7 +1,7 @@
import { z } from "zod";
import { PrismaClient, prisma } from "~/db.server";
import { sortEnvironments } from "~/utils/environmentSort";
import { httpEndpointUrl } from "~/services/httpendpoint/HandleHttpEndpointService";
import { httpEndpointUrl } from "~/services/httpendpoint/HandleHttpEndpointService.server";
import { getSecretStore } from "~/services/secrets/secretStore.server";
import { projectPath } from "~/utils/pathBuilder";
@@ -11,7 +11,7 @@ import { eventRepository } from "~/v3/eventRepository.server";
import { machinePresetFromName } from "~/v3/machinePresets.server";
import { FINAL_ATTEMPT_STATUSES, isFailedRunStatus, isFinalRunStatus } from "~/v3/taskStatus";
import { BasePresenter } from "./basePresenter.server";
import { getMaxDuration } from "@trigger.dev/core/v3/apps";
import { getMaxDuration } from "@trigger.dev/core/v3/isomorphic";
import { logger } from "~/services/logger.server";
import { getTaskEventStoreTableForRun, TaskEventStoreTable } from "~/v3/taskEventStore.server";
import { Pi } from "lucide-react";
@@ -20,7 +20,7 @@ import { logger } from "~/services/logger.server";
import { BasePresenter } from "./basePresenter.server";
import { TaskRunStatus } from "~/database-types";
import { concurrencyTracker } from "~/v3/services/taskRunConcurrencyTracker.server";
import { CURRENT_DEPLOYMENT_LABEL } from "@trigger.dev/core/v3/apps";
import { CURRENT_DEPLOYMENT_LABEL } from "@trigger.dev/core/v3/isomorphic";
export type Task = {
slug: string;
@@ -1,6 +1,9 @@
import type { ActionFunctionArgs } from "@remix-run/server-runtime";
import { json } from "@remix-run/server-runtime";
import { RegisterScheduleBodySchema, RegisterScheduleResponseBodySchema } from "@trigger.dev/core";
import {
RegisterScheduleBodySchema,
RegisterScheduleResponseBodySchema,
} from "@trigger.dev/core/schemas";
import { z } from "zod";
import { authenticateApiRequest } from "~/services/apiAuth.server";
import { logger } from "~/services/logger.server";
@@ -1,6 +1,6 @@
import type { ActionFunctionArgs } from "@remix-run/server-runtime";
import { json } from "@remix-run/server-runtime";
import { UpdateTriggerSourceBodyV1Schema } from "@trigger.dev/core";
import { UpdateTriggerSourceBodyV1Schema } from "@trigger.dev/core/schemas";
import { z } from "zod";
import { authenticateApiRequest } from "~/services/apiAuth.server";
import { logger } from "~/services/logger.server";
@@ -1,6 +1,6 @@
import type { ActionFunctionArgs } from "@remix-run/server-runtime";
import { json } from "@remix-run/server-runtime";
import { RegisterTriggerBodySchemaV1 } from "@trigger.dev/core";
import { RegisterTriggerBodySchemaV1 } from "@trigger.dev/core/schemas";
import { z } from "zod";
import { authenticateApiRequest } from "~/services/apiAuth.server";
import { logger } from "~/services/logger.server";
@@ -1,6 +1,6 @@
import type { ActionFunctionArgs } from "@remix-run/server-runtime";
import { json } from "@remix-run/server-runtime";
import { InitializeTriggerBodySchema } from "@trigger.dev/core";
import { InitializeTriggerBodySchema } from "@trigger.dev/core/schemas";
import { z } from "zod";
import { authenticateApiRequest } from "~/services/apiAuth.server";
import { logger } from "~/services/logger.server";
@@ -1,6 +1,9 @@
import type { ActionFunctionArgs } from "@remix-run/server-runtime";
import { json } from "@remix-run/server-runtime";
import { CreateExternalConnectionBodySchema, ErrorWithStackSchema } from "@trigger.dev/core";
import {
CreateExternalConnectionBodySchema,
ErrorWithStackSchema,
} from "@trigger.dev/core/schemas";
import { z } from "zod";
import { generateErrorMessage } from "zod-error";
import { authenticateApiRequest } from "~/services/apiAuth.server";
@@ -1,4 +1,7 @@
import { GetEndpointIndexResponse, GetEndpointIndexResponseSchema } from "@trigger.dev/core";
import {
GetEndpointIndexResponse,
GetEndpointIndexResponseSchema,
} from "@trigger.dev/core/schemas";
import { ActionFunctionArgs, json } from "@remix-run/server-runtime";
import { z } from "zod";
import { prisma } from "~/db.server";
@@ -1,14 +1,9 @@
import type { ActionFunctionArgs } from "@remix-run/server-runtime";
import { json } from "@remix-run/server-runtime";
import {
EphemeralEventDispatcherRequestBodySchema,
InvokeJobRequestBodySchema,
} from "@trigger.dev/core";
import { z } from "zod";
import { EphemeralEventDispatcherRequestBodySchema } from "@trigger.dev/core/schemas";
import { PrismaErrorSchema } from "~/db.server";
import { authenticateApiRequest } from "~/services/apiAuth.server";
import { CreateEphemeralEventDispatcherService } from "~/services/dispatchers/createEphemeralEventDispatcher.server";
import { InvokeJobService } from "~/services/jobs/invokeJob.server";
import { logger } from "~/services/logger.server";
export async function action({ request, params }: ActionFunctionArgs) {
@@ -1,6 +1,6 @@
import type { LoaderFunctionArgs } from "@remix-run/server-runtime";
import { json } from "@remix-run/server-runtime";
import { GetEvent } from "@trigger.dev/core";
import { GetEvent } from "@trigger.dev/core/schemas";
import { z } from "zod";
import { prisma } from "~/db.server";
import { runOriginalStatus } from "~/models/jobRun.server";
+1 -1
View File
@@ -1,6 +1,6 @@
import type { ActionFunctionArgs } from "@remix-run/server-runtime";
import { json } from "@remix-run/server-runtime";
import { SendBulkEventsBodySchema } from "@trigger.dev/core";
import { SendBulkEventsBodySchema } from "@trigger.dev/core/schemas";
import { generateErrorMessage } from "zod-error";
import { authenticateApiRequest } from "~/services/apiAuth.server";
import { IngestSendEvent } from "~/services/events/ingestSendEvent.server";
+1 -1
View File
@@ -1,6 +1,6 @@
import type { ActionFunctionArgs } from "@remix-run/server-runtime";
import { json } from "@remix-run/server-runtime";
import { SendEventBodySchema } from "@trigger.dev/core";
import { SendEventBodySchema } from "@trigger.dev/core/schemas";
import { generateErrorMessage } from "zod-error";
import { authenticateApiRequest } from "~/services/apiAuth.server";
import { IngestSendEvent } from "~/services/events/ingestSendEvent.server";
@@ -2,7 +2,7 @@ import type { ActionFunctionArgs, LoaderFunctionArgs } from "@remix-run/server-r
import {
HandleHttpEndpointService,
HttpEndpointParamsSchema,
} from "~/services/httpendpoint/HandleHttpEndpointService";
} from "~/services/httpendpoint/HandleHttpEndpointService.server";
import { logger } from "~/services/logger.server";
export async function action({ request, params }: ActionFunctionArgs) {
@@ -1,6 +1,6 @@
import type { ActionFunctionArgs } from "@remix-run/server-runtime";
import { json } from "@remix-run/server-runtime";
import { InvokeJobRequestBodySchema } from "@trigger.dev/core";
import { InvokeJobRequestBodySchema } from "@trigger.dev/core/schemas";
import { z } from "zod";
import { PrismaErrorSchema } from "~/db.server";
import { authenticateApiRequest } from "~/services/apiAuth.server";
@@ -1,6 +1,6 @@
import type { ActionFunctionArgs } from "@remix-run/server-runtime";
import { json } from "@remix-run/server-runtime";
import { LogMessageSchema } from "@trigger.dev/core";
import { LogMessageSchema } from "@trigger.dev/core/schemas";
import { z } from "zod";
import { authenticateApiRequest } from "~/services/apiAuth.server";
import { CreateRunLogService } from "./CreateRunLogService.server";
@@ -1,6 +1,6 @@
import type { ActionFunctionArgs } from "@remix-run/server-runtime";
import { json } from "@remix-run/server-runtime";
import { JobRunStatusRecordSchema, StatusUpdateSchema } from "@trigger.dev/core";
import { JobRunStatusRecordSchema, StatusUpdateSchema } from "@trigger.dev/core/schemas";
import { z } from "zod";
import { authenticateApiRequest } from "~/services/apiAuth.server";
import { logger } from "~/services/logger.server";
@@ -1,6 +1,6 @@
import type { LoaderFunctionArgs } from "@remix-run/server-runtime";
import { json } from "@remix-run/server-runtime";
import { JobRunStatusRecordSchema } from "@trigger.dev/core";
import { JobRunStatusRecordSchema } from "@trigger.dev/core/schemas";
import { z } from "zod";
import { prisma } from "~/db.server";
import { runOriginalStatus } from "~/models/jobRun.server";
@@ -1,11 +1,10 @@
import type { ActionFunctionArgs } from "@remix-run/server-runtime";
import { json } from "@remix-run/server-runtime";
import type { CompleteTaskBodyOutput } from "@trigger.dev/core";
import type { CompleteTaskBodyOutput } from "@trigger.dev/core/schemas";
import {
API_VERSIONS,
CompleteTaskBodyInputSchema,
CompleteTaskBodyV2InputSchema,
} from "@trigger.dev/core";
} from "@trigger.dev/core/schemas";
import { z } from "zod";
import type { AuthenticatedEnvironment } from "~/services/apiAuth.server";
import { authenticateApiRequest } from "~/services/apiAuth.server";
@@ -14,6 +13,11 @@ import { startActiveSpan } from "~/v3/tracer.server";
import { parseRequestJsonAsync } from "~/utils/parseRequestJson.server";
import { FailRunTaskService } from "../api.v1.runs.$runId.tasks.$id.fail/FailRunTaskService.server";
const API_VERSIONS = {
LAZY_LOADED_CACHED_TASKS: "2023-09-29",
SERIALIZED_TASK_OUTPUT: "2023-11-01",
};
const ParamsSchema = z.object({
runId: z.string(),
id: z.string(),
@@ -1,6 +1,6 @@
import type { ActionFunctionArgs } from "@remix-run/server-runtime";
import { json } from "@remix-run/server-runtime";
import { FailTaskBodyInputSchema } from "@trigger.dev/core";
import { FailTaskBodyInputSchema } from "@trigger.dev/core/schemas";
import { z } from "zod";
import { authenticateApiRequest } from "~/services/apiAuth.server";
import { logger } from "~/services/logger.server";
@@ -1,6 +1,6 @@
import type { ActionFunctionArgs } from "@remix-run/server-runtime";
import { json } from "@remix-run/server-runtime";
import { API_VERSIONS, RunTaskBodyOutputSchema } from "@trigger.dev/core";
import { RunTaskBodyOutputSchema } from "@trigger.dev/core/schemas";
import { z } from "zod";
import { authenticateApiRequest } from "~/services/apiAuth.server";
import { logger } from "~/services/logger.server";
@@ -9,6 +9,11 @@ import { ChangeRequestLazyLoadedCachedTasks } from "./ChangeRequestLazyLoadedCac
import { startActiveSpan } from "~/v3/tracer.server";
import { parseRequestJsonAsync } from "~/utils/parseRequestJson.server";
const API_VERSIONS = {
LAZY_LOADED_CACHED_TASKS: "2023-09-29",
SERIALIZED_TASK_OUTPUT: "2023-11-01",
};
const ParamsSchema = z.object({
runId: z.string(),
});
@@ -5,7 +5,7 @@ import {
conditionallyExportPacket,
stringifyIO,
} from "@trigger.dev/core/v3";
import { WaitpointId } from "@trigger.dev/core/v3/apps";
import { WaitpointId } from "@trigger.dev/core/v3/isomorphic";
import { z } from "zod";
import { $replica } from "~/db.server";
import { env } from "~/env.server";
@@ -3,7 +3,7 @@ import {
CreateWaitpointTokenRequestBody,
CreateWaitpointTokenResponseBody,
} from "@trigger.dev/core/v3";
import { WaitpointId } from "@trigger.dev/core/v3/apps";
import { WaitpointId } from "@trigger.dev/core/v3/isomorphic";
import { createActionApiRoute } from "~/services/routeBuilders/apiBuilder.server";
import { parseDelay } from "~/utils/delays";
import { resolveIdempotencyKeyTTL } from "~/utils/idempotencyKeys.server";
@@ -1,6 +1,6 @@
import type { ActionFunctionArgs } from "@remix-run/server-runtime";
import { json } from "@remix-run/server-runtime";
import { UpdateWebhookBodySchema } from "@trigger.dev/core";
import { UpdateWebhookBodySchema } from "@trigger.dev/core/schemas";
import { z } from "zod";
import { authenticateApiRequest } from "~/services/apiAuth.server";
import { logger } from "~/services/logger.server";
@@ -1,6 +1,6 @@
import type { ActionFunctionArgs } from "@remix-run/server-runtime";
import { json } from "@remix-run/server-runtime";
import { UpdateTriggerSourceBodyV2Schema } from "@trigger.dev/core";
import { UpdateTriggerSourceBodyV2Schema } from "@trigger.dev/core/schemas";
import { z } from "zod";
import { authenticateApiRequest } from "~/services/apiAuth.server";
import { logger } from "~/services/logger.server";
@@ -4,7 +4,7 @@ import {
REGISTER_SOURCE_EVENT_V2,
RegisterSourceEventV2,
RegisterTriggerBodySchemaV2,
} from "@trigger.dev/core";
} from "@trigger.dev/core/schemas";
import { z } from "zod";
import { authenticateApiRequest } from "~/services/apiAuth.server";
import { IngestSendEvent } from "~/services/events/ingestSendEvent.server";
@@ -1,6 +1,6 @@
import type { LoaderFunctionArgs } from "@remix-run/server-runtime";
import { json } from "@remix-run/server-runtime";
import { GetEvent } from "@trigger.dev/core";
import { GetEvent } from "@trigger.dev/core/schemas";
import { z } from "zod";
import { $replica } from "~/db.server";
import { authenticateApiRequest } from "~/services/apiAuth.server";
@@ -1,6 +1,6 @@
import type { LoaderFunctionArgs } from "@remix-run/server-runtime";
import { json } from "@remix-run/server-runtime";
import { JobRunStatusRecordSchema } from "@trigger.dev/core";
import { JobRunStatusRecordSchema } from "@trigger.dev/core/schemas";
import { z } from "zod";
import { prisma } from "~/db.server";
import { authenticateApiRequest } from "~/services/apiAuth.server";
@@ -1,6 +1,6 @@
import { json } from "@remix-run/server-runtime";
import { DequeuedMessage, DevDequeueRequestBody, MachineResources } from "@trigger.dev/core/v3";
import { BackgroundWorkerId } from "@trigger.dev/core/v3/apps";
import { BackgroundWorkerId } from "@trigger.dev/core/v3/isomorphic";
import { env } from "~/env.server";
import { createActionApiRoute } from "~/services/routeBuilders/apiBuilder.server";
import { engine } from "~/v3/runEngine.server";
@@ -1,6 +1,6 @@
import { TypedResponse } from "@remix-run/server-runtime";
import { assertExhaustive } from "@trigger.dev/core";
import { RunId } from "@trigger.dev/core/v3/apps";
import { assertExhaustive } from "@trigger.dev/core/utils";
import { RunId } from "@trigger.dev/core/v3/isomorphic";
import {
WorkerApiDebugLogBody,
WorkerApiRunAttemptStartResponseBody,
@@ -1,18 +1,13 @@
import { json, TypedResponse } from "@remix-run/server-runtime";
import { assertExhaustive } from "@trigger.dev/core";
import { RunId, SnapshotId } from "@trigger.dev/core/v3/apps";
import { RunId, SnapshotId } from "@trigger.dev/core/v3/isomorphic";
import {
WorkerApiDebugLogBody,
WorkerApiRunAttemptCompleteRequestBody,
WorkerApiRunAttemptCompleteResponseBody,
WorkerApiRunAttemptStartResponseBody,
WorkloadHeartbeatResponseBody,
} from "@trigger.dev/core/v3/workers";
import { z } from "zod";
import { prisma } from "~/db.server";
import { logger } from "~/services/logger.server";
import { createActionApiRoute } from "~/services/routeBuilders/apiBuilder.server";
import { recordRunDebugLog } from "~/v3/eventRepository.server";
import { engine } from "~/v3/runEngine.server";
const { action } = createActionApiRoute(
@@ -1,6 +1,6 @@
import { json, TypedResponse } from "@remix-run/server-runtime";
import { MachinePreset } from "@trigger.dev/core/v3";
import { RunId, SnapshotId } from "@trigger.dev/core/v3/apps";
import { RunId, SnapshotId } from "@trigger.dev/core/v3/isomorphic";
import {
WorkerApiRunAttemptStartRequestBody,
WorkerApiRunAttemptStartResponseBody,
@@ -1,16 +1,10 @@
import { json, TypedResponse } from "@remix-run/server-runtime";
import { assertExhaustive } from "@trigger.dev/core";
import { RunId, SnapshotId } from "@trigger.dev/core/v3/apps";
import {
WorkerApiDebugLogBody,
WorkerApiRunAttemptStartResponseBody,
WorkloadHeartbeatResponseBody,
} from "@trigger.dev/core/v3/workers";
import { RunId, SnapshotId } from "@trigger.dev/core/v3/isomorphic";
import { WorkloadHeartbeatResponseBody } from "@trigger.dev/core/v3/workers";
import { z } from "zod";
import { prisma } from "~/db.server";
import { logger } from "~/services/logger.server";
import { createActionApiRoute } from "~/services/routeBuilders/apiBuilder.server";
import { recordRunDebugLog } from "~/v3/eventRepository.server";
import { engine } from "~/v3/runEngine.server";
const { action } = createActionApiRoute(
@@ -1,5 +1,5 @@
import { json, TypedResponse } from "@remix-run/server-runtime";
import { RunId } from "@trigger.dev/core/v3/apps";
import { RunId } from "@trigger.dev/core/v3/isomorphic";
import { WorkerApiRunLatestSnapshotResponseBody } from "@trigger.dev/core/v3/workers";
import { z } from "zod";
import { prisma } from "~/db.server";
@@ -1,6 +1,6 @@
import { json, TypedResponse } from "@remix-run/server-runtime";
import { WaitForDurationRequestBody, WaitForDurationResponseBody } from "@trigger.dev/core/v3";
import { RunId } from "@trigger.dev/core/v3/apps";
import { RunId } from "@trigger.dev/core/v3/isomorphic";
import { z } from "zod";
import { prisma } from "~/db.server";
@@ -1,6 +1,6 @@
import { json } from "@remix-run/server-runtime";
import { WaitForWaitpointTokenResponseBody } from "@trigger.dev/core/v3";
import { RunId, WaitpointId } from "@trigger.dev/core/v3/apps";
import { RunId, WaitpointId } from "@trigger.dev/core/v3/isomorphic";
import { z } from "zod";
import { $replica } from "~/db.server";
import { logger } from "~/services/logger.server";
@@ -1,5 +1,5 @@
import { json, TypedResponse } from "@remix-run/server-runtime";
import { CURRENT_DEPLOYMENT_LABEL } from "@trigger.dev/core/v3/apps";
import { CURRENT_DEPLOYMENT_LABEL } from "@trigger.dev/core/v3/isomorphic";
import { WorkerApiDequeueResponseBody } from "@trigger.dev/core/v3/workers";
import { z } from "zod";
import { $replica, prisma } from "~/db.server";
@@ -1,5 +1,5 @@
import { assertExhaustive } from "@trigger.dev/core";
import { RunId } from "@trigger.dev/core/v3/apps";
import { assertExhaustive } from "@trigger.dev/core/utils";
import { RunId } from "@trigger.dev/core/v3/isomorphic";
import { WorkerApiDebugLogBody } from "@trigger.dev/core/v3/runEngineWorker";
import { z } from "zod";
import { createActionWorkerApiRoute } from "~/services/routeBuilders/apiBuilder.server";
@@ -1,6 +1,6 @@
import { parse } from "@conform-to/zod";
import { ActionFunction, json } from "@remix-run/node";
import { assertExhaustive } from "@trigger.dev/core";
import { assertExhaustive } from "@trigger.dev/core/utils";
import { z } from "zod";
import { redirectWithErrorMessage, redirectWithSuccessMessage } from "~/models/message.server";
import { logger } from "~/services/logger.server";
@@ -8,8 +8,8 @@ import {
stringifyIO,
timeoutError,
} from "@trigger.dev/core/v3";
import { WaitpointId } from "@trigger.dev/core/v3/apps";
import { Waitpoint } from "@trigger.dev/database";
import { WaitpointId } from "@trigger.dev/core/v3/isomorphic";
import type { Waitpoint } from "@trigger.dev/database";
import { useCallback, useRef } from "react";
import { z } from "zod";
import { AnimatedHourglassIcon } from "~/assets/icons/AnimatedHourglassIcon";
@@ -1,6 +1,6 @@
import type { EndpointIndexSource } from "@trigger.dev/database";
import { PrismaClient, prisma } from "~/db.server";
import { PerformEndpointIndexService } from "./performEndpointIndexService";
import { PerformEndpointIndexService } from "./performEndpointIndexService.server";
export class IndexEndpointService {
#prismaClient: PrismaClient;
@@ -6,7 +6,7 @@ import { Prisma, WebhookEnvironment } from "@trigger.dev/database";
import { ulid } from "../ulid.server";
import { getSecretStore } from "../secrets/secretStore.server";
import { z } from "zod";
import { httpEndpointUrl } from "../httpendpoint/HandleHttpEndpointService";
import { httpEndpointUrl } from "../httpendpoint/HandleHttpEndpointService.server";
import { isEqual } from "ohash";
type ExtendedWebhook = Prisma.WebhookGetPayload<{
+1 -1
View File
@@ -25,7 +25,7 @@ import { ExpireDispatcherService } from "./dispatchers/expireDispatcher.server";
import { InvokeEphemeralDispatcherService } from "./dispatchers/invokeEphemeralEventDispatcher.server";
import { sendEmail } from "./email.server";
import { IndexEndpointService } from "./endpoints/indexEndpoint.server";
import { PerformEndpointIndexService } from "./endpoints/performEndpointIndexService";
import { PerformEndpointIndexService } from "./endpoints/performEndpointIndexService.server";
import { ProbeEndpointService } from "./endpoints/probeEndpoint.server";
import { RecurringEndpointIndexService } from "./endpoints/recurringEndpointIndex.server";
import { DeliverEventService } from "./events/deliverEvent.server";
+1 -1
View File
@@ -1,4 +1,4 @@
import { parseNaturalLanguageDuration } from "@trigger.dev/core/v3/apps";
import { parseNaturalLanguageDuration } from "@trigger.dev/core/v3/isomorphic";
export const calculateDurationInMs = (options: {
seconds?: number;
+1 -1
View File
@@ -1 +1 @@
export { generateFriendlyId } from "@trigger.dev/core/v3/apps";
export { generateFriendlyId } from "@trigger.dev/core/v3/isomorphic";
+1 -1
View File
@@ -9,7 +9,7 @@ import {
ProviderToPlatformMessages,
SharedQueueToClientMessages,
} from "@trigger.dev/core/v3";
import { RunId } from "@trigger.dev/core/v3/apps";
import { RunId } from "@trigger.dev/core/v3/isomorphic";
import type {
WorkerClientToServerEvents,
WorkerServerToClientEvents,
@@ -19,7 +19,7 @@ import { FailedTaskRunService } from "../failedTaskRun.server";
import { CancelDevSessionRunsService } from "../services/cancelDevSessionRuns.server";
import { CompleteAttemptService } from "../services/completeAttempt.server";
import { attributesFromAuthenticatedEnv, tracer } from "../tracer.server";
import { getMaxDuration } from "@trigger.dev/core/v3/apps";
import { getMaxDuration } from "@trigger.dev/core/v3/isomorphic";
import { DevSubscriber, devPubSub } from "./devPubSub.server";
import { findQueueInEnvironment, sanitizeQueueName } from "~/models/taskQueue.server";
import { createRedisClient, RedisClient } from "~/redis.server";
@@ -1,4 +1,3 @@
import { flattenAttributes } from "@trigger.dev/core/v3";
import { createCache, DefaultStatefulContext, Namespace, Cache as UnkeyCache } from "@unkey/cache";
import { MemoryStore } from "@unkey/cache/stores";
import { randomUUID } from "crypto";
@@ -3,7 +3,7 @@ import { BackgroundWorker, WorkerDeployment } from "@trigger.dev/database";
import {
CURRENT_DEPLOYMENT_LABEL,
CURRENT_UNMANAGED_DEPLOYMENT_LABEL,
} from "@trigger.dev/core/v3/apps";
} from "@trigger.dev/core/v3/isomorphic";
import { Prisma, prisma } from "~/db.server";
import { AuthenticatedEnvironment } from "~/services/apiAuth.server";
+11
View File
@@ -43,6 +43,17 @@ function createRunEngine() {
enableAutoPipelining: true,
...(env.RUN_ENGINE_RUN_QUEUE_REDIS_TLS_DISABLED === "true" ? {} : { tls: {} }),
},
queueSelectionStrategyOptions: {
parentQueueLimit: env.RUN_ENGINE_PARENT_QUEUE_LIMIT,
biases: {
concurrencyLimitBias: env.RUN_ENGINE_CONCURRENCY_LIMIT_BIAS,
availableCapacityBias: env.RUN_ENGINE_AVAILABLE_CAPACITY_BIAS,
queueAgeRandomization: env.RUN_ENGINE_QUEUE_AGE_RANDOMIZATION_BIAS,
},
reuseSnapshotCount: env.RUN_ENGINE_REUSE_SNAPSHOT_COUNT,
maximumEnvCount: env.RUN_ENGINE_MAXIMUM_ENV_COUNT,
tracer,
},
},
runLock: {
redis: {
@@ -12,7 +12,7 @@ import { reportInvocationUsage } from "~/services/platform.v3.server";
import { roomFromFriendlyRunId, socketIo } from "./handleSocketIo.server";
import { engine } from "./runEngine.server";
import { PerformTaskRunAlertsService } from "./services/alerts/performTaskRunAlerts.server";
import { RunId } from "@trigger.dev/core/v3/apps";
import { RunId } from "@trigger.dev/core/v3/isomorphic";
import { updateMetadataService } from "~/services/metadata/updateMetadata.server";
import { findEnvironmentFromRun } from "~/models/runtimeEnvironment.server";
import { env } from "~/env.server";
@@ -6,7 +6,7 @@ import {
packetRequiresOffloading,
parsePacket,
} from "@trigger.dev/core/v3";
import { BatchId, RunId } from "@trigger.dev/core/v3/apps";
import { BatchId, RunId } from "@trigger.dev/core/v3/isomorphic";
import { BatchTaskRun, Prisma } from "@trigger.dev/database";
import { z } from "zod";
import { $transaction, prisma, PrismaClientOrTransaction } from "~/db.server";
@@ -2,7 +2,7 @@ import { WorkerDeployment } from "@trigger.dev/database";
import { BaseService, ServiceValidationError } from "./baseService.server";
import { ExecuteTasksWaitingForDeployService } from "./executeTasksWaitingForDeploy";
import { compareDeploymentVersions } from "../utils/deploymentVersions";
import { CURRENT_DEPLOYMENT_LABEL } from "@trigger.dev/core/v3/apps";
import { CURRENT_DEPLOYMENT_LABEL } from "@trigger.dev/core/v3/isomorphic";
export type ChangeCurrentDeploymentDirection = "promote" | "rollback";
@@ -21,7 +21,7 @@ import {
updateEnvConcurrencyLimits,
updateQueueConcurrencyLimits,
} from "../runQueue.server";
import { BackgroundWorkerId } from "@trigger.dev/core/v3/apps";
import { BackgroundWorkerId } from "@trigger.dev/core/v3/isomorphic";
import { sanitizeQueueName } from "~/models/taskQueue.server";
export class CreateBackgroundWorkerService extends BaseService {
@@ -8,7 +8,7 @@ import { BaseService } from "./baseService.server";
import { CreateCheckpointRestoreEventService } from "./createCheckpointRestoreEvent.server";
import { ResumeBatchRunService } from "./resumeBatchRun.server";
import { ResumeDependentParentsService } from "./resumeDependentParents.server";
import { CheckpointId } from "@trigger.dev/core/v3/apps";
import { CheckpointId } from "@trigger.dev/core/v3/isomorphic";
export class CreateCheckpointService extends BaseService {
public async call(
@@ -1,6 +1,5 @@
import { CreateBackgroundWorkerRequestBody } from "@trigger.dev/core/v3";
import type { BackgroundWorker } from "@trigger.dev/database";
import { CURRENT_DEPLOYMENT_LABEL } from "@trigger.dev/core/v3/apps";
import { AuthenticatedEnvironment } from "~/services/apiAuth.server";
import { logger } from "~/services/logger.server";
import { socketIo } from "../handleSocketIo.server";
@@ -11,7 +10,7 @@ import { createBackgroundTasks, syncDeclarativeSchedules } from "./createBackgro
import { ExecuteTasksWaitingForDeployService } from "./executeTasksWaitingForDeploy";
import { projectPubSub } from "./projectPubSub.server";
import { TimeoutDeploymentService } from "./timeoutDeployment.server";
import { BackgroundWorkerId } from "@trigger.dev/core/v3/apps";
import { CURRENT_DEPLOYMENT_LABEL, BackgroundWorkerId } from "@trigger.dev/core/v3/isomorphic";
export class CreateDeployedBackgroundWorkerService extends BaseService {
public async call(
@@ -9,7 +9,7 @@ import {
syncDeclarativeSchedules,
} from "./createBackgroundWorker.server";
import { TimeoutDeploymentService } from "./timeoutDeployment.server";
import { BackgroundWorkerId } from "@trigger.dev/core/v3/apps";
import { BackgroundWorkerId } from "@trigger.dev/core/v3/isomorphic";
export class CreateDeploymentBackgroundWorkerService extends BaseService {
public async call(
@@ -1,4 +1,4 @@
import { parseNaturalLanguageDuration } from "@trigger.dev/core/v3/apps";
import { parseNaturalLanguageDuration } from "@trigger.dev/core/v3/isomorphic";
import { $transaction } from "~/db.server";
import { logger } from "~/services/logger.server";
import { marqs } from "~/v3/marqs/index.server";
@@ -1,5 +1,4 @@
import { FinalizeDeploymentRequestBody } from "@trigger.dev/core/v3/schemas";
import { CURRENT_DEPLOYMENT_LABEL } from "@trigger.dev/core/v3/apps";
import { AuthenticatedEnvironment } from "~/services/apiAuth.server";
import { logger } from "~/services/logger.server";
import { socketIo } from "../handleSocketIo.server";
@@ -11,7 +11,7 @@ import {
parseNaturalLanguageDuration,
sanitizeQueueName,
stringifyDuration,
} from "@trigger.dev/core/v3/apps";
} from "@trigger.dev/core/v3/isomorphic";
import { Prisma } from "@trigger.dev/database";
import { env } from "~/env.server";
import { createTag, MAX_TAGS_PER_RUN } from "~/models/taskRunTag.server";
@@ -6,7 +6,12 @@ import {
SemanticInternalAttributes,
TriggerTaskRequestBody,
} from "@trigger.dev/core/v3";
import { BatchId, RunId, sanitizeQueueName, stringifyDuration } from "@trigger.dev/core/v3/apps";
import {
BatchId,
RunId,
sanitizeQueueName,
stringifyDuration,
} from "@trigger.dev/core/v3/isomorphic";
import { Prisma, TaskRun } from "@trigger.dev/database";
import { env } from "~/env.server";
import { createTag, MAX_TAGS_PER_RUN } from "~/models/taskRunTag.server";
@@ -24,7 +24,10 @@ import { env } from "~/env.server";
import { $transaction } from "~/db.server";
import { resolveVariablesForEnvironment } from "~/v3/environmentVariables/environmentVariablesRepository.server";
import { generateJWTTokenForEnvironment } from "~/services/apiAuth.server";
import { CURRENT_UNMANAGED_DEPLOYMENT_LABEL, fromFriendlyId } from "@trigger.dev/core/v3/apps";
import {
CURRENT_UNMANAGED_DEPLOYMENT_LABEL,
fromFriendlyId,
} from "@trigger.dev/core/v3/isomorphic";
import { machinePresetFromName } from "~/v3/machinePresets.server";
import { defaultMachine } from "@trigger.dev/platform/v3";
+2 -2
View File
@@ -14,7 +14,7 @@
"lint": "eslint --cache --cache-location ./node_modules/.cache/eslint .",
"start": "cross-env NODE_ENV=production node --max-old-space-size=8192 ./build/server.js",
"start:local": "cross-env node --max-old-space-size=8192 ./build/server.js",
"typecheck": "tsc -p ./tsconfig.check.json",
"typecheck": "tsc --noEmit -p ./tsconfig.check.json",
"db:seed": "node prisma/seed.js",
"db:seed:local": "ts-node prisma/seed.ts",
"build:db:populate": "esbuild --platform=node --bundle --minify --format=cjs ./prisma/populate.ts --outdir=prisma",
@@ -259,4 +259,4 @@
"engines": {
"node": ">=16.0.0"
}
}
}
+2 -8
View File
@@ -9,15 +9,12 @@ module.exports = {
serverModuleFormat: "cjs",
serverDependenciesToBundle: [
/^remix-utils.*/,
/^@internal\//, // Bundle all internal packages
/^@trigger\.dev\//, // Bundle all trigger packages
"marked",
"axios",
"@internal/redis-worker",
"p-limit",
"yocto-queue",
"@trigger.dev/core",
"@trigger.dev/sdk",
"@trigger.dev/platform",
"@trigger.dev/yalt",
"@unkey/cache",
"@unkey/cache/stores",
"emails",
@@ -28,7 +25,4 @@ module.exports = {
"prismjs/components/prism-typescript",
],
browserNodeBuiltinsPolyfill: { modules: { path: true, os: true, crypto: true } },
watchPaths: async () => {
return ["../../packages/core/src/**/*", "../../packages/emails/src/**/*"];
},
};
+2 -1
View File
@@ -5,6 +5,7 @@
"paths": {
"~/*": ["./app/*"],
"@/*": ["./*"]
}
},
"customConditions": []
}
}
+3 -22
View File
@@ -20,28 +20,9 @@
"baseUrl": ".",
"paths": {
"~/*": ["./app/*"],
"@/*": ["./*"],
"@trigger.dev/sdk": ["../../packages/trigger-sdk/src/index"],
"@trigger.dev/sdk/*": ["../../packages/trigger-sdk/src/*"],
"@trigger.dev/core": ["../../packages/core/src/index"],
"@trigger.dev/core/*": ["../../packages/core/src/*"],
"@trigger.dev/database": ["../../internal-packages/database/src/index"],
"@trigger.dev/database/*": ["../../internal-packages/database/src/*"],
"@trigger.dev/yalt": ["../../packages/yalt/src/index"],
"@trigger.dev/yalt/*": ["../../packages/yalt/src/*"],
"@trigger.dev/otlp-importer": ["../../internal-packages/otlp-importer/src/index"],
"@trigger.dev/otlp-importer/*": ["../../internal-packages/otlp-importer/src/*"],
"emails": ["../../internal-packages/emails/src/index"],
"emails/*": ["../../internal-packages/emails/src/*"],
"@internal/zod-worker": ["../../internal-packages/zod-worker/src/index"],
"@internal/zod-worker/*": ["../../internal-packages/zod-worker/src/*"],
"@internal/run-engine": ["../../internal-packages/run-engine/src/index"],
"@internal/run-engine/*": ["../../internal-packages/run-engine/src/*"],
"@internal/redis-worker": ["../../internal-packages/redis-worker/src/index"],
"@internal/redis-worker/*": ["../../internal-packages/redis-worker/src/*"],
"@internal/redis": ["../../internal-packages/redis/src/index"],
"@internal/redis/*": ["../../internal-packages/redis/src/*"]
"@/*": ["./*"]
},
"noEmit": true
"noEmit": true,
"customConditions": ["@triggerdotdev/source"]
}
}
+9 -5
View File
@@ -2,21 +2,25 @@
"name": "@trigger.dev/database",
"private": true,
"version": "0.0.2",
"main": "./src/index.ts",
"types": "./src/index.ts",
"main": "./dist/index.js",
"types": "./dist/index.d.ts",
"dependencies": {
"@prisma/client": "5.4.1"
},
"devDependencies": {
"prisma": "5.4.1"
"prisma": "5.4.1",
"rimraf": "6.0.1"
},
"scripts": {
"clean": "rimraf dist",
"generate": "prisma generate",
"db:migrate:dev:create": "prisma migrate dev --create-only",
"db:migrate:deploy": "prisma migrate deploy",
"db:push": "prisma db push",
"db:studio": "prisma studio",
"db:reset": "prisma migrate reset",
"typecheck": "tsc --noEmit"
"typecheck": "tsc --noEmit",
"build": "pnpm run clean && tsc --noEmit false --outDir dist --declaration",
"dev": "tsc --noEmit false --outDir dist --declaration --watch"
}
}
}
+19 -8
View File
@@ -2,14 +2,21 @@
"name": "@internal/redis-worker",
"private": true,
"version": "0.0.1",
"main": "./src/index.ts",
"types": "./src/index.ts",
"main": "./dist/src/index.js",
"types": "./dist/src/index.d.ts",
"type": "module",
"exports": {
".": {
"@triggerdotdev/source": "./src/index.ts",
"import": "./dist/src/index.js",
"types": "./dist/src/index.d.ts",
"default": "./dist/src/index.js"
}
},
"dependencies": {
"@opentelemetry/api": "^1.9.0",
"@internal/tracing": "workspace:*",
"@internal/redis": "workspace:*",
"@trigger.dev/core": "workspace:*",
"ioredis": "^5.3.2",
"lodash.omit": "^4.5.0",
"nanoid": "^5.0.7",
"p-limit": "^6.2.0",
@@ -18,10 +25,14 @@
"devDependencies": {
"@internal/testcontainers": "workspace:*",
"@types/lodash.omit": "^4.5.7",
"vitest": "^1.4.0"
"vitest": "^1.4.0",
"rimraf": "6.0.1"
},
"scripts": {
"typecheck": "tsc --noEmit",
"test": "vitest --no-file-parallelism"
"clean": "rimraf dist",
"typecheck": "tsc --noEmit -p tsconfig.build.json",
"test": "vitest --sequence.concurrent=false --no-file-parallelism",
"build": "pnpm run clean && tsc -p tsconfig.build.json",
"dev": "tsc --watch -p tsconfig.build.json"
}
}
}
+2 -2
View File
@@ -1,2 +1,2 @@
export * from "./queue";
export * from "./worker";
export * from "./queue.js";
export * from "./worker.js";
+8 -3
View File
@@ -1,6 +1,11 @@
import { createRedisClient } from "@internal/redis";
import {
createRedisClient,
type Redis,
type Callback,
type RedisOptions,
type Result,
} from "@internal/redis";
import { Logger } from "@trigger.dev/core/logger";
import Redis, { type Callback, type RedisOptions, type Result } from "ioredis";
import { nanoid } from "nanoid";
import { z } from "zod";
@@ -436,7 +441,7 @@ export class SimpleQueue<TMessageCatalog extends MessageCatalogSchema> {
}
}
declare module "ioredis" {
declare module "@internal/redis" {
interface RedisCommander<Context> {
enqueueItem(
//keys
@@ -1,31 +0,0 @@
import { SpanOptions, SpanStatusCode, Span, Tracer } from "@opentelemetry/api";
export async function startSpan<T>(
tracer: Tracer,
name: string,
fn: (span: Span) => Promise<T>,
options?: SpanOptions
): Promise<T> {
return tracer.startActiveSpan(name, options ?? {}, async (span) => {
try {
return await fn(span);
} catch (error) {
if (error instanceof Error) {
span.recordException(error);
} else if (typeof error === "string") {
span.recordException(new Error(error));
} else {
span.recordException(new Error(String(error)));
}
span.setStatus({
code: SpanStatusCode.ERROR,
message: error instanceof Error ? error.message : String(error),
});
throw error;
} finally {
span.end();
}
});
}
@@ -4,7 +4,6 @@ import { describe } from "node:test";
import { expect } from "vitest";
import { z } from "zod";
import { Worker } from "./worker.js";
import Redis from "ioredis";
import { createRedisClient } from "@internal/redis";
describe("Worker", () => {
+2 -4
View File
@@ -1,13 +1,11 @@
import { SpanKind, trace, Tracer } from "@opentelemetry/api";
import { SpanKind, startSpan, trace, Tracer } from "@internal/tracing";
import { Logger } from "@trigger.dev/core/logger";
import { calculateNextRetryDelay } from "@trigger.dev/core/v3";
import { type RetryOptions } from "@trigger.dev/core/v3/schemas";
import { type RedisOptions } from "ioredis";
import { Redis, type RedisOptions } from "@internal/redis";
import { z } from "zod";
import { AnyQueueItem, SimpleQueue } from "./queue.js";
import Redis from "ioredis";
import { nanoid } from "nanoid";
import { startSpan } from "./telemetry.js";
import pLimit from "p-limit";
import { createRedisClient } from "@internal/redis";
@@ -0,0 +1,21 @@
{
"include": ["src/**/*.ts"],
"exclude": ["src/**/*.test.ts"],
"compilerOptions": {
"composite": true,
"target": "ES2019",
"lib": ["ES2019", "DOM", "DOM.Iterable", "DOM.AsyncIterable"],
"outDir": "dist",
"module": "Node16",
"moduleResolution": "Node16",
"moduleDetection": "force",
"verbatimModuleSyntax": false,
"esModuleInterop": true,
"forceConsistentCasingInFileNames": true,
"isolatedModules": true,
"preserveWatchOutput": true,
"skipLibCheck": true,
"strict": true,
"declaration": true
}
}
+5 -24
View File
@@ -1,27 +1,8 @@
{
"references": [{ "path": "./tsconfig.src.json" }, { "path": "./tsconfig.test.json" }],
"compilerOptions": {
"target": "ES2019",
"lib": ["ES2019", "DOM", "DOM.Iterable", "DOM.AsyncIterable"],
"module": "CommonJS",
"moduleResolution": "Node",
"moduleDetection": "force",
"verbatimModuleSyntax": false,
"types": ["vitest/globals"],
"esModuleInterop": true,
"forceConsistentCasingInFileNames": true,
"isolatedModules": true,
"preserveWatchOutput": true,
"skipLibCheck": true,
"noEmit": true,
"strict": true,
"paths": {
"@internal/testcontainers": ["../../internal-packages/testcontainers/src/index"],
"@internal/testcontainers/*": ["../../internal-packages/testcontainers/src/*"],
"@trigger.dev/core": ["../../packages/core/src/index"],
"@trigger.dev/core/*": ["../../packages/core/src/*"],
"@internal/redis": ["../../internal-packages/redis/src/index"],
"@internal/redis/*": ["../../internal-packages/redis/src/*"]
}
},
"exclude": ["node_modules"]
"moduleResolution": "Node16",
"module": "Node16",
"customConditions": ["@triggerdotdev/source"]
}
}
@@ -0,0 +1,19 @@
{
"include": ["src/**/*.ts"],
"exclude": ["node_modules", "src/**/*.test.ts"],
"compilerOptions": {
"composite": true,
"target": "ES2019",
"lib": ["ES2019", "DOM", "DOM.Iterable", "DOM.AsyncIterable"],
"module": "Node16",
"moduleResolution": "Node16",
"moduleDetection": "force",
"verbatimModuleSyntax": false,
"esModuleInterop": true,
"forceConsistentCasingInFileNames": true,
"isolatedModules": true,
"preserveWatchOutput": true,
"skipLibCheck": true,
"strict": true
}
}
@@ -0,0 +1,20 @@
{
"include": ["src/**/*.test.ts"],
"references": [{ "path": "./tsconfig.src.json" }],
"compilerOptions": {
"composite": true,
"target": "ES2019",
"lib": ["ES2019", "DOM", "DOM.Iterable", "DOM.AsyncIterable"],
"module": "Node16",
"moduleResolution": "Node16",
"moduleDetection": "force",
"verbatimModuleSyntax": false,
"types": ["vitest/globals"],
"esModuleInterop": true,
"forceConsistentCasingInFileNames": true,
"isolatedModules": true,
"preserveWatchOutput": true,
"skipLibCheck": true,
"strict": true
}
}
+1 -1
View File
@@ -15,4 +15,4 @@
"scripts": {
"typecheck": "tsc --noEmit"
}
}
}
+2
View File
@@ -1,6 +1,8 @@
import { Redis, RedisOptions } from "ioredis";
import { Logger } from "@trigger.dev/core/logger";
export { Redis, type Callback, type RedisOptions, type Result, type RedisCommander } from "ioredis";
const defaultOptions: Partial<RedisOptions> = {
retryStrategy: (times: number) => {
const delay = Math.min(times * 50, 1000);
+2 -2
View File
@@ -2,8 +2,8 @@
"compilerOptions": {
"target": "ES2019",
"lib": ["ES2019", "DOM", "DOM.Iterable", "DOM.AsyncIterable"],
"module": "CommonJS",
"moduleResolution": "Node",
"module": "Node16",
"moduleResolution": "Node16",
"moduleDetection": "force",
"verbatimModuleSyntax": false,
"types": ["vitest/globals"],
+24 -10
View File
@@ -2,27 +2,41 @@
"name": "@internal/run-engine",
"private": true,
"version": "0.0.1",
"main": "./src/index.ts",
"types": "./src/index.ts",
"main": "./dist/src/index.js",
"types": "./dist/src/index.d.ts",
"type": "module",
"exports": {
".": {
"@triggerdotdev/source": "./src/index.ts",
"import": "./dist/src/index.js",
"types": "./dist/src/index.d.ts",
"default": "./dist/src/index.js"
}
},
"dependencies": {
"@internal/redis": "workspace:*",
"@internal/redis-worker": "workspace:*",
"@opentelemetry/api": "^1.9.0",
"@opentelemetry/semantic-conventions": "^1.27.0",
"@internal/tracing": "workspace:*",
"@trigger.dev/core": "workspace:*",
"@trigger.dev/database": "workspace:*",
"assert-never": "^1.2.1",
"ioredis": "^5.3.2",
"nanoid": "^3.3.4",
"redlock": "5.0.0-beta.2",
"zod": "3.23.8"
"zod": "3.23.8",
"@unkey/cache": "^1.5.0",
"seedrandom": "^3.0.5"
},
"devDependencies": {
"@internal/testcontainers": "workspace:*",
"vitest": "^1.4.0"
"vitest": "^1.4.0",
"@types/seedrandom": "^3.0.8",
"rimraf": "6.0.1"
},
"scripts": {
"typecheck": "tsc --noEmit",
"test": "vitest --sequence.concurrent=false"
"clean": "rimraf dist",
"typecheck": "tsc --noEmit -p tsconfig.build.json",
"test": "vitest --sequence.concurrent=false --no-file-parallelism",
"build": "pnpm run clean && tsc -p tsconfig.build.json",
"dev": "tsc --watch -p tsconfig.build.json"
}
}
}
@@ -5,7 +5,7 @@ import {
PrismaClientOrTransaction,
WorkerDeployment,
} from "@trigger.dev/database";
import { CURRENT_DEPLOYMENT_LABEL } from "@trigger.dev/core/v3/apps";
import { CURRENT_DEPLOYMENT_LABEL } from "@trigger.dev/core/v3/isomorphic";
type RunWithMininimalEnvironment = Prisma.TaskRunGetPayload<{
include: {
@@ -1,5 +1,5 @@
import { TaskRunExecutionStatus, TaskRunStatus } from "@trigger.dev/database";
import { AuthenticatedEnvironment } from "../shared";
import { AuthenticatedEnvironment } from "../shared/index.js";
import { FlushedRunMetadata, TaskRunError } from "@trigger.dev/core/v3";
export type EventBusEvents = {
@@ -1,5 +1,5 @@
import { CompletedWaitpoint, ExecutionResult } from "@trigger.dev/core/v3";
import { BatchId, RunId, SnapshotId } from "@trigger.dev/core/v3/apps";
import { BatchId, RunId, SnapshotId } from "@trigger.dev/core/v3/isomorphic";
import {
PrismaClientOrTransaction,
TaskRunCheckpoint,
@@ -1,6 +1,6 @@
import { createRedisClient } from "@internal/redis";
import { createRedisClient, Redis } from "@internal/redis";
import { Worker } from "@internal/redis-worker";
import { Attributes, Span, SpanKind, trace, Tracer } from "@opentelemetry/api";
import { Attributes, Span, SpanKind, trace, Tracer } from "@internal/tracing";
import { assertExhaustive } from "@trigger.dev/core";
import { Logger } from "@trigger.dev/core/logger";
import {
@@ -35,7 +35,7 @@ import {
sanitizeQueueName,
SnapshotId,
WaitpointId,
} from "@trigger.dev/core/v3/apps";
} from "@trigger.dev/core/v3/isomorphic";
import {
$transaction,
Prisma,
@@ -48,30 +48,30 @@ import {
TaskRunStatus,
Waitpoint,
} from "@trigger.dev/database";
import assertNever from "assert-never";
import { Redis } from "ioredis";
import { assertNever } from "assert-never";
import { nanoid } from "nanoid";
import { EventEmitter } from "node:events";
import { z } from "zod";
import { RunQueue } from "../run-queue";
import { SimpleWeightedChoiceStrategy } from "../run-queue/simpleWeightedPriorityStrategy";
import { MinimalAuthenticatedEnvironment } from "../shared";
import { MAX_TASK_RUN_ATTEMPTS } from "./consts";
import { getRunWithBackgroundWorkerTasks } from "./db/worker";
import { runStatusFromError } from "./errors";
import { EventBusEvents } from "./eventBus";
import { executionResultFromSnapshot, getLatestExecutionSnapshot } from "./executionSnapshots";
import { RunLocker } from "./locking";
import { getMachinePreset } from "./machinePresets";
import { RunQueue } from "../run-queue/index.js";
import { FairQueueSelectionStrategy } from "../run-queue/fairQueueSelectionStrategy.js";
import { MinimalAuthenticatedEnvironment } from "../shared/index.js";
import { MAX_TASK_RUN_ATTEMPTS } from "./consts.js";
import { getRunWithBackgroundWorkerTasks } from "./db/worker.js";
import { runStatusFromError } from "./errors.js";
import { EventBusEvents } from "./eventBus.js";
import { executionResultFromSnapshot, getLatestExecutionSnapshot } from "./executionSnapshots.js";
import { RunLocker } from "./locking.js";
import { getMachinePreset } from "./machinePresets.js";
import {
isCheckpointable,
isDequeueableExecutionStatus,
isExecuting,
isFinalRunStatus,
isPendingExecuting,
} from "./statuses";
import { HeartbeatTimeouts, RunEngineOptions, TriggerParams } from "./types";
import { retryOutcomeFromCompletion } from "./retrying";
} from "./statuses.js";
import { HeartbeatTimeouts, RunEngineOptions, TriggerParams } from "./types.js";
import { RunQueueFullKeyProducer } from "../run-queue/keyProducer.js";
import { retryOutcomeFromCompletion } from "./retrying.js";
const workerCatalog = {
finishWaitpoint: {
@@ -153,11 +153,17 @@ export class RunEngine {
);
this.runLock = new RunLocker({ redis: this.runLockRedis });
const keys = new RunQueueFullKeyProducer();
this.runQueue = new RunQueue({
name: "rq",
tracer: trace.getTracer("rq"),
queuePriorityStrategy: new SimpleWeightedChoiceStrategy({ queueSelectionCount: 36 }),
envQueuePriorityStrategy: new SimpleWeightedChoiceStrategy({ queueSelectionCount: 12 }),
keys,
queueSelectionStrategy: new FairQueueSelectionStrategy({
keys,
redis: { ...options.queue.redis, keyPrefix: `${options.queue.redis.keyPrefix}runqueue:` },
defaultEnvConcurrencyLimit: options.queue?.defaultEnvConcurrency ?? 10,
}),
defaultEnvConcurrency: options.queue?.defaultEnvConcurrency ?? 10,
logger: new Logger("RunQueue", "debug"),
redis: { ...options.queue.redis, keyPrefix: `${options.queue.redis.keyPrefix}runqueue:` },
@@ -1,14 +1,16 @@
import Redis from "ioredis";
import Redlock, { RedlockAbortSignal } from "redlock";
// import { default: Redlock } from "redlock";
const { default: Redlock } = require("redlock");
import { AsyncLocalStorage } from "async_hooks";
import { Redis } from "@internal/redis";
import * as redlock from "redlock";
interface LockContext {
resources: string;
signal: RedlockAbortSignal;
signal: redlock.RedlockAbortSignal;
}
export class RunLocker {
private redlock: Redlock;
private redlock: InstanceType<typeof redlock.default>;
private asyncLocalStorage: AsyncLocalStorage<LockContext>;
constructor(options: { redis: Redis }) {
@@ -26,7 +28,7 @@ export class RunLocker {
async lock<T>(
resources: string[],
duration: number,
routine: (signal: RedlockAbortSignal) => Promise<T>
routine: (signal: redlock.RedlockAbortSignal) => Promise<T>
): Promise<T> {
const currentContext = this.asyncLocalStorage.getStore();
const joinedResources = resources.sort().join(",");
@@ -1,16 +1,16 @@
import {
calculateNextRetryDelay,
isOOMRunError,
RetryOptions,
sanitizeError,
shouldRetryError,
TaskRunError,
TaskRunExecutionRetry,
taskRunErrorEnhancer,
sanitizeError,
calculateNextRetryDelay,
TaskRunExecutionRetry,
} from "@trigger.dev/core/v3";
import { PrismaClientOrTransaction, TaskRunStatus } from "@trigger.dev/database";
import { MAX_TASK_RUN_ATTEMPTS } from "./consts";
import { ServiceValidationError } from ".";
import { PrismaClientOrTransaction } from "@trigger.dev/database";
import { MAX_TASK_RUN_ATTEMPTS } from "./consts.js";
import { ServiceValidationError } from "./index.js";
type Params = {
runId: string;
@@ -4,11 +4,10 @@ import {
setupAuthenticatedEnvironment,
setupBackgroundWorker,
} from "@internal/testcontainers";
import { trace } from "@opentelemetry/api";
import { expect } from "vitest";
import { EventBusEventArgs } from "../eventBus.js";
import { RunEngine } from "../index.js";
import { trace } from "@internal/tracing";
import { setTimeout } from "node:timers/promises";
import { expect } from "vitest";
import { RunEngine } from "../index.js";
describe("RunEngine attempt failures", () => {
containerTest(
@@ -3,180 +3,178 @@ import {
setupAuthenticatedEnvironment,
setupBackgroundWorker,
} from "@internal/testcontainers";
import { trace } from "@opentelemetry/api";
import { generateFriendlyId } from "@trigger.dev/core/v3/apps";
import { trace } from "@internal/tracing";
import { generateFriendlyId } from "@trigger.dev/core/v3/isomorphic";
import { expect } from "vitest";
import { RunEngine } from "../index.js";
import { setTimeout } from "node:timers/promises";
describe("RunEngine batchTrigger", () => {
containerTest(
"Batch trigger shares a batch",
{ timeout: 15_000 },
async ({ prisma, redisOptions }) => {
//create environment
const authenticatedEnvironment = await setupAuthenticatedEnvironment(prisma, "PRODUCTION");
vi.setConfig({ testTimeout: 60_000 });
const engine = new RunEngine({
prisma,
worker: {
redis: redisOptions,
workers: 1,
tasksPerWorker: 10,
pollIntervalMs: 100,
},
queue: {
redis: redisOptions,
},
runLock: {
redis: redisOptions,
},
describe("RunEngine batchTrigger", () => {
containerTest("Batch trigger shares a batch", async ({ prisma, redisOptions }) => {
//create environment
const authenticatedEnvironment = await setupAuthenticatedEnvironment(prisma, "PRODUCTION");
const engine = new RunEngine({
prisma,
worker: {
redis: redisOptions,
workers: 1,
tasksPerWorker: 10,
pollIntervalMs: 100,
},
queue: {
redis: redisOptions,
},
runLock: {
redis: redisOptions,
},
machines: {
defaultMachine: "small-1x",
machines: {
defaultMachine: "small-1x",
machines: {
"small-1x": {
name: "small-1x" as const,
cpu: 0.5,
memory: 0.5,
centsPerMs: 0.0001,
},
"small-1x": {
name: "small-1x" as const,
cpu: 0.5,
memory: 0.5,
centsPerMs: 0.0001,
},
baseCostInCents: 0.0005,
},
tracer: trace.getTracer("test", "0.0.0"),
baseCostInCents: 0.0005,
},
tracer: trace.getTracer("test", "0.0.0"),
});
try {
const taskIdentifier = "test-task";
//create background worker
const backgroundWorker = await setupBackgroundWorker(
prisma,
authenticatedEnvironment,
taskIdentifier
);
const batch = await prisma.batchTaskRun.create({
data: {
friendlyId: generateFriendlyId("batch"),
runtimeEnvironmentId: authenticatedEnvironment.id,
},
});
try {
const taskIdentifier = "test-task";
//trigger the runs
const run1 = await engine.trigger(
{
number: 1,
friendlyId: "run_1234",
environment: authenticatedEnvironment,
taskIdentifier,
payload: "{}",
payloadType: "application/json",
context: {},
traceContext: {},
traceId: "t12345",
spanId: "s12345",
masterQueue: "main",
queueName: "task/test-task",
isTest: false,
tags: [],
batch: { id: batch.id, index: 0 },
},
prisma
);
//create background worker
const backgroundWorker = await setupBackgroundWorker(
prisma,
authenticatedEnvironment,
taskIdentifier
);
const run2 = await engine.trigger(
{
number: 2,
friendlyId: "run_1235",
environment: authenticatedEnvironment,
taskIdentifier,
payload: "{}",
payloadType: "application/json",
context: {},
traceContext: {},
traceId: "t12345",
spanId: "s12345",
masterQueue: "main",
queueName: "task/test-task",
isTest: false,
tags: [],
batch: { id: batch.id, index: 1 },
},
prisma
);
const batch = await prisma.batchTaskRun.create({
data: {
friendlyId: generateFriendlyId("batch"),
runtimeEnvironmentId: authenticatedEnvironment.id,
},
});
expect(run1).toBeDefined();
expect(run1.friendlyId).toBe("run_1234");
expect(run1.batchId).toBe(batch.id);
//trigger the runs
const run1 = await engine.trigger(
{
number: 1,
friendlyId: "run_1234",
environment: authenticatedEnvironment,
taskIdentifier,
payload: "{}",
payloadType: "application/json",
context: {},
traceContext: {},
traceId: "t12345",
spanId: "s12345",
masterQueue: "main",
queueName: "task/test-task",
isTest: false,
tags: [],
batch: { id: batch.id, index: 0 },
},
prisma
);
expect(run2).toBeDefined();
expect(run2.friendlyId).toBe("run_1235");
expect(run2.batchId).toBe(batch.id);
const run2 = await engine.trigger(
{
number: 2,
friendlyId: "run_1235",
environment: authenticatedEnvironment,
taskIdentifier,
payload: "{}",
payloadType: "application/json",
context: {},
traceContext: {},
traceId: "t12345",
spanId: "s12345",
masterQueue: "main",
queueName: "task/test-task",
isTest: false,
tags: [],
batch: { id: batch.id, index: 1 },
},
prisma
);
//check the queue length
const queueLength = await engine.runQueue.lengthOfEnvQueue(authenticatedEnvironment);
expect(queueLength).toBe(2);
expect(run1).toBeDefined();
expect(run1.friendlyId).toBe("run_1234");
expect(run1.batchId).toBe(batch.id);
//dequeue
const [d1, d2] = await engine.dequeueFromMasterQueue({
consumerId: "test_12345",
masterQueue: run1.masterQueue,
maxRunCount: 10,
});
expect(run2).toBeDefined();
expect(run2.friendlyId).toBe("run_1235");
expect(run2.batchId).toBe(batch.id);
//attempts
const attempt1 = await engine.startRunAttempt({
runId: d1.run.id,
snapshotId: d1.snapshot.id,
});
const attempt2 = await engine.startRunAttempt({
runId: d2.run.id,
snapshotId: d2.snapshot.id,
});
//check the queue length
const queueLength = await engine.runQueue.lengthOfEnvQueue(authenticatedEnvironment);
expect(queueLength).toBe(2);
//complete the runs
const result1 = await engine.completeRunAttempt({
runId: attempt1.run.id,
snapshotId: attempt1.snapshot.id,
completion: {
ok: true,
id: attempt1.run.id,
output: `{"foo":"bar"}`,
outputType: "application/json",
},
});
const result2 = await engine.completeRunAttempt({
runId: attempt2.run.id,
snapshotId: attempt2.snapshot.id,
completion: {
ok: true,
id: attempt2.run.id,
output: `{"baz":"qux"}`,
outputType: "application/json",
},
});
//dequeue
const [d1, d2] = await engine.dequeueFromMasterQueue({
consumerId: "test_12345",
masterQueue: run1.masterQueue,
maxRunCount: 10,
});
//the batch won't complete immediately
const batchAfter1 = await prisma.batchTaskRun.findUnique({
where: {
id: batch.id,
},
});
expect(batchAfter1?.status).toBe("PENDING");
//attempts
const attempt1 = await engine.startRunAttempt({
runId: d1.run.id,
snapshotId: d1.snapshot.id,
});
const attempt2 = await engine.startRunAttempt({
runId: d2.run.id,
snapshotId: d2.snapshot.id,
});
await setTimeout(3_000);
//complete the runs
const result1 = await engine.completeRunAttempt({
runId: attempt1.run.id,
snapshotId: attempt1.snapshot.id,
completion: {
ok: true,
id: attempt1.run.id,
output: `{"foo":"bar"}`,
outputType: "application/json",
},
});
const result2 = await engine.completeRunAttempt({
runId: attempt2.run.id,
snapshotId: attempt2.snapshot.id,
completion: {
ok: true,
id: attempt2.run.id,
output: `{"baz":"qux"}`,
outputType: "application/json",
},
});
//the batch won't complete immediately
const batchAfter1 = await prisma.batchTaskRun.findUnique({
where: {
id: batch.id,
},
});
expect(batchAfter1?.status).toBe("PENDING");
await setTimeout(3_000);
//the batch should complete
const batchAfter2 = await prisma.batchTaskRun.findUnique({
where: {
id: batch.id,
},
});
expect(batchAfter2?.status).toBe("COMPLETED");
} finally {
engine.quit();
}
//the batch should complete
const batchAfter2 = await prisma.batchTaskRun.findUnique({
where: {
id: batch.id,
},
});
expect(batchAfter2?.status).toBe("COMPLETED");
} finally {
engine.quit();
}
);
});
});

Some files were not shown because too many files have changed in this diff Show More