diff --git a/api/platform/src/internal.ts b/api/platform/src/internal.ts index 24f0cfaebe..649112e4a1 100644 --- a/api/platform/src/internal.ts +++ b/api/platform/src/internal.ts @@ -3,26 +3,17 @@ * * NOT served over HTTP. Accessed exclusively via createInternalCaller(). * All procedures use internalProcedure (observability middleware, no auth). - * - * Sub-routers will be added as business logic is migrated from: - * - Inngest function steps (DB calls, provider APIs, token management) - * - Route handler lib calls (OAuth, webhook ingestion) */ -import { createTRPCRouter, internalProcedure } from "./trpc"; + +import { oauthInternalRouter } from "./router/internal/oauth"; +import { webhooksInternalRouter } from "./router/internal/webhooks"; +import { createTRPCRouter } from "./trpc"; // -- Internal Router ---------------------------------------------------------- export const internalRouter = createTRPCRouter({ - /** - * Proof-of-concept procedure. - * Validates the full chain: caller -> router -> procedure -> middleware -> response. - * Remove once real sub-routers are added. - */ - ping: internalProcedure.query(({ ctx }) => ({ - ok: true as const, - timestamp: new Date().toISOString(), - source: ctx.auth.type === "internal" ? ctx.auth.source : "unknown", - })), + webhooks: webhooksInternalRouter, + oauth: oauthInternalRouter, }); export type InternalRouter = typeof internalRouter; diff --git a/api/platform/src/lib/oauth/callback.ts b/api/platform/src/lib/oauth/callback.ts index ac7e6002bc..11bbcf4192 100644 --- a/api/platform/src/lib/oauth/callback.ts +++ b/api/platform/src/lib/oauth/callback.ts @@ -10,6 +10,7 @@ import { gatewayInstallations } from "@db/app/schema"; import type { SourceType } from "@repo/app-providers"; import { getProvider, providerAccountInfoSchema } from "@repo/app-providers"; import { and, eq } from "@vendor/db"; +import { parseError } from "@vendor/observability/error/next"; import { log } from "@vendor/observability/log/next"; import { providerConfigs } from "../provider-configs"; import { writeTokenRecord } from "../token-store"; @@ -348,7 +349,7 @@ export async function processOAuthCallback( reactivated, }); } catch (err) { - const message = err instanceof Error ? err.message : "unknown"; + const message = parseError(err); log.error("[oauth/callback] oauth callback failed", { provider: providerName, error: message, diff --git a/api/platform/src/router/internal/oauth.ts b/api/platform/src/router/internal/oauth.ts new file mode 100644 index 0000000000..705c03829d --- /dev/null +++ b/api/platform/src/router/internal/oauth.ts @@ -0,0 +1,94 @@ +/** + * Internal OAuth sub-router. + * + * Handles OAuth authorize URL generation, callback processing, and CLI polling. + * Moved from apps/platform/src/app/api/connect/ route handlers. + */ + +import type { SourceType } from "@repo/app-providers"; +import type { TRPCRouterRecord } from "@trpc/server"; +import { TRPCError } from "@trpc/server"; +import { z } from "zod"; +import { buildAuthorizeUrl } from "../../lib/oauth/authorize"; +import { + type CallbackProcessResult, + processOAuthCallback, +} from "../../lib/oauth/callback"; +import { getOAuthResult } from "../../lib/oauth/state"; +import { internalProcedure } from "../../trpc"; + +// ── Router ────────────────────────────────────────────────────────────────── + +export const oauthInternalRouter = { + /** + * Build OAuth authorize URL for a provider. + * + * Generates a cryptographically random state token, stores it in Redis, + * and returns the authorization URL. + */ + buildAuthorizeUrl: internalProcedure + .input( + z.object({ + provider: z.string(), + orgId: z.string(), + connectedBy: z.string(), + redirectTo: z.string().optional(), + }) + ) + .mutation(async ({ input }) => { + const result = await buildAuthorizeUrl({ + provider: input.provider as SourceType, + orgId: input.orgId, + connectedBy: input.connectedBy, + redirectTo: input.redirectTo, + }); + + if (!result.ok) { + throw new TRPCError({ + code: "BAD_REQUEST", + message: result.error, + }); + } + + return { url: result.url, state: result.state }; + }), + + /** + * Process OAuth callback: validate state, exchange code, upsert installation, + * persist tokens, store result for CLI polling. + * + * Returns CallbackProcessResult — the route handler maps this to HTTP responses + * (redirect, inline HTML, or error JSON). + */ + processCallback: internalProcedure + .input( + z.object({ + provider: z.string(), + state: z.string(), + query: z.record(z.string(), z.string()), + }) + ) + .mutation(async ({ input }): Promise => { + return processOAuthCallback({ + provider: input.provider as SourceType, + state: input.state, + query: input.query as Record, + }); + }), + + /** + * Poll for OAuth completion result. + * + * Returns the result hash from Redis if the OAuth flow has completed, + * or null if still pending. + */ + pollResult: internalProcedure + .input( + z.object({ + state: z.string(), + }) + ) + .query(async ({ input }) => { + return getOAuthResult(input.state); + }), +} satisfies TRPCRouterRecord; diff --git a/api/platform/src/router/internal/webhooks.ts b/api/platform/src/router/internal/webhooks.ts new file mode 100644 index 0000000000..701d2aa00e --- /dev/null +++ b/api/platform/src/router/internal/webhooks.ts @@ -0,0 +1,210 @@ +/** + * Internal webhooks sub-router. + * + * Handles webhook ingestion: HMAC verification, DB persistence, Inngest dispatch. + * Moved from apps/platform/src/app/api/ingest/[provider]/route.ts. + */ + +import { db } from "@db/app/client"; +import { gatewayWebhookDeliveries } from "@db/app/schema"; +import type { WebhookDef } from "@repo/app-providers"; +import { + deriveVerifySignature, + getProvider, + hasInboundWebhooks, + isWebhookProvider, +} from "@repo/app-providers"; +import type { TRPCRouterRecord } from "@trpc/server"; +import { TRPCError } from "@trpc/server"; +import { log } from "@vendor/observability/log/next"; +import { z } from "zod"; +import { inngest } from "../../inngest/client"; +import { getProviderConfigs } from "../../lib/provider-configs"; +import { internalProcedure } from "../../trpc"; + +// ── Helpers (moved from route handler) ────────────────────────────────────── + +function getWebhookDef( + providerDef: NonNullable> +): WebhookDef | null { + if (isWebhookProvider(providerDef)) { + return providerDef.webhook as WebhookDef; + } + if (providerDef.kind === "managed") { + return providerDef.inbound.webhook as WebhookDef; + } + if ( + providerDef.kind === "api" && + "inbound" in providerDef && + providerDef.inbound + ) { + return providerDef.inbound.webhook as WebhookDef; + } + return null; +} + +// ── Router ────────────────────────────────────────────────────────────────── + +export const webhooksInternalRouter = { + /** + * Ingest a webhook delivery: verify HMAC, persist to DB, dispatch to Inngest. + * + * The route handler extracts rawBody and headers from the HTTP request + * and passes them here. This procedure handles everything else. + * + * Returns the response shape for the route handler to forward as JSON. + */ + ingest: internalProcedure + .input( + z.object({ + provider: z.string(), + rawBody: z.string(), + headers: z.record(z.string(), z.string()), + receivedAt: z.number(), + }) + ) + .mutation(async ({ input }) => { + const { provider: providerSlug, rawBody, headers, receivedAt } = input; + + // Provider guard + const providerDef = getProvider(providerSlug); + if (!providerDef) { + throw new TRPCError({ + code: "NOT_FOUND", + message: `unknown_provider: ${providerSlug}`, + }); + } + + if (!hasInboundWebhooks(providerDef)) { + throw new TRPCError({ + code: "BAD_REQUEST", + message: `not_webhook_provider: ${providerSlug}`, + }); + } + + const webhookDef = getWebhookDef(providerDef); + if (!webhookDef) { + throw new TRPCError({ + code: "BAD_REQUEST", + message: `no_webhook_def: ${providerSlug}`, + }); + } + + // Webhook header validation + const headersObj: Record = {}; + for (const key of Object.keys( + (webhookDef.headersSchema as { shape: Record }).shape + )) { + headersObj[key] = (headers as Record)[key] ?? undefined; + } + const headersParsed = webhookDef.headersSchema.safeParse(headersObj); + if (!headersParsed.success) { + throw new TRPCError({ + code: "BAD_REQUEST", + message: "missing_headers", + }); + } + + // Signature verification + const configs = getProviderConfigs(); + const providerConfig = configs[providerSlug]; + if (!providerConfig) { + log.error("[webhooks.ingest] provider config not found", { + provider: providerSlug, + }); + throw new TRPCError({ + code: "INTERNAL_SERVER_ERROR", + message: `provider_not_configured: ${providerSlug}`, + }); + } + + const secret = (webhookDef.extractSecret as (config: unknown) => string)( + providerConfig + ); + const verify = + webhookDef.verifySignature ?? + deriveVerifySignature(webhookDef.signatureScheme); + + // Build a Headers object for the verify function (it expects Headers, not Record) + const reqHeaders = new Headers(Object.entries(headers)); + const isValid = await verify(rawBody, reqHeaders, secret); + if (!isValid) { + log.warn("[webhooks.ingest] signature verification failed", { + provider: providerSlug, + }); + throw new TRPCError({ + code: "UNAUTHORIZED", + message: "signature_invalid", + }); + } + + // Payload parse + metadata extraction + let jsonPayload: unknown; + try { + jsonPayload = JSON.parse(rawBody); + } catch { + throw new TRPCError({ + code: "BAD_REQUEST", + message: "invalid_json", + }); + } + + let parsedPayload: unknown; + try { + parsedPayload = webhookDef.parsePayload(jsonPayload); + } catch { + throw new TRPCError({ + code: "BAD_REQUEST", + message: `payload_validation_failed: ${providerSlug}`, + }); + } + + const deliveryId = webhookDef.extractDeliveryId( + reqHeaders, + parsedPayload + ); + const eventType = webhookDef.extractEventType(reqHeaders, parsedPayload); + const resourceId = webhookDef.extractResourceId(parsedPayload); + + // Persist to DB + await db + .insert(gatewayWebhookDeliveries) + .values({ + provider: providerSlug, + deliveryId, + eventType, + installationId: null, + status: "received", + payload: JSON.stringify(parsedPayload), + receivedAt: new Date(receivedAt).toISOString(), + }) + .onConflictDoNothing(); + + // Dispatch Inngest event + const correlationId = crypto.randomUUID(); + + await inngest.send({ + id: `wh-${providerSlug}-${deliveryId}`, + name: "platform/webhook.received", + data: { + provider: providerSlug, + deliveryId, + eventType, + resourceId, + payload: parsedPayload, + receivedAt, + correlationId, + }, + }); + + log.info("[webhooks.ingest] webhook received", { + provider: providerSlug, + deliveryId, + eventType, + resourceId, + correlationId, + }); + + return { status: "accepted" as const, deliveryId }; + }), +} satisfies TRPCRouterRecord; diff --git a/apps/platform/src/app/api/connect/[provider]/authorize/route.ts b/apps/platform/src/app/api/connect/[provider]/authorize/route.ts index b10ad5b287..442df40290 100644 --- a/apps/platform/src/app/api/connect/[provider]/authorize/route.ts +++ b/apps/platform/src/app/api/connect/[provider]/authorize/route.ts @@ -1,17 +1,14 @@ /** * GET /api/connect/:provider/authorize * - * Initiate OAuth flow. Ported from apps/gateway/src/routes/connections.ts (lines 79-141). - * - * NOT tRPC — returns authorize URL + state for browser OAuth. - * The main flow uses tRPC `connections.getAuthorizeUrl` instead; - * this route supports direct browser navigation as a fallback. + * Initiate OAuth flow. Returns authorize URL + state for browser OAuth. + * All business logic lives in platform.oauth.buildAuthorizeUrl(). */ -import { buildAuthorizeUrl } from "@api/platform/lib/oauth/authorize"; -import type { SourceType } from "@repo/app-providers"; +import { TRPCError } from "@trpc/server"; import { log } from "@vendor/observability/log/next"; import type { NextRequest } from "next/server"; +import { platform } from "~/lib/internal-caller"; export const runtime = "nodejs"; @@ -20,7 +17,6 @@ export async function GET( { params }: { params: Promise<{ provider: string }> } ) { const { provider } = await params; - const providerName = provider as SourceType; const orgId = req.nextUrl.searchParams.get("org_id"); const connectedBy = @@ -30,25 +26,24 @@ export async function GET( const redirectTo = req.nextUrl.searchParams.get("redirect_to") ?? undefined; if (!orgId) { - log.warn("[oauth/authorize] missing org_id", { provider: providerName }); + log.warn("[oauth/authorize] missing org_id", { provider }); return Response.json({ error: "missing_org_id" }, { status: 400 }); } - const result = await buildAuthorizeUrl({ - provider: providerName, - orgId, - connectedBy, - redirectTo, - }); - - if (!result.ok) { - log.warn("[oauth/authorize] failed to build authorize URL", { - provider: providerName, - error: result.error, + try { + const result = await platform.oauth.buildAuthorizeUrl({ + provider, + orgId, + connectedBy, + redirectTo, }); - return Response.json({ error: result.error }, { status: 400 }); - } - log.info("[oauth/authorize] authorize URL built", { provider: providerName }); - return Response.json({ url: result.url, state: result.state }); + log.info("[oauth/authorize] authorize URL built", { provider }); + return Response.json(result); + } catch (err) { + if (err instanceof TRPCError) { + return Response.json({ error: err.message }, { status: 400 }); + } + throw err; + } } diff --git a/apps/platform/src/app/api/connect/[provider]/callback/route.ts b/apps/platform/src/app/api/connect/[provider]/callback/route.ts index a029121f85..9dec2abb15 100644 --- a/apps/platform/src/app/api/connect/[provider]/callback/route.ts +++ b/apps/platform/src/app/api/connect/[provider]/callback/route.ts @@ -1,19 +1,14 @@ /** * GET /api/connect/:provider/callback * - * OAuth callback. Ported from apps/gateway/src/routes/connections.ts (lines 208-489). - * - * NOT tRPC — OAuth provider redirects here directly (browser). - * Maps CallbackProcessResult from the lib layer to HTTP responses. + * OAuth callback. Provider redirects here after authorization. + * All business logic lives in platform.oauth.processCallback(). + * Maps CallbackProcessResult to HTTP responses. */ -import { - type CallbackProcessResult, - processOAuthCallback, -} from "@api/platform/lib/oauth/callback"; -import type { SourceType } from "@repo/app-providers"; import { log } from "@vendor/observability/log/next"; import { type NextRequest, NextResponse } from "next/server"; +import { platform } from "~/lib/internal-caller"; export const runtime = "nodejs"; @@ -22,9 +17,7 @@ export async function GET( { params }: { params: Promise<{ provider: string }> } ) { const { provider } = await params; - const providerName = provider as SourceType; - // Build query dict from all URL search params (null-prototype to prevent prototype pollution) const query: Record = Object.assign( Object.create(null) as Record, Object.fromEntries(req.nextUrl.searchParams) @@ -32,20 +25,20 @@ export async function GET( const state = query.state ?? ""; - const result: CallbackProcessResult = await processOAuthCallback({ - provider: providerName, + const result = await platform.oauth.processCallback({ + provider, state, query, }); switch (result.kind) { case "redirect": - log.info("[oauth/callback] redirecting", { provider: providerName }); + log.info("[oauth/callback] redirecting", { provider }); return NextResponse.redirect(result.url); case "inline_html": log.info("[oauth/callback] inline html response", { - provider: providerName, + provider, status: result.status ?? 200, }); return new Response(result.html, { @@ -55,7 +48,7 @@ export async function GET( case "error": log.warn("[oauth/callback] error result", { - provider: providerName, + provider, error: result.error, status: result.status, }); diff --git a/apps/platform/src/app/api/connect/oauth/poll/route.ts b/apps/platform/src/app/api/connect/oauth/poll/route.ts index 80dc5e1489..2bdb000f3c 100644 --- a/apps/platform/src/app/api/connect/oauth/poll/route.ts +++ b/apps/platform/src/app/api/connect/oauth/poll/route.ts @@ -1,16 +1,13 @@ /** * GET /api/connect/oauth/poll * - * Poll for OAuth completion. Ported from gateway connections.ts (lines 180-196). - * - * NOT tRPC — CLI polling with state token as auth. - * The state token itself is the secret (cryptographically random nanoid, - * known only to the initiator). + * Poll for OAuth completion. CLI polling with state token as auth. + * All business logic lives in platform.oauth.pollResult(). */ -import { getOAuthResult } from "@api/platform/lib/oauth/state"; import { log } from "@vendor/observability/log/next"; import type { NextRequest } from "next/server"; +import { platform } from "~/lib/internal-caller"; export const runtime = "nodejs"; @@ -22,7 +19,7 @@ export async function GET(req: NextRequest) { return Response.json({ error: "missing_state" }, { status: 400 }); } - const result = await getOAuthResult(state); + const result = await platform.oauth.pollResult({ state }); if (!result) { return Response.json({ status: "pending" }); diff --git a/apps/platform/src/app/api/ingest/[provider]/route.ts b/apps/platform/src/app/api/ingest/[provider]/route.ts index b7a13e6069..3f950c9406 100644 --- a/apps/platform/src/app/api/ingest/[provider]/route.ts +++ b/apps/platform/src/app/api/ingest/[provider]/route.ts @@ -5,201 +5,49 @@ * Validates HMAC signatures and dispatches to the Inngest pipeline. * * NOT tRPC — external providers send raw HTTP with HMAC signatures. + * All business logic lives in platform.webhooks.ingest(). */ -import { inngest } from "@api/platform/inngest/client"; -import { getProviderConfigs } from "@api/platform/lib/provider-configs"; -import { db } from "@db/app/client"; -import { gatewayWebhookDeliveries } from "@db/app/schema"; -import type { WebhookDef } from "@repo/app-providers"; -import { - deriveVerifySignature, - getProvider, - hasInboundWebhooks, - isWebhookProvider, -} from "@repo/app-providers"; -import { parseError } from "@vendor/observability/error/next"; -import { log } from "@vendor/observability/log/next"; +import { TRPCError } from "@trpc/server"; import type { NextRequest } from "next/server"; +import { platform } from "~/lib/internal-caller"; export const runtime = "nodejs"; -// ── Helpers ────────────────────────────────────────────────────────────────── - -/** - * Resolve WebhookDef from any provider kind that supports inbound webhooks. - */ -function getWebhookDef( - providerDef: NonNullable> -): WebhookDef | null { - if (isWebhookProvider(providerDef)) { - return providerDef.webhook as WebhookDef; - } - if (providerDef.kind === "managed") { - return providerDef.inbound.webhook as WebhookDef; - } - if ( - providerDef.kind === "api" && - "inbound" in providerDef && - providerDef.inbound - ) { - return providerDef.inbound.webhook as WebhookDef; - } - return null; -} - -// ── Route Handler ──────────────────────────────────────────────────────────── - export async function POST( req: NextRequest, { params }: { params: Promise<{ provider: string }> } ) { - const { provider: providerSlug } = await params; + const { provider } = await params; const receivedAt = Date.now(); - - // Step 1: Provider guard - const providerDef = getProvider(providerSlug); - if (!providerDef) { - return Response.json( - { error: "unknown_provider", provider: providerSlug }, - { status: 404 } - ); - } - - if (!hasInboundWebhooks(providerDef)) { - return Response.json( - { error: "not_webhook_provider", provider: providerSlug }, - { status: 400 } - ); - } - - return handleStandardWebhook(req, providerSlug, providerDef, receivedAt); -} - -// ── Standard Webhook Path ──────────────────────────────────────────────────── - -async function handleStandardWebhook( - req: NextRequest, - providerSlug: string, - providerDef: NonNullable>, - receivedAt: number -): Promise { - const webhookDef = getWebhookDef(providerDef); - if (!webhookDef) { - return Response.json( - { error: "no_webhook_def", provider: providerSlug }, - { status: 400 } - ); - } - - // Step 4: Webhook header guard - const headersObj: Record = {}; - for (const key of Object.keys( - (webhookDef.headersSchema as { shape: Record }).shape - )) { - headersObj[key] = req.headers.get(key) ?? undefined; - } - const headersParsed = webhookDef.headersSchema.safeParse(headersObj); - if (!headersParsed.success) { - return Response.json( - { error: "missing_headers", issues: headersParsed.error.issues }, - { status: 400 } - ); - } - - // Step 5: Raw body capture (MUST be req.text() — HMAC needs raw bytes) const rawBody = await req.text(); - // Step 6: Signature verification - const configs = getProviderConfigs(); - const providerConfig = configs[providerSlug]; - if (!providerConfig) { - log.error("[ingest] provider config not found", { - provider: providerSlug, - }); - return Response.json( - { error: "provider_not_configured", provider: providerSlug }, - { status: 500 } - ); - } - - const secret = (webhookDef.extractSecret as (config: unknown) => string)( - providerConfig - ); - const verify = - webhookDef.verifySignature ?? - deriveVerifySignature(webhookDef.signatureScheme); - const isValid = await verify(rawBody, req.headers, secret); - if (!isValid) { - log.warn("[ingest] signature verification failed", { - provider: providerSlug, - }); - return Response.json({ error: "signature_invalid" }, { status: 401 }); - } + // Collect headers as a plain Record for the tRPC procedure + const headers: Record = {}; + req.headers.forEach((value, key) => { + headers[key] = value; + }); - // Step 7: Payload parse + metadata extraction - let jsonPayload: unknown; try { - jsonPayload = JSON.parse(rawBody); - } catch { - return Response.json({ error: "invalid_json" }, { status: 400 }); - } + const result = await platform.webhooks.ingest({ + provider, + rawBody, + headers, + receivedAt, + }); - let parsedPayload: unknown; - try { - parsedPayload = webhookDef.parsePayload(jsonPayload); + return Response.json(result, { status: 202 }); } catch (err) { - log.warn("[ingest] payload schema validation failed", { - provider: providerSlug, - error: parseError(err), - }); - return Response.json( - { error: "payload_validation_failed", provider: providerSlug }, - { status: 400 } - ); + if (err instanceof TRPCError) { + const statusMap: Record = { + NOT_FOUND: 404, + BAD_REQUEST: 400, + UNAUTHORIZED: 401, + INTERNAL_SERVER_ERROR: 500, + }; + const status = statusMap[err.code] ?? 500; + return Response.json({ error: err.message }, { status }); + } + throw err; } - const deliveryId = webhookDef.extractDeliveryId(req.headers, parsedPayload); - const eventType = webhookDef.extractEventType(req.headers, parsedPayload); - const resourceId = webhookDef.extractResourceId(parsedPayload); - - // Persist to DB - await db - .insert(gatewayWebhookDeliveries) - .values({ - provider: providerSlug, - deliveryId, - eventType, - installationId: null, - status: "received", - payload: JSON.stringify(parsedPayload), - receivedAt: new Date(receivedAt).toISOString(), - }) - .onConflictDoNothing(); - - // Dispatch Inngest event - const correlationId = crypto.randomUUID(); - - await inngest.send({ - id: `wh-${providerSlug}-${deliveryId}`, - name: "platform/webhook.received", - data: { - provider: providerSlug, - deliveryId, - eventType, - resourceId, - payload: parsedPayload, - receivedAt, - correlationId, - }, - }); - - log.info("[ingest] webhook received", { - provider: providerSlug, - deliveryId, - eventType, - resourceId, - correlationId, - }); - - return Response.json({ status: "accepted", deliveryId }); }