diff --git a/apps/mobile/src/features/settings/SettingsScheduledTasksRouteScreen.tsx b/apps/mobile/src/features/settings/SettingsScheduledTasksRouteScreen.tsx index a3219888241a..2effae896fc6 100644 --- a/apps/mobile/src/features/settings/SettingsScheduledTasksRouteScreen.tsx +++ b/apps/mobile/src/features/settings/SettingsScheduledTasksRouteScreen.tsx @@ -87,6 +87,7 @@ const DAYS = [ ] as const; function describeSchedule(task: ScheduledTask): string { + if (task.schedule.type === "webhook") return "On webhook"; if (task.schedule.type === "interval") return formatScheduledTaskInterval(task.schedule.everyMs); const days = task.schedule.weekdays?.length ? repeatLabel(task.schedule.weekdays) : "Every day"; return `${days} at ${formatTime(task.schedule.timeOfDay)}`; @@ -1024,6 +1025,8 @@ function EnvironmentTasks({ { onEdit(task); }} @@ -1048,7 +1051,7 @@ function EnvironmentTasks({ ): ScheduleDraft { + // Webhook tasks cannot be edited here yet; show them as the default schedule. + if (task.schedule.type === "webhook") return DEFAULT_SCHEDULE; return task.schedule.type === "fixed_time" ? { ...DEFAULT_SCHEDULE, diff --git a/apps/server/src/auth/RpcAuthorization.test.ts b/apps/server/src/auth/RpcAuthorization.test.ts index 1af8fd7c6214..1afb2fb75e3f 100644 --- a/apps/server/src/auth/RpcAuthorization.test.ts +++ b/apps/server/src/auth/RpcAuthorization.test.ts @@ -38,6 +38,16 @@ describe("RPC authorization scopes", () => { ); }); + it("keeps webhook delivery logs, which hold request bodies, behind operate scope", () => { + for (const method of [ + WS_METHODS.scheduledTasksListWebhookDeliveries, + WS_METHODS.scheduledTasksGetWebhookDelivery, + WS_METHODS.scheduledTasksRotateWebhookToken, + ]) { + expect(requiredScopeForRpcMethod(method)).toBe(AuthOrchestrationOperateScope); + } + }); + it("allows relay status reads without granting relay installation access", () => { expect(requiredScopeForRpcMethod(WS_METHODS.cloudGetRelayClientStatus)).toBe( AuthRelayReadScope, diff --git a/apps/server/src/auth/RpcAuthorization.ts b/apps/server/src/auth/RpcAuthorization.ts index fb4a43185ef0..61d7ea279bf4 100644 --- a/apps/server/src/auth/RpcAuthorization.ts +++ b/apps/server/src/auth/RpcAuthorization.ts @@ -95,6 +95,10 @@ export const RPC_REQUIRED_SCOPES = { [WS_METHODS.scheduledTasksSetEnabled]: AuthOrchestrationOperateScope, [WS_METHODS.scheduledTasksDelete]: AuthOrchestrationOperateScope, [WS_METHODS.scheduledTasksRunNow]: AuthOrchestrationOperateScope, + [WS_METHODS.scheduledTasksRotateWebhookToken]: AuthOrchestrationOperateScope, + // Delivery logs hold request bodies, so they need the same scope as the URL. + [WS_METHODS.scheduledTasksListWebhookDeliveries]: AuthOrchestrationOperateScope, + [WS_METHODS.scheduledTasksGetWebhookDelivery]: AuthOrchestrationOperateScope, [WS_METHODS.cloudGetRelayClientStatus]: AuthRelayReadScope, [WS_METHODS.cloudInstallRelayClient]: AuthRelayWriteScope, [WS_METHODS.pullRequestsList]: AuthOrchestrationReadScope, diff --git a/apps/server/src/http.ts b/apps/server/src/http.ts index 7be9c6adad0b..535654fd20b4 100644 --- a/apps/server/src/http.ts +++ b/apps/server/src/http.ts @@ -46,6 +46,7 @@ import { failEnvironmentInternal, } from "./auth/http.ts"; import * as ServerEnvironment from "./environment/ServerEnvironment.ts"; +import { WEBHOOK_ROUTE_PREFIX } from "./scheduledTasks/ScheduledTaskService.ts"; import { browserApiCorsAllowedHeaders, browserApiCorsAllowedMethods } from "./httpCors.ts"; const OTLP_TRACES_PROXY_PATH = "/api/observability/v1/traces"; @@ -381,9 +382,9 @@ const UNTRACED_REQUEST_PATHS: ReadonlySet = new Set([OTLP_TRACES_PROXY_P // ignored, as in routing. const untracedRequestsLayer = Layer.succeed(HttpMiddleware.TracerDisabledWhen)((request) => { const queryIndex = request.url.indexOf("?"); - return UNTRACED_REQUEST_PATHS.has( - queryIndex === -1 ? request.url : request.url.slice(0, queryIndex), - ); + const path = queryIndex === -1 ? request.url : request.url.slice(0, queryIndex); + // Webhook URLs carry their secret token in the path, so they never reach a trace. + return UNTRACED_REQUEST_PATHS.has(path) || path.startsWith(`${WEBHOOK_ROUTE_PREFIX}/`); }); export const withUntracedRequests = Layer.provide(untracedRequestsLayer); diff --git a/apps/server/src/mcp/OrchestratorMcpService.ts b/apps/server/src/mcp/OrchestratorMcpService.ts index 9287dda35642..d8175f18d62d 100644 --- a/apps/server/src/mcp/OrchestratorMcpService.ts +++ b/apps/server/src/mcp/OrchestratorMcpService.ts @@ -212,6 +212,7 @@ function scheduledTaskSummary(task: ScheduledTask): OrchestratorMcpScheduledTask schedule: task.schedule, nextRunAt: task.nextRunAt, lastRunStatus: task.lastRunStatus, + ...(task.webhook === undefined ? {} : { webhookUrl: task.webhook.url ?? task.webhook.path }), }; } diff --git a/apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts b/apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts index cd4a0df423a7..a847c8d081f5 100644 --- a/apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts +++ b/apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts @@ -443,7 +443,8 @@ function scheduledTaskFromUpsert(input: ScheduledTaskUpsertInput): ScheduledTask title: input.title, prompt: input.prompt, enabled: input.enabled, - schedule: input.schedule, + schedule: + input.schedule.type === "webhook" ? { type: "webhook", signature: null } : input.schedule, projectId: input.projectId, threadId: input.threadId ?? null, workspaceStrategy: input.workspaceStrategy, @@ -471,6 +472,10 @@ const unusedScheduledTaskStubLayer = Layer.succeed( setEnabled: () => Effect.die("ScheduledTaskService.setEnabled is unused in this test"), delete: () => Effect.die("ScheduledTaskService.delete is unused in this test"), runNow: () => Effect.die("ScheduledTaskService.runNow is unused in this test"), + rotateWebhookToken: () => Effect.die("unused in this test"), + listWebhookDeliveries: () => Effect.die("unused in this test"), + getWebhookDelivery: () => Effect.die("unused in this test"), + triggerWebhook: () => Effect.die("unused in this test"), }), ); @@ -628,6 +633,10 @@ describe("orchestrator MCP toolkit", () => { all.filter((candidate) => candidate.id !== input.id), ).pipe(Effect.as({ id: input.id })), runNow: () => Effect.die("ScheduledTaskService.runNow is unused in this test"), + rotateWebhookToken: () => Effect.die("unused in this test"), + listWebhookDeliveries: () => Effect.die("unused in this test"), + getWebhookDelivery: () => Effect.die("unused in this test"), + triggerWebhook: () => Effect.die("unused in this test"), }), ); const testLayer = Layer.merge( diff --git a/apps/server/src/mcp/toolkits/orchestrator/tools.ts b/apps/server/src/mcp/toolkits/orchestrator/tools.ts index 80fec3103162..a4a4ba728fec 100644 --- a/apps/server/src/mcp/toolkits/orchestrator/tools.ts +++ b/apps/server/src/mcp/toolkits/orchestrator/tools.ts @@ -97,7 +97,7 @@ const TaskCancelTool = Tool.make("task_cancel", { export const ScheduleTaskTool = Tool.make("schedule_task", { description: - "Create persistent recurring work in the app scheduler, which runs even when no turn is active. Pass schedule as a STRUCTURED OBJECT, never JSON text: {type:'interval', everyMs:3600000} means hourly; {type:'fixed_time', timeOfDay:'09:00', weekdays:[1,2,3,4,5]} means weekday mornings. Omit projectId for this thread's project. In this thread's project, runs post into THIS thread by default (bindToCurrentThread=true); use false only when the user wants a fresh top-level thread per run. Elsewhere each run launches a fresh thread. Provider, model, and runtime settings inherit from this thread, or from the project default when there is no calling thread. Report the returned schedule and nextRunAt after success.", + "Create persistent recurring work in the app scheduler, which runs even when no turn is active. Pass schedule as a STRUCTURED OBJECT, never JSON text: {type:'interval', everyMs:3600000} means hourly; {type:'fixed_time', timeOfDay:'09:00', weekdays:[1,2,3,4,5]} means weekday mornings; {type:'webhook'} runs on each request to a generated URL (returned as webhookUrl), and its prompt may use {{body.path}}, {{headers.name}}, {{query.name}}, {{body}} or {{request}} placeholders, which are the only request data the run sees. Omit projectId for this thread's project. In this thread's project, runs post into THIS thread by default (bindToCurrentThread=true); use false only when the user wants a fresh top-level thread per run. Elsewhere each run launches a fresh thread. Provider, model, and runtime settings inherit from this thread, or from the project default when there is no calling thread. Report the returned schedule and nextRunAt after success.", parameters: OrchestratorMcpScheduleTaskInput, success: OrchestratorMcpScheduleTaskResult, failure: OrchestratorMcpFailure, diff --git a/apps/server/src/persistence/Migrations.ts b/apps/server/src/persistence/Migrations.ts index 561aca8bb5ef..1067c2416085 100644 --- a/apps/server/src/persistence/Migrations.ts +++ b/apps/server/src/persistence/Migrations.ts @@ -70,6 +70,7 @@ import Migration0053 from "./Migrations/053_PullRequestFilesViewed.ts"; import Migration0054 from "./Migrations/054_ProjectionThreadsAutoSettleDisabledAt.ts"; import Migration0055 from "./Migrations/055_OrchestrationV2.ts"; import Migration0056 from "./Migrations/056_RemoveRedundantProjectionIndexes.ts"; +import Migration0057 from "./Migrations/057_ScheduledTaskWebhooks.ts"; /** * Migration loader with all migrations defined inline. @@ -140,6 +141,7 @@ export const migrationEntries = [ // Preserve this migration's schema. Future V2 schema changes need new migrations. [55, "OrchestrationV2", Migration0055], [56, "RemoveRedundantProjectionIndexes", Migration0056], + [57, "ScheduledTaskWebhooks", Migration0057], ] as const; export const migrationManifest = migrationEntries.map(([id, name]) => [id, name] as const); diff --git a/apps/server/src/persistence/Migrations/055_OrchestrationV2.test.ts b/apps/server/src/persistence/Migrations/055_OrchestrationV2.test.ts index 71a323f99ddf..e3839e3f0d17 100644 --- a/apps/server/src/persistence/Migrations/055_OrchestrationV2.test.ts +++ b/apps/server/src/persistence/Migrations/055_OrchestrationV2.test.ts @@ -13,7 +13,7 @@ layer("055_OrchestrationV2", (it) => { Effect.sync(() => { assert.deepStrictEqual( migrationEntries.map(([id]) => id), - Array.from({ length: 56 }, (_, index) => index + 1), + Array.from({ length: 57 }, (_, index) => index + 1), ); }), ); @@ -28,6 +28,7 @@ layer("055_OrchestrationV2", (it) => { [54, "ProjectionThreadsAutoSettleDisabledAt"], [55, "OrchestrationV2"], [56, "RemoveRedundantProjectionIndexes"], + [57, "ScheduledTaskWebhooks"], ]); assert.deepStrictEqual(yield* runMigrations(), []); @@ -50,6 +51,7 @@ layer("055_OrchestrationV2", (it) => { { migration_id: 54, name: "ProjectionThreadsAutoSettleDisabledAt" }, { migration_id: 55, name: "OrchestrationV2" }, { migration_id: 56, name: "RemoveRedundantProjectionIndexes" }, + { migration_id: 57, name: "ScheduledTaskWebhooks" }, ]); const tables = yield* sql<{ readonly name: string }>` diff --git a/apps/server/src/persistence/Migrations/057_ScheduledTaskWebhooks.test.ts b/apps/server/src/persistence/Migrations/057_ScheduledTaskWebhooks.test.ts new file mode 100644 index 000000000000..3abd9c7e8c01 --- /dev/null +++ b/apps/server/src/persistence/Migrations/057_ScheduledTaskWebhooks.test.ts @@ -0,0 +1,51 @@ +import { assert, it } from "@effect/vitest"; +import * as NodeSqliteClient from "@t3tools/shared/nodeSqliteClient"; +import * as Effect from "effect/Effect"; +import * as Layer from "effect/Layer"; +import * as SqlClient from "effect/sql/SqlClient"; + +import { runMigrations } from "../Migrations.ts"; + +const layer = it.layer(Layer.mergeAll(NodeSqliteClient.layer({ filename: ":memory:" }))); + +layer("057_ScheduledTaskWebhooks", (it) => { + it.effect("keeps existing scheduled tasks and adds webhook storage", () => + Effect.gen(function* () { + const sql = yield* SqlClient.SqlClient; + yield* runMigrations({ toMigrationInclusive: 56 }); + yield* sql`INSERT INTO scheduled_tasks ${sql.insert({ + task_id: "existing", + title: "task", + prompt: "Run", + enabled: 1, + schedule_json: '{"type":"interval","everyMs":60000}', + project_id: "project", + thread_id: null, + workspace_strategy_json: '{"type":"root"}', + model_selection_json: '{"instanceId":"codex","model":"gpt-5"}', + runtime_mode: "full-access", + interaction_mode: "default", + created_by: "user", + creation_source: "web", + created_at: "2026-10-01T00:00:00.000Z", + updated_at: "2026-10-01T00:00:00.000Z", + next_run_at: null, + last_run_at: null, + last_run_status: "never", + last_run_error: null, + run_count: 0, + })}`; + yield* runMigrations({ toMigrationInclusive: 57 }); + + const rows = yield* sql<{ + task_id: string; + webhook_token: string | null; + }>`SELECT task_id, webhook_token FROM scheduled_tasks`; + assert.deepEqual(rows, [{ task_id: "existing", webhook_token: null }]); + const deliveries = yield* sql<{ + count: number; + }>`SELECT COUNT(*) AS count FROM scheduled_task_webhook_deliveries`; + assert.equal(deliveries[0]?.count, 0); + }), + ); +}); diff --git a/apps/server/src/persistence/Migrations/057_ScheduledTaskWebhooks.ts b/apps/server/src/persistence/Migrations/057_ScheduledTaskWebhooks.ts new file mode 100644 index 000000000000..4ef116634c03 --- /dev/null +++ b/apps/server/src/persistence/Migrations/057_ScheduledTaskWebhooks.ts @@ -0,0 +1,34 @@ +import * as Effect from "effect/Effect"; +import * as SqlClient from "effect/sql/SqlClient"; + +export default Effect.gen(function* () { + const sql = yield* SqlClient.SqlClient; + + // Kept out of schedule_json so they never decode into the read model. + yield* sql`ALTER TABLE scheduled_tasks ADD COLUMN webhook_token TEXT`; + yield* sql`ALTER TABLE scheduled_tasks ADD COLUMN webhook_secret TEXT`; + + yield* sql` + CREATE TABLE IF NOT EXISTS scheduled_task_webhook_deliveries ( + delivery_id TEXT PRIMARY KEY, + task_id TEXT NOT NULL, + received_at TEXT NOT NULL, + method TEXT NOT NULL, + query TEXT NOT NULL, + headers_json TEXT NOT NULL, + body TEXT NOT NULL, + body_bytes INTEGER NOT NULL, + body_truncated INTEGER NOT NULL, + outcome TEXT NOT NULL, + signature_verified INTEGER NOT NULL, + missing_fields_json TEXT NOT NULL, + rendered_prompt TEXT, + error TEXT + ) + `; + + yield* sql` + CREATE INDEX IF NOT EXISTS idx_scheduled_task_webhook_deliveries_task + ON scheduled_task_webhook_deliveries(task_id, received_at) + `; +}); diff --git a/apps/server/src/persistence/reconcileV2PreviewMigration.test.ts b/apps/server/src/persistence/reconcileV2PreviewMigration.test.ts index c0778bea2e9a..e872cb553bfa 100644 --- a/apps/server/src/persistence/reconcileV2PreviewMigration.test.ts +++ b/apps/server/src/persistence/reconcileV2PreviewMigration.test.ts @@ -37,6 +37,7 @@ describe("V2 preview upgrade", () => { [53, "PullRequestFilesViewed"], [54, "ProjectionThreadsAutoSettleDisabledAt"], [56, "RemoveRedundantProjectionIndexes"], + [57, "ScheduledTaskWebhooks"], ]); assert.deepStrictEqual(yield* runMigrations(), []); assert.deepStrictEqual(yield* sql`SELECT * FROM orchestration_v2_legacy_imports`, imports); @@ -116,6 +117,7 @@ describe("V2 preview upgrade", () => { [53, "PullRequestFilesViewed"], [54, "ProjectionThreadsAutoSettleDisabledAt"], [56, "RemoveRedundantProjectionIndexes"], + [57, "ScheduledTaskWebhooks"], ]); }).pipe(Effect.provide(NodeSqliteClient.layer({ filename: ":memory:" }))), ); diff --git a/apps/server/src/provider/T3OrchestrationInstructions.ts b/apps/server/src/provider/T3OrchestrationInstructions.ts index a70ef35bffcd..a04a84c7633c 100644 --- a/apps/server/src/provider/T3OrchestrationInstructions.ts +++ b/apps/server/src/provider/T3OrchestrationInstructions.ts @@ -9,7 +9,7 @@ The \`t3-code\` MCP server provides app-owned orchestration. Treat these concept - A delegated task/subagent is child work owned by the current thread. Use \`orchestrator_capabilities\` to discover the current provider/model IDs from the same live catalog as the composer, including configured custom models. Do not treat a native tool's model list as the full list of available subagent models. Prefer native subagent tools for same-provider work only when they support the chosen model. Use \`delegate_task\` with that provider instance and model when native tools cannot, including for same-provider work. Also use \`delegate_task\` for cross-provider or explicitly T3-owned child tasks. Retain each returned \`taskId\`, and use \`task_status\` or \`task_cancel\` to manage it. The returned \`childThreadId\` is backing storage for the subagent, not the target for starting another delegated review round. - \`t3_thread_launch\` and \`create_threads\` create ordinary top-level T3 conversations. Use them only when the user explicitly asks for separate/new/top-level threads or conversations. Never use them merely because the user said "subagent" or requested parallel delegated work. - For every T3 delegated review round, call \`delegate_task\` again. Include the original brief, prior findings, responses, and unresolved objections in each new task prompt. Track each round by its own \`taskId\`. Use a distinct \`clientRequestId\` per round, stable across retries of that round. Do not use \`t3_thread_send\` on \`childThreadId\` to continue a delegated review. -- \`schedule_task\` creates persistent recurring work in the app scheduler. Pass \`schedule\` as a structured object, never as JSON text: \`{"type":"interval","everyMs":3600000}\` for an interval, or \`{"type":"fixed_time","timeOfDay":"09:00","weekdays":[1,2,3,4,5]}\` for a wall-clock schedule. By default runs return to the current thread; set \`bindToCurrentThread=false\` only when the user wants a fresh thread for every run. After scheduling, report the returned cadence and next run time. +- \`schedule_task\` creates persistent recurring work in the app scheduler. Pass \`schedule\` as a structured object, never as JSON text: \`{"type":"interval","everyMs":3600000}\` for an interval, or \`{"type":"fixed_time","timeOfDay":"09:00","weekdays":[1,2,3,4,5]}\` for a wall-clock schedule, or \`{"type":"webhook"}\` to run on each request to the returned \`webhookUrl\` (the prompt may use \`{{body.path}}\`-style placeholders). By default runs return to the current thread; set \`bindToCurrentThread=false\` only when the user wants a fresh thread for every run. After scheduling, report the returned cadence and next run time. ### Choose the workspace before starting a new thread diff --git a/apps/server/src/scheduledTasks/Schedule.ts b/apps/server/src/scheduledTasks/Schedule.ts index fcf9ea81f50f..c6ec41e75143 100644 --- a/apps/server/src/scheduledTasks/Schedule.ts +++ b/apps/server/src/scheduledTasks/Schedule.ts @@ -13,6 +13,7 @@ export function nextScheduledRunAt( schedule: ScheduledTaskSchedule, from: DateTime.DateTime, ): DateTime.DateTime | null { + if (schedule.type === "webhook") return null; if (schedule.type === "interval") { // Persisted rows created before the one-minute floor remain readable, but // they must not retain their old high-frequency execution rate. @@ -55,6 +56,8 @@ export function isSameSchedule(a: ScheduledTaskSchedule, b: ScheduledTaskSchedul if (a.type === "interval") { return b.type === "interval" && a.everyMs === b.everyMs; } + // Webhook tasks never have a next run, whatever their signature settings. + if (a.type === "webhook") return b.type === "webhook"; if (b.type !== "fixed_time") return false; // The contract accepts padded and unpadded hours ("9:00" and "09:00"), so // compare the parsed time — string equality would treat a format-only edit @@ -93,6 +96,7 @@ export function isMissedFixedTimeRun( } function describeSchedule(schedule: ScheduledTaskSchedule): string { + if (schedule.type === "webhook") return "On webhook"; if (schedule.type === "interval") { const minutes = schedule.everyMs / MINUTE_MS; if (Number.isInteger(minutes)) { diff --git a/apps/server/src/scheduledTasks/ScheduledTaskService.ts b/apps/server/src/scheduledTasks/ScheduledTaskService.ts index f6fd7e9cb8c5..dca44f1e4112 100644 --- a/apps/server/src/scheduledTasks/ScheduledTaskService.ts +++ b/apps/server/src/scheduledTasks/ScheduledTaskService.ts @@ -5,8 +5,15 @@ import { ScheduledTaskError, ScheduledTaskId, ThreadId, + ScheduledTaskWebhookDeliveryId, type ScheduledTaskDeleteInput, type ScheduledTaskDeleteResult, + type ScheduledTaskGetWebhookDeliveryInput, + type ScheduledTaskGetWebhookDeliveryResult, + type ScheduledTaskListWebhookDeliveriesInput, + type ScheduledTaskListWebhookDeliveriesResult, + type ScheduledTaskRotateWebhookTokenInput, + type ScheduledTaskWebhookDeliveryOutcome, type ScheduledTaskListResult, type ScheduledTaskMutationResult, type ScheduledTaskRunNowInput, @@ -17,6 +24,7 @@ import { import * as Cause from "effect/Cause"; import * as Context from "effect/Context"; import * as Crypto from "effect/Crypto"; +import * as Data from "effect/Data"; import * as DateTime from "effect/DateTime"; import * as Effect from "effect/Effect"; import * as Layer from "effect/Layer"; @@ -25,6 +33,7 @@ import * as PubSub from "effect/PubSub"; import * as Ref from "effect/Ref"; import * as Result from "effect/Result"; import * as Schema from "effect/Schema"; +import * as Semaphore from "effect/Semaphore"; import * as Stream from "effect/Stream"; import * as SqlClient from "effect/sql/SqlClient"; @@ -32,6 +41,57 @@ import * as ThreadLaunchService from "../orchestration-v2/ThreadLaunchService.ts import * as ThreadManagementService from "../orchestration-v2/ThreadManagementService.ts"; import * as Scheduler from "../scheduling/Scheduler.ts"; import { isMissedFixedTimeRun, isSameSchedule, nextScheduledRunAt } from "./Schedule.ts"; +import { renderWebhookPrompt, type WebhookRequest } from "./webhookTemplate.ts"; +import { constantTimeEquals, verifyWebhookSignature } from "./webhookVerification.ts"; + +/** Path prefix of the environment route that receives webhook requests. */ +export const WEBHOOK_ROUTE_PREFIX = "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/api/hooks"; +/** Deliveries kept per task; older ones are pruned on insert. */ +const WEBHOOK_DELIVERY_RETENTION = 50; +/** Body text kept in the delivery log. Larger bodies are cut and flagged. */ +const WEBHOOK_DELIVERY_LOG_BODY_LIMIT = 64 * 1024; +/** Deliveries one task may hold at once, running or waiting their turn. */ +const WEBHOOK_MAX_QUEUED_PER_TASK = 20; +/** Accepted deliveries per task per minute, enforced here as well as on the relay because the tunnel hostname is public too. */ +const WEBHOOK_RATE_LIMIT_PER_MINUTE = 60; + +/** Where a webhook task's public URL points. `relayUrl` is null when the environment is not linked to T3 Connect. */ +interface WebhookOrigin { + readonly environmentId: string; + readonly relayUrl: string | null; +} + +export class ScheduledTaskWebhookOrigin extends Context.Reference>( + "t3/scheduledTasks/ScheduledTaskWebhookOrigin", + { + defaultValue: () => Effect.succeed({ environmentId: "local", relayUrl: null }), + }, +) {} + +/** A queued webhook delivery that no longer applies to its task; `reason` is shown in the delivery log. */ +class WebhookDeliverySkipped extends Data.TaggedError("WebhookDeliverySkipped")<{ + readonly reason: string; +}> {} + +interface RateWindow { + readonly accepted: ReadonlyArray; + /** Whether a rejection was already logged in this window. */ + readonly rejectedLogged: boolean; +} + +export interface WebhookTriggerRequest extends WebhookRequest { + readonly hookId: string; + readonly token: string; + readonly body: Uint8Array; +} + +/** What the HTTP route should answer. `not_found` covers unknown hooks and wrong tokens alike. */ +export type WebhookTriggerResult = + | { readonly _tag: "accepted"; readonly deliveryId: ScheduledTaskWebhookDeliveryId } + | { readonly _tag: "not_found" } + | { readonly _tag: "rejected_signature" } + | { readonly _tag: "disabled" } + | { readonly _tag: "rate_limited" }; const decodeTask = Schema.decodeUnknownEffect(ScheduledTask); const decodeTaskId = Schema.decodeUnknownOption(ScheduledTaskId); @@ -44,6 +104,12 @@ const decodeWorkspaceStrategyJson = Schema.decodeUnknownEffect( const decodeModelSelectionJson = Schema.decodeUnknownEffect( Schema.fromJsonString(ScheduledTask.fields.modelSelection), ); +const HeadersJson = Schema.fromJsonString(Schema.Record(Schema.String, Schema.String)); +const MissingFieldsJson = Schema.fromJsonString(Schema.Array(Schema.String)); +const decodeHeadersJson = Schema.decodeUnknownOption(HeadersJson); +const decodeMissingFieldsJson = Schema.decodeUnknownOption(MissingFieldsJson); +const encodeHeadersJson = Schema.encodeSync(HeadersJson); +const encodeMissingFieldsJson = Schema.encodeSync(MissingFieldsJson); interface ScheduledTaskRow { readonly task_id: string; @@ -66,6 +132,25 @@ interface ScheduledTaskRow { readonly last_run_status: string; readonly last_run_error: string | null; readonly run_count: number; + readonly webhook_token: string | null; + readonly webhook_secret: string | null; +} + +interface WebhookDeliveryRow { + readonly delivery_id: string; + readonly task_id: string; + readonly received_at: string; + readonly method: string; + readonly query: string; + readonly headers_json: string; + readonly body: string; + readonly body_bytes: number; + readonly body_truncated: number; + readonly outcome: string; + readonly signature_verified: number; + readonly missing_fields_json: string; + readonly rendered_prompt: string | null; + readonly error: string | null; } export class ScheduledTaskService extends Context.Service< @@ -87,6 +172,24 @@ export class ScheduledTaskService extends Context.Service< readonly runNow: ( input: ScheduledTaskRunNowInput, ) => Effect.Effect; + /** Issues a new URL token for a webhook task; the old URL stops working at once. */ + readonly rotateWebhookToken: ( + input: ScheduledTaskRotateWebhookTokenInput, + ) => Effect.Effect; + readonly listWebhookDeliveries: ( + input: ScheduledTaskListWebhookDeliveriesInput, + ) => Effect.Effect; + readonly getWebhookDelivery: ( + input: ScheduledTaskGetWebhookDeliveryInput, + ) => Effect.Effect; + /** + * Verifies, logs and dispatches one webhook request. Returns as soon as the + * delivery is logged; the run itself continues in the background so senders + * with short timeouts get their answer immediately. + */ + readonly triggerWebhook: ( + request: WebhookTriggerRequest, + ) => Effect.Effect; } >()("t3/scheduledTasks/ScheduledTaskService") {} @@ -125,9 +228,43 @@ function errorMessage(error: unknown): string { return String(error); } -const decodeRow = (row: ScheduledTaskRow) => +/** Headers kept out of the delivery log because they commonly carry credentials. */ +const REDACTED_HEADER = + /^(authorization|proxy-authorization|cookie|set-cookie)$|token|secret|signature|key|password|auth/i; + +function redactHeaders(headers: Readonly>): Record { + return Object.fromEntries( + Object.entries(headers).map(([name, value]) => [ + name, + REDACTED_HEADER.test(name) ? "[redacted]" : value, + ]), + ); +} + +function webhookPath(taskId: string, token: string): string { + return `${WEBHOOK_ROUTE_PREFIX}/${encodeURIComponent(taskId)}/${token}`; +} + +function webhookEndpoint( + row: Pick, + origin: WebhookOrigin | null, +): ScheduledTask["webhook"] { + if (row.webhook_token === null) return undefined; + const relayUrl = origin?.relayUrl?.replace(/\/+$/, "") ?? null; + return { + path: webhookPath(row.task_id, row.webhook_token), + url: + relayUrl === null || origin === null + ? null + : `${relayUrl}/v1/hooks/${encodeURIComponent(origin.environmentId)}/${encodeURIComponent(row.task_id)}/${row.webhook_token}`, + hasSecret: row.webhook_secret !== null, + }; +} + +const decodeRow = (row: ScheduledTaskRow, origin: WebhookOrigin | null = null) => Effect.gen(function* () { const schedule = yield* decodeScheduleJson(row.schedule_json); + const webhook = schedule.type === "webhook" ? webhookEndpoint(row, origin) : undefined; const workspaceStrategy = yield* decodeWorkspaceStrategyJson(row.workspace_strategy_json); const modelSelection = yield* decodeModelSelectionJson(row.model_selection_json); return yield* decodeTask({ @@ -153,6 +290,7 @@ const decodeRow = (row: ScheduledTaskRow) => lastRunStatus: row.last_run_status, lastRunError: row.last_run_error, runCount: row.run_count, + ...(webhook === undefined ? {} : { webhook }), }); }).pipe( Effect.mapError((cause) => { @@ -210,6 +348,16 @@ export const layer = Layer.effect( const threadLaunch = yield* ThreadLaunchService.ThreadLaunchService; const threadManagement = yield* ThreadManagementService.ThreadManagementService; const scheduler = yield* Scheduler.Scheduler; + const readWebhookOrigin = yield* ScheduledTaskWebhookOrigin; + // Webhook deliveries for one task dispatch in arrival order rather than + // being dropped while an earlier delivery is still dispatching. + const webhookPermits = yield* Ref.make>( + new Map(), + ); + // Keyed by task id and creation time, so deliveries of a deleted task that + // finish late release their own count, never a recreated task's. + const webhookQueued = yield* Ref.make>(new Map()); + const webhookRateWindows = yield* Ref.make>(new Map()); const activeRuns = yield* Ref.make>(new Set()); // Sliding(1) coalesces the dirty-signal: every notification triggers a // full list() re-emit anyway, so a slow subscriber only ever needs the @@ -238,7 +386,9 @@ export const layer = Layer.effect( last_run_at, last_run_status, last_run_error, - run_count + run_count, + webhook_token, + webhook_secret FROM scheduled_tasks ORDER BY updated_at DESC, task_id ASC `; @@ -246,7 +396,8 @@ export const layer = Layer.effect( // Strict decode for the API surface: a corrupt row is a visible error. const listRows = Effect.fn("ScheduledTaskService.listRows")(function* () { const rows = yield* selectAllRows(); - return yield* Effect.forEach(rows, decodeRow, { concurrency: 1 }); + const origin = yield* readWebhookOrigin; + return yield* Effect.forEach(rows, (row) => decodeRow(row, origin), { concurrency: 1 }); }); const getRows = (id: ScheduledTaskId) => sql` @@ -270,7 +421,9 @@ export const layer = Layer.effect( last_run_at, last_run_status, last_run_error, - run_count + run_count, + webhook_token, + webhook_secret FROM scheduled_tasks WHERE task_id = ${id} `; @@ -284,9 +437,26 @@ export const layer = Layer.effect( ); const row = rows[0]; if (row === undefined) return null; - return yield* decodeRow(row); + return yield* decodeRow(row, yield* readWebhookOrigin); }); + const findWebhookCredentials = (id: ScheduledTaskId) => + getRows(id).pipe( + Effect.map((rows) => + rows[0] === undefined + ? null + : { token: rows[0].webhook_token, secret: rows[0].webhook_secret }, + ), + Effect.mapError((cause) => + taskError("Could not load schedule task.", { taskId: id, cause }), + ), + ); + + const newWebhookToken = crypto.randomBytes(32).pipe( + Effect.map((bytes) => Buffer.from(bytes).toString("base64url")), + Effect.mapError((cause) => taskError("Could not generate webhook token.", { cause })), + ); + const loadTask = Effect.fn("ScheduledTaskService.loadTask")(function* (id: ScheduledTaskId) { const task = yield* findTask(id); if (task === null) { @@ -300,7 +470,16 @@ export const layer = Layer.effect( // concurrent settings save must not overwrite an in-flight increment. // Check existence in the write itself so an edit cannot undo a deletion // that landed after upsert loaded the previous task. - const saveTask = (task: ScheduledTask, requireExisting: boolean) => + const saveTask = ( + task: ScheduledTask, + requireExisting: boolean, + webhook: { + readonly token: string | null; + readonly secret: string | null; + /** False when the save carried no new secret, so a concurrent change survives. */ + readonly secretChanged: boolean; + }, + ) => sql<{ task_id: string }>` INSERT INTO scheduled_tasks ( task_id, @@ -322,7 +501,9 @@ export const layer = Layer.effect( last_run_at, last_run_status, last_run_error, - run_count + run_count, + webhook_token, + webhook_secret ) SELECT ${task.id}, @@ -344,7 +525,9 @@ export const layer = Layer.effect( ${task.lastRunAt}, ${task.lastRunStatus}, ${task.lastRunError}, - ${task.runCount} + ${task.runCount}, + ${webhook.token}, + ${webhook.secret} WHERE ${requireExisting ? 0 : 1} = 1 OR EXISTS (SELECT 1 FROM scheduled_tasks WHERE task_id = ${task.id}) ON CONFLICT (task_id) @@ -361,7 +544,17 @@ export const layer = Layer.effect( interaction_mode = excluded.interaction_mode, creation_source = excluded.creation_source, updated_at = excluded.updated_at, - next_run_at = excluded.next_run_at + next_run_at = excluded.next_run_at, + -- Only rotate changes a live token, so a save racing a rotation + -- cannot bring the old URL back. + webhook_token = CASE + WHEN excluded.webhook_token IS NULL THEN NULL + ELSE COALESCE(scheduled_tasks.webhook_token, excluded.webhook_token) + END, + webhook_secret = CASE + WHEN ${webhook.secretChanged ? 1 : 0} = 1 THEN excluded.webhook_secret + ELSE scheduled_tasks.webhook_secret + END RETURNING task_id `.pipe( Effect.mapError((cause) => @@ -375,27 +568,38 @@ export const layer = Layer.effect( ); const deleteRow = (id: ScheduledTaskId) => - sql`DELETE FROM scheduled_tasks WHERE task_id = ${id}`.pipe( - Effect.mapError((cause) => - taskError("Could not delete schedule task.", { taskId: id, cause }), - ), - ); + sql + .withTransaction( + sql`DELETE FROM scheduled_task_webhook_deliveries WHERE task_id = ${id}`.pipe( + Effect.andThen(sql`DELETE FROM scheduled_tasks WHERE task_id = ${id}`), + ), + ) + .pipe( + Effect.mapError((cause) => + taskError("Could not delete schedule task.", { taskId: id, cause }), + ), + ); // Run-state transitions use targeted UPDATEs (never the full-row upsert) so // a completing run cannot resurrect a deleted task or clobber concurrent // edits to the task definition. const markRunning = (id: ScheduledTaskId, startedAtIso: string) => - sql` + sql<{ task_id: string }>` UPDATE scheduled_tasks SET updated_at = ${startedAtIso}, last_run_at = ${startedAtIso}, last_run_status = 'running', last_run_error = NULL WHERE task_id = ${id} + RETURNING task_id `.pipe( Effect.mapError((cause) => taskError("Could not mark schedule task as running.", { taskId: id, cause }), ), + // A task deleted after the re-read must not be dispatched from the stale snapshot. + Effect.flatMap((rows) => + rows.length > 0 ? Effect.void : taskError("Schedule task not found.", { taskId: id }), + ), ); const markCompleted = (input: { @@ -457,7 +661,8 @@ export const layer = Layer.effect( const runTask = Effect.fn("ScheduledTaskService.runTask")(function* ( task: ScheduledTask, - trigger: "scheduled" | "manual", + trigger: "scheduled" | "manual" | "webhook", + webhook?: { readonly deliveryId: string; readonly prompt: string }, ) { const reserved = yield* Ref.modify(activeRuns, (active) => { if (active.has(task.id)) return [false, active] as const; @@ -466,7 +671,7 @@ export const layer = Layer.effect( return [true, next] as const; }); if (!reserved) { - if (trigger === "manual") { + if (trigger !== "scheduled") { return yield* taskError("Schedule task is already running.", { taskId: task.id }); } return task; @@ -481,9 +686,14 @@ export const layer = Layer.effect( // the poll loaded it — none of those may fire. const active = yield* findTask(task.id); if (active === null) { + if (webhook !== undefined) { + return yield* new WebhookDeliverySkipped({ + reason: "The task was deleted before this delivery ran.", + }); + } // A manual run on a just-deleted task must fail loudly, not report // a successful run that never dispatched. - if (trigger === "manual") { + if (trigger !== "scheduled") { return yield* taskError("Schedule task not found.", { taskId: task.id }); } return task; @@ -500,16 +710,36 @@ export const layer = Layer.effect( ) { return active; } + // A queued delivery must not run a task that was paused, deleted and + // recreated under the same id, or switched to another trigger while + // it waited for its turn. + if (webhook !== undefined) { + const reason = + active.createdAt !== task.createdAt + ? "The task was replaced before this delivery ran." + : active.schedule.type !== "webhook" + ? "The task's trigger changed before this delivery ran." + : !active.enabled + ? "The task was paused before this delivery ran." + : null; + if (reason !== null) return yield* new WebhookDeliverySkipped({ reason }); + } yield* markRunning(active.id, startedAtIso); yield* notifyChanged; - const fireKey = `${active.id}:${DateTime.toEpochMillis(startedAt)}:${trigger}`; + // A webhook run is keyed by its delivery so the same delivery can + // never dispatch twice. + const fireKey = + webhook === undefined + ? `${active.id}:${DateTime.toEpochMillis(startedAt)}:${trigger}` + : `${active.id}:webhook:${webhook.deliveryId}`; const commandId = CommandId.make(`scheduled-task:${fireKey}`); const messageId = MessageId.make(`scheduled-task-message:${fireKey}`); // Dispatch from the fresh row so prompt/model/binding edits made - // after the poll read are honoured. - const prompt = active.prompt; + // after the poll read are honoured. A webhook prompt was rendered + // from the row when the request arrived. + const prompt = webhook?.prompt ?? active.prompt; // Effect.exit (not Effect.result) so defects and interruptions in the // dispatch are also captured and recorded as a failed run instead of @@ -632,9 +862,10 @@ export const layer = Layer.effect( yield* Effect.forEach( due, ({ task, dueAt }) => - (isMissedFixedTimeRun(task.schedule, dueAt, now) - ? rescheduleMissedRun(task, now) - : runTask(task, "scheduled") + Effect.suspend(() => + isMissedFixedTimeRun(task.schedule, dueAt, now) + ? rescheduleMissedRun(task, now) + : runTask(task, "scheduled"), ).pipe( Effect.catch((cause) => Effect.logWarning("Scheduled task run failed", { taskId: task.id, cause }), @@ -742,16 +973,50 @@ export const layer = Layer.effect( // Keep the existing next_run_at when the schedule itself is untouched: // editing a title or prompt must not postpone (or resurrect) a due // run — only schedule/enabled changes restart the clock. + const schedule: ScheduledTask["schedule"] = + input.schedule.type === "webhook" + ? { + type: "webhook", + signature: + input.schedule.signature == null + ? null + : { + header: input.schedule.signature.header.toLowerCase(), + encoding: input.schedule.signature.encoding, + prefix: input.schedule.signature.prefix, + }, + } + : input.schedule; + const webhook = + input.schedule.type === "webhook" + ? yield* Effect.gen(function* () { + const existing = existingTask === null ? null : yield* findWebhookCredentials(id); + // Saving a webhook task keeps its URL; only rotate changes it. + const token = existing?.token ?? (yield* newWebhookToken); + const signature = + input.schedule.type === "webhook" ? input.schedule.signature : null; + const secret = + signature == null ? null : (signature.secret ?? existing?.secret ?? null); + const secretChanged = + signature == null || signature.secret !== undefined || existing === null; + if (signature != null && secret === null) { + return yield* taskError("A webhook signature check needs a signing secret.", { + taskId: id, + }); + } + return { token, secret, secretChanged }; + }) + : { token: null, secret: null, secretChanged: true }; const scheduleUnchanged = existingTask !== null && existingTask.enabled === input.enabled && - isSameSchedule(existingTask.schedule, input.schedule); + isSameSchedule(existingTask.schedule, schedule); const task: ScheduledTask = { id, title: input.title, prompt: input.prompt, enabled: input.enabled, - schedule: input.schedule, + schedule, projectId: input.projectId, threadId: input.threadId ?? null, workspaceStrategy: input.workspaceStrategy, @@ -764,15 +1029,15 @@ export const layer = Layer.effect( updatedAt: iso(now), nextRunAt: scheduleUnchanged ? existingTask.nextRunAt - : nextRunAt({ enabled: input.enabled, schedule: input.schedule }, now), + : nextRunAt({ enabled: input.enabled, schedule }, now), lastRunAt: existingTask?.lastRunAt ?? null, lastRunStatus: existingTask?.lastRunStatus ?? "never", lastRunError: existingTask?.lastRunError ?? null, runCount: existingTask?.runCount ?? 0, }; - yield* saveTask(task, input.requireExisting === true); + yield* saveTask(task, input.requireExisting === true, webhook); yield* notifyChanged; - return { task }; + return { task: (yield* findTask(id)) ?? task }; }); const setEnabled: ScheduledTaskService["Service"]["setEnabled"] = (input) => @@ -805,11 +1070,34 @@ export const layer = Layer.effect( }); const deleteTask: ScheduledTaskService["Service"]["delete"] = (input) => - deleteRow(input.id).pipe(Effect.andThen(notifyChanged), Effect.as({ id: input.id })); + deleteRow(input.id).pipe( + Effect.andThen( + Effect.all([ + Ref.update(webhookRateWindows, (windows) => { + const next = new Map(windows); + next.delete(input.id); + return next; + }), + Ref.update(webhookPermits, (permits) => { + const next = new Map(permits); + next.delete(input.id); + return next; + }), + ]), + ), + Effect.andThen(notifyChanged), + Effect.as({ id: input.id }), + ); const runNow: ScheduledTaskService["Service"]["runNow"] = (input: ScheduledTaskRunNowInput) => Effect.gen(function* () { const task = yield* loadTask(input.id); + if (task.schedule.type === "webhook") { + // There is no request to render the prompt from. + return yield* taskError("Webhook tasks run when their URL receives a request.", { + taskId: input.id, + }); + } const next = yield* runTask(task, "manual").pipe( Effect.mapError((cause) => taskError("Could not run schedule task.", { taskId: input.id, cause }), @@ -818,6 +1106,318 @@ export const layer = Layer.effect( return { task: next }; }); + const rotateWebhookToken: ScheduledTaskService["Service"]["rotateWebhookToken"] = (input) => + Effect.gen(function* () { + const task = yield* loadTask(input.id); + if (task.schedule.type !== "webhook") { + return yield* taskError("Only webhook tasks have a URL token.", { taskId: input.id }); + } + const token = yield* newWebhookToken; + const now = yield* localNow; + // Matching created_at keeps a rotation from landing on a task deleted + // and recreated under the same id since it was loaded. + const updated = yield* sql<{ task_id: string }>` + UPDATE scheduled_tasks + SET webhook_token = ${token}, updated_at = ${iso(now)} + WHERE task_id = ${input.id} AND created_at = ${task.createdAt} + RETURNING task_id + `.pipe( + Effect.mapError((cause) => + taskError("Could not rotate webhook token.", { taskId: input.id, cause }), + ), + ); + if (updated.length === 0) { + return yield* taskError("Schedule task was deleted or replaced.", { taskId: input.id }); + } + yield* notifyChanged; + return { task: yield* loadTask(input.id) }; + }); + + const deliveryHeaders = (row: WebhookDeliveryRow): Readonly> => + Option.getOrElse(decodeHeadersJson(row.headers_json), () => ({})); + const decodeDeliverySummary = (row: WebhookDeliveryRow) => ({ + id: ScheduledTaskWebhookDeliveryId.make(row.delivery_id), + taskId: ScheduledTaskId.make(row.task_id), + receivedAt: row.received_at, + method: row.method, + contentType: deliveryHeaders(row)["content-type"] ?? null, + bodyBytes: row.body_bytes, + outcome: row.outcome as ScheduledTaskWebhookDeliveryOutcome, + signatureVerified: row.signature_verified === 1, + missingFields: Option.getOrElse(decodeMissingFieldsJson(row.missing_fields_json), () => []), + error: row.error, + }); + + const listWebhookDeliveries: ScheduledTaskService["Service"]["listWebhookDeliveries"] = ( + input, + ) => + sql` + SELECT * FROM scheduled_task_webhook_deliveries + WHERE task_id = ${input.id} + ORDER BY received_at DESC, rowid DESC + `.pipe( + Effect.map((rows) => ({ deliveries: rows.map(decodeDeliverySummary) })), + Effect.mapError((cause) => + taskError("Could not list webhook deliveries.", { taskId: input.id, cause }), + ), + ); + + const getWebhookDelivery: ScheduledTaskService["Service"]["getWebhookDelivery"] = (input) => + Effect.gen(function* () { + const rows = yield* sql` + SELECT * FROM scheduled_task_webhook_deliveries + WHERE task_id = ${input.id} AND delivery_id = ${input.deliveryId} + `.pipe( + Effect.mapError((cause) => + taskError("Could not load webhook delivery.", { taskId: input.id, cause }), + ), + ); + const row = rows[0]; + if (row === undefined) { + return yield* taskError("Webhook delivery not found.", { taskId: input.id }); + } + return { + delivery: { + ...decodeDeliverySummary(row), + query: row.query, + headers: deliveryHeaders(row), + body: row.body, + bodyTruncated: row.body_truncated === 1, + renderedPrompt: row.rendered_prompt, + }, + }; + }); + + const recordDelivery = (input: { + readonly id: string; + readonly taskId: ScheduledTaskId; + readonly receivedAt: string; + readonly request: WebhookTriggerRequest; + readonly outcome: ScheduledTaskWebhookDeliveryOutcome; + readonly signatureVerified: boolean; + readonly missing: ReadonlyArray; + readonly renderedPrompt: string | null; + }) => { + const truncated = input.request.body.byteLength > WEBHOOK_DELIVERY_LOG_BODY_LIMIT; + const loggedBody = truncated + ? new TextDecoder().decode(input.request.body.subarray(0, WEBHOOK_DELIVERY_LOG_BODY_LIMIT)) + : input.request.bodyText; + return sql + .withTransaction( + Effect.gen(function* () { + // Conditional on the task existing, so a delivery racing a delete + // cannot leave rows that a recreated task with the same id would show. + yield* sql` + INSERT INTO scheduled_task_webhook_deliveries ( + delivery_id, task_id, received_at, method, query, headers_json, body, + body_bytes, body_truncated, outcome, signature_verified, + missing_fields_json, rendered_prompt, error + ) + SELECT + ${input.id}, ${input.taskId}, ${input.receivedAt}, ${input.request.method}, + ${input.request.query}, ${encodeHeadersJson(redactHeaders(input.request.headers))}, + ${loggedBody}, + ${input.request.body.byteLength}, ${truncated ? 1 : 0}, ${input.outcome}, + ${input.signatureVerified ? 1 : 0}, ${encodeMissingFieldsJson(input.missing)}, + ${input.renderedPrompt}, NULL + WHERE EXISTS (SELECT 1 FROM scheduled_tasks WHERE task_id = ${input.taskId}) + `; + yield* sql` + DELETE FROM scheduled_task_webhook_deliveries + WHERE task_id = ${input.taskId} + AND delivery_id NOT IN ( + SELECT delivery_id FROM scheduled_task_webhook_deliveries + WHERE task_id = ${input.taskId} + -- rowid breaks timestamp ties in arrival order; delivery ids are random. + ORDER BY received_at DESC, rowid DESC + LIMIT ${WEBHOOK_DELIVERY_RETENTION} + ) + `; + }), + ) + .pipe( + Effect.mapError((cause) => + taskError("Could not record webhook delivery.", { taskId: input.taskId, cause }), + ), + ); + }; + + const markDeliveryFailed = (deliveryId: string, message: string) => + sql` + UPDATE scheduled_task_webhook_deliveries + SET outcome = 'dispatch_failed', error = ${message} + WHERE delivery_id = ${deliveryId} + `.pipe(Effect.ignore); + + /** + * Sliding one-minute window counting every request with a valid token, + * including ones the signature check later rejects. + */ + const takeRateSlot = (id: ScheduledTaskId, nowMs: number) => + Ref.modify( + webhookRateWindows, + ( + windows, + ): readonly [ + "allowed" | "first_rejected" | "rejected", + ReadonlyMap, + ] => { + const current = windows.get(id); + const recent = (current?.accepted ?? []).filter((at) => nowMs - at < 60_000); + if (recent.length >= WEBHOOK_RATE_LIMIT_PER_MINUTE) { + const first = current?.rejectedLogged !== true; + return [ + first ? "first_rejected" : "rejected", + new Map(windows).set(id, { accepted: recent, rejectedLogged: true }), + ]; + } + return [ + "allowed", + new Map(windows).set(id, { accepted: [...recent, nowMs], rejectedLogged: false }), + ]; + }, + ); + + const webhookPermit = (id: ScheduledTaskId) => + Effect.gen(function* () { + const existing = (yield* Ref.get(webhookPermits)).get(id); + if (existing !== undefined) return existing; + const created = yield* Semaphore.make(1); + return yield* Ref.modify(webhookPermits, (permits) => { + const raced = permits.get(id); + return raced === undefined + ? [created, new Map(permits).set(id, created)] + : [raced, permits]; + }); + }); + + // Detached from the request so the HTTP response does not wait for the + // run; scoped to the service so shutdown interrupts it. + const serviceScope = yield* Effect.scope; + + const triggerWebhook: ScheduledTaskService["Service"]["triggerWebhook"] = (request) => + Effect.gen(function* () { + const taskId = decodeTaskId(request.hookId); + if (Option.isNone(taskId)) return { _tag: "not_found" as const }; + const rows = yield* getRows(taskId.value).pipe( + Effect.mapError((cause) => + taskError("Could not load schedule task.", { taskId: taskId.value, cause }), + ), + ); + const row = rows[0]; + // A wrong token is indistinguishable from an unknown hook, so the URL + // does not reveal which hooks exist. + if ( + row === undefined || + row.webhook_token === null || + !constantTimeEquals(request.token, row.webhook_token) + ) { + return { _tag: "not_found" as const }; + } + const task = yield* decodeRow(row); + if (task.schedule.type !== "webhook") return { _tag: "not_found" as const }; + + const receivedAt = yield* localNow; + const deliveryUuid = yield* crypto.randomUUIDv4.pipe( + Effect.mapError((cause) => taskError("Could not generate delivery id.", { cause })), + ); + const deliveryId = ScheduledTaskWebhookDeliveryId.make(`delivery:${deliveryUuid}`); + const log = ( + outcome: ScheduledTaskWebhookDeliveryOutcome, + details: { + readonly signatureVerified?: boolean; + readonly missing?: ReadonlyArray; + readonly renderedPrompt?: string; + } = {}, + ) => + recordDelivery({ + id: deliveryId, + taskId: task.id, + receivedAt: iso(receivedAt), + request, + outcome, + signatureVerified: details.signatureVerified ?? false, + missing: details.missing ?? [], + renderedPrompt: details.renderedPrompt ?? null, + }); + + // Only the first rejected request in a window is logged, so a flood + // cannot write rows or push the real deliveries out of the log. + const slot = yield* takeRateSlot(task.id, DateTime.toEpochMillis(receivedAt)); + if (slot !== "allowed") { + if (slot === "first_rejected") yield* log("rate_limited"); + return { _tag: "rate_limited" as const }; + } + if (!task.enabled) { + yield* log("disabled"); + return { _tag: "disabled" as const }; + } + const signature = task.schedule.signature; + if (signature !== null) { + const verified = + row.webhook_secret !== null && + verifyWebhookSignature({ + signature, + secret: row.webhook_secret, + headers: request.headers, + body: request.body, + }); + if (!verified) { + yield* log("rejected_signature"); + return { _tag: "rejected_signature" as const }; + } + } + + const rendered = renderWebhookPrompt(task.prompt, request); + // Bound the deliveries one task holds, so steady traffic to a stuck + // task cannot pile up parked fibers. A refused request is not logged, + // so it cannot push real deliveries out of the log. + const queueKey = `${task.id}\u0000${task.createdAt}`; + const queued = yield* Ref.modify(webhookQueued, (counts) => { + const count = counts.get(queueKey) ?? 0; + return count >= WEBHOOK_MAX_QUEUED_PER_TASK + ? ([false, counts] as const) + : ([true, new Map(counts).set(queueKey, count + 1)] as const); + }); + if (!queued) return { _tag: "rate_limited" as const }; + // Entries leave the map when their count reaches zero, so a deleted + // task's key does not linger once its last delivery finishes. + const release = Ref.update(webhookQueued, (counts) => { + const next = new Map(counts); + const count = (next.get(queueKey) ?? 1) - 1; + if (count <= 0) next.delete(queueKey); + else next.set(queueKey, count); + return next; + }); + yield* log("accepted", { + signatureVerified: signature !== null, + missing: rendered.missing, + renderedPrompt: rendered.prompt, + }).pipe(Effect.onError(() => release)); + const permit = yield* webhookPermit(task.id); + yield* runTask(task, "webhook", { deliveryId, prompt: rendered.prompt }).pipe( + Effect.flatMap((completed) => + completed.lastRunStatus === "failed" + ? markDeliveryFailed(deliveryId, "The run failed to start.") + : Effect.void, + ), + Effect.catchTag("WebhookDeliverySkipped", (skipped) => + markDeliveryFailed(deliveryId, skipped.reason), + ), + // The log is readable over RPC, so it gets a fixed reason; the + // cause, which can carry request data, stays in the server log. + Effect.catchCause((cause) => + Effect.logWarning("Webhook dispatch failed", { taskId: task.id, cause }).pipe( + Effect.andThen(markDeliveryFailed(deliveryId, "The run failed to start.")), + ), + ), + permit.withPermits(1), + Effect.ensuring(release), + Effect.forkIn(serviceScope), + ); + return { _tag: "accepted" as const, deliveryId }; + }); + return ScheduledTaskService.of({ list, subscribeList, @@ -825,6 +1425,10 @@ export const layer = Layer.effect( setEnabled, delete: deleteTask, runNow, + rotateWebhookToken, + listWebhookDeliveries, + getWebhookDelivery, + triggerWebhook, }); }), ); diff --git a/apps/server/src/scheduledTasks/ScheduledTaskService.webhook.test.ts b/apps/server/src/scheduledTasks/ScheduledTaskService.webhook.test.ts new file mode 100644 index 000000000000..7108d7b804f5 --- /dev/null +++ b/apps/server/src/scheduledTasks/ScheduledTaskService.webhook.test.ts @@ -0,0 +1,466 @@ +import * as NodeCrypto from "node:crypto"; + +import * as NodePlatformCrypto from "@effect/platform-node/NodeCrypto"; +import { assert, it } from "@effect/vitest"; +import { ScheduledTaskUpsertInput } from "@t3tools/contracts"; +import * as Deferred from "effect/Deferred"; +import * as Effect from "effect/Effect"; +import * as Layer from "effect/Layer"; +import * as Queue from "effect/Queue"; +import * as Schema from "effect/Schema"; +import * as TestClock from "effect/testing/TestClock"; + +import * as ThreadLaunchService from "../orchestration-v2/ThreadLaunchService.ts"; +import * as ThreadManagementService from "../orchestration-v2/ThreadManagementService.ts"; +import { SqlitePersistenceMemory } from "../persistence/Layers/Sqlite.ts"; +import * as Scheduler from "../scheduling/Scheduler.ts"; +import * as ScheduledTaskService from "./ScheduledTaskService.ts"; + +const decodeUpsertInput = Schema.decodeUnknownEffect(ScheduledTaskUpsertInput); + +type LaunchInput = ThreadLaunchService.ThreadLaunchInput; + +const webhookTaskInput = (overrides: Record = {}) => + decodeUpsertInput({ + id: "scheduled-task:hook", + title: "Review PRs", + prompt: "Review this PR: {{body.pull_request.url}}", + enabled: true, + schedule: { type: "webhook" }, + projectId: "project-webhook", + workspaceStrategy: { type: "root" }, + modelSelection: { instanceId: "codex", model: "gpt-5.4" }, + runtimeMode: "full-access", + interactionMode: "default", + ...overrides, + }); + +const pullRequestBody = new TextEncoder().encode( + JSON.stringify({ pull_request: { url: "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/org/repo/pull/45" } }), +); + +const requestFor = ( + task: { readonly id: string; readonly webhook?: { readonly path: string } | undefined }, + overrides: Partial = {}, +): ScheduledTaskService.WebhookTriggerRequest => ({ + hookId: task.id, + token: task.webhook?.path.split("/").at(-1) ?? "", + method: "POST", + path: `/api/hooks/${task.id}`, + query: "", + headers: { "content-type": "application/json" }, + body: pullRequestBody, + bodyText: new TextDecoder().decode(pullRequestBody), + ...overrides, +}); + +/** + * Runs `body` against a service whose launches are pushed to `launches`; + * `gate`, when given, holds each launch until the test releases it. + */ +const withService = ( + body: (input: { + readonly service: ScheduledTaskService.ScheduledTaskService["Service"]; + readonly launches: Queue.Queue; + }) => Effect.Effect, + options: { readonly gate?: Deferred.Deferred } = {}, +) => + Effect.gen(function* () { + const launches = yield* Queue.unbounded(); + const dependencies = Layer.mergeAll( + NodePlatformCrypto.layer, + Scheduler.layer, + Layer.mock(ThreadLaunchService.ThreadLaunchService)({ + launch: (input) => + Queue.offer(launches, input).pipe( + Effect.andThen(options.gate ? Deferred.await(options.gate) : Effect.void), + Effect.as({ threadId: "thread-1", resumed: false } as never), + ), + }), + Layer.mock(ThreadManagementService.ThreadManagementService)({}), + ); + return yield* Effect.gen(function* () { + const service = yield* ScheduledTaskService.ScheduledTaskService; + return yield* body({ service, launches }); + }).pipe(Effect.provide(ScheduledTaskService.layer.pipe(Layer.provide(dependencies)))); + }).pipe(Effect.provide(SqlitePersistenceMemory)); + +it.effect("dispatches exactly the rendered prompt and logs the delivery", () => + withService(({ service, launches }) => + Effect.gen(function* () { + const { task } = yield* service.upsert(yield* webhookTaskInput()); + assert.equal(task.nextRunAt, null); + assert.isDefined(task.webhook); + assert.isTrue(task.webhook!.path.startsWith("/api/hooks/scheduled-task%3Ahook/")); + // Not linked to T3 Connect in tests. + assert.equal(task.webhook!.url, null); + + const result = yield* service.triggerWebhook(requestFor(task)); + assert.equal(result._tag, "accepted"); + const launched = yield* Queue.take(launches); + assert.equal( + launched.initialMessage?.text, + "Review this PR: https://github.com/org/repo/pull/45", + ); + assert.equal( + launched.commandId, + `scheduled-task:${task.id}:webhook:${result._tag === "accepted" ? result.deliveryId : ""}`, + ); + + const { deliveries } = yield* service.listWebhookDeliveries({ id: task.id }); + assert.equal(deliveries.length, 1); + assert.equal(deliveries[0]?.outcome, "accepted"); + const { delivery } = yield* service.getWebhookDelivery({ + id: task.id, + deliveryId: deliveries[0]!.id, + }); + assert.equal(delivery.body, new TextDecoder().decode(pullRequestBody)); + assert.equal(delivery.renderedPrompt, launched.initialMessage?.text); + }), + ), +); + +it.effect("answers not found for a wrong token or unknown hook without logging", () => + withService(({ service }) => + Effect.gen(function* () { + const { task } = yield* service.upsert(yield* webhookTaskInput()); + const wrongToken = yield* service.triggerWebhook(requestFor(task, { token: "nope" })); + assert.equal(wrongToken._tag, "not_found"); + const unknown = yield* service.triggerWebhook(requestFor(task, { hookId: "missing" })); + assert.equal(unknown._tag, "not_found"); + assert.equal((yield* service.listWebhookDeliveries({ id: task.id })).deliveries.length, 0); + }), + ), +); + +it.effect("rotating the token retires the old URL and saving keeps it", () => + withService(({ service }) => + Effect.gen(function* () { + const input = yield* webhookTaskInput(); + const { task } = yield* service.upsert(input); + const saved = yield* service.upsert(yield* webhookTaskInput({ title: "Renamed" })); + assert.equal(saved.task.webhook?.path, task.webhook?.path); + + const rotated = yield* service.rotateWebhookToken({ id: task.id }); + assert.notEqual(rotated.task.webhook?.path, task.webhook?.path); + assert.equal((yield* service.triggerWebhook(requestFor(task)))._tag, "not_found"); + assert.equal((yield* service.triggerWebhook(requestFor(rotated.task)))._tag, "accepted"); + }), + ), +); + +it.effect("checks the configured signature and keeps the secret write-only", () => + withService(({ service }) => + Effect.gen(function* () { + const { task } = yield* service.upsert( + yield* webhookTaskInput({ + schedule: { + type: "webhook", + signature: { + header: "X-Hub-Signature-256", + encoding: "hex", + prefix: "sha256=", + secret: "s3cret", + }, + }, + }), + ); + assert.deepEqual(task.schedule, { + type: "webhook", + signature: { header: "x-hub-signature-256", encoding: "hex", prefix: "sha256=" }, + }); + assert.isTrue(task.webhook?.hasSecret); + + const unsigned = yield* service.triggerWebhook(requestFor(task)); + assert.equal(unsigned._tag, "rejected_signature"); + + const signature = `sha256=${NodeCrypto.createHmac("sha256", "s3cret").update(pullRequestBody).digest("hex")}`; + const signed = yield* service.triggerWebhook( + requestFor(task, { + headers: { "content-type": "application/json", "x-hub-signature-256": signature }, + }), + ); + assert.equal(signed._tag, "accepted"); + + // Saving without a secret keeps the stored one. + const resaved = yield* service.upsert( + yield* webhookTaskInput({ + schedule: { + type: "webhook", + signature: { header: "x-hub-signature-256", encoding: "hex", prefix: "sha256=" }, + }, + }), + ); + assert.isTrue(resaved.task.webhook?.hasSecret); + + const outcomes = (yield* service.listWebhookDeliveries({ id: task.id })).deliveries.map( + (delivery) => [delivery.outcome, delivery.signatureVerified], + ); + assert.deepEqual(outcomes.toSorted(), [ + ["accepted", true], + ["rejected_signature", false], + ]); + }), + ), +); + +it.effect("logs but does not run deliveries to a disabled task, and refuses run now", () => + withService(({ service, launches }) => + Effect.gen(function* () { + const { task } = yield* service.upsert(yield* webhookTaskInput({ enabled: false })); + assert.equal((yield* service.triggerWebhook(requestFor(task)))._tag, "disabled"); + assert.equal(yield* Queue.size(launches), 0); + const runNow = yield* service.runNow({ id: task.id }).pipe(Effect.flip); + assert.equal(runNow.message, "Webhook tasks run when their URL receives a request."); + }), + ), +); + +it.effect("queues a burst of deliveries instead of dropping them", () => + Effect.gen(function* () { + const gate = yield* Deferred.make(); + yield* withService( + ({ service, launches }) => + Effect.gen(function* () { + const { task } = yield* service.upsert(yield* webhookTaskInput()); + const results = yield* Effect.forEach([1, 2, 3], () => + service.triggerWebhook(requestFor(task)), + ); + assert.deepEqual( + results.map((result) => result._tag), + ["accepted", "accepted", "accepted"], + ); + // Only the first is dispatching; the others wait their turn. + yield* Queue.take(launches); + yield* Deferred.succeed(gate, undefined); + const rest = yield* Effect.all([Queue.take(launches), Queue.take(launches)]); + assert.equal(new Set(rest.map((launch) => launch.commandId)).size, 2); + }), + { gate }, + ); + }), +); + +it.effect("rate limits a hook past 60 deliveries a minute", () => + withService(({ service }) => + Effect.gen(function* () { + const { task } = yield* service.upsert(yield* webhookTaskInput({ enabled: false })); + const results = yield* Effect.forEach(Array.from({ length: 61 }), () => + service.triggerWebhook(requestFor(task)), + ); + assert.equal(results.at(-2)?._tag, "disabled"); + assert.equal(results.at(-1)?._tag, "rate_limited"); + // Further rejections in the same window are counted, not logged. + yield* Effect.forEach([1, 2, 3], () => service.triggerWebhook(requestFor(task))); + const outcomes = (yield* service.listWebhookDeliveries({ id: task.id })).deliveries.map( + (delivery) => delivery.outcome, + ); + assert.equal(outcomes.filter((outcome) => outcome === "rate_limited").length, 1); + }), + ), +); + +it.effect("a save carrying a stale token cannot undo a rotation", () => + withService(({ service }) => + Effect.gen(function* () { + const { task } = yield* service.upsert(yield* webhookTaskInput()); + const rotated = yield* service.rotateWebhookToken({ id: task.id }); + // The editor was opened before the rotation and saves afterwards. + const saved = yield* service.upsert(yield* webhookTaskInput({ title: "Edited" })); + assert.equal(saved.task.webhook?.path, rotated.task.webhook?.path); + assert.equal((yield* service.triggerWebhook(requestFor(task)))._tag, "not_found"); + }), + ), +); + +it.effect("keeps the newest 50 deliveries when they share a timestamp", () => + withService(({ service }) => + Effect.gen(function* () { + // Paused, so each request is logged without starting a run. The test + // clock is frozen, so every delivery has the same received_at. + const { task } = yield* service.upsert(yield* webhookTaskInput({ enabled: false })); + const ids = yield* Effect.forEach(Array.from({ length: 55 }), (_, index) => + service.triggerWebhook(requestFor(task, { query: `n=${index}` })).pipe(Effect.as(index)), + ); + const { deliveries } = yield* service.listWebhookDeliveries({ id: task.id }); + assert.equal(deliveries.length, 50); + const first = yield* service.getWebhookDelivery({ + id: task.id, + deliveryId: deliveries[0]!.id, + }); + const last = yield* service.getWebhookDelivery({ + id: task.id, + deliveryId: deliveries.at(-1)!.id, + }); + assert.equal(first.delivery.query, `n=${ids.at(-1)}`); + assert.equal(last.delivery.query, "n=5"); + }), + ), +); + +it.effect("keeps credential headers out of the delivery log", () => + withService(({ service }) => + Effect.gen(function* () { + const { task } = yield* service.upsert(yield* webhookTaskInput({ enabled: false })); + yield* service.triggerWebhook( + requestFor(task, { + headers: { + "content-type": "application/json", + authorization: "Bearer sender-token", + "x-webhook-key": "k", + "x-github-event": "push", + }, + }), + ); + const [summary] = (yield* service.listWebhookDeliveries({ id: task.id })).deliveries; + const { delivery } = yield* service.getWebhookDelivery({ + id: task.id, + deliveryId: summary!.id, + }); + assert.equal(delivery.headers.authorization, "[redacted]"); + assert.equal(delivery.headers["x-webhook-key"], "[redacted]"); + assert.equal(delivery.headers["x-github-event"], "push"); + }), + ), +); + +const queuedDeliveryCases = [ + { + change: "paused", + reason: "The task was paused before this delivery ran.", + apply: (service: ScheduledTaskService.ScheduledTaskService["Service"], id: string) => + service.setEnabled({ id: id as never, enabled: false }), + }, + { + change: "switched to an interval trigger", + reason: "The task's trigger changed before this delivery ran.", + apply: (service: ScheduledTaskService.ScheduledTaskService["Service"]) => + webhookTaskInput({ schedule: { type: "interval", everyMs: 3_600_000 } }).pipe( + Effect.flatMap(service.upsert), + ), + }, +] as const; + +it.effect.each(queuedDeliveryCases)( + "a delivery queued behind a run does not start once the task is $change", + ({ reason, apply }) => + Effect.gen(function* () { + const gate = yield* Deferred.make(); + yield* withService( + ({ service, launches }) => + Effect.gen(function* () { + const { task } = yield* service.upsert(yield* webhookTaskInput()); + yield* service.triggerWebhook(requestFor(task)); + const queued = yield* service.triggerWebhook(requestFor(task)); + yield* Queue.take(launches); + yield* apply(service, task.id); + yield* Deferred.succeed(gate, undefined); + // The queued delivery is marked failed instead of launching. + const deliveryId = queued._tag === "accepted" ? queued.deliveryId : undefined; + let delivery = (yield* service.getWebhookDelivery({ + id: task.id, + deliveryId: deliveryId!, + })).delivery; + while (delivery.outcome === "accepted") { + yield* Effect.yieldNow; + delivery = (yield* service.getWebhookDelivery({ + id: task.id, + deliveryId: deliveryId!, + })).delivery; + } + assert.equal(delivery.outcome, "dispatch_failed"); + assert.equal(delivery.error, reason); + assert.equal(yield* Queue.size(launches), 0); + }), + { gate }, + ); + }), +); + +it.effect("a save without a secret keeps a secret changed after it was read", () => + withService(({ service }) => + Effect.gen(function* () { + const signature = { header: "x-hub-signature-256", encoding: "hex", prefix: "sha256=" }; + const { task } = yield* service.upsert( + yield* webhookTaskInput({ + schedule: { type: "webhook", signature: { ...signature, secret: "old" } }, + }), + ); + yield* service.upsert( + yield* webhookTaskInput({ + schedule: { type: "webhook", signature: { ...signature, secret: "new" } }, + }), + ); + // A form opened before the change saves without sending a secret. + yield* service.upsert( + yield* webhookTaskInput({ title: "Edited", schedule: { type: "webhook", signature } }), + ); + const sign = (secret: string) => + `sha256=${NodeCrypto.createHmac("sha256", secret).update(pullRequestBody).digest("hex")}`; + const withSignature = (secret: string) => + requestFor(task, { + headers: { "content-type": "application/json", "x-hub-signature-256": sign(secret) }, + }); + assert.equal( + (yield* service.triggerWebhook(withSignature("old")))._tag, + "rejected_signature", + ); + assert.equal((yield* service.triggerWebhook(withSignature("new")))._tag, "accepted"); + }), + ), +); + +it.effect("caps deliveries waiting behind a stuck run", () => + Effect.gen(function* () { + const gate = yield* Deferred.make(); + yield* withService( + ({ service, launches }) => + Effect.gen(function* () { + const { task } = yield* service.upsert(yield* webhookTaskInput()); + yield* service.triggerWebhook(requestFor(task)); + yield* Queue.take(launches); + // The cap counts the running delivery too: 19 more wait, the next is refused. + const waiting = yield* Effect.forEach(Array.from({ length: 20 }), () => + service.triggerWebhook(requestFor(task)), + ); + assert.equal(waiting.filter((result) => result._tag === "accepted").length, 19); + assert.equal(waiting.at(-1)?._tag, "rate_limited"); + // The refused request is not logged. + const logged = (yield* service.listWebhookDeliveries({ id: task.id })).deliveries; + assert.equal(logged.length, 20); + yield* Deferred.succeed(gate, undefined); + }), + { gate }, + ); + }), +); + +it.effect("logs a body's first 64 KiB by bytes, not characters", () => + withService(({ service }) => + Effect.gen(function* () { + const { task } = yield* service.upsert(yield* webhookTaskInput({ enabled: false })); + // 30 000 three-byte characters: under 64 Ki characters, over 64 KiB. + const text = "界".repeat(30_000); + const body = new TextEncoder().encode(text); + yield* service.triggerWebhook(requestFor(task, { body, bodyText: text })); + const [summary] = (yield* service.listWebhookDeliveries({ id: task.id })).deliveries; + const { delivery } = yield* service.getWebhookDelivery({ + id: task.id, + deliveryId: summary!.id, + }); + assert.isTrue(delivery.bodyTruncated); + assert.isAtMost(new TextEncoder().encode(delivery.body).byteLength, 64 * 1024 + 3); + }), + ), +); + +it.effect("deleting a task removes its delivery log", () => + withService(({ service }) => + Effect.gen(function* () { + const { task } = yield* service.upsert(yield* webhookTaskInput({ enabled: false })); + yield* service.triggerWebhook(requestFor(task)); + yield* service.delete({ id: task.id }); + assert.equal((yield* service.listWebhookDeliveries({ id: task.id })).deliveries.length, 0); + }), + ), +); diff --git a/apps/server/src/scheduledTasks/webhookRoute.test.ts b/apps/server/src/scheduledTasks/webhookRoute.test.ts new file mode 100644 index 000000000000..e6fb5ff7d1ed --- /dev/null +++ b/apps/server/src/scheduledTasks/webhookRoute.test.ts @@ -0,0 +1,147 @@ +import { describe, expect, it } from "@effect/vitest"; +import * as Effect from "effect/Effect"; +import * as Layer from "effect/Layer"; +import * as NodeServices from "@effect/platform-node/NodeServices"; +import * as Etag from "effect/http/Etag"; +import * as HttpPlatform from "effect/http/HttpPlatform"; +import * as HttpRouter from "effect/http/HttpRouter"; +import * as HttpApi from "effect/http-api/HttpApi"; +import * as HttpApiBuilder from "effect/http-api/HttpApiBuilder"; + +import { + EnvironmentHttpApi, + ScheduledTaskWebhookDeliveryId, + ScheduledTaskError, +} from "@t3tools/contracts"; +import { + ScheduledTaskService, + type WebhookTriggerRequest, + type WebhookTriggerResult, +} from "./ScheduledTaskService.ts"; +import { WEBHOOK_MAX_BODY_BYTES, webhookHttpApiLayer } from "./webhookRoute.ts"; + +class WebhookTestApi extends HttpApi.make("environment").add(EnvironmentHttpApi.groups.webhooks) {} + +const handlerFor = ( + trigger: ( + request: WebhookTriggerRequest, + ) => Effect.Effect, +) => + HttpRouter.toWebHandler( + HttpApiBuilder.layer(WebhookTestApi).pipe( + Layer.provide(webhookHttpApiLayer), + Layer.provide(Layer.mock(ScheduledTaskService)({ triggerWebhook: trigger })), + Layer.provide( + HttpPlatform.layer.pipe( + Layer.provideMerge(NodeServices.layer), + Layer.provideMerge(Etag.layerWeak), + ), + ), + Layer.provide(NodeServices.layer), + ), + { disableLogger: true }, + ); + +const post = ( + path: string, + body: string | Uint8Array, + headers: Record = {}, +) => new Request(`http://env.local${path}`, { method: "POST", body, headers }); + +describe("webhook route", () => { + it("passes the raw request to the service and answers 202 with the delivery id", async () => { + let received: WebhookTriggerRequest | undefined; + const { handler, dispose } = handlerFor((request) => { + received = request; + return Effect.succeed({ + _tag: "accepted", + deliveryId: ScheduledTaskWebhookDeliveryId.make("delivery:1"), + }); + }); + try { + const response = await handler( + post("/api/hooks/scheduled-task%3Ahook/tok?x=1", '{"a":1}', { + "Content-Type": "application/json", + "X-GitHub-Event": "push", + }), + ); + expect(response.status).toBe(202); + expect(await response.json()).toEqual({ deliveryId: "delivery:1" }); + expect(received?.hookId).toBe("scheduled-task:hook"); + expect(received?.token).toBe("tok"); + expect(received?.query).toBe("x=1"); + expect(received?.headers["x-github-event"]).toBe("push"); + expect(received?.bodyText).toBe('{"a":1}'); + } finally { + await dispose(); + } + }); + + it("maps service outcomes to status codes", async () => { + const cases: ReadonlyArray<[WebhookTriggerResult["_tag"], number]> = [ + ["not_found", 404], + ["rejected_signature", 401], + ["disabled", 409], + ["rate_limited", 429], + ]; + for (const [tag, status] of cases) { + const { handler, dispose } = handlerFor(() => + Effect.succeed({ _tag: tag } as WebhookTriggerResult), + ); + try { + expect((await handler(post("/api/hooks/id/tok", "{}"))).status).toBe(status); + } finally { + await dispose(); + } + } + }); + + it("rejects oversized bodies and malformed paths before reaching the service", async () => { + let calls = 0; + const { handler, dispose } = handlerFor(() => { + calls += 1; + return Effect.succeed({ _tag: "not_found" }); + }); + try { + const big = new Uint8Array(WEBHOOK_MAX_BODY_BYTES + 1); + expect((await handler(post("/api/hooks/id/tok", big))).status).toBe(413); + expect((await handler(post("/api/hooks/id", "{}"))).status).toBe(404); + expect((await handler(post("/api/hooks/id/tok/extra", "{}"))).status).toBe(404); + expect((await handler(post("/api/hooks/%E0/tok", "{}"))).status).toBe(404); + // No content-length: the reader cap must still apply. + const chunked = new ReadableStream({ + start(controller) { + for (let sent = 0; sent <= WEBHOOK_MAX_BODY_BYTES; sent += 64 * 1024) { + controller.enqueue(new Uint8Array(64 * 1024)); + } + controller.close(); + }, + }); + const streamed = await handler( + new Request("http://env.local/api/hooks/id/tok", { + method: "POST", + body: chunked, + // Node's fetch needs duplex for streamed bodies. + duplex: "half", + }), + ); + expect(streamed.status).toBe(413); + expect(calls).toBe(0); + } finally { + await dispose(); + } + }); + + it("hides service failures behind a 500", async () => { + const { handler, dispose } = handlerFor(() => + Effect.fail(new ScheduledTaskError({ message: "database locked" })), + ); + try { + const response = await handler(post("/api/hooks/id/tok", "{}")); + expect(response.status).toBe(500); + expect(await response.text()).not.toContain("database"); + } finally { + await dispose(); + } + }); +}); diff --git a/apps/server/src/scheduledTasks/webhookRoute.ts b/apps/server/src/scheduledTasks/webhookRoute.ts new file mode 100644 index 000000000000..63d2772c0c82 --- /dev/null +++ b/apps/server/src/scheduledTasks/webhookRoute.ts @@ -0,0 +1,106 @@ +import { EnvironmentHttpApi } from "@t3tools/contracts"; +import * as ByteSize from "effect/ByteSize"; +import * as Effect from "effect/Effect"; +import * as Option from "effect/Option"; +import * as HttpIncomingMessage from "effect/http/HttpIncomingMessage"; +import type * as HttpServerRequest from "effect/http/HttpServerRequest"; +import * as HttpServerResponse from "effect/http/HttpServerResponse"; +import * as HttpApiBuilder from "effect/http-api/HttpApiBuilder"; + +import * as ScheduledTaskService from "./ScheduledTaskService.ts"; + +/** Largest request body a webhook accepts. The relay enforces the same cap. */ +export const WEBHOOK_MAX_BODY_BYTES = 1024 * 1024; + +const json = (status: number, body: Record) => + HttpServerResponse.jsonUnsafe(body, { status }); + +/** + * Handles `/api/hooks/:hookId/:token` for every accepted method. The endpoint + * is raw so the signature is checked over the exact body bytes; the service + * checks the token and signature. It is reachable directly, over the managed + * tunnel, or through the relay's stable `/v1/hooks/...` URL. + */ +const handleWebhook = + (scheduledTasks: ScheduledTaskService.ScheduledTaskService["Service"]) => + ({ + params, + request, + }: { + readonly params: { readonly hookId: string; readonly token: string }; + readonly request: HttpServerRequest.HttpServerRequest; + }) => + Effect.gen(function* () { + const contentLength = Number(request.headers["content-length"] ?? "0"); + if (!Number.isFinite(contentLength) || contentLength > WEBHOOK_MAX_BODY_BYTES) { + return json(413, { error: "body_too_large" }); + } + // Chunked requests carry no content-length, so the reader itself is capped. + const body = yield* request.arrayBuffer.pipe( + Effect.map((buffer) => new Uint8Array(buffer)), + Effect.provideService( + HttpIncomingMessage.MaxBodySize, + ByteSize.bytes(WEBHOOK_MAX_BODY_BYTES), + ), + Effect.option, + ); + if (Option.isNone(body)) return json(413, { error: "body_too_large_or_unreadable" }); + if (body.value.byteLength > WEBHOOK_MAX_BODY_BYTES) { + return json(413, { error: "body_too_large" }); + } + + const headers: Record = {}; + for (const [name, value] of Object.entries(request.headers)) { + if (typeof value === "string") headers[name.toLowerCase()] = value; + } + const queryIndex = request.url.indexOf("?"); + + const result = yield* scheduledTasks + .triggerWebhook({ + hookId: params.hookId, + token: params.token, + method: request.method, + path: `${ScheduledTaskService.WEBHOOK_ROUTE_PREFIX}/${encodeURIComponent(params.hookId)}`, + query: queryIndex === -1 ? "" : request.url.slice(queryIndex + 1), + headers, + body: body.value, + bodyText: new TextDecoder().decode(body.value), + }) + .pipe( + Effect.catch((cause) => + Effect.logWarning("Webhook delivery failed").pipe( + Effect.annotateLogs({ hookId: params.hookId }), + Effect.andThen(Effect.logDebug("Webhook delivery failure cause", { cause })), + Effect.as({ _tag: "error" as const }), + ), + ), + ); + + switch (result._tag) { + case "accepted": + return json(202, { deliveryId: result.deliveryId }); + case "not_found": + return json(404, { error: "hook_not_found" }); + case "rejected_signature": + return json(401, { error: "invalid_signature" }); + case "disabled": + return json(409, { error: "hook_disabled" }); + case "rate_limited": + return json(429, { error: "rate_limited" }); + case "error": + return json(500, { error: "internal_error" }); + } + }); + +export const webhookHttpApiLayer = HttpApiBuilder.group( + EnvironmentHttpApi, + "webhooks", + Effect.fnUntraced(function* (handlers) { + const handler = handleWebhook(yield* ScheduledTaskService.ScheduledTaskService); + return handlers + .handleRaw("webhookPost", handler) + .handleRaw("webhookPut", handler) + .handleRaw("webhookPatch", handler) + .handleRaw("webhookGet", handler); + }), +); diff --git a/apps/server/src/scheduledTasks/webhookTemplate.test.ts b/apps/server/src/scheduledTasks/webhookTemplate.test.ts new file mode 100644 index 000000000000..c87fa8b064a2 --- /dev/null +++ b/apps/server/src/scheduledTasks/webhookTemplate.test.ts @@ -0,0 +1,79 @@ +import { assert, describe, it } from "@effect/vitest"; + +import { renderWebhookPrompt, type WebhookRequest } from "./webhookTemplate.ts"; + +const githubPullRequest: WebhookRequest = { + method: "POST", + path: "/api/hooks/task-1", + query: "source=github", + headers: { "content-type": "application/json", "x-github-event": "pull_request" }, + bodyText: JSON.stringify({ + action: "opened", + pull_request: { url: "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/org/repo/pull/45", number: 45, labels: ["bug"] }, + }), +}; + +describe("renderWebhookPrompt", () => { + it("sends only what the template names", () => { + const rendered = renderWebhookPrompt( + "Review this PR: {{body.pull_request.url}}", + githubPullRequest, + ); + assert.equal(rendered.prompt, "Review this PR: https://github.com/org/repo/pull/45"); + assert.deepEqual(rendered.missing, []); + }); + + it("reads headers case-insensitively, query parameters and array indexes", () => { + const rendered = renderWebhookPrompt( + "{{ headers.X-GitHub-Event }} #{{body.pull_request.number}} {{body.pull_request.labels.0}} {{query.source}}", + githubPullRequest, + ); + assert.equal(rendered.prompt, "pull_request #45 bug github"); + }); + + it("renders objects as JSON and the raw body and request on request", () => { + const rendered = renderWebhookPrompt( + "{{body.pull_request.labels}}|{{body}}", + githubPullRequest, + ); + assert.equal(rendered.prompt, `[\n "bug"\n]|${githubPullRequest.bodyText}`); + const request = renderWebhookPrompt("{{request}}", githubPullRequest).prompt; + assert.isTrue(request.startsWith("POST /api/hooks/task-1?source=github\n")); + assert.include(request, "x-github-event: pull_request"); + assert.isTrue(request.endsWith(githubPullRequest.bodyText)); + }); + + it("renders missing fields empty and reports them once", () => { + const rendered = renderWebhookPrompt( + "a={{body.nope}} b={{body.nope}} c={{headers.x-missing}} d={{unknown}}", + githubPullRequest, + ); + assert.equal(rendered.prompt, "a= b= c= d="); + assert.deepEqual(rendered.missing, ["body.nope", "headers.x-missing", "unknown"]); + }); + + it("addresses form-encoded bodies and leaves other bodies to {{body}}", () => { + const form = renderWebhookPrompt("{{body.text}} by {{body.user_name}}", { + ...githubPullRequest, + headers: { "content-type": "application/x-www-form-urlencoded" }, + bodyText: "text=deploy+prod&user_name=alice", + }); + assert.equal(form.prompt, "deploy prod by alice"); + + const plain = renderWebhookPrompt("{{body.field}}{{body}}", { + ...githubPullRequest, + headers: { "content-type": "text/plain" }, + bodyText: "build failed", + }); + assert.equal(plain.prompt, "build failed"); + assert.deepEqual(plain.missing, ["body.field"]); + }); + + it("does not resolve inherited object properties", () => { + const rendered = renderWebhookPrompt( + "{{body.constructor}}{{body.__proto__}}", + githubPullRequest, + ); + assert.equal(rendered.prompt, ""); + }); +}); diff --git a/apps/server/src/scheduledTasks/webhookTemplate.ts b/apps/server/src/scheduledTasks/webhookTemplate.ts new file mode 100644 index 000000000000..e365a32858d6 --- /dev/null +++ b/apps/server/src/scheduledTasks/webhookTemplate.ts @@ -0,0 +1,114 @@ +/** + * Renders a webhook task's prompt from the request that triggered it. The + * rendered prompt is the only thing the agent receives, so a user who wants + * the whole request writes `{{request}}` or `{{body}}` themselves. + * + * The syntax is deliberately just `{{path}}` lookups — no conditionals, + * defaults, filters or escaping — so it never grows into a template language: + * + * - `{{body.a.b.0}}` a field of a JSON or form-encoded body + * - `{{headers.name}}` a request header (case-insensitive) + * - `{{query.name}}` a URL query parameter + * - `{{body}}` the raw body text + * - `{{request}}` method, path, query, headers and body + * + * Strings and numbers render as text, objects and arrays as JSON. A path with + * no value renders empty and is reported in `missing`. + */ + +export interface WebhookRequest { + readonly method: string; + readonly path: string; + /** Raw query string without the leading `?`. */ + readonly query: string; + /** Lowercased header names. */ + readonly headers: Readonly>; + readonly bodyText: string; +} + +export interface RenderedWebhookPrompt { + readonly prompt: string; + readonly missing: ReadonlyArray; +} + +const PLACEHOLDER = /\{\{\s*([^{}]*?)\s*\}\}/g; + +function parseBody(request: WebhookRequest): unknown { + const contentType = (request.headers["content-type"] ?? "").toLowerCase(); + if (contentType.includes("application/x-www-form-urlencoded")) { + return Object.fromEntries(new URLSearchParams(request.bodyText)); + } + // Senders are inconsistent about content types, so any body that parses + // as JSON is addressable. + try { + return JSON.parse(request.bodyText) as unknown; + } catch { + return undefined; + } +} + +function lookup(root: unknown, segments: ReadonlyArray): unknown { + let current = root; + for (const segment of segments) { + if (current === null || typeof current !== "object") return undefined; + if (!Object.hasOwn(current, segment)) return undefined; + current = (current as Record)[segment]; + } + return current; +} + +function stringify(value: unknown): string { + if (typeof value === "string") return value; + if (typeof value === "number" || typeof value === "boolean") return String(value); + return JSON.stringify(value, null, 2); +} + +function formatWebhookRequest(request: WebhookRequest): string { + const headerLines = Object.entries(request.headers).map(([name, value]) => `${name}: ${value}`); + return [ + `${request.method} ${request.path}${request.query ? `?${request.query}` : ""}`, + ...headerLines, + "", + request.bodyText, + ].join("\n"); +} + +export function renderWebhookPrompt( + template: string, + request: WebhookRequest, +): RenderedWebhookPrompt { + const missing: string[] = []; + let body: { parsed: unknown } | undefined; + const parsedBody = () => (body ??= { parsed: parseBody(request) }).parsed; + + const resolve = (expression: string): unknown => { + const [root, ...segments] = expression.split("."); + switch (root) { + case "request": + return segments.length === 0 ? formatWebhookRequest(request) : undefined; + case "body": + return segments.length === 0 ? request.bodyText : lookup(parsedBody(), segments); + case "headers": + return segments.length === 0 + ? request.headers + : request.headers[segments.join(".").toLowerCase()]; + case "query": { + const params = new URLSearchParams(request.query); + if (segments.length === 0) return Object.fromEntries(params); + return params.get(segments.join(".")) ?? undefined; + } + default: + return undefined; + } + }; + + const prompt = template.replace(PLACEHOLDER, (_match, expression: string) => { + const value = resolve(expression); + if (value === undefined || value === null) { + if (!missing.includes(expression)) missing.push(expression); + return ""; + } + return stringify(value); + }); + return { prompt, missing }; +} diff --git a/apps/server/src/scheduledTasks/webhookVerification.test.ts b/apps/server/src/scheduledTasks/webhookVerification.test.ts new file mode 100644 index 000000000000..de0bf869c400 --- /dev/null +++ b/apps/server/src/scheduledTasks/webhookVerification.test.ts @@ -0,0 +1,59 @@ +import * as NodeCrypto from "node:crypto"; + +import { assert, describe, it } from "@effect/vitest"; + +import { verifyWebhookSignature } from "./webhookVerification.ts"; + +const body = new TextEncoder().encode('{"action":"opened"}'); +const secret = "shared-secret"; +const hmac = () => NodeCrypto.createHmac("sha256", secret).update(body); + +describe("verifyWebhookSignature", () => { + const github = { header: "x-hub-signature-256", encoding: "hex", prefix: "sha256=" } as const; + + it("accepts a GitHub-style hex signature with prefix", () => { + const headers = { "x-hub-signature-256": `sha256=${hmac().digest("hex")}` }; + assert.isTrue(verifyWebhookSignature({ signature: github, secret, headers, body })); + }); + + it("rejects a tampered body, a wrong secret, a missing header and a missing prefix", () => { + const valid = `sha256=${hmac().digest("hex")}`; + const tampered = new TextEncoder().encode('{"action":"closed"}'); + assert.isFalse( + verifyWebhookSignature({ + signature: github, + secret, + headers: { "x-hub-signature-256": valid }, + body: tampered, + }), + ); + assert.isFalse( + verifyWebhookSignature({ + signature: github, + secret: "other", + headers: { "x-hub-signature-256": valid }, + body, + }), + ); + assert.isFalse(verifyWebhookSignature({ signature: github, secret, headers: {}, body })); + assert.isFalse( + verifyWebhookSignature({ + signature: github, + secret, + headers: { "x-hub-signature-256": hmac().digest("hex") }, + body, + }), + ); + }); + + it("accepts a base64 signature without prefix", () => { + assert.isTrue( + verifyWebhookSignature({ + signature: { header: "X-Signature", encoding: "base64", prefix: "" }, + secret, + headers: { "x-signature": hmac().digest("base64") }, + body, + }), + ); + }); +}); diff --git a/apps/server/src/scheduledTasks/webhookVerification.ts b/apps/server/src/scheduledTasks/webhookVerification.ts new file mode 100644 index 000000000000..6bda529e4834 --- /dev/null +++ b/apps/server/src/scheduledTasks/webhookVerification.ts @@ -0,0 +1,36 @@ +import * as NodeCrypto from "node:crypto"; + +import type { ScheduledTaskWebhookSignature } from "@t3tools/contracts"; + +/** Constant-time string comparison that does not leak length through timing. */ +export function constantTimeEquals(a: string, b: string): boolean { + const digestA = NodeCrypto.createHash("sha256").update(a).digest(); + const digestB = NodeCrypto.createHash("sha256").update(b).digest(); + return NodeCrypto.timingSafeEqual(digestA, digestB); +} + +/** + * Checks an HMAC-SHA256 signature over the raw body bytes, as GitHub, Linear, + * Shopify and most other signing senders do. The header value is + * ``, with the digest hex- or base64-encoded. + */ +export function verifyWebhookSignature(input: { + readonly signature: ScheduledTaskWebhookSignature; + readonly secret: string; + readonly headers: Readonly>; + readonly body: Uint8Array; +}): boolean { + const received = input.headers[input.signature.header.toLowerCase()]; + if (received === undefined) return false; + const value = received.trim(); + const prefix = input.signature.prefix; + if (prefix !== "" && !value.toLowerCase().startsWith(prefix.toLowerCase())) return false; + const digest = NodeCrypto.createHmac("sha256", input.secret).update(input.body).digest(); + const expected = + input.signature.encoding === "hex" ? digest.toString("hex") : digest.toString("base64"); + const candidate = value.slice(prefix.length); + return constantTimeEquals( + input.signature.encoding === "hex" ? candidate.toLowerCase() : candidate, + expected, + ); +} diff --git a/apps/server/src/server.ts b/apps/server/src/server.ts index 4f7db2a7141e..c203fe65c904 100644 --- a/apps/server/src/server.ts +++ b/apps/server/src/server.ts @@ -15,6 +15,7 @@ import * as Cause from "effect/Cause"; import * as Duration from "effect/Duration"; import * as Deferred from "effect/Deferred"; import * as Effect from "effect/Effect"; +import * as Option from "effect/Option"; import * as Layer from "effect/Layer"; import * as Stream from "effect/Stream"; import * as Schedule from "effect/Schedule"; @@ -117,6 +118,9 @@ import * as RemoteOpenTargets from "./environment/RemoteOpenTargets.ts"; import { authHttpApiLayer, environmentAuthenticatedAuthLayer } from "./auth/http.ts"; import * as ReplayMarkers from "./auth/replayMarkers.ts"; import * as ServerSecretStore from "./auth/ServerSecretStore.ts"; +import { webhookHttpApiLayer } from "./scheduledTasks/webhookRoute.ts"; +import { ScheduledTaskWebhookOrigin } from "./scheduledTasks/ScheduledTaskService.ts"; +import { CLOUD_ENDPOINT_RUNTIME_CONFIG, RELAY_URL_SECRET } from "./cloud/config.ts"; import * as EnvironmentAuth from "./auth/EnvironmentAuth.ts"; import { connectHttpApiLayer, @@ -444,7 +448,32 @@ const CloudManagedEndpointRuntimeLive = Layer.mergeAll( ), ); +// Webhook URLs go through the relay only when the managed tunnel it forwards +// to is configured; otherwise clients show the environment-relative path. +const ScheduledTaskWebhookOriginLive = Layer.effect( + ScheduledTaskWebhookOrigin, + Effect.map( + Effect.all([ServerEnvironment.ServerEnvironment, ServerSecretStore.ServerSecretStore]), + ([environment, secrets]) => + // The reference holds an effect so each read sees the current link state. + Effect.gen(function* () { + const [relayUrl, tunnelConfig] = yield* Effect.all([ + secrets.get(RELAY_URL_SECRET), + secrets.get(CLOUD_ENDPOINT_RUNTIME_CONFIG), + ]).pipe(Effect.orElseSucceed(() => [Option.none(), Option.none()] as const)); + return { + environmentId: yield* environment.getEnvironmentId, + relayUrl: + Option.isSome(relayUrl) && Option.isSome(tunnelConfig) + ? new TextDecoder().decode(relayUrl.value) || null + : null, + }; + }), + ), +); + const OrchestrationV2RuntimeLayerLive = OrchestrationV2ProductionLayerLive.pipe( + Layer.provide(ScheduledTaskWebhookOriginLive), Layer.provide(ProviderEventIngestor.analyticsLive), Layer.provide(CheckpointStoreLayerLive), Layer.provide(GitWorkflowLayerLive), @@ -643,6 +672,7 @@ const makeRoutesLayer = Layer.mergeAll( Layer.provide(pullRequestHttpApiLayer), Layer.provide(projectHttpApiLayer), Layer.provide(serverEnvironmentHttpApiLayer), + Layer.provide(webhookHttpApiLayer), Layer.provide(environmentAuthenticatedAuthLayer), ), otlpTracesProxyRouteLayer, diff --git a/apps/server/src/ws.ts b/apps/server/src/ws.ts index f47594e08ff9..36cbba02ff0a 100644 --- a/apps/server/src/ws.ts +++ b/apps/server/src/ws.ts @@ -27,7 +27,9 @@ import { CommandId, AuthAccessStreamError, type AuthAccessStreamEvent, + AuthOrchestrationOperateScope, type AuthEnvironmentScope, + type ScheduledTaskListResult, AuthSessionId, ClientConnectionMethod, ClientDeviceType, @@ -1315,6 +1317,14 @@ const makeWsRpcLayer = ( const processResourceMonitor = yield* ProcessResourceMonitor.ProcessResourceMonitor; const resourceTelemetry = yield* ResourceTelemetry.ResourceTelemetry; const relayClient = yield* RelayClient.RelayClient; + // A webhook URL starts agent runs, so only sessions that may operate + // see it; read-only sessions still see the task itself. + const withVisibleWebhookUrls = (result: ScheduledTaskListResult): ScheduledTaskListResult => + currentSession.scopes.includes(AuthOrchestrationOperateScope) + ? result + : { + tasks: result.tasks.map(({ webhook: _webhook, ...task }) => task), + }; // RpcScopeAuthorization checks each RPC's declared scope before its handler // runs. This covers the one RPC whose scope depends on its input. const authorizeEffect = ( @@ -2054,13 +2064,17 @@ const makeWsRpcLayer = ( }, ), [WS_METHODS.scheduledTasksList]: (_input) => - observeRpcEffect(WS_METHODS.scheduledTasksList, scheduledTasks.list(), { - "rpc.aggregate": "scheduledTasks", - }), + observeRpcEffect( + WS_METHODS.scheduledTasksList, + scheduledTasks.list().pipe(Effect.map(withVisibleWebhookUrls)), + { "rpc.aggregate": "scheduledTasks" }, + ), [WS_METHODS.scheduledTasksSubscribe]: (_input) => - observeRpcStream(WS_METHODS.scheduledTasksSubscribe, scheduledTasks.subscribeList(), { - "rpc.aggregate": "scheduledTasks", - }), + observeRpcStream( + WS_METHODS.scheduledTasksSubscribe, + scheduledTasks.subscribeList().pipe(Stream.map(withVisibleWebhookUrls)), + { "rpc.aggregate": "scheduledTasks" }, + ), [WS_METHODS.scheduledTasksUpsert]: (input) => observeRpcEffect(WS_METHODS.scheduledTasksUpsert, scheduledTasks.upsert(input), { "rpc.aggregate": "scheduledTasks", @@ -2080,6 +2094,24 @@ const makeWsRpcLayer = ( "rpc.aggregate": "scheduledTasks", "scheduled_task.id": input.id, }), + [WS_METHODS.scheduledTasksRotateWebhookToken]: (input) => + observeRpcEffect( + WS_METHODS.scheduledTasksRotateWebhookToken, + scheduledTasks.rotateWebhookToken(input), + { "rpc.aggregate": "scheduledTasks", "scheduled_task.id": input.id }, + ), + [WS_METHODS.scheduledTasksListWebhookDeliveries]: (input) => + observeRpcEffect( + WS_METHODS.scheduledTasksListWebhookDeliveries, + scheduledTasks.listWebhookDeliveries(input), + { "rpc.aggregate": "scheduledTasks", "scheduled_task.id": input.id }, + ), + [WS_METHODS.scheduledTasksGetWebhookDelivery]: (input) => + observeRpcEffect( + WS_METHODS.scheduledTasksGetWebhookDelivery, + scheduledTasks.getWebhookDelivery(input), + { "rpc.aggregate": "scheduledTasks", "scheduled_task.id": input.id }, + ), [WS_METHODS.serverProbe]: (_input) => observeRpcEffect(WS_METHODS.serverProbe, Effect.succeed({}), { "rpc.aggregate": "server", diff --git a/apps/web/src/components/chat/ThreadAutomationsPanel.tsx b/apps/web/src/components/chat/ThreadAutomationsPanel.tsx index b53972f0116e..13fde2d2e38f 100644 --- a/apps/web/src/components/chat/ThreadAutomationsPanel.tsx +++ b/apps/web/src/components/chat/ThreadAutomationsPanel.tsx @@ -194,7 +194,12 @@ export function ThreadAutomationsPanel(props: { variant="ghost" part="icon" aria-label={`Run ${task.title} now`} - disabled={busyTaskId !== null || task.lastRunStatus === "running"} + // A webhook task runs from its URL; there is no request to run it with. + disabled={ + busyTaskId !== null || + task.lastRunStatus === "running" || + task.schedule.type === "webhook" + } onClick={() => void runNow(task)} > diff --git a/apps/web/src/components/settings/ScheduledTasksSettings.tsx b/apps/web/src/components/settings/ScheduledTasksSettings.tsx index 15380c18600b..64ba04cf2b3c 100644 --- a/apps/web/src/components/settings/ScheduledTasksSettings.tsx +++ b/apps/web/src/components/settings/ScheduledTasksSettings.tsx @@ -162,6 +162,7 @@ function scheduleFromDraft(draft: DraftState): ScheduledTaskSchedule { } export function scheduleLabel(schedule: ScheduledTaskSchedule): string { + if (schedule.type === "webhook") return "On webhook"; if (schedule.type === "interval") { const minutes = schedule.everyMs / 60_000; return Number.isInteger(minutes) @@ -312,7 +313,7 @@ function ScheduledTaskEnvironmentSection({ const linkedTask = tasks?.find((task) => task.id === taskId); const openedLink = useRef(false); useEffect(() => { - if (!openedLink.current && linkedTask) { + if (!openedLink.current && linkedTask && linkedTask.schedule.type !== "webhook") { openedLink.current = true; onEdit(environment.environmentId, linkedTask); } @@ -346,6 +347,12 @@ function ScheduledTaskEnvironmentSection({ description="This task no longer exists or is outside the selected project scope." role="status" /> + ) : linkedTask?.schedule.type === "webhook" ? ( + ) : null} {tasks.length === 0 ? ( - + {/* Webhook tasks are not editable here yet; saving would drop their URL. */} + Edit - void act("run")}> + void act("run")} disabled={task.schedule.type === "webhook"}> Run now diff --git a/packages/client-runtime/src/state/server.ts b/packages/client-runtime/src/state/server.ts index a41b64af3be0..86f947894827 100644 --- a/packages/client-runtime/src/state/server.ts +++ b/packages/client-runtime/src/state/server.ts @@ -1285,6 +1285,20 @@ export function createServerEnvironmentAtoms( label: "environment-data:server:scheduled-task:run-now", tag: WS_METHODS.scheduledTasksRunNow, }), + rotateScheduledTaskWebhookToken: createEnvironmentRpcCommand(runtime, { + label: "environment-data:server:scheduled-task:rotate-webhook-token", + tag: WS_METHODS.scheduledTasksRotateWebhookToken, + scheduler: configScheduler, + concurrency: configConcurrency, + }), + listScheduledTaskWebhookDeliveries: createEnvironmentRpcCommand(runtime, { + label: "environment-data:server:scheduled-task:list-webhook-deliveries", + tag: WS_METHODS.scheduledTasksListWebhookDeliveries, + }), + getScheduledTaskWebhookDelivery: createEnvironmentRpcCommand(runtime, { + label: "environment-data:server:scheduled-task:get-webhook-delivery", + tag: WS_METHODS.scheduledTasksGetWebhookDelivery, + }), refreshUsageRates: createEnvironmentRpcCommand(runtime, { label: "environment-data:server:refresh-usage-rates", tag: WS_METHODS.serverRefreshUsageRates, diff --git a/packages/contracts/src/environmentHttp.ts b/packages/contracts/src/environmentHttp.ts index 3e62fb2685a9..411b10574b02 100644 --- a/packages/contracts/src/environmentHttp.ts +++ b/packages/contracts/src/environmentHttp.ts @@ -5,6 +5,7 @@ import * as HttpApi from "effect/http-api/HttpApi"; import * as HttpApiEndpoint from "effect/http-api/HttpApiEndpoint"; import * as HttpApiGroup from "effect/http-api/HttpApiGroup"; import * as HttpApiMiddleware from "effect/http-api/HttpApiMiddleware"; +import * as HttpApiSchema from "effect/http-api/HttpApiSchema"; import * as HttpServerRespondable from "effect/http/HttpServerRespondable"; import * as HttpServerResponse from "effect/http/HttpServerResponse"; @@ -651,10 +652,35 @@ class EnvironmentConnectHttpApi extends HttpApiGroup.make("connect") }), ) {} +/** + * Public entry point for webhook tasks. Unauthenticated by design: the token + * in the path, and an optional body signature, are the credential. The handler + * reads the raw body itself so a signature is checked over the exact bytes. + */ +const WebhookParams = Schema.Struct({ + hookId: TrimmedNonEmptyString, + token: TrimmedNonEmptyString, +}); +const WebhookAccepted = Schema.Struct({ deliveryId: TrimmedNonEmptyString }).pipe( + HttpApiSchema.status(202), +); +const webhookEndpoint = { + params: WebhookParams, + success: WebhookAccepted, +} as const; +const WEBHOOK_PATH = "/api/hooks/:hookId/:token"; + +class EnvironmentWebhooksHttpApi extends HttpApiGroup.make("webhooks") + .add(HttpApiEndpoint.post("webhookPost", WEBHOOK_PATH, webhookEndpoint)) + .add(HttpApiEndpoint.put("webhookPut", WEBHOOK_PATH, webhookEndpoint)) + .add(HttpApiEndpoint.patch("webhookPatch", WEBHOOK_PATH, webhookEndpoint)) + .add(HttpApiEndpoint.get("webhookGet", WEBHOOK_PATH, webhookEndpoint)) {} + export class EnvironmentHttpApi extends HttpApi.make("environment") .add(EnvironmentMetadataHttpApi) .add(EnvironmentAuthHttpApi) .add(EnvironmentOrchestrationHttpApi) .add(EnvironmentPullRequestsHttpApi) .add(EnvironmentProjectsHttpApi) - .add(EnvironmentConnectHttpApi) {} + .add(EnvironmentConnectHttpApi) + .add(EnvironmentWebhooksHttpApi) {} diff --git a/packages/contracts/src/orchestratorMcp.ts b/packages/contracts/src/orchestratorMcp.ts index bed1e7a782e2..776cf4be52d5 100644 --- a/packages/contracts/src/orchestratorMcp.ts +++ b/packages/contracts/src/orchestratorMcp.ts @@ -64,7 +64,7 @@ const OrchestratorMcpSchedule = Schema.Union([ OrchestratorMcpScheduleFromJsonString, ]).annotate({ description: - "Recurring schedule object: {type:'interval', everyMs} or {type:'fixed_time', timeOfDay, weekdays?}. Never stringify it unless the provider requires the compatibility form.", + "Trigger object: {type:'interval', everyMs}, {type:'fixed_time', timeOfDay, weekdays?}, or {type:'webhook'} to run on each request to a generated URL. Never stringify it unless the provider requires the compatibility form.", }); /** @@ -541,6 +541,8 @@ export const OrchestratorMcpScheduledTask = Schema.Struct({ schedule: ScheduledTaskSchedule, nextRunAt: Schema.NullOr(IsoDateTime), lastRunStatus: ScheduledTaskRunStatus, + /** For webhook tasks: the public T3 Connect URL, or the environment-relative path when the environment is not linked. */ + webhookUrl: Schema.optional(Schema.String), }); export type OrchestratorMcpScheduledTask = typeof OrchestratorMcpScheduledTask.Type; diff --git a/packages/contracts/src/rpc.ts b/packages/contracts/src/rpc.ts index c9e2e6f03ee1..575ad6e32517 100644 --- a/packages/contracts/src/rpc.ts +++ b/packages/contracts/src/rpc.ts @@ -310,6 +310,11 @@ import { ScheduledTaskListInput, ScheduledTaskListResult, ScheduledTaskRunNowInput, + ScheduledTaskRotateWebhookTokenInput, + ScheduledTaskListWebhookDeliveriesInput, + ScheduledTaskListWebhookDeliveriesResult, + ScheduledTaskGetWebhookDeliveryInput, + ScheduledTaskGetWebhookDeliveryResult, ScheduledTaskRunNowResult, ScheduledTaskSetEnabledInput, ScheduledTaskUpsertInput, @@ -474,6 +479,9 @@ export const WS_METHODS = { scheduledTasksSetEnabled: "scheduledTasks.setEnabled", scheduledTasksDelete: "scheduledTasks.delete", scheduledTasksRunNow: "scheduledTasks.runNow", + scheduledTasksRotateWebhookToken: "scheduledTasks.rotateWebhookToken", + scheduledTasksListWebhookDeliveries: "scheduledTasks.listWebhookDeliveries", + scheduledTasksGetWebhookDelivery: "scheduledTasks.getWebhookDelivery", // Cloud environment methods cloudGetRelayClientStatus: "cloud.getRelayClientStatus", @@ -1668,6 +1676,33 @@ const WsScheduledTasksRunNowRpc = Rpc.make(WS_METHODS.scheduledTasksRunNow, { error: Schema.Union([ScheduledTaskError, EnvironmentAuthorizationError]), }); +const WsScheduledTasksRotateWebhookTokenRpc = Rpc.make( + WS_METHODS.scheduledTasksRotateWebhookToken, + { + payload: ScheduledTaskRotateWebhookTokenInput, + success: ScheduledTaskMutationResult, + error: Schema.Union([ScheduledTaskError, EnvironmentAuthorizationError]), + }, +); + +const WsScheduledTasksListWebhookDeliveriesRpc = Rpc.make( + WS_METHODS.scheduledTasksListWebhookDeliveries, + { + payload: ScheduledTaskListWebhookDeliveriesInput, + success: ScheduledTaskListWebhookDeliveriesResult, + error: Schema.Union([ScheduledTaskError, EnvironmentAuthorizationError]), + }, +); + +const WsScheduledTasksGetWebhookDeliveryRpc = Rpc.make( + WS_METHODS.scheduledTasksGetWebhookDelivery, + { + payload: ScheduledTaskGetWebhookDeliveryInput, + success: ScheduledTaskGetWebhookDeliveryResult, + error: Schema.Union([ScheduledTaskError, EnvironmentAuthorizationError]), + }, +); + const WsSubscribeAuthAccessRpc = Rpc.make(WS_METHODS.subscribeAuthAccess, { payload: Schema.Struct({}), success: AuthAccessStreamEvent, @@ -1753,6 +1788,9 @@ export const WsRpcGroup = RpcGroup.make( WsScheduledTasksSetEnabledRpc, WsScheduledTasksDeleteRpc, WsScheduledTasksRunNowRpc, + WsScheduledTasksRotateWebhookTokenRpc, + WsScheduledTasksListWebhookDeliveriesRpc, + WsScheduledTasksGetWebhookDeliveryRpc, WsServerReportClientActivityRpc, WsServerReportHostPowerStateRpc, WsServerGetBackgroundPolicyRpc, diff --git a/packages/contracts/src/scheduledTask.ts b/packages/contracts/src/scheduledTask.ts index 068fe7542236..f00ce5da93b8 100644 --- a/packages/contracts/src/scheduledTask.ts +++ b/packages/contracts/src/scheduledTask.ts @@ -2,6 +2,7 @@ import * as Schema from "effect/Schema"; import { CommandId, + ForwardCompatibleArray, IsoDateTime, ProjectId, ScheduledTaskId, @@ -54,6 +55,55 @@ const ScheduledTaskFixedTimeSchedule = Schema.Struct({ description: "Run at a fixed local wall-clock time on selected weekdays.", }); +const ScheduledTaskWebhookSignatureFields = { + header: TrimmedNonEmptyString.annotate({ + description: "Request header carrying the signature, such as x-hub-signature-256.", + }), + encoding: Schema.Literals(["hex", "base64"]).annotate({ + description: "How the HMAC-SHA256 digest is encoded in the header.", + }), + prefix: Schema.String.annotate({ + description: "Text before the digest in the header value, such as 'sha256='. Empty for none.", + }), +}; + +/** HMAC-SHA256 over the raw request body. The secret is never part of the read model. */ +export const ScheduledTaskWebhookSignature = Schema.Struct( + ScheduledTaskWebhookSignatureFields, +).annotate({ description: "Optional HMAC-SHA256 signature check over the raw request body." }); +export type ScheduledTaskWebhookSignature = typeof ScheduledTaskWebhookSignature.Type; + +const ScheduledTaskWebhookSchedule = Schema.Struct({ + type: Schema.Literal("webhook").annotate({ + description: "Run when the task's webhook URL receives a request.", + }), + signature: Schema.NullOr(ScheduledTaskWebhookSignature), +}).annotate({ + description: + "Run on each request to the task's webhook URL. The prompt may use {{body.path}}, {{headers.name}}, {{query.name}}, {{body}} and {{request}} placeholders.", +}); + +const ScheduledTaskUpsertWebhookSchedule = Schema.Struct({ + type: Schema.Literal("webhook").annotate({ + description: "Run when the task's webhook URL receives a request.", + }), + signature: Schema.optional( + Schema.NullOr( + Schema.Struct({ + ...ScheduledTaskWebhookSignatureFields, + secret: Schema.optional(TrimmedNonEmptyString).annotate({ + description: "Shared signing secret. Omit to keep the stored secret.", + }), + }), + ), + ).annotate({ + description: "Signature check; omit or null to accept requests by URL token only.", + }), +}).annotate({ + description: + "Run on each request to the task's webhook URL. The prompt may use {{body.path}}, {{headers.name}}, {{query.name}}, {{body}} and {{request}} placeholders.", +}); + /** * Read model for persisted schedules. Keep accepting legacy sub-minute rows so * users can list, disable, edit, or delete them after the write minimum changes. @@ -61,9 +111,10 @@ const ScheduledTaskFixedTimeSchedule = Schema.Struct({ export const ScheduledTaskSchedule = Schema.Union([ ScheduledTaskIntervalSchedule, ScheduledTaskFixedTimeSchedule, + ScheduledTaskWebhookSchedule, ]).annotate({ description: - "Structured recurring schedule. Pass an object with type 'interval' or 'fixed_time'.", + "Structured trigger. Pass an object with type 'interval', 'fixed_time' or 'webhook'.", }); export type ScheduledTaskSchedule = typeof ScheduledTaskSchedule.Type; @@ -82,14 +133,25 @@ export const ScheduledTaskUpsertSchedule = Schema.Union([ description: "Run repeatedly after a fixed number of milliseconds.", }), ScheduledTaskFixedTimeSchedule, + ScheduledTaskUpsertWebhookSchedule, ]).annotate({ - description: "Writable recurring schedule. Pass an object with type 'interval' or 'fixed_time'.", + description: "Writable trigger. Pass an object with type 'interval', 'fixed_time' or 'webhook'.", }); export type ScheduledTaskUpsertSchedule = typeof ScheduledTaskUpsertSchedule.Type; export const ScheduledTaskRunStatus = Schema.Literals(["never", "running", "succeeded", "failed"]); export type ScheduledTaskRunStatus = typeof ScheduledTaskRunStatus.Type; +/** Where a webhook task receives requests. Present only on webhook tasks. */ +export const ScheduledTaskWebhookEndpoint = Schema.Struct({ + /** Environment-relative path including the secret token; works on any origin that reaches the environment. */ + path: TrimmedNonEmptyString, + /** Public T3 Connect URL, or null when the environment is not linked to T3 Connect. */ + url: Schema.NullOr(TrimmedNonEmptyString), + hasSecret: Schema.Boolean, +}); +export type ScheduledTaskWebhookEndpoint = typeof ScheduledTaskWebhookEndpoint.Type; + export const ScheduledTask = Schema.Struct({ id: ScheduledTaskId, title: TrimmedNonEmptyString, @@ -111,6 +173,7 @@ export const ScheduledTask = Schema.Struct({ lastRunStatus: ScheduledTaskRunStatus, lastRunError: Schema.NullOr(Schema.String), runCount: Schema.Int.check(Schema.isGreaterThanOrEqualTo(0)), + webhook: Schema.optional(ScheduledTaskWebhookEndpoint), }); export type ScheduledTask = typeof ScheduledTask.Type; @@ -118,7 +181,9 @@ export const ScheduledTaskListInput = Schema.Struct({}); export type ScheduledTaskListInput = typeof ScheduledTaskListInput.Type; export const ScheduledTaskListResult = Schema.Struct({ - tasks: Schema.Array(ScheduledTask), + // Trigger types grow over time; a client must not lose the whole list over + // one task it cannot decode. + tasks: ForwardCompatibleArray(ScheduledTask), }); export type ScheduledTaskListResult = typeof ScheduledTaskListResult.Type; @@ -160,6 +225,76 @@ export const ScheduledTaskRunNowInput = Schema.Struct({ }); export type ScheduledTaskRunNowInput = typeof ScheduledTaskRunNowInput.Type; +export const ScheduledTaskRotateWebhookTokenInput = Schema.Struct({ + id: ScheduledTaskId, +}); +export type ScheduledTaskRotateWebhookTokenInput = typeof ScheduledTaskRotateWebhookTokenInput.Type; + +export const ScheduledTaskWebhookDeliveryId = TrimmedNonEmptyString.pipe( + Schema.brand("ScheduledTaskWebhookDeliveryId"), +); +export type ScheduledTaskWebhookDeliveryId = typeof ScheduledTaskWebhookDeliveryId.Type; + +export const ScheduledTaskWebhookDeliveryOutcome = Schema.Literals([ + "accepted", + "dispatch_failed", + "rejected_signature", + "disabled", + "rate_limited", +]); +export type ScheduledTaskWebhookDeliveryOutcome = typeof ScheduledTaskWebhookDeliveryOutcome.Type; + +export const ScheduledTaskWebhookDeliverySummary = Schema.Struct({ + id: ScheduledTaskWebhookDeliveryId, + taskId: ScheduledTaskId, + receivedAt: IsoDateTime, + method: TrimmedNonEmptyString, + contentType: Schema.NullOr(Schema.String), + bodyBytes: Schema.Int.check(Schema.isGreaterThanOrEqualTo(0)), + outcome: ScheduledTaskWebhookDeliveryOutcome, + /** True when a configured signature matched; false when none was configured. */ + signatureVerified: Schema.Boolean, + /** Template placeholders that had no value in this request and rendered empty. */ + missingFields: Schema.Array(Schema.String), + error: Schema.NullOr(Schema.String), +}); +export type ScheduledTaskWebhookDeliverySummary = typeof ScheduledTaskWebhookDeliverySummary.Type; + +export const ScheduledTaskWebhookDelivery = Schema.Struct({ + ...ScheduledTaskWebhookDeliverySummary.fields, + query: Schema.String, + headers: Schema.Record(Schema.String, Schema.String), + /** Body as UTF-8 text, cut at the log limit; see bodyTruncated. */ + body: Schema.String, + bodyTruncated: Schema.Boolean, + renderedPrompt: Schema.NullOr(Schema.String), +}); +export type ScheduledTaskWebhookDelivery = typeof ScheduledTaskWebhookDelivery.Type; + +export const ScheduledTaskListWebhookDeliveriesInput = Schema.Struct({ + id: ScheduledTaskId, +}); +export type ScheduledTaskListWebhookDeliveriesInput = + typeof ScheduledTaskListWebhookDeliveriesInput.Type; + +export const ScheduledTaskListWebhookDeliveriesResult = Schema.Struct({ + deliveries: Schema.Array(ScheduledTaskWebhookDeliverySummary), +}); +export type ScheduledTaskListWebhookDeliveriesResult = + typeof ScheduledTaskListWebhookDeliveriesResult.Type; + +export const ScheduledTaskGetWebhookDeliveryInput = Schema.Struct({ + id: ScheduledTaskId, + deliveryId: ScheduledTaskWebhookDeliveryId, +}); +export type ScheduledTaskGetWebhookDeliveryInput = typeof ScheduledTaskGetWebhookDeliveryInput.Type; + +export const ScheduledTaskGetWebhookDeliveryResult = Schema.Struct({ + delivery: ScheduledTaskWebhookDelivery, +}); +export type ScheduledTaskGetWebhookDeliveryResult = + typeof ScheduledTaskGetWebhookDeliveryResult.Type; + export const ScheduledTaskMutationResult = Schema.Struct({ task: ScheduledTask, });