Files
triggerdotdev--trigger.dev/apps/webapp/app/presenters/RunStreamPresenter.server.ts
Eric Allam 9a7c08c26a
🚀 Publish Trigger.dev Docker / typecheck (push) Failing after 1s
🚀 Publish Trigger.dev Docker / units (push) Failing after 0s
🚀 Publish Trigger.dev Docker / e2e (push) Failing after 0s
🚀 Publish Trigger.dev Docker / publish (push) Has been skipped
Improvements: Fix dangling SSE issue and compression memory leak (#733)
* Downgrade to remix-auth-email-link to remove yarn dependency

* Turn off the pg listen service for now

* Add snapshot admin route

* A couple logger fixes

* Add ability to disable compression

* Add ability to disable SSE

* Fix SSE memory leak + DB load issue
2023-11-10 16:02:53 +00:00

72 lines
1.8 KiB
TypeScript

import { JobRun } from "@trigger.dev/database";
import { PrismaClient, prisma } from "~/db.server";
import { sse } from "~/utils/sse.server";
export class RunStreamPresenter {
#prismaClient: PrismaClient;
constructor(prismaClient: PrismaClient = prisma) {
this.#prismaClient = prismaClient;
}
public async call({ request, runId }: { request: Request; runId: JobRun["id"] }) {
const run = await this.#runForUpdates(runId);
if (!run) {
return new Response("Not found", { status: 404 });
}
let lastUpdatedAt: number = run.updatedAt.getTime();
let lastTotalTaskUpdatedTime = run.tasks.reduce(
(prev, task) => prev + task.updatedAt.getTime(),
0
);
return sse({
request,
run: async (send, stop) => {
const result = await this.#runForUpdates(runId);
if (!result) {
return stop();
}
if (result.completedAt) {
send({ data: new Date().toISOString() });
return stop();
}
const totalRunUpdated = result.tasks.reduce(
(prev, task) => prev + task.updatedAt.getTime(),
0
);
if (lastUpdatedAt !== result.updatedAt.getTime()) {
send({ data: result.updatedAt.toISOString() });
} else if (lastTotalTaskUpdatedTime !== totalRunUpdated) {
send({ data: new Date().toISOString() });
}
lastUpdatedAt = result.updatedAt.getTime();
lastTotalTaskUpdatedTime = totalRunUpdated;
},
});
}
#runForUpdates(id: string) {
return this.#prismaClient.jobRun.findUnique({
where: {
id,
},
select: {
updatedAt: true,
completedAt: true,
tasks: {
select: {
updatedAt: true,
},
},
},
});
}
}