feat: allow cancelling job runs for an event-id
This commit is contained in:
@@ -0,0 +1,6 @@
|
||||
---
|
||||
"@trigger.dev/sdk": patch
|
||||
"@trigger.dev/core": patch
|
||||
---
|
||||
|
||||
implement functionality to cancel job runs triggered by a given eventId.
|
||||
@@ -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 });
|
||||
}
|
||||
}
|
||||
@@ -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<CancelRunsForEvent | undefined>(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,
|
||||
};
|
||||
});
|
||||
}
|
||||
}
|
||||
@@ -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",
|
||||
|
||||
@@ -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
|
||||
|
||||
<ResponseField name="eventId" type="string" required>
|
||||
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).
|
||||
</ResponseField>
|
||||
|
||||
## Returns
|
||||
|
||||
<ResponseField type="object">
|
||||
<Expandable title="properties" defaultOpen>
|
||||
<ResponseField name="cancelled_run_ids" type="array" required>
|
||||
List of Job Run IDs that are cancelled.
|
||||
</ResponseField>
|
||||
<ResponseField name="failed_to_cancel_run_ids" type="array" required>
|
||||
List of Job Run IDs that have been failed to be cancelled.
|
||||
</ResponseField>
|
||||
</Expandable>
|
||||
</ResponseField>
|
||||
|
||||
<RequestExample>
|
||||
|
||||
```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);
|
||||
```
|
||||
|
||||
</RequestExample>
|
||||
@@ -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.
|
||||
|
||||
@@ -26,3 +26,10 @@ export const GetEventSchema = z.object({
|
||||
});
|
||||
|
||||
export type GetEvent = z.infer<typeof GetEventSchema>;
|
||||
|
||||
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<typeof CancelRunsForEventSchema>;
|
||||
|
||||
@@ -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();
|
||||
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user