From 6afda40c356d19866f32f734a7c63c0779b4fc6d Mon Sep 17 00:00:00 2001 From: Eric Allam Date: Tue, 19 Mar 2024 09:24:33 +0000 Subject: [PATCH 1/9] Add depot setup action --- .github/workflows/publish-docker.yml | 3 +++ 1 file changed, 3 insertions(+) diff --git a/.github/workflows/publish-docker.yml b/.github/workflows/publish-docker.yml index ea6c149cd..fbb3c2c2f 100644 --- a/.github/workflows/publish-docker.yml +++ b/.github/workflows/publish-docker.yml @@ -8,6 +8,9 @@ jobs: version: ${{ steps.get_version.outputs.version }} short_sha: ${{ steps.get_commit.outputs.sha_short }} steps: + - name: Setup Depot CLI + uses: depot/setup-action@v1 + - name: ⬇️ Checkout repo uses: actions/checkout@v3 with: From 86c24afb44c003153df8c1a0b25df401b5200238 Mon Sep 17 00:00:00 2001 From: Eric Allam Date: Tue, 19 Mar 2024 10:59:50 +0000 Subject: [PATCH 2/9] Support json traces and logs in the otel endpoints (#954) --- apps/webapp/app/routes/otel.v1.logs.ts | 22 ++++-- apps/webapp/app/routes/otel.v1.traces.ts | 22 ++++-- apps/webapp/app/v3/otlpExporter.server.ts | 75 ++++++++++--------- .../src/workers/dev/backgroundWorker.ts | 22 +++--- .../otlp-importer/scripts/generate-protos.mjs | 1 - 5 files changed, 85 insertions(+), 57 deletions(-) diff --git a/apps/webapp/app/routes/otel.v1.logs.ts b/apps/webapp/app/routes/otel.v1.logs.ts index 0554a2942..88cf28589 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); - 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..9dfbfb64a 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); - 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/v3/otlpExporter.server.ts b/apps/webapp/app/v3/otlpExporter.server.ts index 0696995c3..e233c6f80 100644 --- a/apps/webapp/app/v3/otlpExporter.server.ts +++ b/apps/webapp/app/v3/otlpExporter.server.ts @@ -35,7 +35,7 @@ class OTLPExporter { constructor( private readonly _eventRepository: EventRepository, private readonly _verbose: boolean - ) {} + ) { } async exportTraces(request: ExportTraceServiceRequest): Promise { this.#logExportTracesVerbose(request); @@ -109,7 +109,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 +123,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,7 +141,7 @@ function convertLogsToCreateableEvents(resourceLog: ResourceLogs): Array Date: Tue, 19 Mar 2024 11:38:48 +0000 Subject: [PATCH 3/9] Trying beefier buildjet machines for CI --- .github/workflows/e2e.yml | 2 +- .github/workflows/publish-dev.yml | 2 +- .github/workflows/release.yml | 2 +- .github/workflows/typecheck.yml | 2 +- .github/workflows/unit-tests.yml | 2 +- 5 files changed, 5 insertions(+), 5 deletions(-) diff --git a/.github/workflows/e2e.yml b/.github/workflows/e2e.yml index 0b8e39782..0e513a814 100644 --- a/.github/workflows/e2e.yml +++ b/.github/workflows/e2e.yml @@ -4,7 +4,7 @@ on: jobs: e2e: name: "🧪 E2E Tests" - runs-on: buildjet-4vcpu-ubuntu-2204 + runs-on: buildjet-16vcpu-ubuntu-2204 steps: - name: 🐳 Login to Docker Hub uses: docker/login-action@v2 diff --git a/.github/workflows/publish-dev.yml b/.github/workflows/publish-dev.yml index 58e4d5cb4..34afea272 100644 --- a/.github/workflows/publish-dev.yml +++ b/.github/workflows/publish-dev.yml @@ -38,7 +38,7 @@ jobs: strategy: matrix: package: [coordinator, kubernetes-provider] - runs-on: buildjet-4vcpu-ubuntu-2204 + runs-on: buildjet-16vcpu-ubuntu-2204 env: DOCKER_BUILDKIT: "1" steps: diff --git a/.github/workflows/release.yml b/.github/workflows/release.yml index e90d3751a..96e67a21a 100644 --- a/.github/workflows/release.yml +++ b/.github/workflows/release.yml @@ -12,7 +12,7 @@ on: jobs: release: name: 🦋 Changesets Release - runs-on: buildjet-4vcpu-ubuntu-2204 + runs-on: buildjet-16vcpu-ubuntu-2204 if: | github.repository == 'triggerdotdev/trigger.dev' outputs: diff --git a/.github/workflows/typecheck.yml b/.github/workflows/typecheck.yml index 468222d45..73add70d9 100644 --- a/.github/workflows/typecheck.yml +++ b/.github/workflows/typecheck.yml @@ -3,7 +3,7 @@ on: workflow_call: jobs: typecheck: - runs-on: buildjet-4vcpu-ubuntu-2204 + runs-on: buildjet-16vcpu-ubuntu-2204 steps: - name: ⬇️ Checkout repo diff --git a/.github/workflows/unit-tests.yml b/.github/workflows/unit-tests.yml index 893dea5ce..a0070e1ea 100644 --- a/.github/workflows/unit-tests.yml +++ b/.github/workflows/unit-tests.yml @@ -4,7 +4,7 @@ on: jobs: unitTests: name: "🧪 Unit Tests" - runs-on: buildjet-4vcpu-ubuntu-2204 + runs-on: buildjet-16vcpu-ubuntu-2204 steps: - name: ⬇️ Checkout repo uses: actions/checkout@v3 From c5aaabdb4ec46c19cbcf5717be4a4d2999da2f44 Mon Sep 17 00:00:00 2001 From: Matt Aitken Date: Tue, 19 Mar 2024 12:13:48 +0000 Subject: [PATCH 4/9] The run timeline updates on the client if there's no new data (#955) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * Use “@v3” instead of “@latest” for the npx commands for v3 * Split the Timeline into a component so it can live refresh without re-rendering everything * Timeline now live refreshes on the client every 500ms when run is executing * The timeline bars now animate their position/width when it changes * Export some more types from TreeView --- apps/webapp/app/components/SetupCommands.tsx | 13 +- .../primitives/TreeView/TreeView.tsx | 5 +- .../app/presenters/v3/RunPresenter.server.ts | 24 +- .../route.tsx | 463 +++++++++++------- 4 files changed, 301 insertions(+), 204 deletions(-) diff --git a/apps/webapp/app/components/SetupCommands.tsx b/apps/webapp/app/components/SetupCommands.tsx index fe73742b8..8d94aa629 100644 --- a/apps/webapp/app/components/SetupCommands.tsx +++ b/apps/webapp/app/components/SetupCommands.tsx @@ -131,6 +131,7 @@ export function TriggerDevStep({ extra }: { extra?: string }) { } // Trigger.dev version 3 setup commands +const v3PackageTag = "v3"; export function InitCommandV3() { const project = useProject(); @@ -147,7 +148,7 @@ export function InitCommandV3() { variant="primary/medium" iconButton className="mb-4" - value={`npx trigger.dev@latest init -p ${projectRef}`} + value={`npx trigger.dev@${v3PackageTag} init -p ${projectRef}`} /> @@ -155,7 +156,7 @@ export function InitCommandV3() { variant="primary/medium" iconButton className="mb-4" - value={`pnpm dlx trigger.dev@latest init -p ${projectRef}`} + value={`pnpm dlx trigger.dev@${v3PackageTag} init -p ${projectRef}`} /> @@ -163,7 +164,7 @@ export function InitCommandV3() { variant="primary/medium" iconButton className="mb-4" - value={`yarn dlx trigger.dev@latest init -p ${projectRef}`} + value={`yarn dlx trigger.dev@${v3PackageTag} init -p ${projectRef}`} /> @@ -183,7 +184,7 @@ export function TriggerDevStepV3() { variant="primary/medium" iconButton className="mb-4" - value={`npx trigger.dev@latest dev`} + value={`npx trigger.dev@${v3PackageTag} dev`} /> @@ -191,7 +192,7 @@ export function TriggerDevStepV3() { variant="primary/medium" iconButton className="mb-4" - value={`pnpm dlx trigger.dev@latest dev`} + value={`pnpm dlx trigger.dev@${v3PackageTag} dev`} /> @@ -199,7 +200,7 @@ export function TriggerDevStepV3() { variant="primary/medium" iconButton className="mb-4" - value={`yarn dlx trigger.dev@latest dev`} + value={`yarn dlx trigger.dev@${v3PackageTag} dev`} /> diff --git a/apps/webapp/app/components/primitives/TreeView/TreeView.tsx b/apps/webapp/app/components/primitives/TreeView/TreeView.tsx index 69b437e13..b35917d54 100644 --- a/apps/webapp/app/components/primitives/TreeView/TreeView.tsx +++ b/apps/webapp/app/components/primitives/TreeView/TreeView.tsx @@ -23,6 +23,9 @@ export type TreeViewProps = { onScroll?: (scrollTop: number) => void; } & Pick; +export type GetTreePropsFn = UseTreeStateOutput["getTreeProps"]; +export type GetNodePropsFn = UseTreeStateOutput["getNodeProps"]; + export function TreeView({ tree, renderNode, @@ -144,7 +147,7 @@ type HTMLAttributes = Omit< "onAnimationStart" | "onDragStart" | "onDragEnd" | "onDrag" >; -type UseTreeStateOutput = { +export type UseTreeStateOutput = { selected: string | undefined; nodes: NodesState; virtualizer: Virtualizer; diff --git a/apps/webapp/app/presenters/v3/RunPresenter.server.ts b/apps/webapp/app/presenters/v3/RunPresenter.server.ts index c55ddc4b2..34d5c87b1 100644 --- a/apps/webapp/app/presenters/v3/RunPresenter.server.ts +++ b/apps/webapp/app/presenters/v3/RunPresenter.server.ts @@ -79,25 +79,22 @@ export class RunPresenter { n.data.startTime.getTime() - treeRootStartTimeMs ); totalDuration = Math.max(totalDuration, offset + n.data.duration); - return { ...n, data: { ...n.data, offset, isRoot: n.id === traceSummary.rootSpan.id } }; + return { + ...n, + data: { + ...n.data, + //set partial nodes to null duration + duration: n.data.isPartial ? null : n.data.duration, + offset, + isRoot: n.id === traceSummary.rootSpan.id, + }, + }; }) : []; - //if any elements are partial we want the total duration to represent all of the time until now - totalDuration = events.some((e) => e.data.isPartial) - ? millisecondsToNanoseconds(Date.now() - treeRootStartTimeMs) - : totalDuration; - //total duration should be a minimum of 1ms totalDuration = Math.max(totalDuration, millisecondsToNanoseconds(1)); - //we need to adjust any partial nodes so they run the full duration - for (const event of events) { - if (event.data.isPartial) { - event.data.duration = totalDuration - event.data.offset; - } - } - let rootSpanStatus: "executing" | "completed" | "failed" = "executing"; if (events[0]) { if (events[0].data.isError) { @@ -123,6 +120,7 @@ export class RunPresenter { parentRunFriendlyId: tree?.id === traceSummary.rootSpan.id ? undefined : traceSummary.rootSpan.runId, duration: totalDuration, + rootStartedAt: tree?.data.startTime, }; } } diff --git a/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.v3.$projectParam.runs.$runParam/route.tsx b/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.v3.$projectParam.runs.$runParam/route.tsx index 4ec9abbdb..b9323eed1 100644 --- a/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.v3.$projectParam.runs.$runParam/route.tsx +++ b/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.v3.$projectParam.runs.$runParam/route.tsx @@ -4,9 +4,16 @@ import { MagnifyingGlassMinusIcon, MagnifyingGlassPlusIcon, } from "@heroicons/react/20/solid"; +import { Time } from "@internationalized/date"; import { Link, Outlet, useNavigate, useParams, useRevalidator } from "@remix-run/react"; import { LoaderFunctionArgs } from "@remix-run/server-runtime"; -import { formatDurationMilliseconds, nanosecondsToMilliseconds } from "@trigger.dev/core/v3"; +import { Virtualizer } from "@tanstack/react-virtual"; +import { + formatDurationMilliseconds, + millisecondsToNanoseconds, + nanosecondsToMilliseconds, +} from "@trigger.dev/core/v3"; +import { motion } from "framer-motion"; import { useEffect, useRef, useState } from "react"; import { typedjson, useTypedLoaderData } from "remix-typedjson"; import { ShowParentIcon, ShowParentIconSelected } from "~/assets/icons/ShowParentIcon"; @@ -26,7 +33,15 @@ import { import { Slider } from "~/components/primitives/Slider"; import { Switch } from "~/components/primitives/Switch"; import * as Timeline from "~/components/primitives/Timeline"; -import { TreeView, useTree } from "~/components/primitives/TreeView/TreeView"; +import { + GetNodePropsFn, + GetTreePropsFn, + TreeView, + TreeViewProps, + UseTreeStateOutput, + useTree, +} from "~/components/primitives/TreeView/TreeView"; +import { NodesState } from "~/components/primitives/TreeView/reducer"; import { RunIcon } from "~/components/runs/v3/RunIcon"; import { SpanTitle, eventBackgroundClassName } from "~/components/runs/v3/SpanTitle"; import { TaskRunStatusIcon, runStatusClassNameColor } from "~/components/runs/v3/TaskRunStatus"; @@ -55,7 +70,7 @@ export const loader = async ({ request, params }: LoaderFunctionArgs) => { const { projectParam, organizationSlug, runParam } = v3RunParamsSchema.parse(params); const presenter = new RunPresenter(); - const { run, events, parentRunFriendlyId, duration, rootSpanStatus } = await presenter.call({ + const result = await presenter.call({ userId, organizationSlug, projectSlug: projectParam, @@ -66,12 +81,8 @@ export const loader = async ({ request, params }: LoaderFunctionArgs) => { const resizeSettings = await getResizableRunSettings(request); return typedjson({ - run, - events, - parentRunFriendlyId, + ...result, resizeSettings, - duration, - rootSpanStatus, }); }; @@ -82,8 +93,15 @@ function getSpanId(path: string): string | undefined { } export default function Page() { - const { run, events, parentRunFriendlyId, resizeSettings, duration, rootSpanStatus } = - useTypedLoaderData(); + const { + run, + events, + parentRunFriendlyId, + resizeSettings, + duration, + rootSpanStatus, + rootStartedAt, + } = useTypedLoaderData(); const navigate = useNavigate(); const organization = useOrganization(); const pathName = usePathName(); @@ -142,6 +160,7 @@ export default function Page() { }} totalDuration={duration} rootSpanStatus={rootSpanStatus} + rootStartedAt={rootStartedAt} /> ) : ( @@ -183,7 +203,15 @@ export default function Page() { ); } -const tickCount = 5; +type TasksTreeViewProps = { + events: RunEvent[]; + selectedId?: string; + parentRunFriendlyId?: string; + onSelectedIdChanged: (selectedId: string | undefined) => void; + totalDuration: number; + rootSpanStatus: "executing" | "completed" | "failed"; + rootStartedAt: Date | undefined; +}; function TasksTreeView({ events, @@ -192,14 +220,8 @@ function TasksTreeView({ onSelectedIdChanged, totalDuration, rootSpanStatus, -}: { - events: RunEvent[]; - selectedId?: string; - parentRunFriendlyId?: string; - onSelectedIdChanged: (selectedId: string | undefined) => void; - totalDuration: number; - rootSpanStatus: "executing" | "completed" | "failed"; -}) { + rootStartedAt, +}: TasksTreeViewProps) { const [filterText, setFilterText] = useState(""); const [errorsOnly, setErrorsOnly] = useState(false); const [showDurations, setShowDurations] = useState(false); @@ -207,8 +229,6 @@ function TasksTreeView({ const parentRef = useRef(null); const treeScrollRef = useRef(null); const timelineScrollRef = useRef(null); - const timelineContainerRef = useRef(null); - const initialTimelineDimensions = useInitialDimensions(timelineContainerRef); const { nodes, @@ -238,9 +258,6 @@ function TasksTreeView({ }, }); - const minTimelineWidth = initialTimelineDimensions?.width ?? 300; - const maxTimelineWidth = minTimelineWidth * 10; - return (
@@ -378,170 +395,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 +753,12 @@ function SpanWithDuration({ }: Timeline.SpanProps & { node: RunEvent; showDuration: boolean }) { return ( -
{node.data.isPartial && (
-
+
); } From 41e1ac9e0b4eb5b7296e5e1d011d386de184b151 Mon Sep 17 00:00:00 2001 From: Eric Allam Date: Tue, 19 Mar 2024 15:06:47 +0000 Subject: [PATCH 5/9] v3: nanosecond precision task events and immediate mode for dev events (#956) * Store startTime in nanoseconds since epoch, ensure logs using nanosecond precision for their timestamp (so they are not out of order) * export task events immediately in dev --- apps/webapp/app/routes/otel.v1.logs.ts | 2 +- apps/webapp/app/routes/otel.v1.traces.ts | 2 +- apps/webapp/app/v3/eventRepository.server.ts | 61 ++++++++++++------- apps/webapp/app/v3/otlpExporter.server.ts | 20 ++++-- packages/cli-v3/src/commands/deploy.ts | 4 +- packages/cli-v3/src/commands/dev.tsx | 8 +-- .../cli-v3/src/utilities/installPackages.ts | 4 +- .../cli-v3/src/workers/dev/worker-facade.ts | 7 ++- .../cli-v3/src/workers/prod/worker-facade.ts | 6 +- packages/core/package.json | 1 + packages/core/src/v3/consoleInterceptor.ts | 15 ++++- packages/core/src/v3/index.ts | 2 + packages/core/src/v3/logger/taskLogger.ts | 23 ++++--- .../src/v3/utils/detectDependencyVersion.ts | 5 ++ packages/core/src/v3/utils/preciseDate.ts | 23 +++++++ .../migration.sql | 34 +++++++++++ .../migration.sql | 8 +++ packages/database/prisma/schema.prisma | 3 +- pnpm-lock.yaml | 9 ++- references/v3-catalog/src/trigger/simple.ts | 4 ++ 20 files changed, 188 insertions(+), 53 deletions(-) create mode 100644 packages/core/src/v3/utils/detectDependencyVersion.ts create mode 100644 packages/core/src/v3/utils/preciseDate.ts create mode 100644 packages/database/prisma/migrations/20240319120645_convert_start_time_to_nanoseconds_since_epoch/migration.sql create mode 100644 packages/database/prisma/migrations/20240319121124_what_migration_is_this/migration.sql diff --git a/apps/webapp/app/routes/otel.v1.logs.ts b/apps/webapp/app/routes/otel.v1.logs.ts index 88cf28589..3bbafab9b 100644 --- a/apps/webapp/app/routes/otel.v1.logs.ts +++ b/apps/webapp/app/routes/otel.v1.logs.ts @@ -8,7 +8,7 @@ export async function action({ request }: ActionFunctionArgs) { if (contentType === "application/json") { const body = await request.json(); - const exportResponse = await otlpExporter.exportLogs(body as ExportLogsServiceRequest); + const exportResponse = await otlpExporter.exportLogs(body as ExportLogsServiceRequest, true); return json(exportResponse, { status: 200 }) } else if (contentType === "application/x-protobuf") { diff --git a/apps/webapp/app/routes/otel.v1.traces.ts b/apps/webapp/app/routes/otel.v1.traces.ts index 9dfbfb64a..6dfc8cb34 100644 --- a/apps/webapp/app/routes/otel.v1.traces.ts +++ b/apps/webapp/app/routes/otel.v1.traces.ts @@ -8,7 +8,7 @@ export async function action({ request }: ActionFunctionArgs) { if (contentType === "application/json") { const body = await request.json(); - const exportResponse = await otlpExporter.exportTraces(body as ExportTraceServiceRequest); + const exportResponse = await otlpExporter.exportTraces(body as ExportTraceServiceRequest, true); return json(exportResponse, { status: 200 }) } else if (contentType === "application/x-protobuf") { 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/otlpExporter.server.ts b/apps/webapp/app/v3/otlpExporter.server.ts index e233c6f80..b09fbc403 100644 --- a/apps/webapp/app/v3/otlpExporter.server.ts +++ b/apps/webapp/app/v3/otlpExporter.server.ts @@ -37,7 +37,7 @@ class OTLPExporter { 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(); } @@ -147,7 +155,7 @@ function convertLogsToCreateableEvents(resourceLog: ResourceLogs): Array)[ 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/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/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 8a59cdf9c..12e224f4e 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/pnpm-lock.yaml b/pnpm-lock.yaml index efa350912..2fafa4084 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -1196,6 +1196,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 @@ -1228,6 +1229,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 @@ -7601,6 +7603,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 @@ -37077,7 +37084,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 }); From 891278ab1f9de0a72151523595f62674d6e6a1c6 Mon Sep 17 00:00:00 2001 From: Eric Allam Date: Tue, 19 Mar 2024 15:24:01 +0000 Subject: [PATCH 6/9] Enable loading package.json inside of core --- packages/core/tsconfig.json | 1 + 1 file changed, 1 insertion(+) 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": { From dc5eb68a0b2d16fa7867b04e549b6ca70ddae17d Mon Sep 17 00:00:00 2001 From: Eric Allam Date: Tue, 19 Mar 2024 15:33:09 +0000 Subject: [PATCH 7/9] Fix typecheck errors because of resolveJsonModule --- apps/coordinator/tsconfig.json | 1 + apps/docker-provider/tsconfig.json | 1 + apps/webapp/server.ts | 1 + 3 files changed, 3 insertions(+) diff --git a/apps/coordinator/tsconfig.json b/apps/coordinator/tsconfig.json index 967c7281e..2e1257735 100644 --- a/apps/coordinator/tsconfig.json +++ b/apps/coordinator/tsconfig.json @@ -5,6 +5,7 @@ "target": "es2016", "module": "commonjs", "esModuleInterop": true, + "resolveJsonModule": true, "forceConsistentCasingInFileNames": true, "strict": true, "skipLibCheck": true, diff --git a/apps/docker-provider/tsconfig.json b/apps/docker-provider/tsconfig.json index 345326da8..8e6c54f74 100644 --- a/apps/docker-provider/tsconfig.json +++ b/apps/docker-provider/tsconfig.json @@ -4,6 +4,7 @@ "module": "commonjs", "esModuleInterop": true, "forceConsistentCasingInFileNames": true, + "resolveJsonModule": true, "strict": true, "skipLibCheck": true, "paths": { 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, From 2e56af974291ec77eb36cdef11e490b9cbf5a298 Mon Sep 17 00:00:00 2001 From: Matt Aitken Date: Tue, 19 Mar 2024 15:42:57 +0000 Subject: [PATCH 8/9] Tasks page live-reloads with new tasks (#957) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * DevPubSub: use the schema that’s already been defined * Publish a message to projectPubSub when a new DEV worker is created * Live reloading of the tasks table (for dev tasks) * Live notifications for deployed workers as well as local workers --- .../v3/TasksStreamPresenter.server.ts | 115 ++++++++++++++++++ .../route.tsx | 20 ++- .../route.tsx | 12 +- .../route.tsx | 13 ++ apps/webapp/app/utils/pathBuilder.ts | 4 + apps/webapp/app/v3/marqs/devPubSub.server.ts | 9 +- .../services/createBackgroundWorker.server.ts | 10 ++ .../createDeployedBackgroundWorker.server.ts | 14 +++ .../app/v3/services/projectPubSub.server.ts | 32 +++++ 9 files changed, 208 insertions(+), 21 deletions(-) create mode 100644 apps/webapp/app/presenters/v3/TasksStreamPresenter.server.ts create mode 100644 apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.v3.$projectParam.tasks.stream/route.tsx create mode 100644 apps/webapp/app/v3/services/projectPubSub.server.ts diff --git a/apps/webapp/app/presenters/v3/TasksStreamPresenter.server.ts b/apps/webapp/app/presenters/v3/TasksStreamPresenter.server.ts new file mode 100644 index 000000000..c7318c340 --- /dev/null +++ b/apps/webapp/app/presenters/v3/TasksStreamPresenter.server.ts @@ -0,0 +1,115 @@ +import { TaskRun, TaskRunAttempt } from "@trigger.dev/database"; +import { eventStream } from "remix-utils/sse/server"; +import { PrismaClient, prisma } from "~/db.server"; +import { logger } from "~/services/logger.server"; +import { eventRepository } from "~/v3/eventRepository.server"; +import { projectPubSub } from "~/v3/services/projectPubSub.server"; + +type RunWithAttempts = { + updatedAt: Date; + attempts: { + status: TaskRunAttempt["status"]; + updatedAt: Date; + }[]; +}; + +const pingInterval = 1000; + +export class TasksStreamPresenter { + #prismaClient: PrismaClient; + + constructor(prismaClient: PrismaClient = prisma) { + this.#prismaClient = prismaClient; + } + + public async call({ + request, + organizationSlug, + projectSlug, + userId, + }: { + request: Request; + organizationSlug: string; + projectSlug: string; + userId: string; + }) { + const project = await this.#prismaClient.project.findUnique({ + where: { + slug: projectSlug, + organization: { + slug: organizationSlug, + members: { + some: { + userId, + }, + }, + }, + }, + select: { + id: true, + }, + }); + + if (!project) { + return new Response("Not found", { status: 404 }); + } + + logger.info("TasksStreamPresenter.call", { + projectSlug, + }); + + let pinger: NodeJS.Timer | undefined = undefined; + + const subscriber = await projectPubSub.subscribe(`project:${project.id}:*`); + + return eventStream(request.signal, (send, close) => { + const safeSend = (args: { event?: string; data: string }) => { + try { + send(args); + } catch (error) { + if (error instanceof Error) { + if (error.name !== "TypeError") { + logger.debug("Error sending SSE, aborting", { + error: { + name: error.name, + message: error.message, + stack: error.stack, + }, + args, + }); + } + } else { + logger.debug("Unknown error sending SSE, aborting", { + error, + args, + }); + } + + close(); + } + }; + + subscriber.on("WORKER_CREATED", async (message) => { + safeSend({ data: message.createdAt.toISOString() }); + }); + + pinger = setInterval(() => { + if (request.signal.aborted) { + return close(); + } + + safeSend({ event: "ping", data: new Date().toISOString() }); + }, pingInterval); + + return async function clear() { + logger.info("TasksStreamPresenter.abort", { + projectSlug, + }); + + clearInterval(pinger); + + await subscriber.stopListening(); + }; + }); + } +} diff --git a/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.v3.$projectParam._index/route.tsx b/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.v3.$projectParam._index/route.tsx index 14100b24e..fd29dca27 100644 --- a/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.v3.$projectParam._index/route.tsx +++ b/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.v3.$projectParam._index/route.tsx @@ -1,6 +1,8 @@ import { ChatBubbleLeftRightIcon } from "@heroicons/react/20/solid"; +import { useRevalidator } from "@remix-run/react"; import { LoaderFunctionArgs } from "@remix-run/server-runtime"; import { TaskRunAttemptStatus, TaskRunStatus } from "@trigger.dev/database"; +import { useEffect } from "react"; import { typedjson, useTypedLoaderData } from "remix-typedjson"; import invariant from "tiny-invariant"; import { Feedback } from "~/components/Feedback"; @@ -28,13 +30,14 @@ import { import { TaskFunctionName, TaskPath } from "~/components/runs/v3/TaskPath"; import { TaskRunStatusCombo } from "~/components/runs/v3/TaskRunStatus"; import { useDevEnvironment } from "~/hooks/useEnvironments"; +import { useEventSource } from "~/hooks/useEventSource"; import { useOrganization } from "~/hooks/useOrganizations"; import { useProject } from "~/hooks/useProject"; import { useUser } from "~/hooks/useUser"; import { TaskListPresenter } from "~/presenters/v3/TaskListPresenter.server"; import { requireUserId } from "~/services/session.server"; import { cn } from "~/utils/cn"; -import { ProjectParamSchema, v3RunsPath } from "~/utils/pathBuilder"; +import { ProjectParamSchema, v3RunsPath, v3TasksStreamingPath } from "~/utils/pathBuilder"; export const loader = async ({ request, params }: LoaderFunctionArgs) => { const userId = await requireUserId(request); @@ -67,6 +70,19 @@ export default function Page() { const { tasks } = useTypedLoaderData(); const hasTasks = tasks.length > 0; + //live reload the page when the tasks change + const revalidator = useRevalidator(); + const streamedEvents = useEventSource(v3TasksStreamingPath(organization, project), { + event: "message", + }); + + useEffect(() => { + if (streamedEvents !== null) { + revalidator.revalidate(); + } + // WARNING Don't put the revalidator in the useEffect deps array or bad things will happen + }, [streamedEvents]); // eslint-disable-line react-hooks/exhaustive-deps + return ( @@ -179,8 +195,6 @@ function classForTaskRunStatus(status: TaskRunStatus) { } function CreateTaskInstructions() { - const devEnvironment = useDevEnvironment(); - invariant(devEnvironment, "Dev environment must be defined"); return (
diff --git a/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.v3.$projectParam.runs.$runParam/route.tsx b/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.v3.$projectParam.runs.$runParam/route.tsx index b9323eed1..71b6687bb 100644 --- a/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.v3.$projectParam.runs.$runParam/route.tsx +++ b/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.v3.$projectParam.runs.$runParam/route.tsx @@ -4,8 +4,7 @@ import { MagnifyingGlassMinusIcon, MagnifyingGlassPlusIcon, } from "@heroicons/react/20/solid"; -import { Time } from "@internationalized/date"; -import { Link, Outlet, useNavigate, useParams, useRevalidator } from "@remix-run/react"; +import { Outlet, useNavigate, useParams, useRevalidator } from "@remix-run/react"; import { LoaderFunctionArgs } from "@remix-run/server-runtime"; import { Virtualizer } from "@tanstack/react-virtual"; import { @@ -33,14 +32,7 @@ import { import { Slider } from "~/components/primitives/Slider"; import { Switch } from "~/components/primitives/Switch"; import * as Timeline from "~/components/primitives/Timeline"; -import { - GetNodePropsFn, - GetTreePropsFn, - TreeView, - TreeViewProps, - UseTreeStateOutput, - useTree, -} from "~/components/primitives/TreeView/TreeView"; +import { TreeView, UseTreeStateOutput, useTree } from "~/components/primitives/TreeView/TreeView"; import { NodesState } from "~/components/primitives/TreeView/reducer"; import { RunIcon } from "~/components/runs/v3/RunIcon"; import { SpanTitle, eventBackgroundClassName } from "~/components/runs/v3/SpanTitle"; 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/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/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/services/createBackgroundWorker.server.ts b/apps/webapp/app/v3/services/createBackgroundWorker.server.ts index 289c744a7..1c937ff3f 100644 --- a/apps/webapp/app/v3/services/createBackgroundWorker.server.ts +++ b/apps/webapp/app/v3/services/createBackgroundWorker.server.ts @@ -7,6 +7,7 @@ import { generateFriendlyId } from "../friendlyIdentifiers"; import { marqs } from "../marqs.server"; import { calculateNextBuildVersion } from "../utils/calculateNextBuildVersion"; import { BaseService } from "./baseService.server"; +import { projectPubSub } from "./projectPubSub.server"; export class CreateBackgroundWorkerService extends BaseService { public async call( @@ -67,6 +68,15 @@ export class CreateBackgroundWorkerService extends BaseService { await createBackgroundTasks(body.metadata.tasks, backgroundWorker, environment, this._prisma); + //send a notification that a new worker has been created + await projectPubSub.publish(`project:${project.id}:env:${environment.id}`, "WORKER_CREATED", { + environmentId: environment.id, + environmentType: environment.type, + createdAt: backgroundWorker.createdAt, + taskCount: body.metadata.tasks.length, + type: "local", + }); + return backgroundWorker; }); } diff --git a/apps/webapp/app/v3/services/createDeployedBackgroundWorker.server.ts b/apps/webapp/app/v3/services/createDeployedBackgroundWorker.server.ts index 175774d24..1acaf9379 100644 --- a/apps/webapp/app/v3/services/createDeployedBackgroundWorker.server.ts +++ b/apps/webapp/app/v3/services/createDeployedBackgroundWorker.server.ts @@ -5,6 +5,7 @@ import { generateFriendlyId } from "../friendlyIdentifiers"; import { BaseService } from "./baseService.server"; import { createBackgroundTasks } from "./createBackgroundWorker.server"; import { CURRENT_DEPLOYMENT_LABEL } from "~/consts"; +import { projectPubSub } from "./projectPubSub.server"; export class CreateDeployedBackgroundWorkerService extends BaseService { public async call( @@ -71,6 +72,19 @@ export class CreateDeployedBackgroundWorkerService extends BaseService { }, }); + //send a notification that a new worker has been created + await projectPubSub.publish( + `project:${environment.projectId}:env:${environment.id}`, + "WORKER_CREATED", + { + environmentId: environment.id, + environmentType: environment.type, + createdAt: backgroundWorker.createdAt, + taskCount: body.metadata.tasks.length, + type: "deployed", + } + ); + return backgroundWorker; }); } diff --git a/apps/webapp/app/v3/services/projectPubSub.server.ts b/apps/webapp/app/v3/services/projectPubSub.server.ts new file mode 100644 index 000000000..826e04ad8 --- /dev/null +++ b/apps/webapp/app/v3/services/projectPubSub.server.ts @@ -0,0 +1,32 @@ +import { z } from "zod"; +import { singleton } from "~/utils/singleton"; +import { ZodPubSub, ZodSubscriber } from "../utils/zodPubSub.server"; +import { env } from "~/env.server"; + +const messageCatalog = { + WORKER_CREATED: z.object({ + environmentId: z.string(), + environmentType: z.string(), + createdAt: z.coerce.date(), + taskCount: z.number(), + type: z.union([z.literal("local"), z.literal("deployed")]), + }), +}; + +export type ProjectSubscriber = ZodSubscriber; + +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, + }); +} From 0e347b001b2f3602a7c322aa2aa7099124d0c45d Mon Sep 17 00:00:00 2001 From: Eric Allam Date: Tue, 19 Mar 2024 16:04:26 +0000 Subject: [PATCH 9/9] =?UTF-8?q?Don=E2=80=99t=20override=20TRIGGER=5FSECRET?= =?UTF-8?q?=5FKEY=20and=20TRIGGER=5FAPI=5FURL=20with=20local=20env=20vars?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- packages/cli-v3/src/workers/dev/backgroundWorker.ts | 10 ++++++++-- 1 file changed, 8 insertions(+), 2 deletions(-) diff --git a/packages/cli-v3/src/workers/dev/backgroundWorker.ts b/packages/cli-v3/src/workers/dev/backgroundWorker.ts index 26b4823df..391b162b8 100644 --- a/packages/cli-v3/src/workers/dev/backgroundWorker.ts +++ b/packages/cli-v3/src/workers/dev/backgroundWorker.ts @@ -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; }