diff --git a/.changeset/soft-ties-turn.md b/.changeset/soft-ties-turn.md new file mode 100644 index 00000000000..0f9393f1e8e --- /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 00000000000..ad47ab18360 --- /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 00000000000..8d479b47fed --- /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 58228c57e05..ef5e50bcad0 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 00000000000..0e56015ef28 --- /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 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 44b9e331194..f57cb50dd82 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 f7c980b6d30..4f5a0e8d8f9 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 2f0efc2e7d5..f46f3384dd1 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 bbb5e8b21be..4743dd2ccfc 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); }