@@ -378,170 +387,247 @@ function TasksTreeView({
{/* Timeline */}
-
-
- {/* Follows the cursor */}
-
+
+
+
+
+ );
+}
-
- {/* The duration labels */}
-
-
-
- {(ms: number, index: number) => {
- if (index === tickCount - 1) return null;
- return (
-
- {(ms) => (
-
- {formatDurationMilliseconds(ms, {
- style: "short",
- maxDecimalPoints: ms < 1000 ? 0 : 1,
- })}
-
- )}
-
- );
- }}
-
- {rootSpanStatus !== "executing" && (
-
+type TimelineViewProps = Pick<
+ TasksTreeViewProps,
+ "totalDuration" | "rootSpanStatus" | "events" | "rootStartedAt"
+> & {
+ scale: number;
+ parentRef: React.RefObject;
+ timelineScrollRef: React.RefObject;
+ virtualizer: Virtualizer;
+ nodes: NodesState;
+ getNodeProps: UseTreeStateOutput["getNodeProps"];
+ getTreeProps: UseTreeStateOutput["getTreeProps"];
+ toggleNodeSelection: UseTreeStateOutput["toggleNodeSelection"];
+ showDurations: boolean;
+ treeScrollRef: React.RefObject;
+};
+
+const tickCount = 5;
+
+function TimelineView({
+ totalDuration,
+ scale,
+ rootSpanStatus,
+ rootStartedAt,
+ parentRef,
+ timelineScrollRef,
+ virtualizer,
+ events,
+ nodes,
+ getNodeProps,
+ getTreeProps,
+ toggleNodeSelection,
+ showDurations,
+ treeScrollRef,
+}: TimelineViewProps) {
+ const timelineContainerRef = useRef(null);
+ const initialTimelineDimensions = useInitialDimensions(timelineContainerRef);
+ const minTimelineWidth = initialTimelineDimensions?.width ?? 300;
+ const maxTimelineWidth = minTimelineWidth * 10;
+
+ //we want to live-update the duration if the root span is still executing
+ const [duration, setDuration] = useState(totalDuration);
+ useEffect(() => {
+ if (rootSpanStatus !== "executing" || !rootStartedAt) {
+ setDuration(totalDuration);
+ return;
+ }
+
+ const interval = setInterval(() => {
+ setDuration(millisecondsToNanoseconds(Date.now() - rootStartedAt.getTime()));
+ }, 500);
+
+ return () => clearInterval(interval);
+ }, [totalDuration, rootSpanStatus]);
+
+ return (
+
+
+ {/* Follows the cursor */}
+
+
+
+ {/* The duration labels */}
+
+
+
+ {(ms: number, index: number) => {
+ if (index === tickCount - 1) return null;
+ return (
+
+ {(ms) => (
+
+ {formatDurationMilliseconds(ms, {
+ style: "short",
+ maxDecimalPoints: ms < 1000 ? 0 : 1,
+ })}
+
+ )}
+
+ );
+ }}
+
+ {rootSpanStatus !== "executing" && (
+
+ {(ms) => (
+
+ {formatDurationMilliseconds(ms, {
+ style: "short",
+ maxDecimalPoints: ms < 1000 ? 0 : 1,
+ })}
+
+ )}
+
+ )}
+
+
+
+ {(ms: number, index: number) => {
+ if (index === 0 || index === tickCount - 1) return null;
+ return (
+
+ );
+ }}
+
+
+
+
+ {/* Main timeline body */}
+
+ {/* The vertical tick lines */}
+
+ {(ms: number, index: number) => {
+ if (index === 0) return null;
+ return ;
+ }}
+
+ {/* The completed line */}
+ {rootSpanStatus !== "executing" && (
+
+ )}
+ {
+ return (
+ console.log(`hover ${index}`)}
+ onClick={(e) => {
+ toggleNodeSelection(node.id);
+ }}
+ >
+ {node.data.level === "TRACE" ? (
+
+ ) : (
+
{(ms) => (
-
- {formatDurationMilliseconds(ms, {
- style: "short",
- maxDecimalPoints: ms < 1000 ? 0 : 1,
- })}
-
+
)}
)}
-
-
- {(ms: number, index: number) => {
- if (index === 0 || index === tickCount - 1) return null;
- return (
-
- );
- }}
-
-
-
-
- {/* Main timeline body */}
-
- {/* The vertical tick lines */}
-
- {(ms: number, index: number) => {
- if (index === 0) return null;
- return (
-
- );
- }}
-
- {/* The completed line */}
- {rootSpanStatus !== "executing" && (
-
- )}
- {
- return (
- console.log(`hover ${index}`)}
- onClick={(e) => {
- toggleNodeSelection(node.id);
- }}
- >
- {node.data.level === "TRACE" ? (
-
- ) : (
-
- )}
-
- );
- }}
- onScroll={(scrollTop) => {
- //sync the scroll to the tree
- if (treeScrollRef.current && treeScrollRef.current.scrollTop !== scrollTop) {
- treeScrollRef.current.scrollTop = scrollTop;
- }
- }}
- />
-
-
-
-
-
-
+ );
+ }}
+ onScroll={(scrollTop) => {
+ //sync the scroll to the tree
+ if (treeScrollRef.current && treeScrollRef.current.scrollTop !== scrollTop) {
+ treeScrollRef.current.scrollTop = scrollTop;
+ }
+ }}
+ />
+
+
+
);
}
@@ -659,11 +745,12 @@ function SpanWithDuration({
}: Timeline.SpanProps & { node: RunEvent; showDuration: boolean }) {
return (
-
{node.data.isPartial && (
-
+
);
}
diff --git a/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.v3.$projectParam.tasks.stream/route.tsx b/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.v3.$projectParam.tasks.stream/route.tsx
new file mode 100644
index 000000000..e16ba2bac
--- /dev/null
+++ b/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.v3.$projectParam.tasks.stream/route.tsx
@@ -0,0 +1,13 @@
+import type { LoaderFunctionArgs } from "@remix-run/server-runtime";
+import { TasksStreamPresenter } from "~/presenters/v3/TasksStreamPresenter.server";
+import { requireUserId } from "~/services/session.server";
+import { ProjectParamSchema } from "~/utils/pathBuilder";
+
+export async function loader({ request, params }: LoaderFunctionArgs) {
+ const userId = await requireUserId(request);
+
+ const { organizationSlug, projectParam } = ProjectParamSchema.parse(params);
+
+ const presenter = new TasksStreamPresenter();
+ return presenter.call({ request, projectSlug: projectParam, organizationSlug, userId });
+}
diff --git a/apps/webapp/app/routes/otel.v1.logs.ts b/apps/webapp/app/routes/otel.v1.logs.ts
index 0554a2942..3bbafab9b 100644
--- a/apps/webapp/app/routes/otel.v1.logs.ts
+++ b/apps/webapp/app/routes/otel.v1.logs.ts
@@ -1,13 +1,25 @@
-import { ActionFunctionArgs } from "@remix-run/server-runtime";
+import { ActionFunctionArgs, json } from "@remix-run/server-runtime";
import { ExportLogsServiceRequest, ExportLogsServiceResponse } from "@trigger.dev/otlp-importer";
import { otlpExporter } from "~/v3/otlpExporter.server";
export async function action({ request }: ActionFunctionArgs) {
- const buffer = await request.arrayBuffer();
+ const contentType = request.headers.get("content-type");
- const exportRequest = ExportLogsServiceRequest.decode(new Uint8Array(buffer));
+ if (contentType === "application/json") {
+ const body = await request.json();
- const exportResponse = await otlpExporter.exportLogs(exportRequest);
+ const exportResponse = await otlpExporter.exportLogs(body as ExportLogsServiceRequest, true);
- return new Response(ExportLogsServiceResponse.encode(exportResponse).finish(), { status: 200 });
+ return json(exportResponse, { status: 200 })
+ } else if (contentType === "application/x-protobuf") {
+ const buffer = await request.arrayBuffer();
+
+ const exportRequest = ExportLogsServiceRequest.decode(new Uint8Array(buffer));
+
+ const exportResponse = await otlpExporter.exportLogs(exportRequest);
+
+ return new Response(ExportLogsServiceResponse.encode(exportResponse).finish(), { status: 200 });
+ } else {
+ return new Response("Unsupported content type. Must be either application/x-protobuf or application/json", { status: 400 });
+ }
}
diff --git a/apps/webapp/app/routes/otel.v1.traces.ts b/apps/webapp/app/routes/otel.v1.traces.ts
index 2e0677877..6dfc8cb34 100644
--- a/apps/webapp/app/routes/otel.v1.traces.ts
+++ b/apps/webapp/app/routes/otel.v1.traces.ts
@@ -1,13 +1,25 @@
-import { ActionFunctionArgs } from "@remix-run/server-runtime";
+import { ActionFunctionArgs, json } from "@remix-run/server-runtime";
import { ExportTraceServiceRequest, ExportTraceServiceResponse } from "@trigger.dev/otlp-importer";
import { otlpExporter } from "~/v3/otlpExporter.server";
export async function action({ request }: ActionFunctionArgs) {
- const buffer = await request.arrayBuffer();
+ const contentType = request.headers.get("content-type");
- const exportRequest = ExportTraceServiceRequest.decode(new Uint8Array(buffer));
+ if (contentType === "application/json") {
+ const body = await request.json();
- const exportResponse = await otlpExporter.exportTraces(exportRequest);
+ const exportResponse = await otlpExporter.exportTraces(body as ExportTraceServiceRequest, true);
- return new Response(ExportTraceServiceResponse.encode(exportResponse).finish(), { status: 200 });
+ return json(exportResponse, { status: 200 })
+ } else if (contentType === "application/x-protobuf") {
+ const buffer = await request.arrayBuffer();
+
+ const exportRequest = ExportTraceServiceRequest.decode(new Uint8Array(buffer));
+
+ const exportResponse = await otlpExporter.exportTraces(exportRequest);
+
+ return new Response(ExportTraceServiceResponse.encode(exportResponse).finish(), { status: 200 });
+ } else {
+ return new Response("Unsupported content type. Must be either application/x-protobuf or application/json", { status: 400 });
+ }
}
diff --git a/apps/webapp/app/utils/pathBuilder.ts b/apps/webapp/app/utils/pathBuilder.ts
index 8961e4eb2..67f396a4c 100644
--- a/apps/webapp/app/utils/pathBuilder.ts
+++ b/apps/webapp/app/utils/pathBuilder.ts
@@ -301,6 +301,10 @@ export function v3ProjectPath(organization: OrgForPath, project: ProjectForPath)
return `/orgs/${organizationParam(organization)}/projects/v3/${projectParam(project)}`;
}
+export function v3TasksStreamingPath(organization: OrgForPath, project: ProjectForPath) {
+ return `${v3ProjectPath(organization, project)}/tasks/stream`;
+}
+
export function v3ApiKeysPath(organization: OrgForPath, project: ProjectForPath) {
return `${v3ProjectPath(organization, project)}/apikeys`;
}
diff --git a/apps/webapp/app/v3/eventRepository.server.ts b/apps/webapp/app/v3/eventRepository.server.ts
index 7f560ab3c..a87ab0b9c 100644
--- a/apps/webapp/app/v3/eventRepository.server.ts
+++ b/apps/webapp/app/v3/eventRepository.server.ts
@@ -67,7 +67,7 @@ export type TraceEventOptions = {
attributes: TraceAttributes;
environment: AuthenticatedEnvironment;
taskSlug: string;
- startTime?: Date;
+ startTime?: bigint;
endTime?: Date;
immediate?: boolean;
};
@@ -153,6 +153,10 @@ export class EventRepository {
this._flushScheduler.addToBatch(events);
}
+ async insertManyImmediate(events: CreatableEvent[]) {
+ return await this.#flushBatch(events);
+ }
+
async completeEvent(spanId: string, options?: UpdateEventOptions) {
const events = await this.queryIncompleteEvents({ spanId });
@@ -172,8 +176,7 @@ export class EventRepository {
status: options?.attributes.isError ? "ERROR" : "OK",
links: event.links ?? [],
events: event.events ?? [],
- duration:
- ((options?.endTime ?? new Date()).getTime() - event.startTime.getTime()) * 1_000_000, // convert to nanoseconds
+ duration: calculateDurationFromStart(event.startTime, options?.endTime),
properties: event.properties as Attributes,
metadata: event.metadata as Attributes,
style: event.style as Attributes,
@@ -205,7 +208,7 @@ export class EventRepository {
},
...((event.events as any[]) ?? []),
],
- duration: (cancelledAt.getTime() - event.startTime.getTime()) * 1_000_000, // convert to nanoseconds
+ duration: calculateDurationFromStart(event.startTime, cancelledAt),
properties: event.properties as Attributes,
metadata: event.metadata as Attributes,
style: event.style as Attributes,
@@ -285,7 +288,7 @@ export class EventRepository {
isError: event.isError,
isPartial: ancestorCancelled ? false : event.isPartial,
isCancelled: event.isCancelled === true ? true : event.isPartial && ancestorCancelled,
- startTime: event.startTime,
+ startTime: getDateFromNanoseconds(event.startTime),
level: event.level,
events: event.events,
},
@@ -368,8 +371,8 @@ export class EventRepository {
public async recordEvent(message: string, options: TraceEventOptions) {
const propagatedContext = extractContextFromCarrier(options.context ?? {});
- const startTime = options.startTime ?? new Date();
- const durationInMs = options.endTime ? options.endTime.getTime() - startTime.getTime() : 100;
+ const startTime = options.startTime ?? getNowInNanoseconds();
+ const duration = options.endTime ? calculateDurationFromStart(startTime, options.endTime) : 100;
const traceId = propagatedContext?.traceparent?.traceId ?? this.generateTraceId();
const parentId = propagatedContext?.traceparent?.spanId;
@@ -414,7 +417,7 @@ export class EventRepository {
status: "OK",
startTime,
isPartial: false,
- duration: durationInMs * 1_000_000, // convert to nanoseconds
+ duration, // convert to nanoseconds
environmentId: options.environment.id,
environmentType: options.environment.type,
organizationId: options.environment.organizationId,
@@ -459,7 +462,7 @@ export class EventRepository {
const propagatedContext = extractContextFromCarrier(options.context ?? {});
const start = process.hrtime.bigint();
- const startTime = new Date();
+ const startTime = getNowInNanoseconds();
const traceId = options.spanParentAsLink
? this.generateTraceId()
@@ -477,14 +480,14 @@ export class EventRepository {
const links: Link[] =
options.spanParentAsLink && propagatedContext?.traceparent
? [
- {
- context: {
- traceId: propagatedContext.traceparent.traceId,
- spanId: propagatedContext.traceparent.spanId,
- traceFlags: TraceFlags.SAMPLED,
- },
+ {
+ context: {
+ traceId: propagatedContext.traceparent.traceId,
+ spanId: propagatedContext.traceparent.spanId,
+ traceFlags: TraceFlags.SAMPLED,
},
- ]
+ },
+ ]
: [];
const eventBuilder = {
@@ -548,7 +551,7 @@ export class EventRepository {
level: "TRACE",
kind: options.kind,
status: "OK",
- startTime: startTime,
+ startTime,
environmentId: options.environment.id,
environmentType: options.environment.type,
organizationId: options.environment.organizationId,
@@ -739,9 +742,9 @@ function prepareEvent(event: QueriedEvent): PreparedEvent {
function parseEventsField(events: Prisma.JsonValue): SpanEvents {
const eventsUnflattened = events
? (events as any[]).map((e) => ({
- ...e,
- properties: unflattenAttributes(e.properties as Attributes),
- }))
+ ...e,
+ properties: unflattenAttributes(e.properties as Attributes),
+ }))
: undefined;
const spanEvents = SpanEvents.safeParse(eventsUnflattened);
@@ -810,7 +813,7 @@ function calculateDurationIfAncestorIsCancelled(
);
if (cancellationEvent) {
- return (cancellationEvent.time.getTime() - event.startTime.getTime()) * 1_000_000;
+ return calculateDurationFromStart(event.startTime, cancellationEvent.time);
}
}
}
@@ -935,8 +938,8 @@ function transformException(
...exception,
stacktrace: exception.stacktrace
? correctErrorStackTrace(exception.stacktrace, projectDirAttributeValue, {
- removeFirstLine: true,
- })
+ removeFirstLine: true,
+ })
: undefined,
};
}
@@ -952,3 +955,15 @@ function filteredAttributes(attributes: Attributes, prefix: string): Attributes
return result;
}
+
+function calculateDurationFromStart(startTime: bigint, endTime: Date = new Date()) {
+ return Number(BigInt(endTime.getTime() * 1_000_000) - startTime);
+}
+
+function getNowInNanoseconds(): bigint {
+ return BigInt(new Date().getTime() * 1_000_000);
+}
+
+function getDateFromNanoseconds(nanoseconds: bigint) {
+ return new Date(Number(nanoseconds) / 1_000_000);
+}
\ No newline at end of file
diff --git a/apps/webapp/app/v3/marqs/devPubSub.server.ts b/apps/webapp/app/v3/marqs/devPubSub.server.ts
index 7d3de2ff6..a9ecafe87 100644
--- a/apps/webapp/app/v3/marqs/devPubSub.server.ts
+++ b/apps/webapp/app/v3/marqs/devPubSub.server.ts
@@ -26,13 +26,6 @@ function initializeDevPubSub() {
enableAutoPipelining: true,
...(env.REDIS_TLS_DISABLED === "true" ? {} : { tls: {} }),
},
- schema: {
- CANCEL_ATTEMPT: z.object({
- version: z.literal("v1").default("v1"),
- backgroundWorkerId: z.string(),
- attemptId: z.string(),
- taskRunId: z.string(),
- }),
- },
+ schema: messageCatalog,
});
}
diff --git a/apps/webapp/app/v3/otlpExporter.server.ts b/apps/webapp/app/v3/otlpExporter.server.ts
index 0696995c3..b09fbc403 100644
--- a/apps/webapp/app/v3/otlpExporter.server.ts
+++ b/apps/webapp/app/v3/otlpExporter.server.ts
@@ -35,9 +35,9 @@ class OTLPExporter {
constructor(
private readonly _eventRepository: EventRepository,
private readonly _verbose: boolean
- ) {}
+ ) { }
- async exportTraces(request: ExportTraceServiceRequest): Promise
{
+ async exportTraces(request: ExportTraceServiceRequest, immediate: boolean = false): Promise {
this.#logExportTracesVerbose(request);
const events = this.#filterResourceSpans(request.resourceSpans).flatMap((resourceSpan) => {
@@ -46,12 +46,16 @@ class OTLPExporter {
this.#logEventsVerbose(events);
- this._eventRepository.insertMany(events);
+ if (immediate) {
+ await this._eventRepository.insertManyImmediate(events);
+ } else {
+ await this._eventRepository.insertMany(events);
+ }
return ExportTraceServiceResponse.create();
}
- async exportLogs(request: ExportLogsServiceRequest): Promise {
+ async exportLogs(request: ExportLogsServiceRequest, immediate: boolean = false): Promise {
this.#logExportLogsVerbose(request);
const events = this.#filterResourceLogs(request.resourceLogs).flatMap((resourceLog) => {
@@ -60,7 +64,11 @@ class OTLPExporter {
this.#logEventsVerbose(events);
- this._eventRepository.insertMany(events);
+ if (immediate) {
+ await this._eventRepository.insertManyImmediate(events);
+ } else {
+ await this._eventRepository.insertMany(events);
+ }
return ExportLogsServiceResponse.create();
}
@@ -109,7 +117,7 @@ class OTLPExporter {
if (!triggerAttribute) return false;
- return isBoolValue(triggerAttribute.value) ? triggerAttribute.value.value.boolValue : false;
+ return isBoolValue(triggerAttribute.value) ? triggerAttribute.value.boolValue : false;
});
}
@@ -123,7 +131,7 @@ class OTLPExporter {
if (!attribute) return false;
- return isBoolValue(attribute.value) ? attribute.value.value.boolValue : false;
+ return isBoolValue(attribute.value) ? attribute.value.boolValue : false;
});
}
}
@@ -141,13 +149,13 @@ function convertLogsToCreateableEvents(resourceLog: ResourceLogs): Array;
+
+export const projectPubSub = singleton("projectPubSub", initializeProjectPubSub);
+
+function initializeProjectPubSub() {
+ return new ZodPubSub({
+ redis: {
+ port: env.REDIS_PORT,
+ host: env.REDIS_HOST,
+ username: env.REDIS_USERNAME,
+ password: env.REDIS_PASSWORD,
+ enableAutoPipelining: true,
+ ...(env.REDIS_TLS_DISABLED === "true" ? {} : { tls: {} }),
+ },
+ schema: messageCatalog,
+ });
+}
diff --git a/apps/webapp/server.ts b/apps/webapp/server.ts
index 02e86aec2..a09316f90 100644
--- a/apps/webapp/server.ts
+++ b/apps/webapp/server.ts
@@ -71,6 +71,7 @@ if (process.env.HTTP_SERVER_DISABLED !== "true") {
if (process.env.DASHBOARD_AND_API_DISABLED !== "true") {
app.all(
"*",
+ // @ts-ignore
createRequestHandler({
build,
mode: MODE,
diff --git a/packages/cli-v3/src/commands/deploy.ts b/packages/cli-v3/src/commands/deploy.ts
index c2fb3b8da..6ef01a7f2 100644
--- a/packages/cli-v3/src/commands/deploy.ts
+++ b/packages/cli-v3/src/commands/deploy.ts
@@ -1,7 +1,7 @@
import { intro, log, outro, spinner } from "@clack/prompts";
import { depot } from "@depot/cli";
import { context, trace } from "@opentelemetry/api";
-import { ResolvedConfig, flattenAttributes, recordSpanException } from "@trigger.dev/core/v3";
+import { ResolvedConfig, detectDependencyVersion, flattenAttributes, recordSpanException } from "@trigger.dev/core/v3";
import chalk from "chalk";
import { Command, Option as CommandOption } from "commander";
import { Metafile, build } from "esbuild";
@@ -1225,7 +1225,7 @@ function gatherRequiredDependencies(
const internalDependencyVersion = (packageJson.dependencies as Record)[
packageName
- ];
+ ] ?? detectDependencyVersion(packageName);
if (internalDependencyVersion) {
dependencies[packageName] = internalDependencyVersion;
diff --git a/packages/cli-v3/src/commands/dev.tsx b/packages/cli-v3/src/commands/dev.tsx
index 6cda6cd37..c11a077c5 100644
--- a/packages/cli-v3/src/commands/dev.tsx
+++ b/packages/cli-v3/src/commands/dev.tsx
@@ -5,6 +5,7 @@ import {
ZodMessageHandler,
ZodMessageSender,
clientWebsocketMessages,
+ detectDependencyVersion,
serverWebsocketMessages,
} from "@trigger.dev/core/v3";
import chalk from "chalk";
@@ -33,7 +34,6 @@ import { isLoggedIn } from "../utilities/session.js";
import { createTaskFileImports, gatherTaskFiles } from "../utilities/taskFiles";
import { UncaughtExceptionError } from "../workers/common/errors";
import { BackgroundWorker, BackgroundWorkerCoordinator } from "../workers/dev/backgroundWorker.js";
-import { fromZodError } from "zod-validation-error";
let apiClient: CliApiClient | undefined;
@@ -692,9 +692,9 @@ function gatherRequiredDependencies(outputMeta: Metafile["outputs"][string]) {
continue;
}
- const internalDependencyVersion = (packageJson.dependencies as Record)[
- packageName
- ];
+ const internalDependencyVersion =
+ (packageJson.dependencies as Record)[packageName] ??
+ detectDependencyVersion(packageName);
if (internalDependencyVersion) {
dependencies[packageName] = internalDependencyVersion;
diff --git a/packages/cli-v3/src/utilities/installPackages.ts b/packages/cli-v3/src/utilities/installPackages.ts
index da4b48cfc..fdb678fe4 100644
--- a/packages/cli-v3/src/utilities/installPackages.ts
+++ b/packages/cli-v3/src/utilities/installPackages.ts
@@ -12,6 +12,8 @@ export async function installPackages(
) {
const cwd = options?.cwd ?? process.cwd();
+ logger.debug(`Installing packages at ${cwd}:`, { packages });
+
// Make sure the cwd has a package.json file (if not create a barebones one)
try {
await readJSONFile(join(cwd, "package.json"));
@@ -49,7 +51,7 @@ export async function installPackages(
return;
}
- logger.debug(`Installing packages at ${cwd}:`);
+ logger.debug(`Found installable packages`);
logger.table(
Object.entries(installablePackages).map(([name, version]) => ({ name, version })),
"debug"
diff --git a/packages/cli-v3/src/workers/dev/backgroundWorker.ts b/packages/cli-v3/src/workers/dev/backgroundWorker.ts
index 339d83fce..391b162b8 100644
--- a/packages/cli-v3/src/workers/dev/backgroundWorker.ts
+++ b/packages/cli-v3/src/workers/dev/backgroundWorker.ts
@@ -167,8 +167,8 @@ export class BackgroundWorkerCoordinator {
!completion.ok && completion.skippedRetrying
? " (retrying skipped)"
: !completion.ok && completion.retry !== undefined
- ? ` (retrying in ${completion.retry.delay}ms)`
- : "";
+ ? ` (retrying in ${completion.retry.delay}ms)`
+ : "";
const resultText = !completion.ok
? completion.error.type === "INTERNAL_ERROR" &&
@@ -181,8 +181,8 @@ export class BackgroundWorkerCoordinator {
const errorText = !completion.ok
? this.#formatErrorLog(completion.error)
: "retry" in completion
- ? `retry in ${completion.retry}ms`
- : "";
+ ? `retry in ${completion.retry}ms`
+ : "";
const elapsedText = chalk.dim(`(${elapsed.toFixed(2)}ms)`);
@@ -263,7 +263,7 @@ export class BackgroundWorker {
constructor(
public path: string,
private params: BackgroundWorkerParams
- ) {}
+ ) { }
close() {
if (this._closed) {
@@ -302,8 +302,8 @@ export class BackgroundWorker {
const child = fork(this.path, {
stdio: [/*stdin*/ "ignore", /*stdout*/ "pipe", /*stderr*/ "pipe", "ipc"],
env: {
- ...this.#readEnvVars(),
...this.params.env,
+ ...this.#readEnvVars(),
},
});
@@ -475,13 +475,19 @@ export class BackgroundWorker {
}
#readEnvVars() {
- const result = {};
+ const result: { [key: string]: string } = {};
dotenv.config({
processEnv: result,
path: [".env", ".env.local", ".env.development.local"].map((p) => resolve(process.cwd(), p)),
});
+ process.env.TRIGGER_API_URL && (result.TRIGGER_API_URL = process.env.TRIGGER_API_URL);
+
+ // remove TRIGGER_API_URL and TRIGGER_SECRET_KEY, since those should be coming from the worker
+ delete result.TRIGGER_API_URL;
+ delete result.TRIGGER_SECRET_KEY;
+
return result;
}
@@ -536,6 +542,11 @@ class TaskRunProcess {
}
async initialize() {
+ logger.debug("initializing task run process", {
+ env: this.env,
+ path: this.path,
+ })
+
this._child = fork(this.path, {
stdio: [/*stdin*/ "ignore", /*stdout*/ "pipe", /*stderr*/ "pipe", "ipc"],
cwd: dirname(this.path),
@@ -544,6 +555,7 @@ class TaskRunProcess {
OTEL_RESOURCE_ATTRIBUTES: JSON.stringify({
[SemanticInternalAttributes.PROJECT_DIR]: this.worker.projectConfig.projectDir,
}),
+ OTEL_EXPORTER_OTLP_COMPRESSION: "none",
...(this.worker.debugOtel ? { OTEL_LOG_LEVEL: "debug" } : {}),
},
execArgv: this.worker.debuggerOn
@@ -688,8 +700,7 @@ class TaskRunProcess {
}
logger.log(
- `[${this.metadata.version}][${this._currentExecution.run.id}.${
- this._currentExecution.attempt.number
+ `[${this.metadata.version}][${this._currentExecution.run.id}.${this._currentExecution.attempt.number
}] ${data.toString()}`
);
}
@@ -706,8 +717,7 @@ class TaskRunProcess {
}
logger.error(
- `[${this.metadata.version}][${this._currentExecution.run.id}.${
- this._currentExecution.attempt.number
+ `[${this.metadata.version}][${this._currentExecution.run.id}.${this._currentExecution.attempt.number
}] ${data.toString()}`
);
}
diff --git a/packages/cli-v3/src/workers/dev/worker-facade.ts b/packages/cli-v3/src/workers/dev/worker-facade.ts
index da7943bac..9498ac17a 100644
--- a/packages/cli-v3/src/workers/dev/worker-facade.ts
+++ b/packages/cli-v3/src/workers/dev/worker-facade.ts
@@ -1,4 +1,4 @@
-import { Config, ProjectConfig, TaskExecutor, type TracingSDK } from "@trigger.dev/core/v3";
+import { Config, ProjectConfig, TaskExecutor, preciseDateOriginNow, type TracingSDK } from "@trigger.dev/core/v3";
import "source-map-support/register.js";
__WORKER_SETUP__;
@@ -35,8 +35,10 @@ import { TaskMetadataWithFunctions } from "../../types.js";
declare const sender: ZodMessageSender;
+const preciseDateOrigin = preciseDateOriginNow();
+
const tracer = new TriggerTracer({ tracer: otelTracer, logger: otelLogger });
-const consoleInterceptor = new ConsoleInterceptor(otelLogger);
+const consoleInterceptor = new ConsoleInterceptor(otelLogger, preciseDateOrigin);
const devRuntimeManager = new DevRuntimeManager();
@@ -46,6 +48,7 @@ const otelTaskLogger = new OtelTaskLogger({
logger: otelLogger,
tracer: tracer,
level: "info",
+ preciseDateOrigin
});
logger.setGlobalTaskLogger(otelTaskLogger);
diff --git a/packages/cli-v3/src/workers/prod/worker-facade.ts b/packages/cli-v3/src/workers/prod/worker-facade.ts
index 0e4ed89f6..ad492e031 100644
--- a/packages/cli-v3/src/workers/prod/worker-facade.ts
+++ b/packages/cli-v3/src/workers/prod/worker-facade.ts
@@ -6,6 +6,7 @@ import {
TaskExecutor,
ZodIpcConnection,
type TracingSDK,
+ preciseDateOriginNow,
} from "@trigger.dev/core/v3";
import "source-map-support/register.js";
@@ -37,13 +38,16 @@ import * as packageJson from "../../../package.json";
import { TaskMetadataWithFunctions } from "../../types";
+const preciseDateOrigin = preciseDateOriginNow();
+
const tracer = new TriggerTracer({ tracer: otelTracer, logger: otelLogger });
-const consoleInterceptor = new ConsoleInterceptor(otelLogger);
+const consoleInterceptor = new ConsoleInterceptor(otelLogger, preciseDateOrigin);
const otelTaskLogger = new OtelTaskLogger({
logger: otelLogger,
tracer: tracer,
level: "info",
+ preciseDateOrigin
});
logger.setGlobalTaskLogger(otelTaskLogger);
diff --git a/packages/core/package.json b/packages/core/package.json
index cb0f7a05e..17d2ea737 100644
--- a/packages/core/package.json
+++ b/packages/core/package.json
@@ -58,6 +58,7 @@
"test": "jest"
},
"dependencies": {
+ "@google-cloud/precise-date": "^4.0.0",
"@opentelemetry/api": "^1.7.0",
"@opentelemetry/api-logs": "^0.48.0",
"@opentelemetry/auto-instrumentations-node": "^0.40.3",
diff --git a/packages/core/src/v3/consoleInterceptor.ts b/packages/core/src/v3/consoleInterceptor.ts
index 0ee04a5a6..40c9c5e20 100644
--- a/packages/core/src/v3/consoleInterceptor.ts
+++ b/packages/core/src/v3/consoleInterceptor.ts
@@ -1,12 +1,14 @@
import type * as logsAPI from "@opentelemetry/api-logs";
import { SeverityNumber } from "@opentelemetry/api-logs";
import util from "node:util";
-import { flattenAttributes } from "./utils/flattenAttributes";
-import { SemanticInternalAttributes } from "./semanticInternalAttributes";
import { iconStringForSeverity } from "./icons";
+import { SemanticInternalAttributes } from "./semanticInternalAttributes";
+import { flattenAttributes } from "./utils/flattenAttributes";
+import { type PreciseDateOrigin, calculatePreciseDateHrTime } from "./utils/preciseDate";
+
export class ConsoleInterceptor {
- constructor(private readonly logger: logsAPI.Logger) {}
+ constructor(private readonly logger: logsAPI.Logger, private readonly preciseDateOrigin: PreciseDateOrigin) { }
// Intercept the console and send logs to the OpenTelemetry logger
// during the execution of the callback
@@ -54,6 +56,7 @@ export class ConsoleInterceptor {
#handleLog(severityNumber: SeverityNumber, severityText: string, ...args: unknown[]): void {
const body = util.format(...args);
+ const timestamp = this.#getTimestampInHrTime();
const parsed = tryParseJSON(body);
@@ -63,6 +66,7 @@ export class ConsoleInterceptor {
severityText,
body: getLogMessage(parsed.value, severityText),
attributes: { ...this.#getAttributes(severityNumber), ...flattenAttributes(parsed.value) },
+ timestamp,
});
return;
@@ -73,9 +77,14 @@ export class ConsoleInterceptor {
severityText,
body,
attributes: this.#getAttributes(severityNumber),
+ timestamp,
});
}
+ #getTimestampInHrTime(): [number, number] {
+ return calculatePreciseDateHrTime(this.preciseDateOrigin);
+ }
+
#getAttributes(severityNumber: SeverityNumber): logsAPI.LogAttributes {
const icon = iconStringForSeverity(severityNumber);
let result: logsAPI.LogAttributes = {};
diff --git a/packages/core/src/v3/index.ts b/packages/core/src/v3/index.ts
index 3f6c31abc..1ee6cda6c 100644
--- a/packages/core/src/v3/index.ts
+++ b/packages/core/src/v3/index.ts
@@ -50,3 +50,5 @@ export { eventFilterMatches } from "../eventFilterMatches";
export { omit } from "./utils/omit";
export { TracingSDK, type TracingDiagnosticLogLevel, recordSpanException } from "./otel";
export { TaskExecutor, type TaskExecutorOptions } from "./workers/taskExecutor";
+export { detectDependencyVersion } from "./utils/detectDependencyVersion";
+export { type PreciseDateOrigin, calculatePreciseDateHrTime, preciseDateOriginNow } from "./utils/preciseDate";
diff --git a/packages/core/src/v3/logger/taskLogger.ts b/packages/core/src/v3/logger/taskLogger.ts
index 471b84230..0cd36582e 100644
--- a/packages/core/src/v3/logger/taskLogger.ts
+++ b/packages/core/src/v3/logger/taskLogger.ts
@@ -1,9 +1,10 @@
-import { Logger, SeverityNumber } from "@opentelemetry/api-logs";
-import { flattenAttributes } from "../utils/flattenAttributes";
import { Attributes, Span, SpanOptions } from "@opentelemetry/api";
+import { Logger, SeverityNumber } from "@opentelemetry/api-logs";
import { iconStringForSeverity } from "../icons";
import { SemanticInternalAttributes } from "../semanticInternalAttributes";
import { TriggerTracer } from "../tracer";
+import { flattenAttributes } from "../utils/flattenAttributes";
+import { PreciseDateOrigin, calculatePreciseDateHrTime } from "../utils/preciseDate";
export type LogLevel = "log" | "error" | "warn" | "info" | "debug";
@@ -13,6 +14,7 @@ export type TaskLoggerConfig = {
logger: Logger;
tracer: TriggerTracer;
level: LogLevel;
+ preciseDateOrigin: PreciseDateOrigin;
};
export interface TaskLogger {
@@ -67,6 +69,8 @@ export class OtelTaskLogger implements TaskLogger {
severityNumber: SeverityNumber,
properties?: Record
) {
+ const timestamp = this.#getTimestampInHrTime();
+
let attributes: Attributes = { ...flattenAttributes(properties) };
const icon = iconStringForSeverity(severityNumber);
@@ -79,20 +83,25 @@ export class OtelTaskLogger implements TaskLogger {
severityText,
body: message,
attributes,
+ timestamp
});
}
trace(name: string, fn: (span: Span) => Promise, options?: SpanOptions): Promise {
return this._config.tracer.startActiveSpan(name, fn, options);
}
+
+ #getTimestampInHrTime(): [number, number] {
+ return calculatePreciseDateHrTime(this._config.preciseDateOrigin);
+ }
}
export class NoopTaskLogger implements TaskLogger {
- debug() {}
- log() {}
- info() {}
- warn() {}
- error() {}
+ debug() { }
+ log() { }
+ info() { }
+ warn() { }
+ error() { }
trace(name: string, fn: (span: Span) => Promise): Promise {
return fn({} as Span);
}
diff --git a/packages/core/src/v3/utils/detectDependencyVersion.ts b/packages/core/src/v3/utils/detectDependencyVersion.ts
new file mode 100644
index 000000000..fa67fb21a
--- /dev/null
+++ b/packages/core/src/v3/utils/detectDependencyVersion.ts
@@ -0,0 +1,5 @@
+import { dependencies } from "../../../package.json"
+
+export function detectDependencyVersion(dependency: string): string | undefined {
+ return (dependencies as Record)[dependency]
+}
\ No newline at end of file
diff --git a/packages/core/src/v3/utils/preciseDate.ts b/packages/core/src/v3/utils/preciseDate.ts
new file mode 100644
index 000000000..3efc92ecf
--- /dev/null
+++ b/packages/core/src/v3/utils/preciseDate.ts
@@ -0,0 +1,23 @@
+import { PreciseDate } from "@google-cloud/precise-date";
+
+export type PreciseDateOrigin = {
+ hrtime: [number, number];
+ timestamp: PreciseDate
+}
+
+export function preciseDateOriginNow(): PreciseDateOrigin {
+ return {
+ hrtime: process.hrtime(),
+ timestamp: new PreciseDate()
+ }
+}
+
+export function calculatePreciseDateHrTime(origin: PreciseDateOrigin): [number, number] {
+ const elapsedHrTime = process.hrtime(origin.hrtime);
+ const elapsedNanoseconds = BigInt(elapsedHrTime[0]) * BigInt(1e9) + BigInt(elapsedHrTime[1]);
+
+ const preciseDate = new PreciseDate(origin.timestamp.getFullTime() + elapsedNanoseconds)
+ const dateStruct = preciseDate.toStruct();
+
+ return [dateStruct.seconds, dateStruct.nanos];
+}
\ No newline at end of file
diff --git a/packages/core/tsconfig.json b/packages/core/tsconfig.json
index aa110b9ab..d6d3f05e0 100644
--- a/packages/core/tsconfig.json
+++ b/packages/core/tsconfig.json
@@ -6,6 +6,7 @@
"emitDecoratorMetadata": true,
"declaration": false,
"declarationMap": false,
+ "resolveJsonModule": true,
"types": ["jest"],
"lib": ["DOM", "DOM.Iterable"],
"paths": {
diff --git a/packages/database/prisma/migrations/20240319120645_convert_start_time_to_nanoseconds_since_epoch/migration.sql b/packages/database/prisma/migrations/20240319120645_convert_start_time_to_nanoseconds_since_epoch/migration.sql
new file mode 100644
index 000000000..89c58423f
--- /dev/null
+++ b/packages/database/prisma/migrations/20240319120645_convert_start_time_to_nanoseconds_since_epoch/migration.sql
@@ -0,0 +1,34 @@
+/*
+ Warnings:
+
+ - Changed the type of `startTime` on the `TaskEvent` table. No cast exists, the column would be dropped and recreated, which cannot be done if there is data, since the column is required.
+
+ */
+-- AlterTable
+BEGIN;
+
+-- Step 1: Add a new column
+ALTER TABLE
+ "TaskEvent"
+ADD
+ COLUMN "temporary_startTime" BIGINT;
+
+-- Step 2: Convert the data in the current "startTime" column to nanoseconds
+UPDATE
+ "TaskEvent"
+SET
+ "temporary_startTime" = EXTRACT(
+ EPOCH
+ FROM
+ "startTime"
+ ) * 1000000000;
+
+-- Step 3: Drop the original column
+ALTER TABLE
+ "TaskEvent" DROP COLUMN "startTime";
+
+-- Step 4: Rename the new column
+ALTER TABLE
+ "TaskEvent" RENAME COLUMN "temporary_startTime" TO "startTime";
+
+COMMIT;
\ No newline at end of file
diff --git a/packages/database/prisma/migrations/20240319121124_what_migration_is_this/migration.sql b/packages/database/prisma/migrations/20240319121124_what_migration_is_this/migration.sql
new file mode 100644
index 000000000..e17ff572b
--- /dev/null
+++ b/packages/database/prisma/migrations/20240319121124_what_migration_is_this/migration.sql
@@ -0,0 +1,8 @@
+/*
+ Warnings:
+
+ - Made the column `startTime` on table `TaskEvent` required. This step will fail if there are existing NULL values in that column.
+
+*/
+-- AlterTable
+ALTER TABLE "TaskEvent" ALTER COLUMN "startTime" SET NOT NULL;
diff --git a/packages/database/prisma/schema.prisma b/packages/database/prisma/schema.prisma
index 982e05e86..27828cafc 100644
--- a/packages/database/prisma/schema.prisma
+++ b/packages/database/prisma/schema.prisma
@@ -1762,7 +1762,8 @@ model TaskEvent {
links Json?
events Json?
- startTime DateTime
+ /// This is the time the event started in nanoseconds since the epoch
+ startTime BigInt
/// This is the duration of the event in nanoseconds
duration BigInt @default(0)
diff --git a/packages/otlp-importer/scripts/generate-protos.mjs b/packages/otlp-importer/scripts/generate-protos.mjs
index 6689fc3c8..17ad437d6 100644
--- a/packages/otlp-importer/scripts/generate-protos.mjs
+++ b/packages/otlp-importer/scripts/generate-protos.mjs
@@ -46,7 +46,6 @@ for (const proto of protos) {
`--ts_proto_opt=env=node ` +
`--ts_proto_opt=removeEnumPrefix=true ` +
`--ts_proto_opt=lowerCaseServiceMethods=true ` +
- `--ts_proto_opt=oneof=unions ` +
`--experimental_allow_proto3_optional ` +
`"${path.join(protosPath, proto)}"`;
try {
diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml
index 06c94158a..d114e4e94 100644
--- a/pnpm-lock.yaml
+++ b/pnpm-lock.yaml
@@ -1198,6 +1198,7 @@ importers:
packages/core:
specifiers:
+ '@google-cloud/precise-date': ^4.0.0
'@opentelemetry/api': ^1.7.0
'@opentelemetry/api-logs': ^0.48.0
'@opentelemetry/auto-instrumentations-node': ^0.40.3
@@ -1230,6 +1231,7 @@ importers:
zod: 3.22.3
zod-error: 1.5.0
dependencies:
+ '@google-cloud/precise-date': 4.0.0
'@opentelemetry/api': 1.7.0
'@opentelemetry/api-logs': 0.48.0
'@opentelemetry/auto-instrumentations-node': 0.40.3_@opentelemetry+api@1.7.0
@@ -7599,6 +7601,11 @@ packages:
resolution: {integrity: sha512-k2Ty1JcVojjJFwrg/ThKi2ujJ7XNLYaFGNB/bWT9wGR+oSMJHMa5w+CUq6p/pVrKeNNgA7pCqEcjSnHVoqJQFw==}
dev: true
+ /@google-cloud/precise-date/4.0.0:
+ resolution: {integrity: sha512-1TUx3KdaU3cN7nfCdNf+UVqA/PSX29Cjcox3fZZBtINlRrXVTmUkQnCKv2MbBUbCopbK4olAT1IHl76uZyCiVA==}
+ engines: {node: '>=14.0.0'}
+ dev: false
+
/@graphile/logger/0.2.0:
resolution: {integrity: sha512-jjcWBokl9eb1gVJ85QmoaQ73CQ52xAaOCF29ukRbYNl6lY+ts0ErTaDYOBlejcbUs2OpaiqYLO5uDhyLFzWw4w==}
dev: false
@@ -37081,7 +37088,7 @@ packages:
dependencies:
bs-logger: 0.2.6
fast-json-stable-stringify: 2.1.0
- jest: 29.6.2_@types+node@18.15.13
+ jest: 29.6.2_@types+node@18.17.1
jest-util: 29.6.2
json5: 2.2.3
lodash.memoize: 4.1.2
diff --git a/references/v3-catalog/src/trigger/simple.ts b/references/v3-catalog/src/trigger/simple.ts
index 9b31bca0e..1acae5c81 100644
--- a/references/v3-catalog/src/trigger/simple.ts
+++ b/references/v3-catalog/src/trigger/simple.ts
@@ -72,9 +72,13 @@ export const parentTask = task({
logger.info("Parent task payload", { payload });
console.info("This is an info message");
+ logger.info("This is an info message from logger.info");
console.log(JSON.stringify({ ctx, message: "This is the parent task context" }));
+ logger.log(JSON.stringify({ ctx, message: "This is the parent task context from logger.log" }));
console.warn("You've been warned buddy");
+ logger.warn("You've been warned buddy from logger.warn");
console.error("This is an error message");
+ logger.error("This is an error message from logger.error");
await wait.for({ seconds: 5 });