From ce436b01a793912b5405588ebc1ee42e5cefa435 Mon Sep 17 00:00:00 2001 From: HugoCasa Date: Fri, 7 Mar 2025 19:49:43 +0100 Subject: [PATCH] capture nits (#5456) --- .../src/lib/components/ScriptBuilder.svelte | 12 +- .../components/triggers/CaptureTable.svelte | 12 +- .../http/RouteEditorConfigSection.svelte | 7 +- .../webhook/WebhooksConfigSection.svelte | 28 +- frontend/src/lib/script_helpers.ts | 367 ++++++++++-------- 5 files changed, 232 insertions(+), 194 deletions(-) diff --git a/frontend/src/lib/components/ScriptBuilder.svelte b/frontend/src/lib/components/ScriptBuilder.svelte index a31a8d992f..d112455e72 100644 --- a/frontend/src/lib/components/ScriptBuilder.svelte +++ b/frontend/src/lib/components/ScriptBuilder.svelte @@ -71,8 +71,10 @@ import TriggersEditor from './triggers/TriggersEditor.svelte' import type { ScheduleTrigger, TriggerContext, TriggerKind } from './triggers' import { - BUN_PREPROCESSOR_MODULE_CODE, - PYTHON_PREPROCESSOR_MODULE_CODE + TS_PREPROCESSOR_MODULE_CODE, + TS_PREPROCESSOR_SCRIPT_INTRO, + PYTHON_PREPROCESSOR_MODULE_CODE, + PYTHON_PREPROCESSOR_SCRIPT_INTRO } from '$lib/script_helpers' import CaptureTable from './triggers/CaptureTable.svelte' import type { SavedAndModifiedValue } from './common/confirmationModal/unsavedTypes' @@ -733,8 +735,8 @@ if (code) { const preprocessorCode = script.language === 'python3' - ? PYTHON_PREPROCESSOR_MODULE_CODE - : BUN_PREPROCESSOR_MODULE_CODE + ? PYTHON_PREPROCESSOR_SCRIPT_INTRO + PYTHON_PREPROCESSOR_MODULE_CODE + : TS_PREPROCESSOR_SCRIPT_INTRO + TS_PREPROCESSOR_MODULE_CODE const mainIndex = code.indexOf( script.language === 'python3' ? 'def main' : 'export async function main' ) @@ -1390,7 +1392,7 @@ on:exitTriggers={() => { captureTable?.loadCaptures(true) }} - {args} + args={hasPreprocessor && selectedInputTab !== 'preprocessor' ? {} : args} {initialPath} schema={script.schema} noEditor={true} diff --git a/frontend/src/lib/components/triggers/CaptureTable.svelte b/frontend/src/lib/components/triggers/CaptureTable.svelte index fa77f161ef..1865ec4fbd 100644 --- a/frontend/src/lib/components/triggers/CaptureTable.svelte +++ b/frontend/src/lib/components/triggers/CaptureTable.svelte @@ -60,7 +60,7 @@ }>() interface CaptureWithPayload extends Capture { - getFullCapture?: () => Promise + getFullCapture?: () => Promise payloadData?: any } @@ -68,7 +68,7 @@ await infiniteList?.loadData(refresh ? 'refresh' : 'loadMore') } - function initLoadCaptures(testKind: 'preprocessor' | 'main' = 'main') { + function initLoadCaptures(kind: 'preprocessor' | 'main' = testKind) { const loadInputsPageFn = async (page: number, perPage: number) => { const captures = await CaptureService.listCaptures({ workspace: $workspaceStore!, @@ -93,7 +93,7 @@ } const trigger_extra = isObject(capture.trigger_extra) ? capture.trigger_extra : {} newCapture.payloadData = - testKind === 'preprocessor' + kind === 'preprocessor' ? capture.payload === 'WINDMILL_TOO_BIG' ? { payload: capture.payload, @@ -122,7 +122,7 @@ infiniteList?.setDeleteItemFn(deleteInputFn) } - async function handleSelect(capture: any) { + async function handleSelect(capture: Capture) { if (selected === capture.id) { resetSelected() } else { @@ -132,8 +132,8 @@ } } - async function getPayload(capture: any) { - let payloadData = {} + async function getPayload(capture: CaptureWithPayload) { + let payloadData: any = {} if (capture.getFullCapture) { const fullCapture = await capture.getFullCapture() payloadData = diff --git a/frontend/src/lib/components/triggers/http/RouteEditorConfigSection.svelte b/frontend/src/lib/components/triggers/http/RouteEditorConfigSection.svelte index 30cba3abee..4837d1d561 100644 --- a/frontend/src/lib/components/triggers/http/RouteEditorConfigSection.svelte +++ b/frontend/src/lib/components/triggers/http/RouteEditorConfigSection.svelte @@ -16,6 +16,7 @@ import CaptureTable from '../CaptureTable.svelte' import ClipboardPanel from '../../details/ClipboardPanel.svelte' import { isCloudHosted } from '$lib/cloud' + import { isObject } from '$lib/utils' export let initialTriggerPath: string | undefined = undefined export let dirtyRoutePath: boolean = false @@ -84,6 +85,10 @@ {@const captureURL = `${location.origin}${base}/api/w/${$workspaceStore}/capture_u/http/${ captureInfo.isFlow ? 'flow' : 'script' }/${captureInfo.path.replaceAll('/', '.')}/${route_path}`} + {@const cleanedRunnableArgs = + isObject(runnableArgs) && 'wm_trigger' in runnableArgs + ? Object.fromEntries(Object.entries(runnableArgs).filter(([key]) => key !== 'wm_trigger')) + : runnableArgs} diff --git a/frontend/src/lib/components/triggers/webhook/WebhooksConfigSection.svelte b/frontend/src/lib/components/triggers/webhook/WebhooksConfigSection.svelte index d589c709fb..188750b0a9 100644 --- a/frontend/src/lib/components/triggers/webhook/WebhooksConfigSection.svelte +++ b/frontend/src/lib/components/triggers/webhook/WebhooksConfigSection.svelte @@ -14,7 +14,7 @@ import { Highlight } from 'svelte-highlight' import { typescript } from 'svelte-highlight/languages' import ClipboardPanel from '../../details/ClipboardPanel.svelte' - import { copyToClipboard } from '$lib/utils' + import { copyToClipboard, isObject } from '$lib/utils' // import { page } from '$app/stores' import { base } from '$lib/base' import TriggerTokens from '../TriggerTokens.svelte' @@ -104,6 +104,11 @@ return headers } + $: cleanedRunnableArgs = + isObject(runnableArgs) && 'wm_trigger' in runnableArgs + ? Object.fromEntries(Object.entries(runnableArgs).filter(([key]) => key !== 'wm_trigger')) + : runnableArgs + function fetchCode() { if (webhookType === 'sync') { return ` @@ -117,10 +122,11 @@ async function triggerJob() { ${ requestType === 'get_path' ? '// Payload is a base64 encoded string of the arguments' - : `const body = JSON.stringify(${JSON.stringify(runnableArgs ?? {}, null, 2).replaceAll( - '\n', - '\n\t' - )});` + : `const body = JSON.stringify(${JSON.stringify( + cleanedRunnableArgs ?? {}, + null, + 2 + ).replaceAll('\n', '\n\t')});` } const endpoint = \`${url}\`; @@ -145,7 +151,7 @@ export async function main() { // triggerJob function let triggerJobFunction = ` async function triggerJob() { - const body = JSON.stringify(${JSON.stringify(runnableArgs ?? {}, null, 2).replaceAll( + const body = JSON.stringify(${JSON.stringify(cleanedRunnableArgs ?? {}, null, 2).replaceAll( '\n', '\n\t' )}); @@ -200,12 +206,12 @@ function waitForJobCompletion(UUID) { return `curl \\ -X POST ${captureUrl} \\ -H 'Content-Type: application/json' \\ --d '{"foo": 42}'` +-d '${JSON.stringify(cleanedRunnableArgs ?? {}, null, 2)}'` } function curlCode() { return `TOKEN='${token}' -${requestType !== 'get_path' ? `BODY='${JSON.stringify(runnableArgs ?? {})}'` : ''} +${requestType !== 'get_path' ? `BODY='${JSON.stringify(cleanedRunnableArgs ?? {})}'` : ''} URL='${url}' ${webhookType === 'sync' ? 'RESULT' : 'UUID'}=$(curl -s ${ requestType != 'get_path' ? "-H 'Content-Type: application/json'" : '' @@ -236,12 +242,12 @@ done` (tokenType === 'query' ? `?token=${token}${ requestType === 'get_path' - ? `&payload=${encodeURIComponent(btoa(JSON.stringify(runnableArgs ?? {})))}` + ? `&payload=${encodeURIComponent(btoa(JSON.stringify(cleanedRunnableArgs ?? {})))}` : '' }` : `${ requestType === 'get_path' - ? `?payload=${encodeURIComponent(btoa(JSON.stringify(runnableArgs ?? {})))}` + ? `?payload=${encodeURIComponent(btoa(JSON.stringify(cleanedRunnableArgs ?? {})))}` : '' }`) @@ -384,7 +390,7 @@ done` {#if requestType !== 'get_path'} {/if} {#key requestType} diff --git a/frontend/src/lib/script_helpers.ts b/frontend/src/lib/script_helpers.ts index ce16d1e9ab..1f1755a5a7 100644 --- a/frontend/src/lib/script_helpers.ts +++ b/frontend/src/lib/script_helpers.ts @@ -577,127 +577,111 @@ export async function main(approver?: string) { // add a form in Advanced - Suspend // all on approval steps: https://www.windmill.dev/docs/flows/flow_approval` -export const BUN_PREPROCESSOR_MODULE_CODE = ` -export async function preprocessor( - wm_trigger: { - kind: 'http' | 'email' | 'webhook' | 'websocket' | 'kafka' | 'nats' | 'postgres' | 'sqs' | 'mqtt', - http?: { - route: string // The route path, e.g. "/users/:id" - path: string // The actual path called, e.g. "/users/123" - method: string - params: Record // path parameters - query: Record // query parameters - headers: Record - }, - websocket?: { - url: string // The websocket url - }, - kafka?: { - brokers: string[] - topic: string - group_id: string - }, - nats?: { - servers: string[] - subject: string - headers?: Record - status?: number - description?: string - length: number - }, - sqs?: { - queue_url: string, - message_id?: string, - receipt_handle?: string, - attributes: Record, - message_attributes?: Record - }, - mqtt?: { - topic: string, - retain: boolean, - pkid: number, - qos: number, - v5?: { - payload_format_indicator?: number, - topic_alias?: number, - response_topic?: string, - correlation_data?: Array, - user_properties?: Array<[string, string]>, - subscription_identifiers?: Array, - content_type?: string - } - } - }, - /* your other args */ -) { - return { - // return the args to be passed to the runnable - } -} -` +export const TS_PREPROCESSOR_SCRIPT_INTRO = `/** + * Trigger preprocessor + * + * ⚠️ This function runs BEFORE the main function. + * + * It processes raw trigger data from various sources (webhook, custom HTTP route, SQS, WebSocket, Kafka, NATS, MQTT, Postgres, or email) + * before passing it to \`main\`. This separates the trigger logic from the main logic and keeps the auto-generated runnable UI clean. + * + * The preprocessor receives the same data \`main\` would if no preprocessor was used, + * plus trigger metadata in the \`wm_trigger\` object: + * - Webhook/HTTP: \`{ wm_trigger, bodyKey1, bodyKey2, ... }\` + * - Postgres: \`{ transaction_type, schema_name, table_name, row, wm_trigger }\` + * - WebSocket/Kafka/NATS/SQS/MQTT: \`{ msg, wm_trigger }\` + * - Email: \`{ raw_email, parsed_email, wm_trigger }\` + * + * The returned object defines the parameter values passed to \`main()\`. + * e.g., { b: 1, a: 2 } → Calls \`main(2, 1)\`, assuming \`main\` is defined as \`main(a: number, b: number)\`. + * Ensure that the parameter names in \`main\` match the keys in the returned object. + * + * Learn more: https://www.windmill.dev/docs/core_concepts/preprocessors + */\n\n` -const DENO_PREPROCESSOR_MODULE_CODE = ` -export async function preprocessor( - wm_trigger: { - kind: 'http' | 'email' | 'webhook' | 'websocket' | 'kafka' | 'nats' | 'postgres' | 'sqs' | 'mqtt', - http?: { - route: string // The route path, e.g. "/users/:id" - path: string // The actual path called, e.g. "/users/123" - method: string - params: Record // path parameters - query: Record // query parameters - headers: Record - }, - websocket?: { - url: string // The websocket url - }, - kafka?: { - brokers: string[] - topic: string - group_id: string - }, - nats?: { - servers: string[] - subject: string - headers?: Record - status?: number - description?: string - length: number - }, - sqs?: { - queue_url: string, - message_id?: string, - receipt_handle?: string, - attributes: Record, - message_attributes?: Record - }, - mqtt?: { - topic: string, - retain: boolean, - pkid: number, - qos: number, - v5?: { - payload_format_indicator?: number, - topic_alias?: number, - response_topic?: string, - correlation_data?: Array, - user_properties?: Array<[string, string]>, - subscription_identifiers?: Array, - content_type?: string - } - } - }, - /* your other args */ +export const TS_PREPROCESSOR_FLOW_INTRO = `/** + * Trigger preprocessor + * + * It processes raw trigger data from various sources (webhook, custom HTTP route, SQS, WebSocket, Kafka, NATS, MQTT, Postgres, or email) + * before passing it to the flow. This separates the trigger logic from the flow logic and keeps the auto-generated UI clean. + * + * The preprocessor receives the same data the flow would if no preprocessor was used, + * plus trigger metadata in the \`wm_trigger\` object: + * - Webhook/HTTP: \`{ wm_trigger, bodyKey1, bodyKey2, ... }\` + * - Postgres: \`{ transaction_type, schema_name, table_name, row, wm_trigger }\` + * - WebSocket/Kafka/NATS/SQS/MQTT: \`{ msg, wm_trigger }\` + * - Email: \`{ raw_email, parsed_email, wm_trigger }\` + * + * The returned object determines the parameter values passed to the flow. + * e.g., \`{ b: 1, a: 2 }\` → Calls the flow with \`a = 2\` and \`b = 1\`, assuming the flow has two inputs called \`a\` and \`b\`. + * Ensure that the input names of the flow match the keys in the returned object. + * + * Learn more: https://www.windmill.dev/docs/core_concepts/preprocessors + */\n\n` + +export const TS_PREPROCESSOR_MODULE_CODE = `export async function preprocessor( + /* + * Replace this comment with the parameters received from the trigger. + * Examples: \`bodyKey1\`, \`bodyKey2\` for Webhook/HTTP, \`msg\` for WebSocket, etc. + */ + + // The trigger metadata + wm_trigger: { + kind: 'http' | 'email' | 'webhook' | 'websocket' | 'kafka' | 'nats' | 'postgres' | 'sqs' | 'mqtt', + http?: { + route: string // The route path, e.g. "/users/:id" + path: string // The actual path called, e.g. "/users/123" + method: string + params: Record // path parameters + query: Record // query parameters + headers: Record + }, + websocket?: { + url: string // The websocket url + }, + kafka?: { + brokers: string[] + topic: string + group_id: string + }, + nats?: { + servers: string[] + subject: string + headers?: Record + status?: number + description?: string + length: number + }, + sqs?: { + queue_url: string, + message_id?: string, + receipt_handle?: string, + attributes: Record, + message_attributes?: Record + }, + mqtt?: { + topic: string, + retain: boolean, + pkid: number, + qos: number, + v5?: { + payload_format_indicator?: number, + topic_alias?: number, + response_topic?: string, + correlation_data?: Array, + user_properties?: Array<[string, string]>, + subscription_identifiers?: Array, + content_type?: string + } + } + } ) { - return { - // return the args to be passed to the runnable - } + return { + // return the args to be passed to the runnable + } } ` @@ -728,74 +712,115 @@ def main(): # add a form in Advanced - Suspend # all on approval steps: https://www.windmill.dev/docs/flows/flow_approval` +export const PYTHON_PREPROCESSOR_SCRIPT_INTRO = `# Trigger preprocessor +# +# ⚠️ This function runs BEFORE the main function. +# +# It processes raw trigger data from various sources (webhook, custom HTTP route, SQS, WebSocket, Kafka, NATS, MQTT, Postgres, or email) +# before passing it to \`main\`. This separates the trigger logic from the main logic and keeps the auto-generated UI clean. +# +# The preprocessor receives the same data \`main\` would if no preprocessor was used, +# plus trigger metadata in the \`wm_trigger\` object: +# - Webhook/HTTP: \`{ wm_trigger, bodyKey1, bodyKey2, ... }\` +# - Postgres: \`{ transaction_type, schema_name, table_name, row, wm_trigger }\` +# - WebSocket/Kafka/NATS/SQS/MQTT: \`{ msg, wm_trigger }\` +# - Email: \`{ raw_email, parsed_email, wm_trigger }\` +# +# The returned object defines the parameter values passed to \`main()\`. +# e.g., { b: 1, a: 2 } → Calls \`main(2, 1)\`, assuming \`main\` is defined as \`main(a: int, b: int)\`. +# Ensure that the parameter names in \`main\` match the keys in the returned object. +# +# Learn more: https://www.windmill.dev/docs/core_concepts/preprocessors\n\n` + +export const PYTHON_PREPROCESSOR_FLOW_INTRO = `# Trigger preprocessor +# +# It processes raw trigger data from various sources (webhook, custom HTTP route, SQS, WebSocket, Kafka, NATS, MQTT, Postgres, or email) +# before passing it to the flow. This separates the trigger logic from the flow logic and keeps the auto-generated UI clean. +# +# The preprocessor receives the same data the flow would if no preprocessor was used, +# plus trigger metadata in the \`wm_trigger\` object: +# - Webhook/HTTP: \`{ wm_trigger, bodyKey1, bodyKey2, ... }\` +# - Postgres: \`{ transaction_type, schema_name, table_name, row, wm_trigger }\` +# - WebSocket/Kafka/NATS/SQS/MQTT: \`{ msg, wm_trigger }\` +# - Email: \`{ raw_email, parsed_email, wm_trigger }\` +# +# The returned object determines the parameter values passed to the flow. +# e.g., \`{ b: 1, a: 2 }\` → Calls the flow with \`a = 2\` and \`b = 1\`, assuming the flow has two inputs called \`a\` and \`b\`. +# Ensure that the input names of the flow match the keys in the returned object. +# +# Learn more: https://www.windmill.dev/docs/core_concepts/preprocessors\n\n` + export const PYTHON_PREPROCESSOR_MODULE_CODE = `from typing import TypedDict, Literal class Http(TypedDict): - route: str # The route path, e.g. "/users/:id" - path: str # The actual path called, e.g. "/users/123" - method: str - params: dict[str, str] - query: dict[str, str] - headers: dict[str, str] + route: str # The route path, e.g. "/users/:id" + path: str # The actual path called, e.g. "/users/123" + method: str + params: dict[str, str] + query: dict[str, str] + headers: dict[str, str] class Websocket(TypedDict): - url: str # The websocket url + url: str # The websocket url class Kafka(TypedDict): - topic: str - brokers: list[str] - group_id: str + topic: str + brokers: list[str] + group_id: str class Nats(TypedDict): - servers: list[str] - subject: str - headers: dict[str, list[str]] | None - status: int | None - description: str | None - length: int + servers: list[str] + subject: str + headers: dict[str, list[str]] | None + status: int | None + description: str | None + length: int class MessageAttribute(TypedDict): - string_value: str | None - data_type: str + string_value: str | None + data_type: str class Sqs(TypedDict): - queue_url: str - message_id: str | None - receipt_handle: str | None - attributes: dict[str, str] - message_attributes: dict[str, MessageAttribute] | None + queue_url: str + message_id: str | None + receipt_handle: str | None + attributes: dict[str, str] + message_attributes: dict[str, MessageAttribute] | None class MqttV5Properties: - payload_format_indicator: int | None - topic_alias: int | None - response_topic: str | None - correlation_data: list[int] | None - user_properties: list[tuple[str, str]] | None - subscription_identifiers: list[int] | None - content_type: str | None + payload_format_indicator: int | None + topic_alias: int | None + response_topic: str | None + correlation_data: list[int] | None + user_properties: list[tuple[str, str]] | None + subscription_identifiers: list[int] | None + content_type: str | None -class Mqtt(TypeDict): - topic: str - retain: bool - pkid: int - qos: int - v5: MqttV5Properties | None +class Mqtt(TypedDict): + topic: str + retain: bool + pkid: int + qos: int + v5: MqttV5Properties | None class WmTrigger(TypedDict): - kind: Literal["http", "email", "webhook", "websocket", "kafka", "nats", "postgres", "sqs", "mqtt"] - http: Http | None - websocket: Websocket | None - kafka: Kafka | None - nats: Nats | None - sqs: Sqs | None - mqtt: Mqtt | None + kind: Literal["http", "email", "webhook", "websocket", "kafka", "nats", "postgres", "sqs", "mqtt"] + http: Http | None + websocket: Websocket | None + kafka: Kafka | None + nats: Nats | None + sqs: Sqs | None + mqtt: Mqtt | None def preprocessor( - wm_trigger: WmTrigger, - # your other args + # Replace this comment with the parameters received from the trigger. + # Examples: \`bodyKey1\`, \`bodyKey2\` for Webhook/HTTP, \`msg\` for WebSocket, etc. + + # Trigger metadata + wm_trigger: WmTrigger, ): - return { - # return the args to be passed to the runnable - } + return { + # return the args to be passed to the runnable + } ` const DOCKER_INIT_CODE = `# shellcheck shell=bash @@ -881,7 +906,7 @@ export const INITIAL_CODE = { trigger: BUN_INIT_CODE_TRIGGER, approval: BUN_INIT_CODE_APPROVAL, failure: BUN_FAILURE_MODULE_CODE, - preprocessor: BUN_PREPROCESSOR_MODULE_CODE, + preprocessor: TS_PREPROCESSOR_FLOW_INTRO + TS_PREPROCESSOR_MODULE_CODE, clear: BUN_INIT_CODE_CLEAR }, python3: { @@ -889,7 +914,7 @@ export const INITIAL_CODE = { trigger: PYTHON_INIT_CODE_TRIGGER, approval: PYTHON_INIT_CODE_APPROVAL, failure: PYTHON_FAILURE_MODULE_CODE, - preprocessor: PYTHON_PREPROCESSOR_MODULE_CODE, + preprocessor: PYTHON_PREPROCESSOR_FLOW_INTRO + PYTHON_PREPROCESSOR_MODULE_CODE, clear: PYTHON_INIT_CODE_CLEAR }, deno: { @@ -898,7 +923,7 @@ export const INITIAL_CODE = { trigger: DENO_INIT_CODE_TRIGGER, approval: DENO_INIT_CODE_APPROVAL, failure: DENO_FAILURE_MODULE_CODE, - preprocessor: DENO_PREPROCESSOR_MODULE_CODE, + preprocessor: TS_PREPROCESSOR_FLOW_INTRO + TS_PREPROCESSOR_MODULE_CODE, fetch: FETCH_INIT_CODE, clear: DENO_INIT_CODE_CLEAR }, @@ -1112,4 +1137,4 @@ export function getResetCode( } else { return initialCode(language, kind, subkind) } -} \ No newline at end of file +}