From 6e1b8a11d464617201d889325f3fe2ead2e5bf6e Mon Sep 17 00:00:00 2001 From: hmacr Date: Sun, 15 Oct 2023 19:06:03 +0530 Subject: [PATCH] feat: allow cancelling job runs for an event-id --- .changeset/soft-ties-turn.md | 6 ++ .../api.v1.events.$eventId.cancel-runs.ts | 49 +++++++++++++ .../events/cancelRunsForEvent.server.ts | 70 +++++++++++++++++++ docs/mint.json | 1 + .../instancemethods/cancel-runs-for-event.mdx | 41 +++++++++++ docs/sdk/triggerclient/overview.mdx | 4 ++ packages/core/src/schemas/events.ts | 7 ++ packages/trigger-sdk/src/apiClient.ts | 21 ++++++ packages/trigger-sdk/src/triggerClient.ts | 4 ++ 9 files changed, 203 insertions(+) create mode 100644 .changeset/soft-ties-turn.md create mode 100644 apps/webapp/app/routes/api.v1.events.$eventId.cancel-runs.ts create mode 100644 apps/webapp/app/services/events/cancelRunsForEvent.server.ts create mode 100644 docs/sdk/triggerclient/instancemethods/cancel-runs-for-event.mdx diff --git a/.changeset/soft-ties-turn.md b/.changeset/soft-ties-turn.md new file mode 100644 index 000000000..0f9393f1e --- /dev/null +++ b/.changeset/soft-ties-turn.md @@ -0,0 +1,6 @@ +--- +"@trigger.dev/sdk": patch +"@trigger.dev/core": patch +--- + +implement functionality to cancel job runs triggered by a given eventId. diff --git a/apps/webapp/app/routes/api.v1.events.$eventId.cancel-runs.ts b/apps/webapp/app/routes/api.v1.events.$eventId.cancel-runs.ts new file mode 100644 index 000000000..ad47ab183 --- /dev/null +++ b/apps/webapp/app/routes/api.v1.events.$eventId.cancel-runs.ts @@ -0,0 +1,49 @@ +import type { ActionArgs } from "@remix-run/server-runtime"; +import { json } from "@remix-run/server-runtime"; +import { z } from "zod"; +import { authenticateApiRequest } from "~/services/apiAuth.server"; +import { CancelEventService } from "~/services/events/cancelEvent.server"; +import { logger } from "~/services/logger.server"; +import { CancelRunsForEventService } from "~/services/events/cancelRunsForEvent.server"; + +const ParamsSchema = z.object({ + eventId: z.string(), +}); + +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" }; + } + + // Next authenticate the request + const authenticationResult = await authenticateApiRequest(request); + + if (!authenticationResult) { + return json({ error: "Invalid or Missing API key" }, { status: 401 }); + } + + const authenticatedEnv = authenticationResult.environment; + + const parsed = ParamsSchema.safeParse(params); + + if (!parsed.success) { + return json({ error: "Invalid or Missing eventId" }, { status: 400 }); + } + + const { eventId } = parsed.data; + + const service = new CancelRunsForEventService(); + try { + const res = await service.call(authenticatedEnv, eventId); + + if (!res) { + return json({ error: "Event not found" }, { status: 404 }); + } + + return json(res); + } catch (err) { + logger.error("CancelRunsForEventService.call() error", { error: err }); + return json({ error: "Internal Server Error" }, { status: 500 }); + } +} diff --git a/apps/webapp/app/services/events/cancelRunsForEvent.server.ts b/apps/webapp/app/services/events/cancelRunsForEvent.server.ts new file mode 100644 index 000000000..8d479b47f --- /dev/null +++ b/apps/webapp/app/services/events/cancelRunsForEvent.server.ts @@ -0,0 +1,70 @@ +import { $transaction, PrismaClient, prisma } from "~/db.server"; +import { AuthenticatedEnvironment } from "../apiAuth.server"; +import { JobRunStatus } from "@trigger.dev/database"; +import { CancelRunService } from "../runs/cancelRun.server"; +import { logger } from "../logger.server"; +import { CancelRunsForEvent } from "@trigger.dev/core/schemas/events"; + +const CANCELLABLE_JOB_RUN_STATUS: JobRunStatus[] = [ + JobRunStatus.PENDING, + JobRunStatus.QUEUED, + JobRunStatus.WAITING_ON_CONNECTIONS, + JobRunStatus.PREPROCESSING, + JobRunStatus.STARTED, +]; + +export class CancelRunsForEventService { + #prismaClient: PrismaClient; + + constructor(prismaClient: PrismaClient = prisma) { + this.#prismaClient = prismaClient; + } + + public async call(environment: AuthenticatedEnvironment, eventId: string) { + return await $transaction(this.#prismaClient, async (tx) => { + const event = await tx.eventRecord.findUnique({ + where: { + eventId_environmentId: { + eventId: eventId, + environmentId: environment.id, + }, + }, + }); + + if (!event) { + return; + } + + const jobRuns = await tx.jobRun.findMany({ + where: { + eventId: event.id, + status: { + in: CANCELLABLE_JOB_RUN_STATUS, + }, + }, + select: { + id: true, + }, + }); + + const cancelRunService = new CancelRunService(this.#prismaClient); + const cancelledRunIds: string[] = []; + const failedToCancelRunIds: string[] = []; + + for (const jobRun of jobRuns) { + try { + await cancelRunService.call({ runId: jobRun.id }); + cancelledRunIds.push(jobRun.id); + } catch (err) { + logger.debug(`failed to cancel job run with id ${jobRun.id} for event id ${eventId}`); + failedToCancelRunIds.push(jobRun.id); + } + } + + return { + cancelled_run_ids: cancelledRunIds, + failed_to_cancel_run_ids: failedToCancelRunIds, + }; + }); + } +} diff --git a/docs/mint.json b/docs/mint.json index 58228c57e..ef5e50bca 100644 --- a/docs/mint.json +++ b/docs/mint.json @@ -312,6 +312,7 @@ "sdk/triggerclient/instancemethods/sendevent", "sdk/triggerclient/instancemethods/getevent", "sdk/triggerclient/instancemethods/cancel-event", + "sdk/triggerclient/instancemethods/cancel-runs-for-event", "sdk/triggerclient/instancemethods/getruns", "sdk/triggerclient/instancemethods/getrun", "sdk/triggerclient/instancemethods/define-job", diff --git a/docs/sdk/triggerclient/instancemethods/cancel-runs-for-event.mdx b/docs/sdk/triggerclient/instancemethods/cancel-runs-for-event.mdx new file mode 100644 index 000000000..790846a5d --- /dev/null +++ b/docs/sdk/triggerclient/instancemethods/cancel-runs-for-event.mdx @@ -0,0 +1,41 @@ +--- +title: "TriggerClient: cancelRunsForEvent() Instance Method" +sidebarTitle: "cancelRunsForEvent()" +description: "The `cancelRunsForEvent()` instance method will cancel all the job runs (yet to be executed) that are triggered by a given eventId." +--- + +## Parameters + + + The event ID to cancel the job runs for. This is returned when calling either + [client.sendEvent()](/sdk/triggerclient/instancemethods/sendevent) or + [io.sendEvent()](/sdk/io/sendevent). + + +## Returns + + + + + List of Job Run IDs that are cancelled. + + + List of Job Run IDs that have been failed to be cancelled. + + + + + + +```ts Cancelling Runs for an Event +const event = client.sendEvent({ + name: "test.job", +}); + +const res = client.cancelRunsForEvent(event.id); + +console.log(res.cancelled_run_ids); +console.log(res.failed_to_cancel_run_ids); +``` + + diff --git a/docs/sdk/triggerclient/overview.mdx b/docs/sdk/triggerclient/overview.mdx index 44b9e3311..f57cb50dd 100644 --- a/docs/sdk/triggerclient/overview.mdx +++ b/docs/sdk/triggerclient/overview.mdx @@ -46,6 +46,10 @@ The `getEvent()` method gets the event details for a given eventId. The `cancelEvent()` method cancels an event that is scheduled to be delivered in the future. +#### [cancelRunsForEvent()](/sdk/triggerclient/instancemethods/cancel-runs-for-event) + +The `cancelRunsForEvent()` method cancels the job runs (yet to be executed) that are triggered by a given eventId. + #### [getRuns()](/sdk/triggerclient/instancemethods/getruns) The `getRuns()` method gets runs for a Job. diff --git a/packages/core/src/schemas/events.ts b/packages/core/src/schemas/events.ts index f7c980b6d..4f5a0e8d8 100644 --- a/packages/core/src/schemas/events.ts +++ b/packages/core/src/schemas/events.ts @@ -26,3 +26,10 @@ export const GetEventSchema = z.object({ }); export type GetEvent = z.infer; + +export const CancelRunsForEventSchema = z.object({ + cancelled_run_ids: z.array(z.string()), + failed_to_cancel_run_ids: z.array(z.string()), +}); + +export type CancelRunsForEvent = z.infer; diff --git a/packages/trigger-sdk/src/apiClient.ts b/packages/trigger-sdk/src/apiClient.ts index 2f0efc2e7..f46f3384d 100644 --- a/packages/trigger-sdk/src/apiClient.ts +++ b/packages/trigger-sdk/src/apiClient.ts @@ -1,6 +1,7 @@ import { ApiEventLog, ApiEventLogSchema, + CancelRunsForEventSchema, CompleteTaskBodyInput, ConnectionAuthSchema, FailTaskBodyInput, @@ -215,6 +216,26 @@ export class ApiClient { }); } + async cancelRunsForEvent(eventId: string) { + const apiKey = await this.#apiKey(); + + this.#logger.debug("Cancelling runs for event", { + eventId, + }); + + return await zodfetch( + CancelRunsForEventSchema, + `${this.#apiUrl}/api/v1/events/${eventId}/cancel-runs`, + { + method: "POST", + headers: { + "Content-Type": "application/json", + Authorization: `Bearer ${apiKey}`, + }, + } + ); + } + async updateStatus(runId: string, id: string, status: StatusUpdate) { const apiKey = await this.#apiKey(); diff --git a/packages/trigger-sdk/src/triggerClient.ts b/packages/trigger-sdk/src/triggerClient.ts index bbb5e8b21..4743dd2cc 100644 --- a/packages/trigger-sdk/src/triggerClient.ts +++ b/packages/trigger-sdk/src/triggerClient.ts @@ -644,6 +644,10 @@ export class TriggerClient { return this.#client.cancelEvent(eventId); } + async cancelRunsForEvent(eventId: string) { + return this.#client.cancelRunsForEvent(eventId); + } + async updateStatus(runId: string, id: string, status: StatusUpdate) { return this.#client.updateStatus(runId, id, status); }