diff --git a/.github/workflows/docs-embeddings.yml b/.github/workflows/docs-embeddings.yml index c42b3d96eb..5f51d4b0cc 100644 --- a/.github/workflows/docs-embeddings.yml +++ b/.github/workflows/docs-embeddings.yml @@ -17,7 +17,7 @@ jobs: - name: Setup Bun uses: oven-sh/setup-bun@v2 with: - bun-version: latest + bun-version: 1.2.22 - name: Setup Node uses: actions/setup-node@v4 diff --git a/.github/workflows/i18n.yml b/.github/workflows/i18n.yml index 5275b9444b..cad647e3df 100644 --- a/.github/workflows/i18n.yml +++ b/.github/workflows/i18n.yml @@ -27,7 +27,7 @@ jobs: - name: Setup Bun uses: oven-sh/setup-bun@v2 with: - bun-version: latest + bun-version: 1.2.22 - name: Run Lingo.dev translations env: @@ -116,7 +116,7 @@ jobs: - name: Setup Bun uses: oven-sh/setup-bun@v2 with: - bun-version: latest + bun-version: 1.2.22 - name: Install dependencies run: | diff --git a/.github/workflows/migrations.yml b/.github/workflows/migrations.yml index 191ee0fea2..2bfb6ca1c6 100644 --- a/.github/workflows/migrations.yml +++ b/.github/workflows/migrations.yml @@ -16,7 +16,7 @@ jobs: - name: Setup Bun uses: oven-sh/setup-bun@v2 with: - bun-version: latest + bun-version: 1.2.22 - name: Install dependencies run: bun install diff --git a/.github/workflows/publish-cli.yml b/.github/workflows/publish-cli.yml index 88209378c9..3e48ac6dc2 100644 --- a/.github/workflows/publish-cli.yml +++ b/.github/workflows/publish-cli.yml @@ -16,7 +16,7 @@ jobs: - name: Setup Bun uses: oven-sh/setup-bun@v2 with: - bun-version: latest + bun-version: 1.2.22 - name: Setup Node.js for npm publishing uses: actions/setup-node@v4 diff --git a/.github/workflows/publish-ts-sdk.yml b/.github/workflows/publish-ts-sdk.yml index cf6c53972f..187393ea8f 100644 --- a/.github/workflows/publish-ts-sdk.yml +++ b/.github/workflows/publish-ts-sdk.yml @@ -16,7 +16,7 @@ jobs: - name: Setup Bun uses: oven-sh/setup-bun@v2 with: - bun-version: latest + bun-version: 1.2.22 - name: Setup Node.js for npm publishing uses: actions/setup-node@v4 diff --git a/.github/workflows/test-build.yml b/.github/workflows/test-build.yml index 2022aed361..443c4d6a30 100644 --- a/.github/workflows/test-build.yml +++ b/.github/workflows/test-build.yml @@ -16,7 +16,7 @@ jobs: - name: Setup Bun uses: oven-sh/setup-bun@v2 with: - bun-version: latest + bun-version: 1.2.22 - name: Setup Node uses: actions/setup-node@v4 diff --git a/.github/workflows/trigger-deploy.yml b/.github/workflows/trigger-deploy.yml index ad8b0675af..88a059b282 100644 --- a/.github/workflows/trigger-deploy.yml +++ b/.github/workflows/trigger-deploy.yml @@ -27,7 +27,7 @@ jobs: - name: Setup Bun uses: oven-sh/setup-bun@v2 with: - bun-version: latest + bun-version: 1.2.22 - name: Install dependencies run: bun install diff --git a/apps/docs/content/docs/en/execution/costs.mdx b/apps/docs/content/docs/en/execution/costs.mdx index 890af4b1c5..d35423c409 100644 --- a/apps/docs/content/docs/en/execution/costs.mdx +++ b/apps/docs/content/docs/en/execution/costs.mdx @@ -166,6 +166,38 @@ Different subscription plans have different usage limits: | **Team** | $500 (pooled) | 50 sync, 100 async | | **Enterprise** | Custom | Custom | +## Billing Model + +Sim uses a **base subscription + overage** billing model: + +### How It Works + +**Pro Plan ($20/month):** +- Monthly subscription includes $20 of usage +- Usage under $20 → No additional charges +- Usage over $20 → Pay the overage at month end +- Example: $35 usage = $20 (subscription) + $15 (overage) + +**Team Plan ($40/seat/month):** +- Pooled usage across all team members +- Overage calculated from total team usage +- Organization owner receives one bill + +**Enterprise Plans:** +- Fixed monthly price, no overages +- Custom usage limits per agreement + +### Threshold Billing + +When unbilled overage reaches $50, Sim automatically bills the full unbilled amount. + +**Example:** +- Day 10: $70 overage → Bill $70 immediately +- Day 15: Additional $35 usage ($105 total) → Already billed, no action +- Day 20: Another $50 usage ($155 total, $85 unbilled) → Bill $85 immediately + +This spreads large overage charges throughout the month instead of one large bill at period end. + ## Cost Management Best Practices 1. **Monitor Regularly**: Check your usage dashboard frequently to avoid surprises diff --git a/apps/docs/content/docs/en/sdks/typescript.mdx b/apps/docs/content/docs/en/sdks/typescript.mdx index 1e8ea8782a..974e15442d 100644 --- a/apps/docs/content/docs/en/sdks/typescript.mdx +++ b/apps/docs/content/docs/en/sdks/typescript.mdx @@ -593,14 +593,91 @@ async function executeClientSideWorkflow() { }); console.log('Workflow result:', result); - + // Update UI with result - document.getElementById('result')!.textContent = + document.getElementById('result')!.textContent = JSON.stringify(result.output, null, 2); } catch (error) { console.error('Error:', error); } } +``` + +### File Upload + +File objects are automatically detected and converted to base64 format. Include them in your input under the field name matching your workflow's API trigger input format. + +The SDK converts File objects to this format: +```typescript +{ + type: 'file', + data: 'data:mime/type;base64,base64data', + name: 'filename', + mime: 'mime/type' +} +``` + +Alternatively, you can manually provide files using the URL format: +```typescript +{ + type: 'url', + data: 'https://example.com/file.pdf', + name: 'file.pdf', + mime: 'application/pdf' +} +``` + + + + ```typescript + import { SimStudioClient } from 'simstudio-ts-sdk'; + + const client = new SimStudioClient({ + apiKey: process.env.NEXT_PUBLIC_SIM_API_KEY! + }); + + // From file input + async function handleFileUpload(event: Event) { + const input = event.target as HTMLInputElement; + const files = Array.from(input.files || []); + + // Include files under the field name from your API trigger's input format + const result = await client.executeWorkflow('workflow-id', { + input: { + documents: files, // Must match your workflow's "files" field name + instructions: 'Analyze these documents' + } + }); + + console.log('Result:', result); + } + ``` + + + ```typescript + import { SimStudioClient } from 'simstudio-ts-sdk'; + import fs from 'fs'; + + const client = new SimStudioClient({ + apiKey: process.env.SIM_API_KEY! + }); + + // Read file and create File object + const fileBuffer = fs.readFileSync('./document.pdf'); + const file = new File([fileBuffer], 'document.pdf', { + type: 'application/pdf' + }); + + // Include files under the field name from your API trigger's input format + const result = await client.executeWorkflow('workflow-id', { + input: { + documents: [file], // Must match your workflow's "files" field name + query: 'Summarize this document' + } + }); + ``` + + // Attach to button click document.getElementById('executeBtn')?.addEventListener('click', executeClientSideWorkflow); diff --git a/apps/docs/content/docs/en/triggers/api.mdx b/apps/docs/content/docs/en/triggers/api.mdx index bc16047016..bb9c14f45a 100644 --- a/apps/docs/content/docs/en/triggers/api.mdx +++ b/apps/docs/content/docs/en/triggers/api.mdx @@ -22,9 +22,17 @@ The API trigger exposes your workflow as a secure HTTP endpoint. Send JSON data /> -Add an **Input Format** field for each parameter. Runtime output keys mirror the schema and are also available under ``. +Add an **Input Format** field for each parameter. Supported types: -Manual runs in the editor use the `value` column so you can test without sending a request. During execution the resolver populates both `` and ``. +- **string** - Text values +- **number** - Numeric values +- **boolean** - True/false values +- **json** - JSON objects +- **files** - File uploads (access via ``, ``, etc.) + +Runtime output keys mirror the schema and are available under ``. + +Manual runs in the editor use the `value` column so you can test without sending a request. During execution the resolver populates both `` and ``. ## Request Example @@ -123,6 +131,53 @@ data: {"blockId":"agent1-uuid","chunk":" complete"} | `` | Field defined in the Input Format | | `` | Entire structured request body | +### File Upload Format + +The API accepts files in two formats: + +**1. Base64-encoded files** (recommended for SDKs): +```json +{ + "documents": [{ + "type": "file", + "data": "data:application/pdf;base64,JVBERi0xLjQK...", + "name": "document.pdf", + "mime": "application/pdf" + }] +} +``` +- Maximum file size: 20MB per file +- Files are uploaded to cloud storage and converted to UserFile objects with all properties + +**2. Direct URL references**: +```json +{ + "documents": [{ + "type": "url", + "data": "https://example.com/document.pdf", + "name": "document.pdf", + "mime": "application/pdf" + }] +} +``` +- File is not uploaded, URL is passed through directly +- Useful for referencing existing files + +### File Properties + +For files, access all properties: + +| Property | Description | Type | +|----------|-------------|------| +| `` | Signed download URL | string | +| `` | Original filename | string | +| `` | File size in bytes | number | +| `` | MIME type | string | +| `` | Upload timestamp (ISO 8601) | string | +| `` | URL expiry timestamp (ISO 8601) | string | + +For URL-referenced files, the same properties are available except `uploadedAt` and `expiresAt` since the file is not uploaded to our storage. + If no Input Format is defined, the executor exposes the raw JSON at `` only. diff --git a/apps/docs/content/docs/en/triggers/chat.mdx b/apps/docs/content/docs/en/triggers/chat.mdx index 4d6d7b7161..12bd55c2ef 100644 --- a/apps/docs/content/docs/en/triggers/chat.mdx +++ b/apps/docs/content/docs/en/triggers/chat.mdx @@ -24,13 +24,24 @@ The Chat trigger creates a conversational interface for your workflow. Deploy yo The trigger writes three fields that downstream blocks can reference: -| Reference | Description | -|-----------|-------------| -| `` | Latest user message | -| `` | Conversation thread ID | -| `` | Optional uploaded files | +| Reference | Description | Type | +|-----------|-------------|------| +| `` | Latest user message | string | +| `` | Conversation thread ID | string | +| `` | Optional uploaded files | files array | -Files include `name`, `mimeType`, and a signed download `url`. +### File Properties + +Access individual file properties using array indexing: + +| Property | Description | Type | +|----------|-------------|------| +| `` | Signed download URL | string | +| `` | Original filename | string | +| `` | File size in bytes | number | +| `` | MIME type | string | +| `` | Upload timestamp (ISO 8601) | string | +| `` | URL expiry timestamp (ISO 8601) | string | ## Usage Notes diff --git a/apps/sim/app/api/__test-utils__/utils.ts b/apps/sim/app/api/__test-utils__/utils.ts index ce02592a9b..da96feedd8 100644 --- a/apps/sim/app/api/__test-utils__/utils.ts +++ b/apps/sim/app/api/__test-utils__/utils.ts @@ -1116,12 +1116,20 @@ export function createMockDatabase(options: MockDatabaseOptions = {}) { const createUpdateChain = () => ({ set: vi.fn().mockImplementation(() => ({ - where: vi.fn().mockImplementation(() => { - if (updateOptions.throwError) { - return Promise.reject(createDbError('update', updateOptions.errorMessage)) - } - return Promise.resolve(updateOptions.results) - }), + where: vi.fn().mockImplementation(() => ({ + returning: vi.fn().mockImplementation(() => { + if (updateOptions.throwError) { + return Promise.reject(createDbError('update', updateOptions.errorMessage)) + } + return Promise.resolve(updateOptions.results) + }), + then: vi.fn().mockImplementation((resolve) => { + if (updateOptions.throwError) { + return Promise.reject(createDbError('update', updateOptions.errorMessage)) + } + return Promise.resolve(updateOptions.results).then(resolve) + }), + })), })), }) diff --git a/apps/sim/app/api/billing/update-cost/route.ts b/apps/sim/app/api/billing/update-cost/route.ts index fd897c59bd..418691a97b 100644 --- a/apps/sim/app/api/billing/update-cost/route.ts +++ b/apps/sim/app/api/billing/update-cost/route.ts @@ -3,6 +3,7 @@ import { userStats } from '@sim/db/schema' import { eq, sql } from 'drizzle-orm' import { type NextRequest, NextResponse } from 'next/server' import { z } from 'zod' +import { checkAndBillOverageThreshold } from '@/lib/billing/threshold-billing' import { checkInternalApiKey } from '@/lib/copilot/utils' import { isBillingEnabled } from '@/lib/environment' import { createLogger } from '@/lib/logs/console/logger' @@ -148,6 +149,9 @@ export async function POST(req: NextRequest) { addedTokens: totalTokens, }) + // Check if user has hit overage threshold and bill incrementally + await checkAndBillOverageThreshold(userId) + const duration = Date.now() - startTime logger.info(`[${requestId}] Cost update completed successfully`, { diff --git a/apps/sim/app/api/chat/[identifier]/otp/route.ts b/apps/sim/app/api/chat/[identifier]/otp/route.ts index 8ba425972d..3cb30844d3 100644 --- a/apps/sim/app/api/chat/[identifier]/otp/route.ts +++ b/apps/sim/app/api/chat/[identifier]/otp/route.ts @@ -6,7 +6,7 @@ import { z } from 'zod' import { renderOTPEmail } from '@/components/emails/render-email' import { sendEmail } from '@/lib/email/mailer' import { createLogger } from '@/lib/logs/console/logger' -import { getRedisClient, markMessageAsProcessed, releaseLock } from '@/lib/redis' +import { getRedisClient } from '@/lib/redis' import { generateRequestId } from '@/lib/utils' import { addCorsHeaders, setChatAuthCookie } from '@/app/api/chat/utils' import { createErrorResponse, createSuccessResponse } from '@/app/api/workflows/utils' @@ -21,83 +21,52 @@ function generateOTP() { // We use 15 minutes (900 seconds) expiry for OTPs const OTP_EXPIRY = 15 * 60 -// Store OTP in Redis -async function storeOTP(email: string, chatId: string, otp: string): Promise { +async function storeOTP(email: string, chatId: string, otp: string): Promise { const key = `otp:${email}:${chatId}` const redis = getRedisClient() - if (redis) { - // Use Redis if available - await redis.set(key, otp, 'EX', OTP_EXPIRY) - } else { - // Use the existing function as fallback to mark that an OTP exists - await markMessageAsProcessed(key, OTP_EXPIRY) + if (!redis) { + logger.warn('Redis not available, OTP functionality requires Redis') + return false + } - // For the fallback case, we need to handle storing the OTP value separately - // since markMessageAsProcessed only stores "1" - const valueKey = `${key}:value` - try { - // Access the in-memory cache directly - hacky but works for fallback - const inMemoryCache = (global as any).inMemoryCache - if (inMemoryCache) { - const fullKey = `processed:${valueKey}` - const expiry = OTP_EXPIRY ? Date.now() + OTP_EXPIRY * 1000 : null - inMemoryCache.set(fullKey, { value: otp, expiry }) - } - } catch (error) { - logger.error('Error storing OTP in fallback cache:', error) - } + try { + await redis.set(key, otp, 'EX', OTP_EXPIRY) + return true + } catch (error) { + logger.error('Error storing OTP in Redis:', error) + return false } } -// Get OTP from Redis async function getOTP(email: string, chatId: string): Promise { const key = `otp:${email}:${chatId}` const redis = getRedisClient() - if (redis) { - // Use Redis if available - return await redis.get(key) + if (!redis) { + return null } - // Use the existing function as fallback - check if it exists - const exists = await new Promise((resolve) => { - try { - // Check the in-memory cache directly - hacky but works for fallback - const inMemoryCache = (global as any).inMemoryCache - const fullKey = `processed:${key}` - const cacheEntry = inMemoryCache?.get(fullKey) - resolve(!!cacheEntry) - } catch { - resolve(false) - } - }) - if (!exists) return null - - // Try to get the value key - const valueKey = `${key}:value` try { - const inMemoryCache = (global as any).inMemoryCache - const fullKey = `processed:${valueKey}` - const cacheEntry = inMemoryCache?.get(fullKey) - return cacheEntry?.value || null - } catch { + return await redis.get(key) + } catch (error) { + logger.error('Error getting OTP from Redis:', error) return null } } -// Delete OTP from Redis async function deleteOTP(email: string, chatId: string): Promise { const key = `otp:${email}:${chatId}` const redis = getRedisClient() - if (redis) { - // Use Redis if available + if (!redis) { + return + } + + try { await redis.del(key) - } else { - // Use the existing function as fallback - await releaseLock(`processed:${key}`) - await releaseLock(`processed:${key}:value`) + } catch (error) { + logger.error('Error deleting OTP from Redis:', error) } } @@ -177,7 +146,17 @@ export async function POST( const otp = generateOTP() - await storeOTP(email, deployment.id, otp) + const stored = await storeOTP(email, deployment.id, otp) + if (!stored) { + logger.error(`[${requestId}] Failed to store OTP - Redis unavailable`) + return addCorsHeaders( + createErrorResponse( + 'Email verification temporarily unavailable, please try again later', + 503 + ), + request + ) + } const emailHtml = await renderOTPEmail( otp, diff --git a/apps/sim/app/api/chat/[identifier]/route.ts b/apps/sim/app/api/chat/[identifier]/route.ts index e349dfe74c..3ddb926d3e 100644 --- a/apps/sim/app/api/chat/[identifier]/route.ts +++ b/apps/sim/app/api/chat/[identifier]/route.ts @@ -6,6 +6,7 @@ import { createLogger } from '@/lib/logs/console/logger' import { generateRequestId } from '@/lib/utils' import { addCorsHeaders, + processChatFiles, setChatAuthCookie, validateAuthToken, validateChatAuth, @@ -75,7 +76,7 @@ export async function POST( } // Use the already parsed body - const { input, password, email, conversationId } = parsedBody + const { input, password, email, conversationId, files } = parsedBody // If this is an authentication request (has password or email but no input), // set auth cookie and return success @@ -88,8 +89,8 @@ export async function POST( return response } - // For chat messages, create regular response - if (!input) { + // For chat messages, create regular response (allow empty input if files are present) + if (!input && (!files || files.length === 0)) { return addCorsHeaders(createErrorResponse('No input provided', 400), request) } @@ -108,7 +109,6 @@ export async function POST( } try { - // Transform outputConfigs to selectedOutputs format (blockId_attribute format) const selectedOutputs: string[] = [] if (deployment.outputConfigs && Array.isArray(deployment.outputConfigs)) { for (const config of deployment.outputConfigs) { @@ -123,11 +123,30 @@ export async function POST( const { SSE_HEADERS } = await import('@/lib/utils') const { createFilteredResult } = await import('@/app/api/workflows/[id]/execute/route') + 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, + executionId, + } + + const uploadedFiles = await processChatFiles(files, executionContext, requestId) + + if (uploadedFiles.length > 0) { + workflowInput.files = uploadedFiles + logger.info(`[${requestId}] Successfully processed ${uploadedFiles.length} files`) + } + } + const stream = await createStreamingResponse({ requestId, workflow: { id: deployment.workflowId, userId: deployment.userId, isDeployed: true }, - input: { input, conversationId }, // Format for chat_trigger - executingUserId: deployment.userId, // Use workflow owner's ID for chat deployments + input: workflowInput, + executingUserId: deployment.userId, streamConfig: { selectedOutputs, isSecureMode: true, diff --git a/apps/sim/app/api/chat/utils.ts b/apps/sim/app/api/chat/utils.ts index d3c33b474c..b505aeac57 100644 --- a/apps/sim/app/api/chat/utils.ts +++ b/apps/sim/app/api/chat/utils.ts @@ -6,6 +6,8 @@ import { isDev } from '@/lib/environment' 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') @@ -263,3 +265,61 @@ export async function validateChatAuth( // Unknown auth type return { authorized: false, error: 'Unsupported authentication type' } } + +/** + * Process and upload chat files to execution storage + * Handles both base64 dataUrl format and direct URL pass-through + */ +export async function processChatFiles( + files: Array<{ dataUrl?: string; url?: string; name: string; type: string }>, + executionContext: { workspaceId: string; workflowId: string; executionId: string }, + requestId: string +): Promise { + const uploadedFiles: UserFile[] = [] + + 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 +} diff --git a/apps/sim/app/api/files/serve/[...path]/route.ts b/apps/sim/app/api/files/serve/[...path]/route.ts index 5365eab0f7..8f1462fed8 100644 --- a/apps/sim/app/api/files/serve/[...path]/route.ts +++ b/apps/sim/app/api/files/serve/[...path]/route.ts @@ -1,10 +1,11 @@ import { readFile } from 'fs/promises' -import type { NextRequest, NextResponse } from 'next/server' +import type { NextRequest } from 'next/server' +import { NextResponse } from 'next/server' +import { checkHybridAuth } from '@/lib/auth/hybrid' import { createLogger } from '@/lib/logs/console/logger' import { downloadFile, getStorageProvider, isUsingCloudStorage } from '@/lib/uploads' import { S3_KB_CONFIG } from '@/lib/uploads/setup' import '@/lib/uploads/setup.server' - import { createErrorResponse, createFileResponse, @@ -15,9 +16,6 @@ import { const logger = createLogger('FilesServeAPI') -/** - * Main API route handler for serving files - */ export async function GET( request: NextRequest, { params }: { params: Promise<{ path: string[] }> } @@ -31,27 +29,26 @@ export async function GET( logger.info('File serve request:', { path }) - // Join the path segments to get the filename or cloud key - const fullPath = path.join('/') + const authResult = await checkHybridAuth(request, { requireWorkflowId: false }) - // Check if this is a cloud file (path starts with 's3/' or 'blob/') + if (!authResult.success) { + logger.warn('Unauthorized file access attempt', { path, error: authResult.error }) + return NextResponse.json({ error: 'Unauthorized' }, { status: 401 }) + } + + const userId = authResult.userId + const fullPath = path.join('/') const isS3Path = path[0] === 's3' const isBlobPath = path[0] === 'blob' const isCloudPath = isS3Path || isBlobPath + const cloudKey = isCloudPath ? path.slice(1).join('/') : fullPath - // Use cloud handler if in production, path explicitly specifies cloud storage, or we're using cloud storage if (isUsingCloudStorage() || isCloudPath) { - // Extract the actual key (remove 's3/' or 'blob/' prefix if present) - const cloudKey = isCloudPath ? path.slice(1).join('/') : fullPath - - // Get bucket type from query parameter const bucketType = request.nextUrl.searchParams.get('bucket') - - return await handleCloudProxy(cloudKey, bucketType) + return await handleCloudProxy(cloudKey, bucketType, userId) } - // Use local handler for local files - return await handleLocalFile(fullPath) + return await handleLocalFile(fullPath, userId) } catch (error) { logger.error('Error serving file:', error) @@ -63,10 +60,7 @@ export async function GET( } } -/** - * Handle local file serving - */ -async function handleLocalFile(filename: string): Promise { +async function handleLocalFile(filename: string, userId?: string): Promise { try { const filePath = findLocalFile(filename) @@ -77,6 +71,8 @@ async function handleLocalFile(filename: string): Promise { const fileBuffer = await readFile(filePath) const contentType = getContentType(filename) + logger.info('Local file served', { userId, filename, size: fileBuffer.length }) + return createFileResponse({ buffer: fileBuffer, contentType, @@ -112,12 +108,10 @@ async function downloadKBFile(cloudKey: string): Promise { throw new Error(`Unsupported storage provider for KB files: ${storageProvider}`) } -/** - * Proxy cloud file through our server - */ async function handleCloudProxy( cloudKey: string, - bucketType?: string | null + bucketType?: string | null, + userId?: string ): Promise { try { // Check if this is a KB file (starts with 'kb/') @@ -156,6 +150,13 @@ async function handleCloudProxy( const originalFilename = cloudKey.split('/').pop() || 'download' const contentType = getContentType(originalFilename) + logger.info('Cloud file served', { + userId, + key: cloudKey, + size: fileBuffer.length, + bucket: bucketType || 'default', + }) + return createFileResponse({ buffer: fileBuffer, contentType, diff --git a/apps/sim/app/api/files/upload/route.ts b/apps/sim/app/api/files/upload/route.ts index d0824c7e2b..68cd149327 100644 --- a/apps/sim/app/api/files/upload/route.ts +++ b/apps/sim/app/api/files/upload/route.ts @@ -123,8 +123,7 @@ export async function POST(request: NextRequest) { } } - // Create the serve path - const servePath = `/api/files/serve/${result.key}` + const servePath = result.path const uploadResult = { name: originalName, diff --git a/apps/sim/app/api/files/utils.ts b/apps/sim/app/api/files/utils.ts index 2e88f0c9c4..22460e4c74 100644 --- a/apps/sim/app/api/files/utils.ts +++ b/apps/sim/app/api/files/utils.ts @@ -307,6 +307,22 @@ function getSecureFileHeaders(filename: string, originalContentType: string) { } } +/** + * Encode filename for Content-Disposition header to support non-ASCII characters + * Uses RFC 5987 encoding for international characters + */ +function encodeFilenameForHeader(filename: string): string { + const hasNonAscii = /[^\x00-\x7F]/.test(filename) + + if (!hasNonAscii) { + return `filename="${filename}"` + } + + const encodedFilename = encodeURIComponent(filename) + const asciiSafe = filename.replace(/[^\x00-\x7F]/g, '_') + return `filename="${asciiSafe}"; filename*=UTF-8''${encodedFilename}` +} + /** * Create a file response with appropriate security headers */ @@ -317,7 +333,7 @@ export function createFileResponse(file: FileResponse): NextResponse { status: 200, headers: { 'Content-Type': contentType, - 'Content-Disposition': `${disposition}; filename="${file.filename}"`, + 'Content-Disposition': `${disposition}; ${encodeFilenameForHeader(file.filename)}`, 'Cache-Control': 'public, max-age=31536000', // Cache for 1 year 'X-Content-Type-Options': 'nosniff', 'Content-Security-Policy': "default-src 'none'; style-src 'unsafe-inline'; sandbox;", diff --git a/apps/sim/app/api/wand-generate/route.ts b/apps/sim/app/api/wand-generate/route.ts index 099bc3ea55..88cb8442c4 100644 --- a/apps/sim/app/api/wand-generate/route.ts +++ b/apps/sim/app/api/wand-generate/route.ts @@ -3,6 +3,7 @@ import { userStats, workflow } from '@sim/db/schema' import { eq, sql } from 'drizzle-orm' import { type NextRequest, NextResponse } from 'next/server' import OpenAI, { AzureOpenAI } from 'openai' +import { checkAndBillOverageThreshold } from '@/lib/billing/threshold-billing' import { env } from '@/lib/env' import { getCostMultiplier, isBillingEnabled } from '@/lib/environment' import { createLogger } from '@/lib/logs/console/logger' @@ -133,6 +134,9 @@ async function updateUserStatsForWand( tokensUsed: totalTokens, costAdded: costToStore, }) + + // Check if user has hit overage threshold and bill incrementally + await checkAndBillOverageThreshold(userId) } catch (error) { logger.error(`[${requestId}] Failed to update user stats for wand usage`, error) } diff --git a/apps/sim/app/api/webhooks/[id]/route.ts b/apps/sim/app/api/webhooks/[id]/route.ts index 672d04a2d3..6561b4532f 100644 --- a/apps/sim/app/api/webhooks/[id]/route.ts +++ b/apps/sim/app/api/webhooks/[id]/route.ts @@ -282,11 +282,13 @@ export async function DELETE( if (!resolvedExternalId) { try { - const requestOrigin = new URL(request.url).origin - const effectiveOrigin = requestOrigin.includes('localhost') - ? env.NEXT_PUBLIC_APP_URL || requestOrigin - : requestOrigin - const expectedNotificationUrl = `${effectiveOrigin}/api/webhooks/trigger/${foundWebhook.path}` + if (!env.NEXT_PUBLIC_APP_URL) { + logger.error( + `[${requestId}] NEXT_PUBLIC_APP_URL not configured, cannot match Airtable webhook` + ) + throw new Error('NEXT_PUBLIC_APP_URL must be configured') + } + const expectedNotificationUrl = `${env.NEXT_PUBLIC_APP_URL}/api/webhooks/trigger/${foundWebhook.path}` const listUrl = `https://api.airtable.com/v0/bases/${baseId}/webhooks` const listResp = await fetch(listUrl, { diff --git a/apps/sim/app/api/webhooks/[id]/test-url/route.ts b/apps/sim/app/api/webhooks/[id]/test-url/route.ts index c6bea06752..b4da4142d6 100644 --- a/apps/sim/app/api/webhooks/[id]/test-url/route.ts +++ b/apps/sim/app/api/webhooks/[id]/test-url/route.ts @@ -64,13 +64,13 @@ export async function POST(request: NextRequest, { params }: { params: Promise<{ return NextResponse.json({ error: 'Forbidden' }, { status: 403 }) } - const origin = new URL(request.url).origin - const effectiveOrigin = origin.includes('localhost') - ? env.NEXT_PUBLIC_APP_URL || origin - : origin + if (!env.NEXT_PUBLIC_APP_URL) { + logger.error(`[${requestId}] NEXT_PUBLIC_APP_URL not configured`) + return NextResponse.json({ error: 'Server configuration error' }, { status: 500 }) + } const token = await signTestWebhookToken(id, ttlSeconds) - const url = `${effectiveOrigin}/api/webhooks/test/${id}?token=${encodeURIComponent(token)}` + const url = `${env.NEXT_PUBLIC_APP_URL}/api/webhooks/test/${id}?token=${encodeURIComponent(token)}` logger.info(`[${requestId}] Minted test URL for webhook ${id}`) return NextResponse.json({ diff --git a/apps/sim/app/api/webhooks/route.ts b/apps/sim/app/api/webhooks/route.ts index d0bca2da67..912406eced 100644 --- a/apps/sim/app/api/webhooks/route.ts +++ b/apps/sim/app/api/webhooks/route.ts @@ -432,25 +432,20 @@ async function createAirtableWebhookSubscription( logger.warn( `[${requestId}] Could not retrieve Airtable access token for user ${userId}. Cannot create webhook in Airtable.` ) - // Instead of silently returning, throw an error with clear user guidance throw new Error( 'Airtable account connection required. Please connect your Airtable account in the trigger configuration and try again.' ) } - const requestOrigin = new URL(request.url).origin - // Ensure origin does not point to localhost for external API calls - const effectiveOrigin = requestOrigin.includes('localhost') - ? env.NEXT_PUBLIC_APP_URL || requestOrigin // Use env var if available, fallback to original - : requestOrigin - - const notificationUrl = `${effectiveOrigin}/api/webhooks/trigger/${path}` - if (effectiveOrigin !== requestOrigin) { - logger.debug( - `[${requestId}] Remapped localhost origin to ${effectiveOrigin} for notificationUrl` + if (!env.NEXT_PUBLIC_APP_URL) { + logger.error( + `[${requestId}] NEXT_PUBLIC_APP_URL not configured, cannot register Airtable webhook` ) + throw new Error('NEXT_PUBLIC_APP_URL must be configured for Airtable webhook registration') } + const notificationUrl = `${env.NEXT_PUBLIC_APP_URL}/api/webhooks/trigger/${path}` + const airtableApiUrl = `https://api.airtable.com/v0/bases/${baseId}/webhooks` const specification: any = { @@ -549,19 +544,15 @@ async function createTelegramWebhookSubscription( return // Cannot proceed without botToken } - const requestOrigin = new URL(request.url).origin - // Ensure origin does not point to localhost for external API calls - const effectiveOrigin = requestOrigin.includes('localhost') - ? env.NEXT_PUBLIC_APP_URL || requestOrigin // Use env var if available, fallback to original - : requestOrigin - - const notificationUrl = `${effectiveOrigin}/api/webhooks/trigger/${path}` - if (effectiveOrigin !== requestOrigin) { - logger.debug( - `[${requestId}] Remapped localhost origin to ${effectiveOrigin} for notificationUrl` + if (!env.NEXT_PUBLIC_APP_URL) { + logger.error( + `[${requestId}] NEXT_PUBLIC_APP_URL not configured, cannot register Telegram webhook` ) + throw new Error('NEXT_PUBLIC_APP_URL must be configured for Telegram webhook registration') } + const notificationUrl = `${env.NEXT_PUBLIC_APP_URL}/api/webhooks/trigger/${path}` + const telegramApiUrl = `https://api.telegram.org/bot${botToken}/setWebhook` const requestBody: any = { diff --git a/apps/sim/app/api/webhooks/test/route.ts b/apps/sim/app/api/webhooks/test/route.ts index 3bd05846c8..5b8886bf23 100644 --- a/apps/sim/app/api/webhooks/test/route.ts +++ b/apps/sim/app/api/webhooks/test/route.ts @@ -2,6 +2,7 @@ import { db } from '@sim/db' import { webhook } from '@sim/db/schema' import { eq } from 'drizzle-orm' import { type NextRequest, NextResponse } from 'next/server' +import { env } from '@/lib/env' import { createLogger } from '@/lib/logs/console/logger' import { generateRequestId } from '@/lib/utils' @@ -13,7 +14,6 @@ export async function GET(request: NextRequest) { const requestId = generateRequestId() try { - // Get the webhook ID and provider from the query parameters const { searchParams } = new URL(request.url) const webhookId = searchParams.get('id') @@ -24,7 +24,6 @@ export async function GET(request: NextRequest) { logger.debug(`[${requestId}] Testing webhook with ID: ${webhookId}`) - // Find the webhook in the database const webhooks = await db.select().from(webhook).where(eq(webhook.id, webhookId)).limit(1) if (webhooks.length === 0) { @@ -36,8 +35,14 @@ export async function GET(request: NextRequest) { const provider = foundWebhook.provider || 'generic' const providerConfig = (foundWebhook.providerConfig as Record) || {} - // Construct the webhook URL - const baseUrl = new URL(request.url).origin + if (!env.NEXT_PUBLIC_APP_URL) { + logger.error(`[${requestId}] NEXT_PUBLIC_APP_URL not configured, cannot test webhook`) + return NextResponse.json( + { success: false, error: 'NEXT_PUBLIC_APP_URL must be configured' }, + { status: 500 } + ) + } + const baseUrl = env.NEXT_PUBLIC_APP_URL const webhookUrl = `${baseUrl}/api/webhooks/trigger/${foundWebhook.path}` logger.info(`[${requestId}] Testing webhook for provider: ${provider}`, { @@ -46,7 +51,6 @@ export async function GET(request: NextRequest) { isActive: foundWebhook.isActive, }) - // Provider-specific test logic switch (provider) { case 'whatsapp': { const verificationToken = providerConfig.verificationToken @@ -59,10 +63,8 @@ export async function GET(request: NextRequest) { ) } - // Generate a test challenge const challenge = `test_${Date.now()}` - // Construct the WhatsApp verification URL const whatsappUrl = `${webhookUrl}?hub.mode=subscribe&hub.verify_token=${verificationToken}&hub.challenge=${challenge}` logger.debug(`[${requestId}] Testing WhatsApp webhook verification`, { @@ -70,19 +72,16 @@ export async function GET(request: NextRequest) { challenge, }) - // Make a request to the webhook endpoint const response = await fetch(whatsappUrl, { headers: { 'User-Agent': 'facebookplatform/1.0', }, }) - // Get the response details const status = response.status const contentType = response.headers.get('content-type') const responseText = await response.text() - // Check if the test was successful const success = status === 200 && responseText === challenge if (success) { @@ -139,7 +138,6 @@ export async function GET(request: NextRequest) { ) } - // Test the webhook endpoint with a simple message to check if it's reachable const testMessage = { update_id: 12345, message: { @@ -165,7 +163,6 @@ export async function GET(request: NextRequest) { url: webhookUrl, }) - // Make a test request to the webhook endpoint const response = await fetch(webhookUrl, { method: 'POST', headers: { @@ -175,16 +172,12 @@ export async function GET(request: NextRequest) { body: JSON.stringify(testMessage), }) - // Get the response details const status = response.status let responseText = '' try { responseText = await response.text() - } catch (_e) { - // Ignore if we can't get response text - } + } catch (_e) {} - // Consider success if we get a 2xx response const success = status >= 200 && status < 300 if (success) { @@ -196,7 +189,6 @@ export async function GET(request: NextRequest) { }) } - // Get webhook info from Telegram API let webhookInfo = null try { const webhookInfoUrl = `https://api.telegram.org/bot${botToken}/getWebhookInfo` @@ -215,7 +207,6 @@ export async function GET(request: NextRequest) { logger.warn(`[${requestId}] Failed to get Telegram webhook info`, e) } - // Format the curl command for testing const curlCommand = [ `curl -X POST "${webhookUrl}"`, `-H "Content-Type: application/json"`, @@ -288,16 +279,13 @@ export async function GET(request: NextRequest) { } case 'generic': { - // Get the general webhook configuration const token = providerConfig.token const secretHeaderName = providerConfig.secretHeaderName const requireAuth = providerConfig.requireAuth const allowedIps = providerConfig.allowedIps - // Generate sample curl command for testing let curlCommand = `curl -X POST "${webhookUrl}" -H "Content-Type: application/json"` - // Add auth headers to the curl command if required if (requireAuth && token) { if (secretHeaderName) { curlCommand += ` -H "${secretHeaderName}: ${token}"` @@ -306,7 +294,6 @@ export async function GET(request: NextRequest) { } } - // Add a sample payload curlCommand += ` -d '{"event":"test_event","timestamp":"${new Date().toISOString()}"}'` logger.info(`[${requestId}] General webhook test successful: ${webhookId}`) @@ -391,7 +378,6 @@ export async function GET(request: NextRequest) { }) } - // Add the Airtable test case case 'airtable': { const baseId = providerConfig.baseId const tableId = providerConfig.tableId @@ -408,7 +394,6 @@ export async function GET(request: NextRequest) { ) } - // Define a sample payload structure const samplePayload = { webhook: { id: 'whiYOUR_WEBHOOK_ID', @@ -418,16 +403,15 @@ export async function GET(request: NextRequest) { }, payloadFormat: 'v0', actionMetadata: { - source: 'tableOrViewChange', // Example source + source: 'tableOrViewChange', sourceMetadata: {}, }, payloads: [ { timestamp: new Date().toISOString(), - baseTransactionNumber: Date.now(), // Example transaction number + baseTransactionNumber: Date.now(), changedTablesById: { [tableId]: { - // Example changes - structure may vary based on actual event changedRecordsById: { recSAMPLEID1: { current: { cellValuesByFieldId: { fldSAMPLEID: 'New Value' } }, @@ -442,7 +426,6 @@ export async function GET(request: NextRequest) { ], } - // Generate sample curl command let curlCommand = `curl -X POST "${webhookUrl}" -H "Content-Type: application/json"` curlCommand += ` -d '${JSON.stringify(samplePayload, null, 2)}'` @@ -519,7 +502,6 @@ export async function GET(request: NextRequest) { } default: { - // Generic webhook test logger.info(`[${requestId}] Generic webhook test successful: ${webhookId}`) return NextResponse.json({ success: true, diff --git a/apps/sim/app/api/workflows/[id]/autolayout/route.ts b/apps/sim/app/api/workflows/[id]/autolayout/route.ts index 5d9c896143..4596d15e10 100644 --- a/apps/sim/app/api/workflows/[id]/autolayout/route.ts +++ b/apps/sim/app/api/workflows/[id]/autolayout/route.ts @@ -1,17 +1,14 @@ -import { db } from '@sim/db' -import { workflow as workflowTable } from '@sim/db/schema' -import { eq } from 'drizzle-orm' import { type NextRequest, NextResponse } from 'next/server' import { z } from 'zod' import { getSession } from '@/lib/auth' import { createLogger } from '@/lib/logs/console/logger' -import { getUserEntityPermissions } from '@/lib/permissions/utils' import { generateRequestId } from '@/lib/utils' import { applyAutoLayout } from '@/lib/workflows/autolayout' import { loadWorkflowFromNormalizedTables, type NormalizedWorkflowData, } from '@/lib/workflows/db-helpers' +import { getWorkflowAccessContext } from '@/lib/workflows/utils' export const dynamic = 'force-dynamic' @@ -77,11 +74,8 @@ export async function POST(request: NextRequest, { params }: { params: Promise<{ }) // Fetch the workflow to check ownership/access - const workflowData = await db - .select() - .from(workflowTable) - .where(eq(workflowTable.id, workflowId)) - .then((rows) => rows[0]) + const accessContext = await getWorkflowAccessContext(workflowId, userId) + const workflowData = accessContext?.workflow if (!workflowData) { logger.warn(`[${requestId}] Workflow ${workflowId} not found for autolayout`) @@ -89,24 +83,12 @@ export async function POST(request: NextRequest, { params }: { params: Promise<{ } // Check if user has permission to update this workflow - let canUpdate = false - - // Case 1: User owns the workflow - if (workflowData.userId === userId) { - canUpdate = true - } - - // Case 2: Workflow belongs to a workspace and user has write or admin permission - if (!canUpdate && workflowData.workspaceId) { - const userPermission = await getUserEntityPermissions( - userId, - 'workspace', - workflowData.workspaceId - ) - if (userPermission === 'write' || userPermission === 'admin') { - canUpdate = true - } - } + const canUpdate = + accessContext?.isOwner || + (workflowData.workspaceId + ? accessContext?.workspacePermission === 'write' || + accessContext?.workspacePermission === 'admin' + : false) if (!canUpdate) { logger.warn( diff --git a/apps/sim/app/api/workflows/[id]/duplicate/route.ts b/apps/sim/app/api/workflows/[id]/duplicate/route.ts index a6f61ae80f..9daec4b637 100644 --- a/apps/sim/app/api/workflows/[id]/duplicate/route.ts +++ b/apps/sim/app/api/workflows/[id]/duplicate/route.ts @@ -47,17 +47,17 @@ export async function POST(req: NextRequest, { params }: { params: Promise<{ id: // Duplicate workflow and all related data in a transaction const result = await db.transaction(async (tx) => { // First verify the source workflow exists - const sourceWorkflow = await tx + const sourceWorkflowRow = await tx .select() .from(workflow) .where(eq(workflow.id, sourceWorkflowId)) .limit(1) - if (sourceWorkflow.length === 0) { + if (sourceWorkflowRow.length === 0) { throw new Error('Source workflow not found') } - const source = sourceWorkflow[0] + const source = sourceWorkflowRow[0] // Check if user has permission to access the source workflow let canAccessSource = false diff --git a/apps/sim/app/api/workflows/[id]/execute/route.ts b/apps/sim/app/api/workflows/[id]/execute/route.ts index 4d12560d14..0262a30447 100644 --- a/apps/sim/app/api/workflows/[id]/execute/route.ts +++ b/apps/sim/app/api/workflows/[id]/execute/route.ts @@ -23,7 +23,11 @@ import { workflowHasResponseBlock, } from '@/lib/workflows/utils' import { validateWorkflowAccess } from '@/app/api/workflows/middleware' -import { createErrorResponse, createSuccessResponse } from '@/app/api/workflows/utils' +import { + createErrorResponse, + createSuccessResponse, + processApiWorkflowField, +} from '@/app/api/workflows/utils' import { Executor } from '@/executor' import type { ExecutionResult } from '@/executor/types' import { Serializer } from '@/serializer' @@ -74,16 +78,13 @@ function resolveOutputIds( return selectedOutputs } - // UUID regex to detect if it's already in blockId_attribute format const UUID_REGEX = /^[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}/i return selectedOutputs.map((outputId) => { - // If it starts with a UUID, it's already in blockId_attribute format (from chat deployments) if (UUID_REGEX.test(outputId)) { return outputId } - // Otherwise, it's in blockName.path format from the user/API const dotIndex = outputId.indexOf('.') if (dotIndex === -1) { logger.warn(`Invalid output ID format (missing dot): ${outputId}`) @@ -93,7 +94,6 @@ function resolveOutputIds( const blockName = outputId.substring(0, dotIndex) const path = outputId.substring(dotIndex + 1) - // Find the block by name (case-insensitive, ignoring spaces) const normalizedBlockName = blockName.toLowerCase().replace(/\s+/g, '') const block = Object.values(blocks).find((b: any) => { const normalized = (b.name || '').toLowerCase().replace(/\s+/g, '') @@ -137,9 +137,6 @@ export async function executeWorkflow( const loggingSession = new LoggingSession(workflowId, executionId, 'api', requestId) - // Rate limiting is now handled before entering the sync queue - - // Check if the actor has exceeded their usage limits const usageCheck = await checkServerSideUsageLimits(actorUserId) if (usageCheck.isExceeded) { logger.warn(`[${requestId}] User ${workflow.userId} has exceeded usage limits`, { @@ -151,13 +148,11 @@ export async function executeWorkflow( ) } - // Log input to help debug logger.info( `[${requestId}] Executing workflow with input:`, input ? JSON.stringify(input, null, 2) : 'No input provided' ) - // Use input directly for API workflows const processedInput = input logger.info( `[${requestId}] Using input directly for workflow:`, @@ -168,10 +163,7 @@ export async function executeWorkflow( runningExecutions.add(executionKey) logger.info(`[${requestId}] Starting workflow execution: ${workflowId}`) - // Load workflow data from deployed state for API executions const deployedData = await loadDeployedWorkflowState(workflowId) - - // Use deployed data as primary source for API executions const { blocks, edges, loops, parallels } = deployedData logger.info(`[${requestId}] Using deployed state for workflow execution: ${workflowId}`) logger.debug(`[${requestId}] Deployed data loaded:`, { @@ -181,10 +173,8 @@ export async function executeWorkflow( parallelsCount: Object.keys(parallels || {}).length, }) - // Use the same execution flow as in scheduled executions const mergedStates = mergeSubblockState(blocks) - // Load personal (for the executing user) and workspace env (workspace overrides personal) const { personalEncrypted, workspaceEncrypted } = await getPersonalAndWorkspaceEnv( actorUserId, workflow.workspaceId || undefined @@ -197,7 +187,6 @@ export async function executeWorkflow( variables, }) - // Replace environment variables in the block states const currentBlockStates = await Object.entries(mergedStates).reduce( async (accPromise, [id, block]) => { const acc = await accPromise @@ -206,13 +195,11 @@ export async function executeWorkflow( const subAcc = await subAccPromise let value = subBlock.value - // If the value is a string and contains environment variable syntax if (typeof value === 'string' && value.includes('{{') && value.includes('}}')) { const matches = value.match(/{{([^}]+)}}/g) if (matches) { - // Process all matches sequentially for (const match of matches) { - const varName = match.slice(2, -2) // Remove {{ and }} + const varName = match.slice(2, -2) const encryptedValue = variables[varName] if (!encryptedValue) { throw new Error(`Environment variable "${varName}" was not found`) @@ -244,7 +231,6 @@ export async function executeWorkflow( Promise.resolve({} as Record>) ) - // Create a map of decrypted environment variables const decryptedEnvVars: Record = {} for (const [key, encryptedValue] of Object.entries(variables)) { try { @@ -256,22 +242,17 @@ export async function executeWorkflow( } } - // Process the block states to ensure response formats are properly parsed const processedBlockStates = Object.entries(currentBlockStates).reduce( (acc, [blockId, blockState]) => { - // Check if this block has a responseFormat that needs to be parsed if (blockState.responseFormat && typeof blockState.responseFormat === 'string') { const responseFormatValue = blockState.responseFormat.trim() - // Check for variable references like if (responseFormatValue.startsWith('<') && responseFormatValue.includes('>')) { logger.debug( `[${requestId}] Response format contains variable reference for block ${blockId}` ) - // Keep variable references as-is - they will be resolved during execution acc[blockId] = blockState } else if (responseFormatValue === '') { - // Empty string - remove response format acc[blockId] = { ...blockState, responseFormat: undefined, @@ -279,7 +260,6 @@ export async function executeWorkflow( } else { try { logger.debug(`[${requestId}] Parsing responseFormat for block ${blockId}`) - // Attempt to parse the responseFormat if it's a string const parsedResponseFormat = JSON.parse(responseFormatValue) acc[blockId] = { @@ -291,7 +271,6 @@ export async function executeWorkflow( `[${requestId}] Failed to parse responseFormat for block ${blockId}, using undefined`, error ) - // Set to undefined instead of keeping malformed JSON - this allows execution to continue acc[blockId] = { ...blockState, responseFormat: undefined, @@ -306,7 +285,6 @@ export async function executeWorkflow( {} as Record> ) - // Get workflow variables - they are stored as JSON objects in the database const workflowVariables = (workflow.variables as Record) || {} if (Object.keys(workflowVariables).length > 0) { @@ -317,20 +295,15 @@ export async function executeWorkflow( logger.debug(`[${requestId}] No workflow variables found for: ${workflowId}`) } - // Serialize and execute the workflow logger.debug(`[${requestId}] Serializing workflow: ${workflowId}`) const serializedWorkflow = new Serializer().serializeWorkflow( mergedStates, edges, loops, parallels, - true // Enable validation during execution + true ) - // Determine trigger start block based on execution type - // - 'chat': For chat deployments (looks for chat_trigger block) - // - 'api': For direct API execution (looks for api_trigger block) - // streamConfig is passed from POST handler when using streaming/chat const preferredTriggerType = streamConfig?.workflowTriggerType || 'api' const startBlock = TriggerUtils.findStartBlock(mergedStates, preferredTriggerType, false) @@ -346,8 +319,6 @@ export async function executeWorkflow( const startBlockId = startBlock.blockId const triggerBlock = startBlock.block - // Check if the API trigger has any outgoing connections (except for legacy starter blocks) - // Legacy starter blocks have their own validation in the executor if (triggerBlock.type !== 'starter') { const outgoingConnections = serializedWorkflow.connections.filter( (conn) => conn.source === startBlockId @@ -358,14 +329,12 @@ export async function executeWorkflow( } } - // Build context extensions const contextExtensions: any = { executionId, workspaceId: workflow.workspaceId, isDeployedContext: true, } - // Add streaming configuration if enabled if (streamConfig?.enabled) { contextExtensions.stream = true contextExtensions.selectedOutputs = streamConfig.selectedOutputs || [] @@ -386,10 +355,8 @@ export async function executeWorkflow( contextExtensions, }) - // Set up logging on the executor loggingSession.setupExecutor(executor) - // Execute workflow (will always return ExecutionResult since we don't use onStream) const result = (await executor.execute(workflowId, startBlockId)) as ExecutionResult logger.info(`[${requestId}] Workflow execution completed: ${workflowId}`, { @@ -397,14 +364,11 @@ export async function executeWorkflow( executionTime: result.metadata?.duration, }) - // Build trace spans from execution result (works for both success and failure) const { traceSpans, totalDuration } = buildTraceSpans(result) - // Update workflow run counts if execution was successful if (result.success) { await updateWorkflowRunCounts(workflowId) - // Track API call in user stats await db .update(userStats) .set({ @@ -419,9 +383,9 @@ export async function executeWorkflow( totalDurationMs: totalDuration || 0, finalOutput: result.output || {}, traceSpans: (traceSpans || []) as any, + workflowInput: processedInput, }) - // For non-streaming, return the execution result return result } catch (error: any) { logger.error(`[${requestId}] Workflow execution failed: ${workflowId}`, error) @@ -461,22 +425,16 @@ export async function GET(request: NextRequest, { params }: { params: Promise<{ return createErrorResponse(validation.error.message, validation.error.status) } - // Determine trigger type based on authentication let triggerType: TriggerType = 'manual' const session = await getSession() if (!session?.user?.id) { - // Check for API key const apiKeyHeader = request.headers.get('X-API-Key') if (apiKeyHeader) { triggerType = 'api' } } - // Note: Async execution is now handled in the POST handler below - - // Synchronous execution try { - // Resolve actor user id let actorUserId: string | null = null if (triggerType === 'manual') { actorUserId = session!.user!.id @@ -491,7 +449,6 @@ export async function GET(request: NextRequest, { params }: { params: Promise<{ void updateApiKeyLastUsed(auth.keyId).catch(() => {}) } - // Check rate limits BEFORE entering execution for API requests const userSubscription = await getHighestPrioritySubscription(actorUserId) const rateLimiter = new RateLimiter() const rateLimitCheck = await rateLimiter.checkRateLimitWithSubscription( @@ -514,13 +471,11 @@ export async function GET(request: NextRequest, { params }: { params: Promise<{ actorUserId as string ) - // Check if the workflow execution contains a response block output const hasResponseBlock = workflowHasResponseBlock(result) if (hasResponseBlock) { return createHttpResponseFromBlock(result) } - // Filter out logs and workflowConnections from the API response const filteredResult = createFilteredResult(result) return createSuccessResponse(filteredResult) } catch (error: any) { @@ -536,12 +491,10 @@ export async function GET(request: NextRequest, { params }: { params: Promise<{ } catch (error: any) { logger.error(`[${requestId}] Error executing workflow: ${id}`, error) - // Check if this is a rate limit error if (error instanceof RateLimitError) { return createErrorResponse(error.message, error.statusCode, 'RATE_LIMIT_EXCEEDED') } - // Check if this is a usage limit error if (error instanceof UsageLimitError) { return createErrorResponse(error.message, error.statusCode, 'USAGE_LIMIT_EXCEEDED') } @@ -566,18 +519,15 @@ export async function POST( const workflowId = id try { - // Validate workflow access const validation = await validateWorkflowAccess(request as NextRequest, id) if (validation.error) { logger.warn(`[${requestId}] Workflow access validation failed: ${validation.error.message}`) return createErrorResponse(validation.error.message, validation.error.status) } - // Check execution mode from header const executionMode = request.headers.get('X-Execution-Mode') const isAsync = executionMode === 'async' - // Parse request body first to check for internal parameters const body = await request.text() logger.info(`[${requestId}] ${body ? 'Request body provided' : 'No request body provided'}`) @@ -616,17 +566,78 @@ export async function POST( streamResponse, selectedOutputs, workflowTriggerType, - input, + input: rawInput, } = extractExecutionParams(request as NextRequest, parsedBody) - // Get authenticated user and determine trigger type + let processedInput = rawInput + logger.info(`[${requestId}] Raw input received:`, JSON.stringify(rawInput, null, 2)) + + try { + const deployedData = await loadDeployedWorkflowState(workflowId) + const blocks = deployedData.blocks || {} + logger.info(`[${requestId}] Loaded ${Object.keys(blocks).length} blocks from workflow`) + + const apiTriggerBlock = Object.values(blocks).find( + (block: any) => block.type === 'api_trigger' + ) as any + logger.info(`[${requestId}] API trigger block found:`, !!apiTriggerBlock) + + if (apiTriggerBlock?.subBlocks?.inputFormat?.value) { + const inputFormat = apiTriggerBlock.subBlocks.inputFormat.value as Array<{ + name: string + type: 'string' | 'number' | 'boolean' | 'object' | 'array' | 'files' + }> + logger.info( + `[${requestId}] Input format fields:`, + inputFormat.map((f) => `${f.name}:${f.type}`).join(', ') + ) + + const fileFields = inputFormat.filter((field) => field.type === 'files') + logger.info(`[${requestId}] Found ${fileFields.length} file-type fields`) + + if (fileFields.length > 0 && typeof rawInput === 'object' && rawInput !== null) { + const executionContext = { + workspaceId: validation.workflow.workspaceId, + workflowId, + } + + for (const fileField of fileFields) { + const fieldValue = rawInput[fileField.name] + + if (fieldValue && typeof fieldValue === 'object') { + const uploadedFiles = await processApiWorkflowField( + fieldValue, + executionContext, + requestId + ) + + if (uploadedFiles.length > 0) { + processedInput = { + ...processedInput, + [fileField.name]: uploadedFiles, + } + logger.info( + `[${requestId}] Successfully processed ${uploadedFiles.length} file(s) for field: ${fileField.name}` + ) + } + } + } + } + } + } catch (error) { + logger.error(`[${requestId}] Failed to process file uploads:`, error) + const errorMessage = error instanceof Error ? error.message : 'Failed to process file uploads' + return createErrorResponse(errorMessage, 400) + } + + const input = processedInput + let authenticatedUserId: string let triggerType: TriggerType = 'manual' - // For internal calls (chat deployments), use the workflow owner's ID if (finalIsSecureMode) { authenticatedUserId = validation.workflow.userId - triggerType = 'manual' // Chat deployments use manual trigger type (no rate limit) + triggerType = 'manual' } else { const session = await getSession() const apiKeyHeader = request.headers.get('X-API-Key') @@ -649,7 +660,6 @@ export async function POST( } } - // Get user subscription (checks both personal and org subscriptions) const userSubscription = await getHighestPrioritySubscription(authenticatedUserId) if (isAsync) { @@ -659,7 +669,7 @@ export async function POST( authenticatedUserId, userSubscription, 'api', - true // isAsync = true + true ) if (!rateLimitCheck.allowed) { @@ -683,7 +693,6 @@ export async function POST( ) } - // Rate limit passed - always use Trigger.dev for async executions const handle = await tasks.trigger('workflow-execution', { workflowId, userId: authenticatedUserId, @@ -723,7 +732,7 @@ export async function POST( authenticatedUserId, userSubscription, triggerType, - false // isAsync = false for sync calls + false ) if (!rateLimitCheck.allowed) { @@ -732,15 +741,12 @@ export async function POST( ) } - // Handle streaming response - wrap execution in SSE stream if (streamResponse) { - // Load workflow blocks to resolve output IDs from blockName.attribute to blockId_attribute format const deployedData = await loadDeployedWorkflowState(workflowId) const resolvedSelectedOutputs = selectedOutputs ? resolveOutputIds(selectedOutputs, deployedData.blocks || {}) : selectedOutputs - // Use shared streaming response creator const { createStreamingResponse } = await import('@/lib/workflows/streaming') const { SSE_HEADERS } = await import('@/lib/utils') @@ -763,7 +769,6 @@ export async function POST( }) } - // Non-streaming execution const result = await executeWorkflow( validation.workflow, requestId, @@ -772,13 +777,11 @@ export async function POST( undefined ) - // Non-streaming response const hasResponseBlock = workflowHasResponseBlock(result) if (hasResponseBlock) { return createHttpResponseFromBlock(result) } - // Filter out logs and workflowConnections from the API response const filteredResult = createFilteredResult(result) return createSuccessResponse(filteredResult) } catch (error: any) { @@ -794,17 +797,14 @@ export async function POST( } catch (error: any) { logger.error(`[${requestId}] Error executing workflow: ${workflowId}`, error) - // Check if this is a rate limit error if (error instanceof RateLimitError) { return createErrorResponse(error.message, error.statusCode, 'RATE_LIMIT_EXCEEDED') } - // Check if this is a usage limit error if (error instanceof UsageLimitError) { return createErrorResponse(error.message, error.statusCode, 'USAGE_LIMIT_EXCEEDED') } - // Check if this is a rate limit error (string match for backward compatibility) if (error.message?.includes('Rate limit exceeded')) { return createErrorResponse(error.message, 429, 'RATE_LIMIT_EXCEEDED') } diff --git a/apps/sim/app/api/workflows/[id]/route.test.ts b/apps/sim/app/api/workflows/[id]/route.test.ts index fb099ac7f8..c03d50baaf 100644 --- a/apps/sim/app/api/workflows/[id]/route.test.ts +++ b/apps/sim/app/api/workflows/[id]/route.test.ts @@ -16,6 +16,9 @@ describe('Workflow By ID API Route', () => { error: vi.fn(), } + const mockGetWorkflowById = vi.fn() + const mockGetWorkflowAccessContext = vi.fn() + beforeEach(() => { vi.resetModules() @@ -30,6 +33,20 @@ describe('Workflow By ID API Route', () => { vi.doMock('@/lib/workflows/db-helpers', () => ({ loadWorkflowFromNormalizedTables: vi.fn().mockResolvedValue(null), })) + + mockGetWorkflowById.mockReset() + mockGetWorkflowAccessContext.mockReset() + + vi.doMock('@/lib/workflows/utils', async () => { + const actual = + await vi.importActual('@/lib/workflows/utils') + + return { + ...actual, + getWorkflowById: mockGetWorkflowById, + getWorkflowAccessContext: mockGetWorkflowAccessContext, + } + }) }) afterEach(() => { @@ -60,17 +77,14 @@ describe('Workflow By ID API Route', () => { }), })) - vi.doMock('@sim/db', () => ({ - db: { - select: vi.fn().mockReturnValue({ - from: vi.fn().mockReturnValue({ - where: vi.fn().mockReturnValue({ - then: vi.fn().mockResolvedValue(undefined), - }), - }), - }), - }, - })) + mockGetWorkflowById.mockResolvedValueOnce(null) + mockGetWorkflowAccessContext.mockResolvedValueOnce({ + workflow: null, + workspaceOwnerId: null, + workspacePermission: null, + isOwner: false, + isWorkspaceOwner: false, + }) const req = new NextRequest('http://localhost:3000/api/workflows/nonexistent') const params = Promise.resolve({ id: 'nonexistent' }) @@ -105,22 +119,28 @@ describe('Workflow By ID API Route', () => { }), })) - vi.doMock('@sim/db', () => ({ - db: { - select: vi.fn().mockReturnValue({ - from: vi.fn().mockReturnValue({ - where: vi.fn().mockReturnValue({ - then: vi.fn().mockResolvedValue(mockWorkflow), - }), - }), - }), - }, - })) + mockGetWorkflowById.mockResolvedValueOnce(mockWorkflow) + mockGetWorkflowAccessContext.mockResolvedValueOnce({ + workflow: mockWorkflow, + workspaceOwnerId: null, + workspacePermission: null, + isOwner: true, + isWorkspaceOwner: false, + }) vi.doMock('@/lib/workflows/db-helpers', () => ({ loadWorkflowFromNormalizedTables: vi.fn().mockResolvedValue(mockNormalizedData), })) + mockGetWorkflowById.mockResolvedValueOnce(mockWorkflow) + mockGetWorkflowAccessContext.mockResolvedValueOnce({ + workflow: mockWorkflow, + workspaceOwnerId: null, + workspacePermission: null, + isOwner: true, + isWorkspaceOwner: false, + }) + const req = new NextRequest('http://localhost:3000/api/workflows/workflow-123') const params = Promise.resolve({ id: 'workflow-123' }) @@ -154,22 +174,28 @@ describe('Workflow By ID API Route', () => { }), })) - vi.doMock('@sim/db', () => ({ - db: { - select: vi.fn().mockReturnValue({ - from: vi.fn().mockReturnValue({ - where: vi.fn().mockReturnValue({ - then: vi.fn().mockResolvedValue(mockWorkflow), - }), - }), - }), - }, - })) + mockGetWorkflowById.mockResolvedValueOnce(mockWorkflow) + mockGetWorkflowAccessContext.mockResolvedValueOnce({ + workflow: mockWorkflow, + workspaceOwnerId: 'workspace-456', + workspacePermission: 'admin', + isOwner: false, + isWorkspaceOwner: false, + }) vi.doMock('@/lib/workflows/db-helpers', () => ({ loadWorkflowFromNormalizedTables: vi.fn().mockResolvedValue(mockNormalizedData), })) + mockGetWorkflowById.mockResolvedValueOnce(mockWorkflow) + mockGetWorkflowAccessContext.mockResolvedValueOnce({ + workflow: mockWorkflow, + workspaceOwnerId: 'workspace-456', + workspacePermission: 'read', + isOwner: false, + isWorkspaceOwner: false, + }) + vi.doMock('@/lib/permissions/utils', () => ({ getUserEntityPermissions: vi.fn().mockResolvedValue('read'), hasAdminPermission: vi.fn().mockResolvedValue(false), @@ -200,22 +226,14 @@ describe('Workflow By ID API Route', () => { }), })) - vi.doMock('@sim/db', () => ({ - db: { - select: vi.fn().mockReturnValue({ - from: vi.fn().mockReturnValue({ - where: vi.fn().mockReturnValue({ - then: vi.fn().mockResolvedValue(mockWorkflow), - }), - }), - }), - }, - })) - - vi.doMock('@/lib/permissions/utils', () => ({ - getUserEntityPermissions: vi.fn().mockResolvedValue(null), - hasAdminPermission: vi.fn().mockResolvedValue(false), - })) + mockGetWorkflowById.mockResolvedValueOnce(mockWorkflow) + mockGetWorkflowAccessContext.mockResolvedValueOnce({ + workflow: mockWorkflow, + workspaceOwnerId: 'workspace-456', + workspacePermission: null, + isOwner: false, + isWorkspaceOwner: false, + }) const req = new NextRequest('http://localhost:3000/api/workflows/workflow-123') const params = Promise.resolve({ id: 'workflow-123' }) @@ -250,17 +268,14 @@ describe('Workflow By ID API Route', () => { }), })) - vi.doMock('@sim/db', () => ({ - db: { - select: vi.fn().mockReturnValue({ - from: vi.fn().mockReturnValue({ - where: vi.fn().mockReturnValue({ - then: vi.fn().mockResolvedValue(mockWorkflow), - }), - }), - }), - }, - })) + mockGetWorkflowById.mockResolvedValueOnce(mockWorkflow) + mockGetWorkflowAccessContext.mockResolvedValueOnce({ + workflow: mockWorkflow, + workspaceOwnerId: null, + workspacePermission: null, + isOwner: true, + isWorkspaceOwner: false, + }) vi.doMock('@/lib/workflows/db-helpers', () => ({ loadWorkflowFromNormalizedTables: vi.fn().mockResolvedValue(mockNormalizedData), @@ -294,19 +309,22 @@ describe('Workflow By ID API Route', () => { }), })) + mockGetWorkflowById.mockResolvedValueOnce(mockWorkflow) + mockGetWorkflowAccessContext.mockResolvedValueOnce({ + workflow: mockWorkflow, + workspaceOwnerId: null, + workspacePermission: null, + isOwner: true, + isWorkspaceOwner: false, + }) + vi.doMock('@sim/db', () => ({ db: { - select: vi.fn().mockReturnValue({ - from: vi.fn().mockReturnValue({ - where: vi.fn().mockReturnValue({ - then: vi.fn().mockResolvedValue(mockWorkflow), - }), - }), - }), delete: vi.fn().mockReturnValue({ - where: vi.fn().mockResolvedValue(undefined), + where: vi.fn().mockResolvedValue([{ id: 'workflow-123' }]), }), }, + workflow: {}, })) global.fetch = vi.fn().mockResolvedValue({ @@ -340,24 +358,22 @@ describe('Workflow By ID API Route', () => { }), })) + mockGetWorkflowById.mockResolvedValueOnce(mockWorkflow) + mockGetWorkflowAccessContext.mockResolvedValueOnce({ + workflow: mockWorkflow, + workspaceOwnerId: 'workspace-456', + workspacePermission: 'admin', + isOwner: false, + isWorkspaceOwner: false, + }) + vi.doMock('@sim/db', () => ({ db: { - select: vi.fn().mockReturnValue({ - from: vi.fn().mockReturnValue({ - where: vi.fn().mockReturnValue({ - then: vi.fn().mockResolvedValue(mockWorkflow), - }), - }), - }), delete: vi.fn().mockReturnValue({ - where: vi.fn().mockResolvedValue(undefined), + where: vi.fn().mockResolvedValue([{ id: 'workflow-123' }]), }), }, - })) - - vi.doMock('@/lib/permissions/utils', () => ({ - getUserEntityPermissions: vi.fn().mockResolvedValue('admin'), - hasAdminPermission: vi.fn().mockResolvedValue(true), + workflow: {}, })) global.fetch = vi.fn().mockResolvedValue({ @@ -391,22 +407,14 @@ describe('Workflow By ID API Route', () => { }), })) - vi.doMock('@sim/db', () => ({ - db: { - select: vi.fn().mockReturnValue({ - from: vi.fn().mockReturnValue({ - where: vi.fn().mockReturnValue({ - then: vi.fn().mockResolvedValue(mockWorkflow), - }), - }), - }), - }, - })) - - vi.doMock('@/lib/permissions/utils', () => ({ - getUserEntityPermissions: vi.fn().mockResolvedValue('read'), - hasAdminPermission: vi.fn().mockResolvedValue(false), - })) + mockGetWorkflowById.mockResolvedValueOnce(mockWorkflow) + mockGetWorkflowAccessContext.mockResolvedValueOnce({ + workflow: mockWorkflow, + workspaceOwnerId: 'workspace-456', + workspacePermission: null, + isOwner: false, + isWorkspaceOwner: false, + }) const req = new NextRequest('http://localhost:3000/api/workflows/workflow-123', { method: 'DELETE', @@ -432,6 +440,7 @@ describe('Workflow By ID API Route', () => { } const updateData = { name: 'Updated Workflow' } + const updatedWorkflow = { ...mockWorkflow, ...updateData, updatedAt: new Date() } vi.doMock('@/lib/auth', () => ({ getSession: vi.fn().mockResolvedValue({ @@ -439,23 +448,26 @@ describe('Workflow By ID API Route', () => { }), })) + mockGetWorkflowById.mockResolvedValueOnce(mockWorkflow) + mockGetWorkflowAccessContext.mockResolvedValueOnce({ + workflow: mockWorkflow, + workspaceOwnerId: null, + workspacePermission: null, + isOwner: true, + isWorkspaceOwner: false, + }) + vi.doMock('@sim/db', () => ({ db: { - select: vi.fn().mockReturnValue({ - from: vi.fn().mockReturnValue({ - where: vi.fn().mockReturnValue({ - then: vi.fn().mockResolvedValue(mockWorkflow), - }), - }), - }), update: vi.fn().mockReturnValue({ set: vi.fn().mockReturnValue({ where: vi.fn().mockReturnValue({ - returning: vi.fn().mockResolvedValue([{ ...mockWorkflow, ...updateData }]), + returning: vi.fn().mockResolvedValue([updatedWorkflow]), }), }), }), }, + workflow: {}, })) const req = new NextRequest('http://localhost:3000/api/workflows/workflow-123', { @@ -481,6 +493,7 @@ describe('Workflow By ID API Route', () => { } const updateData = { name: 'Updated Workflow' } + const updatedWorkflow = { ...mockWorkflow, ...updateData, updatedAt: new Date() } vi.doMock('@/lib/auth', () => ({ getSession: vi.fn().mockResolvedValue({ @@ -488,28 +501,26 @@ describe('Workflow By ID API Route', () => { }), })) + mockGetWorkflowById.mockResolvedValueOnce(mockWorkflow) + mockGetWorkflowAccessContext.mockResolvedValueOnce({ + workflow: mockWorkflow, + workspaceOwnerId: 'workspace-456', + workspacePermission: 'write', + isOwner: false, + isWorkspaceOwner: false, + }) + vi.doMock('@sim/db', () => ({ db: { - select: vi.fn().mockReturnValue({ - from: vi.fn().mockReturnValue({ - where: vi.fn().mockReturnValue({ - then: vi.fn().mockResolvedValue(mockWorkflow), - }), - }), - }), update: vi.fn().mockReturnValue({ set: vi.fn().mockReturnValue({ where: vi.fn().mockReturnValue({ - returning: vi.fn().mockResolvedValue([{ ...mockWorkflow, ...updateData }]), + returning: vi.fn().mockResolvedValue([updatedWorkflow]), }), }), }), }, - })) - - vi.doMock('@/lib/permissions/utils', () => ({ - getUserEntityPermissions: vi.fn().mockResolvedValue('write'), - hasAdminPermission: vi.fn().mockResolvedValue(false), + workflow: {}, })) const req = new NextRequest('http://localhost:3000/api/workflows/workflow-123', { @@ -542,22 +553,14 @@ describe('Workflow By ID API Route', () => { }), })) - vi.doMock('@sim/db', () => ({ - db: { - select: vi.fn().mockReturnValue({ - from: vi.fn().mockReturnValue({ - where: vi.fn().mockReturnValue({ - then: vi.fn().mockResolvedValue(mockWorkflow), - }), - }), - }), - }, - })) - - vi.doMock('@/lib/permissions/utils', () => ({ - getUserEntityPermissions: vi.fn().mockResolvedValue('read'), - hasAdminPermission: vi.fn().mockResolvedValue(false), - })) + mockGetWorkflowById.mockResolvedValueOnce(mockWorkflow) + mockGetWorkflowAccessContext.mockResolvedValueOnce({ + workflow: mockWorkflow, + workspaceOwnerId: 'workspace-456', + workspacePermission: 'read', + isOwner: false, + isWorkspaceOwner: false, + }) const req = new NextRequest('http://localhost:3000/api/workflows/workflow-123', { method: 'PUT', @@ -587,17 +590,14 @@ describe('Workflow By ID API Route', () => { }), })) - vi.doMock('@sim/db', () => ({ - db: { - select: vi.fn().mockReturnValue({ - from: vi.fn().mockReturnValue({ - where: vi.fn().mockReturnValue({ - then: vi.fn().mockResolvedValue(mockWorkflow), - }), - }), - }), - }, - })) + mockGetWorkflowById.mockResolvedValueOnce(mockWorkflow) + mockGetWorkflowAccessContext.mockResolvedValueOnce({ + workflow: mockWorkflow, + workspaceOwnerId: null, + workspacePermission: null, + isOwner: true, + isWorkspaceOwner: false, + }) // Invalid data - empty name const invalidData = { name: '' } @@ -625,17 +625,7 @@ describe('Workflow By ID API Route', () => { }), })) - vi.doMock('@sim/db', () => ({ - db: { - select: vi.fn().mockReturnValue({ - from: vi.fn().mockReturnValue({ - where: vi.fn().mockReturnValue({ - then: vi.fn().mockRejectedValue(new Error('Database connection timeout')), - }), - }), - }), - }, - })) + mockGetWorkflowById.mockRejectedValueOnce(new Error('Database connection timeout')) const req = new NextRequest('http://localhost:3000/api/workflows/workflow-123') const params = Promise.resolve({ id: 'workflow-123' }) diff --git a/apps/sim/app/api/workflows/[id]/route.ts b/apps/sim/app/api/workflows/[id]/route.ts index aa7b31d9ba..cd82d9fda5 100644 --- a/apps/sim/app/api/workflows/[id]/route.ts +++ b/apps/sim/app/api/workflows/[id]/route.ts @@ -8,9 +8,9 @@ import { getSession } from '@/lib/auth' import { verifyInternalToken } from '@/lib/auth/internal' import { env } from '@/lib/env' import { createLogger } from '@/lib/logs/console/logger' -import { getUserEntityPermissions, hasAdminPermission } from '@/lib/permissions/utils' import { generateRequestId } from '@/lib/utils' import { loadWorkflowFromNormalizedTables } from '@/lib/workflows/db-helpers' +import { getWorkflowAccessContext, getWorkflowById } from '@/lib/workflows/utils' const logger = createLogger('WorkflowByIdAPI') @@ -74,12 +74,8 @@ export async function GET(request: NextRequest, { params }: { params: Promise<{ userId = authenticatedUserId } - // Fetch the workflow - const workflowData = await db - .select() - .from(workflow) - .where(eq(workflow.id, workflowId)) - .then((rows) => rows[0]) + let accessContext = null + let workflowData = await getWorkflowById(workflowId) if (!workflowData) { logger.warn(`[${requestId}] Workflow ${workflowId} not found`) @@ -94,18 +90,21 @@ export async function GET(request: NextRequest, { params }: { params: Promise<{ hasAccess = true } else { // Case 1: User owns the workflow - if (workflowData.userId === userId) { - hasAccess = true - } + if (workflowData) { + accessContext = await getWorkflowAccessContext(workflowId, userId ?? undefined) - // Case 2: Workflow belongs to a workspace the user has permissions for - if (!hasAccess && workflowData.workspaceId && userId) { - const userPermission = await getUserEntityPermissions( - userId, - 'workspace', - workflowData.workspaceId - ) - if (userPermission !== null) { + if (!accessContext) { + logger.warn(`[${requestId}] Workflow ${workflowId} not found`) + return NextResponse.json({ error: 'Workflow not found' }, { status: 404 }) + } + + workflowData = accessContext.workflow + + if (accessContext.isOwner) { + hasAccess = true + } + + if (!hasAccess && workflowData.workspaceId && accessContext.workspacePermission) { hasAccess = true } } @@ -179,11 +178,8 @@ export async function DELETE( const userId = session.user.id - const workflowData = await db - .select() - .from(workflow) - .where(eq(workflow.id, workflowId)) - .then((rows) => rows[0]) + const accessContext = await getWorkflowAccessContext(workflowId, userId) + const workflowData = accessContext?.workflow || (await getWorkflowById(workflowId)) if (!workflowData) { logger.warn(`[${requestId}] Workflow ${workflowId} not found for deletion`) @@ -200,8 +196,8 @@ export async function DELETE( // Case 2: Workflow belongs to a workspace and user has admin permission if (!canDelete && workflowData.workspaceId) { - const hasAdmin = await hasAdminPermission(userId, workflowData.workspaceId) - if (hasAdmin) { + const context = accessContext || (await getWorkflowAccessContext(workflowId, userId)) + if (context?.workspacePermission === 'admin') { canDelete = true } } @@ -320,11 +316,8 @@ export async function PUT(request: NextRequest, { params }: { params: Promise<{ const updates = UpdateWorkflowSchema.parse(body) // Fetch the workflow to check ownership/access - const workflowData = await db - .select() - .from(workflow) - .where(eq(workflow.id, workflowId)) - .then((rows) => rows[0]) + const accessContext = await getWorkflowAccessContext(workflowId, userId) + const workflowData = accessContext?.workflow || (await getWorkflowById(workflowId)) if (!workflowData) { logger.warn(`[${requestId}] Workflow ${workflowId} not found for update`) @@ -341,12 +334,8 @@ export async function PUT(request: NextRequest, { params }: { params: Promise<{ // Case 2: Workflow belongs to a workspace and user has write or admin permission if (!canUpdate && workflowData.workspaceId) { - const userPermission = await getUserEntityPermissions( - userId, - 'workspace', - workflowData.workspaceId - ) - if (userPermission === 'write' || userPermission === 'admin') { + const context = accessContext || (await getWorkflowAccessContext(workflowId, userId)) + if (context?.workspacePermission === 'write' || context?.workspacePermission === 'admin') { canUpdate = true } } diff --git a/apps/sim/app/api/workflows/[id]/state/route.ts b/apps/sim/app/api/workflows/[id]/state/route.ts index 5ec3903f5e..e197246730 100644 --- a/apps/sim/app/api/workflows/[id]/state/route.ts +++ b/apps/sim/app/api/workflows/[id]/state/route.ts @@ -5,10 +5,10 @@ import { type NextRequest, NextResponse } from 'next/server' import { z } from 'zod' import { getSession } from '@/lib/auth' import { createLogger } from '@/lib/logs/console/logger' -import { getUserEntityPermissions } from '@/lib/permissions/utils' import { generateRequestId } from '@/lib/utils' import { extractAndPersistCustomTools } from '@/lib/workflows/custom-tools-persistence' import { saveWorkflowToNormalizedTables } from '@/lib/workflows/db-helpers' +import { getWorkflowAccessContext } from '@/lib/workflows/utils' import { sanitizeAgentToolsInBlocks } from '@/lib/workflows/validation' const logger = createLogger('WorkflowStateAPI') @@ -124,11 +124,8 @@ export async function PUT(request: NextRequest, { params }: { params: Promise<{ const state = WorkflowStateSchema.parse(body) // Fetch the workflow to check ownership/access - const workflowData = await db - .select() - .from(workflow) - .where(eq(workflow.id, workflowId)) - .then((rows) => rows[0]) + const accessContext = await getWorkflowAccessContext(workflowId, userId) + const workflowData = accessContext?.workflow if (!workflowData) { logger.warn(`[${requestId}] Workflow ${workflowId} not found for state update`) @@ -136,24 +133,12 @@ export async function PUT(request: NextRequest, { params }: { params: Promise<{ } // Check if user has permission to update this workflow - let canUpdate = false - - // Case 1: User owns the workflow - if (workflowData.userId === userId) { - canUpdate = true - } - - // Case 2: Workflow belongs to a workspace and user has write or admin permission - if (!canUpdate && workflowData.workspaceId) { - const userPermission = await getUserEntityPermissions( - userId, - 'workspace', - workflowData.workspaceId - ) - if (userPermission === 'write' || userPermission === 'admin') { - canUpdate = true - } - } + const canUpdate = + accessContext?.isOwner || + (workflowData.workspaceId + ? accessContext?.workspacePermission === 'write' || + accessContext?.workspacePermission === 'admin' + : false) if (!canUpdate) { logger.warn( diff --git a/apps/sim/app/api/workflows/[id]/variables/route.test.ts b/apps/sim/app/api/workflows/[id]/variables/route.test.ts index 7b27dd2b63..f7e105d3c9 100644 --- a/apps/sim/app/api/workflows/[id]/variables/route.test.ts +++ b/apps/sim/app/api/workflows/[id]/variables/route.test.ts @@ -18,12 +18,18 @@ import { describe('Workflow Variables API Route', () => { let authMocks: ReturnType let databaseMocks: ReturnType + const mockGetWorkflowAccessContext = vi.fn() beforeEach(() => { vi.resetModules() setupCommonApiMocks() mockCryptoUuid('mock-request-id-12345678') authMocks = mockAuth(mockUser) + mockGetWorkflowAccessContext.mockReset() + + vi.doMock('@/lib/workflows/utils', () => ({ + getWorkflowAccessContext: mockGetWorkflowAccessContext, + })) }) afterEach(() => { @@ -47,9 +53,7 @@ describe('Workflow Variables API Route', () => { it('should return 404 when workflow does not exist', async () => { authMocks.setAuthenticated({ id: 'user-123', email: 'test@example.com' }) - databaseMocks = createMockDatabase({ - select: { results: [[]] }, // No workflow found - }) + mockGetWorkflowAccessContext.mockResolvedValueOnce(null) const req = new NextRequest('http://localhost:3000/api/workflows/nonexistent/variables') const params = Promise.resolve({ id: 'nonexistent' }) @@ -73,8 +77,12 @@ describe('Workflow Variables API Route', () => { } authMocks.setAuthenticated({ id: 'user-123', email: 'test@example.com' }) - databaseMocks = createMockDatabase({ - select: { results: [[mockWorkflow]] }, + mockGetWorkflowAccessContext.mockResolvedValueOnce({ + workflow: mockWorkflow, + workspaceOwnerId: null, + workspacePermission: null, + isOwner: true, + isWorkspaceOwner: false, }) const req = new NextRequest('http://localhost:3000/api/workflows/workflow-123/variables') @@ -99,14 +107,14 @@ describe('Workflow Variables API Route', () => { } authMocks.setAuthenticated({ id: 'user-123', email: 'test@example.com' }) - databaseMocks = createMockDatabase({ - select: { results: [[mockWorkflow]] }, + mockGetWorkflowAccessContext.mockResolvedValueOnce({ + workflow: mockWorkflow, + workspaceOwnerId: 'workspace-owner', + workspacePermission: 'read', + isOwner: false, + isWorkspaceOwner: false, }) - vi.doMock('@/lib/permissions/utils', () => ({ - getUserEntityPermissions: vi.fn().mockResolvedValue('read'), - })) - const req = new NextRequest('http://localhost:3000/api/workflows/workflow-123/variables') const params = Promise.resolve({ id: 'workflow-123' }) @@ -116,14 +124,6 @@ describe('Workflow Variables API Route', () => { expect(response.status).toBe(200) const data = await response.json() expect(data.data).toEqual(mockWorkflow.variables) - - // Verify permissions check was called - const { getUserEntityPermissions } = await import('@/lib/permissions/utils') - expect(getUserEntityPermissions).toHaveBeenCalledWith( - 'user-123', - 'workspace', - 'workspace-456' - ) }) it('should deny access when user has no workspace permissions', async () => { @@ -135,14 +135,14 @@ describe('Workflow Variables API Route', () => { } authMocks.setAuthenticated({ id: 'user-123', email: 'test@example.com' }) - databaseMocks = createMockDatabase({ - select: { results: [[mockWorkflow]] }, + mockGetWorkflowAccessContext.mockResolvedValueOnce({ + workflow: mockWorkflow, + workspaceOwnerId: 'workspace-owner', + workspacePermission: null, + isOwner: false, + isWorkspaceOwner: false, }) - vi.doMock('@/lib/permissions/utils', () => ({ - getUserEntityPermissions: vi.fn().mockResolvedValue(null), - })) - const req = new NextRequest('http://localhost:3000/api/workflows/workflow-123/variables') const params = Promise.resolve({ id: 'workflow-123' }) @@ -165,8 +165,12 @@ describe('Workflow Variables API Route', () => { } authMocks.setAuthenticated({ id: 'user-123', email: 'test@example.com' }) - databaseMocks = createMockDatabase({ - select: { results: [[mockWorkflow]] }, + mockGetWorkflowAccessContext.mockResolvedValueOnce({ + workflow: mockWorkflow, + workspaceOwnerId: null, + workspacePermission: null, + isOwner: true, + isWorkspaceOwner: false, }) const req = new NextRequest('http://localhost:3000/api/workflows/workflow-123/variables') @@ -191,8 +195,15 @@ describe('Workflow Variables API Route', () => { } authMocks.setAuthenticated({ id: 'user-123', email: 'test@example.com' }) + mockGetWorkflowAccessContext.mockResolvedValueOnce({ + workflow: mockWorkflow, + workspaceOwnerId: null, + workspacePermission: null, + isOwner: true, + isWorkspaceOwner: false, + }) + databaseMocks = createMockDatabase({ - select: { results: [[mockWorkflow]] }, update: { results: [{}] }, }) @@ -223,14 +234,14 @@ describe('Workflow Variables API Route', () => { } authMocks.setAuthenticated({ id: 'user-123', email: 'test@example.com' }) - databaseMocks = createMockDatabase({ - select: { results: [[mockWorkflow]] }, + mockGetWorkflowAccessContext.mockResolvedValueOnce({ + workflow: mockWorkflow, + workspaceOwnerId: 'workspace-owner', + workspacePermission: null, + isOwner: false, + isWorkspaceOwner: false, }) - vi.doMock('@/lib/permissions/utils', () => ({ - getUserEntityPermissions: vi.fn().mockResolvedValue(null), - })) - const variables = [ { id: 'var-1', workflowId: 'workflow-123', name: 'test', type: 'string', value: 'hello' }, ] @@ -258,8 +269,12 @@ describe('Workflow Variables API Route', () => { } authMocks.setAuthenticated({ id: 'user-123', email: 'test@example.com' }) - databaseMocks = createMockDatabase({ - select: { results: [[mockWorkflow]] }, + mockGetWorkflowAccessContext.mockResolvedValueOnce({ + workflow: mockWorkflow, + workspaceOwnerId: null, + workspacePermission: null, + isOwner: true, + isWorkspaceOwner: false, }) // Invalid data - missing required fields @@ -283,9 +298,7 @@ describe('Workflow Variables API Route', () => { describe('Error handling', () => { it.concurrent('should handle database errors gracefully', async () => { authMocks.setAuthenticated({ id: 'user-123', email: 'test@example.com' }) - databaseMocks = createMockDatabase({ - select: { throwError: true, errorMessage: 'Database connection failed' }, - }) + mockGetWorkflowAccessContext.mockRejectedValueOnce(new Error('Database connection failed')) const req = new NextRequest('http://localhost:3000/api/workflows/workflow-123/variables') const params = Promise.resolve({ id: 'workflow-123' }) diff --git a/apps/sim/app/api/workflows/[id]/variables/route.ts b/apps/sim/app/api/workflows/[id]/variables/route.ts index 3c3ea5c22c..c4befb40fa 100644 --- a/apps/sim/app/api/workflows/[id]/variables/route.ts +++ b/apps/sim/app/api/workflows/[id]/variables/route.ts @@ -5,8 +5,8 @@ import { type NextRequest, NextResponse } from 'next/server' import { z } from 'zod' import { getSession } from '@/lib/auth' import { createLogger } from '@/lib/logs/console/logger' -import { getUserEntityPermissions } from '@/lib/permissions/utils' import { generateRequestId } from '@/lib/utils' +import { getWorkflowAccessContext } from '@/lib/workflows/utils' import type { Variable } from '@/stores/panel/variables/types' const logger = createLogger('WorkflowVariablesAPI') @@ -35,32 +35,18 @@ export async function POST(req: NextRequest, { params }: { params: Promise<{ id: } // Get the workflow record - const workflowRecord = await db - .select() - .from(workflow) - .where(eq(workflow.id, workflowId)) - .limit(1) + const accessContext = await getWorkflowAccessContext(workflowId, session.user.id) + const workflowData = accessContext?.workflow - if (!workflowRecord.length) { + if (!workflowData) { logger.warn(`[${requestId}] Workflow not found: ${workflowId}`) return NextResponse.json({ error: 'Workflow not found' }, { status: 404 }) } - - const workflowData = workflowRecord[0] const workspaceId = workflowData.workspaceId // Check authorization - either the user owns the workflow or has workspace permissions - let isAuthorized = workflowData.userId === session.user.id - - // If not authorized by ownership and the workflow belongs to a workspace, check workspace permissions - if (!isAuthorized && workspaceId) { - const userPermission = await getUserEntityPermissions( - session.user.id, - 'workspace', - workspaceId - ) - isAuthorized = userPermission !== null - } + const isAuthorized = + accessContext?.isOwner || (workspaceId ? accessContext?.workspacePermission !== null : false) if (!isAuthorized) { logger.warn( @@ -125,32 +111,18 @@ export async function GET(req: NextRequest, { params }: { params: Promise<{ id: } // Get the workflow record - const workflowRecord = await db - .select() - .from(workflow) - .where(eq(workflow.id, workflowId)) - .limit(1) + const accessContext = await getWorkflowAccessContext(workflowId, session.user.id) + const workflowData = accessContext?.workflow - if (!workflowRecord.length) { + if (!workflowData) { logger.warn(`[${requestId}] Workflow not found: ${workflowId}`) return NextResponse.json({ error: 'Workflow not found' }, { status: 404 }) } - - const workflowData = workflowRecord[0] const workspaceId = workflowData.workspaceId // Check authorization - either the user owns the workflow or has workspace permissions - let isAuthorized = workflowData.userId === session.user.id - - // If not authorized by ownership and the workflow belongs to a workspace, check workspace permissions - if (!isAuthorized && workspaceId) { - const userPermission = await getUserEntityPermissions( - session.user.id, - 'workspace', - workspaceId - ) - isAuthorized = userPermission !== null - } + const isAuthorized = + accessContext?.isOwner || (workspaceId ? accessContext?.workspacePermission !== null : false) if (!isAuthorized) { logger.warn( diff --git a/apps/sim/app/api/workflows/utils.ts b/apps/sim/app/api/workflows/utils.ts index 10478bcfda..27688b03e4 100644 --- a/apps/sim/app/api/workflows/utils.ts +++ b/apps/sim/app/api/workflows/utils.ts @@ -1,9 +1,14 @@ 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( { @@ -37,3 +42,99 @@ 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 { + 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 { + 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 +} diff --git a/apps/sim/app/chat/[identifier]/chat.tsx b/apps/sim/app/chat/[identifier]/chat.tsx index 85321217ce..cb940f7568 100644 --- a/apps/sim/app/chat/[identifier]/chat.tsx +++ b/apps/sim/app/chat/[identifier]/chat.tsx @@ -45,6 +45,18 @@ const DEFAULT_VOICE_SETTINGS = { voiceId: 'EXAVITQu4vr4xnSDxMaL', // Default ElevenLabs voice (Bella) } +/** + * Converts a File object to a base64 data URL + */ +function fileToBase64(file: File): Promise { + return new Promise((resolve, reject) => { + const reader = new FileReader() + reader.onload = () => resolve(reader.result as string) + reader.onerror = reject + reader.readAsDataURL(file) + }) +} + /** * Creates an audio stream handler for text-to-speech conversion * @param streamTextToAudio - Function to stream text to audio @@ -265,20 +277,43 @@ export default function ChatClient({ identifier }: { identifier: string }) { } // Handle sending a message - const handleSendMessage = async (messageParam?: string, isVoiceInput = false) => { + const handleSendMessage = async ( + messageParam?: string, + isVoiceInput = false, + files?: Array<{ + id: string + name: string + size: number + type: string + file: File + dataUrl?: string + }> + ) => { const messageToSend = messageParam ?? inputValue - if (!messageToSend.trim() || isLoading) return + if ((!messageToSend.trim() && (!files || files.length === 0)) || isLoading) return - logger.info('Sending message:', { messageToSend, isVoiceInput, conversationId }) + logger.info('Sending message:', { + messageToSend, + isVoiceInput, + conversationId, + filesCount: files?.length, + }) // Reset userHasScrolled when sending a new message setUserHasScrolled(false) const userMessage: ChatMessage = { id: crypto.randomUUID(), - content: messageToSend, + content: messageToSend || (files && files.length > 0 ? `Sent ${files.length} file(s)` : ''), type: 'user', timestamp: new Date(), + attachments: files?.map((file) => ({ + id: file.id, + name: file.name, + type: file.type, + size: file.size, + dataUrl: file.dataUrl || '', + })), } // Add the user's message to the chat @@ -299,7 +334,7 @@ export default function ChatClient({ identifier }: { identifier: string }) { try { // Send structured payload to maintain chat context - const payload = { + const payload: any = { input: typeof userMessage.content === 'string' ? userMessage.content @@ -307,7 +342,22 @@ export default function ChatClient({ identifier }: { identifier: string }) { conversationId, } - logger.info('API payload:', payload) + // Add files if present (convert to base64 for JSON transmission) + if (files && files.length > 0) { + payload.files = await Promise.all( + files.map(async (file) => ({ + name: file.name, + size: file.size, + type: file.type, + dataUrl: file.dataUrl || (await fileToBase64(file.file)), + })) + ) + } + + logger.info('API payload:', { + ...payload, + files: payload.files ? `${payload.files.length} files` : undefined, + }) const response = await fetch(`/api/chat/${identifier}`, { method: 'POST', @@ -499,8 +549,8 @@ export default function ChatClient({ identifier }: { identifier: string }) {
{ - void handleSendMessage(value, isVoiceInput) + onSubmit={(value, isVoiceInput, files) => { + void handleSendMessage(value, isVoiceInput, files) }} isStreaming={isStreamingResponse} onStopStreaming={() => stopStreaming(setMessages)} diff --git a/apps/sim/app/chat/components/input/input.tsx b/apps/sim/app/chat/components/input/input.tsx index 4f57dca188..81ec20e46c 100644 --- a/apps/sim/app/chat/components/input/input.tsx +++ b/apps/sim/app/chat/components/input/input.tsx @@ -3,7 +3,7 @@ import type React from 'react' import { useEffect, useRef, useState } from 'react' import { motion } from 'framer-motion' -import { Send, Square } from 'lucide-react' +import { AlertCircle, Paperclip, Send, Square, X } from 'lucide-react' import { Tooltip, TooltipContent, TooltipProvider, TooltipTrigger } from '@/components/ui/tooltip' import { VoiceInput } from '@/app/chat/components/input/voice-input' @@ -12,8 +12,17 @@ const PLACEHOLDER_DESKTOP = 'Enter a message or click the mic to speak' const MAX_TEXTAREA_HEIGHT = 120 // Max height in pixels (e.g., for about 3-4 lines) const MAX_TEXTAREA_HEIGHT_MOBILE = 100 // Smaller for mobile +interface AttachedFile { + id: string + name: string + size: number + type: string + file: File + dataUrl?: string +} + export const ChatInput: React.FC<{ - onSubmit?: (value: string, isVoiceInput?: boolean) => void + onSubmit?: (value: string, isVoiceInput?: boolean, files?: AttachedFile[]) => void isStreaming?: boolean onStopStreaming?: () => void onVoiceStart?: () => void @@ -21,8 +30,11 @@ export const ChatInput: React.FC<{ }> = ({ onSubmit, isStreaming = false, onStopStreaming, onVoiceStart, voiceOnly = false }) => { const wrapperRef = useRef(null) const textareaRef = useRef(null) // Ref for the textarea + const fileInputRef = useRef(null) const [isActive, setIsActive] = useState(false) const [inputValue, setInputValue] = useState('') + const [attachedFiles, setAttachedFiles] = useState([]) + const [uploadErrors, setUploadErrors] = useState([]) // Check if speech-to-text is available in the browser const isSttAvailable = @@ -85,10 +97,75 @@ export const ChatInput: React.FC<{ // Focus is now handled by the useEffect above } + // Handle file selection + const handleFileSelect = async (selectedFiles: FileList | null) => { + if (!selectedFiles) return + + const newFiles: AttachedFile[] = [] + const maxSize = 10 * 1024 * 1024 // 10MB limit + const maxFiles = 5 + + for (let i = 0; i < selectedFiles.length; i++) { + if (attachedFiles.length + newFiles.length >= maxFiles) break + + const file = selectedFiles[i] + + // Check file size + if (file.size > maxSize) { + setUploadErrors((prev) => [...prev, `${file.name} is too large (max 10MB)`]) + continue + } + + // Check for duplicates + const isDuplicate = attachedFiles.some( + (existingFile) => existingFile.name === file.name && existingFile.size === file.size + ) + if (isDuplicate) { + setUploadErrors((prev) => [...prev, `${file.name} already added`]) + continue + } + + // Read file as data URL if it's an image + let dataUrl: string | undefined + if (file.type.startsWith('image/')) { + try { + dataUrl = await new Promise((resolve, reject) => { + const reader = new FileReader() + reader.onload = () => resolve(reader.result as string) + reader.onerror = reject + reader.readAsDataURL(file) + }) + } catch (error) { + console.error('Error reading file:', error) + } + } + + newFiles.push({ + id: crypto.randomUUID(), + name: file.name, + size: file.size, + type: file.type, + file, + dataUrl, + }) + } + + if (newFiles.length > 0) { + setAttachedFiles([...attachedFiles, ...newFiles]) + setUploadErrors([]) // Clear errors when files are successfully added + } + } + + const handleRemoveFile = (fileId: string) => { + setAttachedFiles(attachedFiles.filter((f) => f.id !== fileId)) + } + const handleSubmit = () => { - if (!inputValue.trim()) return - onSubmit?.(inputValue.trim(), false) // false = not voice input + if (!inputValue.trim() && attachedFiles.length === 0) return + onSubmit?.(inputValue.trim(), false, attachedFiles) // false = not voice input setInputValue('') + setAttachedFiles([]) + setUploadErrors([]) // Clear errors when sending message if (textareaRef.current) { textareaRef.current.style.height = 'auto' // Reset height after submit textareaRef.current.style.overflowY = 'hidden' // Ensure overflow is hidden @@ -132,6 +209,29 @@ export const ChatInput: React.FC<{ <>
+ {/* Error Messages */} + {uploadErrors.length > 0 && ( +
+
+
+ +
+
+ File upload error +
+
+ {uploadErrors.map((error, idx) => ( +
+ {error} +
+ ))} +
+
+
+
+
+ )} + {/* Text Input Area with Controls */} -
- {/* Voice Input */} - {isSttAvailable && ( - - - -
- 0 && ( +
+ {attachedFiles.map((file) => { + const formatFileSize = (bytes: number) => { + if (bytes === 0) return '0 B' + const k = 1024 + const sizes = ['B', 'KB', 'MB', 'GB'] + const i = Math.floor(Math.log(bytes) / Math.log(k)) + return `${Math.round((bytes / k ** i) * 10) / 10} ${sizes[i]}` + } + + return ( +
+ {file.dataUrl ? ( + {file.name} -
- - -

Start voice conversation

- Click to enter voice mode -
- - - )} + ) : ( + <> +
+ +
+
+
+ {file.name} +
+
+ {formatFileSize(file.size)} +
+
+ + )} + +
+ ) + })} +
+ )} + +
+ {/* Paperclip Button */} + + + + + + +

Attach files

+
+
+
+ + {/* Hidden file input */} + { + handleFileSelect(e.target.files) + if (fileInputRef.current) { + fileInputRef.current.value = '' + } + }} + className='hidden' + disabled={isStreaming} + /> {/* Text Input Container */}
@@ -208,10 +380,30 @@ export const ChatInput: React.FC<{
+ {/* Voice Input */} + {isSttAvailable && ( + + + +
+ +
+
+ +

Start voice conversation

+
+
+
+ )} + {/* Send Button */} ) } diff --git a/apps/sim/app/chat/components/message/message.tsx b/apps/sim/app/chat/components/message/message.tsx index 4565ac4c57..a22bad3bbd 100644 --- a/apps/sim/app/chat/components/message/message.tsx +++ b/apps/sim/app/chat/components/message/message.tsx @@ -1,10 +1,18 @@ 'use client' import { memo, useMemo, useState } from 'react' -import { Check, Copy } from 'lucide-react' +import { Check, Copy, File as FileIcon, FileText, Image as ImageIcon } from 'lucide-react' import { Tooltip, TooltipContent, TooltipProvider, TooltipTrigger } from '@/components/ui/tooltip' import MarkdownRenderer from './components/markdown-renderer' +export interface ChatAttachment { + id: string + name: string + type: string + dataUrl: string + size?: number +} + export interface ChatMessage { id: string content: string | Record @@ -12,6 +20,7 @@ export interface ChatMessage { timestamp: Date isInitialMessage?: boolean isStreaming?: boolean + attachments?: ChatAttachment[] } function EnhancedMarkdownRenderer({ content }: { content: string }) { @@ -39,15 +48,96 @@ export const ClientChatMessage = memo( return (
+ {/* File attachments displayed above the message */} + {message.attachments && message.attachments.length > 0 && ( +
+
+ {message.attachments.map((attachment) => { + const isImage = attachment.type.startsWith('image/') + const getFileIcon = (type: string) => { + if (type.includes('pdf')) + return ( + + ) + if (type.startsWith('image/')) + return ( + + ) + if (type.includes('text') || type.includes('json')) + return ( + + ) + return ( + + ) + } + const formatFileSize = (bytes?: number) => { + if (!bytes || bytes === 0) return '' + const k = 1024 + const sizes = ['B', 'KB', 'MB', 'GB'] + const i = Math.floor(Math.log(bytes) / Math.log(k)) + return `${Math.round((bytes / k ** i) * 10) / 10} ${sizes[i]}` + } + + return ( +
{ + if (attachment.dataUrl?.trim()) { + e.preventDefault() + window.open(attachment.dataUrl, '_blank') + } + }} + > + {isImage ? ( + {attachment.name} + ) : ( + <> +
+ {getFileIcon(attachment.type)} +
+
+
+ {attachment.name} +
+ {attachment.size && ( +
+ {formatFileSize(attachment.size)} +
+ )} +
+ + )} +
+ ) + })} +
+
+ )} +
-
- {isJsonObject ? ( -
{JSON.stringify(message.content, null, 2)}
- ) : ( - {message.content as string} - )} -
+ {/* Render text content if present and not just file count message */} + {message.content && !String(message.content).startsWith('Sent') && ( +
+ {isJsonObject ? ( +
{JSON.stringify(message.content, null, 2)}
+ ) : ( + {message.content as string} + )} +
+ )}
diff --git a/apps/sim/app/workspace/[workspaceId]/w/[workflowId]/components/control-bar/components/deploy-modal/components/deployment-info/components/example-command/example-command.tsx b/apps/sim/app/workspace/[workspaceId]/w/[workflowId]/components/control-bar/components/deploy-modal/components/deployment-info/components/example-command/example-command.tsx index ce3ee3d8fc..4ba1a4cb4f 100644 --- a/apps/sim/app/workspace/[workspaceId]/w/[workflowId]/components/control-bar/components/deploy-modal/components/deployment-info/components/example-command/example-command.tsx +++ b/apps/sim/app/workspace/[workspaceId]/w/[workflowId]/components/control-bar/components/deploy-modal/components/deployment-info/components/example-command/example-command.tsx @@ -236,8 +236,8 @@ export function ExampleCommand({
)} -
-
{getDisplayCommand()}
+
+
{getDisplayCommand()}
diff --git a/apps/sim/app/workspace/[workspaceId]/w/[workflowId]/components/control-bar/components/deploy-modal/deploy-modal.tsx b/apps/sim/app/workspace/[workspaceId]/w/[workflowId]/components/control-bar/components/deploy-modal/deploy-modal.tsx index 310dbd341a..82cf7cfca3 100644 --- a/apps/sim/app/workspace/[workspaceId]/w/[workflowId]/components/control-bar/components/deploy-modal/deploy-modal.tsx +++ b/apps/sim/app/workspace/[workspaceId]/w/[workflowId]/components/control-bar/components/deploy-modal/deploy-modal.tsx @@ -139,6 +139,16 @@ export function DeployModal({ case 'array': exampleData[field.name] = [1, 2, 3] break + case 'files': + exampleData[field.name] = [ + { + data: 'data:application/pdf;base64,...', + type: 'file', + name: 'document.pdf', + mime: 'application/pdf', + }, + ] + break } } }) diff --git a/apps/sim/app/workspace/[workspaceId]/w/[workflowId]/components/panel/components/chat/chat.tsx b/apps/sim/app/workspace/[workspaceId]/w/[workflowId]/components/panel/components/chat/chat.tsx index 441d8625ea..9ecbc00426 100644 --- a/apps/sim/app/workspace/[workspaceId]/w/[workflowId]/components/panel/components/chat/chat.tsx +++ b/apps/sim/app/workspace/[workspaceId]/w/[workflowId]/components/panel/components/chat/chat.tsx @@ -1,10 +1,9 @@ 'use client' import { type KeyboardEvent, useCallback, useEffect, useMemo, useRef, useState } from 'react' -import { ArrowDown, ArrowUp } from 'lucide-react' +import { AlertCircle, ArrowDown, ArrowUp, File, FileText, Image, Paperclip, X } from 'lucide-react' import { Button } from '@/components/ui/button' import { Input } from '@/components/ui/input' -import { Notice } from '@/components/ui/notice' import { ScrollArea } from '@/components/ui/scroll-area' import { createLogger } from '@/lib/logs/console/logger' import { @@ -13,7 +12,6 @@ import { parseOutputContentSafely, } from '@/lib/response-format' import { - ChatFileUpload, ChatMessage, OutputSelect, } from '@/app/workspace/[workspaceId]/w/[workflowId]/components/panel/components/chat/components' @@ -261,12 +259,40 @@ export function Chat({ chatMessage, setChatMessage }: ChatProps) { let result: any = null try { - // Add user message + // Read files as data URLs for display in chat (only images to avoid localStorage quota issues) + const attachmentsWithData = await Promise.all( + chatFiles.map(async (file) => { + let dataUrl = '' + // Only read images as data URLs to avoid storing large files in localStorage + if (file.type.startsWith('image/')) { + try { + dataUrl = await new Promise((resolve, reject) => { + const reader = new FileReader() + reader.onload = () => resolve(reader.result as string) + reader.onerror = reject + reader.readAsDataURL(file.file) + }) + } catch (error) { + logger.error('Error reading file as data URL:', error) + } + } + return { + id: file.id, + name: file.name, + type: file.type, + size: file.size, + dataUrl, + } + }) + ) + + // Add user message with attachments (include all files, even non-images without dataUrl) addMessage({ content: sentMessage || (chatFiles.length > 0 ? `Uploaded ${chatFiles.length} file(s)` : ''), workflowId: activeWorkflowId, type: 'user', + attachments: attachmentsWithData, }) // Prepare workflow input @@ -626,66 +652,212 @@ export function Chat({ chatMessage, setChatMessage }: ChatProps) { if (validNewFiles.length > 0) { setChatFiles([...chatFiles, ...validNewFiles]) + setUploadErrors([]) // Clear errors when files are successfully added } } } }} > - {/* File upload section */} -
- {uploadErrors.length > 0 && ( -
- -
    - {uploadErrors.map((err, idx) => ( -
  • {err}
  • - ))} -
-
+ {/* Error messages */} + {uploadErrors.length > 0 && ( +
+
+
+ +
+
+ File upload error +
+
+ {uploadErrors.map((err, idx) => ( +
+ {err} +
+ ))} +
+
+
+
+
+ )} + + {/* Combined input container matching copilot style */} +
+ {/* File thumbnails */} + {chatFiles.length > 0 && ( +
+ {chatFiles.map((file) => { + const isImage = file.type.startsWith('image/') + const previewUrl = isImage ? URL.createObjectURL(file.file) : null + const getFileIcon = (type: string) => { + if (type.includes('pdf')) + return + if (type.startsWith('image/')) + return + if (type.includes('text') || type.includes('json')) + return + return + } + const formatFileSize = (bytes: number) => { + if (bytes === 0) return '0 B' + const k = 1024 + const sizes = ['B', 'KB', 'MB', 'GB'] + const i = Math.floor(Math.log(bytes) / Math.log(k)) + return `${Math.round((bytes / k ** i) * 10) / 10} ${sizes[i]}` + } + + return ( +
+ {previewUrl ? ( + {file.name} + ) : ( + <> +
+ {getFileIcon(file.type)} +
+
+
+ {file.name} +
+
+ {formatFileSize(file.size)} +
+
+ + )} + + {/* Remove button */} + +
+ ) + })}
)} - { - setChatFiles(files) - }} - maxFiles={5} - maxSize={10} - disabled={!activeWorkflowId || isExecuting || isUploadingFiles} - onError={(errors) => setUploadErrors(errors)} - /> -
-
- { - setChatMessage(e.target.value) - setHistoryIndex(-1) // Reset history index when typing - }} - onKeyDown={handleKeyPress} - placeholder={isDragOver ? 'Drop files here...' : 'Type a message...'} - className={`h-9 flex-1 rounded-lg border-[#E5E5E5] bg-[#FFFFFF] text-muted-foreground shadow-xs focus-visible:ring-0 focus-visible:ring-offset-0 dark:border-[#414141] dark:bg-[var(--surface-elevated)] ${ - isDragOver - ? 'border-[var(--brand-primary-hover-hex)] bg-purple-50/50 dark:border-[var(--brand-primary-hover-hex)] dark:bg-purple-950/20' - : '' - }`} - disabled={!activeWorkflowId || isExecuting || isUploadingFiles} - /> - + {/* Input row */} +
+ {/* Attach button */} + + + {/* Hidden file input */} + { + const files = e.target.files + if (!files) return + + const newFiles: ChatFile[] = [] + const errors: string[] = [] + for (let i = 0; i < files.length; i++) { + if (chatFiles.length + newFiles.length >= 5) { + errors.push('Maximum 5 files allowed') + break + } + const file = files[i] + if (file.size > 10 * 1024 * 1024) { + errors.push(`${file.name} is too large (max 10MB)`) + continue + } + + // Check for duplicates + const isDuplicate = chatFiles.some( + (existingFile) => + existingFile.name === file.name && existingFile.size === file.size + ) + if (isDuplicate) { + errors.push(`${file.name} already added`) + continue + } + + newFiles.push({ + id: crypto.randomUUID(), + name: file.name, + size: file.size, + type: file.type, + file, + }) + } + if (errors.length > 0) setUploadErrors(errors) + if (newFiles.length > 0) { + setChatFiles([...chatFiles, ...newFiles]) + setUploadErrors([]) // Clear errors when files are successfully added + } + e.target.value = '' + }} + className='hidden' + disabled={!activeWorkflowId || isExecuting || isUploadingFiles} + /> + + {/* Text input */} + { + setChatMessage(e.target.value) + setHistoryIndex(-1) + }} + onKeyDown={handleKeyPress} + placeholder={isDragOver ? 'Drop files here...' : 'Type a message...'} + className='h-8 flex-1 border-0 bg-transparent font-sans text-foreground text-sm shadow-none placeholder:text-muted-foreground focus-visible:ring-0 focus-visible:ring-offset-0' + disabled={!activeWorkflowId || isExecuting || isUploadingFiles} + /> + + {/* Send button */} + +
diff --git a/apps/sim/app/workspace/[workspaceId]/w/[workflowId]/components/panel/components/chat/components/chat-message/chat-message.tsx b/apps/sim/app/workspace/[workspaceId]/w/[workflowId]/components/panel/components/chat/components/chat-message/chat-message.tsx index 2709da199f..28c07c2e6d 100644 --- a/apps/sim/app/workspace/[workspaceId]/w/[workflowId]/components/panel/components/chat/components/chat-message/chat-message.tsx +++ b/apps/sim/app/workspace/[workspaceId]/w/[workflowId]/components/panel/components/chat/components/chat-message/chat-message.tsx @@ -1,4 +1,13 @@ import { useMemo } from 'react' +import { File, FileText, Image as ImageIcon } from 'lucide-react' + +interface ChatAttachment { + id: string + name: string + type: string + dataUrl: string + size?: number +} interface ChatMessageProps { message: { @@ -7,6 +16,7 @@ interface ChatMessageProps { timestamp: string | Date type: 'user' | 'workflow' isStreaming?: boolean + attachments?: ChatAttachment[] } } @@ -58,12 +68,81 @@ export function ChatMessage({ message }: ChatMessageProps) { if (message.type === 'user') { return (
+ {/* File attachments displayed above the message, completely separate from message box */} + {message.attachments && message.attachments.length > 0 && ( +
+
+ {message.attachments.map((attachment) => { + const isImage = attachment.type.startsWith('image/') + const getFileIcon = (type: string) => { + if (type.includes('pdf')) + return + if (type.startsWith('image/')) + return + if (type.includes('text') || type.includes('json')) + return + return + } + const formatFileSize = (bytes?: number) => { + if (!bytes || bytes === 0) return '' + const k = 1024 + const sizes = ['B', 'KB', 'MB', 'GB'] + const i = Math.floor(Math.log(bytes) / Math.log(k)) + return `${Math.round((bytes / k ** i) * 10) / 10} ${sizes[i]}` + } + + return ( +
{ + if (attachment.dataUrl?.trim()) { + e.preventDefault() + window.open(attachment.dataUrl, '_blank') + } + }} + > + {isImage && attachment.dataUrl ? ( + {attachment.name} + ) : ( + <> +
+ {getFileIcon(attachment.type)} +
+
+
+ {attachment.name} +
+ {attachment.size && ( +
+ {formatFileSize(attachment.size)} +
+ )} +
+ + )} +
+ ) + })} +
+
+ )} +
-
- -
+ {/* Render text content if present and not just file count message */} + {formattedContent && !formattedContent.startsWith('Uploaded') && ( +
+ +
+ )}
@@ -78,7 +157,7 @@ export function ChatMessage({ message }: ChatMessageProps) {
{message.isStreaming && ( - + )}
diff --git a/apps/sim/app/workspace/[workspaceId]/w/[workflowId]/components/workflow-block/components/sub-block/components/document-tag-entry/document-tag-entry.tsx b/apps/sim/app/workspace/[workspaceId]/w/[workflowId]/components/workflow-block/components/sub-block/components/document-tag-entry/document-tag-entry.tsx index 0f25671850..fccca2b5e4 100644 --- a/apps/sim/app/workspace/[workspaceId]/w/[workflowId]/components/workflow-block/components/sub-block/components/document-tag-entry/document-tag-entry.tsx +++ b/apps/sim/app/workspace/[workspaceId]/w/[workflowId]/components/workflow-block/components/sub-block/components/document-tag-entry/document-tag-entry.tsx @@ -9,6 +9,7 @@ import { checkTagTrigger, TagDropdown } from '@/components/ui/tag-dropdown' import { MAX_TAG_SLOTS } from '@/lib/knowledge/consts' import { cn } from '@/lib/utils' import { useSubBlockValue } from '@/app/workspace/[workspaceId]/w/[workflowId]/components/workflow-block/components/sub-block/hooks/use-sub-block-value' +import { useAccessibleReferencePrefixes } from '@/app/workspace/[workspaceId]/w/[workflowId]/hooks/use-accessible-reference-prefixes' import type { SubBlockConfig } from '@/blocks/types' import { useKnowledgeBaseTagDefinitions } from '@/hooks/use-knowledge-base-tag-definitions' import { useTagSelection } from '@/hooks/use-tag-selection' @@ -40,6 +41,7 @@ export function DocumentTagEntry({ isConnecting = false, }: DocumentTagEntryProps) { const [storeValue, setStoreValue] = useSubBlockValue(blockId, subBlock.id) + const accessiblePrefixes = useAccessibleReferencePrefixes(blockId) // Get the knowledge base ID from other sub-blocks const [knowledgeBaseIdValue] = useSubBlockValue(blockId, 'knowledgeBaseId') @@ -301,7 +303,12 @@ export function DocumentTagEntry({ )} />
-
{formatDisplayText(cellValue)}
+
+ {formatDisplayText(cellValue, { + accessiblePrefixes, + highlightAll: !accessiblePrefixes, + })} +
{showDropdown && availableTagDefinitions.length > 0 && (
@@ -389,7 +396,10 @@ export function DocumentTagEntry({ />
- {formatDisplayText(cellValue)} + {formatDisplayText(cellValue, { + accessiblePrefixes, + highlightAll: !accessiblePrefixes, + })}
{showTypeDropdown && !isReadOnly && ( @@ -469,7 +479,12 @@ export function DocumentTagEntry({ className='w-full border-0 text-transparent caret-foreground placeholder:text-muted-foreground/50 focus-visible:ring-0 focus-visible:ring-offset-0' />
-
{formatDisplayText(cellValue)}
+
+ {formatDisplayText(cellValue, { + accessiblePrefixes, + highlightAll: !accessiblePrefixes, + })} +
diff --git a/apps/sim/app/workspace/[workspaceId]/w/[workflowId]/components/workflow-block/components/sub-block/components/knowledge-tag-filters/knowledge-tag-filters.tsx b/apps/sim/app/workspace/[workspaceId]/w/[workflowId]/components/workflow-block/components/sub-block/components/knowledge-tag-filters/knowledge-tag-filters.tsx index 6df8d6582c..9b06b8649c 100644 --- a/apps/sim/app/workspace/[workspaceId]/w/[workflowId]/components/workflow-block/components/sub-block/components/knowledge-tag-filters/knowledge-tag-filters.tsx +++ b/apps/sim/app/workspace/[workspaceId]/w/[workflowId]/components/workflow-block/components/sub-block/components/knowledge-tag-filters/knowledge-tag-filters.tsx @@ -239,7 +239,12 @@ export function KnowledgeTagFilters({ onBlur={handleBlur} />
-
{formatDisplayText(cellValue || 'Select tag')}
+
+ {formatDisplayText(cellValue || 'Select tag', { + accessiblePrefixes, + highlightAll: !accessiblePrefixes, + })} +
{showDropdown && tagDefinitions.length > 0 && (
diff --git a/apps/sim/app/workspace/[workspaceId]/w/[workflowId]/components/workflow-block/components/sub-block/components/mcp-dynamic-args/mcp-dynamic-args.tsx b/apps/sim/app/workspace/[workspaceId]/w/[workflowId]/components/workflow-block/components/sub-block/components/mcp-dynamic-args/mcp-dynamic-args.tsx index 1eb5b15cd9..d40f8d0302 100644 --- a/apps/sim/app/workspace/[workspaceId]/w/[workflowId]/components/workflow-block/components/sub-block/components/mcp-dynamic-args/mcp-dynamic-args.tsx +++ b/apps/sim/app/workspace/[workspaceId]/w/[workflowId]/components/workflow-block/components/sub-block/components/mcp-dynamic-args/mcp-dynamic-args.tsx @@ -1,5 +1,6 @@ import { useCallback } from 'react' import { useParams } from 'next/navigation' +import { formatDisplayText } from '@/components/ui/formatted-text' import { Input } from '@/components/ui/input' import { Label } from '@/components/ui/label' import { @@ -14,6 +15,7 @@ import { Switch } from '@/components/ui/switch' import { Textarea } from '@/components/ui/textarea' import { cn } from '@/lib/utils' import { useSubBlockValue } from '@/app/workspace/[workspaceId]/w/[workflowId]/components/workflow-block/components/sub-block/hooks/use-sub-block-value' +import { useAccessibleReferencePrefixes } from '@/app/workspace/[workspaceId]/w/[workflowId]/hooks/use-accessible-reference-prefixes' import { useMcpTools } from '@/hooks/use-mcp-tools' import { formatParameterLabel } from '@/tools/params' @@ -37,6 +39,7 @@ export function McpDynamicArgs({ const { mcpTools } = useMcpTools(workspaceId) const [selectedTool] = useSubBlockValue(blockId, 'tool') const [toolArgs, setToolArgs] = useSubBlockValue(blockId, subBlockId) + const accessiblePrefixes = useAccessibleReferencePrefixes(blockId) const selectedToolConfig = mcpTools.find((tool) => tool.id === selectedTool) const toolSchema = selectedToolConfig?.inputSchema @@ -180,7 +183,7 @@ export function McpDynamicArgs({ case 'long-input': return ( -
+