import { encodableSchema } from "@automate.ax/codec" import { automationInvocationOutputSchema, automationInvocationScheduleSchema, platformEmailAttachmentsSchema, } from "@automate.ax/api-contract/runtime" import { httpRequestTriggerDefinition } from "@automate.ax/catalog/triggers/core-http" import { getMailhookAddress } from "@automate.ax/catalog/triggers/core-mailhook" import type { DistributedOmit } from "type-fest" import * as z from "zod" import { defineAction, type ActionCallInput, type ActionObjectInput, } from "../../automation/actions" import { isAutomationDescriptor, normalizeAutomationSelector, type AutomationDescriptor, } from "../../automation/automation-descriptor" import { isSignal, signalContextOrigins, type Signal, } from "../../automation/signal-protocol" import { formatHookLocation, getAutomationInvocationEnvironment, getAutomationPlanningContext, } from "../../automation/runtime" import { userInvocationEntrypointSchema } from "./invocation" const HTTP_RESPONSE_OUTPUT_TYPE = "http.request.response" const JSON_CONTENT_TYPE = "application/json" const TEXT_CONTENT_TYPE = "text/plain;charset=UTF-8" const BINARY_CONTENT_TYPE = "application/octet-stream" const HTTP_METHOD_PATTERN = /^[!#$%&'*+\-.^_`|~0-9A-Za-z]+$/ // Node's Fetch implementation follows the WHATWG blocked-port list. const FETCH_BLOCKED_PORTS = new Set([ "1", "7", "9", "11", "13", "15", "17", "19", "20", "21", "22", "23", "25", "37", "42", "43", "53", "69", "77", "79", "87", "95", "101", "102", "103", "104", "109", "110", "111", "113", "115", "117", "119", "123", "135", "137", "139", "143", "161", "179", "389", "427", "465", "512", "513", "514", "515", "526", "530", "531", "532", "540", "548", "554", "556", "563", "587", "601", "636", "989", "990", "993", "995", "1719", "1720", "1723", "2049", "3659", "4045", "4190", "5060", "5061", "6000", "6566", "6665", "6666", "6667", "6668", "6669", "6679", "6697", "10080", ]) const FETCH_MANAGED_HEADER_NAMES = new Set([ "connection", "content-length", "expect", "host", "keep-alive", "sec-fetch-mode", "transfer-encoding", "upgrade", ]) const HTTP_QUERY_VALUE_SCHEMA = z.union([z.boolean(), z.number(), z.string()]) const FETCH_HTTP_INPUT_SCHEMA = z .object({ /** JSON value, text, or binary request body. */ body: z.union([z.string(), z.instanceof(Blob), z.json()]).optional(), /** HTTP request headers. */ headers: z.record(z.string(), z.string()).default({}), /** HTTP method. */ method: z .string() .trim() .min(1) .transform((method) => method.toUpperCase()) .refine( (method) => HTTP_METHOD_PATTERN.test(method) && method !== "CONNECT" && method !== "TRACE" && method !== "TRACK", "Invalid or unsupported HTTP method.", ) .default("GET"), /** Query parameters merged into the URL. */ query: z .record( z.string(), z .union([HTTP_QUERY_VALUE_SCHEMA, HTTP_QUERY_VALUE_SCHEMA.array()]) .nullish(), ) .default({}), /** Native fetch redirect behavior. */ redirect: z.enum(["error", "follow", "manual"]).default("follow"), /** How to read the response body. */ responseType: z.enum(["auto", "blob", "json", "text"]).default("auto"), /** Absolute HTTP or HTTPS URL. */ url: z.url({ protocol: /^https?$/ }).superRefine((url, context) => { const parsedUrl = new URL(url) if (parsedUrl.username !== "" || parsedUrl.password !== "") { context.addIssue({ code: "custom", message: "HTTP URLs cannot include credentials.", }) } if (FETCH_BLOCKED_PORTS.has(parsedUrl.port)) { context.addIssue({ code: "custom", message: "HTTP URL uses a port blocked by Fetch.", }) } }), }) .superRefine((input, context) => { try { // oxlint-disable-next-line eslint/no-new -- REVIEW: Construction is the native Fetch header validation boundary. new Headers(input.headers) } catch { context.addIssue({ code: "custom", message: "Invalid HTTP headers.", path: ["headers"], }) } const managedHeader = Object.keys(input.headers).find((name) => FETCH_MANAGED_HEADER_NAMES.has(name.toLowerCase()), ) if (managedHeader !== undefined) { context.addIssue({ code: "custom", message: `HTTP header "${managedHeader}" is managed by Fetch.`, path: ["headers", managedHeader], }) } if ( input.body !== undefined && (input.method === "GET" || input.method === "HEAD") ) { context.addIssue({ code: "custom", message: `${input.method} requests cannot include a body.`, path: ["body"], }) } }) const FETCH_HTTP_RESPONSE_SCHEMA = z .object({ /** Response headers normalized by the Fetch API. Set-Cookie is an array. */ headers: z.record(z.string(), z.union([z.string(), z.string().array()])), /** Whether the status is in the 200-299 range. */ ok: z.boolean(), /** Whether the request followed a redirect. */ redirected: z.boolean(), /** HTTP response status. */ status: z.number().int().min(100).max(599), /** HTTP response status text. */ statusText: z.string(), /** Final response URL after redirects. */ url: z.string(), }) .and( z.discriminatedUnion("bodyType", [ z.object({ body: z.instanceof(Blob), bodyType: z.literal("blob") }), z.object({ body: z.json(), bodyType: z.literal("json") }), z.object({ body: z.string(), bodyType: z.literal("text") }), z.object({ body: z.null(), bodyType: z.literal("empty") }), ]), ) const EMAIL_BODY_SCHEMA = z.object({ /** HTML message body. */ html: z.string().min(1).optional(), /** Markdown message body. */ markdown: z.string().min(1).optional(), /** Plain-text message body. */ text: z.string().min(1).optional(), }) const EMAIL_PLUS_PATH_SCHEMA = z .string() .min(1) .refine( (plusPath) => z.email().safeParse(`mailhook+${plusPath}@example.com`).success, "Invalid email plus path.", ) const SCOPED_EMAIL_REPLY_TO_SCHEMA = z.object({ /** Optional suffix appended through email plus addressing. */ plusPath: EMAIL_PLUS_PATH_SCHEMA.optional(), /** Managed mailbox shared by the current automation or project. */ scope: z.enum(["automation", "project"]), }) const SEND_EMAIL_INPUT_SCHEMA = z .object({ /** Files attached to the message. */ attachments: platformEmailAttachmentsSchema.optional(), /** Custom email headers. */ headers: z.record(z.string().min(1), z.string()).optional(), /** Address that should receive replies to this email. */ replyTo: z.email().optional(), /** Email subject line. */ subject: z.string(), /** Recipient email address, or `"*"` for all organization members. */ to: z.union([z.email(), z.literal("*")]).default("*"), }) .and( z.union([ EMAIL_BODY_SCHEMA.required({ html: true }), EMAIL_BODY_SCHEMA.required({ markdown: true }), EMAIL_BODY_SCHEMA.required({ text: true }), ]), ) type InvocationInputValue = T | Signal /** Result emitted by {@link fetchHttp}. */ export type FetchHttpResponse = z.output /** Enforces mutually exclusive absolute and relative invocation scheduling. */ type InvocationTimingInput = | { runAt: InvocationInputValue runIn?: never } | { runAt?: never runIn: InvocationInputValue } | { runAt?: never runIn?: never } export type InvokeAutomationInput = { automation: AutomationDescriptor | InvocationInputValue entrypoint?: InvocationInputValue payload: InvocationInputValue } & InvocationTimingInput const invokeAutomationAction = defineAction("Invoke automation") .input( z .object({ automation: z.string().min(1), entrypoint: userInvocationEntrypointSchema, payload: z.unknown().pipe(encodableSchema), }) .and(automationInvocationScheduleSchema), ) .output(automationInvocationOutputSchema) .retry({ replaySafety: "safe" }) .handler(async ({ input, runtime }) => await runtime.invokeAutomation(input)) /** * Invokes one entrypoint immediately or after an absolute/relative delay. * * Pass an automation descriptor, description, or ID in the current project. * Omit `entrypoint` to use `"default"`. Entrypoints beginning with `__$` are * reserved by Automate.ax. * * @param input - Typed payload, target automation, and optional delivery time. */ export function invokeAutomation( input: InvokeAutomationInput, ) { const { automation, ...invocation } = input return invokeAutomationAction({ ...invocation, automation: isAutomationDescriptor(automation) ? normalizeAutomationSelector(automation) : automation, }) } const sendEmailAction = defineAction("Send email") .input(SEND_EMAIL_INPUT_SCHEMA) .output( z.object({ /** * Resend internal email resource ID, or `null` when unavailable. This is * not an RFC 5322 Message-ID and is not a valid threading target. */ id: z.string().nullable(), }), ) .retry({ replaySafety: "unsafe" }) .handler(async ({ input, runtime }) => await runtime.sendEmail(input)) export type ScopedEmailReplyTo = Omit< z.input, "plusPath" > & { plusPath?: Signal | string } /** Inputs accepted by the built-in organization email action. */ export type SendEmailInput = DistributedOmit< ActionObjectInput, "replyTo" > & { replyTo?: ScopedEmailReplyTo | string | Signal } /** * Sends an email to one or every member of the project's organization. * * Provide at least one of `html`, `markdown`, or `text`. The platform derives * missing HTML and plain-text alternatives before delivery. * * The result's `id` is Resend's internal email resource ID, not the RFC 5322 * Message-ID used for threading. To reply to a received mailhook message, use * its `messageId` for `In-Reply-To`. Build `References` from its `references`, * or its `inReplyTo` when references are empty, followed by its `messageId`. * * @param input - Recipient, subject, attachments, and one or more message * bodies. * @throws When a managed Reply-To scope or plus path is invalid. */ export function sendEmail(input: SendEmailInput): Signal<{ /** * Resend internal email resource ID, or `null` when unavailable. This is not * an RFC 5322 Message-ID and is not a valid threading target. */ id: string | null }> { const { replyTo: replyToInput, ...email } = input if ( replyToInput === undefined || typeof replyToInput !== "object" || isSignal(replyToInput) ) { return sendEmailAction({ ...email, ...(replyToInput !== undefined && { replyTo: replyToInput }), }) } const environment = getAutomationInvocationEnvironment() if (!environment) { throw new Error( "Scoped email replies must be declared inside an automation invocation.", ) } const { plusPath, ...scopedReplyTo } = replyToInput const parsedReplyTo = SCOPED_EMAIL_REPLY_TO_SCHEMA.parse({ ...scopedReplyTo, ...(!isSignal(plusPath) && { plusPath }), }) return sendEmailAction({ ...email, replyTo: isSignal(plusPath) ? plusPath.transform((value) => getMailhookAddress({ ...environment, hookSlot: 0, ...parsedReplyTo, plusPath: EMAIL_PLUS_PATH_SCHEMA.parse(value), }), ) : getMailhookAddress({ ...environment, hookSlot: 0, ...parsedReplyTo, }), }) } /** * Makes a durable HTTP request to an arbitrary HTTP or HTTPS endpoint. * * JSON request bodies are serialized automatically. The response body is read * according to `responseType`, or inferred from its Content-Type with `auto`. * HTTP error statuses are returned with `ok: false`; network and body parsing * failures throw through the ordinary action retry lifecycle. * * @param input - URL, request options, and response parsing behavior. */ export const fetchHttp = defineAction("Fetch HTTP request") .input(FETCH_HTTP_INPUT_SCHEMA) .output(FETCH_HTTP_RESPONSE_SCHEMA) .retry({ replaySafety: ({ method }) => method === "GET" || method === "HEAD" || method === "OPTIONS" ? "safe" : "unsafe", }) .handler(async ({ input }) => { const url = new URL(input.url) for (const [name, value] of Object.entries(input.query)) { if (value == null) continue if (Array.isArray(value)) { url.searchParams.delete(name) for (const item of value) url.searchParams.append(name, String(item)) } else { url.searchParams.set(name, String(value)) } } const headers = new Headers(input.headers) if ( input.body !== undefined && typeof input.body !== "string" && !(input.body instanceof Blob) && !headers.has("content-type") ) { headers.set("content-type", JSON_CONTENT_TYPE) } const response = await fetch(url, { body: input.body === undefined || typeof input.body === "string" || input.body instanceof Blob ? input.body : JSON.stringify(input.body), headers, method: input.method, redirect: input.redirect, }) const responseBlob = await response.blob() const mediaType = response.headers .get("content-type") ?.split(";", 1)[0] ?.trim() .toLowerCase() const responseType = input.responseType !== "auto" ? input.responseType : mediaType === JSON_CONTENT_TYPE || mediaType?.endsWith("+json") ? "json" : mediaType?.startsWith("text/") || mediaType === "application/javascript" || mediaType === "application/x-www-form-urlencoded" || mediaType === "application/xml" || mediaType?.endsWith("+xml") ? "text" : "blob" // Consume the response stream once before combining its body with metadata. const bodyOutput = responseBlob.size === 0 ? { body: null, bodyType: "empty" as const } : responseType === "json" ? { body: z.json().parse(JSON.parse(await responseBlob.text())), bodyType: "json" as const, } : responseType === "text" ? { body: await responseBlob.text(), bodyType: "text" as const, } : { body: responseBlob, bodyType: "blob" as const } const responseHeaders: Record = Object.fromEntries(response.headers) const setCookies = response.headers.getSetCookie() if (setCookies.length > 0) responseHeaders["set-cookie"] = setCookies return { ...bodyOutput, headers: responseHeaders, ok: response.ok, redirected: response.redirected, status: response.status, statusText: response.statusText, url: response.url, } }) const RESPOND_TO_HTTP_REQUEST_INPUT_SCHEMA = z.object({ body: z .union([z.string(), z.instanceof(Blob), z.json()]) .nullable() .optional(), headers: z.record(z.string(), z.string()).default({}), requestId: z.string().min(1), status: z.number().int().min(100).max(599).optional(), }) const respondToHttpRequestAction = defineAction("Respond to HTTP request") .input(RESPOND_TO_HTTP_REQUEST_INPUT_SCHEMA) .output(z.void()) .retry({ replaySafety: "safe" }) .handler(async ({ input, runtime }) => { const contentType = input.body instanceof Blob ? input.body.type || BINARY_CONTENT_TYPE : typeof input.body === "string" ? TEXT_CONTENT_TYPE : input.body == null ? undefined : JSON_CONTENT_TYPE await runtime.sendOutput({ data: { body: input.body == null || typeof input.body === "string" || input.body instanceof Blob ? input.body : JSON.stringify(input.body), headers: contentType && !new Headers(input.headers).has("content-type") ? { ...input.headers, "content-type": contentType } : input.headers, requestId: input.requestId, status: input.status ?? (input.body == null ? 204 : 200), }, type: HTTP_RESPONSE_OUTPUT_TYPE, }) }) export type RespondToHttpRequestInput = ActionCallInput< typeof RESPOND_TO_HTTP_REQUEST_INPUT_SCHEMA > /** * Responds for an HTTP request occurrence and emits its stable request ID. * * Omitting `waitForResponse` on the originating trigger enables waiting when * its request ID is passed directly or through a derived signal. * * @param input - Originating request ID and response fields. */ export const respondToHttpRequest = Object.assign( function respondToHttpRequest(input: RespondToHttpRequestInput) { inferHttpResponseWaiting(isSignal(input) ? input : input.requestId) return respondToHttpRequestAction(input) }, { meta: respondToHttpRequestAction.meta }, ) satisfies typeof respondToHttpRequestAction /** * Enables waiting for HTTP subscriptions targeted by one response action. * * @param requestId - Literal or signal-backed response target. */ function inferHttpResponseWaiting(requestId: unknown) { const planningContext = getAutomationPlanningContext() if (!planningContext || !isSignal(requestId)) return const origins = new Set(requestId[signalContextOrigins]) planningContext.subscriptions = planningContext.subscriptions.map( (subscription) => { const config = subscription.config const origin = `event:${formatHookLocation({ scopePath: subscription.scopePath ?? [], slot: subscription.hookSlot, })}` if ( subscription.eventType !== httpRequestTriggerDefinition.type || !origins.has(origin) || typeof config !== "object" || config === null || Array.isArray(config) || Reflect.get(config, "waitForResponse") !== undefined ) { return subscription } return { ...subscription, config: { ...config, waitForResponse: true }, } }, ) }