diff --git a/apps/webapp/app/platform/zodWorker.server.ts b/apps/webapp/app/platform/zodWorker.server.ts index 092902836..fb38c95fe 100644 --- a/apps/webapp/app/platform/zodWorker.server.ts +++ b/apps/webapp/app/platform/zodWorker.server.ts @@ -175,10 +175,9 @@ export class ZodWorker { if (!this.#runner) { throw new Error("Worker not initialized"); } + const job = await this.#removeJob(jobKey, option?.tx ?? this.#prisma); - logger.debug("dequeued worker task", { - job, - }); + logger.debug("dequeued worker task", { job }); return job; } @@ -206,7 +205,7 @@ export class ZodWorker { spec.queueName || null, spec.runAt || null, spec.maxAttempts || null, - identifier, + spec.jobKey, spec.priority || null, spec.jobKeyMode || null, spec.flags || null @@ -226,22 +225,25 @@ export class ZodWorker { } async #removeJob(jobKey: string, tx: PrismaClientOrTransaction) { - const result = await tx.$queryRawUnsafe( - `SELECT * FROM graphile_worker.remove_job( + try { + const result = await tx.$queryRawUnsafe( + `SELECT * FROM graphile_worker.remove_job( job_key => $1::text )`, - jobKey - ); - - const job = GraphileJobSchema.safeParse(result); - - if (!job.success) { - throw new Error( - `Failed to remove job from queue, zod parsing error: ${JSON.stringify(job.error)}` + jobKey ); - } + const job = AddJobResultsSchema.safeParse(result); - return job.data as GraphileJob; + if (!job.success) { + throw new Error( + `Failed to remove job from queue, zod parsing error: ${JSON.stringify(job.error)}` + ); + } + + return job.data[0] as GraphileJob; + } catch (e) { + throw new Error(`Failed to remove job from queue, ${e}}`); + } } #createTaskListFromTasks() { diff --git a/apps/webapp/app/routes/api.v1.events.$eventId.cancel.ts b/apps/webapp/app/routes/api.v1.events.$eventId.cancel.ts index da9d3d427..6ca26a901 100644 --- a/apps/webapp/app/routes/api.v1.events.$eventId.cancel.ts +++ b/apps/webapp/app/routes/api.v1.events.$eventId.cancel.ts @@ -1,4 +1,4 @@ -import type { ActionArgs, LoaderArgs } from "@remix-run/server-runtime"; +import type { ActionArgs } from "@remix-run/server-runtime"; import { json } from "@remix-run/server-runtime"; import { z } from "zod"; import { prisma } from "~/db.server"; @@ -10,7 +10,7 @@ const ParamsSchema = z.object({ eventId: z.string(), }); -export async function loader({ request, params }: LoaderArgs) { +export async function action({ request, params }: ActionArgs) { // Ensure this is a POST request if (request.method.toUpperCase() !== "POST") { return { status: 405, body: "Method Not Allowed" }; @@ -57,22 +57,15 @@ export async function loader({ request, params }: LoaderArgs) { if (!event) { return json({ error: "Event not found" }, { status: 404 }); } - // Update the event with cancelledAt - let updatedEvent; - try { - updatedEvent = await prisma.eventRecord.update({ - where: { id: event.id }, - data: { cancelledAt: new Date() }, - }); +//update the cancelledAt column in the eventRecord table + const updatedEvent = await prisma.eventRecord.update({ + where: { id: event.id }, + data: { cancelledAt: new Date() }, + }); - // Dequeue the event only if the update is successful - await workerQueue.dequeue(updatedEvent.id, { tx: prisma }); + // Dequeue the event after the db has been updated + await workerQueue.dequeue(event.id, { tx: prisma }); - return json(updatedEvent); - } catch (error) { - logger.error("Failed to update event", { error, event }); - - return json({ error: "Failed to update event" }, { status: 500 }); - } + return json(updatedEvent); } diff --git a/apps/webapp/app/services/events/ingestSendEvent.server.ts b/apps/webapp/app/services/events/ingestSendEvent.server.ts index e32576cb5..ff8de199f 100644 --- a/apps/webapp/app/services/events/ingestSendEvent.server.ts +++ b/apps/webapp/app/services/events/ingestSendEvent.server.ts @@ -6,10 +6,7 @@ import { workerQueue } from "~/services/worker.server"; export class IngestSendEvent { #prismaClient: PrismaClientOrTransaction; - constructor( - prismaClient: PrismaClientOrTransaction = prisma, - private deliverEvents = true - ) { + constructor(prismaClient: PrismaClientOrTransaction = prisma, private deliverEvents = true) { this.#prismaClient = prismaClient; } @@ -91,7 +88,7 @@ export class IngestSendEvent { { id: eventLog.id, }, - { runAt: eventLog.deliverAt, tx } + { runAt: eventLog.deliverAt, tx, jobKey: eventLog.id } ); }