From ce0cd1d643b2a3164f99a49ce2ed3d4c36b92c25 Mon Sep 17 00:00:00 2001 From: Matt Aitken Date: Thu, 9 Apr 2026 19:42:32 +0100 Subject: [PATCH] Use resolveEventRepositoryForStore everywhere --- .../app/routes/api.v1.runs.$runId.events.ts | 5 ++- .../api.v1.runs.$runId.spans.$spanId.ts | 5 ++- .../app/routes/api.v1.runs.$runId.trace.ts | 5 ++- .../resources.runs.$runParam.logs.download.ts | 5 ++- .../app/v3/eventRepository/index.server.ts | 15 +++++++ .../webapp/app/v3/runEngineHandlers.server.ts | 45 ++++++++++++++----- .../app/v3/services/cancelTaskRunV1.server.ts | 5 ++- .../app/v3/services/completeAttempt.server.ts | 15 +++++-- .../app/v3/services/crashTaskRun.server.ts | 5 ++- .../v3/services/expireEnqueuedRun.server.ts | 5 ++- .../app/v3/services/triggerTaskV1.server.ts | 1 + 11 files changed, 89 insertions(+), 22 deletions(-) diff --git a/apps/webapp/app/routes/api.v1.runs.$runId.events.ts b/apps/webapp/app/routes/api.v1.runs.$runId.events.ts index ac96c9ddb..92288cd3d 100644 --- a/apps/webapp/app/routes/api.v1.runs.$runId.events.ts +++ b/apps/webapp/app/routes/api.v1.runs.$runId.events.ts @@ -31,7 +31,10 @@ export const loader = createLoaderApiRoute( }, }, async ({ resource: run, authentication }) => { - const eventRepository = resolveEventRepositoryForStore(run.taskEventStore); + const eventRepository = resolveEventRepositoryForStore( + run.taskEventStore, + authentication.environment.organization.id + ); const runEvents = await eventRepository.getRunEvents( getTaskEventStoreTableForRun(run), diff --git a/apps/webapp/app/routes/api.v1.runs.$runId.spans.$spanId.ts b/apps/webapp/app/routes/api.v1.runs.$runId.spans.$spanId.ts index 7c093efd9..5f12a931f 100644 --- a/apps/webapp/app/routes/api.v1.runs.$runId.spans.$spanId.ts +++ b/apps/webapp/app/routes/api.v1.runs.$runId.spans.$spanId.ts @@ -38,7 +38,10 @@ export const loader = createLoaderApiRoute( }, }, async ({ params, resource: run, authentication }) => { - const eventRepository = resolveEventRepositoryForStore(run.taskEventStore); + const eventRepository = resolveEventRepositoryForStore( + run.taskEventStore, + authentication.environment.organization.id + ); const eventStore = getTaskEventStoreTableForRun(run); const span = await eventRepository.getSpan( diff --git a/apps/webapp/app/routes/api.v1.runs.$runId.trace.ts b/apps/webapp/app/routes/api.v1.runs.$runId.trace.ts index cc35836bf..d1964d8e7 100644 --- a/apps/webapp/app/routes/api.v1.runs.$runId.trace.ts +++ b/apps/webapp/app/routes/api.v1.runs.$runId.trace.ts @@ -36,7 +36,10 @@ export const loader = createLoaderApiRoute( }, }, async ({ resource: run, authentication }) => { - const eventRepository = resolveEventRepositoryForStore(run.taskEventStore); + const eventRepository = resolveEventRepositoryForStore( + run.taskEventStore, + authentication.environment.organization.id + ); const traceSummary = await eventRepository.getTraceDetailedSummary( getTaskEventStoreTableForRun(run), diff --git a/apps/webapp/app/routes/resources.runs.$runParam.logs.download.ts b/apps/webapp/app/routes/resources.runs.$runParam.logs.download.ts index f3f21fc15..5da10a8a8 100644 --- a/apps/webapp/app/routes/resources.runs.$runParam.logs.download.ts +++ b/apps/webapp/app/routes/resources.runs.$runParam.logs.download.ts @@ -33,7 +33,10 @@ export async function loader({ params, request }: LoaderFunctionArgs) { return new Response("Not found", { status: 404 }); } - const eventRepository = resolveEventRepositoryForStore(run.taskEventStore); + const eventRepository = resolveEventRepositoryForStore( + run.taskEventStore, + run.organizationId ?? "" + ); const runEvents = await eventRepository.getRunEvents( getTaskEventStoreTableForRun(run), diff --git a/apps/webapp/app/v3/eventRepository/index.server.ts b/apps/webapp/app/v3/eventRepository/index.server.ts index 5c9026572..98d6858ed 100644 --- a/apps/webapp/app/v3/eventRepository/index.server.ts +++ b/apps/webapp/app/v3/eventRepository/index.server.ts @@ -16,6 +16,21 @@ export const EVENT_STORE_TYPES = { export type EventStoreType = (typeof EVENT_STORE_TYPES)[keyof typeof EVENT_STORE_TYPES]; +/** + * Resolve the event repository for a run's persisted `taskEventStore` value and org. + * Postgres-backed runs use the Prisma `eventRepository`; ClickHouse-backed runs use + * `clickhouseFactory.getEventRepositoryForOrganizationSync`. + */ +export function resolveEventRepositoryForStore( + store: string, + organizationId: string +): IEventRepository { + if (store === EVENT_STORE_TYPES.CLICKHOUSE || store === EVENT_STORE_TYPES.CLICKHOUSE_V2) { + return clickhouseFactory.getEventRepositoryForOrganizationSync(store, organizationId).repository; + } + return eventRepository; +} + export async function getConfiguredEventRepository( organizationId: string ): Promise<{ repository: IEventRepository; store: EventStoreType }> { diff --git a/apps/webapp/app/v3/runEngineHandlers.server.ts b/apps/webapp/app/v3/runEngineHandlers.server.ts index 9e69d4ba0..eced6c758 100644 --- a/apps/webapp/app/v3/runEngineHandlers.server.ts +++ b/apps/webapp/app/v3/runEngineHandlers.server.ts @@ -27,7 +27,7 @@ import { PerformTaskRunAlertsService } from "./services/alerts/performTaskRunAle import { TaskRunErrorCodes } from "@trigger.dev/core/v3"; export function registerRunEngineEventBusHandlers() { - engine.eventBus.on("runSucceeded", async ({ time, run }) => { + engine.eventBus.on("runSucceeded", async ({ time, run, organization }) => { const [taskRunError, taskRun] = await tryCatch( $replica.taskRun.findFirstOrThrow({ where: { @@ -60,7 +60,10 @@ export function registerRunEngineEventBusHandlers() { return; } - const eventRepository = resolveEventRepositoryForStore(run.taskEventStore); + const eventRepository = resolveEventRepositoryForStore( + run.taskEventStore, + taskRun.organizationId ?? organization.id + ); const [completeSuccessfulRunEventError] = await tryCatch( eventRepository.completeSuccessfulRunEvent({ @@ -91,7 +94,7 @@ export function registerRunEngineEventBusHandlers() { }); // Handle events - engine.eventBus.on("runFailed", async ({ time, run }) => { + engine.eventBus.on("runFailed", async ({ time, run, organization }) => { const sanitizedError = sanitizeError(run.error); const exception = createExceptionPropertiesFromError(sanitizedError); @@ -127,7 +130,10 @@ export function registerRunEngineEventBusHandlers() { return; } - const eventRepository = resolveEventRepositoryForStore(taskRun.taskEventStore); + const eventRepository = resolveEventRepositoryForStore( + run.taskEventStore, + taskRun.organizationId ?? organization.id + ); const [completeFailedRunEventError] = await tryCatch( eventRepository.completeFailedRunEvent({ @@ -181,7 +187,10 @@ export function registerRunEngineEventBusHandlers() { return; } - const eventRepository = resolveEventRepositoryForStore(taskRun.taskEventStore); + const eventRepository = resolveEventRepositoryForStore( + run.taskEventStore, + taskRun.organizationId ?? "" + ); const [createAttemptFailedRunEventError] = await tryCatch( eventRepository.createAttemptFailedRunEvent({ @@ -282,7 +291,10 @@ export function registerRunEngineEventBusHandlers() { return; } - const eventRepository = resolveEventRepositoryForStore(blockedRun.taskEventStore); + const eventRepository = resolveEventRepositoryForStore( + blockedRun.taskEventStore, + blockedRun.organizationId ?? "" + ); const [completeCachedRunEventError] = await tryCatch( eventRepository.completeCachedRunEvent({ @@ -305,7 +317,7 @@ export function registerRunEngineEventBusHandlers() { } ); - engine.eventBus.on("runExpired", async ({ time, run }) => { + engine.eventBus.on("runExpired", async ({ time, run, organization }) => { if (!run.ttl) { return; } @@ -342,7 +354,10 @@ export function registerRunEngineEventBusHandlers() { return; } - const eventRepository = resolveEventRepositoryForStore(taskRun.taskEventStore); + const eventRepository = resolveEventRepositoryForStore( + taskRun.taskEventStore, + taskRun.organizationId ?? organization.id + ); const [completeExpiredRunEventError] = await tryCatch( eventRepository.completeExpiredRunEvent({ @@ -360,7 +375,7 @@ export function registerRunEngineEventBusHandlers() { } }); - engine.eventBus.on("runCancelled", async ({ time, run }) => { + engine.eventBus.on("runCancelled", async ({ time, run, organization }) => { const [taskRunError, taskRun] = await tryCatch( $replica.taskRun.findFirstOrThrow({ where: { @@ -393,7 +408,10 @@ export function registerRunEngineEventBusHandlers() { return; } - const eventRepository = resolveEventRepositoryForStore(taskRun.taskEventStore); + const eventRepository = resolveEventRepositoryForStore( + taskRun.taskEventStore, + taskRun.organizationId ?? organization.id + ); const error = createJsonErrorObject(run.error); @@ -413,7 +431,7 @@ export function registerRunEngineEventBusHandlers() { } }); - engine.eventBus.on("runRetryScheduled", async ({ time, run, environment, retryAt }) => { + engine.eventBus.on("runRetryScheduled", async ({ time, run, environment, retryAt, organization }) => { try { if (retryAt && time && time >= retryAt) { return; @@ -426,7 +444,10 @@ export function registerRunEngineEventBusHandlers() { retryMessage += ` after OOM`; } - const eventRepository = resolveEventRepositoryForStore(run.taskEventStore); + const eventRepository = resolveEventRepositoryForStore( + run.taskEventStore ?? "taskEvent", + organization.id + ); await eventRepository.recordEvent(retryMessage, { startTime: BigInt(time.getTime() * 1000000), diff --git a/apps/webapp/app/v3/services/cancelTaskRunV1.server.ts b/apps/webapp/app/v3/services/cancelTaskRunV1.server.ts index 4b4482a1d..6a5b5ef8c 100644 --- a/apps/webapp/app/v3/services/cancelTaskRunV1.server.ts +++ b/apps/webapp/app/v3/services/cancelTaskRunV1.server.ts @@ -101,7 +101,10 @@ export class CancelTaskRunServiceV1 extends BaseService { }, }); - const eventRepository = resolveEventRepositoryForStore(cancelledTaskRun.taskEventStore); + const eventRepository = resolveEventRepositoryForStore( + cancelledTaskRun.taskEventStore, + cancelledTaskRun.runtimeEnvironment.organizationId + ); const [cancelRunEventError] = await tryCatch( eventRepository.cancelRunEvent({ diff --git a/apps/webapp/app/v3/services/completeAttempt.server.ts b/apps/webapp/app/v3/services/completeAttempt.server.ts index 99169331a..79647b2f1 100644 --- a/apps/webapp/app/v3/services/completeAttempt.server.ts +++ b/apps/webapp/app/v3/services/completeAttempt.server.ts @@ -163,7 +163,10 @@ export class CompleteAttemptService extends BaseService { env, }); - const eventRepository = resolveEventRepositoryForStore(taskRunAttempt.taskRun.taskEventStore); + const eventRepository = resolveEventRepositoryForStore( + taskRunAttempt.taskRun.taskEventStore, + taskRunAttempt.taskRun.organizationId ?? "" + ); const [completeSuccessfulRunEventError] = await tryCatch( eventRepository.completeSuccessfulRunEvent({ @@ -316,7 +319,10 @@ export class CompleteAttemptService extends BaseService { exitRun(taskRunAttempt.taskRunId); } - const eventRepository = resolveEventRepositoryForStore(taskRunAttempt.taskRun.taskEventStore); + const eventRepository = resolveEventRepositoryForStore( + taskRunAttempt.taskRun.taskEventStore, + taskRunAttempt.taskRun.organizationId ?? "" + ); const [completeFailedRunEventError] = await tryCatch( eventRepository.completeFailedRunEvent({ @@ -538,7 +544,10 @@ export class CompleteAttemptService extends BaseService { }) { const retryAt = new Date(executionRetry.timestamp); - const eventRepository = resolveEventRepositoryForStore(taskRunAttempt.taskRun.taskEventStore); + const eventRepository = resolveEventRepositoryForStore( + taskRunAttempt.taskRun.taskEventStore, + taskRunAttempt.taskRun.organizationId ?? "" + ); // Retry the task run await eventRepository.recordEvent( diff --git a/apps/webapp/app/v3/services/crashTaskRun.server.ts b/apps/webapp/app/v3/services/crashTaskRun.server.ts index 9ed7d8b7a..333a488f7 100644 --- a/apps/webapp/app/v3/services/crashTaskRun.server.ts +++ b/apps/webapp/app/v3/services/crashTaskRun.server.ts @@ -120,7 +120,10 @@ export class CrashTaskRunService extends BaseService { }, }); - const eventRepository = resolveEventRepositoryForStore(crashedTaskRun.taskEventStore); + const eventRepository = resolveEventRepositoryForStore( + crashedTaskRun.taskEventStore, + crashedTaskRun.runtimeEnvironment.organizationId + ); const [createAttemptFailedEventError] = await tryCatch( eventRepository.completeFailedRunEvent({ diff --git a/apps/webapp/app/v3/services/expireEnqueuedRun.server.ts b/apps/webapp/app/v3/services/expireEnqueuedRun.server.ts index aa69f4c9f..3fa1c356c 100644 --- a/apps/webapp/app/v3/services/expireEnqueuedRun.server.ts +++ b/apps/webapp/app/v3/services/expireEnqueuedRun.server.ts @@ -78,7 +78,10 @@ export class ExpireEnqueuedRunService extends BaseService { }, }); - const eventRepository = resolveEventRepositoryForStore(run.taskEventStore); + const eventRepository = resolveEventRepositoryForStore( + run.taskEventStore, + run.runtimeEnvironment.organization.id + ); if (run.ttl) { const [completeExpiredRunEventError] = await tryCatch( diff --git a/apps/webapp/app/v3/services/triggerTaskV1.server.ts b/apps/webapp/app/v3/services/triggerTaskV1.server.ts index 6a927765b..a2f0b9de2 100644 --- a/apps/webapp/app/v3/services/triggerTaskV1.server.ts +++ b/apps/webapp/app/v3/services/triggerTaskV1.server.ts @@ -294,6 +294,7 @@ export class TriggerTaskServiceV1 extends BaseService { : undefined; const { repository, store } = await getV3EventRepository( + environment.organization.id, dependentAttempt?.taskRun.taskEventStore ?? parentAttempt?.taskRun.taskEventStore ?? dependentBatchRun?.dependentTaskAttempt?.taskRun.taskEventStore