Compare commits
37 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| d22a460555 | |||
| 9ba2a217a4 | |||
| 39885a427f | |||
| ccb0bc510a | |||
| 56d66ee07c | |||
| 4ca8887972 | |||
| 89bffc066c | |||
| 34ca7667d3 | |||
| 3e327acc0f | |||
| 77ad4127cb | |||
| 8a5076aacf | |||
| 5399f6bfb7 | |||
| ecef199660 | |||
| 4acfb8f4bb | |||
| 2ef278db67 | |||
| c7a55804d9 | |||
| da6a66efff | |||
| 7c36a1a4b0 | |||
| 225effb599 | |||
| 98eb6ed4f9 | |||
| 3069ebf0d8 | |||
| e133e628ca | |||
| 098932ea96 | |||
| 65f960e883 | |||
| ccbeff47e6 | |||
| 7c8f2df105 | |||
| fd44dabfe0 | |||
| 5daed3f69d | |||
| 596bf78e55 | |||
| 6ca66b76f4 | |||
| 29ef0395ce | |||
| 55d1f8c677 | |||
| 9835f4ec55 | |||
| 7fae10db23 | |||
| 8cf1f0a37d | |||
| dba4313c5c | |||
| 506613dc92 |
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"@trigger.dev/sdk": patch
|
||||
---
|
||||
|
||||
Use global setTimeout to ensure cross-runtime support
|
||||
@@ -99,6 +99,7 @@
|
||||
"pink-pumas-rhyme",
|
||||
"plenty-ducks-beam",
|
||||
"polite-ducks-switch",
|
||||
"polite-pears-grow",
|
||||
"polite-rockets-matter",
|
||||
"poor-flowers-cross",
|
||||
"purple-garlics-shop",
|
||||
@@ -115,14 +116,18 @@
|
||||
"sharp-emus-compare",
|
||||
"sharp-zebras-serve",
|
||||
"shiny-coats-cry",
|
||||
"silly-buses-obey",
|
||||
"silly-suits-switch",
|
||||
"silver-doors-juggle",
|
||||
"six-ligers-exist",
|
||||
"six-rats-hunt",
|
||||
"sixty-insects-watch",
|
||||
"slow-buses-own",
|
||||
"slow-kiwis-hide",
|
||||
"slow-sloths-retire",
|
||||
"smart-needles-move",
|
||||
"smart-olives-eat",
|
||||
"sour-pugs-teach",
|
||||
"spicy-lamps-smoke",
|
||||
"spicy-terms-bow",
|
||||
"strange-ghosts-matter",
|
||||
@@ -132,6 +137,7 @@
|
||||
"strong-phones-smoke",
|
||||
"stupid-adults-sniff",
|
||||
"stupid-bulldogs-applaud",
|
||||
"sweet-ducks-remember",
|
||||
"sweet-lizards-press",
|
||||
"swift-dragons-peel",
|
||||
"tall-bees-wave",
|
||||
@@ -139,6 +145,7 @@
|
||||
"tender-moose-tell",
|
||||
"tender-oranges-rhyme",
|
||||
"tender-turkeys-compete",
|
||||
"thick-carrots-sneeze",
|
||||
"thin-parents-heal",
|
||||
"thirty-islands-kiss",
|
||||
"tidy-balloons-suffer",
|
||||
@@ -149,7 +156,9 @@
|
||||
"tricky-bulldogs-heal",
|
||||
"tricky-keys-attack",
|
||||
"tricky-ladybugs-unite",
|
||||
"twelve-knives-notice",
|
||||
"two-pumas-wait",
|
||||
"violet-clocks-notice",
|
||||
"warm-olives-provide",
|
||||
"warm-planes-taste",
|
||||
"young-snails-sell"
|
||||
|
||||
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"@trigger.dev/core": patch
|
||||
---
|
||||
|
||||
Add callback to checkpoint created message
|
||||
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"trigger.dev": patch
|
||||
---
|
||||
|
||||
Improved ESM module require error detection logic
|
||||
@@ -0,0 +1,6 @@
|
||||
---
|
||||
"trigger.dev": patch
|
||||
"@trigger.dev/core": patch
|
||||
---
|
||||
|
||||
v3: vercel edge runtime support
|
||||
@@ -0,0 +1,6 @@
|
||||
---
|
||||
"trigger.dev": patch
|
||||
"@trigger.dev/core": patch
|
||||
---
|
||||
|
||||
v3: fix otel flushing causing CLEANUP ack timeout errors by always setting a forceFlushTimeoutMillis value
|
||||
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"@trigger.dev/sdk": patch
|
||||
---
|
||||
|
||||
v3: Adding SDK functions for triggering tasks in a typesafe way, without importing task file
|
||||
@@ -0,0 +1,6 @@
|
||||
---
|
||||
"@trigger.dev/core": patch
|
||||
"@trigger.dev/sdk": patch
|
||||
---
|
||||
|
||||
v3: Include presigned urls for downloading large payloads and outputs when using runs.retrieve
|
||||
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"@trigger.dev/core": patch
|
||||
---
|
||||
|
||||
v3: fix missing init output in task run function when no middleware is defined
|
||||
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"trigger.dev": patch
|
||||
---
|
||||
|
||||
Fix jsonc-parser import
|
||||
@@ -1162,13 +1162,7 @@ class TaskCoordinator {
|
||||
return;
|
||||
}
|
||||
|
||||
if (!checkpoint.docker || !willSimulate) {
|
||||
socket.emit("REQUEST_EXIT", {
|
||||
version: "v1",
|
||||
});
|
||||
}
|
||||
|
||||
this.#platformSocket?.send("CHECKPOINT_CREATED", {
|
||||
const ack = await this.#platformSocket?.sendWithAck("CHECKPOINT_CREATED", {
|
||||
version: "v1",
|
||||
attemptFriendlyId: message.attemptFriendlyId,
|
||||
docker: checkpoint.docker,
|
||||
@@ -1179,6 +1173,17 @@ class TaskCoordinator {
|
||||
now: message.now,
|
||||
},
|
||||
});
|
||||
|
||||
if (ack?.keepRunAlive) {
|
||||
logger.log("keeping run alive after duration checkpoint", { runId: socket.data.runId });
|
||||
return;
|
||||
}
|
||||
|
||||
if (!checkpoint.docker || !willSimulate) {
|
||||
socket.emit("REQUEST_EXIT", {
|
||||
version: "v1",
|
||||
});
|
||||
}
|
||||
});
|
||||
|
||||
socket.on("WAIT_FOR_TASK", async (message, callback) => {
|
||||
@@ -1205,13 +1210,7 @@ class TaskCoordinator {
|
||||
return;
|
||||
}
|
||||
|
||||
if (!checkpoint.docker || !willSimulate) {
|
||||
socket.emit("REQUEST_EXIT", {
|
||||
version: "v1",
|
||||
});
|
||||
}
|
||||
|
||||
this.#platformSocket?.send("CHECKPOINT_CREATED", {
|
||||
const ack = await this.#platformSocket?.sendWithAck("CHECKPOINT_CREATED", {
|
||||
version: "v1",
|
||||
attemptFriendlyId: message.attemptFriendlyId,
|
||||
docker: checkpoint.docker,
|
||||
@@ -1221,6 +1220,17 @@ class TaskCoordinator {
|
||||
friendlyId: message.friendlyId,
|
||||
},
|
||||
});
|
||||
|
||||
if (ack?.keepRunAlive) {
|
||||
logger.log("keeping run alive after task checkpoint", { runId: socket.data.runId });
|
||||
return;
|
||||
}
|
||||
|
||||
if (!checkpoint.docker || !willSimulate) {
|
||||
socket.emit("REQUEST_EXIT", {
|
||||
version: "v1",
|
||||
});
|
||||
}
|
||||
});
|
||||
|
||||
socket.on("WAIT_FOR_BATCH", async (message, callback) => {
|
||||
@@ -1247,13 +1257,7 @@ class TaskCoordinator {
|
||||
return;
|
||||
}
|
||||
|
||||
if (!checkpoint.docker || !willSimulate) {
|
||||
socket.emit("REQUEST_EXIT", {
|
||||
version: "v1",
|
||||
});
|
||||
}
|
||||
|
||||
this.#platformSocket?.send("CHECKPOINT_CREATED", {
|
||||
const ack = await this.#platformSocket?.sendWithAck("CHECKPOINT_CREATED", {
|
||||
version: "v1",
|
||||
attemptFriendlyId: message.attemptFriendlyId,
|
||||
docker: checkpoint.docker,
|
||||
@@ -1264,6 +1268,17 @@ class TaskCoordinator {
|
||||
runFriendlyIds: message.runFriendlyIds,
|
||||
},
|
||||
});
|
||||
|
||||
if (ack?.keepRunAlive) {
|
||||
logger.log("keeping run alive after batch checkpoint", { runId: socket.data.runId });
|
||||
return;
|
||||
}
|
||||
|
||||
if (!checkpoint.docker || !willSimulate) {
|
||||
socket.emit("REQUEST_EXIT", {
|
||||
version: "v1",
|
||||
});
|
||||
}
|
||||
});
|
||||
|
||||
socket.on("INDEX_TASKS", async (message, callback) => {
|
||||
|
||||
@@ -528,7 +528,7 @@ export type Tree<TData> = {
|
||||
/** A tree but flattened so it can easily be used for DOM elements */
|
||||
export type FlatTreeItem<TData> = {
|
||||
id: string;
|
||||
parentId: string | undefined;
|
||||
parentId?: string | undefined;
|
||||
children: string[];
|
||||
hasChildren: boolean;
|
||||
/** The indentation level, the root is 0 */
|
||||
|
||||
@@ -186,3 +186,9 @@ export { apiRateLimiter } from "./services/apiRateLimit.server";
|
||||
export { socketIo } from "./v3/handleSocketIo.server";
|
||||
export { wss } from "./v3/handleWebsockets.server";
|
||||
export { registryProxy } from "./v3/registryProxy.server";
|
||||
import { eventLoopMonitor } from "./eventLoopMonitor.server";
|
||||
import { env } from "./env.server";
|
||||
|
||||
if (env.EVENT_LOOP_MONITOR_ENABLED === "1") {
|
||||
eventLoopMonitor.enable();
|
||||
}
|
||||
|
||||
@@ -204,6 +204,11 @@ const EnvironmentSchema = z.object({
|
||||
|
||||
USAGE_OPEN_METER_API_KEY: z.string().optional(),
|
||||
USAGE_OPEN_METER_BASE_URL: z.string().optional(),
|
||||
EVENT_LOOP_MONITOR_ENABLED: z.string().default("1"),
|
||||
MAXIMUM_LIVE_RELOADING_EVENTS: z.coerce.number().int().default(1000),
|
||||
MAXIMUM_TRACE_SUMMARY_VIEW_COUNT: z.coerce.number().int().default(25_000),
|
||||
TASK_PAYLOAD_OFFLOAD_THRESHOLD: z.coerce.number().int().default(524_288), // 512KB
|
||||
TASK_PAYLOAD_MAXIMUM_SIZE: z.coerce.number().int().default(3_145_728), // 3MB
|
||||
});
|
||||
|
||||
export type Environment = z.infer<typeof EnvironmentSchema>;
|
||||
|
||||
@@ -0,0 +1,77 @@
|
||||
import { createHook } from "node:async_hooks";
|
||||
import { singleton } from "./utils/singleton";
|
||||
import { tracer } from "./v3/tracer.server";
|
||||
|
||||
const THRESHOLD_NS = 1e8; // 100ms
|
||||
|
||||
const cache = new Map<number, { type: string; start?: [number, number] }>();
|
||||
|
||||
function init(asyncId: number, type: string, triggerAsyncId: number, resource: any) {
|
||||
cache.set(asyncId, {
|
||||
type,
|
||||
});
|
||||
}
|
||||
|
||||
function destroy(asyncId: number) {
|
||||
cache.delete(asyncId);
|
||||
}
|
||||
|
||||
function before(asyncId: number) {
|
||||
const cached = cache.get(asyncId);
|
||||
|
||||
if (!cached) {
|
||||
return;
|
||||
}
|
||||
|
||||
cache.set(asyncId, {
|
||||
...cached,
|
||||
start: process.hrtime(),
|
||||
});
|
||||
}
|
||||
|
||||
function after(asyncId: number) {
|
||||
const cached = cache.get(asyncId);
|
||||
|
||||
if (!cached) {
|
||||
return;
|
||||
}
|
||||
|
||||
cache.delete(asyncId);
|
||||
|
||||
if (!cached.start) {
|
||||
return;
|
||||
}
|
||||
|
||||
const diff = process.hrtime(cached.start);
|
||||
const diffNs = diff[0] * 1e9 + diff[1];
|
||||
if (diffNs > THRESHOLD_NS) {
|
||||
const time = diffNs / 1e6; // in ms
|
||||
|
||||
const newSpan = tracer.startSpan("event-loop-blocked", {
|
||||
startTime: new Date(new Date().getTime() - time),
|
||||
attributes: {
|
||||
asyncType: cached.type,
|
||||
label: "EventLoopMonitor",
|
||||
},
|
||||
});
|
||||
|
||||
newSpan.end();
|
||||
}
|
||||
}
|
||||
|
||||
export const eventLoopMonitor = singleton("eventLoopMonitor", () => {
|
||||
const hook = createHook({ init, before, after, destroy });
|
||||
|
||||
return {
|
||||
enable: () => {
|
||||
console.log("🥸 Initializing event loop monitor");
|
||||
|
||||
hook.enable();
|
||||
},
|
||||
disable: () => {
|
||||
console.log("🥸 Disabling event loop monitor");
|
||||
|
||||
hook.disable();
|
||||
},
|
||||
};
|
||||
});
|
||||
@@ -26,7 +26,7 @@ export function useEventSource(
|
||||
const eventSource = new EventSource(url, init);
|
||||
eventSource.addEventListener(event ?? "message", handler);
|
||||
|
||||
// rest data if dependencies change
|
||||
// reset data if dependencies change
|
||||
setData(null);
|
||||
|
||||
function handler(event: MessageEvent) {
|
||||
|
||||
@@ -65,59 +65,55 @@ export async function findEnvironmentById(id: string) {
|
||||
}
|
||||
|
||||
export async function createNewSession(environment: RuntimeEnvironment, ipAddress: string) {
|
||||
return prisma.$transaction(async (tx) => {
|
||||
const session = await tx.runtimeEnvironmentSession.create({
|
||||
data: {
|
||||
environmentId: environment.id,
|
||||
ipAddress,
|
||||
},
|
||||
});
|
||||
|
||||
await tx.runtimeEnvironment.update({
|
||||
where: {
|
||||
id: environment.id,
|
||||
},
|
||||
data: {
|
||||
currentSessionId: session.id,
|
||||
},
|
||||
});
|
||||
|
||||
return session;
|
||||
const session = await prisma.runtimeEnvironmentSession.create({
|
||||
data: {
|
||||
environmentId: environment.id,
|
||||
ipAddress,
|
||||
},
|
||||
});
|
||||
|
||||
await prisma.runtimeEnvironment.update({
|
||||
where: {
|
||||
id: environment.id,
|
||||
},
|
||||
data: {
|
||||
currentSessionId: session.id,
|
||||
},
|
||||
});
|
||||
|
||||
return session;
|
||||
}
|
||||
|
||||
export async function disconnectSession(environmentId: string) {
|
||||
return prisma.$transaction(async (tx) => {
|
||||
const environment = await tx.runtimeEnvironment.findUnique({
|
||||
where: {
|
||||
id: environmentId,
|
||||
},
|
||||
});
|
||||
|
||||
if (!environment || !environment.currentSessionId) {
|
||||
return null;
|
||||
}
|
||||
|
||||
const session = await tx.runtimeEnvironmentSession.update({
|
||||
where: {
|
||||
id: environment.currentSessionId,
|
||||
},
|
||||
data: {
|
||||
disconnectedAt: new Date(),
|
||||
},
|
||||
});
|
||||
|
||||
await tx.runtimeEnvironment.update({
|
||||
where: {
|
||||
id: environment.id,
|
||||
},
|
||||
data: {
|
||||
currentSessionId: null,
|
||||
},
|
||||
});
|
||||
|
||||
return session;
|
||||
const environment = await prisma.runtimeEnvironment.findUnique({
|
||||
where: {
|
||||
id: environmentId,
|
||||
},
|
||||
});
|
||||
|
||||
if (!environment || !environment.currentSessionId) {
|
||||
return null;
|
||||
}
|
||||
|
||||
const session = await prisma.runtimeEnvironmentSession.update({
|
||||
where: {
|
||||
id: environment.currentSessionId,
|
||||
},
|
||||
data: {
|
||||
disconnectedAt: new Date(),
|
||||
},
|
||||
});
|
||||
|
||||
await prisma.runtimeEnvironment.update({
|
||||
where: {
|
||||
id: environment.id,
|
||||
},
|
||||
data: {
|
||||
currentSessionId: null,
|
||||
},
|
||||
});
|
||||
|
||||
return session;
|
||||
}
|
||||
|
||||
type DisplayableInputEnvironment = Prisma.RuntimeEnvironmentGetPayload<{
|
||||
|
||||
@@ -1,11 +1,9 @@
|
||||
import { z } from "zod";
|
||||
import {
|
||||
Direction,
|
||||
FilterableEnvironment,
|
||||
FilterableStatus,
|
||||
filterableStatuses,
|
||||
} from "~/components/runs/RunStatuses";
|
||||
import { PrismaClient, prisma } from "~/db.server";
|
||||
import { getUsername } from "~/utils/username";
|
||||
import { BasePresenter } from "./v3/basePresenter.server";
|
||||
|
||||
@@ -29,8 +27,6 @@ const DEFAULT_PAGE_SIZE = 20;
|
||||
export type RunList = Awaited<ReturnType<RunListPresenter["call"]>>;
|
||||
|
||||
export class RunListPresenter extends BasePresenter {
|
||||
|
||||
|
||||
public async call({
|
||||
userId,
|
||||
eventId,
|
||||
|
||||
@@ -1,27 +1,52 @@
|
||||
import { User } from "@trigger.dev/database";
|
||||
import { ScheduleMetadataSchema } from "@trigger.dev/core";
|
||||
import { PrismaClient, prisma } from "~/db.server";
|
||||
import { User } from "@trigger.dev/database";
|
||||
import { Organization } from "~/models/organization.server";
|
||||
import { Project } from "~/models/project.server";
|
||||
import { calculateNextScheduledEvent } from "~/services/schedules/nextScheduledEvent.server";
|
||||
import { BasePresenter } from "./v3/basePresenter.server";
|
||||
|
||||
export class ScheduledTriggersPresenter {
|
||||
#prismaClient: PrismaClient;
|
||||
|
||||
constructor(prismaClient: PrismaClient = prisma) {
|
||||
this.#prismaClient = prismaClient;
|
||||
}
|
||||
const DEFAULT_PAGE_SIZE = 20;
|
||||
|
||||
export class ScheduledTriggersPresenter extends BasePresenter {
|
||||
public async call({
|
||||
userId,
|
||||
projectSlug,
|
||||
organizationSlug,
|
||||
direction = "forward",
|
||||
pageSize = DEFAULT_PAGE_SIZE,
|
||||
cursor,
|
||||
}: {
|
||||
userId: User["id"];
|
||||
projectSlug: Project["slug"];
|
||||
organizationSlug: Organization["slug"];
|
||||
direction?: "forward" | "backward";
|
||||
pageSize?: number;
|
||||
cursor?: string;
|
||||
}) {
|
||||
const scheduled = await this.#prismaClient.scheduleSource.findMany({
|
||||
const organization = await this._replica.organization.findFirstOrThrow({
|
||||
select: {
|
||||
id: true,
|
||||
},
|
||||
where: {
|
||||
slug: organizationSlug,
|
||||
members: { some: { userId } },
|
||||
},
|
||||
});
|
||||
|
||||
// Find the project scoped to the organization
|
||||
const project = await this._replica.project.findFirstOrThrow({
|
||||
select: {
|
||||
id: true,
|
||||
},
|
||||
where: {
|
||||
slug: projectSlug,
|
||||
organizationId: organization.id,
|
||||
},
|
||||
});
|
||||
|
||||
const directionMultiplier = direction === "forward" ? 1 : -1;
|
||||
|
||||
const scheduled = await this._replica.scheduleSource.findMany({
|
||||
select: {
|
||||
id: true,
|
||||
key: true,
|
||||
@@ -50,23 +75,50 @@ export class ScheduledTriggersPresenter {
|
||||
},
|
||||
},
|
||||
],
|
||||
organization: {
|
||||
slug: organizationSlug,
|
||||
members: {
|
||||
some: {
|
||||
userId,
|
||||
},
|
||||
},
|
||||
},
|
||||
project: {
|
||||
slug: projectSlug,
|
||||
},
|
||||
projectId: project.id,
|
||||
},
|
||||
},
|
||||
orderBy: [{ id: "desc" }],
|
||||
//take an extra record to tell if there are more
|
||||
take: directionMultiplier * (pageSize + 1),
|
||||
//skip the cursor if there is one
|
||||
skip: cursor ? 1 : 0,
|
||||
cursor: cursor
|
||||
? {
|
||||
id: cursor,
|
||||
}
|
||||
: undefined,
|
||||
});
|
||||
|
||||
const hasMore = scheduled.length > pageSize;
|
||||
|
||||
//get cursors for next and previous pages
|
||||
let next: string | undefined;
|
||||
let previous: string | undefined;
|
||||
switch (direction) {
|
||||
case "forward":
|
||||
previous = cursor ? scheduled.at(0)?.id : undefined;
|
||||
if (hasMore) {
|
||||
next = scheduled[pageSize - 1]?.id;
|
||||
}
|
||||
break;
|
||||
case "backward":
|
||||
if (hasMore) {
|
||||
previous = scheduled[1]?.id;
|
||||
next = scheduled[pageSize]?.id;
|
||||
} else {
|
||||
next = scheduled[pageSize - 1]?.id;
|
||||
}
|
||||
break;
|
||||
}
|
||||
|
||||
const scheduledToReturn =
|
||||
direction === "backward" && hasMore
|
||||
? scheduled.slice(1, pageSize + 1)
|
||||
: scheduled.slice(0, pageSize);
|
||||
|
||||
return {
|
||||
scheduled: scheduled.map((s) => {
|
||||
scheduled: scheduledToReturn.map((s) => {
|
||||
const schedule = ScheduleMetadataSchema.parse(s.schedule);
|
||||
const nextEventTimestamp = s.active
|
||||
? calculateNextScheduledEvent(schedule, s.lastEventTimestamp)
|
||||
@@ -78,6 +130,10 @@ export class ScheduledTriggersPresenter {
|
||||
nextEventTimestamp,
|
||||
};
|
||||
}),
|
||||
pagination: {
|
||||
next,
|
||||
previous,
|
||||
},
|
||||
};
|
||||
}
|
||||
}
|
||||
|
||||
@@ -12,6 +12,7 @@ import {
|
||||
import { Prisma, TaskRunAttemptStatus, TaskRunStatus } from "@trigger.dev/database";
|
||||
import assertNever from "assert-never";
|
||||
import { AuthenticatedEnvironment } from "~/services/apiAuth.server";
|
||||
import { generatePresignedUrl } from "~/v3/r2.server";
|
||||
import { BasePresenter } from "./basePresenter.server";
|
||||
|
||||
export class ApiRetrieveRunPresenter extends BasePresenter {
|
||||
@@ -44,7 +45,9 @@ export class ApiRetrieveRunPresenter extends BasePresenter {
|
||||
}
|
||||
|
||||
let $payload: any;
|
||||
let $payloadPresignedUrl: string | undefined;
|
||||
let $output: any;
|
||||
let $outputPresignedUrl: string | undefined;
|
||||
|
||||
if (showSecretDetails) {
|
||||
const payloadPacket = await conditionallyImportPacket({
|
||||
@@ -52,7 +55,19 @@ export class ApiRetrieveRunPresenter extends BasePresenter {
|
||||
dataType: taskRun.payloadType,
|
||||
});
|
||||
|
||||
$payload = await parsePacket(payloadPacket);
|
||||
if (
|
||||
payloadPacket.dataType === "application/store" &&
|
||||
typeof payloadPacket.data === "string"
|
||||
) {
|
||||
$payloadPresignedUrl = await generatePresignedUrl(
|
||||
env.project.externalRef,
|
||||
env.slug,
|
||||
payloadPacket.data,
|
||||
"GET"
|
||||
);
|
||||
} else {
|
||||
$payload = await parsePacket(payloadPacket);
|
||||
}
|
||||
|
||||
if (taskRun.status === "COMPLETED_SUCCESSFULLY") {
|
||||
const completedAttempt = taskRun.attempts.find(
|
||||
@@ -65,7 +80,19 @@ export class ApiRetrieveRunPresenter extends BasePresenter {
|
||||
dataType: completedAttempt.outputType,
|
||||
});
|
||||
|
||||
$output = await parsePacket(outputPacket);
|
||||
if (
|
||||
outputPacket.dataType === "application/store" &&
|
||||
typeof outputPacket.data === "string"
|
||||
) {
|
||||
$outputPresignedUrl = await generatePresignedUrl(
|
||||
env.project.externalRef,
|
||||
env.slug,
|
||||
outputPacket.data,
|
||||
"GET"
|
||||
);
|
||||
} else {
|
||||
$output = await parsePacket(outputPacket);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -85,7 +112,9 @@ export class ApiRetrieveRunPresenter extends BasePresenter {
|
||||
? taskRun.updatedAt
|
||||
: undefined,
|
||||
payload: $payload,
|
||||
payloadPresignedUrl: $payloadPresignedUrl,
|
||||
output: $output,
|
||||
outputPresignedUrl: $outputPresignedUrl,
|
||||
isTest: taskRun.isTest,
|
||||
schedule: taskRun.schedule
|
||||
? {
|
||||
|
||||
@@ -4,8 +4,8 @@ import { Direction } from "~/components/runs/RunStatuses";
|
||||
import { FINISHED_STATUSES } from "~/components/runs/v3/TaskRunStatus";
|
||||
import { sqlDatabaseSchema } from "~/db.server";
|
||||
import { displayableEnvironment } from "~/models/runtimeEnvironment.server";
|
||||
import { CANCELLABLE_STATUSES } from "~/v3/services/cancelTaskRun.server";
|
||||
import { BasePresenter } from "./basePresenter.server";
|
||||
import { isCancellableRunStatus } from "~/v3/taskStatus";
|
||||
|
||||
export type RunListOptions = {
|
||||
userId?: string;
|
||||
@@ -291,7 +291,7 @@ export class RunListPresenter extends BasePresenter {
|
||||
taskIdentifier: run.taskIdentifier,
|
||||
spanId: run.spanId,
|
||||
isReplayable: true,
|
||||
isCancellable: CANCELLABLE_STATUSES.includes(run.status),
|
||||
isCancellable: isCancellableRunStatus(run.status),
|
||||
environment: displayableEnvironment(environment, userId),
|
||||
idempotencyKey: run.idempotencyKey ? run.idempotencyKey : undefined,
|
||||
};
|
||||
|
||||
@@ -94,17 +94,24 @@ export class RunStreamPresenter {
|
||||
|
||||
eventEmitter.removeAllListeners();
|
||||
|
||||
unsubscribe().catch((error) => {
|
||||
logger.error("RunStreamPresenter.abort.unsubscribe", {
|
||||
runFriendlyId,
|
||||
traceId: run.traceId,
|
||||
error: {
|
||||
name: error.name,
|
||||
message: error.message,
|
||||
stack: error.stack,
|
||||
},
|
||||
unsubscribe()
|
||||
.then(() => {
|
||||
logger.info("RunStreamPresenter.abort.unsubscribe succeeded", {
|
||||
runFriendlyId,
|
||||
traceId: run.traceId,
|
||||
});
|
||||
})
|
||||
.catch((error) => {
|
||||
logger.error("RunStreamPresenter.abort.unsubscribe failed", {
|
||||
runFriendlyId,
|
||||
traceId: run.traceId,
|
||||
error: {
|
||||
name: error.name,
|
||||
message: error.message,
|
||||
stack: error.stack,
|
||||
},
|
||||
});
|
||||
});
|
||||
});
|
||||
};
|
||||
});
|
||||
}
|
||||
|
||||
+20
-6
@@ -2,6 +2,8 @@ import { NoSymbolIcon } from "@heroicons/react/20/solid";
|
||||
import { CheckCircleIcon, XCircleIcon } from "@heroicons/react/24/solid";
|
||||
import type { LoaderFunctionArgs } from "@remix-run/server-runtime";
|
||||
import { typedjson, useTypedLoaderData } from "remix-typedjson";
|
||||
import { z } from "zod";
|
||||
import { ListPagination } from "~/components/ListPagination";
|
||||
import { EnvironmentLabel } from "~/components/environments/EnvironmentLabel";
|
||||
import { DateTime } from "~/components/primitives/DateTime";
|
||||
import { LabelValueStack } from "~/components/primitives/LabelValueStack";
|
||||
@@ -17,30 +19,38 @@ import {
|
||||
TableRow,
|
||||
} from "~/components/primitives/Table";
|
||||
import { TextLink } from "~/components/primitives/TextLink";
|
||||
import { useOrganization } from "~/hooks/useOrganizations";
|
||||
import { useProject } from "~/hooks/useProject";
|
||||
import { DirectionSchema } from "~/components/runs/RunStatuses";
|
||||
import { ScheduledTriggersPresenter } from "~/presenters/ScheduledTriggersPresenter.server";
|
||||
import { requireUserId } from "~/services/session.server";
|
||||
import { ProjectParamSchema, docsPath } from "~/utils/pathBuilder";
|
||||
|
||||
const SearchSchema = z.object({
|
||||
cursor: z.string().optional(),
|
||||
direction: DirectionSchema.optional(),
|
||||
});
|
||||
|
||||
export const loader = async ({ request, params }: LoaderFunctionArgs) => {
|
||||
const userId = await requireUserId(request);
|
||||
const { organizationSlug, projectParam } = ProjectParamSchema.parse(params);
|
||||
|
||||
const url = new URL(request.url);
|
||||
const s = Object.fromEntries(url.searchParams.entries());
|
||||
const searchParams = SearchSchema.parse(s);
|
||||
|
||||
const presenter = new ScheduledTriggersPresenter();
|
||||
const data = await presenter.call({
|
||||
userId,
|
||||
organizationSlug,
|
||||
projectSlug: projectParam,
|
||||
direction: searchParams.direction,
|
||||
cursor: searchParams.cursor,
|
||||
});
|
||||
|
||||
return typedjson(data);
|
||||
};
|
||||
|
||||
export default function Integrations() {
|
||||
const { scheduled } = useTypedLoaderData<typeof loader>();
|
||||
const organization = useOrganization();
|
||||
const project = useProject();
|
||||
export default function Route() {
|
||||
const { scheduled, pagination } = useTypedLoaderData<typeof loader>();
|
||||
|
||||
return (
|
||||
<>
|
||||
@@ -49,6 +59,10 @@ export default function Integrations() {
|
||||
expression or an interval.
|
||||
</Paragraph>
|
||||
|
||||
{scheduled.length > 0 && (
|
||||
<ListPagination list={{ pagination }} className="mt-2 justify-end" />
|
||||
)}
|
||||
|
||||
<Table containerClassName="mt-4">
|
||||
<TableHeader>
|
||||
<TableRow>
|
||||
|
||||
+66
-23
@@ -1,12 +1,13 @@
|
||||
import {
|
||||
BoltSlashIcon,
|
||||
ChevronDownIcon,
|
||||
ChevronRightIcon,
|
||||
MagnifyingGlassMinusIcon,
|
||||
MagnifyingGlassPlusIcon,
|
||||
} from "@heroicons/react/20/solid";
|
||||
import type { Location } from "@remix-run/react";
|
||||
import { useParams, useRevalidator } from "@remix-run/react";
|
||||
import { LoaderFunctionArgs } from "@remix-run/server-runtime";
|
||||
import { useLoaderData, useParams, useRevalidator } from "@remix-run/react";
|
||||
import { LoaderFunctionArgs, SerializeFrom } from "@remix-run/server-runtime";
|
||||
import { Virtualizer } from "@tanstack/react-virtual";
|
||||
import {
|
||||
formatDurationMilliseconds,
|
||||
@@ -17,10 +18,10 @@ import { RuntimeEnvironmentType } from "@trigger.dev/database";
|
||||
import { motion } from "framer-motion";
|
||||
import { useCallback, useEffect, useRef, useState } from "react";
|
||||
import { useHotkeys } from "react-hotkeys-hook";
|
||||
import { typedjson, useTypedLoaderData } from "remix-typedjson";
|
||||
import { ShowParentIcon, ShowParentIconSelected } from "~/assets/icons/ShowParentIcon";
|
||||
import tileBgPath from "~/assets/images/error-banner-tile@2x.png";
|
||||
import { BlankstateInstructions } from "~/components/BlankstateInstructions";
|
||||
import { AdminDebugTooltip } from "~/components/admin/debugTooltip";
|
||||
import { InlineCode } from "~/components/code/InlineCode";
|
||||
import { EnvironmentLabel } from "~/components/environments/EnvironmentLabel";
|
||||
import { MainCenteredContainer, PageBody } from "~/components/layout/AppLayout";
|
||||
@@ -32,6 +33,7 @@ import { Input } from "~/components/primitives/Input";
|
||||
import { NavBar, PageAccessories, PageTitle } from "~/components/primitives/PageHeader";
|
||||
import { Paragraph } from "~/components/primitives/Paragraph";
|
||||
import { Popover, PopoverArrowTrigger, PopoverContent } from "~/components/primitives/Popover";
|
||||
import { Property, PropertyTable } from "~/components/primitives/PropertyTable";
|
||||
import {
|
||||
ResizableHandle,
|
||||
ResizablePanel,
|
||||
@@ -54,7 +56,7 @@ import { useProject } from "~/hooks/useProject";
|
||||
import { useReplaceLocation } from "~/hooks/useReplaceLocation";
|
||||
import { Shortcut, useShortcutKeys } from "~/hooks/useShortcutKeys";
|
||||
import { useUser } from "~/hooks/useUser";
|
||||
import { RunEvent, RunPresenter } from "~/presenters/v3/RunPresenter.server";
|
||||
import { RunPresenter } from "~/presenters/v3/RunPresenter.server";
|
||||
import { getResizableRunSettings, setResizableRunSettings } from "~/services/resizablePanel";
|
||||
import { requireUserId } from "~/services/session.server";
|
||||
import { cn } from "~/utils/cn";
|
||||
@@ -67,8 +69,10 @@ import {
|
||||
v3RunsPath,
|
||||
} from "~/utils/pathBuilder";
|
||||
import { SpanView } from "../resources.orgs.$organizationSlug.projects.v3.$projectParam.runs.$runParam.spans.$spanParam/route";
|
||||
import { AdminDebugTooltip } from "~/components/admin/debugTooltip";
|
||||
import { Property, PropertyTable } from "~/components/primitives/PropertyTable";
|
||||
import { SimpleTooltip } from "~/components/primitives/Tooltip";
|
||||
import { env } from "~/env.server";
|
||||
|
||||
type TraceEvent = NonNullable<SerializeFrom<typeof loader>["trace"]>["events"][0];
|
||||
|
||||
export const loader = async ({ request, params }: LoaderFunctionArgs) => {
|
||||
const userId = await requireUserId(request);
|
||||
@@ -85,10 +89,12 @@ export const loader = async ({ request, params }: LoaderFunctionArgs) => {
|
||||
//resizable settings
|
||||
const resizeSettings = await getResizableRunSettings(request);
|
||||
|
||||
return typedjson({
|
||||
...result,
|
||||
return {
|
||||
run: result.run,
|
||||
trace: result.trace,
|
||||
maximumLiveReloadingSetting: env.MAXIMUM_LIVE_RELOADING_EVENTS,
|
||||
resizeSettings,
|
||||
});
|
||||
};
|
||||
};
|
||||
|
||||
function getSpanId(location: Location<any>): string | undefined {
|
||||
@@ -97,7 +103,8 @@ function getSpanId(location: Location<any>): string | undefined {
|
||||
}
|
||||
|
||||
export default function Page() {
|
||||
const { run, trace, resizeSettings } = useTypedLoaderData<typeof loader>();
|
||||
const { run, trace, resizeSettings, maximumLiveReloadingSetting } =
|
||||
useLoaderData<typeof loader>();
|
||||
const organization = useOrganization();
|
||||
const project = useProject();
|
||||
const user = useUser();
|
||||
@@ -167,6 +174,7 @@ export default function Page() {
|
||||
}
|
||||
|
||||
const { events, parentRunFriendlyId, duration, rootSpanStatus, rootStartedAt } = trace;
|
||||
const shouldLiveReload = events.length <= maximumLiveReloadingSetting;
|
||||
|
||||
const changeToSpan = useDebounce((selectedSpan: string) => {
|
||||
replaceSearchParam("span", selectedSpan);
|
||||
@@ -175,6 +183,7 @@ export default function Page() {
|
||||
const revalidator = useRevalidator();
|
||||
const streamedEvents = useEventSource(v3RunStreamingPath(organization, project, run), {
|
||||
event: "message",
|
||||
disabled: !shouldLiveReload,
|
||||
});
|
||||
useEffect(() => {
|
||||
if (streamedEvents !== null) {
|
||||
@@ -252,8 +261,10 @@ export default function Page() {
|
||||
}}
|
||||
totalDuration={duration}
|
||||
rootSpanStatus={rootSpanStatus}
|
||||
rootStartedAt={rootStartedAt}
|
||||
rootStartedAt={rootStartedAt ? new Date(rootStartedAt) : undefined}
|
||||
environmentType={run.environment.type}
|
||||
shouldLiveReload={shouldLiveReload}
|
||||
maximumLiveReloadingSetting={maximumLiveReloadingSetting}
|
||||
/>
|
||||
</ResizablePanel>
|
||||
<ResizableHandle withHandle />
|
||||
@@ -274,7 +285,7 @@ export default function Page() {
|
||||
}
|
||||
|
||||
type TasksTreeViewProps = {
|
||||
events: RunEvent[];
|
||||
events: TraceEvent[];
|
||||
selectedId?: string;
|
||||
parentRunFriendlyId?: string;
|
||||
onSelectedIdChanged: (selectedId: string | undefined) => void;
|
||||
@@ -282,6 +293,8 @@ type TasksTreeViewProps = {
|
||||
rootSpanStatus: "executing" | "completed" | "failed";
|
||||
rootStartedAt: Date | undefined;
|
||||
environmentType: RuntimeEnvironmentType;
|
||||
shouldLiveReload: boolean;
|
||||
maximumLiveReloadingSetting: number;
|
||||
};
|
||||
|
||||
function TasksTreeView({
|
||||
@@ -293,6 +306,8 @@ function TasksTreeView({
|
||||
rootSpanStatus,
|
||||
rootStartedAt,
|
||||
environmentType,
|
||||
shouldLiveReload,
|
||||
maximumLiveReloadingSetting,
|
||||
}: TasksTreeViewProps) {
|
||||
const [filterText, setFilterText] = useState("");
|
||||
const [errorsOnly, setErrorsOnly] = useState(false);
|
||||
@@ -367,7 +382,11 @@ function TasksTreeView({
|
||||
This is the root task
|
||||
</Paragraph>
|
||||
)}
|
||||
<LiveReloadingStatus rootSpanCompleted={rootSpanStatus !== "executing"} />
|
||||
<LiveReloadingStatus
|
||||
rootSpanCompleted={rootSpanStatus !== "executing"}
|
||||
isLiveReloading={shouldLiveReload}
|
||||
settingValue={maximumLiveReloadingSetting}
|
||||
/>
|
||||
</div>
|
||||
<TreeView
|
||||
parentRef={parentRef}
|
||||
@@ -750,7 +769,7 @@ function TimelineView({
|
||||
);
|
||||
}
|
||||
|
||||
function NodeText({ node }: { node: RunEvent }) {
|
||||
function NodeText({ node }: { node: TraceEvent }) {
|
||||
const className = "truncate";
|
||||
return (
|
||||
<Paragraph variant="small" className={cn(className)}>
|
||||
@@ -759,7 +778,7 @@ function NodeText({ node }: { node: RunEvent }) {
|
||||
);
|
||||
}
|
||||
|
||||
function NodeStatusIcon({ node }: { node: RunEvent }) {
|
||||
function NodeStatusIcon({ node }: { node: TraceEvent }) {
|
||||
if (node.data.level !== "TRACE") return null;
|
||||
if (node.data.style.variant !== "primary") return null;
|
||||
|
||||
@@ -834,16 +853,40 @@ function ShowParentLink({ runFriendlyId }: { runFriendlyId: string }) {
|
||||
);
|
||||
}
|
||||
|
||||
function LiveReloadingStatus({ rootSpanCompleted }: { rootSpanCompleted: boolean }) {
|
||||
function LiveReloadingStatus({
|
||||
rootSpanCompleted,
|
||||
isLiveReloading,
|
||||
settingValue,
|
||||
}: {
|
||||
rootSpanCompleted: boolean;
|
||||
isLiveReloading: boolean;
|
||||
settingValue: number;
|
||||
}) {
|
||||
if (rootSpanCompleted) return null;
|
||||
|
||||
return (
|
||||
<div className="flex items-center gap-1">
|
||||
<PulsingDot />
|
||||
<Paragraph variant="extra-small" className="whitespace-nowrap text-blue-500">
|
||||
Live reloading
|
||||
</Paragraph>
|
||||
</div>
|
||||
<>
|
||||
{isLiveReloading ? (
|
||||
<div className="flex items-center gap-1">
|
||||
<PulsingDot />
|
||||
<Paragraph variant="extra-small" className="whitespace-nowrap text-blue-500">
|
||||
Live reloading
|
||||
</Paragraph>
|
||||
</div>
|
||||
) : (
|
||||
<SimpleTooltip
|
||||
content={`Live reloading is disabled because you've exceeded ${settingValue} logs.`}
|
||||
button={
|
||||
<div className="flex items-center gap-1">
|
||||
<BoltSlashIcon className="size-3.5 text-text-dimmed" />
|
||||
<Paragraph variant="extra-small" className="whitespace-nowrap text-text-dimmed">
|
||||
Live reloading disabled
|
||||
</Paragraph>
|
||||
</div>
|
||||
}
|
||||
></SimpleTooltip>
|
||||
)}
|
||||
</>
|
||||
);
|
||||
}
|
||||
|
||||
@@ -862,7 +905,7 @@ function SpanWithDuration({
|
||||
showDuration,
|
||||
node,
|
||||
...props
|
||||
}: Timeline.SpanProps & { node: RunEvent; showDuration: boolean }) {
|
||||
}: Timeline.SpanProps & { node: TraceEvent; showDuration: boolean }) {
|
||||
return (
|
||||
<Timeline.Span {...props}>
|
||||
<motion.div
|
||||
|
||||
@@ -1,10 +1,8 @@
|
||||
import type { ActionFunctionArgs } from "@remix-run/server-runtime";
|
||||
import { json } from "@remix-run/server-runtime";
|
||||
import { z } from "zod";
|
||||
import { env } from "~/env.server";
|
||||
import { authenticateApiRequest } from "~/services/apiAuth.server";
|
||||
import { logger } from "~/services/logger.server";
|
||||
import { r2 } from "~/v3/r2.server";
|
||||
import { generatePresignedUrl } from "~/v3/r2.server";
|
||||
|
||||
const ParamsSchema = z.object({
|
||||
"*": z.string(),
|
||||
@@ -26,34 +24,19 @@ export async function action({ request, params }: ActionFunctionArgs) {
|
||||
const parsedParams = ParamsSchema.parse(params);
|
||||
const filename = parsedParams["*"];
|
||||
|
||||
if (!env.OBJECT_STORE_BASE_URL) {
|
||||
return json({ error: "Object store base URL is not set" }, { status: 500 });
|
||||
}
|
||||
|
||||
if (!r2) {
|
||||
return json({ error: "Object store credentials are not set" }, { status: 500 });
|
||||
}
|
||||
|
||||
const url = new URL(env.OBJECT_STORE_BASE_URL);
|
||||
url.pathname = `/packets/${authenticationResult.environment.project.externalRef}/${authenticationResult.environment.slug}/${filename}`;
|
||||
url.searchParams.set("X-Amz-Expires", "300"); // 5 minutes
|
||||
|
||||
const signed = await r2.sign(
|
||||
new Request(url, {
|
||||
method: "PUT",
|
||||
}),
|
||||
{
|
||||
aws: { signQuery: true },
|
||||
}
|
||||
const presignedUrl = await generatePresignedUrl(
|
||||
authenticationResult.environment.project.externalRef,
|
||||
authenticationResult.environment.slug,
|
||||
filename,
|
||||
"PUT"
|
||||
);
|
||||
|
||||
logger.debug("Generated presigned URL", {
|
||||
url: signed.url,
|
||||
headers: Object.fromEntries(signed.headers),
|
||||
});
|
||||
if (!presignedUrl) {
|
||||
return json({ error: "Failed to generate presigned URL" }, { status: 500 });
|
||||
}
|
||||
|
||||
// Caller can now use this URL to upload to that object.
|
||||
return json({ presignedUrl: signed.url });
|
||||
return json({ presignedUrl });
|
||||
}
|
||||
|
||||
export async function loader({ request, params }: ActionFunctionArgs) {
|
||||
@@ -67,35 +50,17 @@ export async function loader({ request, params }: ActionFunctionArgs) {
|
||||
const parsedParams = ParamsSchema.parse(params);
|
||||
const filename = parsedParams["*"];
|
||||
|
||||
if (!env.OBJECT_STORE_BASE_URL) {
|
||||
return json({ error: "Object store base URL is not set" }, { status: 500 });
|
||||
}
|
||||
|
||||
if (!r2) {
|
||||
return json({ error: "Object store credentials are not set" }, { status: 500 });
|
||||
}
|
||||
|
||||
const url = new URL(env.OBJECT_STORE_BASE_URL);
|
||||
url.pathname = `/packets/${authenticationResult.environment.project.externalRef}/${authenticationResult.environment.slug}/${filename}`;
|
||||
url.searchParams.set("X-Amz-Expires", "300"); // 5 minutes
|
||||
|
||||
const signed = await r2.sign(
|
||||
new Request(url, {
|
||||
method: request.method,
|
||||
}),
|
||||
{
|
||||
aws: { signQuery: true },
|
||||
}
|
||||
const presignedUrl = await generatePresignedUrl(
|
||||
authenticationResult.environment.project.externalRef,
|
||||
authenticationResult.environment.slug,
|
||||
filename,
|
||||
"GET"
|
||||
);
|
||||
|
||||
logger.debug("Generated presigned URL", {
|
||||
url: signed.url,
|
||||
headers: Object.fromEntries(signed.headers),
|
||||
});
|
||||
if (!presignedUrl) {
|
||||
return json({ error: "Failed to generate presigned URL" }, { status: 500 });
|
||||
}
|
||||
|
||||
const getUrl = new URL(url.href);
|
||||
getUrl.searchParams.delete("X-Amz-Expires");
|
||||
|
||||
// Caller can now use this URL to upload to that object.
|
||||
return json({ presignedUrl: signed.url });
|
||||
// Caller can now use this URL to fetch that object.
|
||||
return json({ presignedUrl });
|
||||
}
|
||||
|
||||
+65
-59
@@ -3,6 +3,7 @@ import { PrismaClientOrTransaction, prisma } from "~/db.server";
|
||||
import { taskWithAttemptsToServerTask } from "~/models/task.server";
|
||||
import type { AuthenticatedEnvironment } from "~/services/apiAuth.server";
|
||||
import { logger } from "~/services/logger.server";
|
||||
import { startActiveSpan } from "~/v3/tracer.server";
|
||||
|
||||
export class CompleteRunTaskService {
|
||||
#prismaClient: PrismaClientOrTransaction;
|
||||
@@ -17,76 +18,81 @@ export class CompleteRunTaskService {
|
||||
id: string,
|
||||
taskBody: CompleteTaskBodyOutput
|
||||
): Promise<ServerTask | undefined> {
|
||||
const existingTask = await this.#prismaClient.task.findUnique({
|
||||
where: {
|
||||
id,
|
||||
},
|
||||
include: {
|
||||
run: true,
|
||||
attempts: {
|
||||
where: {
|
||||
status: "PENDING",
|
||||
},
|
||||
orderBy: {
|
||||
number: "desc",
|
||||
},
|
||||
take: 1,
|
||||
return startActiveSpan("CompleteRunTaskService.call", async (span) => {
|
||||
span.setAttribute("runId", runId);
|
||||
span.setAttribute("taskId", id);
|
||||
|
||||
const existingTask = await this.#prismaClient.task.findUnique({
|
||||
where: {
|
||||
id,
|
||||
},
|
||||
include: {
|
||||
run: true,
|
||||
attempts: {
|
||||
where: {
|
||||
status: "PENDING",
|
||||
},
|
||||
orderBy: {
|
||||
number: "desc",
|
||||
},
|
||||
take: 1,
|
||||
},
|
||||
},
|
||||
},
|
||||
});
|
||||
|
||||
if (!existingTask) {
|
||||
return;
|
||||
}
|
||||
|
||||
if (existingTask.runId !== runId) {
|
||||
return;
|
||||
}
|
||||
|
||||
if (existingTask.run.environmentId !== environment.id) {
|
||||
return;
|
||||
}
|
||||
|
||||
if (
|
||||
existingTask.status === "COMPLETED" ||
|
||||
existingTask.status === "ERRORED" ||
|
||||
existingTask.status === "CANCELED"
|
||||
) {
|
||||
logger.debug("Task already completed", {
|
||||
existingTask,
|
||||
});
|
||||
|
||||
return taskWithAttemptsToServerTask(existingTask);
|
||||
}
|
||||
if (!existingTask) {
|
||||
return;
|
||||
}
|
||||
|
||||
if (existingTask.attempts.length === 1) {
|
||||
await this.#prismaClient.taskAttempt.update({
|
||||
if (existingTask.runId !== runId) {
|
||||
return;
|
||||
}
|
||||
|
||||
if (existingTask.run.environmentId !== environment.id) {
|
||||
return;
|
||||
}
|
||||
|
||||
if (
|
||||
existingTask.status === "COMPLETED" ||
|
||||
existingTask.status === "ERRORED" ||
|
||||
existingTask.status === "CANCELED"
|
||||
) {
|
||||
logger.debug("Task already completed", {
|
||||
taskId: id,
|
||||
});
|
||||
|
||||
return taskWithAttemptsToServerTask(existingTask);
|
||||
}
|
||||
|
||||
if (existingTask.attempts.length === 1) {
|
||||
await this.#prismaClient.taskAttempt.update({
|
||||
where: {
|
||||
id: existingTask.attempts[0].id,
|
||||
},
|
||||
data: {
|
||||
status: "COMPLETED",
|
||||
},
|
||||
});
|
||||
}
|
||||
|
||||
const updatedTask = await this.#prismaClient.task.update({
|
||||
where: {
|
||||
id: existingTask.attempts[0].id,
|
||||
id,
|
||||
},
|
||||
data: {
|
||||
status: "COMPLETED",
|
||||
output: taskBody.output as any,
|
||||
outputIsUndefined: typeof taskBody.output === "undefined",
|
||||
completedAt: new Date(),
|
||||
outputProperties: taskBody.properties,
|
||||
},
|
||||
include: {
|
||||
attempts: true,
|
||||
run: true,
|
||||
},
|
||||
});
|
||||
}
|
||||
|
||||
const updatedTask = await this.#prismaClient.task.update({
|
||||
where: {
|
||||
id,
|
||||
},
|
||||
data: {
|
||||
status: "COMPLETED",
|
||||
output: taskBody.output as any,
|
||||
outputIsUndefined: typeof taskBody.output === "undefined",
|
||||
completedAt: new Date(),
|
||||
outputProperties: taskBody.properties,
|
||||
},
|
||||
include: {
|
||||
attempts: true,
|
||||
run: true,
|
||||
},
|
||||
return taskWithAttemptsToServerTask(updatedTask);
|
||||
});
|
||||
|
||||
return taskWithAttemptsToServerTask(updatedTask);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -9,8 +9,10 @@ import {
|
||||
import { z } from "zod";
|
||||
import type { AuthenticatedEnvironment } from "~/services/apiAuth.server";
|
||||
import { authenticateApiRequest } from "~/services/apiAuth.server";
|
||||
import { logger } from "~/services/logger.server";
|
||||
import { CompleteRunTaskService } from "./CompleteRunTaskService.server";
|
||||
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 ParamsSchema = z.object({
|
||||
runId: z.string(),
|
||||
@@ -44,46 +46,51 @@ export async function action({ request, params }: ActionFunctionArgs) {
|
||||
return json({ error: "Invalid headers" }, { status: 400 });
|
||||
}
|
||||
|
||||
// Check the content size of the request and make sure it's not too large
|
||||
const contentLength = request.headers.get("content-length");
|
||||
|
||||
if (!contentLength || parseInt(contentLength) > 3 * 1024 * 1024) {
|
||||
const service = new FailRunTaskService();
|
||||
|
||||
await service.call(authenticatedEnv, runId, id, {
|
||||
error: {
|
||||
message: "Task output is too large. The limit is 3MB",
|
||||
},
|
||||
});
|
||||
|
||||
return json({ error: "Task output is too large. The limit is 3MB" }, { status: 413 });
|
||||
}
|
||||
|
||||
const { "trigger-version": triggerVersion } = headers.data;
|
||||
|
||||
// Now parse the request body
|
||||
const anyBody = await request.json();
|
||||
|
||||
logger.debug("CompleteRunTaskService.call() request body", {
|
||||
body: anyBody,
|
||||
runId,
|
||||
id,
|
||||
});
|
||||
const anyBody = await parseRequestJsonAsync(request, { runId });
|
||||
|
||||
if (triggerVersion === API_VERSIONS.SERIALIZED_TASK_OUTPUT) {
|
||||
const body = CompleteTaskBodyV2InputSchema.safeParse(anyBody);
|
||||
const body = await startActiveSpan("CompleteTaskBodyV2InputSchema.safeParse()", async () => {
|
||||
return CompleteTaskBodyV2InputSchema.safeParse(anyBody);
|
||||
});
|
||||
|
||||
if (!body.success) {
|
||||
return json({ error: "Invalid request body" }, { status: 400 });
|
||||
}
|
||||
|
||||
// Make sure the length of the output is less than 3MB
|
||||
if (body.data.output && body.data.output.length > 3 * 1024 * 1024) {
|
||||
return json({ error: "Output must be less than 3MB" }, { status: 400 });
|
||||
}
|
||||
|
||||
return await completeRunTask(authenticatedEnv, runId, id, {
|
||||
...body.data,
|
||||
output: body.data.output ? (JSON.parse(body.data.output) as any) : undefined,
|
||||
});
|
||||
} else {
|
||||
const body = CompleteTaskBodyInputSchema.safeParse(anyBody);
|
||||
const body = await startActiveSpan("CompleteTaskBodyInputSchema.safeParse()", async () => {
|
||||
return CompleteTaskBodyInputSchema.omit({ output: true }).safeParse(anyBody);
|
||||
});
|
||||
|
||||
if (!body.success) {
|
||||
return json({ error: "Invalid request body" }, { status: 400 });
|
||||
}
|
||||
|
||||
// Make sure the length of the output is less than 3MB
|
||||
if (JSON.stringify(body.data.output).length > 3 * 1024 * 1024) {
|
||||
return json({ error: "Output must be less than 3MB" }, { status: 400 });
|
||||
}
|
||||
const output = (anyBody as any).output;
|
||||
|
||||
return await completeRunTask(authenticatedEnv, runId, id, body.data);
|
||||
return await completeRunTask(authenticatedEnv, runId, id, { ...body.data, output });
|
||||
}
|
||||
}
|
||||
|
||||
@@ -98,12 +105,6 @@ async function completeRunTask(
|
||||
try {
|
||||
const task = await service.call(environment, runId, id, taskBody);
|
||||
|
||||
logger.debug("CompleteRunTaskService.call() response body", {
|
||||
runId,
|
||||
id,
|
||||
task,
|
||||
});
|
||||
|
||||
if (!task) {
|
||||
return json({ message: "Task not found" }, { status: 404 });
|
||||
}
|
||||
|
||||
@@ -57,10 +57,6 @@ export class FailRunTaskService {
|
||||
existingTask.status === "ERRORED" ||
|
||||
existingTask.status === "CANCELED"
|
||||
) {
|
||||
logger.debug("Task already completed", {
|
||||
existingTask,
|
||||
});
|
||||
|
||||
return existingTask;
|
||||
}
|
||||
|
||||
|
||||
@@ -6,6 +6,8 @@ import { authenticateApiRequest } from "~/services/apiAuth.server";
|
||||
import { logger } from "~/services/logger.server";
|
||||
import { RunTaskService } from "~/services/tasks/runTask.server";
|
||||
import { ChangeRequestLazyLoadedCachedTasks } from "./ChangeRequestLazyLoadedCachedTasks.server";
|
||||
import { startActiveSpan } from "~/v3/tracer.server";
|
||||
import { parseRequestJsonAsync } from "~/utils/parseRequestJson.server";
|
||||
|
||||
const ParamsSchema = z.object({
|
||||
runId: z.string(),
|
||||
@@ -17,6 +19,8 @@ const HeadersSchema = z.object({
|
||||
"x-cached-tasks-cursor": z.string().optional().nullable(),
|
||||
});
|
||||
|
||||
const BodySchema = RunTaskBodyOutputSchema.omit({ params: true });
|
||||
|
||||
export async function action({ request, params }: ActionFunctionArgs) {
|
||||
// Ensure this is a POST request
|
||||
if (request.method.toUpperCase() !== "POST") {
|
||||
@@ -44,18 +48,26 @@ export async function action({ request, params }: ActionFunctionArgs) {
|
||||
|
||||
const { runId } = ParamsSchema.parse(params);
|
||||
|
||||
const contentLength = request.headers.get("content-length");
|
||||
|
||||
if (!contentLength || parseInt(contentLength) > 3 * 1024 * 1024) {
|
||||
return json({ error: "Request body too large" }, { status: 413 });
|
||||
}
|
||||
|
||||
// Now parse the request body
|
||||
const anyBody = await request.json();
|
||||
const anyBody = await parseRequestJsonAsync(request, { runId });
|
||||
|
||||
logger.debug("RunTaskService.call() request body", {
|
||||
body: anyBody,
|
||||
runId,
|
||||
idempotencyKey,
|
||||
triggerVersion,
|
||||
cachedTasksCursor,
|
||||
});
|
||||
|
||||
const body = RunTaskBodyOutputSchema.safeParse(anyBody);
|
||||
const body = await startActiveSpan(
|
||||
"BodySchema.safeParse",
|
||||
async () => {
|
||||
return BodySchema.safeParse(anyBody);
|
||||
},
|
||||
{
|
||||
attributes: {
|
||||
runId,
|
||||
},
|
||||
}
|
||||
);
|
||||
|
||||
if (!body.success) {
|
||||
return json({ error: "Invalid request body" }, { status: 400 });
|
||||
@@ -64,12 +76,9 @@ export async function action({ request, params }: ActionFunctionArgs) {
|
||||
const service = new RunTaskService();
|
||||
|
||||
try {
|
||||
const task = await service.call(runId, idempotencyKey, body.data);
|
||||
|
||||
logger.debug("RunTaskService.call() response body", {
|
||||
runId,
|
||||
idempotencyKey,
|
||||
task,
|
||||
const task = await service.call(runId, idempotencyKey, {
|
||||
...body.data,
|
||||
params: (anyBody as any).params,
|
||||
});
|
||||
|
||||
if (!task) {
|
||||
@@ -84,7 +93,6 @@ export async function action({ request, params }: ActionFunctionArgs) {
|
||||
logger.debug(
|
||||
"RunTaskService.call() response migrating with ChangeRequestLazyLoadedCachedTasks",
|
||||
{
|
||||
responseBody,
|
||||
cachedTasksCursor,
|
||||
}
|
||||
);
|
||||
|
||||
@@ -7,6 +7,7 @@ import { authenticateApiRequest } from "~/services/apiAuth.server";
|
||||
import { logger } from "~/services/logger.server";
|
||||
import { BatchTriggerTaskService } from "~/v3/services/batchTriggerTask.server";
|
||||
import { HeadersSchema } from "./api.v1.tasks.$taskId.trigger";
|
||||
import { env } from "~/env.server";
|
||||
|
||||
const ParamsSchema = z.object({
|
||||
taskId: z.string(),
|
||||
@@ -43,6 +44,12 @@ export async function action({ request, params }: ActionFunctionArgs) {
|
||||
|
||||
const { taskId } = ParamsSchema.parse(params);
|
||||
|
||||
const contentLength = request.headers.get("content-length");
|
||||
|
||||
if (!contentLength || parseInt(contentLength) > env.TASK_PAYLOAD_MAXIMUM_SIZE) {
|
||||
return json({ error: "Request body too large" }, { status: 413 });
|
||||
}
|
||||
|
||||
// Now parse the request body
|
||||
const anyBody = await request.json();
|
||||
|
||||
|
||||
@@ -2,9 +2,12 @@ import type { ActionFunctionArgs } from "@remix-run/server-runtime";
|
||||
import { json } from "@remix-run/server-runtime";
|
||||
import { TriggerTaskRequestBody } from "@trigger.dev/core/v3";
|
||||
import { z } from "zod";
|
||||
import { env } from "~/env.server";
|
||||
import { authenticateApiRequest } from "~/services/apiAuth.server";
|
||||
import { logger } from "~/services/logger.server";
|
||||
import { parseRequestJsonAsync } from "~/utils/parseRequestJson.server";
|
||||
import { TriggerTaskService } from "~/v3/services/triggerTask.server";
|
||||
import { startActiveSpan } from "~/v3/tracer.server";
|
||||
|
||||
const ParamsSchema = z.object({
|
||||
taskId: z.string(),
|
||||
@@ -32,6 +35,12 @@ export async function action({ request, params }: ActionFunctionArgs) {
|
||||
return json({ error: "Invalid or Missing API key" }, { status: 401 });
|
||||
}
|
||||
|
||||
const contentLength = request.headers.get("content-length");
|
||||
|
||||
if (!contentLength || parseInt(contentLength) > env.TASK_PAYLOAD_MAXIMUM_SIZE) {
|
||||
return json({ error: "Request body too large" }, { status: 413 });
|
||||
}
|
||||
|
||||
const rawHeaders = Object.fromEntries(request.headers);
|
||||
|
||||
const headers = HeadersSchema.safeParse(rawHeaders);
|
||||
@@ -52,9 +61,11 @@ export async function action({ request, params }: ActionFunctionArgs) {
|
||||
const { taskId } = ParamsSchema.parse(params);
|
||||
|
||||
// Now parse the request body
|
||||
const anyBody = await request.json();
|
||||
const anyBody = await parseRequestJsonAsync(request, { taskId });
|
||||
|
||||
const body = TriggerTaskRequestBody.safeParse(anyBody);
|
||||
const body = await startActiveSpan("TriggerTaskRequestBody.safeParse()", async (span) => {
|
||||
return TriggerTaskRequestBody.safeParse(anyBody);
|
||||
});
|
||||
|
||||
if (!body.success) {
|
||||
return json({ error: "Invalid request body" }, { status: 400 });
|
||||
@@ -76,17 +87,23 @@ export async function action({ request, params }: ActionFunctionArgs) {
|
||||
idempotencyKey,
|
||||
triggerVersion,
|
||||
headers: Object.fromEntries(request.headers),
|
||||
body: body.data,
|
||||
options: body.data.options,
|
||||
isFromWorker,
|
||||
traceContext,
|
||||
});
|
||||
|
||||
const run = await service.call(taskId, authenticationResult.environment, body.data, {
|
||||
idempotencyKey: idempotencyKey ?? undefined,
|
||||
triggerVersion: triggerVersion ?? undefined,
|
||||
traceContext,
|
||||
spanParentAsLink: spanParentAsLink === 1,
|
||||
});
|
||||
const run = await service.call(
|
||||
taskId,
|
||||
authenticationResult.environment,
|
||||
{ ...body.data },
|
||||
// { ...body.data, payload: (anyBody as any).payload },
|
||||
{
|
||||
idempotencyKey: idempotencyKey ?? undefined,
|
||||
triggerVersion: triggerVersion ?? undefined,
|
||||
traceContext,
|
||||
spanParentAsLink: spanParentAsLink === 1,
|
||||
}
|
||||
);
|
||||
|
||||
if (!run) {
|
||||
return json({ error: "Task not found" }, { status: 404 });
|
||||
|
||||
+25
-1
@@ -33,7 +33,13 @@ import { redirectWithErrorMessage } from "~/models/message.server";
|
||||
import { Span, SpanPresenter } from "~/presenters/v3/SpanPresenter.server";
|
||||
import { requireUserId } from "~/services/session.server";
|
||||
import { cn } from "~/utils/cn";
|
||||
import { v3RunPath, v3RunSpanPath, v3SpanParamsSchema, v3TraceSpanPath } from "~/utils/pathBuilder";
|
||||
import {
|
||||
v3RunDownloadLogsPath,
|
||||
v3RunPath,
|
||||
v3RunSpanPath,
|
||||
v3SpanParamsSchema,
|
||||
v3TraceSpanPath,
|
||||
} from "~/utils/pathBuilder";
|
||||
import { SpanLink } from "~/v3/eventRepository.server";
|
||||
|
||||
export const loader = async ({ request, params }: LoaderFunctionArgs) => {
|
||||
@@ -256,6 +262,15 @@ function RunActionButtons({ span }: { span: Span }) {
|
||||
if (span.isPartial) {
|
||||
return (
|
||||
<Dialog>
|
||||
<LinkButton
|
||||
to={v3RunDownloadLogsPath({ friendlyId: runParam })}
|
||||
LeadingIcon={CloudArrowDownIcon}
|
||||
variant="tertiary/medium"
|
||||
target="_blank"
|
||||
download
|
||||
>
|
||||
Download logs
|
||||
</LinkButton>
|
||||
<DialogTrigger asChild>
|
||||
<Button variant="danger/medium" LeadingIcon={StopCircleIcon}>
|
||||
Cancel run
|
||||
@@ -276,6 +291,15 @@ function RunActionButtons({ span }: { span: Span }) {
|
||||
|
||||
return (
|
||||
<Dialog>
|
||||
<LinkButton
|
||||
to={v3RunDownloadLogsPath({ friendlyId: runParam })}
|
||||
LeadingIcon={CloudArrowDownIcon}
|
||||
variant="tertiary/medium"
|
||||
target="_blank"
|
||||
download
|
||||
>
|
||||
Download logs
|
||||
</LinkButton>
|
||||
<DialogTrigger asChild>
|
||||
<Button variant="tertiary/medium" LeadingIcon={ArrowPathIcon}>
|
||||
Replay run
|
||||
|
||||
@@ -2,9 +2,8 @@ import { LoaderFunctionArgs } from "@remix-run/node";
|
||||
import { basename } from "node:path";
|
||||
import { z } from "zod";
|
||||
import { prisma } from "~/db.server";
|
||||
import { env } from "~/env.server";
|
||||
import { requireUserId } from "~/services/session.server";
|
||||
import { r2 } from "~/v3/r2.server";
|
||||
import { generatePresignedRequest } from "~/v3/r2.server";
|
||||
|
||||
const ParamSchema = z.object({
|
||||
environmentId: z.string(),
|
||||
@@ -35,27 +34,17 @@ export async function loader({ request, params }: LoaderFunctionArgs) {
|
||||
return new Response("Not found", { status: 404 });
|
||||
}
|
||||
|
||||
if (!env.OBJECT_STORE_BASE_URL) {
|
||||
return new Response("Object store base URL is not set", { status: 500 });
|
||||
}
|
||||
|
||||
if (!r2) {
|
||||
return new Response("Object store credentials are not set", { status: 500 });
|
||||
}
|
||||
|
||||
const url = new URL(env.OBJECT_STORE_BASE_URL);
|
||||
url.pathname = `/packets/${environment.project.externalRef}/${environment.slug}/${filename}`;
|
||||
url.searchParams.set("X-Amz-Expires", "30"); // 30 seconds
|
||||
|
||||
const signed = await r2.sign(
|
||||
new Request(url, {
|
||||
method: "GET",
|
||||
}),
|
||||
{
|
||||
aws: { signQuery: true },
|
||||
}
|
||||
const signed = await generatePresignedRequest(
|
||||
environment.project.externalRef,
|
||||
environment.slug,
|
||||
filename,
|
||||
"GET"
|
||||
);
|
||||
|
||||
if (!signed) {
|
||||
return new Response("Failed to generate presigned URL", { status: 500 });
|
||||
}
|
||||
|
||||
const response = await fetch(signed.url, {
|
||||
headers: signed.headers,
|
||||
});
|
||||
@@ -64,7 +53,7 @@ export async function loader({ request, params }: LoaderFunctionArgs) {
|
||||
status: 200,
|
||||
headers: {
|
||||
"Content-Type": "application/octet-stream",
|
||||
"Content-Disposition": `attachment; filename="${basename(url.pathname)}"`,
|
||||
"Content-Disposition": `attachment; filename="${basename(filename)}"`,
|
||||
},
|
||||
});
|
||||
}
|
||||
|
||||
@@ -0,0 +1,112 @@
|
||||
import { LoaderFunctionArgs } from "@remix-run/server-runtime";
|
||||
import { prisma } from "~/db.server";
|
||||
import { requireUserId } from "~/services/session.server";
|
||||
import { v3RunParamsSchema } from "~/utils/pathBuilder";
|
||||
import {
|
||||
PreparedEvent,
|
||||
RunPreparedEvent,
|
||||
eventRepository,
|
||||
getDateFromNanoseconds,
|
||||
} from "~/v3/eventRepository.server";
|
||||
import { createGzip } from "zlib";
|
||||
import { Readable } from "stream";
|
||||
import { formatDurationMilliseconds } from "@trigger.dev/core/v3/utils/durations";
|
||||
|
||||
export async function loader({ params, request }: LoaderFunctionArgs) {
|
||||
const userId = await requireUserId(request);
|
||||
const parsedParams = v3RunParamsSchema.pick({ runParam: true }).parse(params);
|
||||
|
||||
const run = await prisma.taskRun.findFirst({
|
||||
where: {
|
||||
friendlyId: parsedParams.runParam,
|
||||
project: {
|
||||
organization: {
|
||||
members: {
|
||||
some: {
|
||||
userId,
|
||||
},
|
||||
},
|
||||
},
|
||||
},
|
||||
},
|
||||
});
|
||||
|
||||
if (!run) {
|
||||
return new Response("Not found", { status: 404 });
|
||||
}
|
||||
|
||||
const runEvents = await eventRepository.getRunEvents(run.friendlyId);
|
||||
|
||||
// Create a Readable stream from the runEvents array
|
||||
const readable = new Readable({
|
||||
read() {
|
||||
runEvents.forEach((event) => {
|
||||
try {
|
||||
this.push(formatRunEvent(event) + "\n");
|
||||
} catch {}
|
||||
});
|
||||
this.push(null); // End of stream
|
||||
},
|
||||
});
|
||||
|
||||
// Create a gzip transform stream
|
||||
const gzip = createGzip();
|
||||
|
||||
// Pipe the readable stream into the gzip stream
|
||||
const compressedStream = readable.pipe(gzip);
|
||||
|
||||
// Return the response with the compressed stream
|
||||
return new Response(compressedStream as any, {
|
||||
status: 200,
|
||||
headers: {
|
||||
"Content-Type": "application/octet-stream",
|
||||
"Content-Disposition": `attachment; filename="${parsedParams.runParam}.log"`,
|
||||
"Content-Encoding": "gzip",
|
||||
},
|
||||
});
|
||||
}
|
||||
|
||||
function formatRunEvent(event: RunPreparedEvent): string {
|
||||
const entries = [];
|
||||
const parts: string[] = [];
|
||||
|
||||
parts.push(getDateFromNanoseconds(event.startTime).toISOString());
|
||||
|
||||
if (event.taskSlug) {
|
||||
parts.push(event.taskSlug);
|
||||
}
|
||||
|
||||
parts.push(event.level);
|
||||
parts.push(event.message);
|
||||
|
||||
if (event.level === "TRACE") {
|
||||
parts.push(`(${formatDurationMilliseconds(event.duration / 1_000_000)})`);
|
||||
}
|
||||
|
||||
entries.push(parts.join(" "));
|
||||
|
||||
if (event.events) {
|
||||
for (const subEvent of event.events) {
|
||||
if (subEvent.name === "exception") {
|
||||
const subEventParts: string[] = [];
|
||||
|
||||
subEventParts.push(subEvent.time as unknown as string);
|
||||
|
||||
if (event.taskSlug) {
|
||||
subEventParts.push(event.taskSlug);
|
||||
}
|
||||
|
||||
subEventParts.push(subEvent.name);
|
||||
subEventParts.push((subEvent.properties as any).exception.message);
|
||||
|
||||
if ((subEvent.properties as any).exception.stack) {
|
||||
subEventParts.push((subEvent.properties as any).exception.stack);
|
||||
}
|
||||
|
||||
entries.push(subEventParts.join(" "));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
return entries.join("\n");
|
||||
}
|
||||
@@ -136,7 +136,7 @@ export class EndpointApi {
|
||||
};
|
||||
}
|
||||
|
||||
async executeJobRequest(options: RunJobBody) {
|
||||
async executeJobRequest(options: RunJobBody, timeoutInMs?: number) {
|
||||
const startTimeInMs = performance.now();
|
||||
|
||||
const response = await safeFetch(this.url, {
|
||||
@@ -147,8 +147,18 @@ export class EndpointApi {
|
||||
"x-trigger-action": "EXECUTE_JOB",
|
||||
},
|
||||
body: JSON.stringify(options),
|
||||
signal: timeoutInMs ? AbortSignal.timeout(timeoutInMs) : undefined,
|
||||
});
|
||||
|
||||
if (response) {
|
||||
logger.debug("executeJobRequest() response from endpoint", {
|
||||
status: response.status,
|
||||
headers: Object.fromEntries(response.headers.entries()),
|
||||
});
|
||||
} else {
|
||||
logger.debug("executeJobRequest() no response from endpoint");
|
||||
}
|
||||
|
||||
return {
|
||||
response,
|
||||
parser: RunJobResponseSchema,
|
||||
@@ -434,7 +444,10 @@ async function safeFetch(url: string, options: RequestInit) {
|
||||
} catch (error) {
|
||||
logger.debug("Error while trying to connect to endpoint", {
|
||||
url,
|
||||
error,
|
||||
error:
|
||||
error instanceof Error
|
||||
? { name: error.name, message: error.message, stack: error.stack }
|
||||
: String(error),
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
@@ -269,7 +269,10 @@ export class PerformRunExecutionV3Service {
|
||||
|
||||
// TODO: add the ability to abort the execution from any server using Redis pub/sub
|
||||
const { response, parser, errorParser, headersParser, durationInMs } =
|
||||
await client.executeJobRequest(executionBody);
|
||||
await client.executeJobRequest(
|
||||
executionBody,
|
||||
run.environment.type === "DEVELOPMENT" ? 60_000 * 5 : undefined
|
||||
);
|
||||
|
||||
await createExecutionEvent({
|
||||
eventType: "finish",
|
||||
@@ -929,6 +932,25 @@ export class PerformRunExecutionV3Service {
|
||||
executionCount: number = 1
|
||||
) {
|
||||
await $transaction(this.#prismaClient, async (tx) => {
|
||||
const service = new CompleteRunTaskService(tx);
|
||||
|
||||
const task = await service.call(run.environment, run.id, data.id, {
|
||||
properties: data.properties,
|
||||
output: data.output ? (JSON.parse(data.output) as any) : undefined,
|
||||
});
|
||||
|
||||
if (!task || task.status === "ERRORED") {
|
||||
return await this.#failRunExecution(
|
||||
tx,
|
||||
run,
|
||||
{
|
||||
message: task ? `Task '${task.name}' failed to complete` : "Task failed to complete",
|
||||
},
|
||||
"FAILURE",
|
||||
durationInMs
|
||||
);
|
||||
}
|
||||
|
||||
await tx.jobRun.update({
|
||||
where: {
|
||||
id: run.id,
|
||||
@@ -958,13 +980,6 @@ export class PerformRunExecutionV3Service {
|
||||
},
|
||||
});
|
||||
|
||||
const service = new CompleteRunTaskService(tx);
|
||||
|
||||
await service.call(run.environment, run.id, data.id, {
|
||||
properties: data.properties,
|
||||
output: data.output ? (JSON.parse(data.output) as any) : undefined,
|
||||
});
|
||||
|
||||
await ResumeRunService.enqueue(run, tx);
|
||||
});
|
||||
}
|
||||
|
||||
@@ -6,6 +6,7 @@ import { taskWithAttemptsToServerTask } from "~/models/task.server";
|
||||
import { generateSecret } from "~/services/sources/utils.server";
|
||||
import { ulid } from "~/services/ulid.server";
|
||||
import { taskOperationWorker, workerQueue } from "~/services/worker.server";
|
||||
import { startActiveSpan } from "~/v3/tracer.server";
|
||||
|
||||
export class RunTaskService {
|
||||
#prismaClient: PrismaClient;
|
||||
@@ -19,142 +20,154 @@ export class RunTaskService {
|
||||
idempotencyKey: string,
|
||||
taskBody: RunTaskBodyOutput
|
||||
): Promise<ServerTask | undefined> {
|
||||
const delayUntilInFuture = taskBody.delayUntil
|
||||
? taskBody.delayUntil.getTime() > Date.now()
|
||||
: false;
|
||||
const callbackEnabled = taskBody.callback?.enabled ?? false;
|
||||
return startActiveSpan("RunTaskService.call", async (span) => {
|
||||
span.setAttribute("runId", runId);
|
||||
|
||||
// First
|
||||
const existingTask = await this.#handleExistingTask(
|
||||
runId,
|
||||
idempotencyKey,
|
||||
taskBody,
|
||||
delayUntilInFuture,
|
||||
callbackEnabled
|
||||
);
|
||||
const delayUntilInFuture = taskBody.delayUntil
|
||||
? taskBody.delayUntil.getTime() > Date.now()
|
||||
: false;
|
||||
const callbackEnabled = taskBody.callback?.enabled ?? false;
|
||||
|
||||
if (existingTask) {
|
||||
return taskWithAttemptsToServerTask(existingTask);
|
||||
}
|
||||
// First
|
||||
const existingTask = await this.#handleExistingTask(
|
||||
runId,
|
||||
idempotencyKey,
|
||||
taskBody,
|
||||
delayUntilInFuture,
|
||||
callbackEnabled
|
||||
);
|
||||
|
||||
const run = await this.#prismaClient.jobRun.findUnique({
|
||||
where: {
|
||||
id: runId,
|
||||
},
|
||||
select: {
|
||||
status: true,
|
||||
forceYieldImmediately: true,
|
||||
},
|
||||
});
|
||||
if (existingTask) {
|
||||
span.setAttribute("taskId", existingTask.id);
|
||||
|
||||
if (!run) throw new Error("Run not found");
|
||||
|
||||
const runConnection = taskBody.connectionKey
|
||||
? await this.#prismaClient.runConnection.findUnique({
|
||||
where: {
|
||||
runId_key: {
|
||||
runId,
|
||||
key: taskBody.connectionKey,
|
||||
},
|
||||
},
|
||||
select: {
|
||||
id: true,
|
||||
},
|
||||
})
|
||||
: undefined;
|
||||
|
||||
const results = await $transaction(this.#prismaClient, async (tx) => {
|
||||
// If task.delayUntil is set and is in the future, we'll set the task's status to "WAITING", else set it to RUNNING
|
||||
let status: TaskStatus;
|
||||
|
||||
if (run.status === "CANCELED") {
|
||||
status = "CANCELED";
|
||||
} else {
|
||||
status =
|
||||
delayUntilInFuture || callbackEnabled
|
||||
? "WAITING"
|
||||
: taskBody.noop
|
||||
? "COMPLETED"
|
||||
: "RUNNING";
|
||||
return taskWithAttemptsToServerTask(existingTask);
|
||||
}
|
||||
|
||||
const taskId = ulid();
|
||||
const callbackUrl = callbackEnabled
|
||||
? `${env.APP_ORIGIN}/api/v1/tasks/${taskId}/callback/${generateSecret(12)}`
|
||||
const run = await this.#prismaClient.jobRun.findUnique({
|
||||
where: {
|
||||
id: runId,
|
||||
},
|
||||
select: {
|
||||
status: true,
|
||||
forceYieldImmediately: true,
|
||||
},
|
||||
});
|
||||
|
||||
if (!run) throw new Error("Run not found");
|
||||
|
||||
const runConnection = taskBody.connectionKey
|
||||
? await this.#prismaClient.runConnection.findUnique({
|
||||
where: {
|
||||
runId_key: {
|
||||
runId,
|
||||
key: taskBody.connectionKey,
|
||||
},
|
||||
},
|
||||
select: {
|
||||
id: true,
|
||||
},
|
||||
})
|
||||
: undefined;
|
||||
|
||||
const task = await tx.task.create({
|
||||
data: {
|
||||
id: taskId,
|
||||
idempotencyKey,
|
||||
displayKey: taskBody.displayKey,
|
||||
runConnectionId: runConnection ? runConnection.id : undefined,
|
||||
icon: taskBody.icon,
|
||||
runId,
|
||||
parentId: taskBody.parentId,
|
||||
name: taskBody.name ?? "Task",
|
||||
description: taskBody.description,
|
||||
status,
|
||||
startedAt: new Date(),
|
||||
completedAt: status === "COMPLETED" || status === "CANCELED" ? new Date() : undefined,
|
||||
noop: taskBody.noop,
|
||||
delayUntil: taskBody.delayUntil,
|
||||
params: taskBody.params ?? undefined,
|
||||
properties: this.#filterProperties(taskBody.properties) ?? undefined,
|
||||
redact: taskBody.redact ?? undefined,
|
||||
operation: taskBody.operation,
|
||||
callbackUrl,
|
||||
style: taskBody.style ?? { style: "normal" },
|
||||
childExecutionMode: taskBody.parallel ? "PARALLEL" : "SEQUENTIAL",
|
||||
},
|
||||
});
|
||||
const results = await $transaction(
|
||||
this.#prismaClient,
|
||||
async (tx) => {
|
||||
// If task.delayUntil is set and is in the future, we'll set the task's status to "WAITING", else set it to RUNNING
|
||||
let status: TaskStatus;
|
||||
|
||||
const taskAttempt = await tx.taskAttempt.create({
|
||||
data: {
|
||||
number: 1,
|
||||
taskId: task.id,
|
||||
status: "PENDING",
|
||||
},
|
||||
});
|
||||
if (run.status === "CANCELED") {
|
||||
status = "CANCELED";
|
||||
} else {
|
||||
status =
|
||||
delayUntilInFuture || callbackEnabled
|
||||
? "WAITING"
|
||||
: taskBody.noop
|
||||
? "COMPLETED"
|
||||
: "RUNNING";
|
||||
}
|
||||
|
||||
if (task.status === "RUNNING" && typeof taskBody.operation === "string") {
|
||||
// We need to schedule the operation
|
||||
await taskOperationWorker.enqueue(
|
||||
"performTaskOperation",
|
||||
{
|
||||
id: task.id,
|
||||
},
|
||||
{ tx, runAt: task.delayUntil ?? undefined, jobKey: `operation:${task.id}` }
|
||||
);
|
||||
} else if (task.status === "WAITING" && callbackUrl && taskBody.callback) {
|
||||
if (taskBody.callback.timeoutInSeconds > 0) {
|
||||
// We need to schedule the callback timeout
|
||||
await workerQueue.enqueue(
|
||||
"processCallbackTimeout",
|
||||
{
|
||||
id: task.id,
|
||||
const taskId = ulid();
|
||||
const callbackUrl = callbackEnabled
|
||||
? `${env.APP_ORIGIN}/api/v1/tasks/${taskId}/callback/${generateSecret(12)}`
|
||||
: undefined;
|
||||
|
||||
const task = await tx.task.create({
|
||||
data: {
|
||||
id: taskId,
|
||||
idempotencyKey,
|
||||
displayKey: taskBody.displayKey,
|
||||
runConnectionId: runConnection ? runConnection.id : undefined,
|
||||
icon: taskBody.icon,
|
||||
runId,
|
||||
parentId: taskBody.parentId,
|
||||
name: taskBody.name ?? "Task",
|
||||
description: taskBody.description,
|
||||
status,
|
||||
startedAt: new Date(),
|
||||
completedAt: status === "COMPLETED" || status === "CANCELED" ? new Date() : undefined,
|
||||
noop: taskBody.noop,
|
||||
delayUntil: taskBody.delayUntil,
|
||||
params: taskBody.params ?? undefined,
|
||||
properties: this.#filterProperties(taskBody.properties) ?? undefined,
|
||||
redact: taskBody.redact ?? undefined,
|
||||
operation: taskBody.operation,
|
||||
callbackUrl,
|
||||
style: taskBody.style ?? { style: "normal" },
|
||||
childExecutionMode: taskBody.parallel ? "PARALLEL" : "SEQUENTIAL",
|
||||
},
|
||||
{
|
||||
tx,
|
||||
runAt: new Date(Date.now() + taskBody.callback.timeoutInSeconds * 1000),
|
||||
jobKey: `process-callback:${task.id}`,
|
||||
});
|
||||
|
||||
span.setAttribute("taskId", task.id);
|
||||
|
||||
const taskAttempt = await tx.taskAttempt.create({
|
||||
data: {
|
||||
number: 1,
|
||||
taskId: task.id,
|
||||
status: "PENDING",
|
||||
},
|
||||
});
|
||||
|
||||
if (task.status === "RUNNING" && typeof taskBody.operation === "string") {
|
||||
// We need to schedule the operation
|
||||
await taskOperationWorker.enqueue(
|
||||
"performTaskOperation",
|
||||
{
|
||||
id: task.id,
|
||||
},
|
||||
{ tx, runAt: task.delayUntil ?? undefined, jobKey: `operation:${task.id}` }
|
||||
);
|
||||
} else if (task.status === "WAITING" && callbackUrl && taskBody.callback) {
|
||||
if (taskBody.callback.timeoutInSeconds > 0) {
|
||||
// We need to schedule the callback timeout
|
||||
await workerQueue.enqueue(
|
||||
"processCallbackTimeout",
|
||||
{
|
||||
id: task.id,
|
||||
},
|
||||
{
|
||||
tx,
|
||||
runAt: new Date(Date.now() + taskBody.callback.timeoutInSeconds * 1000),
|
||||
jobKey: `process-callback:${task.id}`,
|
||||
}
|
||||
);
|
||||
}
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
return { task, taskAttempt };
|
||||
},
|
||||
{ timeout: 10000 }
|
||||
);
|
||||
|
||||
if (!results) {
|
||||
return;
|
||||
}
|
||||
|
||||
return { task, taskAttempt };
|
||||
const { task, taskAttempt } = results;
|
||||
|
||||
return task
|
||||
? taskWithAttemptsToServerTask({ ...task, attempts: [taskAttempt], run })
|
||||
: undefined;
|
||||
});
|
||||
|
||||
if (!results) {
|
||||
return;
|
||||
}
|
||||
|
||||
const { task, taskAttempt } = results;
|
||||
|
||||
return task
|
||||
? taskWithAttemptsToServerTask({ ...task, attempts: [taskAttempt], run })
|
||||
: undefined;
|
||||
}
|
||||
|
||||
async #handleExistingTask(
|
||||
|
||||
@@ -0,0 +1,29 @@
|
||||
import { Attributes } from "@opentelemetry/api";
|
||||
import { startActiveSpan } from "~/v3/tracer.server";
|
||||
|
||||
export async function parseRequestJsonAsync(
|
||||
request: Request,
|
||||
attributes?: Attributes
|
||||
): Promise<unknown> {
|
||||
return await startActiveSpan(
|
||||
"parseRequestJsonAsync()",
|
||||
async (span) => {
|
||||
span.setAttribute("content-length", parseInt(request.headers.get("content-length") ?? "0"));
|
||||
span.setAttribute("content-type", request.headers.get("content-type") ?? "application/json");
|
||||
span.setAttribute("experiment.async", false);
|
||||
|
||||
const rawText = await startActiveSpan("request.text()", async () => {
|
||||
return await request.text();
|
||||
});
|
||||
|
||||
if (rawText.length === 0) {
|
||||
return;
|
||||
}
|
||||
|
||||
return JSON.parse(rawText);
|
||||
},
|
||||
{
|
||||
attributes,
|
||||
}
|
||||
);
|
||||
}
|
||||
@@ -375,6 +375,10 @@ export function v3RunPath(organization: OrgForPath, project: ProjectForPath, run
|
||||
return `${v3RunsPath(organization, project)}/${run.friendlyId}`;
|
||||
}
|
||||
|
||||
export function v3RunDownloadLogsPath(run: v3RunForPath) {
|
||||
return `/resources/runs/${run.friendlyId}/logs/download`;
|
||||
}
|
||||
|
||||
export function v3RunSpanPath(
|
||||
organization: OrgForPath,
|
||||
project: ProjectForPath,
|
||||
|
||||
@@ -4,6 +4,7 @@ import { SemanticResourceAttributes } from "@opentelemetry/semantic-conventions"
|
||||
import {
|
||||
ExceptionEventProperties,
|
||||
ExceptionSpanEvent,
|
||||
NULL_SENTINEL,
|
||||
PRIMARY_VARIANT,
|
||||
SemanticInternalAttributes,
|
||||
SpanEvent,
|
||||
@@ -14,7 +15,6 @@ import {
|
||||
correctErrorStackTrace,
|
||||
createPacketAttributesAsJson,
|
||||
flattenAttributes,
|
||||
NULL_SENTINEL,
|
||||
isExceptionSpanEvent,
|
||||
omit,
|
||||
unflattenAttributes,
|
||||
@@ -23,14 +23,15 @@ import { Prisma, TaskEvent, TaskEventStatus, type TaskEventKind } from "@trigger
|
||||
import Redis, { RedisOptions } from "ioredis";
|
||||
import { createHash } from "node:crypto";
|
||||
import { EventEmitter } from "node:stream";
|
||||
import { Gauge } from "prom-client";
|
||||
import { $replica, PrismaClient, PrismaReplicaClient, prisma } from "~/db.server";
|
||||
import { env } from "~/env.server";
|
||||
import { metricsRegister } from "~/metrics.server";
|
||||
import { AuthenticatedEnvironment } from "~/services/apiAuth.server";
|
||||
import { logger } from "~/services/logger.server";
|
||||
import { DynamicFlushScheduler } from "./dynamicFlushScheduler.server";
|
||||
import { singleton } from "~/utils/singleton";
|
||||
import { Gauge } from "prom-client";
|
||||
import { metricsRegister } from "~/metrics.server";
|
||||
import { DynamicFlushScheduler } from "./dynamicFlushScheduler.server";
|
||||
import { startActiveSpan } from "./tracer.server";
|
||||
|
||||
export type CreatableEvent = Omit<
|
||||
Prisma.TaskEventCreateInput,
|
||||
@@ -128,6 +129,10 @@ export type PreparedEvent = Omit<QueriedEvent, "events" | "style" | "duration">
|
||||
style: TaskEventStyle;
|
||||
};
|
||||
|
||||
export type RunPreparedEvent = PreparedEvent & {
|
||||
taskSlug?: string;
|
||||
};
|
||||
|
||||
export type SpanLink =
|
||||
| {
|
||||
type: "run";
|
||||
@@ -374,43 +379,261 @@ export class EventRepository {
|
||||
}
|
||||
|
||||
public async getTraceSummary(traceId: string): Promise<TraceSummary | undefined> {
|
||||
const events = await this.readReplica.taskEvent.findMany({
|
||||
select: {
|
||||
id: true,
|
||||
spanId: true,
|
||||
parentId: true,
|
||||
runId: true,
|
||||
idempotencyKey: true,
|
||||
message: true,
|
||||
style: true,
|
||||
startTime: true,
|
||||
duration: true,
|
||||
isError: true,
|
||||
isPartial: true,
|
||||
isCancelled: true,
|
||||
level: true,
|
||||
events: true,
|
||||
environmentType: true,
|
||||
},
|
||||
where: {
|
||||
traceId,
|
||||
},
|
||||
orderBy: {
|
||||
startTime: "asc",
|
||||
},
|
||||
return await startActiveSpan("getTraceSummary", async (span) => {
|
||||
const events = await this.readReplica.taskEvent.findMany({
|
||||
select: {
|
||||
id: true,
|
||||
spanId: true,
|
||||
parentId: true,
|
||||
runId: true,
|
||||
idempotencyKey: true,
|
||||
message: true,
|
||||
style: true,
|
||||
startTime: true,
|
||||
duration: true,
|
||||
isError: true,
|
||||
isPartial: true,
|
||||
isCancelled: true,
|
||||
level: true,
|
||||
events: true,
|
||||
environmentType: true,
|
||||
},
|
||||
where: {
|
||||
traceId,
|
||||
},
|
||||
orderBy: {
|
||||
startTime: "asc",
|
||||
},
|
||||
take: env.MAXIMUM_TRACE_SUMMARY_VIEW_COUNT,
|
||||
});
|
||||
|
||||
let preparedEvents: Array<PreparedEvent> = [];
|
||||
let rootSpanId: string | undefined;
|
||||
const eventsBySpanId = new Map<string, PreparedEvent>();
|
||||
|
||||
for (const event of events) {
|
||||
preparedEvents.push(prepareEvent(event));
|
||||
|
||||
if (!rootSpanId && !event.parentId) {
|
||||
rootSpanId = event.spanId;
|
||||
}
|
||||
}
|
||||
|
||||
for (const event of preparedEvents) {
|
||||
const existingEvent = eventsBySpanId.get(event.spanId);
|
||||
|
||||
if (!existingEvent) {
|
||||
eventsBySpanId.set(event.spanId, event);
|
||||
continue;
|
||||
}
|
||||
|
||||
if (event.isCancelled || !event.isPartial) {
|
||||
eventsBySpanId.set(event.spanId, event);
|
||||
}
|
||||
}
|
||||
|
||||
preparedEvents = Array.from(eventsBySpanId.values());
|
||||
|
||||
const spansBySpanId = new Map<string, SpanSummary>();
|
||||
|
||||
const spans = preparedEvents.map((event) => {
|
||||
const ancestorCancelled = isAncestorCancelled(eventsBySpanId, event.spanId);
|
||||
const duration = calculateDurationIfAncestorIsCancelled(
|
||||
eventsBySpanId,
|
||||
event.spanId,
|
||||
event.duration
|
||||
);
|
||||
|
||||
const span = {
|
||||
recordId: event.id,
|
||||
id: event.spanId,
|
||||
parentId: event.parentId ?? undefined,
|
||||
runId: event.runId,
|
||||
idempotencyKey: event.idempotencyKey,
|
||||
data: {
|
||||
message: event.message,
|
||||
style: event.style,
|
||||
duration,
|
||||
isError: event.isError,
|
||||
isPartial: ancestorCancelled ? false : event.isPartial,
|
||||
isCancelled: event.isCancelled === true ? true : event.isPartial && ancestorCancelled,
|
||||
startTime: getDateFromNanoseconds(event.startTime),
|
||||
level: event.level,
|
||||
events: event.events,
|
||||
environmentType: event.environmentType,
|
||||
},
|
||||
};
|
||||
|
||||
spansBySpanId.set(event.spanId, span);
|
||||
|
||||
return span;
|
||||
});
|
||||
|
||||
if (!rootSpanId) {
|
||||
return;
|
||||
}
|
||||
|
||||
const rootSpan = spansBySpanId.get(rootSpanId);
|
||||
|
||||
if (!rootSpan) {
|
||||
return;
|
||||
}
|
||||
|
||||
return {
|
||||
rootSpan,
|
||||
spans,
|
||||
};
|
||||
});
|
||||
}
|
||||
|
||||
const preparedEvents = removeDuplicateEvents(events.map(prepareEvent));
|
||||
public async getRunEvents(runId: string): Promise<RunPreparedEvent[]> {
|
||||
return await startActiveSpan("getRunEvents", async (span) => {
|
||||
const events = await this.readReplica.taskEvent.findMany({
|
||||
select: {
|
||||
id: true,
|
||||
spanId: true,
|
||||
parentId: true,
|
||||
runId: true,
|
||||
idempotencyKey: true,
|
||||
message: true,
|
||||
style: true,
|
||||
startTime: true,
|
||||
duration: true,
|
||||
isError: true,
|
||||
isPartial: true,
|
||||
isCancelled: true,
|
||||
level: true,
|
||||
events: true,
|
||||
environmentType: true,
|
||||
taskSlug: true,
|
||||
},
|
||||
where: {
|
||||
runId,
|
||||
isPartial: false,
|
||||
},
|
||||
orderBy: {
|
||||
startTime: "asc",
|
||||
},
|
||||
});
|
||||
|
||||
const spans = preparedEvents.map((event) => {
|
||||
const ancestorCancelled = isAncestorCancelled(preparedEvents, event.spanId);
|
||||
const duration = calculateDurationIfAncestorIsCancelled(
|
||||
preparedEvents,
|
||||
event.spanId,
|
||||
event.duration
|
||||
let preparedEvents: Array<PreparedEvent> = [];
|
||||
|
||||
for (const event of events) {
|
||||
preparedEvents.push(prepareEvent(event));
|
||||
}
|
||||
|
||||
return preparedEvents;
|
||||
});
|
||||
}
|
||||
|
||||
// A Span can be cancelled if it is partial and has a parent that is cancelled
|
||||
// And a span's duration, if it is partial and has a cancelled parent, is the time between the start of the span and the time of the cancellation event of the parent
|
||||
public async getSpan(spanId: string, traceId: string) {
|
||||
return await startActiveSpan("getSpan", async (s) => {
|
||||
const spanEvent = await this.#getSpanEvent(spanId);
|
||||
|
||||
if (!spanEvent) {
|
||||
return;
|
||||
}
|
||||
|
||||
const preparedEvent = prepareEvent(spanEvent);
|
||||
|
||||
const span = await this.#createSpanFromEvent(preparedEvent);
|
||||
|
||||
const output = rehydrateJson(spanEvent.output);
|
||||
const payload = rehydrateJson(spanEvent.payload);
|
||||
|
||||
const show = rehydrateShow(spanEvent.properties);
|
||||
|
||||
const properties = sanitizedAttributes(spanEvent.properties);
|
||||
|
||||
const messagingEvent = SpanMessagingEvent.optional().safeParse(
|
||||
(properties as any)?.messaging
|
||||
);
|
||||
|
||||
const links: SpanLink[] = [];
|
||||
|
||||
if (messagingEvent.success && messagingEvent.data) {
|
||||
if (messagingEvent.data.message && "id" in messagingEvent.data.message) {
|
||||
if (messagingEvent.data.message.id.startsWith("run_")) {
|
||||
links.push({
|
||||
type: "run",
|
||||
icon: "runs",
|
||||
title: `Run ${messagingEvent.data.message.id}`,
|
||||
runId: messagingEvent.data.message.id,
|
||||
});
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
const backLinks = spanEvent.links as any as Link[] | undefined;
|
||||
|
||||
if (backLinks && backLinks.length > 0) {
|
||||
backLinks.forEach((l) => {
|
||||
const title = String(
|
||||
l.attributes?.[SemanticInternalAttributes.LINK_TITLE] ?? "Triggered by"
|
||||
);
|
||||
|
||||
links.push({
|
||||
type: "span",
|
||||
icon: "trigger",
|
||||
title,
|
||||
traceId: l.context.traceId,
|
||||
spanId: l.context.spanId,
|
||||
});
|
||||
});
|
||||
}
|
||||
|
||||
const spanEvents = transformEvents(
|
||||
preparedEvent.events,
|
||||
spanEvent.metadata as Attributes,
|
||||
spanEvent.environmentType === "DEVELOPMENT"
|
||||
);
|
||||
|
||||
return {
|
||||
...spanEvent,
|
||||
...span.data,
|
||||
payload,
|
||||
output,
|
||||
properties,
|
||||
events: spanEvents,
|
||||
show,
|
||||
links,
|
||||
};
|
||||
});
|
||||
}
|
||||
|
||||
async #createSpanFromEvent(event: PreparedEvent) {
|
||||
return await startActiveSpan("createSpanFromEvent", async (s) => {
|
||||
let ancestorCancelled = false;
|
||||
let duration = event.duration;
|
||||
|
||||
if (!event.isCancelled && event.isPartial) {
|
||||
await this.#walkSpanAncestors(event, (ancestorEvent, level) => {
|
||||
if (level >= 8) {
|
||||
return { stop: true };
|
||||
}
|
||||
|
||||
if (ancestorEvent.isCancelled) {
|
||||
ancestorCancelled = true;
|
||||
|
||||
// We need to get the cancellation time from the cancellation span event
|
||||
const cancellationEvent = ancestorEvent.events.find(
|
||||
(event) => event.name === "cancellation"
|
||||
);
|
||||
|
||||
if (cancellationEvent) {
|
||||
duration = calculateDurationFromStart(event.startTime, cancellationEvent.time);
|
||||
}
|
||||
|
||||
return { stop: true };
|
||||
}
|
||||
|
||||
return { stop: false };
|
||||
});
|
||||
}
|
||||
|
||||
const span = {
|
||||
recordId: event.id,
|
||||
id: event.spanId,
|
||||
parentId: event.parentId ?? undefined,
|
||||
@@ -429,104 +652,93 @@ export class EventRepository {
|
||||
environmentType: event.environmentType,
|
||||
},
|
||||
};
|
||||
|
||||
return span;
|
||||
});
|
||||
|
||||
const rootSpanId = events.find((event) => !event.parentId);
|
||||
if (!rootSpanId) {
|
||||
return;
|
||||
}
|
||||
|
||||
const rootSpan = spans.find((span) => span.id === rootSpanId.spanId);
|
||||
|
||||
if (!rootSpan) {
|
||||
return;
|
||||
}
|
||||
|
||||
return {
|
||||
rootSpan,
|
||||
spans,
|
||||
};
|
||||
}
|
||||
|
||||
// A Span can be cancelled if it is partial and has a parent that is cancelled
|
||||
// And a span's duration, if it is partial and has a cancelled parent, is the time between the start of the span and the time of the cancellation event of the parent
|
||||
public async getSpan(spanId: string, traceId: string) {
|
||||
const traceSummary = await this.getTraceSummary(traceId);
|
||||
|
||||
const span = traceSummary?.spans.find((span) => span.id === spanId);
|
||||
|
||||
if (!span) {
|
||||
async #walkSpanAncestors(
|
||||
event: PreparedEvent,
|
||||
callback: (event: PreparedEvent, level: number) => { stop: boolean }
|
||||
) {
|
||||
const parentId = event.parentId;
|
||||
if (!parentId) {
|
||||
return;
|
||||
}
|
||||
|
||||
const fullEvent = await this.readReplica.taskEvent.findUnique({
|
||||
where: {
|
||||
id: span.recordId,
|
||||
},
|
||||
});
|
||||
await startActiveSpan("walkSpanAncestors", async (s) => {
|
||||
let parentEvent = await this.#getSpanEvent(parentId);
|
||||
let level = 1;
|
||||
|
||||
if (!fullEvent) {
|
||||
return;
|
||||
}
|
||||
while (parentEvent) {
|
||||
const preparedParentEvent = prepareEvent(parentEvent);
|
||||
|
||||
const output = rehydrateJson(fullEvent.output);
|
||||
const payload = rehydrateJson(fullEvent.payload);
|
||||
const result = callback(preparedParentEvent, level);
|
||||
|
||||
const show = rehydrateShow(fullEvent.properties);
|
||||
|
||||
const properties = sanitizedAttributes(fullEvent.properties);
|
||||
|
||||
const messagingEvent = SpanMessagingEvent.optional().safeParse((properties as any)?.messaging);
|
||||
|
||||
const links: SpanLink[] = [];
|
||||
|
||||
if (messagingEvent.success && messagingEvent.data) {
|
||||
if (messagingEvent.data.message && "id" in messagingEvent.data.message) {
|
||||
if (messagingEvent.data.message.id.startsWith("run_")) {
|
||||
links.push({
|
||||
type: "run",
|
||||
icon: "runs",
|
||||
title: `Run ${messagingEvent.data.message.id}`,
|
||||
runId: messagingEvent.data.message.id,
|
||||
});
|
||||
if (result.stop) {
|
||||
return;
|
||||
}
|
||||
|
||||
if (!preparedParentEvent.parentId) {
|
||||
return;
|
||||
}
|
||||
|
||||
parentEvent = await this.#getSpanEvent(preparedParentEvent.parentId);
|
||||
|
||||
level++;
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
async #getSpanAncestors(event: PreparedEvent, levels = 1): Promise<Array<PreparedEvent>> {
|
||||
if (levels >= 8) {
|
||||
return [];
|
||||
}
|
||||
|
||||
const backLinks = fullEvent.links as any as Link[] | undefined;
|
||||
if (!event.parentId) {
|
||||
return [];
|
||||
}
|
||||
|
||||
if (backLinks && backLinks.length > 0) {
|
||||
backLinks.forEach((l) => {
|
||||
const title = String(
|
||||
l.attributes?.[SemanticInternalAttributes.LINK_TITLE] ?? "Triggered by"
|
||||
);
|
||||
const parentEvent = await this.#getSpanEvent(event.parentId);
|
||||
|
||||
links.push({
|
||||
type: "span",
|
||||
icon: "trigger",
|
||||
title,
|
||||
traceId: l.context.traceId,
|
||||
spanId: l.context.spanId,
|
||||
});
|
||||
if (!parentEvent) {
|
||||
return [];
|
||||
}
|
||||
|
||||
const preparedParentEvent = prepareEvent(parentEvent);
|
||||
|
||||
if (!preparedParentEvent.parentId) {
|
||||
return [preparedParentEvent];
|
||||
}
|
||||
|
||||
const moreAncestors = await this.#getSpanAncestors(preparedParentEvent, levels + 1);
|
||||
|
||||
return [preparedParentEvent, ...moreAncestors];
|
||||
}
|
||||
|
||||
async #getSpanEvent(spanId: string) {
|
||||
return await startActiveSpan("getSpanEvent", async (s) => {
|
||||
const events = await this.readReplica.taskEvent.findMany({
|
||||
where: {
|
||||
spanId,
|
||||
},
|
||||
orderBy: {
|
||||
startTime: "asc",
|
||||
},
|
||||
});
|
||||
}
|
||||
|
||||
const events = transformEvents(
|
||||
span.data.events,
|
||||
fullEvent.metadata as Attributes,
|
||||
traceSummary?.rootSpan.data.environmentType === "DEVELOPMENT"
|
||||
);
|
||||
let finalEvent: TaskEvent | undefined;
|
||||
|
||||
return {
|
||||
...fullEvent,
|
||||
...span.data,
|
||||
payload,
|
||||
output,
|
||||
properties,
|
||||
events,
|
||||
show,
|
||||
links,
|
||||
};
|
||||
for (const event of events) {
|
||||
if (event.isPartial && finalEvent) {
|
||||
continue;
|
||||
}
|
||||
|
||||
finalEvent = event;
|
||||
}
|
||||
|
||||
return finalEvent;
|
||||
});
|
||||
}
|
||||
|
||||
public async recordEvent(message: string, options: TraceEventOptions) {
|
||||
@@ -973,34 +1185,38 @@ function prepareEvent(event: QueriedEvent): PreparedEvent {
|
||||
}
|
||||
|
||||
function parseEventsField(events: Prisma.JsonValue): SpanEvents {
|
||||
const eventsUnflattened = events
|
||||
const unsafe = events
|
||||
? (events as any[]).map((e) => ({
|
||||
...e,
|
||||
properties: unflattenAttributes(e.properties as Attributes),
|
||||
}))
|
||||
: undefined;
|
||||
|
||||
const spanEvents = SpanEvents.safeParse(eventsUnflattened);
|
||||
|
||||
if (spanEvents.success) {
|
||||
return spanEvents.data;
|
||||
}
|
||||
|
||||
return [];
|
||||
return unsafe as SpanEvents;
|
||||
}
|
||||
|
||||
function parseStyleField(style: Prisma.JsonValue): TaskEventStyle {
|
||||
const parsedStyle = TaskEventStyle.safeParse(unflattenAttributes(style as Attributes));
|
||||
const unsafe = unflattenAttributes(style as Attributes);
|
||||
|
||||
if (parsedStyle.success) {
|
||||
return parsedStyle.data;
|
||||
if (!unsafe) {
|
||||
return {};
|
||||
}
|
||||
|
||||
if (typeof unsafe === "object") {
|
||||
return Object.assign(
|
||||
{
|
||||
icon: undefined,
|
||||
variant: undefined,
|
||||
},
|
||||
unsafe
|
||||
) as TaskEventStyle;
|
||||
}
|
||||
|
||||
return {};
|
||||
}
|
||||
|
||||
function isAncestorCancelled(events: PreparedEvent[], spanId: string) {
|
||||
const event = events.find((event) => event.spanId === spanId);
|
||||
function isAncestorCancelled(events: Map<string, PreparedEvent>, spanId: string) {
|
||||
const event = events.get(spanId);
|
||||
|
||||
if (!event) {
|
||||
return false;
|
||||
@@ -1018,11 +1234,11 @@ function isAncestorCancelled(events: PreparedEvent[], spanId: string) {
|
||||
}
|
||||
|
||||
function calculateDurationIfAncestorIsCancelled(
|
||||
events: PreparedEvent[],
|
||||
events: Map<string, PreparedEvent>,
|
||||
spanId: string,
|
||||
defaultDuration: number
|
||||
) {
|
||||
const event = events.find((event) => event.spanId === spanId);
|
||||
const event = events.get(spanId);
|
||||
|
||||
if (!event) {
|
||||
return defaultDuration;
|
||||
@@ -1054,8 +1270,9 @@ function calculateDurationIfAncestorIsCancelled(
|
||||
return defaultDuration;
|
||||
}
|
||||
|
||||
function findFirstCancelledAncestor(events: PreparedEvent[], spanId: string) {
|
||||
const event = events.find((event) => event.spanId === spanId);
|
||||
function findFirstCancelledAncestor(events: Map<string, PreparedEvent>, spanId: string) {
|
||||
const event = events.get(spanId);
|
||||
|
||||
if (!event) {
|
||||
return;
|
||||
}
|
||||
@@ -1196,7 +1413,7 @@ function getNowInNanoseconds(): bigint {
|
||||
return BigInt(new Date().getTime() * 1_000_000);
|
||||
}
|
||||
|
||||
function getDateFromNanoseconds(nanoseconds: bigint) {
|
||||
export function getDateFromNanoseconds(nanoseconds: bigint) {
|
||||
return new Date(Number(nanoseconds) / 1_000_000);
|
||||
}
|
||||
|
||||
|
||||
@@ -138,8 +138,19 @@ function createCoordinatorNamespace(io: Server) {
|
||||
await sharedQueueTasks.taskRunHeartbeat(message.runId);
|
||||
},
|
||||
CHECKPOINT_CREATED: async (message) => {
|
||||
const createCheckpoint = new CreateCheckpointService();
|
||||
await createCheckpoint.call(message);
|
||||
try {
|
||||
const createCheckpoint = new CreateCheckpointService();
|
||||
const result = await createCheckpoint.call(message);
|
||||
|
||||
return { keepRunAlive: result?.keepRunAlive ?? false };
|
||||
} catch (error) {
|
||||
logger.error("Error while creating checkpoint", {
|
||||
rawMessage: message,
|
||||
error: error instanceof Error ? error.message : error,
|
||||
});
|
||||
|
||||
return { keepRunAlive: false };
|
||||
}
|
||||
},
|
||||
CREATE_WORKER: async (message) => {
|
||||
try {
|
||||
|
||||
@@ -1625,7 +1625,7 @@ function getMarQSClient() {
|
||||
defaultEnvConcurrency: env.DEFAULT_ENV_EXECUTION_CONCURRENCY_LIMIT,
|
||||
defaultOrgConcurrency: env.DEFAULT_ORG_EXECUTION_CONCURRENCY_LIMIT,
|
||||
visibilityTimeoutInMs: 120 * 1000, // 2 minutes,
|
||||
enableRebalancing: !env.MARQS_DISABLE_REBALANCING,
|
||||
enableRebalancing: false,
|
||||
});
|
||||
} else {
|
||||
console.warn(
|
||||
|
||||
@@ -17,7 +17,6 @@ import {
|
||||
BackgroundWorkerTask,
|
||||
RuntimeEnvironment,
|
||||
TaskRun,
|
||||
TaskRunAttemptStatus,
|
||||
TaskRunStatus,
|
||||
} from "@trigger.dev/database";
|
||||
import { z } from "zod";
|
||||
@@ -43,6 +42,7 @@ import { generateJWTTokenForEnvironment } from "~/services/apiAuth.server";
|
||||
import { EnvironmentVariable } from "../environmentVariables/repository";
|
||||
import { machinePresetFromConfig } from "../machinePresets.server";
|
||||
import { env } from "~/env.server";
|
||||
import { isFinalAttemptStatus, isFinalRunStatus } from "../taskStatus";
|
||||
|
||||
const WithTraceContext = z.object({
|
||||
traceparent: z.string().optional(),
|
||||
@@ -962,19 +962,7 @@ class SharedQueueTasks {
|
||||
}
|
||||
|
||||
if (setToExecuting) {
|
||||
const FINAL_RUN_STATUSES: TaskRunStatus[] = [
|
||||
"CANCELED",
|
||||
"COMPLETED_SUCCESSFULLY",
|
||||
"COMPLETED_WITH_ERRORS",
|
||||
"INTERRUPTED",
|
||||
"SYSTEM_FAILURE",
|
||||
];
|
||||
const FINAL_ATTEMPT_STATUSES: TaskRunAttemptStatus[] = ["CANCELED", "COMPLETED", "FAILED"];
|
||||
|
||||
if (
|
||||
FINAL_ATTEMPT_STATUSES.includes(attempt.status) ||
|
||||
FINAL_RUN_STATUSES.includes(attempt.taskRun.status)
|
||||
) {
|
||||
if (isFinalAttemptStatus(attempt.status) || isFinalRunStatus(attempt.taskRun.status)) {
|
||||
logger.error("Status already in final state", {
|
||||
attempt: {
|
||||
id: attempt.id,
|
||||
|
||||
@@ -82,7 +82,7 @@ function getMarQSClient() {
|
||||
defaultEnvConcurrency: env.V2_MARQS_DEFAULT_ENV_CONCURRENCY, // this is so we aren't limited by the environment concurrency
|
||||
defaultOrgConcurrency: env.DEFAULT_ORG_EXECUTION_CONCURRENCY_LIMIT,
|
||||
visibilityTimeoutInMs: env.V2_MARQS_VISIBILITY_TIMEOUT_MS, // 15 minutes
|
||||
enableRebalancing: env.V2_MARQS_CONSUMER_POOL_ENABLED === "1",
|
||||
enableRebalancing: false,
|
||||
});
|
||||
}
|
||||
|
||||
|
||||
@@ -3,6 +3,7 @@ import { env } from "~/env.server";
|
||||
import { AuthenticatedEnvironment } from "~/services/apiAuth.server";
|
||||
import { logger } from "~/services/logger.server";
|
||||
import { singleton } from "~/utils/singleton";
|
||||
import { startActiveSpan } from "./tracer.server";
|
||||
|
||||
export const r2 = singleton("r2", initializeR2);
|
||||
|
||||
@@ -23,30 +24,87 @@ export async function uploadToObjectStore(
|
||||
contentType: string,
|
||||
environment: AuthenticatedEnvironment
|
||||
): Promise<string> {
|
||||
if (!r2) {
|
||||
throw new Error("Object store credentials are not set");
|
||||
return await startActiveSpan("uploadToObjectStore()", async (span) => {
|
||||
if (!r2) {
|
||||
throw new Error("Object store credentials are not set");
|
||||
}
|
||||
|
||||
if (!env.OBJECT_STORE_BASE_URL) {
|
||||
throw new Error("Object store base URL is not set");
|
||||
}
|
||||
|
||||
span.setAttributes({
|
||||
projectRef: environment.project.externalRef,
|
||||
environmentSlug: environment.slug,
|
||||
filename: filename,
|
||||
});
|
||||
|
||||
const url = new URL(env.OBJECT_STORE_BASE_URL);
|
||||
url.pathname = `/packets/${environment.project.externalRef}/${environment.slug}/${filename}`;
|
||||
|
||||
logger.debug("Uploading to object store", { url: url.href });
|
||||
|
||||
const response = await r2.fetch(url.toString(), {
|
||||
method: "PUT",
|
||||
headers: {
|
||||
"Content-Type": contentType,
|
||||
},
|
||||
body: data,
|
||||
});
|
||||
|
||||
if (!response.ok) {
|
||||
throw new Error(`Failed to upload output to ${url}: ${response.statusText}`);
|
||||
}
|
||||
|
||||
return url.href;
|
||||
});
|
||||
}
|
||||
|
||||
export async function generatePresignedRequest(
|
||||
projectRef: string,
|
||||
envSlug: string,
|
||||
filename: string,
|
||||
method: "PUT" | "GET" = "PUT"
|
||||
) {
|
||||
if (!env.OBJECT_STORE_BASE_URL) {
|
||||
return;
|
||||
}
|
||||
|
||||
if (!env.OBJECT_STORE_BASE_URL) {
|
||||
throw new Error("Object store base URL is not set");
|
||||
if (!r2) {
|
||||
return;
|
||||
}
|
||||
|
||||
const url = new URL(env.OBJECT_STORE_BASE_URL);
|
||||
url.pathname = `/packets/${environment.project.externalRef}/${environment.slug}/${filename}`;
|
||||
url.pathname = `/packets/${projectRef}/${envSlug}/${filename}`;
|
||||
url.searchParams.set("X-Amz-Expires", "300"); // 5 minutes
|
||||
|
||||
logger.debug("Uploading to object store", { url: url.href });
|
||||
const signed = await r2.sign(
|
||||
new Request(url, {
|
||||
method,
|
||||
}),
|
||||
{
|
||||
aws: { signQuery: true },
|
||||
}
|
||||
);
|
||||
|
||||
const response = await r2.fetch(url.toString(), {
|
||||
method: "PUT",
|
||||
headers: {
|
||||
"Content-Type": contentType,
|
||||
},
|
||||
body: data,
|
||||
logger.debug("Generated presigned URL", {
|
||||
url: signed.url,
|
||||
headers: Object.fromEntries(signed.headers),
|
||||
projectRef,
|
||||
envSlug,
|
||||
filename,
|
||||
});
|
||||
|
||||
if (!response.ok) {
|
||||
throw new Error(`Failed to upload output to ${url}: ${response.statusText}`);
|
||||
}
|
||||
|
||||
return url.href;
|
||||
return signed;
|
||||
}
|
||||
|
||||
export async function generatePresignedUrl(
|
||||
projectRef: string,
|
||||
envSlug: string,
|
||||
filename: string,
|
||||
method: "PUT" | "GET" = "PUT"
|
||||
) {
|
||||
const signed = await generatePresignedRequest(projectRef, envSlug, filename, method);
|
||||
|
||||
return signed?.url;
|
||||
}
|
||||
|
||||
@@ -6,7 +6,7 @@ import { logger } from "~/services/logger.server";
|
||||
|
||||
import { PrismaClientOrTransaction, prisma } from "~/db.server";
|
||||
import { ResumeTaskRunDependenciesService } from "./resumeTaskRunDependencies.server";
|
||||
import { CANCELLABLE_STATUSES } from "./cancelTaskRun.server";
|
||||
import { isCancellableRunStatus } from "../taskStatus";
|
||||
|
||||
export class CancelAttemptService extends BaseService {
|
||||
public async call(
|
||||
@@ -55,7 +55,7 @@ export class CancelAttemptService extends BaseService {
|
||||
taskRun: {
|
||||
update: {
|
||||
data: {
|
||||
status: CANCELLABLE_STATUSES.includes(taskRunAttempt.taskRun.status)
|
||||
status: isCancellableRunStatus(taskRunAttempt.taskRun.status)
|
||||
? "INTERRUPTED"
|
||||
: undefined,
|
||||
},
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
import { Prisma, TaskRun, TaskRunAttemptStatus, TaskRunStatus } from "@trigger.dev/database";
|
||||
import { Prisma, TaskRun } from "@trigger.dev/database";
|
||||
import assertNever from "assert-never";
|
||||
import { logger } from "~/services/logger.server";
|
||||
import { marqs } from "~/v3/marqs/index.server";
|
||||
@@ -7,22 +7,7 @@ import { socketIo } from "../handleSocketIo.server";
|
||||
import { devPubSub } from "../marqs/devPubSub.server";
|
||||
import { BaseService } from "./baseService.server";
|
||||
import { CancelAttemptService } from "./cancelAttempt.server";
|
||||
|
||||
export const CANCELLABLE_STATUSES: Array<TaskRunStatus> = [
|
||||
"PENDING",
|
||||
"WAITING_FOR_DEPLOY",
|
||||
"EXECUTING",
|
||||
"PAUSED",
|
||||
"WAITING_TO_RESUME",
|
||||
"PAUSED",
|
||||
"RETRYING_AFTER_FAILURE",
|
||||
];
|
||||
|
||||
const CANCELLABLE_ATTEMPT_STATUSES: Array<TaskRunAttemptStatus> = [
|
||||
"EXECUTING",
|
||||
"PAUSED",
|
||||
"PENDING",
|
||||
];
|
||||
import { CANCELLABLE_ATTEMPT_STATUSES, isCancellableRunStatus } from "../taskStatus";
|
||||
|
||||
type ExtendedTaskRun = Prisma.TaskRunGetPayload<{
|
||||
include: {
|
||||
@@ -53,7 +38,11 @@ export class CancelTaskRunService extends BaseService {
|
||||
};
|
||||
|
||||
// Make sure the task run is in a cancellable state
|
||||
if (!CANCELLABLE_STATUSES.includes(taskRun.status)) {
|
||||
if (!isCancellableRunStatus(taskRun.status)) {
|
||||
logger.error("Task run is not in a cancellable state", {
|
||||
runId: taskRun.id,
|
||||
status: taskRun.status,
|
||||
});
|
||||
return;
|
||||
}
|
||||
|
||||
|
||||
@@ -1,27 +1,11 @@
|
||||
import {
|
||||
TaskRun,
|
||||
TaskRunAttempt,
|
||||
TaskRunAttemptStatus,
|
||||
TaskRunStatus,
|
||||
} from "@trigger.dev/database";
|
||||
import { TaskRun, TaskRunAttempt } from "@trigger.dev/database";
|
||||
import { eventRepository } from "../eventRepository.server";
|
||||
import { marqs } from "~/v3/marqs/index.server";
|
||||
import { BaseService } from "./baseService.server";
|
||||
import { logger } from "~/services/logger.server";
|
||||
import { AuthenticatedEnvironment } from "~/services/apiAuth.server";
|
||||
import { ResumeTaskRunDependenciesService } from "./resumeTaskRunDependencies.server";
|
||||
|
||||
export const CRASHABLE_RUN_STATUSES: Array<TaskRunStatus> = [
|
||||
"PENDING",
|
||||
"WAITING_FOR_DEPLOY",
|
||||
"EXECUTING",
|
||||
"PAUSED",
|
||||
"WAITING_TO_RESUME",
|
||||
"PAUSED",
|
||||
"RETRYING_AFTER_FAILURE",
|
||||
];
|
||||
|
||||
const CRASHABLE_ATTEMPT_STATUSES: Array<TaskRunAttemptStatus> = ["EXECUTING", "PAUSED", "PENDING"];
|
||||
import { CRASHABLE_ATTEMPT_STATUSES, isCrashableRunStatus } from "../taskStatus";
|
||||
|
||||
export type CrashTaskRunServiceOptions = {
|
||||
reason?: string;
|
||||
@@ -52,7 +36,8 @@ export class CrashTaskRunService extends BaseService {
|
||||
}
|
||||
|
||||
// Make sure the task run is in a crashable state
|
||||
if (!CRASHABLE_RUN_STATUSES.includes(taskRun.status)) {
|
||||
if (!isCrashableRunStatus(taskRun.status)) {
|
||||
logger.error("Task run is not in a crashable state", { runId, status: taskRun.status });
|
||||
return;
|
||||
}
|
||||
|
||||
|
||||
@@ -1,19 +1,13 @@
|
||||
import { CoordinatorToPlatformMessages } from "@trigger.dev/core/v3";
|
||||
import type { InferSocketMessageSchema } from "@trigger.dev/core/v3/zodSocket";
|
||||
import type {
|
||||
CheckpointRestoreEvent,
|
||||
TaskRunAttemptStatus,
|
||||
TaskRunStatus,
|
||||
} from "@trigger.dev/database";
|
||||
import type { Checkpoint, CheckpointRestoreEvent } from "@trigger.dev/database";
|
||||
import { logger } from "~/services/logger.server";
|
||||
import { generateFriendlyId } from "../friendlyIdentifiers";
|
||||
import { marqs } from "~/v3/marqs/index.server";
|
||||
import { CreateCheckpointRestoreEventService } from "./createCheckpointRestoreEvent.server";
|
||||
import { BaseService } from "./baseService.server";
|
||||
import { CrashTaskRunService } from "./crashTaskRun.server";
|
||||
|
||||
const FREEZABLE_RUN_STATUSES: TaskRunStatus[] = ["EXECUTING", "RETRYING_AFTER_FAILURE"];
|
||||
const FREEZABLE_ATTEMPT_STATUSES: TaskRunAttemptStatus[] = ["EXECUTING", "FAILED"];
|
||||
import { isFinalRunStatus, isFreezableAttemptStatus, isFreezableRunStatus } from "../taskStatus";
|
||||
|
||||
export class CreateCheckpointService extends BaseService {
|
||||
public async call(
|
||||
@@ -21,7 +15,14 @@ export class CreateCheckpointService extends BaseService {
|
||||
InferSocketMessageSchema<typeof CoordinatorToPlatformMessages, "CHECKPOINT_CREATED">,
|
||||
"version"
|
||||
>
|
||||
) {
|
||||
): Promise<
|
||||
| {
|
||||
checkpoint: Checkpoint;
|
||||
event: CheckpointRestoreEvent;
|
||||
keepRunAlive: boolean;
|
||||
}
|
||||
| undefined
|
||||
> {
|
||||
logger.debug(`Creating checkpoint`, params);
|
||||
|
||||
const attempt = await this._prisma.taskRunAttempt.findUnique({
|
||||
@@ -49,8 +50,8 @@ export class CreateCheckpointService extends BaseService {
|
||||
}
|
||||
|
||||
if (
|
||||
!FREEZABLE_ATTEMPT_STATUSES.includes(attempt.status) ||
|
||||
!FREEZABLE_RUN_STATUSES.includes(attempt.taskRun.status)
|
||||
!isFreezableAttemptStatus(attempt.status) ||
|
||||
!isFreezableRunStatus(attempt.taskRun.status)
|
||||
) {
|
||||
logger.error("Unfreezable state", {
|
||||
attempt: {
|
||||
@@ -115,7 +116,9 @@ export class CreateCheckpointService extends BaseService {
|
||||
});
|
||||
|
||||
const { reason } = params;
|
||||
|
||||
let checkpointEvent: CheckpointRestoreEvent | undefined;
|
||||
let keepRunAlive = false;
|
||||
|
||||
switch (reason.type) {
|
||||
case "WAIT_FOR_DURATION": {
|
||||
@@ -131,7 +134,12 @@ export class CreateCheckpointService extends BaseService {
|
||||
dependencyFriendlyRunId: reason.friendlyId,
|
||||
});
|
||||
|
||||
await marqs?.acknowledgeMessage(attempt.taskRunId);
|
||||
keepRunAlive = await this.#isRunCompleted(reason.friendlyId);
|
||||
|
||||
if (!keepRunAlive) {
|
||||
await marqs?.acknowledgeMessage(attempt.taskRunId);
|
||||
}
|
||||
|
||||
break;
|
||||
}
|
||||
case "WAIT_FOR_BATCH": {
|
||||
@@ -140,7 +148,12 @@ export class CreateCheckpointService extends BaseService {
|
||||
batchDependencyFriendlyId: reason.batchFriendlyId,
|
||||
});
|
||||
|
||||
await marqs?.acknowledgeMessage(attempt.taskRunId);
|
||||
keepRunAlive = await this.#isBatchCompleted(reason.batchFriendlyId);
|
||||
|
||||
if (!keepRunAlive) {
|
||||
await marqs?.acknowledgeMessage(attempt.taskRunId);
|
||||
}
|
||||
|
||||
break;
|
||||
}
|
||||
case "RETRYING_AFTER_FAILURE": {
|
||||
@@ -180,6 +193,37 @@ export class CreateCheckpointService extends BaseService {
|
||||
return {
|
||||
checkpoint,
|
||||
event: checkpointEvent,
|
||||
keepRunAlive,
|
||||
};
|
||||
}
|
||||
|
||||
async #isBatchCompleted(friendlyId: string): Promise<boolean> {
|
||||
const batch = await this._prisma.batchTaskRun.findUnique({
|
||||
where: {
|
||||
friendlyId,
|
||||
},
|
||||
});
|
||||
|
||||
if (!batch) {
|
||||
logger.error("Batch not found", { friendlyId });
|
||||
return false;
|
||||
}
|
||||
|
||||
return batch.status === "COMPLETED";
|
||||
}
|
||||
|
||||
async #isRunCompleted(friendlyId: string): Promise<boolean> {
|
||||
const run = await this._prisma.taskRun.findUnique({
|
||||
where: {
|
||||
friendlyId,
|
||||
},
|
||||
});
|
||||
|
||||
if (!run) {
|
||||
logger.error("Run not found", { friendlyId });
|
||||
return false;
|
||||
}
|
||||
|
||||
return isFinalRunStatus(run.status);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,12 +1,10 @@
|
||||
import { TaskRunAttemptStatus, TaskRunStatus, type Checkpoint } from "@trigger.dev/database";
|
||||
import { type Checkpoint } from "@trigger.dev/database";
|
||||
import { logger } from "~/services/logger.server";
|
||||
import { socketIo } from "../handleSocketIo.server";
|
||||
import { machinePresetFromConfig } from "../machinePresets.server";
|
||||
import { BaseService } from "./baseService.server";
|
||||
import { CreateCheckpointRestoreEventService } from "./createCheckpointRestoreEvent.server";
|
||||
|
||||
const RESTORABLE_RUN_STATUSES: TaskRunStatus[] = ["WAITING_TO_RESUME"];
|
||||
const RESTORABLE_ATTEMPT_STATUSES: TaskRunAttemptStatus[] = ["PAUSED"];
|
||||
import { isRestorableAttemptStatus, isRestorableRunStatus } from "../taskStatus";
|
||||
|
||||
export class RestoreCheckpointService extends BaseService {
|
||||
public async call(params: {
|
||||
@@ -51,10 +49,7 @@ export class RestoreCheckpointService extends BaseService {
|
||||
|
||||
const checkpoint = checkpointEvent.checkpoint;
|
||||
|
||||
const runIsRestorable = RESTORABLE_RUN_STATUSES.includes(checkpoint.run.status);
|
||||
const attemptIsRestorable = RESTORABLE_ATTEMPT_STATUSES.includes(checkpoint.attempt.status);
|
||||
|
||||
if (!runIsRestorable) {
|
||||
if (!isRestorableRunStatus(checkpoint.run.status)) {
|
||||
logger.error("Run is unrestorable", {
|
||||
eventId: params.eventId,
|
||||
runId: checkpoint.runId,
|
||||
@@ -64,7 +59,7 @@ export class RestoreCheckpointService extends BaseService {
|
||||
return;
|
||||
}
|
||||
|
||||
if (!attemptIsRestorable && !params.isRetry) {
|
||||
if (!isRestorableAttemptStatus(checkpoint.attempt.status) && !params.isRetry) {
|
||||
logger.error("Attempt is unrestorable", {
|
||||
eventId: params.eventId,
|
||||
runId: checkpoint.runId,
|
||||
|
||||
@@ -4,14 +4,15 @@ import {
|
||||
TriggerTaskRequestBody,
|
||||
packetRequiresOffloading,
|
||||
} from "@trigger.dev/core/v3";
|
||||
import { prisma } from "~/db.server";
|
||||
import { AuthenticatedEnvironment } from "~/services/apiAuth.server";
|
||||
import { autoIncrementCounter } from "~/services/autoIncrementCounter.server";
|
||||
import { marqs, sanitizeQueueName } from "~/v3/marqs/index.server";
|
||||
import { eventRepository } from "../eventRepository.server";
|
||||
import { generateFriendlyId } from "../friendlyIdentifiers";
|
||||
import { uploadToObjectStore } from "../r2.server";
|
||||
import { startActiveSpan } from "../tracer.server";
|
||||
import { BaseService } from "./baseService.server";
|
||||
import { env } from "~/env.server";
|
||||
|
||||
export type TriggerTaskServiceOptions = {
|
||||
idempotencyKey?: string;
|
||||
@@ -176,10 +177,12 @@ export class TriggerTaskService extends BaseService {
|
||||
}
|
||||
|
||||
if (body.options?.queue) {
|
||||
const concurrencyLimit = body.options.queue.concurrencyLimit
|
||||
? Math.max(0, body.options.queue.concurrencyLimit)
|
||||
: null;
|
||||
const taskQueue = await prisma.taskQueue.upsert({
|
||||
const concurrencyLimit =
|
||||
typeof body.options.queue.concurrencyLimit === "number"
|
||||
? Math.max(0, body.options.queue.concurrencyLimit)
|
||||
: undefined;
|
||||
|
||||
const taskQueue = await tx.taskQueue.upsert({
|
||||
where: {
|
||||
runtimeEnvironmentId_name: {
|
||||
runtimeEnvironmentId: environment.id,
|
||||
@@ -255,26 +258,31 @@ export class TriggerTaskService extends BaseService {
|
||||
pathPrefix: string,
|
||||
environment: AuthenticatedEnvironment
|
||||
) {
|
||||
const packet = this.#createPayloadPacket(payload, payloadType);
|
||||
return await startActiveSpan("handlePayloadPacket()", async (span) => {
|
||||
const packet = this.#createPayloadPacket(payload, payloadType);
|
||||
|
||||
if (!packet.data) {
|
||||
return packet;
|
||||
}
|
||||
if (!packet.data) {
|
||||
return packet;
|
||||
}
|
||||
|
||||
const { needsOffloading, size } = packetRequiresOffloading(packet);
|
||||
const { needsOffloading, size } = packetRequiresOffloading(
|
||||
packet,
|
||||
env.TASK_PAYLOAD_OFFLOAD_THRESHOLD
|
||||
);
|
||||
|
||||
if (!needsOffloading) {
|
||||
return packet;
|
||||
}
|
||||
if (!needsOffloading) {
|
||||
return packet;
|
||||
}
|
||||
|
||||
const filename = `${pathPrefix}/payload.json`;
|
||||
const filename = `${pathPrefix}/payload.json`;
|
||||
|
||||
await uploadToObjectStore(filename, packet.data, packet.dataType, environment);
|
||||
await uploadToObjectStore(filename, packet.data, packet.dataType, environment);
|
||||
|
||||
return {
|
||||
data: filename,
|
||||
dataType: "application/store",
|
||||
};
|
||||
return {
|
||||
data: filename,
|
||||
dataType: "application/store",
|
||||
};
|
||||
});
|
||||
}
|
||||
|
||||
#createPayloadPacket(payload: any, payloadType: string): IOPacket {
|
||||
|
||||
@@ -0,0 +1,69 @@
|
||||
import type { TaskRunAttemptStatus, TaskRunStatus } from "@trigger.dev/database";
|
||||
|
||||
export const CANCELLABLE_RUN_STATUSES: TaskRunStatus[] = [
|
||||
"PENDING",
|
||||
"WAITING_FOR_DEPLOY",
|
||||
"EXECUTING",
|
||||
"PAUSED",
|
||||
"WAITING_TO_RESUME",
|
||||
"PAUSED",
|
||||
"RETRYING_AFTER_FAILURE",
|
||||
];
|
||||
export const CANCELLABLE_ATTEMPT_STATUSES: TaskRunAttemptStatus[] = [
|
||||
"EXECUTING",
|
||||
"PAUSED",
|
||||
"PENDING",
|
||||
];
|
||||
|
||||
export function isCancellableRunStatus(status: TaskRunStatus): boolean {
|
||||
return CANCELLABLE_RUN_STATUSES.includes(status);
|
||||
}
|
||||
export function isCancellableAttemptStatus(status: TaskRunAttemptStatus): boolean {
|
||||
return CANCELLABLE_ATTEMPT_STATUSES.includes(status);
|
||||
}
|
||||
|
||||
export const CRASHABLE_RUN_STATUSES: TaskRunStatus[] = CANCELLABLE_RUN_STATUSES;
|
||||
export const CRASHABLE_ATTEMPT_STATUSES: TaskRunAttemptStatus[] = CANCELLABLE_ATTEMPT_STATUSES;
|
||||
|
||||
export function isCrashableRunStatus(status: TaskRunStatus): boolean {
|
||||
return CRASHABLE_RUN_STATUSES.includes(status);
|
||||
}
|
||||
export function isCrashableAttemptStatus(status: TaskRunAttemptStatus): boolean {
|
||||
return CRASHABLE_ATTEMPT_STATUSES.includes(status);
|
||||
}
|
||||
|
||||
export const FINAL_RUN_STATUSES: TaskRunStatus[] = [
|
||||
"CANCELED",
|
||||
"COMPLETED_SUCCESSFULLY",
|
||||
"COMPLETED_WITH_ERRORS",
|
||||
"INTERRUPTED",
|
||||
"SYSTEM_FAILURE",
|
||||
];
|
||||
export const FINAL_ATTEMPT_STATUSES: TaskRunAttemptStatus[] = ["CANCELED", "COMPLETED", "FAILED"];
|
||||
|
||||
export const FREEZABLE_RUN_STATUSES: TaskRunStatus[] = ["EXECUTING", "RETRYING_AFTER_FAILURE"];
|
||||
export const FREEZABLE_ATTEMPT_STATUSES: TaskRunAttemptStatus[] = ["EXECUTING", "FAILED"];
|
||||
|
||||
export function isFreezableRunStatus(status: TaskRunStatus): boolean {
|
||||
return FREEZABLE_RUN_STATUSES.includes(status);
|
||||
}
|
||||
export function isFreezableAttemptStatus(status: TaskRunAttemptStatus): boolean {
|
||||
return FREEZABLE_ATTEMPT_STATUSES.includes(status);
|
||||
}
|
||||
|
||||
export function isFinalRunStatus(status: TaskRunStatus): boolean {
|
||||
return FINAL_RUN_STATUSES.includes(status);
|
||||
}
|
||||
export function isFinalAttemptStatus(status: TaskRunAttemptStatus): boolean {
|
||||
return FINAL_ATTEMPT_STATUSES.includes(status);
|
||||
}
|
||||
|
||||
export const RESTORABLE_RUN_STATUSES: TaskRunStatus[] = ["WAITING_TO_RESUME"];
|
||||
export const RESTORABLE_ATTEMPT_STATUSES: TaskRunAttemptStatus[] = ["PAUSED"];
|
||||
|
||||
export function isRestorableRunStatus(status: TaskRunStatus): boolean {
|
||||
return RESTORABLE_RUN_STATUSES.includes(status);
|
||||
}
|
||||
export function isRestorableAttemptStatus(status: TaskRunAttemptStatus): boolean {
|
||||
return RESTORABLE_ATTEMPT_STATUSES.includes(status);
|
||||
}
|
||||
@@ -4,7 +4,10 @@ import {
|
||||
DiagConsoleLogger,
|
||||
DiagLogLevel,
|
||||
Link,
|
||||
Span,
|
||||
SpanKind,
|
||||
SpanOptions,
|
||||
SpanStatusCode,
|
||||
diag,
|
||||
trace,
|
||||
} from "@opentelemetry/api";
|
||||
@@ -76,6 +79,32 @@ class CustomWebappSampler implements Sampler {
|
||||
|
||||
export const tracer = singleton("tracer", getTracer);
|
||||
|
||||
export async function startActiveSpan<T>(
|
||||
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 });
|
||||
|
||||
throw error;
|
||||
} finally {
|
||||
span.end();
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
function getTracer() {
|
||||
if (env.INTERNAL_OTEL_TRACE_DISABLED === "1") {
|
||||
console.log(`🔦 Tracer disabled, returning a noop tracer`);
|
||||
@@ -88,7 +117,7 @@ function getTracer() {
|
||||
const samplingRate = 1.0 / Math.max(parseInt(env.INTERNAL_OTEL_TRACE_SAMPLING_RATE, 10), 1);
|
||||
|
||||
const provider = new NodeTracerProvider({
|
||||
forceFlushTimeoutMillis: 500,
|
||||
forceFlushTimeoutMillis: 15_000,
|
||||
resource: new Resource({
|
||||
[SEMRESATTRS_SERVICE_NAME]: env.SERVICE_NAME,
|
||||
}),
|
||||
@@ -100,7 +129,7 @@ function getTracer() {
|
||||
if (env.INTERNAL_OTEL_TRACE_EXPORTER_URL) {
|
||||
const exporter = new OTLPTraceExporter({
|
||||
url: env.INTERNAL_OTEL_TRACE_EXPORTER_URL,
|
||||
timeoutMillis: 10_000,
|
||||
timeoutMillis: 15_000,
|
||||
headers:
|
||||
env.INTERNAL_OTEL_TRACE_EXPORTER_AUTH_HEADER_NAME &&
|
||||
env.INTERNAL_OTEL_TRACE_EXPORTER_AUTH_HEADER_VALUE
|
||||
|
||||
+80
-20
@@ -1,8 +1,14 @@
|
||||
{
|
||||
"$schema": "https://mintlify.com/schema.json",
|
||||
"name": "Trigger.dev",
|
||||
"openapi": ["/openapi.yml", "/v3-openapi.yaml"],
|
||||
"versions": ["v3 (Developer Preview)", "v2"],
|
||||
"openapi": [
|
||||
"/openapi.yml",
|
||||
"/v3-openapi.yaml"
|
||||
],
|
||||
"versions": [
|
||||
"v3 (Developer Preview)",
|
||||
"v2"
|
||||
],
|
||||
"api": {
|
||||
"playground": {
|
||||
"mode": "simple"
|
||||
@@ -96,7 +102,9 @@
|
||||
{
|
||||
"group": "",
|
||||
"version": "v3 (Developer Preview)",
|
||||
"pages": ["v3/introduction"]
|
||||
"pages": [
|
||||
"v3/introduction"
|
||||
]
|
||||
},
|
||||
{
|
||||
"group": "Getting Started",
|
||||
@@ -119,7 +127,10 @@
|
||||
"v3/apikeys",
|
||||
{
|
||||
"group": "Task types",
|
||||
"pages": ["v3/tasks-regular", "v3/tasks-scheduled"]
|
||||
"pages": [
|
||||
"v3/tasks-regular",
|
||||
"v3/tasks-scheduled"
|
||||
]
|
||||
},
|
||||
"v3/trigger-config"
|
||||
]
|
||||
@@ -127,7 +138,10 @@
|
||||
{
|
||||
"group": "Development",
|
||||
"version": "v3 (Developer Preview)",
|
||||
"pages": ["v3/cli-dev", "v3/run-tests"]
|
||||
"pages": [
|
||||
"v3/cli-dev",
|
||||
"v3/run-tests"
|
||||
]
|
||||
},
|
||||
{
|
||||
"group": "Deployment",
|
||||
@@ -138,7 +152,9 @@
|
||||
"v3/github-actions",
|
||||
{
|
||||
"group": "Deployment integrations",
|
||||
"pages": ["v3/vercel-integration"]
|
||||
"pages": [
|
||||
"v3/vercel-integration"
|
||||
]
|
||||
}
|
||||
]
|
||||
},
|
||||
@@ -172,6 +188,13 @@
|
||||
"version": "v3 (Developer Preview)",
|
||||
"pages": [
|
||||
"v3/management/overview",
|
||||
{
|
||||
"group": "Tasks API",
|
||||
"pages": [
|
||||
"v3/management/tasks/trigger",
|
||||
"v3/management/tasks/batch-trigger"
|
||||
]
|
||||
},
|
||||
{
|
||||
"group": "Runs API",
|
||||
"pages": [
|
||||
@@ -207,14 +230,20 @@
|
||||
},
|
||||
{
|
||||
"group": "Projects API",
|
||||
"pages": ["v3/management/projects/runs"]
|
||||
"pages": [
|
||||
"v3/management/projects/runs"
|
||||
]
|
||||
}
|
||||
]
|
||||
},
|
||||
{
|
||||
"group": "Open source",
|
||||
"version": "v3 (Developer Preview)",
|
||||
"pages": ["v3/github-repo", "v3/open-source-self-hosting", "v3/open-source-contributing"]
|
||||
"pages": [
|
||||
"v3/github-repo",
|
||||
"v3/open-source-self-hosting",
|
||||
"v3/open-source-contributing"
|
||||
]
|
||||
},
|
||||
{
|
||||
"group": "Troubleshooting",
|
||||
@@ -230,7 +259,11 @@
|
||||
{
|
||||
"group": "Help",
|
||||
"version": "v3 (Developer Preview)",
|
||||
"pages": ["v3/community", "v3/help-slack", "v3/help-email"]
|
||||
"pages": [
|
||||
"v3/community",
|
||||
"v3/help-slack",
|
||||
"v3/help-email"
|
||||
]
|
||||
},
|
||||
{
|
||||
"group": "Getting Started",
|
||||
@@ -422,7 +455,10 @@
|
||||
"pages": [
|
||||
{
|
||||
"group": "Airtable",
|
||||
"pages": ["integrations/apis/airtable", "integrations/apis/airtable-tasks"]
|
||||
"pages": [
|
||||
"integrations/apis/airtable",
|
||||
"integrations/apis/airtable-tasks"
|
||||
]
|
||||
},
|
||||
{
|
||||
"group": "GitHub",
|
||||
@@ -448,16 +484,25 @@
|
||||
},
|
||||
{
|
||||
"group": "Plain",
|
||||
"pages": ["integrations/apis/plain", "integrations/apis/plain-tasks"]
|
||||
"pages": [
|
||||
"integrations/apis/plain",
|
||||
"integrations/apis/plain-tasks"
|
||||
]
|
||||
},
|
||||
"integrations/apis/replicate",
|
||||
{
|
||||
"group": "SendGrid",
|
||||
"pages": ["integrations/apis/sendgrid", "integrations/apis/sendgrid-tasks"]
|
||||
"pages": [
|
||||
"integrations/apis/sendgrid",
|
||||
"integrations/apis/sendgrid-tasks"
|
||||
]
|
||||
},
|
||||
{
|
||||
"group": "Resend",
|
||||
"pages": ["integrations/apis/resend", "integrations/apis/resend-tasks"]
|
||||
"pages": [
|
||||
"integrations/apis/resend",
|
||||
"integrations/apis/resend-tasks"
|
||||
]
|
||||
},
|
||||
{
|
||||
"group": "Shopify",
|
||||
@@ -469,7 +514,10 @@
|
||||
},
|
||||
{
|
||||
"group": "Slack",
|
||||
"pages": ["integrations/apis/slack", "integrations/apis/slack-tasks"]
|
||||
"pages": [
|
||||
"integrations/apis/slack",
|
||||
"integrations/apis/slack-tasks"
|
||||
]
|
||||
},
|
||||
"integrations/apis/stripe",
|
||||
{
|
||||
@@ -495,7 +543,9 @@
|
||||
"sdk/triggerclient/constructor",
|
||||
{
|
||||
"group": "Instance properties",
|
||||
"pages": ["sdk/triggerclient/store"]
|
||||
"pages": [
|
||||
"sdk/triggerclient/store"
|
||||
]
|
||||
},
|
||||
{
|
||||
"group": "Instance methods",
|
||||
@@ -558,7 +608,10 @@
|
||||
"sdk/dynamictrigger/constructor",
|
||||
{
|
||||
"group": "Instance methods",
|
||||
"pages": ["sdk/dynamictrigger/register", "sdk/dynamictrigger/unregister"]
|
||||
"pages": [
|
||||
"sdk/dynamictrigger/register",
|
||||
"sdk/dynamictrigger/unregister"
|
||||
]
|
||||
}
|
||||
]
|
||||
},
|
||||
@@ -569,7 +622,10 @@
|
||||
"sdk/dynamicschedule/constructor",
|
||||
{
|
||||
"group": "Instance methods",
|
||||
"pages": ["sdk/dynamicschedule/register", "sdk/dynamicschedule/unregister"]
|
||||
"pages": [
|
||||
"sdk/dynamicschedule/register",
|
||||
"sdk/dynamicschedule/unregister"
|
||||
]
|
||||
}
|
||||
]
|
||||
},
|
||||
@@ -582,7 +638,9 @@
|
||||
{
|
||||
"group": "HTTP Reference",
|
||||
"version": "v2",
|
||||
"pages": ["sdk/api-reference/events/create-an-event"]
|
||||
"pages": [
|
||||
"sdk/api-reference/events/create-an-event"
|
||||
]
|
||||
},
|
||||
{
|
||||
"group": "React SDK",
|
||||
@@ -598,7 +656,9 @@
|
||||
{
|
||||
"group": "Overview",
|
||||
"version": "v2",
|
||||
"pages": ["examples/introduction"]
|
||||
"pages": [
|
||||
"examples/introduction"
|
||||
]
|
||||
}
|
||||
],
|
||||
"footerSocials": {
|
||||
@@ -606,4 +666,4 @@
|
||||
"github": "https://github.com/triggerdotdev",
|
||||
"linkedin": "https://www.linkedin.com/company/triggerdotdev"
|
||||
}
|
||||
}
|
||||
}
|
||||
+290
-58
@@ -811,17 +811,7 @@ paths:
|
||||
description: Whether to override existing variables or not
|
||||
default: false
|
||||
required: ["variables"]
|
||||
multipart/form-data:
|
||||
schema:
|
||||
type: object
|
||||
properties:
|
||||
variables:
|
||||
type: string
|
||||
format: binary
|
||||
override:
|
||||
type: boolean
|
||||
required:
|
||||
- variables
|
||||
|
||||
responses:
|
||||
"200":
|
||||
description: Environment variables imported successfully
|
||||
@@ -864,57 +854,11 @@ paths:
|
||||
source: |-
|
||||
import { envvars } from "@trigger.dev/sdk/v3";
|
||||
|
||||
// Import variables from an array
|
||||
await envvars.upload("proj_yubjwjsfkxnylobaqvqz", "dev", {
|
||||
variables: [
|
||||
{
|
||||
name: "SLACK_API_KEY",
|
||||
value: "slack_123456"
|
||||
}
|
||||
],
|
||||
variables: { SLACK_API_KEY: "slack_key_1234" },
|
||||
override: false
|
||||
});
|
||||
- lang: typescript
|
||||
label: Import variables from a read stream
|
||||
source: |-
|
||||
import { envvars } from "@trigger.dev/sdk/v3";
|
||||
import { createReadStream } from "node:fs";
|
||||
|
||||
// Import variables in dotenv format from a file
|
||||
await envvars.upload("proj_yubjwjsfkxnylobaqvqz", "dev", {
|
||||
variables: createReadStream(".env"),
|
||||
override: false
|
||||
});
|
||||
- lang: typescript
|
||||
label: Import variables from a response
|
||||
source: |-
|
||||
import { envvars } from "@trigger.dev/sdk/v3";
|
||||
|
||||
// Import variables in dotenv format from a response
|
||||
await envvars.upload("proj_yubjwjsfkxnylobaqvqz", "dev", {
|
||||
variables: await fetch("https://example.com/.env"),
|
||||
override: false
|
||||
});
|
||||
- lang: typescript
|
||||
label: Import variables from a Buffer
|
||||
source: |-
|
||||
import { envvars } from "@trigger.dev/sdk/v3";
|
||||
|
||||
// Import variables in dotenv format from a buffer
|
||||
await envvars.upload("proj_yubjwjsfkxnylobaqvqz", "dev", {
|
||||
variables: Buffer.from("SLACK_API_KEY=slack_1234"),
|
||||
override: false
|
||||
});
|
||||
- lang: typescript
|
||||
label: Import variables from a File
|
||||
source: |-
|
||||
import { envvars } from "@trigger.dev/sdk/v3";
|
||||
|
||||
// Import variables in dotenv format from a file
|
||||
await envvars.upload("proj_yubjwjsfkxnylobaqvqz", "dev", {
|
||||
variables: new File(["SLACK_API_KEY=slack_1234"], ".env"),
|
||||
override: false
|
||||
});
|
||||
|
||||
"/api/v1/projects/{projectRef}/envvars/{env}/{name}":
|
||||
parameters:
|
||||
@@ -1095,9 +1039,229 @@ paths:
|
||||
});
|
||||
}
|
||||
})
|
||||
"/api/v1/tasks/{taskIdentifier}/trigger":
|
||||
parameters:
|
||||
- $ref: "#/components/parameters/taskIdentifier"
|
||||
post:
|
||||
operationId: trigger_task_v1
|
||||
summary: Trigger a task
|
||||
description: Trigger a task by its identifier.
|
||||
requestBody:
|
||||
required: true
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
"$ref": "#/components/schemas/TriggerTaskRequestBody"
|
||||
responses:
|
||||
"200":
|
||||
description: Task triggered successfully
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
"$ref": "#/components/schemas/TriggerTaskResponse"
|
||||
"400":
|
||||
description: Invalid request parameters or body
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
"$ref": "#/components/schemas/ErrorResponse"
|
||||
"401":
|
||||
description: Unauthorized request
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
"$ref": "#/components/schemas/ErrorResponse"
|
||||
"404":
|
||||
description: Resource not found
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
"$ref": "#/components/schemas/ErrorResponse"
|
||||
tags:
|
||||
- tasks
|
||||
security:
|
||||
- secretKey: []
|
||||
x-codeSamples:
|
||||
- lang: typescript
|
||||
source: |-
|
||||
import { task } from "@trigger.dev/sdk/v3";
|
||||
|
||||
export const myTask = await task({
|
||||
id: "my-task",
|
||||
run: async (payload: { message: string }) => {
|
||||
console.log("Hello, world!");
|
||||
}
|
||||
});
|
||||
|
||||
// Somewhere else in your code
|
||||
await myTask.trigger({ message: "Hello, world!" }, {
|
||||
idempotencyKey: "unique-key-123",
|
||||
concurrencyKey: "user123-task",
|
||||
queue: {
|
||||
name: "my-task-queue",
|
||||
concurrencyLimit: 5
|
||||
},
|
||||
});
|
||||
- lang: curl
|
||||
source: |-
|
||||
curl -X POST "https://api.trigger.dev/api/v1/tasks/my-task/trigger" \
|
||||
-H "Content-Type: application/json" \
|
||||
-H "Authorization: Bearer tr_dev_1234" \
|
||||
-d '{
|
||||
"payload": {
|
||||
"message": "Hello, world!"
|
||||
},
|
||||
"context": {
|
||||
"user": "user123"
|
||||
},
|
||||
"options": {
|
||||
"queue": {
|
||||
"name": "default",
|
||||
"concurrencyLimit": 5
|
||||
},
|
||||
"concurrencyKey": "user123-task",
|
||||
"idempotencyKey": "unique-key-123"
|
||||
}
|
||||
}'
|
||||
- lang: python
|
||||
source: |-
|
||||
import requests
|
||||
|
||||
url = "https://api.trigger.dev/api/v1/tasks/my-task/trigger"
|
||||
headers = {
|
||||
"Content-Type": "application/json",
|
||||
"Authorization": "Bearer tr_dev_1234"
|
||||
}
|
||||
data = {
|
||||
"payload": {
|
||||
"message": "Hello, world!"
|
||||
},
|
||||
"context": {
|
||||
"user": "user123"
|
||||
},
|
||||
"options": {
|
||||
"queue": {
|
||||
"name": "default",
|
||||
"concurrencyLimit": 5
|
||||
},
|
||||
"concurrencyKey": "user123-task",
|
||||
"idempotencyKey": "unique-key-123"
|
||||
}
|
||||
}
|
||||
|
||||
response = requests.post(url, headers=headers, json=data)
|
||||
print(response.json())
|
||||
|
||||
"/api/v1/tasks/{taskIdentifier}/batch":
|
||||
parameters:
|
||||
- $ref: "#/components/parameters/taskIdentifier"
|
||||
post:
|
||||
operationId: batch_trigger_task_v1
|
||||
summary: Batch trigger a task
|
||||
description: Batch trigger a task with up to 100 payloads.
|
||||
requestBody:
|
||||
required: true
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
"$ref": "#/components/schemas/BatchTriggerRequestBody"
|
||||
responses:
|
||||
"200":
|
||||
description: Task batch triggered successfully
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
"$ref": "#/components/schemas/BatchTriggerTaskResponse"
|
||||
"400":
|
||||
description: Invalid request parameters or body
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
"$ref": "#/components/schemas/ErrorResponse"
|
||||
"401":
|
||||
description: Unauthorized request
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
"$ref": "#/components/schemas/ErrorResponse"
|
||||
"404":
|
||||
description: Resource not found
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
"$ref": "#/components/schemas/ErrorResponse"
|
||||
tags:
|
||||
- tasks
|
||||
security:
|
||||
- secretKey: []
|
||||
x-codeSamples:
|
||||
- lang: typescript
|
||||
source: |-
|
||||
import { task } from "@trigger.dev/sdk/v3";
|
||||
|
||||
export const myTask = await task({
|
||||
id: "my-task",
|
||||
run: async (payload: { message: string }) => {
|
||||
console.log("Hello, world!");
|
||||
}
|
||||
});
|
||||
|
||||
// Somewhere else in your code
|
||||
await myTask.batchTrigger({
|
||||
items: [
|
||||
{
|
||||
payload: { message: "Hello, world!" },
|
||||
options: {
|
||||
idempotencyKey: "unique-key-123",
|
||||
concurrencyKey: "user-123-task",
|
||||
queue: {
|
||||
name: "my-task-queue",
|
||||
concurrencyLimit: 5
|
||||
}
|
||||
}
|
||||
}
|
||||
]
|
||||
});
|
||||
- lang: curl
|
||||
source: |-
|
||||
curl -X POST "https://api.trigger.dev/api/v1/tasks/my-task/batch" \
|
||||
-H "Content-Type: application/json" \
|
||||
-H "Authorization: Bearer tr_dev_1234" \
|
||||
-d '{
|
||||
"items": [
|
||||
{
|
||||
"payload": {
|
||||
"message": "Hello, world!"
|
||||
},
|
||||
"context": {
|
||||
"user": "user123"
|
||||
},
|
||||
"options": {
|
||||
"queue": {
|
||||
"name": "default",
|
||||
"concurrencyLimit": 5
|
||||
},
|
||||
"concurrencyKey": "user123-task",
|
||||
"idempotencyKey": "unique-key-123"
|
||||
}
|
||||
}
|
||||
]
|
||||
}'
|
||||
|
||||
|
||||
|
||||
|
||||
|
||||
components:
|
||||
parameters:
|
||||
taskIdentifier:
|
||||
in: path
|
||||
name: taskIdentifier
|
||||
required: true
|
||||
schema:
|
||||
type: string
|
||||
description: The id of a task
|
||||
example: my-task
|
||||
runsFilterWithEnv:
|
||||
in: query
|
||||
name: filter
|
||||
@@ -1215,6 +1379,65 @@ components:
|
||||
configure({ secretKey: "tr_pat_1234" });
|
||||
```
|
||||
schemas:
|
||||
TriggerTaskResponse:
|
||||
type: object
|
||||
properties:
|
||||
id:
|
||||
type: string
|
||||
description: The ID of the run that was triggered.
|
||||
example: run_1234
|
||||
QueueOptions:
|
||||
type: object
|
||||
properties:
|
||||
name:
|
||||
type: string
|
||||
description: You can define a shared queue and then pass the name in to your task.
|
||||
concurrencyLimit:
|
||||
type: integer
|
||||
minimum: 0
|
||||
maximum: 1000
|
||||
description: An optional property that specifies the maximum number of concurrent run executions. If this property is omitted, the task can potentially use up the full concurrency of an environment.
|
||||
BatchTriggerRequestBody:
|
||||
type: object
|
||||
properties:
|
||||
items:
|
||||
type: array
|
||||
items:
|
||||
"$ref": "#/components/schemas/TriggerTaskRequestBody"
|
||||
description: An array of payloads to trigger the task with
|
||||
required: ["items"]
|
||||
BatchTriggerTaskResponse:
|
||||
type: object
|
||||
required: ["batchId", "runs"]
|
||||
properties:
|
||||
batchId:
|
||||
type: string
|
||||
description: The ID of the batch that was triggered
|
||||
example: batch_1234
|
||||
runs:
|
||||
type: array
|
||||
items:
|
||||
type: string
|
||||
description: An array of run IDs that were triggered
|
||||
TriggerTaskRequestBody:
|
||||
type: object
|
||||
properties:
|
||||
payload:
|
||||
description: The payload can include any valid JSON
|
||||
context:
|
||||
description: The context can include any valid JSON
|
||||
options:
|
||||
type: object
|
||||
properties:
|
||||
queue:
|
||||
$ref: '#/components/schemas/QueueOptions'
|
||||
concurrencyKey:
|
||||
type: string
|
||||
description: Scope the concurrency limit to a specific key.
|
||||
idempotencyKey:
|
||||
type: string
|
||||
description: An optional property that specifies the idempotency key used to prevent creating duplicate runs. If you provide an existing idempotency key, we will return the existing run ID.
|
||||
|
||||
EnvFilter:
|
||||
type: object
|
||||
properties:
|
||||
@@ -1406,6 +1629,7 @@ components:
|
||||
properties:
|
||||
error:
|
||||
type: string
|
||||
example: Something went wrong
|
||||
required: ["error"]
|
||||
ErrorWithDetailsResponse:
|
||||
type: object
|
||||
@@ -1498,10 +1722,18 @@ components:
|
||||
type: object
|
||||
description: The payload that was sent to the task. Will be omitted if the request was made with a Public API key
|
||||
example: { "foo": "bar" }
|
||||
payloadPresignedUrl:
|
||||
type: string
|
||||
description: The presigned URL to download the payload. Will only be included if the payload is too large to be included in the response. Expires in 5 minutes.
|
||||
example: "https://r2.cloudflarestorage.com/packets/yubjwjsfkxnylobaqvqz/dev/run_p4omhh45hgxxnq1re6ovy/payload.json?X-Amz-Expires=300&X-Amz-Date=20240625T154526Z&X-Amz-Algorithm=AWS4-HMAC-SHA256&X-Amz-Credential=10b064e58a0680db5b5e077be2be3b2a%2F20240625%2Fauto%2Fs3%2Faws4_request&X-Amz-SignedHeaders=host&X-Amz-Signature=88604cb993ffc151b4d73f2439da431d9928488e4b3dcfa4a7c8f1819"
|
||||
output:
|
||||
type: object
|
||||
description: The output of the run. Will be omitted if the request was made with a Public API key
|
||||
example: { "foo": "bar" }
|
||||
outputPresignedUrl:
|
||||
type: string
|
||||
description: The presigned URL to download the output. Will only be included if the output is too large to be included in the response. Expires in 5 minutes.
|
||||
example: "https://r2.cloudflarestorage.com/packets/yubjwjsfkxnylobaqvqz/dev/run_p4omhh45hgxxnq1re6ovy/payload.json?X-Amz-Expires=300&X-Amz-Date=20240625T154526Z&X-Amz-Algorithm=AWS4-HMAC-SHA256&X-Amz-Credential=10b064e58a0680db5b5e077be2be3b2a%2F20240625%2Fauto%2Fs3%2Faws4_request&X-Amz-SignedHeaders=host&X-Amz-Signature=88604cb993ffc151b4d73f2439da431d9928488e4b3dcfa4a7c8f1819"
|
||||
idempotencyKey:
|
||||
type: string
|
||||
description: The idempotency key used to prevent creating duplicate runs, if provided
|
||||
|
||||
@@ -37,3 +37,13 @@ If you add them dynamically using code make sure you add a `deduplicationKey` so
|
||||
If you're creating schedules for your user you will definitely need to request more schedules from us.
|
||||
|
||||
<Snippet file="v3/soft-limit.mdx" />
|
||||
|
||||
## Task payloads and outputs
|
||||
|
||||
| Limit | Details |
|
||||
| ---------------------- | ---------------------------------------------- |
|
||||
| Single trigger payload | Must not exceed 10MB |
|
||||
| Batch trigger payload | The total of all payloads must not exceed 10MB |
|
||||
| Task outputs | Must not exceed 10MB |
|
||||
|
||||
Payloads and outputs that exceed 512KB will be offloaded to object storage and a presigned URL will be provided to download the data when calling `runs.retrieve`. You don't need to do anything to handle this in your tasks however, as we will transparently upload/download these during operation.
|
||||
|
||||
@@ -94,6 +94,8 @@ function personalAccessTokenExample() {
|
||||
|
||||
| Endpoint | Secret key | Personal Access Token |
|
||||
| ---------------------- | ---------- | --------------------- |
|
||||
| `task.trigger` | ✅ | |
|
||||
| `task.batchTrigger` | ✅ | |
|
||||
| `runs.list` | ✅ | ✅ |
|
||||
| `runs.retrieve` | ✅ | |
|
||||
| `runs.cancel` | ✅ | |
|
||||
|
||||
@@ -0,0 +1,4 @@
|
||||
---
|
||||
title: "Batch trigger"
|
||||
openapi: "v3-openapi POST /api/v1/tasks/{taskIdentifier}/batch"
|
||||
---
|
||||
@@ -0,0 +1,4 @@
|
||||
---
|
||||
title: "Trigger"
|
||||
openapi: "v3-openapi POST /api/v1/tasks/{taskIdentifier}/trigger"
|
||||
---
|
||||
+317
-47
@@ -3,7 +3,7 @@ title: "Triggering"
|
||||
description: "Tasks need to be triggered to run."
|
||||
---
|
||||
|
||||
There are currently four ways you can trigger any task from your own code:
|
||||
There are currently six ways you can trigger tasks:
|
||||
|
||||
| Function | Where does this work? | What it does |
|
||||
| -------------------------------- | --------------------- | ---------------------------------------------------------------------------------------------------------------------------------- |
|
||||
@@ -11,6 +11,8 @@ There are currently four ways you can trigger any task from your own code:
|
||||
| `yourTask.batchTrigger()` | Anywhere | Triggers a task multiple times and gets a handle you can use to monitor and manage the runs. It does not wait for the results. |
|
||||
| `yourTask.triggerAndWait()` | Inside a task | Triggers a task and then waits until it's complete. You get the result data to continue with. |
|
||||
| `yourTask.batchTriggerAndWait()` | Inside a task | Triggers a task multiple times in parallel and then waits until they're all complete. You get the resulting data to continue with. |
|
||||
| `tasks.trigger()` | Outside of a task | Triggers a task and gets a handle you can use to fetch and manage the run. |
|
||||
| `tasks.batchTrigger()` | Outside of a task | Triggers a task multiple times and gets a handle you can use to fetch and manage the runs. |
|
||||
|
||||
Additionally, [scheduled tasks](/v3/tasks-scheduled) get automatically triggered on their schedule and [webhooks](/v3/tasks-webhooks) when receiving a webhook.
|
||||
|
||||
@@ -18,25 +20,20 @@ Additionally, [scheduled tasks](/v3/tasks-scheduled) get automatically triggered
|
||||
|
||||
You should attach one or more schedules to your `schedules.task()` to trigger it on a recurring schedule. [Read the scheduled tasks docs](/v3/tasks-scheduled).
|
||||
|
||||
## From outside of a task
|
||||
|
||||
You can trigger any task from your backend code, using either `trigger()` or `batchTrigger()`.
|
||||
|
||||
<Note>
|
||||
Do not trigger tasks directly from your frontend. If you do, you will leak your private
|
||||
Trigger.dev API key to the world.
|
||||
</Note>
|
||||
|
||||
You can use Next.js Server Actions but [you need to be careful with bundling](#next-js-server-actions).
|
||||
|
||||
### Authentication
|
||||
## Authentication
|
||||
|
||||
When you trigger a task from your backend code, you need to set the `TRIGGER_SECRET_KEY` environment variable. You can find the value on the API keys page in the Trigger.dev dashboard. [More info on API keys](/v3/apikeys).
|
||||
|
||||
### trigger()
|
||||
## Task instance methods
|
||||
|
||||
Task instance methods are available on the `Task` object you receive when you define a task. They can be called from your backend code or from inside another task.
|
||||
|
||||
### Task.trigger()
|
||||
|
||||
Triggers a single run of a task with the payload you pass in, and any options you specify. It does NOT wait for the result, you cannot do that from outside a task.
|
||||
|
||||
If called from within a task, you can use the `AndWait` version to pause execution until the triggered run is complete.
|
||||
|
||||
<CodeGroup>
|
||||
|
||||
```ts Next.js API route
|
||||
@@ -74,11 +71,24 @@ export async function action({ request, params }: ActionFunctionArgs) {
|
||||
}
|
||||
```
|
||||
|
||||
```ts /trigger/my-task.ts
|
||||
import { myOtherTask } from "~/trigger/my-other-task";
|
||||
|
||||
export const myTask = task({
|
||||
id: "my-task",
|
||||
run: async (payload: string) => {
|
||||
const handle = await myOtherTask.trigger("some data");
|
||||
|
||||
//...do other stuff
|
||||
},
|
||||
});
|
||||
```
|
||||
|
||||
</CodeGroup>
|
||||
|
||||
### batchTrigger()
|
||||
### Task.batchTrigger()
|
||||
|
||||
Triggers multiples runs of a task with the payloads you pass in, and any options you specify. It does NOT wait for the results, you cannot do that from outside a task.
|
||||
Triggers multiples runs of a task with the payloads you pass in, and any options you specify.
|
||||
|
||||
<CodeGroup>
|
||||
|
||||
@@ -121,33 +131,6 @@ export async function action({ request, params }: ActionFunctionArgs) {
|
||||
}
|
||||
```
|
||||
|
||||
</CodeGroup>
|
||||
|
||||
## From inside a task
|
||||
|
||||
You can trigger tasks from other tasks using `trigger()` or `batchTrigger()`. You can also trigger and wait for the result of triggered tasks using `triggerAndWait()` and `batchTriggerAndWait()`. This is a powerful way to build complex tasks.
|
||||
|
||||
### trigger()
|
||||
|
||||
This works the same as from outside a task. You call it and you get a handle back, but it does not wait for the result.
|
||||
|
||||
```ts /trigger/my-task.ts
|
||||
import { myOtherTask } from "~/trigger/my-other-task";
|
||||
|
||||
export const myTask = task({
|
||||
id: "my-task",
|
||||
run: async (payload: string) => {
|
||||
const handle = await myOtherTask.trigger("some data");
|
||||
|
||||
//...do other stuff
|
||||
},
|
||||
});
|
||||
```
|
||||
|
||||
### batchTrigger()
|
||||
|
||||
This works the same as from outside a task. You call it and you get a handle back, but it does not wait for the results.
|
||||
|
||||
```ts /trigger/my-task.ts
|
||||
import { myOtherTask } from "~/trigger/my-other-task";
|
||||
|
||||
@@ -161,7 +144,9 @@ export const myTask = task({
|
||||
});
|
||||
```
|
||||
|
||||
### triggerAndWait()
|
||||
</CodeGroup>
|
||||
|
||||
### Task.triggerAndWait()
|
||||
|
||||
This is where it gets interesting. You can trigger a task and then wait for the result. This is useful when you need to call a different task and then use the result to continue with your task.
|
||||
|
||||
@@ -219,7 +204,7 @@ export const parentTask = task({
|
||||
});
|
||||
```
|
||||
|
||||
### batchTriggerAndWait()
|
||||
### Task.batchTriggerAndWait()
|
||||
|
||||
You can batch trigger a task and wait for all the results. This is useful for the fan-out pattern, where you need to call a task multiple times and then wait for all the results to continue with your task.
|
||||
|
||||
@@ -284,11 +269,213 @@ export const batchParentTask = task({
|
||||
});
|
||||
```
|
||||
|
||||
## SDK functions
|
||||
|
||||
You can trigger any task from your backend code using the `tasks.trigger()` or `tasks.batchTrigger()` SDK functions.
|
||||
|
||||
<Note>
|
||||
Do not trigger tasks directly from your frontend. If you do, you will leak your private
|
||||
Trigger.dev API key.
|
||||
</Note>
|
||||
|
||||
You can use Next.js Server Actions but [you need to be careful with bundling](#next-js-server-actions).
|
||||
|
||||
### tasks.trigger()
|
||||
|
||||
Triggers a single run of a task with the payload you pass in, and any options you specify, without needing to import the task.
|
||||
|
||||
<Note>
|
||||
Why would you use this instead of the `Task.trigger()` instance method? Tasks can import
|
||||
dependencies/modules that you might not want included in your application code or cause problems
|
||||
with building.
|
||||
</Note>
|
||||
|
||||
<CodeGroup>
|
||||
|
||||
```ts Next.js API route
|
||||
import { tasks } from "@trigger.dev/sdk/v3";
|
||||
import type { emailSequence } from "~/trigger/emails";
|
||||
// 👆 **type-only** import
|
||||
|
||||
//app/email/route.ts
|
||||
export async function POST(request: Request) {
|
||||
//get the JSON from the request
|
||||
const data = await request.json();
|
||||
|
||||
// Pass the task type to `trigger()` as a generic argument, giving you full type checking
|
||||
const handle = await tasks.trigger<typeof emailSequence>("email-sequence", {
|
||||
to: data.email,
|
||||
name: data.name,
|
||||
});
|
||||
|
||||
//return a success response with the handle
|
||||
return Response.json(handle);
|
||||
}
|
||||
```
|
||||
|
||||
```ts Remix
|
||||
import { tasks } from "@trigger.dev/sdk/v3";
|
||||
|
||||
export async function action({ request, params }: ActionFunctionArgs) {
|
||||
if (request.method.toUpperCase() !== "POST") {
|
||||
return json("Method Not Allowed", { status: 405 });
|
||||
}
|
||||
|
||||
//get the JSON from the request
|
||||
const data = await request.json();
|
||||
|
||||
// The generic argument is optional, but recommended for full type checking
|
||||
const handle = await tasks.trigger("email-sequence", {
|
||||
to: data.email,
|
||||
name: data.name,
|
||||
});
|
||||
|
||||
//return a success response with the handle
|
||||
return json(handle);
|
||||
}
|
||||
```
|
||||
|
||||
</CodeGroup>
|
||||
|
||||
<Tip>
|
||||
By importing the task with the type modifier, the import of `"~/trigger/emails"` is a type-only
|
||||
import. This means that the task code is not included in your application at build time.
|
||||
</Tip>
|
||||
|
||||
### tasks.batchTrigger()
|
||||
|
||||
Triggers multiples runs of a task with the payloads you pass in, and any options you specify, without needing to import the task.
|
||||
|
||||
<CodeGroup>
|
||||
|
||||
```ts Next.js API route
|
||||
import { tasks } from "@trigger.dev/sdk/v3";
|
||||
import type { emailSequence } from "~/trigger/emails";
|
||||
// 👆 **type-only** import
|
||||
|
||||
//app/email/route.ts
|
||||
export async function POST(request: Request) {
|
||||
//get the JSON from the request
|
||||
const data = await request.json();
|
||||
|
||||
// Pass the task type to `batchTrigger()` as a generic argument, giving you full type checking
|
||||
const batchHandle = await tasks.batchTrigger<typeof emailSequence>(
|
||||
"email-sequence",
|
||||
data.users.map((u) => ({ payload: { to: u.email, name: u.name } }))
|
||||
);
|
||||
|
||||
//return a success response with the handle
|
||||
return Response.json(batchHandle);
|
||||
}
|
||||
```
|
||||
|
||||
</CodeGroup>
|
||||
|
||||
### tasks.triggerAndPoll()
|
||||
|
||||
Triggers a single run of a task with the payload you pass in, and any options you specify, and then polls the run until it's complete.
|
||||
|
||||
<CodeGroup>
|
||||
|
||||
```ts Next.js API route
|
||||
import { tasks } from "@trigger.dev/sdk/v3";
|
||||
import type { emailSequence } from "~/trigger/emails";
|
||||
|
||||
//app/email/route.ts
|
||||
export async function POST(request: Request) {
|
||||
//get the JSON from the request
|
||||
const data = await request.json();
|
||||
|
||||
// Pass the task type to `triggerAndPoll()` as a generic argument, giving you full type checking
|
||||
const result = await tasks.triggerAndPoll<typeof emailSequence>(
|
||||
"email-sequence",
|
||||
{
|
||||
to: data.email,
|
||||
name: data.name,
|
||||
},
|
||||
{ pollIntervalMs: 5000 }
|
||||
);
|
||||
|
||||
//return a success response with the result
|
||||
return Response.json(result);
|
||||
}
|
||||
```
|
||||
|
||||
</CodeGroup>
|
||||
|
||||
<Note>
|
||||
The above code is just a demonstration of the API and is not recommended to use in an API route
|
||||
this way as it will block the request until the task is complete.
|
||||
</Note>
|
||||
|
||||
### runs.retrieve()
|
||||
|
||||
You can retrieve a run by its handle using the `runs.retrieve()` function.
|
||||
|
||||
<CodeGroup>
|
||||
|
||||
```ts Next.js API route
|
||||
import { tasks, runs } from "@trigger.dev/sdk/v3";
|
||||
import type { emailSequence } from "~/trigger/emails";
|
||||
// 👆 **type-only** import
|
||||
|
||||
//app/email/route.ts
|
||||
export async function POST(request: Request) {
|
||||
//get the JSON from the request
|
||||
const data = await request.json();
|
||||
|
||||
// Pass the task type to `trigger()` as a generic argument, giving you full type checking
|
||||
const handle = await tasks.trigger<typeof emailSequence>("email-sequence", {
|
||||
to: data.email,
|
||||
name: data.name,
|
||||
});
|
||||
|
||||
const run = await runs.retrieve(handle);
|
||||
|
||||
// run.output will be correctly typed as the return value of the task
|
||||
return Response.json(run.output);
|
||||
}
|
||||
```
|
||||
|
||||
</CodeGroup>
|
||||
|
||||
### runs.poll()
|
||||
|
||||
You can poll a run by its handle using the `runs.poll()` function.
|
||||
|
||||
<CodeGroup>
|
||||
|
||||
```ts Next.js API route
|
||||
import { tasks, runs } from "@trigger.dev/sdk/v3";
|
||||
import type { emailSequence } from "~/trigger/emails";
|
||||
// 👆 **type-only** import
|
||||
|
||||
//app/email/route.ts
|
||||
export async function POST(request: Request) {
|
||||
//get the JSON from the request
|
||||
const data = await request.json();
|
||||
|
||||
// Pass the task type to `trigger()` as a generic argument, giving you full type checking
|
||||
const handle = await tasks.trigger<typeof emailSequence>("email-sequence", {
|
||||
to: data.email,
|
||||
name: data.name,
|
||||
});
|
||||
|
||||
// Poll the run until it's complete
|
||||
const run = await runs.poll(handle, { pollIntervalMs: 5000 });
|
||||
|
||||
// run.output will be correctly typed as the return value of the task
|
||||
return Response.json(run.output);
|
||||
}
|
||||
```
|
||||
|
||||
</CodeGroup>
|
||||
|
||||
## Next.js Server Actions
|
||||
|
||||
Server Actions allow you to call your backend code without creating API routes. This is very useful for triggering tasks but you need to be careful you don't accidentally bundle the Trigger.dev SDK into your frontend code.
|
||||
|
||||
If you see an error like this then you've bundled `@trigger.dev/sdk` into your frontend code:
|
||||
If you see an error like this then you've bundled `@trigger.dev/sdk/v3` into your frontend code:
|
||||
|
||||
```bash
|
||||
Module build failed: UnhandledSchemeError: Reading from "node:crypto" is not handled by plugins (Unhandled scheme).
|
||||
@@ -297,7 +484,7 @@ Webpack supports "data:" and "file:" URIs by default.
|
||||
You may need an additional plugin to handle "node:" URIs.
|
||||
```
|
||||
|
||||
When you use server actions that use `@trigger.dev/sdk`:
|
||||
When you use server actions that use `@trigger.dev/sdk/v3`:
|
||||
|
||||
- The file can't have any React components in it.
|
||||
- The file should have `"use server"` on the first line.
|
||||
@@ -349,3 +536,86 @@ export async function create() {
|
||||
```
|
||||
|
||||
</CodeGroup>
|
||||
|
||||
## Large Payloads
|
||||
|
||||
We recommend keeping your task payloads as small as possible. We currently have a hard limit on task payloads above 10MB.
|
||||
|
||||
If your payload size is larger than 512KB, instead of saving the payload to the database, we will upload it to an S3-compatible object store and store the URL in the database.
|
||||
|
||||
When your task runs, we automatically download the payload from the object store and pass it to your task function. We also will return to you a `payloadPresignedUrl` from the `runs.retrieve` SDK function so you can download the payload if needed:
|
||||
|
||||
```ts
|
||||
import { runs } from "@trigger.dev/sdk/v3";
|
||||
|
||||
const run = await runs.retrieve(handle);
|
||||
|
||||
if (run.payloadPresignedUrl) {
|
||||
const response = await fetch(run.payloadPresignedUrl);
|
||||
const payload = await response.json();
|
||||
|
||||
console.log("Payload", payload);
|
||||
}
|
||||
```
|
||||
|
||||
<Note>
|
||||
We also use this same system for dealing with large task outputs, and subsequently will return a
|
||||
corresponding `outputPresignedUrl`. Task outputs are limited to 100MB.
|
||||
</Note>
|
||||
|
||||
If you need to pass larger payloads, you'll need to upload the payload to your own storage and pass a URL to the file in the payload instead. For example, uploading to S3 and then sending a presigned URL that expires in URL:
|
||||
|
||||
<CodeGroup>
|
||||
|
||||
```ts /yourServer.ts
|
||||
import { myTask } from "./trigger/myTasks";
|
||||
import { s3Client, getSignedUrl, PutObjectCommand, GetObjectCommand } from "./s3";
|
||||
import { createReadStream } from "node:fs";
|
||||
|
||||
// Upload file to S3
|
||||
await s3Client.send(
|
||||
new PutObjectCommand({
|
||||
Bucket: "my-bucket",
|
||||
Key: "myfile.json",
|
||||
Body: createReadStream("large-payload.json"),
|
||||
})
|
||||
);
|
||||
|
||||
// Create presigned URL
|
||||
const presignedUrl = await getSignedUrl(
|
||||
s3Client,
|
||||
new GetObjectCommand({
|
||||
Bucket: "my-bucket",
|
||||
Key: "my-file.json",
|
||||
}),
|
||||
{
|
||||
expiresIn: 3600, // expires in 1 hour
|
||||
}
|
||||
);
|
||||
|
||||
// Now send the URL to the task
|
||||
const handle = await myTask.trigger({
|
||||
url: presignedUrl,
|
||||
});
|
||||
```
|
||||
|
||||
```ts /trigger/myTasks.ts
|
||||
import { task } from "@trigger.dev/sdk/v3";
|
||||
|
||||
export const myTask = task({
|
||||
id: "my-task",
|
||||
run: async (payload: { url: string }) => {
|
||||
// Download the file from the URL
|
||||
const response = await fetch(payload.url);
|
||||
const data = await response.json();
|
||||
|
||||
// Do something with the data
|
||||
},
|
||||
});
|
||||
```
|
||||
|
||||
</CodeGroup>
|
||||
|
||||
### Batch Triggering
|
||||
|
||||
When using `batchTrigger` or `batchTriggerAndWait`, the total size of all payloads cannot exceed 10MB. This means if you are doing a batch of 100 runs, each payload should be less than 100KB.
|
||||
|
||||
@@ -1,5 +1,43 @@
|
||||
# @trigger.dev/airtable
|
||||
|
||||
## 3.0.0-beta.44
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.44
|
||||
- @trigger.dev/sdk@3.0.0-beta.44
|
||||
|
||||
## 3.0.0-beta.43
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- Updated dependencies [34ca7667d]
|
||||
- @trigger.dev/sdk@3.0.0-beta.43
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.43
|
||||
|
||||
## 3.0.0-beta.42
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- Updated dependencies [ecef19966]
|
||||
- @trigger.dev/sdk@3.0.0-beta.42
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.42
|
||||
|
||||
## 3.0.0-beta.41
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- Updated dependencies [7c36a1a4b]
|
||||
- @trigger.dev/sdk@3.0.0-beta.41
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.41
|
||||
|
||||
## 3.0.0-beta.40
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.40
|
||||
- @trigger.dev/sdk@3.0.0-beta.40
|
||||
|
||||
## 3.0.0-beta.39
|
||||
|
||||
### Patch Changes
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"name": "@trigger.dev/airtable",
|
||||
"version": "3.0.0-beta.39",
|
||||
"version": "3.0.0-beta.44",
|
||||
"description": "Trigger.dev integration for airtable",
|
||||
"main": "./dist/index.js",
|
||||
"types": "./dist/index.d.ts",
|
||||
@@ -25,8 +25,8 @@
|
||||
"typecheck": "tsc --noEmit"
|
||||
},
|
||||
"dependencies": {
|
||||
"@trigger.dev/integration-kit": "workspace:^3.0.0-beta.39",
|
||||
"@trigger.dev/sdk": "workspace:^3.0.0-beta.39",
|
||||
"@trigger.dev/integration-kit": "workspace:^3.0.0-beta.44",
|
||||
"@trigger.dev/sdk": "workspace:^3.0.0-beta.44",
|
||||
"airtable": "^0.12.1",
|
||||
"zod": "3.22.3"
|
||||
},
|
||||
|
||||
@@ -1,5 +1,43 @@
|
||||
# @trigger.dev/github
|
||||
|
||||
## 3.0.0-beta.44
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.44
|
||||
- @trigger.dev/sdk@3.0.0-beta.44
|
||||
|
||||
## 3.0.0-beta.43
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- Updated dependencies [34ca7667d]
|
||||
- @trigger.dev/sdk@3.0.0-beta.43
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.43
|
||||
|
||||
## 3.0.0-beta.42
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- Updated dependencies [ecef19966]
|
||||
- @trigger.dev/sdk@3.0.0-beta.42
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.42
|
||||
|
||||
## 3.0.0-beta.41
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- Updated dependencies [7c36a1a4b]
|
||||
- @trigger.dev/sdk@3.0.0-beta.41
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.41
|
||||
|
||||
## 3.0.0-beta.40
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.40
|
||||
- @trigger.dev/sdk@3.0.0-beta.40
|
||||
|
||||
## 3.0.0-beta.39
|
||||
|
||||
### Patch Changes
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"name": "@trigger.dev/github",
|
||||
"version": "3.0.0-beta.39",
|
||||
"version": "3.0.0-beta.44",
|
||||
"description": "The official GitHub integration for Trigger.dev",
|
||||
"main": "./dist/index.js",
|
||||
"types": "./dist/index.d.ts",
|
||||
@@ -30,8 +30,8 @@
|
||||
"@octokit/request-error": "^5.0.1",
|
||||
"@octokit/webhooks": "^12.0.10",
|
||||
"octokit": "^3.1.2",
|
||||
"@trigger.dev/integration-kit": "workspace:^3.0.0-beta.39",
|
||||
"@trigger.dev/sdk": "workspace:^3.0.0-beta.39",
|
||||
"@trigger.dev/integration-kit": "workspace:^3.0.0-beta.44",
|
||||
"@trigger.dev/sdk": "workspace:^3.0.0-beta.44",
|
||||
"zod": "3.22.3"
|
||||
},
|
||||
"engines": {
|
||||
|
||||
@@ -1,5 +1,43 @@
|
||||
# @trigger.dev/linear
|
||||
|
||||
## 3.0.0-beta.44
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.44
|
||||
- @trigger.dev/sdk@3.0.0-beta.44
|
||||
|
||||
## 3.0.0-beta.43
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- Updated dependencies [34ca7667d]
|
||||
- @trigger.dev/sdk@3.0.0-beta.43
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.43
|
||||
|
||||
## 3.0.0-beta.42
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- Updated dependencies [ecef19966]
|
||||
- @trigger.dev/sdk@3.0.0-beta.42
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.42
|
||||
|
||||
## 3.0.0-beta.41
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- Updated dependencies [7c36a1a4b]
|
||||
- @trigger.dev/sdk@3.0.0-beta.41
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.41
|
||||
|
||||
## 3.0.0-beta.40
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.40
|
||||
- @trigger.dev/sdk@3.0.0-beta.40
|
||||
|
||||
## 3.0.0-beta.39
|
||||
|
||||
### Patch Changes
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"name": "@trigger.dev/linear",
|
||||
"version": "3.0.0-beta.39",
|
||||
"version": "3.0.0-beta.44",
|
||||
"description": "Trigger.dev integration for @linear/sdk",
|
||||
"main": "./dist/index.js",
|
||||
"types": "./dist/index.d.ts",
|
||||
@@ -26,8 +26,8 @@
|
||||
},
|
||||
"dependencies": {
|
||||
"@linear/sdk": "^8.0.0",
|
||||
"@trigger.dev/integration-kit": "workspace:^3.0.0-beta.39",
|
||||
"@trigger.dev/sdk": "workspace:^3.0.0-beta.39",
|
||||
"@trigger.dev/integration-kit": "workspace:^3.0.0-beta.44",
|
||||
"@trigger.dev/sdk": "workspace:^3.0.0-beta.44",
|
||||
"zod": "3.22.3"
|
||||
},
|
||||
"engines": {
|
||||
|
||||
@@ -1,5 +1,43 @@
|
||||
# @trigger.dev/slack
|
||||
|
||||
## 3.0.0-beta.44
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.44
|
||||
- @trigger.dev/sdk@3.0.0-beta.44
|
||||
|
||||
## 3.0.0-beta.43
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- Updated dependencies [34ca7667d]
|
||||
- @trigger.dev/sdk@3.0.0-beta.43
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.43
|
||||
|
||||
## 3.0.0-beta.42
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- Updated dependencies [ecef19966]
|
||||
- @trigger.dev/sdk@3.0.0-beta.42
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.42
|
||||
|
||||
## 3.0.0-beta.41
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- Updated dependencies [7c36a1a4b]
|
||||
- @trigger.dev/sdk@3.0.0-beta.41
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.41
|
||||
|
||||
## 3.0.0-beta.40
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.40
|
||||
- @trigger.dev/sdk@3.0.0-beta.40
|
||||
|
||||
## 3.0.0-beta.39
|
||||
|
||||
### Patch Changes
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"name": "@trigger.dev/openai",
|
||||
"version": "3.0.0-beta.39",
|
||||
"version": "3.0.0-beta.44",
|
||||
"description": "The official OpenAI integration for Trigger.dev",
|
||||
"main": "./dist/index.js",
|
||||
"types": "./dist/index.d.ts",
|
||||
@@ -42,8 +42,8 @@
|
||||
},
|
||||
"dependencies": {
|
||||
"openai": "^4.16.1",
|
||||
"@trigger.dev/sdk": "workspace:^3.0.0-beta.39",
|
||||
"@trigger.dev/integration-kit": "workspace:^3.0.0-beta.39"
|
||||
"@trigger.dev/sdk": "workspace:^3.0.0-beta.44",
|
||||
"@trigger.dev/integration-kit": "workspace:^3.0.0-beta.44"
|
||||
},
|
||||
"engines": {
|
||||
"node": ">=18.0.0"
|
||||
|
||||
@@ -1,5 +1,43 @@
|
||||
# @trigger.dev/plain
|
||||
|
||||
## 3.0.0-beta.44
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.44
|
||||
- @trigger.dev/sdk@3.0.0-beta.44
|
||||
|
||||
## 3.0.0-beta.43
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- Updated dependencies [34ca7667d]
|
||||
- @trigger.dev/sdk@3.0.0-beta.43
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.43
|
||||
|
||||
## 3.0.0-beta.42
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- Updated dependencies [ecef19966]
|
||||
- @trigger.dev/sdk@3.0.0-beta.42
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.42
|
||||
|
||||
## 3.0.0-beta.41
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- Updated dependencies [7c36a1a4b]
|
||||
- @trigger.dev/sdk@3.0.0-beta.41
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.41
|
||||
|
||||
## 3.0.0-beta.40
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.40
|
||||
- @trigger.dev/sdk@3.0.0-beta.40
|
||||
|
||||
## 3.0.0-beta.39
|
||||
|
||||
### Patch Changes
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"name": "@trigger.dev/plain",
|
||||
"version": "3.0.0-beta.39",
|
||||
"version": "3.0.0-beta.44",
|
||||
"description": "The official Plain.com integration for Trigger.dev",
|
||||
"main": "./dist/index.js",
|
||||
"types": "./dist/index.d.ts",
|
||||
@@ -24,8 +24,8 @@
|
||||
"build:tsup": "tsup"
|
||||
},
|
||||
"dependencies": {
|
||||
"@trigger.dev/integration-kit": "workspace:^3.0.0-beta.39",
|
||||
"@trigger.dev/sdk": "workspace:^3.0.0-beta.39",
|
||||
"@trigger.dev/integration-kit": "workspace:^3.0.0-beta.44",
|
||||
"@trigger.dev/sdk": "workspace:^3.0.0-beta.44",
|
||||
"@team-plain/typescript-sdk": "^2.7.0"
|
||||
},
|
||||
"engines": {
|
||||
|
||||
@@ -1,5 +1,43 @@
|
||||
# @trigger.dev/replicate
|
||||
|
||||
## 3.0.0-beta.44
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.44
|
||||
- @trigger.dev/sdk@3.0.0-beta.44
|
||||
|
||||
## 3.0.0-beta.43
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- Updated dependencies [34ca7667d]
|
||||
- @trigger.dev/sdk@3.0.0-beta.43
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.43
|
||||
|
||||
## 3.0.0-beta.42
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- Updated dependencies [ecef19966]
|
||||
- @trigger.dev/sdk@3.0.0-beta.42
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.42
|
||||
|
||||
## 3.0.0-beta.41
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- Updated dependencies [7c36a1a4b]
|
||||
- @trigger.dev/sdk@3.0.0-beta.41
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.41
|
||||
|
||||
## 3.0.0-beta.40
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.40
|
||||
- @trigger.dev/sdk@3.0.0-beta.40
|
||||
|
||||
## 3.0.0-beta.39
|
||||
|
||||
### Patch Changes
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"name": "@trigger.dev/replicate",
|
||||
"version": "3.0.0-beta.39",
|
||||
"version": "3.0.0-beta.44",
|
||||
"description": "Trigger.dev integration for replicate",
|
||||
"main": "./dist/index.js",
|
||||
"types": "./dist/index.d.ts",
|
||||
@@ -25,8 +25,8 @@
|
||||
"typecheck": "tsc --noEmit"
|
||||
},
|
||||
"dependencies": {
|
||||
"@trigger.dev/integration-kit": "workspace:^3.0.0-beta.39",
|
||||
"@trigger.dev/sdk": "workspace:^3.0.0-beta.39",
|
||||
"@trigger.dev/integration-kit": "workspace:^3.0.0-beta.44",
|
||||
"@trigger.dev/sdk": "workspace:^3.0.0-beta.44",
|
||||
"replicate": "^0.18.1",
|
||||
"zod": "3.22.3"
|
||||
},
|
||||
|
||||
@@ -1,5 +1,43 @@
|
||||
# @trigger.dev/resend
|
||||
|
||||
## 3.0.0-beta.44
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.44
|
||||
- @trigger.dev/sdk@3.0.0-beta.44
|
||||
|
||||
## 3.0.0-beta.43
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- Updated dependencies [34ca7667d]
|
||||
- @trigger.dev/sdk@3.0.0-beta.43
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.43
|
||||
|
||||
## 3.0.0-beta.42
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- Updated dependencies [ecef19966]
|
||||
- @trigger.dev/sdk@3.0.0-beta.42
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.42
|
||||
|
||||
## 3.0.0-beta.41
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- Updated dependencies [7c36a1a4b]
|
||||
- @trigger.dev/sdk@3.0.0-beta.41
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.41
|
||||
|
||||
## 3.0.0-beta.40
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.40
|
||||
- @trigger.dev/sdk@3.0.0-beta.40
|
||||
|
||||
## 3.0.0-beta.39
|
||||
|
||||
### Patch Changes
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"name": "@trigger.dev/resend",
|
||||
"version": "3.0.0-beta.39",
|
||||
"version": "3.0.0-beta.44",
|
||||
"description": "The official Resend.com integration for Trigger.dev",
|
||||
"main": "./dist/index.js",
|
||||
"types": "./dist/index.d.ts",
|
||||
@@ -24,8 +24,8 @@
|
||||
"build:tsup": "tsup"
|
||||
},
|
||||
"dependencies": {
|
||||
"@trigger.dev/integration-kit": "workspace:^3.0.0-beta.39",
|
||||
"@trigger.dev/sdk": "workspace:^3.0.0-beta.39",
|
||||
"@trigger.dev/integration-kit": "workspace:^3.0.0-beta.44",
|
||||
"@trigger.dev/sdk": "workspace:^3.0.0-beta.44",
|
||||
"resend": "^2.1.0"
|
||||
},
|
||||
"engines": {
|
||||
|
||||
@@ -1,5 +1,43 @@
|
||||
# @trigger.dev/sendgrid
|
||||
|
||||
## 3.0.0-beta.44
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.44
|
||||
- @trigger.dev/sdk@3.0.0-beta.44
|
||||
|
||||
## 3.0.0-beta.43
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- Updated dependencies [34ca7667d]
|
||||
- @trigger.dev/sdk@3.0.0-beta.43
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.43
|
||||
|
||||
## 3.0.0-beta.42
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- Updated dependencies [ecef19966]
|
||||
- @trigger.dev/sdk@3.0.0-beta.42
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.42
|
||||
|
||||
## 3.0.0-beta.41
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- Updated dependencies [7c36a1a4b]
|
||||
- @trigger.dev/sdk@3.0.0-beta.41
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.41
|
||||
|
||||
## 3.0.0-beta.40
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.40
|
||||
- @trigger.dev/sdk@3.0.0-beta.40
|
||||
|
||||
## 3.0.0-beta.39
|
||||
|
||||
### Patch Changes
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"name": "@trigger.dev/sendgrid",
|
||||
"version": "3.0.0-beta.39",
|
||||
"version": "3.0.0-beta.44",
|
||||
"description": "Trigger.dev integration for @sendgrid/mail",
|
||||
"main": "./dist/index.js",
|
||||
"types": "./dist/index.d.ts",
|
||||
@@ -26,8 +26,8 @@
|
||||
},
|
||||
"dependencies": {
|
||||
"@sendgrid/mail": "^7.7.0",
|
||||
"@trigger.dev/sdk": "workspace:^3.0.0-beta.39",
|
||||
"@trigger.dev/integration-kit": "workspace:^3.0.0-beta.39"
|
||||
"@trigger.dev/sdk": "workspace:^3.0.0-beta.44",
|
||||
"@trigger.dev/integration-kit": "workspace:^3.0.0-beta.44"
|
||||
},
|
||||
"engines": {
|
||||
"node": ">=16.8.0"
|
||||
|
||||
@@ -1,5 +1,43 @@
|
||||
# @trigger.dev/shopify
|
||||
|
||||
## 3.0.0-beta.44
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.44
|
||||
- @trigger.dev/sdk@3.0.0-beta.44
|
||||
|
||||
## 3.0.0-beta.43
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- Updated dependencies [34ca7667d]
|
||||
- @trigger.dev/sdk@3.0.0-beta.43
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.43
|
||||
|
||||
## 3.0.0-beta.42
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- Updated dependencies [ecef19966]
|
||||
- @trigger.dev/sdk@3.0.0-beta.42
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.42
|
||||
|
||||
## 3.0.0-beta.41
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- Updated dependencies [7c36a1a4b]
|
||||
- @trigger.dev/sdk@3.0.0-beta.41
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.41
|
||||
|
||||
## 3.0.0-beta.40
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.40
|
||||
- @trigger.dev/sdk@3.0.0-beta.40
|
||||
|
||||
## 3.0.0-beta.39
|
||||
|
||||
### Patch Changes
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"name": "@trigger.dev/shopify",
|
||||
"version": "3.0.0-beta.39",
|
||||
"version": "3.0.0-beta.44",
|
||||
"description": "Trigger.dev integration for @shopify/shopify-api",
|
||||
"main": "./dist/index.js",
|
||||
"types": "./dist/index.d.ts",
|
||||
@@ -26,8 +26,8 @@
|
||||
},
|
||||
"dependencies": {
|
||||
"@shopify/shopify-api": "^8.0.2",
|
||||
"@trigger.dev/sdk": "workspace:^3.0.0-beta.39",
|
||||
"@trigger.dev/integration-kit": "workspace:^3.0.0-beta.39",
|
||||
"@trigger.dev/sdk": "workspace:^3.0.0-beta.44",
|
||||
"@trigger.dev/integration-kit": "workspace:^3.0.0-beta.44",
|
||||
"zod": "3.22.3"
|
||||
},
|
||||
"engines": {
|
||||
|
||||
@@ -1,5 +1,38 @@
|
||||
# @trigger.dev/slack
|
||||
|
||||
## 3.0.0-beta.44
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- @trigger.dev/sdk@3.0.0-beta.44
|
||||
|
||||
## 3.0.0-beta.43
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- Updated dependencies [34ca7667d]
|
||||
- @trigger.dev/sdk@3.0.0-beta.43
|
||||
|
||||
## 3.0.0-beta.42
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- Updated dependencies [ecef19966]
|
||||
- @trigger.dev/sdk@3.0.0-beta.42
|
||||
|
||||
## 3.0.0-beta.41
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- Updated dependencies [7c36a1a4b]
|
||||
- @trigger.dev/sdk@3.0.0-beta.41
|
||||
|
||||
## 3.0.0-beta.40
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- @trigger.dev/sdk@3.0.0-beta.40
|
||||
|
||||
## 3.0.0-beta.39
|
||||
|
||||
### Patch Changes
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"name": "@trigger.dev/slack",
|
||||
"version": "3.0.0-beta.39",
|
||||
"version": "3.0.0-beta.44",
|
||||
"description": "The official Slack integration for Trigger.dev",
|
||||
"main": "./dist/index.js",
|
||||
"types": "./dist/index.d.ts",
|
||||
@@ -25,7 +25,7 @@
|
||||
},
|
||||
"dependencies": {
|
||||
"@slack/web-api": "^6.8.1",
|
||||
"@trigger.dev/sdk": "workspace:^3.0.0-beta.39",
|
||||
"@trigger.dev/sdk": "workspace:^3.0.0-beta.44",
|
||||
"zod": "3.22.3"
|
||||
},
|
||||
"engines": {
|
||||
|
||||
@@ -1,5 +1,43 @@
|
||||
# @trigger.dev/stripe
|
||||
|
||||
## 3.0.0-beta.44
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.44
|
||||
- @trigger.dev/sdk@3.0.0-beta.44
|
||||
|
||||
## 3.0.0-beta.43
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- Updated dependencies [34ca7667d]
|
||||
- @trigger.dev/sdk@3.0.0-beta.43
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.43
|
||||
|
||||
## 3.0.0-beta.42
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- Updated dependencies [ecef19966]
|
||||
- @trigger.dev/sdk@3.0.0-beta.42
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.42
|
||||
|
||||
## 3.0.0-beta.41
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- Updated dependencies [7c36a1a4b]
|
||||
- @trigger.dev/sdk@3.0.0-beta.41
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.41
|
||||
|
||||
## 3.0.0-beta.40
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.40
|
||||
- @trigger.dev/sdk@3.0.0-beta.40
|
||||
|
||||
## 3.0.0-beta.39
|
||||
|
||||
### Patch Changes
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"name": "@trigger.dev/stripe",
|
||||
"version": "3.0.0-beta.39",
|
||||
"version": "3.0.0-beta.44",
|
||||
"description": "Trigger.dev integration for stripe",
|
||||
"main": "./dist/index.js",
|
||||
"types": "./dist/index.d.ts",
|
||||
@@ -25,8 +25,8 @@
|
||||
"typecheck": "tsc --noEmit"
|
||||
},
|
||||
"dependencies": {
|
||||
"@trigger.dev/integration-kit": "workspace:^3.0.0-beta.39",
|
||||
"@trigger.dev/sdk": "workspace:^3.0.0-beta.39",
|
||||
"@trigger.dev/integration-kit": "workspace:^3.0.0-beta.44",
|
||||
"@trigger.dev/sdk": "workspace:^3.0.0-beta.44",
|
||||
"stripe": "^12.14.0",
|
||||
"zod": "3.22.3"
|
||||
},
|
||||
|
||||
@@ -1,5 +1,43 @@
|
||||
# @trigger.dev/supabase
|
||||
|
||||
## 3.0.0-beta.44
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.44
|
||||
- @trigger.dev/sdk@3.0.0-beta.44
|
||||
|
||||
## 3.0.0-beta.43
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- Updated dependencies [34ca7667d]
|
||||
- @trigger.dev/sdk@3.0.0-beta.43
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.43
|
||||
|
||||
## 3.0.0-beta.42
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- Updated dependencies [ecef19966]
|
||||
- @trigger.dev/sdk@3.0.0-beta.42
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.42
|
||||
|
||||
## 3.0.0-beta.41
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- Updated dependencies [7c36a1a4b]
|
||||
- @trigger.dev/sdk@3.0.0-beta.41
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.41
|
||||
|
||||
## 3.0.0-beta.40
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.40
|
||||
- @trigger.dev/sdk@3.0.0-beta.40
|
||||
|
||||
## 3.0.0-beta.39
|
||||
|
||||
### Patch Changes
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"name": "@trigger.dev/supabase",
|
||||
"version": "3.0.0-beta.39",
|
||||
"version": "3.0.0-beta.44",
|
||||
"description": "Trigger.dev integration for @supabase/supabase-js",
|
||||
"main": "./dist/index.js",
|
||||
"types": "./dist/index.d.ts",
|
||||
@@ -26,8 +26,8 @@
|
||||
},
|
||||
"dependencies": {
|
||||
"@supabase/supabase-js": "^2.26.0",
|
||||
"@trigger.dev/integration-kit": "workspace:^3.0.0-beta.39",
|
||||
"@trigger.dev/sdk": "workspace:^3.0.0-beta.39",
|
||||
"@trigger.dev/integration-kit": "workspace:^3.0.0-beta.44",
|
||||
"@trigger.dev/sdk": "workspace:^3.0.0-beta.44",
|
||||
"supabase-management-js": "^1.0.0",
|
||||
"zod": "3.22.3"
|
||||
},
|
||||
|
||||
@@ -1,5 +1,43 @@
|
||||
# @trigger.dev/typeform
|
||||
|
||||
## 3.0.0-beta.44
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.44
|
||||
- @trigger.dev/sdk@3.0.0-beta.44
|
||||
|
||||
## 3.0.0-beta.43
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- Updated dependencies [34ca7667d]
|
||||
- @trigger.dev/sdk@3.0.0-beta.43
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.43
|
||||
|
||||
## 3.0.0-beta.42
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- Updated dependencies [ecef19966]
|
||||
- @trigger.dev/sdk@3.0.0-beta.42
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.42
|
||||
|
||||
## 3.0.0-beta.41
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- Updated dependencies [7c36a1a4b]
|
||||
- @trigger.dev/sdk@3.0.0-beta.41
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.41
|
||||
|
||||
## 3.0.0-beta.40
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- @trigger.dev/integration-kit@3.0.0-beta.40
|
||||
- @trigger.dev/sdk@3.0.0-beta.40
|
||||
|
||||
## 3.0.0-beta.39
|
||||
|
||||
### Patch Changes
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"name": "@trigger.dev/typeform",
|
||||
"version": "3.0.0-beta.39",
|
||||
"version": "3.0.0-beta.44",
|
||||
"description": "The official Typeform integration for Trigger.dev",
|
||||
"main": "./dist/index.js",
|
||||
"types": "./dist/index.d.ts",
|
||||
@@ -24,8 +24,8 @@
|
||||
"typecheck": "tsc --noEmit"
|
||||
},
|
||||
"dependencies": {
|
||||
"@trigger.dev/integration-kit": "workspace:^3.0.0-beta.39",
|
||||
"@trigger.dev/sdk": "workspace:^3.0.0-beta.39",
|
||||
"@trigger.dev/integration-kit": "workspace:^3.0.0-beta.44",
|
||||
"@trigger.dev/sdk": "workspace:^3.0.0-beta.44",
|
||||
"@typeform/api-client": "^1.8.0",
|
||||
"zod": "3.22.3"
|
||||
},
|
||||
|
||||
@@ -1,5 +1,38 @@
|
||||
# @trigger.dev/astro
|
||||
|
||||
## 3.0.0-beta.44
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- @trigger.dev/sdk@3.0.0-beta.44
|
||||
|
||||
## 3.0.0-beta.43
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- Updated dependencies [34ca7667d]
|
||||
- @trigger.dev/sdk@3.0.0-beta.43
|
||||
|
||||
## 3.0.0-beta.42
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- Updated dependencies [ecef19966]
|
||||
- @trigger.dev/sdk@3.0.0-beta.42
|
||||
|
||||
## 3.0.0-beta.41
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- Updated dependencies [7c36a1a4b]
|
||||
- @trigger.dev/sdk@3.0.0-beta.41
|
||||
|
||||
## 3.0.0-beta.40
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- @trigger.dev/sdk@3.0.0-beta.40
|
||||
|
||||
## 3.0.0-beta.39
|
||||
|
||||
### Patch Changes
|
||||
|
||||
@@ -1,7 +1,7 @@
|
||||
{
|
||||
"name": "@trigger.dev/astro",
|
||||
"description": "An Astro-native integration for Trigger.dev background jobs platform",
|
||||
"version": "3.0.0-beta.39",
|
||||
"version": "3.0.0-beta.44",
|
||||
"main": "./dist/index.js",
|
||||
"types": "./dist/index.d.ts",
|
||||
"files": [
|
||||
@@ -20,7 +20,7 @@
|
||||
"build:tsup": "tsup"
|
||||
},
|
||||
"peerDependencies": {
|
||||
"@trigger.dev/sdk": "workspace:^3.0.0-beta.39"
|
||||
"@trigger.dev/sdk": "workspace:^3.0.0-beta.44"
|
||||
},
|
||||
"devDependencies": {
|
||||
"astro": "^3.0.12",
|
||||
|
||||
@@ -1,5 +1,44 @@
|
||||
# trigger.dev
|
||||
|
||||
## 3.0.0-beta.44
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- Updated dependencies [39885a427]
|
||||
- @trigger.dev/core@3.0.0-beta.44
|
||||
|
||||
## 3.0.0-beta.43
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- 77ad4127c: Improved ESM module require error detection logic
|
||||
- Updated dependencies [34ca7667d]
|
||||
- @trigger.dev/core@3.0.0-beta.43
|
||||
|
||||
## 3.0.0-beta.42
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- @trigger.dev/core@3.0.0-beta.42
|
||||
|
||||
## 3.0.0-beta.41
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- c7a55804d: Fix jsonc-parser import
|
||||
- @trigger.dev/core@3.0.0-beta.41
|
||||
|
||||
## 3.0.0-beta.40
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- 098932ea9: v3: vercel edge runtime support
|
||||
- 9835f4ec5: v3: fix otel flushing causing CLEANUP ack timeout errors by always setting a forceFlushTimeoutMillis value
|
||||
- Updated dependencies [55d1f8c67]
|
||||
- Updated dependencies [098932ea9]
|
||||
- Updated dependencies [9835f4ec5]
|
||||
- @trigger.dev/core@3.0.0-beta.40
|
||||
|
||||
## 3.0.0-beta.39
|
||||
|
||||
### Patch Changes
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"name": "trigger.dev",
|
||||
"version": "3.0.0-beta.39",
|
||||
"version": "3.0.0-beta.44",
|
||||
"description": "A Command-Line Interface for Trigger.dev (v3) projects",
|
||||
"main": "./dist/index.js",
|
||||
"types": "./dist/index.d.ts",
|
||||
@@ -87,7 +87,7 @@
|
||||
"@opentelemetry/sdk-trace-base": "^1.22.0",
|
||||
"@opentelemetry/sdk-trace-node": "^1.22.0",
|
||||
"@opentelemetry/semantic-conventions": "^1.22.0",
|
||||
"@trigger.dev/core": "workspace:3.0.0-beta.39",
|
||||
"@trigger.dev/core": "workspace:3.0.0-beta.44",
|
||||
"@types/degit": "^2.8.3",
|
||||
"chalk": "^5.2.0",
|
||||
"chokidar": "^3.5.3",
|
||||
@@ -103,7 +103,7 @@
|
||||
"gradient-string": "^2.0.2",
|
||||
"import-meta-resolve": "^4.0.0",
|
||||
"ink": "^4.4.1",
|
||||
"jsonc-parser": "^3.2.1",
|
||||
"jsonc-parser": "3.2.1",
|
||||
"liquidjs": "^10.9.2",
|
||||
"mock-fs": "^5.2.0",
|
||||
"nanoid": "^4.0.2",
|
||||
@@ -132,4 +132,4 @@
|
||||
"engines": {
|
||||
"node": ">=18.0.0"
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -640,6 +640,10 @@ function useDev({
|
||||
backgroundWorker
|
||||
);
|
||||
} catch (e) {
|
||||
logger.debug("Error starting background worker", {
|
||||
error: e,
|
||||
});
|
||||
|
||||
if (e instanceof TaskMetadataParseError) {
|
||||
logTaskMetadataParseError(e.zodIssues, e.tasks);
|
||||
return;
|
||||
|
||||
@@ -178,7 +178,7 @@ export async function readConfig(
|
||||
write: true,
|
||||
format: "cjs",
|
||||
platform: "node",
|
||||
target: ["es2018", "node18"],
|
||||
target: ["es2020", "node18"],
|
||||
outfile: builtConfigFilePath,
|
||||
logLevel: "silent",
|
||||
plugins: [
|
||||
|
||||
@@ -28,17 +28,10 @@ export function parseBuildErrorStack(error: unknown): BuildError | undefined {
|
||||
|
||||
if (errorIsErrorLike(error)) {
|
||||
if (typeof error.stack === "string") {
|
||||
const isErrRequireEsm = error.stack.includes("ERR_REQUIRE_ESM");
|
||||
|
||||
let moduleName = null;
|
||||
|
||||
if (isErrRequireEsm) {
|
||||
// Regular expression to match the module path
|
||||
const moduleRegex = /node_modules\/(@[^\/]+\/[^\/]+|[^\/]+)\/[^\/]+\s/;
|
||||
const match = moduleRegex.exec(error.stack);
|
||||
if (match) {
|
||||
moduleName = match[1] as string; // Capture the module name
|
||||
if (error.stack.includes("ERR_REQUIRE_ESM")) {
|
||||
const moduleName = getPackageNameFromEsmRequireError(error.stack);
|
||||
|
||||
if (moduleName) {
|
||||
return {
|
||||
type: "esm-require-error",
|
||||
moduleName,
|
||||
@@ -51,6 +44,38 @@ export function parseBuildErrorStack(error: unknown): BuildError | undefined {
|
||||
}
|
||||
}
|
||||
|
||||
function getPackageNameFromEsmRequireError(stack: string): string | undefined {
|
||||
const pathRegex = /require\(\) of ES Module (.*) from/;
|
||||
const pathMatch = pathRegex.exec(stack);
|
||||
|
||||
if (!pathMatch) {
|
||||
return;
|
||||
}
|
||||
|
||||
const filePath = pathMatch[1];
|
||||
|
||||
if (!filePath) {
|
||||
return;
|
||||
}
|
||||
|
||||
const lastPart = filePath.split("node_modules/").pop();
|
||||
|
||||
if (!lastPart) {
|
||||
return;
|
||||
}
|
||||
|
||||
// regular expression to match the package name
|
||||
const moduleRegex = /(@[^\/]+\/[^\/]+|[^\/]+)/;
|
||||
|
||||
const match = moduleRegex.exec(lastPart);
|
||||
|
||||
if (!match) {
|
||||
return;
|
||||
}
|
||||
|
||||
return match[1];
|
||||
}
|
||||
|
||||
export function logESMRequireError(parsedError: ESMRequireError, resolvedConfig: ReadConfigResult) {
|
||||
logger.log(
|
||||
`\n${chalkError("X Error:")} The ${chalkPurple(
|
||||
|
||||
@@ -18,6 +18,7 @@ export const tracingSDK = new TracingSDK({
|
||||
url: process.env.OTEL_EXPORTER_OTLP_ENDPOINT ?? "http://0.0.0.0:4318",
|
||||
instrumentations: setupImportedConfig?.instrumentations ?? [],
|
||||
diagLogLevel: (process.env.OTEL_LOG_LEVEL as TracingDiagnosticLogLevel) ?? "none",
|
||||
forceFlushTimeoutMillis: 5_000,
|
||||
});
|
||||
|
||||
export const otelTracer: Tracer = tracingSDK.getTracer("trigger-dev-worker", packageJson.version);
|
||||
|
||||
@@ -863,33 +863,11 @@ class TaskRunProcess {
|
||||
}
|
||||
|
||||
#handleLog(data: Buffer) {
|
||||
if (!this._currentExecution) {
|
||||
return;
|
||||
}
|
||||
|
||||
console.log(
|
||||
`[${this.metadata.version}][${this._currentExecution.run.id}.${
|
||||
this._currentExecution.attempt.number
|
||||
}] ${data.toString()}`
|
||||
);
|
||||
console.log(data.toString());
|
||||
}
|
||||
|
||||
#handleStdErr(data: Buffer) {
|
||||
if (this._isBeingKilled) {
|
||||
return;
|
||||
}
|
||||
|
||||
if (!this._currentExecution) {
|
||||
console.error(`[${this.metadata.version}] ${data.toString()}`);
|
||||
|
||||
return;
|
||||
}
|
||||
|
||||
console.error(
|
||||
`[${this.metadata.version}][${this._currentExecution.run.id}.${
|
||||
this._currentExecution.attempt.number
|
||||
}] ${data.toString()}`
|
||||
);
|
||||
console.error(data.toString());
|
||||
}
|
||||
|
||||
async kill(signal?: number | NodeJS.Signals, timeoutInMs?: number) {
|
||||
|
||||
@@ -202,18 +202,54 @@ const zodIpc = new ZodIpcConnection({
|
||||
},
|
||||
CLEANUP: async ({ flush, kill }, sender) => {
|
||||
if (kill) {
|
||||
await Promise.all([prodUsageManager.flush(), tracingSDK.flush()]);
|
||||
await flushAll();
|
||||
// Now we need to exit the process
|
||||
await sender.send("READY_TO_DISPOSE", undefined);
|
||||
} else {
|
||||
if (flush) {
|
||||
await Promise.all([prodUsageManager.flush(), tracingSDK.flush()]);
|
||||
await flushAll();
|
||||
}
|
||||
}
|
||||
},
|
||||
},
|
||||
});
|
||||
|
||||
async function flushAll(timeoutInMs: number = 10_000) {
|
||||
const now = performance.now();
|
||||
|
||||
console.log(`Flushing at ${now}`);
|
||||
|
||||
await Promise.all([flushUsage(), flushTracingSDK()]);
|
||||
|
||||
const duration = performance.now() - now;
|
||||
|
||||
console.log(`Flushed in ${duration}ms`);
|
||||
}
|
||||
|
||||
async function flushUsage() {
|
||||
const now = performance.now();
|
||||
|
||||
console.log(`Flushing usage at ${now}`);
|
||||
|
||||
await prodUsageManager.flush();
|
||||
|
||||
const duration = performance.now() - now;
|
||||
|
||||
console.log(`Flushed usage in ${duration}ms`);
|
||||
}
|
||||
|
||||
async function flushTracingSDK() {
|
||||
const now = performance.now();
|
||||
|
||||
console.log(`Flushing tracingSDK at ${now}`);
|
||||
|
||||
await tracingSDK.flush();
|
||||
|
||||
const duration = performance.now() - now;
|
||||
|
||||
console.log(`Flushed tracingSDK in ${duration}ms`);
|
||||
}
|
||||
|
||||
// Ignore SIGTERM, handled by entry point
|
||||
process.on("SIGTERM", async () => {});
|
||||
|
||||
|
||||
@@ -16,7 +16,9 @@ export const tracingSDK = new TracingSDK({
|
||||
url: process.env.OTEL_EXPORTER_OTLP_ENDPOINT ?? "http://0.0.0.0:4318",
|
||||
instrumentations: setupImportedConfig?.instrumentations ?? [],
|
||||
diagLogLevel: (process.env.OTEL_LOG_LEVEL as TracingDiagnosticLogLevel) ?? "none",
|
||||
forceFlushTimeoutMillis: 1_000,
|
||||
forceFlushTimeoutMillis: process.env.OTEL_FORCE_FLUSH_TIMEOUT
|
||||
? parseInt(process.env.OTEL_FORCE_FLUSH_TIMEOUT, 10)
|
||||
: 5_000,
|
||||
});
|
||||
|
||||
export const otelTracer: Tracer = tracingSDK.getTracer("trigger-prod-worker", packageJson.version);
|
||||
|
||||
@@ -1,5 +1,45 @@
|
||||
# create-trigger
|
||||
|
||||
## 3.0.0-beta.44
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- Updated dependencies [39885a427]
|
||||
- @trigger.dev/core@3.0.0-beta.44
|
||||
- @trigger.dev/yalt@3.0.0-beta.44
|
||||
|
||||
## 3.0.0-beta.43
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- Updated dependencies [34ca7667d]
|
||||
- @trigger.dev/core@3.0.0-beta.43
|
||||
- @trigger.dev/yalt@3.0.0-beta.43
|
||||
|
||||
## 3.0.0-beta.42
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- @trigger.dev/core@3.0.0-beta.42
|
||||
- @trigger.dev/yalt@3.0.0-beta.42
|
||||
|
||||
## 3.0.0-beta.41
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- @trigger.dev/core@3.0.0-beta.41
|
||||
- @trigger.dev/yalt@3.0.0-beta.41
|
||||
|
||||
## 3.0.0-beta.40
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- Updated dependencies [55d1f8c67]
|
||||
- Updated dependencies [098932ea9]
|
||||
- Updated dependencies [9835f4ec5]
|
||||
- @trigger.dev/core@3.0.0-beta.40
|
||||
- @trigger.dev/yalt@3.0.0-beta.40
|
||||
|
||||
## 3.0.0-beta.39
|
||||
|
||||
### Patch Changes
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"name": "@trigger.dev/cli",
|
||||
"version": "3.0.0-beta.39",
|
||||
"version": "3.0.0-beta.44",
|
||||
"description": "The Trigger.dev CLI",
|
||||
"main": "./dist/index.js",
|
||||
"types": "./dist/index.d.ts",
|
||||
|
||||
@@ -1,5 +1,15 @@
|
||||
# @trigger.dev/core-apps
|
||||
|
||||
## 3.0.0-beta.44
|
||||
|
||||
## 3.0.0-beta.43
|
||||
|
||||
## 3.0.0-beta.42
|
||||
|
||||
## 3.0.0-beta.41
|
||||
|
||||
## 3.0.0-beta.40
|
||||
|
||||
## 3.0.0-beta.39
|
||||
|
||||
## 3.0.0-beta.38
|
||||
|
||||
Some files were not shown because too many files have changed in this diff Show More
Reference in New Issue
Block a user