Files
triggerdotdev--trigger.dev/apps/webapp/app/services/taskIdentifierRegistry.server.ts
Eric Allam c0b84595a3 feat(webapp): hosted webhook ingress, delivery pipeline, and dashboard (#4344)
## Summary

The server half of hosted webhooks: the public ingress endpoint,
signature verification, the delivery pipeline (Postgres partitioned
storage + ClickHouse for ordering), the in-app partition manager, the
HTTP API, and the dashboard (Deliveries, Endpoints, and the in-app test
console).

The public SDK and docs half is #4537. That PR carries the user-facing
API (`webhook()`, `chat.event` / `chat.channels`, the
`@trigger.dev/slack` connector) and builds on the shared
`@trigger.dev/core` schemas that ship here.

## Shipping behind a flag

A `WEBHOOK_ENABLED` env var (default off) gates the public ingress route
and the engine worker plus partition cron, so merging and deploying this
changes nothing in production until it is flipped on per environment.
The dashboard is separately gated per org by the `hasWebhooksAccess`
feature flag.

## Note on packages

This PR includes the `@trigger.dev/core` schema additions the server
compiles against, but carries no changeset. Core is not consumed
independently of the SDK, so it is released together with the SDK via
#4537. Keeping its changeset off `main` means no release cut from `main`
publishes it early.
2026-08-16 14:33:42 +01:00

164 lines
4.8 KiB
TypeScript

import {
type TaskTriggerSource,
type PrismaClient,
type PrismaClientOrTransaction,
boundedIn,
} from "@trigger.dev/database";
import { $replica, prisma } from "~/db.server";
import { getAllTaskIdentifiers } from "~/models/task.server";
import { logger } from "./logger.server";
import {
getTaskIdentifiersFromCache,
populateTaskIdentifierCache,
type TaskIdentifierEntry,
} from "./taskIdentifierCache.server";
function toTriggerSource(source: string | undefined): TaskTriggerSource {
const normalized = source?.toUpperCase();
if (normalized === "AGENT") return "AGENT";
if (normalized === "WEBHOOK") return "WEBHOOK";
if (normalized === "SCHEDULED" || normalized === "SCHEDULE") return "SCHEDULED";
return "STANDARD";
}
export async function syncTaskIdentifiers(
environmentId: string,
projectId: string,
workerId: string,
tasks: { id: string; triggerSource?: string }[],
db: PrismaClient = prisma
): Promise<void> {
const slugs = tasks.map((t) => t.id);
const now = new Date();
// Group slugs by resolved triggerSource for bulk updates
const slugsBySource = new Map<TaskTriggerSource, string[]>();
for (const task of tasks) {
const source = toTriggerSource(task.triggerSource);
const existing = slugsBySource.get(source);
if (existing) {
existing.push(task.id);
} else {
slugsBySource.set(source, [task.id]);
}
}
// Batch: insert new rows, update existing rows per source group, archive removed tasks
await db.$transaction([
// Insert any new task identifiers (skips rows that already exist)
db.taskIdentifier.createMany({
data: tasks.map((task) => ({
runtimeEnvironmentId: environmentId,
projectId,
slug: task.id,
currentTriggerSource: toTriggerSource(task.triggerSource),
currentWorkerId: workerId,
})),
skipDuplicates: true,
}),
// Update existing rows — one updateMany per distinct triggerSource value
...Array.from(slugsBySource.entries()).map(([source, taskSlugs]) =>
db.taskIdentifier.updateMany({
where: {
runtimeEnvironmentId: environmentId,
slug: { in: boundedIn(taskSlugs) },
},
data: {
currentTriggerSource: source,
currentWorkerId: workerId,
lastSeenAt: now,
isInLatestDeployment: true,
},
})
),
// Archive tasks no longer in this deploy
db.taskIdentifier.updateMany({
where: {
runtimeEnvironmentId: environmentId,
slug: { notIn: boundedIn(slugs) },
isInLatestDeployment: true,
},
data: { isInLatestDeployment: false },
}),
]);
const allIdentifiers = await db.taskIdentifier.findMany({
where: { runtimeEnvironmentId: environmentId },
select: {
slug: true,
currentTriggerSource: true,
isInLatestDeployment: true,
},
});
populateTaskIdentifierCache(
environmentId,
allIdentifiers.map((t) => ({
slug: t.slug,
triggerSource: t.currentTriggerSource,
isInLatestDeployment: t.isInLatestDeployment,
}))
).catch((error) => {
logger.error("Failed to populate task identifier cache after sync", { environmentId, error });
});
}
function sortEntries(entries: TaskIdentifierEntry[]): TaskIdentifierEntry[] {
return entries.sort((a, b) => {
if (a.isInLatestDeployment !== b.isInLatestDeployment) return a.isInLatestDeployment ? -1 : 1;
return a.slug.localeCompare(b.slug);
});
}
export async function getTaskIdentifiers(
environmentId: string,
db: PrismaClientOrTransaction = $replica
): Promise<TaskIdentifierEntry[]> {
const cached = await getTaskIdentifiersFromCache(environmentId);
if (cached) return sortEntries(cached);
const dbRows = await db.taskIdentifier.findMany({
where: { runtimeEnvironmentId: environmentId },
select: {
slug: true,
currentTriggerSource: true,
isInLatestDeployment: true,
},
});
if (dbRows.length > 0) {
const entries: TaskIdentifierEntry[] = dbRows.map((t) => ({
slug: t.slug,
triggerSource: t.currentTriggerSource,
isInLatestDeployment: t.isInLatestDeployment,
}));
populateTaskIdentifierCache(environmentId, entries).catch((error) => {
logger.error("Failed to populate task identifier cache after DB read", {
environmentId,
error,
});
});
return sortEntries(entries);
}
const legacyRows = await getAllTaskIdentifiers(db, environmentId);
const entries: TaskIdentifierEntry[] = legacyRows.map((t) => ({
slug: t.slug,
triggerSource: t.triggerSource,
isInLatestDeployment: true,
}));
if (entries.length > 0) {
populateTaskIdentifierCache(environmentId, entries).catch((error) => {
logger.error("Failed to populate task identifier cache after legacy fallback", {
environmentId,
error,
});
});
}
return sortEntries(entries);
}