mirror of
https://github.com/simstudioai/sim.git
synced 2026-09-24 15:45:35 +08:00
feat(webhooks): added optioanl input format to webhooks, added support for file uploads (#1654)
* feat(webhooks): added optioanl input format to webhooks, added support for file uploads * feat(webhooks): added input format component to generic webhook trigger, added file support * consolidated execution files utils, extended presigned URL duration for async tasks
This commit is contained in:
@@ -26,7 +26,208 @@ import { BlockInfoCard } from "@/components/ui/block-info-card"
|
||||
|
||||
|
||||
|
||||
## Overview
|
||||
|
||||
The Generic Webhook block allows you to receive webhooks from any external service. This is a flexible trigger that can handle any JSON payload, making it ideal for integrating with services that don't have a dedicated Sim block.
|
||||
|
||||
## Basic Usage
|
||||
|
||||
### Simple Passthrough Mode
|
||||
|
||||
Without defining an input format, the webhook passes through the entire request body as-is:
|
||||
|
||||
```bash
|
||||
curl -X POST https://sim.ai/api/webhooks/trigger/{webhook-path} \
|
||||
-H "Content-Type: application/json" \
|
||||
-H "X-Sim-Secret: your-secret" \
|
||||
-d '{
|
||||
"message": "Test webhook trigger",
|
||||
"data": {
|
||||
"key": "value"
|
||||
}
|
||||
}'
|
||||
```
|
||||
|
||||
Access the data in downstream blocks using:
|
||||
- `<webhook1.message>` → "Test webhook trigger"
|
||||
- `<webhook1.data.key>` → "value"
|
||||
|
||||
### Structured Input Format (Optional)
|
||||
|
||||
Define an input schema to get typed fields and enable advanced features like file uploads:
|
||||
|
||||
**Input Format Configuration:**
|
||||
```json
|
||||
[
|
||||
{ "name": "message", "type": "string" },
|
||||
{ "name": "priority", "type": "number" },
|
||||
{ "name": "documents", "type": "files" }
|
||||
]
|
||||
```
|
||||
|
||||
**Webhook Request:**
|
||||
```bash
|
||||
curl -X POST https://sim.ai/api/webhooks/trigger/{webhook-path} \
|
||||
-H "Content-Type: application/json" \
|
||||
-H "X-Sim-Secret: your-secret" \
|
||||
-d '{
|
||||
"message": "Invoice submission",
|
||||
"priority": 1,
|
||||
"documents": [
|
||||
{
|
||||
"type": "file",
|
||||
"data": "data:application/pdf;base64,JVBERi0xLjQK...",
|
||||
"name": "invoice.pdf",
|
||||
"mime": "application/pdf"
|
||||
}
|
||||
]
|
||||
}'
|
||||
```
|
||||
|
||||
## File Uploads
|
||||
|
||||
### Supported File Formats
|
||||
|
||||
The webhook supports two file input formats:
|
||||
|
||||
#### 1. Base64 Encoded Files
|
||||
For uploading file content directly:
|
||||
|
||||
```json
|
||||
{
|
||||
"documents": [
|
||||
{
|
||||
"type": "file",
|
||||
"data": "data:image/png;base64,iVBORw0KGgoAAAANSUhEUgA...",
|
||||
"name": "screenshot.png",
|
||||
"mime": "image/png"
|
||||
}
|
||||
]
|
||||
}
|
||||
```
|
||||
|
||||
- **Max size**: 20MB per file
|
||||
- **Format**: Standard data URL with base64 encoding
|
||||
- **Storage**: Files are uploaded to secure execution storage
|
||||
|
||||
#### 2. URL References
|
||||
For passing existing file URLs:
|
||||
|
||||
```json
|
||||
{
|
||||
"documents": [
|
||||
{
|
||||
"type": "url",
|
||||
"data": "https://example.com/files/document.pdf",
|
||||
"name": "document.pdf",
|
||||
"mime": "application/pdf"
|
||||
}
|
||||
]
|
||||
}
|
||||
```
|
||||
|
||||
### Accessing Files in Downstream Blocks
|
||||
|
||||
Files are processed into `UserFile` objects with the following properties:
|
||||
|
||||
```typescript
|
||||
{
|
||||
id: string, // Unique file identifier
|
||||
name: string, // Original filename
|
||||
url: string, // Presigned URL (valid for 5 minutes)
|
||||
size: number, // File size in bytes
|
||||
type: string, // MIME type
|
||||
key: string, // Storage key
|
||||
uploadedAt: string, // ISO timestamp
|
||||
expiresAt: string // ISO timestamp (5 minutes)
|
||||
}
|
||||
```
|
||||
|
||||
**Access in blocks:**
|
||||
- `<webhook1.documents[0].url>` → Download URL
|
||||
- `<webhook1.documents[0].name>` → "invoice.pdf"
|
||||
- `<webhook1.documents[0].size>` → 524288
|
||||
- `<webhook1.documents[0].type>` → "application/pdf"
|
||||
|
||||
### Complete File Upload Example
|
||||
|
||||
```bash
|
||||
# Create a base64-encoded file
|
||||
echo "Hello World" | base64
|
||||
# SGVsbG8gV29ybGQK
|
||||
|
||||
# Send webhook with file
|
||||
curl -X POST https://sim.ai/api/webhooks/trigger/{webhook-path} \
|
||||
-H "Content-Type: application/json" \
|
||||
-H "X-Sim-Secret: your-secret" \
|
||||
-d '{
|
||||
"subject": "Document for review",
|
||||
"attachments": [
|
||||
{
|
||||
"type": "file",
|
||||
"data": "data:text/plain;base64,SGVsbG8gV29ybGQK",
|
||||
"name": "sample.txt",
|
||||
"mime": "text/plain"
|
||||
}
|
||||
]
|
||||
}'
|
||||
```
|
||||
|
||||
## Authentication
|
||||
|
||||
### Configure Authentication (Optional)
|
||||
|
||||
In the webhook configuration:
|
||||
1. Enable "Require Authentication"
|
||||
2. Set a secret token
|
||||
3. Choose header type:
|
||||
- **Custom Header**: `X-Sim-Secret: your-token`
|
||||
- **Authorization Bearer**: `Authorization: Bearer your-token`
|
||||
|
||||
### Using Authentication
|
||||
|
||||
```bash
|
||||
# With custom header
|
||||
curl -X POST https://sim.ai/api/webhooks/trigger/{webhook-path} \
|
||||
-H "Content-Type: application/json" \
|
||||
-H "X-Sim-Secret: your-secret-token" \
|
||||
-d '{"message": "Authenticated request"}'
|
||||
|
||||
# With bearer token
|
||||
curl -X POST https://sim.ai/api/webhooks/trigger/{webhook-path} \
|
||||
-H "Content-Type: application/json" \
|
||||
-H "Authorization: Bearer your-secret-token" \
|
||||
-d '{"message": "Authenticated request"}'
|
||||
```
|
||||
|
||||
## Best Practices
|
||||
|
||||
1. **Use Input Format for Structure**: Define an input format when you know the expected schema. This provides:
|
||||
- Type validation
|
||||
- Better autocomplete in the editor
|
||||
- File upload capabilities
|
||||
|
||||
2. **Authentication**: Always enable authentication for production webhooks to prevent unauthorized access.
|
||||
|
||||
3. **File Size Limits**: Keep files under 20MB. For larger files, use URL references instead.
|
||||
|
||||
4. **File Expiration**: Downloaded files have 5-minute expiration URLs. Process them promptly or store them elsewhere if needed longer.
|
||||
|
||||
5. **Error Handling**: Webhook processing is asynchronous. Check execution logs for errors.
|
||||
|
||||
6. **Testing**: Use the "Test Webhook" button in the editor to validate your configuration before deployment.
|
||||
|
||||
## Use Cases
|
||||
|
||||
- **Form Submissions**: Receive data from custom forms with file uploads
|
||||
- **Third-Party Integrations**: Connect with services that send webhooks (Stripe, GitHub, etc.)
|
||||
- **Document Processing**: Accept documents from external systems for processing
|
||||
- **Event Notifications**: Receive event data from various sources
|
||||
- **Custom APIs**: Build custom API endpoints for your applications
|
||||
|
||||
## Notes
|
||||
|
||||
- Category: `triggers`
|
||||
- Type: `generic_webhook`
|
||||
- **File Support**: Available via input format configuration
|
||||
- **Max File Size**: 20MB per file
|
||||
|
||||
@@ -123,11 +123,13 @@ export async function POST(
|
||||
const { SSE_HEADERS } = await import('@/lib/utils')
|
||||
const { createFilteredResult } = await import('@/app/api/workflows/[id]/execute/route')
|
||||
|
||||
// Generate executionId early so it can be used for file uploads and workflow execution
|
||||
const executionId = crypto.randomUUID()
|
||||
|
||||
const workflowInput: any = { input, conversationId }
|
||||
if (files && Array.isArray(files) && files.length > 0) {
|
||||
logger.debug(`[${requestId}] Processing ${files.length} attached files`)
|
||||
|
||||
const executionId = crypto.randomUUID()
|
||||
const executionContext = {
|
||||
workspaceId: deployment.userId,
|
||||
workflowId: deployment.workflowId,
|
||||
@@ -153,6 +155,7 @@ export async function POST(
|
||||
workflowTriggerType: 'chat',
|
||||
},
|
||||
createFilteredResult,
|
||||
executionId,
|
||||
})
|
||||
|
||||
const streamResponse = new NextResponse(stream, {
|
||||
|
||||
@@ -3,10 +3,10 @@ import { chat, workflow } from '@sim/db/schema'
|
||||
import { eq } from 'drizzle-orm'
|
||||
import { type NextRequest, NextResponse } from 'next/server'
|
||||
import { isDev } from '@/lib/environment'
|
||||
import { processExecutionFiles } from '@/lib/execution/files'
|
||||
import { createLogger } from '@/lib/logs/console/logger'
|
||||
import { hasAdminPermission } from '@/lib/permissions/utils'
|
||||
import { decryptSecret } from '@/lib/utils'
|
||||
import { uploadExecutionFile } from '@/lib/workflows/execution-file-storage'
|
||||
import type { UserFile } from '@/executor/types'
|
||||
|
||||
const logger = createLogger('ChatAuthUtils')
|
||||
@@ -269,57 +269,20 @@ export async function validateChatAuth(
|
||||
/**
|
||||
* Process and upload chat files to execution storage
|
||||
* Handles both base64 dataUrl format and direct URL pass-through
|
||||
* Delegates to shared execution file processing logic
|
||||
*/
|
||||
export async function processChatFiles(
|
||||
files: Array<{ dataUrl?: string; url?: string; name: string; type: string }>,
|
||||
executionContext: { workspaceId: string; workflowId: string; executionId: string },
|
||||
requestId: string
|
||||
): Promise<UserFile[]> {
|
||||
const uploadedFiles: UserFile[] = []
|
||||
// Transform chat file format to shared execution file format
|
||||
const transformedFiles = files.map((file) => ({
|
||||
type: file.dataUrl ? 'file' : 'url',
|
||||
data: file.dataUrl || file.url || '',
|
||||
name: file.name,
|
||||
mime: file.type,
|
||||
}))
|
||||
|
||||
for (const file of files) {
|
||||
try {
|
||||
if (file.dataUrl) {
|
||||
const dataUrlPrefix = 'data:'
|
||||
const base64Prefix = ';base64,'
|
||||
|
||||
if (!file.dataUrl.startsWith(dataUrlPrefix)) {
|
||||
logger.warn(`[${requestId}] Invalid dataUrl format for file: ${file.name}`)
|
||||
continue
|
||||
}
|
||||
|
||||
const base64Index = file.dataUrl.indexOf(base64Prefix)
|
||||
if (base64Index === -1) {
|
||||
logger.warn(
|
||||
`[${requestId}] Invalid dataUrl format (no base64 marker) for file: ${file.name}`
|
||||
)
|
||||
continue
|
||||
}
|
||||
|
||||
const mimeType = file.dataUrl.substring(dataUrlPrefix.length, base64Index)
|
||||
const base64Data = file.dataUrl.substring(base64Index + base64Prefix.length)
|
||||
const buffer = Buffer.from(base64Data, 'base64')
|
||||
|
||||
logger.debug(`[${requestId}] Uploading file to S3: ${file.name} (${buffer.length} bytes)`)
|
||||
|
||||
const userFile = await uploadExecutionFile(
|
||||
executionContext,
|
||||
buffer,
|
||||
file.name,
|
||||
mimeType || file.type
|
||||
)
|
||||
|
||||
uploadedFiles.push(userFile)
|
||||
logger.debug(`[${requestId}] Successfully uploaded ${file.name} with URL: ${userFile.url}`)
|
||||
} else if (file.url) {
|
||||
uploadedFiles.push(file as UserFile)
|
||||
logger.debug(`[${requestId}] Using existing URL for file: ${file.name}`)
|
||||
}
|
||||
} catch (error) {
|
||||
logger.error(`[${requestId}] Failed to process file ${file.name}:`, error)
|
||||
throw new Error(`Failed to upload file: ${file.name}`)
|
||||
}
|
||||
}
|
||||
|
||||
return uploadedFiles
|
||||
return processExecutionFiles(transformedFiles, executionContext, requestId)
|
||||
}
|
||||
|
||||
@@ -11,6 +11,7 @@ import { checkServerSideUsageLimits } from '@/lib/billing'
|
||||
import { getHighestPrioritySubscription } from '@/lib/billing/core/subscription'
|
||||
import { env } from '@/lib/env'
|
||||
import { getPersonalAndWorkspaceEnv } from '@/lib/environment/utils'
|
||||
import { processExecutionFiles } from '@/lib/execution/files'
|
||||
import { createLogger } from '@/lib/logs/console/logger'
|
||||
import { LoggingSession } from '@/lib/logs/execution/logging-session'
|
||||
import { buildTraceSpans } from '@/lib/logs/execution/trace-spans/trace-spans'
|
||||
@@ -23,11 +24,7 @@ import {
|
||||
workflowHasResponseBlock,
|
||||
} from '@/lib/workflows/utils'
|
||||
import { validateWorkflowAccess } from '@/app/api/workflows/middleware'
|
||||
import {
|
||||
createErrorResponse,
|
||||
createSuccessResponse,
|
||||
processApiWorkflowField,
|
||||
} from '@/app/api/workflows/utils'
|
||||
import { createErrorResponse, createSuccessResponse } from '@/app/api/workflows/utils'
|
||||
import { Executor } from '@/executor'
|
||||
import type { ExecutionResult } from '@/executor/types'
|
||||
import { Serializer } from '@/serializer'
|
||||
@@ -124,10 +121,11 @@ export async function executeWorkflow(
|
||||
onStream?: (streamingExec: any) => Promise<void> // Callback for streaming agent responses
|
||||
onBlockComplete?: (blockId: string, output: any) => Promise<void> // Callback when any block completes
|
||||
skipLoggingComplete?: boolean // When true, skip calling loggingSession.safeComplete (for streaming)
|
||||
}
|
||||
},
|
||||
providedExecutionId?: string
|
||||
): Promise<ExecutionResult> {
|
||||
const workflowId = workflow.id
|
||||
const executionId = uuidv4()
|
||||
const executionId = providedExecutionId || uuidv4()
|
||||
|
||||
const executionKey = `${workflowId}:${requestId}`
|
||||
|
||||
@@ -577,6 +575,9 @@ export async function POST(
|
||||
input: rawInput,
|
||||
} = extractExecutionParams(request as NextRequest, parsedBody)
|
||||
|
||||
// Generate executionId early so it can be used for file uploads
|
||||
const executionId = uuidv4()
|
||||
|
||||
let processedInput = rawInput
|
||||
logger.info(`[${requestId}] Raw input received:`, JSON.stringify(rawInput, null, 2))
|
||||
|
||||
@@ -607,16 +608,18 @@ export async function POST(
|
||||
const executionContext = {
|
||||
workspaceId: validation.workflow.workspaceId,
|
||||
workflowId,
|
||||
executionId,
|
||||
}
|
||||
|
||||
for (const fileField of fileFields) {
|
||||
const fieldValue = rawInput[fileField.name]
|
||||
|
||||
if (fieldValue && typeof fieldValue === 'object') {
|
||||
const uploadedFiles = await processApiWorkflowField(
|
||||
const uploadedFiles = await processExecutionFiles(
|
||||
fieldValue,
|
||||
executionContext,
|
||||
requestId
|
||||
requestId,
|
||||
isAsync
|
||||
)
|
||||
|
||||
if (uploadedFiles.length > 0) {
|
||||
@@ -769,6 +772,7 @@ export async function POST(
|
||||
workflowTriggerType,
|
||||
},
|
||||
createFilteredResult,
|
||||
executionId,
|
||||
})
|
||||
|
||||
return new NextResponse(stream, {
|
||||
@@ -782,7 +786,8 @@ export async function POST(
|
||||
requestId,
|
||||
input,
|
||||
authenticatedUserId,
|
||||
undefined
|
||||
undefined,
|
||||
executionId
|
||||
)
|
||||
|
||||
const hasResponseBlock = workflowHasResponseBlock(result)
|
||||
|
||||
@@ -1,14 +1,9 @@
|
||||
import { NextResponse } from 'next/server'
|
||||
import { v4 as uuidv4 } from 'uuid'
|
||||
import { createLogger } from '@/lib/logs/console/logger'
|
||||
import { getUserEntityPermissions } from '@/lib/permissions/utils'
|
||||
import { uploadExecutionFile } from '@/lib/workflows/execution-file-storage'
|
||||
import type { UserFile } from '@/executor/types'
|
||||
|
||||
const logger = createLogger('WorkflowUtils')
|
||||
|
||||
const MAX_FILE_SIZE = 20 * 1024 * 1024 // 20MB
|
||||
|
||||
export function createErrorResponse(error: string, status: number, code?: string) {
|
||||
return NextResponse.json(
|
||||
{
|
||||
@@ -42,99 +37,3 @@ export async function verifyWorkspaceMembership(
|
||||
return null
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Process API workflow files - handles both base64 ('file' type) and URL pass-through ('url' type)
|
||||
*/
|
||||
export async function processApiWorkflowFiles(
|
||||
file: { type: string; data: string; name: string; mime?: string },
|
||||
executionContext: { workspaceId: string; workflowId: string; executionId: string },
|
||||
requestId: string
|
||||
): Promise<UserFile | null> {
|
||||
if (file.type === 'file' && file.data && file.name) {
|
||||
const dataUrlPrefix = 'data:'
|
||||
const base64Prefix = ';base64,'
|
||||
|
||||
if (!file.data.startsWith(dataUrlPrefix)) {
|
||||
logger.warn(`[${requestId}] Invalid data format for file: ${file.name}`)
|
||||
return null
|
||||
}
|
||||
|
||||
const base64Index = file.data.indexOf(base64Prefix)
|
||||
if (base64Index === -1) {
|
||||
logger.warn(`[${requestId}] Invalid data format (no base64 marker) for file: ${file.name}`)
|
||||
return null
|
||||
}
|
||||
|
||||
const mimeType = file.data.substring(dataUrlPrefix.length, base64Index)
|
||||
const base64Data = file.data.substring(base64Index + base64Prefix.length)
|
||||
const buffer = Buffer.from(base64Data, 'base64')
|
||||
|
||||
if (buffer.length > MAX_FILE_SIZE) {
|
||||
const fileSizeMB = (buffer.length / (1024 * 1024)).toFixed(2)
|
||||
throw new Error(
|
||||
`File "${file.name}" exceeds the maximum size limit of 20MB (actual size: ${fileSizeMB}MB)`
|
||||
)
|
||||
}
|
||||
|
||||
logger.debug(`[${requestId}] Uploading file: ${file.name} (${buffer.length} bytes)`)
|
||||
|
||||
const userFile = await uploadExecutionFile(
|
||||
executionContext,
|
||||
buffer,
|
||||
file.name,
|
||||
mimeType || file.mime || 'application/octet-stream'
|
||||
)
|
||||
|
||||
logger.debug(`[${requestId}] Successfully uploaded ${file.name}`)
|
||||
return userFile
|
||||
}
|
||||
|
||||
if (file.type === 'url' && file.data) {
|
||||
return {
|
||||
id: uuidv4(),
|
||||
url: file.data,
|
||||
name: file.name,
|
||||
size: 0,
|
||||
type: file.mime || 'application/octet-stream',
|
||||
key: `url/${file.name}`,
|
||||
uploadedAt: new Date().toISOString(),
|
||||
expiresAt: new Date(Date.now() + 7 * 24 * 60 * 60 * 1000).toISOString(),
|
||||
}
|
||||
}
|
||||
|
||||
return null
|
||||
}
|
||||
|
||||
/**
|
||||
* Process all files for a given field in the API workflow input
|
||||
*/
|
||||
export async function processApiWorkflowField(
|
||||
fieldValue: any,
|
||||
executionContext: { workspaceId: string; workflowId: string },
|
||||
requestId: string
|
||||
): Promise<UserFile[]> {
|
||||
if (!fieldValue || typeof fieldValue !== 'object') {
|
||||
return []
|
||||
}
|
||||
|
||||
const files = Array.isArray(fieldValue) ? fieldValue : [fieldValue]
|
||||
const uploadedFiles: UserFile[] = []
|
||||
const executionId = uuidv4()
|
||||
const fullContext = { ...executionContext, executionId }
|
||||
|
||||
for (const file of files) {
|
||||
try {
|
||||
const userFile = await processApiWorkflowFiles(file, fullContext, requestId)
|
||||
|
||||
if (userFile) {
|
||||
uploadedFiles.push(userFile)
|
||||
}
|
||||
} catch (error) {
|
||||
logger.error(`[${requestId}] Failed to process file ${file.name}:`, error)
|
||||
throw new Error(`Failed to upload file: ${file.name}`)
|
||||
}
|
||||
}
|
||||
|
||||
return uploadedFiles
|
||||
}
|
||||
|
||||
@@ -363,6 +363,12 @@ export default function Logs() {
|
||||
|
||||
const fetchLogs = useCallback(async (pageNum: number, append = false) => {
|
||||
try {
|
||||
// Don't fetch if workspaceId is not set
|
||||
const { workspaceId: storeWorkspaceId } = useFilterStore.getState()
|
||||
if (!storeWorkspaceId) {
|
||||
return
|
||||
}
|
||||
|
||||
if (pageNum === 1) {
|
||||
setLoading(true)
|
||||
} else {
|
||||
@@ -497,6 +503,11 @@ export default function Logs() {
|
||||
return
|
||||
}
|
||||
|
||||
// Don't fetch if workspaceId is not set yet
|
||||
if (!workspaceId) {
|
||||
return
|
||||
}
|
||||
|
||||
setPage(1)
|
||||
setHasMore(true)
|
||||
|
||||
|
||||
@@ -5,6 +5,7 @@ import { eq, sql } from 'drizzle-orm'
|
||||
import { v4 as uuidv4 } from 'uuid'
|
||||
import { checkServerSideUsageLimits } from '@/lib/billing'
|
||||
import { getPersonalAndWorkspaceEnv } from '@/lib/environment/utils'
|
||||
import { processExecutionFiles } from '@/lib/execution/files'
|
||||
import { IdempotencyService, webhookIdempotency } from '@/lib/idempotency'
|
||||
import { createLogger } from '@/lib/logs/console/logger'
|
||||
import { LoggingSession } from '@/lib/logs/execution/logging-session'
|
||||
@@ -410,6 +411,53 @@ async function executeWebhookJobInternal(
|
||||
}
|
||||
}
|
||||
|
||||
// Process generic webhook files based on inputFormat
|
||||
if (input && payload.provider === 'generic' && payload.blockId && blocks[payload.blockId]) {
|
||||
try {
|
||||
const triggerBlock = blocks[payload.blockId]
|
||||
|
||||
if (triggerBlock?.subBlocks?.inputFormat?.value) {
|
||||
const inputFormat = triggerBlock.subBlocks.inputFormat.value as unknown as Array<{
|
||||
name: string
|
||||
type: 'string' | 'number' | 'boolean' | 'object' | 'array' | 'files'
|
||||
}>
|
||||
logger.debug(`[${requestId}] Processing generic webhook files from inputFormat`)
|
||||
|
||||
const fileFields = inputFormat.filter((field) => field.type === 'files')
|
||||
|
||||
if (fileFields.length > 0 && typeof input === 'object' && input !== null) {
|
||||
const executionContext = {
|
||||
workspaceId: workspaceId || '',
|
||||
workflowId: payload.workflowId,
|
||||
executionId,
|
||||
}
|
||||
|
||||
for (const fileField of fileFields) {
|
||||
const fieldValue = input[fileField.name]
|
||||
|
||||
if (fieldValue && typeof fieldValue === 'object') {
|
||||
const uploadedFiles = await processExecutionFiles(
|
||||
fieldValue,
|
||||
executionContext,
|
||||
requestId
|
||||
)
|
||||
|
||||
if (uploadedFiles.length > 0) {
|
||||
input[fileField.name] = uploadedFiles
|
||||
logger.info(
|
||||
`[${requestId}] Successfully processed ${uploadedFiles.length} file(s) for field: ${fileField.name}`
|
||||
)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
} catch (error) {
|
||||
logger.error(`[${requestId}] Error processing generic webhook files:`, error)
|
||||
// Continue without processing files rather than failing execution
|
||||
}
|
||||
}
|
||||
|
||||
// Create executor and execute
|
||||
const executor = new Executor({
|
||||
workflow: serializedWorkflow,
|
||||
|
||||
@@ -28,6 +28,15 @@ export const GenericWebhookBlock: BlockConfig = {
|
||||
triggerProvider: 'generic',
|
||||
availableTriggers: ['generic_webhook'],
|
||||
},
|
||||
// Optional input format for structured data including files
|
||||
{
|
||||
id: 'inputFormat',
|
||||
title: 'Input Format',
|
||||
type: 'input-format',
|
||||
layout: 'full',
|
||||
description:
|
||||
'Define the expected JSON input schema for this webhook (optional). Use type "files" for file uploads.',
|
||||
},
|
||||
],
|
||||
|
||||
tools: {
|
||||
|
||||
@@ -0,0 +1,106 @@
|
||||
import { v4 as uuidv4 } from 'uuid'
|
||||
import { createLogger } from '@/lib/logs/console/logger'
|
||||
import { uploadExecutionFile } from '@/lib/workflows/execution-file-storage'
|
||||
import type { UserFile } from '@/executor/types'
|
||||
|
||||
const logger = createLogger('ExecutionFiles')
|
||||
|
||||
const MAX_FILE_SIZE = 20 * 1024 * 1024 // 20MB
|
||||
|
||||
/**
|
||||
* Process a single file for workflow execution - handles both base64 ('file' type) and URL pass-through ('url' type)
|
||||
*/
|
||||
export async function processExecutionFile(
|
||||
file: { type: string; data: string; name: string; mime?: string },
|
||||
executionContext: { workspaceId: string; workflowId: string; executionId: string },
|
||||
requestId: string,
|
||||
isAsync?: boolean
|
||||
): Promise<UserFile | null> {
|
||||
if (file.type === 'file' && file.data && file.name) {
|
||||
const dataUrlPrefix = 'data:'
|
||||
const base64Prefix = ';base64,'
|
||||
|
||||
if (!file.data.startsWith(dataUrlPrefix)) {
|
||||
logger.warn(`[${requestId}] Invalid data format for file: ${file.name}`)
|
||||
return null
|
||||
}
|
||||
|
||||
const base64Index = file.data.indexOf(base64Prefix)
|
||||
if (base64Index === -1) {
|
||||
logger.warn(`[${requestId}] Invalid data format (no base64 marker) for file: ${file.name}`)
|
||||
return null
|
||||
}
|
||||
|
||||
const mimeType = file.data.substring(dataUrlPrefix.length, base64Index)
|
||||
const base64Data = file.data.substring(base64Index + base64Prefix.length)
|
||||
const buffer = Buffer.from(base64Data, 'base64')
|
||||
|
||||
if (buffer.length > MAX_FILE_SIZE) {
|
||||
const fileSizeMB = (buffer.length / (1024 * 1024)).toFixed(2)
|
||||
throw new Error(
|
||||
`File "${file.name}" exceeds the maximum size limit of 20MB (actual size: ${fileSizeMB}MB)`
|
||||
)
|
||||
}
|
||||
|
||||
logger.debug(`[${requestId}] Uploading file: ${file.name} (${buffer.length} bytes)`)
|
||||
|
||||
const userFile = await uploadExecutionFile(
|
||||
executionContext,
|
||||
buffer,
|
||||
file.name,
|
||||
mimeType || file.mime || 'application/octet-stream',
|
||||
isAsync
|
||||
)
|
||||
|
||||
logger.debug(`[${requestId}] Successfully uploaded ${file.name}`)
|
||||
return userFile
|
||||
}
|
||||
|
||||
if (file.type === 'url' && file.data) {
|
||||
return {
|
||||
id: uuidv4(),
|
||||
url: file.data,
|
||||
name: file.name,
|
||||
size: 0,
|
||||
type: file.mime || 'application/octet-stream',
|
||||
key: `url/${file.name}`,
|
||||
uploadedAt: new Date().toISOString(),
|
||||
expiresAt: new Date(Date.now() + 7 * 24 * 60 * 60 * 1000).toISOString(),
|
||||
}
|
||||
}
|
||||
|
||||
return null
|
||||
}
|
||||
|
||||
/**
|
||||
* Process all files for a given field in workflow execution input
|
||||
*/
|
||||
export async function processExecutionFiles(
|
||||
fieldValue: any,
|
||||
executionContext: { workspaceId: string; workflowId: string; executionId: string },
|
||||
requestId: string,
|
||||
isAsync?: boolean
|
||||
): Promise<UserFile[]> {
|
||||
if (!fieldValue || typeof fieldValue !== 'object') {
|
||||
return []
|
||||
}
|
||||
|
||||
const files = Array.isArray(fieldValue) ? fieldValue : [fieldValue]
|
||||
const uploadedFiles: UserFile[] = []
|
||||
const fullContext = { ...executionContext }
|
||||
|
||||
for (const file of files) {
|
||||
try {
|
||||
const userFile = await processExecutionFile(file, fullContext, requestId, isAsync)
|
||||
|
||||
if (userFile) {
|
||||
uploadedFiles.push(userFile)
|
||||
}
|
||||
} catch (error) {
|
||||
logger.error(`[${requestId}] Failed to process file ${file.name}:`, error)
|
||||
throw new Error(`Failed to upload file: ${file.name}`)
|
||||
}
|
||||
}
|
||||
|
||||
return uploadedFiles
|
||||
}
|
||||
@@ -87,9 +87,17 @@ export function getBlockOutputs(
|
||||
}
|
||||
|
||||
if (Array.isArray(inputFormatValue)) {
|
||||
// For API and Input triggers, only use inputFormat fields
|
||||
if (blockType === 'api_trigger' || blockType === 'input_trigger') {
|
||||
outputs = {} // Clear all default outputs
|
||||
// For API, Input triggers, and Generic Webhook, use inputFormat fields
|
||||
if (
|
||||
blockType === 'api_trigger' ||
|
||||
blockType === 'input_trigger' ||
|
||||
blockType === 'generic_webhook'
|
||||
) {
|
||||
// For generic_webhook, only clear outputs if inputFormat has fields
|
||||
// Otherwise keep the default outputs (pass-through body)
|
||||
if (inputFormatValue.length > 0 || blockType !== 'generic_webhook') {
|
||||
outputs = {} // Clear all default outputs
|
||||
}
|
||||
|
||||
// Add each field from inputFormat as an output at root level
|
||||
inputFormatValue.forEach((field: { name?: string; type?: string }) => {
|
||||
|
||||
@@ -36,7 +36,8 @@ export async function uploadExecutionFile(
|
||||
context: ExecutionContext,
|
||||
fileBuffer: Buffer,
|
||||
fileName: string,
|
||||
contentType: string
|
||||
contentType: string,
|
||||
isAsync?: boolean
|
||||
): Promise<UserFile> {
|
||||
logger.info(`Uploading execution file: ${fileName} for execution ${context.executionId}`)
|
||||
logger.debug(`File upload context:`, {
|
||||
@@ -53,6 +54,9 @@ export async function uploadExecutionFile(
|
||||
|
||||
logger.info(`Generated storage key: "${storageKey}" for file: ${fileName}`)
|
||||
|
||||
// Use 10-minute expiration for async executions, 5 minutes for sync
|
||||
const urlExpirationSeconds = isAsync ? 10 * 60 : 5 * 60
|
||||
|
||||
try {
|
||||
let fileInfo: any
|
||||
let directUrl: string | undefined
|
||||
@@ -78,16 +82,18 @@ export async function uploadExecutionFile(
|
||||
logger.info(`Original storage key was: "${storageKey}"`)
|
||||
logger.info(`Keys match: ${fileInfo.key === storageKey}`)
|
||||
|
||||
// Generate presigned URL for execution (5 minutes)
|
||||
// Generate presigned URL for execution (5 or 10 minutes)
|
||||
try {
|
||||
logger.info(`Generating presigned URL with key: "${fileInfo.key}"`)
|
||||
logger.info(
|
||||
`Generating presigned URL with key: "${fileInfo.key}" (expiration: ${urlExpirationSeconds / 60} minutes)`
|
||||
)
|
||||
directUrl = await getPresignedUrlWithConfig(
|
||||
fileInfo.key, // Use the actual uploaded key
|
||||
{
|
||||
bucket: S3_EXECUTION_FILES_CONFIG.bucket,
|
||||
region: S3_EXECUTION_FILES_CONFIG.region,
|
||||
},
|
||||
5 * 60 // 5 minutes
|
||||
urlExpirationSeconds
|
||||
)
|
||||
logger.info(`Generated presigned URL: ${directUrl}`)
|
||||
} catch (error) {
|
||||
@@ -102,7 +108,7 @@ export async function uploadExecutionFile(
|
||||
containerName: BLOB_EXECUTION_FILES_CONFIG.containerName,
|
||||
})
|
||||
|
||||
// Generate presigned URL for execution (5 minutes)
|
||||
// Generate presigned URL for execution (5 or 10 minutes)
|
||||
try {
|
||||
directUrl = await getBlobPresignedUrlWithConfig(
|
||||
fileInfo.key, // Use the actual uploaded key
|
||||
@@ -112,7 +118,7 @@ export async function uploadExecutionFile(
|
||||
connectionString: BLOB_EXECUTION_FILES_CONFIG.connectionString,
|
||||
containerName: BLOB_EXECUTION_FILES_CONFIG.containerName,
|
||||
},
|
||||
5 * 60 // 5 minutes
|
||||
urlExpirationSeconds
|
||||
)
|
||||
} catch (error) {
|
||||
logger.warn(`Failed to generate Blob presigned URL for ${fileName}:`, error)
|
||||
@@ -126,7 +132,7 @@ export async function uploadExecutionFile(
|
||||
name: fileName,
|
||||
size: fileBuffer.length,
|
||||
type: contentType,
|
||||
url: directUrl || `/api/files/serve/${fileInfo.key}`, // Use 5-minute presigned URL, fallback to serve path
|
||||
url: directUrl || `/api/files/serve/${fileInfo.key}`, // Use presigned URL (5 or 10 min), fallback to serve path
|
||||
key: fileInfo.key, // Use the actual uploaded key from S3/Blob
|
||||
uploadedAt: new Date().toISOString(),
|
||||
expiresAt: getFileExpirationDate(),
|
||||
|
||||
@@ -21,13 +21,21 @@ export interface StreamingResponseOptions {
|
||||
executingUserId: string
|
||||
streamConfig: StreamingConfig
|
||||
createFilteredResult: (result: ExecutionResult) => any
|
||||
executionId?: string
|
||||
}
|
||||
|
||||
export async function createStreamingResponse(
|
||||
options: StreamingResponseOptions
|
||||
): Promise<ReadableStream> {
|
||||
const { requestId, workflow, input, executingUserId, streamConfig, createFilteredResult } =
|
||||
options
|
||||
const {
|
||||
requestId,
|
||||
workflow,
|
||||
input,
|
||||
executingUserId,
|
||||
streamConfig,
|
||||
createFilteredResult,
|
||||
executionId,
|
||||
} = options
|
||||
|
||||
const { executeWorkflow, createFilteredResult: defaultFilteredResult } = await import(
|
||||
'@/app/api/workflows/[id]/execute/route'
|
||||
@@ -115,15 +123,22 @@ export async function createStreamingResponse(
|
||||
}
|
||||
}
|
||||
|
||||
const result = await executeWorkflow(workflow, requestId, input, executingUserId, {
|
||||
enabled: true,
|
||||
selectedOutputs: streamConfig.selectedOutputs,
|
||||
isSecureMode: streamConfig.isSecureMode,
|
||||
workflowTriggerType: streamConfig.workflowTriggerType,
|
||||
onStream: onStreamCallback,
|
||||
onBlockComplete: onBlockCompleteCallback,
|
||||
skipLoggingComplete: true, // We'll complete logging after tokenization
|
||||
})
|
||||
const result = await executeWorkflow(
|
||||
workflow,
|
||||
requestId,
|
||||
input,
|
||||
executingUserId,
|
||||
{
|
||||
enabled: true,
|
||||
selectedOutputs: streamConfig.selectedOutputs,
|
||||
isSecureMode: streamConfig.isSecureMode,
|
||||
workflowTriggerType: streamConfig.workflowTriggerType,
|
||||
onStream: onStreamCallback,
|
||||
onBlockComplete: onBlockCompleteCallback,
|
||||
skipLoggingComplete: true, // We'll complete logging after tokenization
|
||||
},
|
||||
executionId
|
||||
)
|
||||
|
||||
if (result.logs && streamedContent.size > 0) {
|
||||
result.logs = result.logs.map((log: any) => {
|
||||
|
||||
@@ -380,7 +380,7 @@ export function hasWorkflowChanged(
|
||||
deployedValue = sanitizeToolsForComparison(deployedValue)
|
||||
}
|
||||
|
||||
// Special handling for 'inputFormat' subBlock - sanitize test-only value fields
|
||||
// Special handling for 'inputFormat' subBlock - sanitize UI-only fields (collapsed state)
|
||||
if (
|
||||
subBlockId === 'inputFormat' &&
|
||||
Array.isArray(currentValue) &&
|
||||
|
||||
@@ -15,10 +15,12 @@ const updateURL = (params: URLSearchParams) => {
|
||||
window.history.replaceState({}, '', url)
|
||||
}
|
||||
|
||||
const DEFAULT_TIME_RANGE: TimeRange = 'Past 12 hours'
|
||||
const DEFAULT_TIME_RANGE: TimeRange = 'All time'
|
||||
|
||||
const parseTimeRangeFromURL = (value: string | null): TimeRange => {
|
||||
switch (value) {
|
||||
case 'all-time':
|
||||
return 'All time'
|
||||
case 'past-30-minutes':
|
||||
return 'Past 30 minutes'
|
||||
case 'past-hour':
|
||||
|
||||
Reference in New Issue
Block a user