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, + }); +}