From 302bd02ff16bdba5dce811e93c605d0f832c7920 Mon Sep 17 00:00:00 2001 From: Eric Allam Date: Tue, 22 Aug 2023 10:46:46 +0100 Subject: [PATCH] Issue #377: only expose the external eventId in the API (#380) * Issue #377: only expose the external eventId in the API * Create eighty-zebras-bow.md --- .changeset/eighty-zebras-bow.md | 6 +++ apps/webapp/app/api.server.ts | 15 ++++++ .../routes/api.v1.events.$eventId.cancel.ts | 3 +- .../app/routes/api.v1.events.$eventId.ts | 47 ++++++++++++++----- apps/webapp/app/routes/api.v1.events.ts | 7 ++- .../runs/performRunExecutionV1.server.ts | 6 +-- .../runs/performRunExecutionV2.server.ts | 6 +-- examples/job-catalog/src/events.ts | 6 ++- packages/core/src/schemas/api.ts | 12 ++--- packages/trigger-sdk/src/apiClient.ts | 19 -------- packages/trigger-sdk/src/io.ts | 19 ++++++++ 11 files changed, 97 insertions(+), 49 deletions(-) create mode 100644 .changeset/eighty-zebras-bow.md create mode 100644 apps/webapp/app/api.server.ts diff --git a/.changeset/eighty-zebras-bow.md b/.changeset/eighty-zebras-bow.md new file mode 100644 index 000000000..80d37577a --- /dev/null +++ b/.changeset/eighty-zebras-bow.md @@ -0,0 +1,6 @@ +--- +"@trigger.dev/core": patch +"@trigger.dev/sdk": patch +--- + +Issue #377: only expose the external eventId in the API diff --git a/apps/webapp/app/api.server.ts b/apps/webapp/app/api.server.ts new file mode 100644 index 000000000..b80891361 --- /dev/null +++ b/apps/webapp/app/api.server.ts @@ -0,0 +1,15 @@ +import { ApiEventLog } from "@trigger.dev/core"; +import { EventRecord } from "@trigger.dev/database"; + +export function eventRecordToApiJson(eventRecord: EventRecord): ApiEventLog { + return { + id: eventRecord.eventId, + name: eventRecord.name, + payload: eventRecord.payload as any, + context: eventRecord.context as any, + timestamp: eventRecord.timestamp, + deliverAt: eventRecord.deliverAt, + deliveredAt: eventRecord.deliveredAt, + cancelledAt: eventRecord.cancelledAt, + }; +} 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 b1cc62d64..28dcdf1ca 100644 --- a/apps/webapp/app/routes/api.v1.events.$eventId.cancel.ts +++ b/apps/webapp/app/routes/api.v1.events.$eventId.cancel.ts @@ -1,6 +1,7 @@ import type { ActionArgs } from "@remix-run/server-runtime"; import { json } from "@remix-run/server-runtime"; import { z } from "zod"; +import { eventRecordToApiJson } from "~/api.server"; import { authenticateApiRequest } from "~/services/apiAuth.server"; import { CancelEventService } from "~/services/events/cancelEvent.server"; import { logger } from "~/services/logger.server"; @@ -40,7 +41,7 @@ export async function action({ request, params }: ActionArgs) { return json({ error: "Event not found" }, { status: 404 }); } - return json(updatedEvent); + return json(eventRecordToApiJson(updatedEvent)); } catch (err) { logger.error("CancelEventService.call() error", { error: err, diff --git a/apps/webapp/app/routes/api.v1.events.$eventId.ts b/apps/webapp/app/routes/api.v1.events.$eventId.ts index 83a63e299..da29947dd 100644 --- a/apps/webapp/app/routes/api.v1.events.$eventId.ts +++ b/apps/webapp/app/routes/api.v1.events.$eventId.ts @@ -1,6 +1,6 @@ -import type { ActionArgs, LoaderArgs } from "@remix-run/server-runtime"; +import type { LoaderArgs } from "@remix-run/server-runtime"; import { json } from "@remix-run/server-runtime"; -import { cors } from "remix-utils"; +import { GetEvent } from "@trigger.dev/core"; import { z } from "zod"; import { prisma } from "~/db.server"; import { authenticateApiRequest } from "~/services/apiAuth.server"; @@ -32,9 +32,36 @@ export async function loader({ request, params }: LoaderArgs) { const { eventId } = parsed.data; - const event = await prisma.eventRecord.findFirst({ + const event = await findEventRecord(eventId, authenticatedEnv.id); + + if (!event) { + return apiCors(request, json({ error: "Event not found" }, { status: 404 })); + } + + return apiCors(request, json(toJSON(event))); +} + +function toJSON(eventRecord: FoundEventRecord): GetEvent { + return { + id: eventRecord.eventId, + name: eventRecord.name, + createdAt: eventRecord.createdAt, + updatedAt: eventRecord.updatedAt, + runs: eventRecord.runs.map((run) => ({ + id: run.id, + status: run.status, + startedAt: run.startedAt, + completedAt: run.completedAt, + })), + }; +} + +type FoundEventRecord = NonNullable>>; + +async function findEventRecord(eventId: string, environmentId: string) { + return await prisma.eventRecord.findUnique({ select: { - id: true, + eventId: true, name: true, createdAt: true, updatedAt: true, @@ -48,14 +75,10 @@ export async function loader({ request, params }: LoaderArgs) { }, }, where: { - id: eventId, - environmentId: authenticatedEnv.id, + eventId_environmentId: { + eventId, + environmentId, + }, }, }); - - if (!event) { - return apiCors(request, json({ error: "Event not found" }, { status: 404 })); - } - - return apiCors(request, json(event)); } diff --git a/apps/webapp/app/routes/api.v1.events.ts b/apps/webapp/app/routes/api.v1.events.ts index 25fe48476..403138bc7 100644 --- a/apps/webapp/app/routes/api.v1.events.ts +++ b/apps/webapp/app/routes/api.v1.events.ts @@ -4,6 +4,7 @@ import { SendEventBodySchema } from "@trigger.dev/core"; import { generateErrorMessage } from "zod-error"; import { authenticateApiRequest } from "~/services/apiAuth.server"; import { IngestSendEvent } from "~/services/events/ingestSendEvent.server"; +import { eventRecordToApiJson } from "~/api.server"; export async function action({ request }: ActionArgs) { // Ensure this is a POST request @@ -33,5 +34,9 @@ export async function action({ request }: ActionArgs) { const event = await service.call(authenticatedEnv, body.data.event, body.data.options); - return json(event); + if (!event) { + return json({ error: "Failed to create event" }, { status: 500 }); + } + + return json(eventRecordToApiJson(event)); } diff --git a/apps/webapp/app/services/runs/performRunExecutionV1.server.ts b/apps/webapp/app/services/runs/performRunExecutionV1.server.ts index b56199a1b..c886427db 100644 --- a/apps/webapp/app/services/runs/performRunExecutionV1.server.ts +++ b/apps/webapp/app/services/runs/performRunExecutionV1.server.ts @@ -1,5 +1,4 @@ import { - ApiEventLogSchema, CachedTaskSchema, RunJobError, RunJobResumeWithTask, @@ -9,6 +8,7 @@ import { } from "@trigger.dev/core"; import type { Task } from "@trigger.dev/database"; import { generateErrorMessage } from "zod-error"; +import { eventRecordToApiJson } from "~/api.server"; import { EXECUTE_JOB_RETRY_LIMIT } from "~/consts"; import { $transaction, PrismaClient, PrismaClientOrTransaction, prisma } from "~/db.server"; import { enqueueRunExecutionV1 } from "~/models/jobRunExecution.server"; @@ -54,7 +54,7 @@ export class PerformRunExecutionV1Service { const { run } = execution; const client = new EndpointApi(run.environment.apiKey, run.endpoint.url); - const event = ApiEventLogSchema.parse({ ...run.event, id: run.eventId }); + const event = eventRecordToApiJson(run.event); const startedAt = new Date(); await this.#prismaClient.jobRunExecution.update({ @@ -174,7 +174,7 @@ export class PerformRunExecutionV1Service { } const client = new EndpointApi(run.environment.apiKey, run.endpoint.url); - const event = ApiEventLogSchema.parse({ ...run.event, id: run.eventId }); + const event = eventRecordToApiJson(run.event); const startedAt = new Date(); diff --git a/apps/webapp/app/services/runs/performRunExecutionV2.server.ts b/apps/webapp/app/services/runs/performRunExecutionV2.server.ts index d27b2df0f..b8e5cfe62 100644 --- a/apps/webapp/app/services/runs/performRunExecutionV2.server.ts +++ b/apps/webapp/app/services/runs/performRunExecutionV2.server.ts @@ -1,5 +1,4 @@ import { - ApiEventLogSchema, CachedTask, RunJobError, RunJobResumeWithTask, @@ -9,6 +8,7 @@ import { } from "@trigger.dev/core"; import type { Task } from "@trigger.dev/database"; import { generateErrorMessage } from "zod-error"; +import { eventRecordToApiJson } from "~/api.server"; import { $transaction, PrismaClient, PrismaClientOrTransaction, prisma } from "~/db.server"; import { enqueueRunExecutionV2 } from "~/models/jobRunExecution.server"; import { resolveRunConnections } from "~/models/runConnection.server"; @@ -57,7 +57,7 @@ export class PerformRunExecutionV2Service { // the run execution will be marked as failed and the run will start async #executePreprocessing(run: FoundRun) { const client = new EndpointApi(run.environment.apiKey, run.endpoint.url); - const event = ApiEventLogSchema.parse({ ...run.event, id: run.eventId }); + const event = eventRecordToApiJson(run.event); const { response, parser } = await client.preprocessRunRequest({ event, @@ -146,7 +146,7 @@ export class PerformRunExecutionV2Service { } const client = new EndpointApi(run.environment.apiKey, run.endpoint.url); - const event = ApiEventLogSchema.parse({ ...run.event, id: run.eventId }); + const event = eventRecordToApiJson(run.event); const startedAt = new Date(); diff --git a/examples/job-catalog/src/events.ts b/examples/job-catalog/src/events.ts index 154955a3e..acc834179 100644 --- a/examples/job-catalog/src/events.ts +++ b/examples/job-catalog/src/events.ts @@ -39,15 +39,19 @@ client.defineJob({ run: async (payload, io, ctx) => { await io.sendEvent( "send-event", - { name: "Cancellable Event", id: payload.id }, + { name: "Cancellable Event", id: payload.id, payload: { payload, ctx } }, { deliverAt: new Date(Date.now() + 1000 * 60 * 60 * 24), // 24 hours from now } ); + await io.getEvent("get-event", payload.id); + await io.wait("wait-1", 60); // 1 minute await io.cancelEvent("cancel-event", payload.id); + + await io.getEvent("get-event-2", payload.id); }, }); diff --git a/packages/core/src/schemas/api.ts b/packages/core/src/schemas/api.ts index d90a06066..f4b78e0ea 100644 --- a/packages/core/src/schemas/api.ts +++ b/packages/core/src/schemas/api.ts @@ -257,6 +257,9 @@ export const ApiEventLogSchema = z.object({ /** The timestamp when the event was delivered. Is `undefined` if `deliverAt` or `deliverAfter` were set when sending the event. */ deliveredAt: z.coerce.date().optional().nullable(), + /** The timestamp when the event was cancelled. Is `undefined` if the event + * wasn't cancelled. */ + cancelledAt: z.coerce.date().optional().nullable(), }); export type ApiEventLog = z.infer; @@ -424,15 +427,6 @@ export const PreprocessRunResponseSchema = z.object({ export type PreprocessRunResponse = z.infer; -export const CreateRunBodySchema = z.object({ - client: z.string(), - job: JobMetadataSchema, - event: ApiEventLogSchema, - properties: z.array(DisplayPropertySchema).optional(), -}); - -export type CreateRunBody = z.infer; - const CreateRunResponseOkSchema = z.object({ ok: z.literal(true), data: z.object({ diff --git a/packages/trigger-sdk/src/apiClient.ts b/packages/trigger-sdk/src/apiClient.ts index 3ea557f79..d5fcc54a1 100644 --- a/packages/trigger-sdk/src/apiClient.ts +++ b/packages/trigger-sdk/src/apiClient.ts @@ -3,8 +3,6 @@ import { ApiEventLogSchema, CompleteTaskBodyInput, ConnectionAuthSchema, - CreateRunBody, - CreateRunResponseBodySchema, FailTaskBodyInput, GetEventSchema, GetRunOptionsWithTaskDetails, @@ -105,23 +103,6 @@ export class ApiClient { return await response.json(); } - async createRun(params: CreateRunBody) { - const apiKey = await this.#apiKey(); - - this.#logger.debug("Creating run", { - params, - }); - - return await zodfetch(CreateRunResponseBodySchema, `${this.#apiUrl}/api/v1/runs`, { - method: "POST", - headers: { - "Content-Type": "application/json", - Authorization: `Bearer ${apiKey}`, - }, - body: JSON.stringify(params), - }); - } - async runTask(runId: string, task: RunTaskBodyInput) { const apiKey = await this.#apiKey(); diff --git a/packages/trigger-sdk/src/io.ts b/packages/trigger-sdk/src/io.ts index 7bcd3c973..cf7769573 100644 --- a/packages/trigger-sdk/src/io.ts +++ b/packages/trigger-sdk/src/io.ts @@ -221,6 +221,25 @@ export class IO { ); } + async getEvent(key: string | any[], id: string) { + return await this.runTask( + key, + { + name: "getEvent", + params: { id }, + properties: [ + { + label: "id", + text: id, + }, + ], + }, + async (task) => { + return await this._triggerClient.getEvent(id); + } + ); + } + /** `io.cancelEvent()` allows you to cancel an event that was previously sent with `io.sendEvent()`. This will prevent any Jobs from running that are listening for that event if the event was sent with a delay * @param key * @param eventId