Various perf improvements to prevent event loop lag (#1186)
* WIP * Handle tasks that have failed but are being auto yielded * Limit trace view to 25k event records, add a download run logs button Also added two new indexes to TaskEvent: ``` /// Used on eventRepository.getTraceSummary() @@index([traceId, startTime]) // Used for getting all logs for a run @@index([runId]) ``` * perf improvements on eventRepository.getSpan() * v2: Add a 5 minute timeout for run execution requests in dev * v3: Include presigned urls for downloading large payloads and outputs when using runs.retrieve * v3: better handle large task payloads and outputs * Change to 512KB * v2: paginate trigger schedules endpoint * v3: add 3MB limit on batch and single payloads * Update task payload and output limits
This commit is contained in:
@@ -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
|
||||
@@ -206,6 +206,9 @@ const EnvironmentSchema = z.object({
|
||||
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>;
|
||||
|
||||
@@ -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
|
||||
? {
|
||||
|
||||
+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>
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -129,6 +129,10 @@ export type PreparedEvent = Omit<QueriedEvent, "events" | "style" | "duration">
|
||||
style: TaskEventStyle;
|
||||
};
|
||||
|
||||
export type RunPreparedEvent = PreparedEvent & {
|
||||
taskSlug?: string;
|
||||
};
|
||||
|
||||
export type SpanLink =
|
||||
| {
|
||||
type: "run";
|
||||
@@ -400,13 +404,19 @@ export class EventRepository {
|
||||
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) {
|
||||
@@ -424,6 +434,8 @@ export class EventRepository {
|
||||
|
||||
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(
|
||||
@@ -432,7 +444,7 @@ export class EventRepository {
|
||||
event.duration
|
||||
);
|
||||
|
||||
return {
|
||||
const span = {
|
||||
recordId: event.id,
|
||||
id: event.spanId,
|
||||
parentId: event.parentId ?? undefined,
|
||||
@@ -451,14 +463,17 @@ export class EventRepository {
|
||||
environmentType: event.environmentType,
|
||||
},
|
||||
};
|
||||
|
||||
spansBySpanId.set(event.spanId, span);
|
||||
|
||||
return span;
|
||||
});
|
||||
|
||||
const rootSpanId = events.find((event) => !event.parentId);
|
||||
if (!rootSpanId) {
|
||||
return;
|
||||
}
|
||||
|
||||
const rootSpan = spans.find((span) => span.id === rootSpanId.spanId);
|
||||
const rootSpan = spansBySpanId.get(rootSpanId);
|
||||
|
||||
if (!rootSpan) {
|
||||
return;
|
||||
@@ -471,85 +486,262 @@ export class EventRepository {
|
||||
});
|
||||
}
|
||||
|
||||
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",
|
||||
},
|
||||
});
|
||||
|
||||
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) {
|
||||
const traceSummary = await this.getTraceSummary(traceId);
|
||||
return await startActiveSpan("getSpan", async (s) => {
|
||||
const spanEvent = await this.#getSpanEvent(spanId);
|
||||
|
||||
const span = traceSummary?.spans.find((span) => span.id === spanId);
|
||||
if (!spanEvent) {
|
||||
return;
|
||||
}
|
||||
|
||||
if (!span) {
|
||||
return;
|
||||
}
|
||||
const preparedEvent = prepareEvent(spanEvent);
|
||||
|
||||
const fullEvent = await this.readReplica.taskEvent.findUnique({
|
||||
where: {
|
||||
id: span.recordId,
|
||||
},
|
||||
});
|
||||
const span = await this.#createSpanFromEvent(preparedEvent);
|
||||
|
||||
if (!fullEvent) {
|
||||
return;
|
||||
}
|
||||
const output = rehydrateJson(spanEvent.output);
|
||||
const payload = rehydrateJson(spanEvent.payload);
|
||||
|
||||
const output = rehydrateJson(fullEvent.output);
|
||||
const payload = rehydrateJson(fullEvent.payload);
|
||||
const show = rehydrateShow(spanEvent.properties);
|
||||
|
||||
const show = rehydrateShow(fullEvent.properties);
|
||||
const properties = sanitizedAttributes(spanEvent.properties);
|
||||
|
||||
const properties = sanitizedAttributes(fullEvent.properties);
|
||||
const messagingEvent = SpanMessagingEvent.optional().safeParse(
|
||||
(properties as any)?.messaging
|
||||
);
|
||||
|
||||
const messagingEvent = SpanMessagingEvent.optional().safeParse((properties as any)?.messaging);
|
||||
const links: SpanLink[] = [];
|
||||
|
||||
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 (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 = fullEvent.links as any as Link[] | undefined;
|
||||
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"
|
||||
);
|
||||
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,
|
||||
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,
|
||||
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,
|
||||
},
|
||||
};
|
||||
|
||||
return span;
|
||||
});
|
||||
}
|
||||
|
||||
async #walkSpanAncestors(
|
||||
event: PreparedEvent,
|
||||
callback: (event: PreparedEvent, level: number) => { stop: boolean }
|
||||
) {
|
||||
const parentId = event.parentId;
|
||||
if (!parentId) {
|
||||
return;
|
||||
}
|
||||
|
||||
const events = transformEvents(
|
||||
span.data.events,
|
||||
fullEvent.metadata as Attributes,
|
||||
traceSummary?.rootSpan.data.environmentType === "DEVELOPMENT"
|
||||
);
|
||||
await startActiveSpan("walkSpanAncestors", async (s) => {
|
||||
let parentEvent = await this.#getSpanEvent(parentId);
|
||||
let level = 1;
|
||||
|
||||
return {
|
||||
...fullEvent,
|
||||
...span.data,
|
||||
payload,
|
||||
output,
|
||||
properties,
|
||||
events,
|
||||
show,
|
||||
links,
|
||||
};
|
||||
while (parentEvent) {
|
||||
const preparedParentEvent = prepareEvent(parentEvent);
|
||||
|
||||
const result = callback(preparedParentEvent, level);
|
||||
|
||||
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 [];
|
||||
}
|
||||
|
||||
if (!event.parentId) {
|
||||
return [];
|
||||
}
|
||||
|
||||
const parentEvent = await this.#getSpanEvent(event.parentId);
|
||||
|
||||
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",
|
||||
},
|
||||
});
|
||||
|
||||
let finalEvent: TaskEvent | undefined;
|
||||
|
||||
for (const event of events) {
|
||||
if (event.isPartial) {
|
||||
continue;
|
||||
}
|
||||
|
||||
if (event.spanId === spanId) {
|
||||
finalEvent = event;
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
return finalEvent;
|
||||
});
|
||||
}
|
||||
|
||||
public async recordEvent(message: string, options: TraceEventOptions) {
|
||||
@@ -1224,7 +1416,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);
|
||||
}
|
||||
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
@@ -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;
|
||||
@@ -257,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 {
|
||||
|
||||
@@ -117,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: 5000,
|
||||
forceFlushTimeoutMillis: 15_000,
|
||||
resource: new Resource({
|
||||
[SEMRESATTRS_SERVICE_NAME]: env.SERVICE_NAME,
|
||||
}),
|
||||
@@ -129,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
|
||||
|
||||
@@ -1722,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.
|
||||
|
||||
@@ -536,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.
|
||||
|
||||
@@ -85,6 +85,14 @@
|
||||
"require": "./dist/v3/utils/structuredLogger.js",
|
||||
"types": "./dist/v3/utils/structuredLogger.d.ts"
|
||||
},
|
||||
"./v3/utils/durations": {
|
||||
"import": {
|
||||
"types": "./dist/v3/utils/durations.d.mts",
|
||||
"default": "./dist/v3/utils/durations.mjs"
|
||||
},
|
||||
"require": "./dist/v3/utils/durations.js",
|
||||
"types": "./dist/v3/utils/durations.d.ts"
|
||||
},
|
||||
"./v3/dev": {
|
||||
"import": {
|
||||
"types": "./dist/v3/dev/index.d.mts",
|
||||
|
||||
@@ -431,7 +431,9 @@ const CommonRunFields = {
|
||||
export const RetrieveRunResponse = z.object({
|
||||
...CommonRunFields,
|
||||
payload: z.any().optional(),
|
||||
payloadPresignedUrl: z.string().optional(),
|
||||
output: z.any().optional(),
|
||||
outputPresignedUrl: z.string().optional(),
|
||||
schedule: RunScheduleDetails.optional(),
|
||||
attempts: z.array(
|
||||
z
|
||||
|
||||
@@ -85,7 +85,10 @@ export async function conditionallyExportPacket(
|
||||
return packet;
|
||||
}
|
||||
|
||||
export function packetRequiresOffloading(packet: IOPacket): {
|
||||
export function packetRequiresOffloading(
|
||||
packet: IOPacket,
|
||||
lengthLimit?: number
|
||||
): {
|
||||
needsOffloading: boolean;
|
||||
size: number;
|
||||
} {
|
||||
@@ -99,7 +102,7 @@ export function packetRequiresOffloading(packet: IOPacket): {
|
||||
const byteSize = Buffer.byteLength(packet.data, "utf8");
|
||||
|
||||
return {
|
||||
needsOffloading: byteSize >= OFFLOAD_IO_PACKET_LENGTH_LIMIT,
|
||||
needsOffloading: byteSize >= (lengthLimit ?? OFFLOAD_IO_PACKET_LENGTH_LIMIT),
|
||||
size: byteSize,
|
||||
};
|
||||
}
|
||||
@@ -128,8 +131,6 @@ async function exportPacket(packet: IOPacket, pathPrefix: string): Promise<IOPac
|
||||
data: filename,
|
||||
dataType: "application/store",
|
||||
};
|
||||
|
||||
return packet;
|
||||
}
|
||||
|
||||
export async function conditionallyImportPacket(
|
||||
@@ -186,8 +187,6 @@ async function importPacket(packet: IOPacket, span?: Span): Promise<IOPacket> {
|
||||
data,
|
||||
dataType: response.headers.get("content-type") ?? "application/json",
|
||||
};
|
||||
|
||||
return packet;
|
||||
}
|
||||
|
||||
export async function createPacketAttributes(
|
||||
|
||||
@@ -12,6 +12,7 @@ export default defineConfig({
|
||||
"./src/v3/zodSocket.ts",
|
||||
"./src/v3/zodIpc.ts",
|
||||
"./src/v3/utils/structuredLogger.ts",
|
||||
"./src/v3/utils/durations.ts",
|
||||
"./src/v3/dev/index.ts",
|
||||
"./src/v3/prod/index.ts",
|
||||
"./src/v3/workers/index.ts",
|
||||
|
||||
+2
@@ -0,0 +1,2 @@
|
||||
-- CreateIndex
|
||||
CREATE INDEX "TaskEvent_traceId_startTime_idx" ON "TaskEvent"("traceId", "startTime");
|
||||
+2
@@ -0,0 +1,2 @@
|
||||
-- CreateIndex
|
||||
CREATE INDEX "TaskEvent_runId_idx" ON "TaskEvent"("runId");
|
||||
@@ -1925,6 +1925,10 @@ model TaskEvent {
|
||||
@@index([traceId])
|
||||
/// Used when looking up span events to complete when a run completes
|
||||
@@index([spanId])
|
||||
/// Used on eventRepository.getTraceSummary()
|
||||
@@index([traceId, startTime])
|
||||
// Used for getting all logs for a run
|
||||
@@index([runId])
|
||||
}
|
||||
|
||||
enum TaskEventLevel {
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
@@ -1,12 +1,13 @@
|
||||
import { createExpressServer } from "@trigger.dev/express";
|
||||
import { TriggerClient, eventTrigger } from "@trigger.dev/sdk";
|
||||
import { readFile } from "node:fs/promises";
|
||||
|
||||
export const client = new TriggerClient({
|
||||
id: "job-catalog",
|
||||
apiKey: process.env["TRIGGER_API_KEY"],
|
||||
apiUrl: process.env["TRIGGER_API_URL"],
|
||||
verbose: true,
|
||||
ioLogLocalEnabled: true,
|
||||
ioLogLocalEnabled: false,
|
||||
});
|
||||
|
||||
client.defineJob({
|
||||
@@ -141,6 +142,40 @@ client.defineJob({
|
||||
},
|
||||
});
|
||||
|
||||
client.defineJob({
|
||||
id: "stress-test-event-loop-lag",
|
||||
name: "Stress Test Event Loop Lag",
|
||||
version: "1.0.0",
|
||||
trigger: eventTrigger({
|
||||
name: "stress.test.event.loop.lag",
|
||||
}),
|
||||
run: async (payload, io, ctx) => {
|
||||
const photos = await io.runTask(
|
||||
`task-1`,
|
||||
async (task) => {
|
||||
const file = await readFile("./fixtures/large.json");
|
||||
|
||||
return JSON.parse(file.toString());
|
||||
},
|
||||
{ name: `Task 1` }
|
||||
);
|
||||
|
||||
await io.runTask(`task-2`, async (task) => {}, { name: `Task 2`, params: { photos } });
|
||||
|
||||
const morePhotos = await io.runTask(
|
||||
`task-3`,
|
||||
async (task) => {
|
||||
const file = await readFile("./fixtures/toolarge.json");
|
||||
|
||||
return JSON.parse(file.toString());
|
||||
},
|
||||
{ name: `Task 3` }
|
||||
);
|
||||
|
||||
await io.runTask(`task-4`, async (task) => {}, { name: `Task 4`, params: { morePhotos } });
|
||||
},
|
||||
});
|
||||
|
||||
client.defineJob({
|
||||
id: "long.running",
|
||||
name: "Long Running Job",
|
||||
|
||||
@@ -7,7 +7,8 @@
|
||||
"management": "ts-node -r tsconfig-paths/register ./src/management.ts",
|
||||
"queues": "ts-node -r tsconfig-paths/register ./src/queues.ts",
|
||||
"build:client": "tsup-node ./src/clientUsage.ts --format esm,cjs",
|
||||
"client": "ts-node -r dotenv/config -r tsconfig-paths/register ./src/clientUsage.ts"
|
||||
"client": "ts-node -r dotenv/config -r tsconfig-paths/register ./src/clientUsage.ts",
|
||||
"triggerWithLargePayload": "ts-node -r dotenv/config -r tsconfig-paths/register ./src/triggerWithLargePayload.ts"
|
||||
},
|
||||
"dependencies": {
|
||||
"@ffmpeg-installer/ffmpeg": "^1.1.0",
|
||||
|
||||
@@ -1,32 +0,0 @@
|
||||
import { createEnv } from "@t3-oss/env-nextjs";
|
||||
import { z } from "zod";
|
||||
|
||||
export const env = createEnv({
|
||||
/*
|
||||
* Serverside Environment variables, not available on the client.
|
||||
* Will throw if you access these variables on the client.
|
||||
*/
|
||||
server: {
|
||||
TRIGGER_SECRET_KEY: z.string(),
|
||||
OPENAI_API_KEY: z.string().min(1),
|
||||
},
|
||||
/*
|
||||
* Environment variables available on the client (and server).
|
||||
*
|
||||
* 💡 You'll get type errors if these are not prefixed with NEXT_PUBLIC_.
|
||||
*/
|
||||
client: {
|
||||
NEXT_PUBLIC_TEST: z.string().min(1).optional(),
|
||||
},
|
||||
/*
|
||||
* Due to how Next.js bundles environment variables on Edge and Client,
|
||||
* we need to manually destructure them to make sure all are included in bundle.
|
||||
*
|
||||
* 💡 You'll get type errors if not all variables from `server` & `client` are included here.
|
||||
*/
|
||||
runtimeEnv: {
|
||||
TRIGGER_SECRET_KEY: process.env.TRIGGER_SECRET_KEY,
|
||||
OPENAI_API_KEY: process.env.OPENAI_API_KEY,
|
||||
NEXT_PUBLIC_TEST: process.env.NEXT_PUBLIC_TEST,
|
||||
},
|
||||
});
|
||||
@@ -1,14 +1,12 @@
|
||||
import { logger, task, wait } from "@trigger.dev/sdk/v3";
|
||||
|
||||
import { env } from "../env";
|
||||
|
||||
export const oneAtATime = task({
|
||||
id: "on-at-a-time",
|
||||
queue: {
|
||||
concurrencyLimit: 1,
|
||||
},
|
||||
run: async (payload: { message: string }) => {
|
||||
logger.info("One at a time task payload", { payload, env });
|
||||
logger.info("One at a time task payload", { payload });
|
||||
|
||||
await wait.for({ seconds: 10 });
|
||||
|
||||
|
||||
@@ -60,3 +60,22 @@ export const longRunningWithLotsOfLogs = task({
|
||||
};
|
||||
},
|
||||
});
|
||||
|
||||
export const longRunningWith100kLogs = task({
|
||||
id: "100k-logs",
|
||||
run: async (payload: { message: string }) => {
|
||||
// Log 100000 times over 60 seconds
|
||||
for (let i = 0; i < 100000; i++) {
|
||||
console.log("Log number " + i);
|
||||
|
||||
if (i % 512 === 0) {
|
||||
// wait every 512 logs
|
||||
await new Promise((resolve) => setTimeout(resolve, 1000));
|
||||
}
|
||||
}
|
||||
|
||||
return {
|
||||
finished: new Date().toISOString(),
|
||||
};
|
||||
},
|
||||
});
|
||||
|
||||
@@ -18,6 +18,15 @@ export const simplestTask = task({
|
||||
},
|
||||
});
|
||||
|
||||
export const anyPayloadTask = task({
|
||||
id: "any-payload-task",
|
||||
run: async (payload: any) => {
|
||||
return {
|
||||
payload,
|
||||
};
|
||||
},
|
||||
});
|
||||
|
||||
export const taskWithSpecialCharacters = task({
|
||||
id: "admin:special-characters",
|
||||
run: async (payload: { url: string }) => {
|
||||
|
||||
@@ -136,3 +136,45 @@ export const taskWithNoPayload = task({
|
||||
return { hello: "world" };
|
||||
},
|
||||
});
|
||||
|
||||
export const deeplyNestedTaskParent = task({
|
||||
id: "deeply-nested-task-parent",
|
||||
run: async ({ message = "test" }: { message?: string }) => {
|
||||
await deeplyNestedTaskChild.triggerAndWait({ message: `${message} - 1` });
|
||||
|
||||
return {
|
||||
hello: "world",
|
||||
};
|
||||
},
|
||||
});
|
||||
|
||||
export const deeplyNestedTaskChild = task({
|
||||
id: "deeply-nested-task-child",
|
||||
run: async ({ message = "test" }: { message?: string }) => {
|
||||
await deeplyNestedTaskGrandchild.triggerAndWait({ message: `${message} - 2` });
|
||||
|
||||
return {
|
||||
hello: "world",
|
||||
};
|
||||
},
|
||||
});
|
||||
|
||||
export const deeplyNestedTaskGrandchild = task({
|
||||
id: "deeply-nested-task-grandchild",
|
||||
run: async ({ message = "test" }: { message?: string }) => {
|
||||
await deeplyNestedTaskGreatGrandchild.triggerAndWait({ message: `${message} - 3` });
|
||||
|
||||
return {
|
||||
hello: "world",
|
||||
};
|
||||
},
|
||||
});
|
||||
|
||||
export const deeplyNestedTaskGreatGrandchild = task({
|
||||
id: "deeply-nested-task-great-grandchild",
|
||||
run: async ({ message = "test" }: { message?: string }) => {
|
||||
return {
|
||||
hello: "world",
|
||||
};
|
||||
},
|
||||
});
|
||||
|
||||
@@ -0,0 +1,37 @@
|
||||
import { tasks, runs } from "@trigger.dev/sdk/v3";
|
||||
import { anyPayloadTask } from "./trigger/simple";
|
||||
|
||||
async function main() {
|
||||
const payload = createLargePayload(100000);
|
||||
|
||||
const anyHandle = await tasks.trigger<typeof anyPayloadTask>("any-payload-task", payload);
|
||||
|
||||
const anyRun = await runs.poll(anyHandle);
|
||||
|
||||
if (anyRun.outputPresignedUrl) {
|
||||
const response = await fetch(anyRun.outputPresignedUrl);
|
||||
|
||||
console.log("## Output");
|
||||
console.log(await response.text());
|
||||
}
|
||||
|
||||
if (anyRun.payloadPresignedUrl) {
|
||||
const response = await fetch(anyRun.payloadPresignedUrl);
|
||||
|
||||
console.log("## Payload");
|
||||
console.log(await response.text());
|
||||
}
|
||||
}
|
||||
|
||||
main().catch(console.error);
|
||||
|
||||
// Creates a large object payload, with many keys and values
|
||||
function createLargePayload(size: number) {
|
||||
const payload: Record<string, any> = {};
|
||||
|
||||
for (let i = 0; i < size; i++) {
|
||||
payload[`key-${i}`] = `value-${i}`;
|
||||
}
|
||||
|
||||
return payload;
|
||||
}
|
||||
@@ -14,8 +14,6 @@
|
||||
},
|
||||
"emitDecoratorMetadata": true,
|
||||
"experimentalDecorators": true,
|
||||
"allowJs": true,
|
||||
"moduleResolution": "Bundler",
|
||||
"module": "ES2015"
|
||||
"allowJs": true
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user