diff --git a/sim/app/api/marketplace/[id]/info/route.ts b/sim/app/api/marketplace/[id]/info/route.ts index bc447c9bda..c3e7d95b0d 100644 --- a/sim/app/api/marketplace/[id]/info/route.ts +++ b/sim/app/api/marketplace/[id]/info/route.ts @@ -1,10 +1,10 @@ import { NextRequest } from 'next/server' import { eq } from 'drizzle-orm' import { createLogger } from '@/lib/logs/console-logger' +import { validateWorkflowAccess } from '@/app/api/workflows/middleware' +import { createErrorResponse, createSuccessResponse } from '@/app/api/workflows/utils' import { db } from '@/db' import * as schema from '@/db/schema' -import { validateWorkflowAccess } from '@/app/api/workflow/middleware' -import { createErrorResponse, createSuccessResponse } from '@/app/api/workflow/utils' const logger = createLogger('MarketplaceInfoAPI') @@ -27,7 +27,7 @@ export async function GET(request: NextRequest, { params }: { params: Promise<{ .from(schema.marketplace) .where(eq(schema.marketplace.workflowId, id)) .limit(1) - .then(rows => rows[0]) + .then((rows) => rows[0]) if (!marketplaceEntry) { logger.warn(`[${requestId}] No marketplace entry found for workflow: ${id}`) @@ -48,7 +48,10 @@ export async function GET(request: NextRequest, { params }: { params: Promise<{ updatedAt: marketplaceEntry.updatedAt, }) } catch (error) { - logger.error(`[${requestId}] Error getting marketplace info for workflow: ${(await params).id}`, error) + logger.error( + `[${requestId}] Error getting marketplace info for workflow: ${(await params).id}`, + error + ) return createErrorResponse('Failed to get marketplace information', 500) } -} \ No newline at end of file +} diff --git a/sim/app/api/marketplace/[id]/star/route.ts b/sim/app/api/marketplace/[id]/star/route.ts index 22b6a0f204..2c115e89fe 100644 --- a/sim/app/api/marketplace/[id]/star/route.ts +++ b/sim/app/api/marketplace/[id]/star/route.ts @@ -2,9 +2,9 @@ import { NextRequest } from 'next/server' import { and, eq } from 'drizzle-orm' import { getSession } from '@/lib/auth' import { createLogger } from '@/lib/logs/console-logger' +import { createErrorResponse, createSuccessResponse } from '@/app/api/workflows/utils' import { db } from '@/db' import * as schema from '@/db/schema' -import { createErrorResponse, createSuccessResponse } from '@/app/api/workflow/utils' const logger = createLogger('MarketplaceStarAPI') @@ -15,7 +15,7 @@ export async function POST(request: NextRequest, { params }: { params: Promise<{ const { id } = await params const session = await getSession() const userId = session?.user?.id - + if (!userId) { return createErrorResponse('Unauthorized', 401) } @@ -26,7 +26,7 @@ export async function POST(request: NextRequest, { params }: { params: Promise<{ .from(schema.marketplace) .where(eq(schema.marketplace.id, id)) .limit(1) - .then(rows => rows[0]) + .then((rows) => rows[0]) if (!marketplaceEntry) { logger.warn(`[${requestId}] No marketplace entry found with ID: ${id}`) @@ -37,28 +37,30 @@ export async function POST(request: NextRequest, { params }: { params: Promise<{ const existingStar = await db .select() .from(schema.marketplaceStar) - .where(and( - eq(schema.marketplaceStar.marketplaceId, id), - eq(schema.marketplaceStar.userId, userId) - )) + .where( + and(eq(schema.marketplaceStar.marketplaceId, id), eq(schema.marketplaceStar.userId, userId)) + ) .limit(1) - .then(rows => rows[0]) + .then((rows) => rows[0]) let action if (existingStar) { // User has already starred, so unstar it - await db.delete(schema.marketplaceStar).where( - and( - eq(schema.marketplaceStar.marketplaceId, id), - eq(schema.marketplaceStar.userId, userId) + await db + .delete(schema.marketplaceStar) + .where( + and( + eq(schema.marketplaceStar.marketplaceId, id), + eq(schema.marketplaceStar.userId, userId) + ) ) - ) - + // Decrement the star count - await db.update(schema.marketplace) + await db + .update(schema.marketplace) .set({ stars: marketplaceEntry.stars - 1 }) .where(eq(schema.marketplace.id, id)) - + action = 'unstarred' } else { // User hasn't starred yet, add a star @@ -68,12 +70,13 @@ export async function POST(request: NextRequest, { params }: { params: Promise<{ userId: userId, createdAt: new Date(), }) - + // Increment the star count - await db.update(schema.marketplace) + await db + .update(schema.marketplace) .set({ stars: marketplaceEntry.stars + 1 }) .where(eq(schema.marketplace.id, id)) - + action = 'starred' } @@ -107,12 +110,11 @@ export async function GET(request: NextRequest, { params }: { params: Promise<{ const existingStar = await db .select() .from(schema.marketplaceStar) - .where(and( - eq(schema.marketplaceStar.marketplaceId, id), - eq(schema.marketplaceStar.userId, userId) - )) + .where( + and(eq(schema.marketplaceStar.marketplaceId, id), eq(schema.marketplaceStar.userId, userId)) + ) .limit(1) - .then(rows => rows[0]) + .then((rows) => rows[0]) return createSuccessResponse({ isStarred: !!existingStar, @@ -121,4 +123,4 @@ export async function GET(request: NextRequest, { params }: { params: Promise<{ logger.error(`[${requestId}] Error checking star status: ${(await params).id}`, error) return createErrorResponse('Failed to check star status', 500) } -} \ No newline at end of file +} diff --git a/sim/app/api/marketplace/[id]/unpublish/route.ts b/sim/app/api/marketplace/[id]/unpublish/route.ts index 9f3bdff48c..f636238d53 100644 --- a/sim/app/api/marketplace/[id]/unpublish/route.ts +++ b/sim/app/api/marketplace/[id]/unpublish/route.ts @@ -1,10 +1,10 @@ import { NextRequest } from 'next/server' import { eq } from 'drizzle-orm' import { createLogger } from '@/lib/logs/console-logger' +import { validateWorkflowAccess } from '@/app/api/workflows/middleware' +import { createErrorResponse, createSuccessResponse } from '@/app/api/workflows/utils' import { db } from '@/db' import * as schema from '@/db/schema' -import { validateWorkflowAccess } from '@/app/api/workflow/middleware' -import { createErrorResponse, createSuccessResponse } from '@/app/api/workflow/utils' const logger = createLogger('MarketplaceUnpublishAPI') @@ -28,7 +28,7 @@ export async function POST(request: NextRequest, { params }: { params: Promise<{ .from(schema.marketplace) .where(eq(schema.marketplace.workflowId, id)) .limit(1) - .then(rows => rows[0]) + .then((rows) => rows[0]) if (!marketplaceEntry) { logger.warn(`[${requestId}] No marketplace entry found for workflow: ${id}`) @@ -51,4 +51,4 @@ export async function POST(request: NextRequest, { params }: { params: Promise<{ logger.error(`[${requestId}] Error unpublishing workflow: ${(await params).id}`, error) return createErrorResponse('Failed to unpublish workflow', 500) } -} \ No newline at end of file +} diff --git a/sim/app/api/scheduled/[id]/route.ts b/sim/app/api/schedules/[id]/route.ts similarity index 100% rename from sim/app/api/scheduled/[id]/route.ts rename to sim/app/api/schedules/[id]/route.ts diff --git a/sim/app/api/scheduled/execute/route.ts b/sim/app/api/schedules/execute/route.ts similarity index 97% rename from sim/app/api/scheduled/execute/route.ts rename to sim/app/api/schedules/execute/route.ts index 28a62f453c..689b0281ba 100644 --- a/sim/app/api/scheduled/execute/route.ts +++ b/sim/app/api/schedules/execute/route.ts @@ -5,6 +5,7 @@ import { v4 as uuidv4 } from 'uuid' import { z } from 'zod' import { createLogger } from '@/lib/logs/console-logger' import { persistExecutionError, persistExecutionLogs } from '@/lib/logs/execution-logger' +import { buildTraceSpans } from '@/lib/logs/trace-spans' import { decryptSecret } from '@/lib/utils' import { mergeSubblockState } from '@/stores/workflows/utils' import { BlockState, WorkflowState } from '@/stores/workflows/workflow/types' @@ -325,8 +326,18 @@ export async function GET(req: NextRequest) { ) const result = await executor.execute(schedule.workflowId) + // Build trace spans from execution logs + const { traceSpans, totalDuration } = buildTraceSpans(result) + + // Add trace spans to the execution result + const enrichedResult = { + ...result, + traceSpans, + totalDuration, + } + // Log each execution step and the final result - await persistExecutionLogs(schedule.workflowId, executionId, result, 'schedule') + await persistExecutionLogs(schedule.workflowId, executionId, enrichedResult, 'schedule') // Only update next_run_at if execution was successful if (result.success) { diff --git a/sim/app/api/scheduled/route.ts b/sim/app/api/schedules/route.ts similarity index 100% rename from sim/app/api/scheduled/route.ts rename to sim/app/api/schedules/route.ts diff --git a/sim/app/api/scheduled/schedule/route.ts b/sim/app/api/schedules/schedule/route.ts similarity index 100% rename from sim/app/api/scheduled/schedule/route.ts rename to sim/app/api/schedules/schedule/route.ts diff --git a/sim/app/api/webhooks/trigger/[path]/route.ts b/sim/app/api/webhooks/trigger/[path]/route.ts index c02f752391..a55c2f93df 100644 --- a/sim/app/api/webhooks/trigger/[path]/route.ts +++ b/sim/app/api/webhooks/trigger/[path]/route.ts @@ -3,6 +3,7 @@ import { and, eq } from 'drizzle-orm' import { v4 as uuidv4 } from 'uuid' import { createLogger } from '@/lib/logs/console-logger' import { persistExecutionError, persistExecutionLogs } from '@/lib/logs/execution-logger' +import { buildTraceSpans } from '@/lib/logs/trace-spans' import { closeRedisConnection, hasProcessedMessage, markMessageAsProcessed } from '@/lib/redis' import { decryptSecret } from '@/lib/utils' import { mergeSubblockStateAsync } from '@/stores/workflows/utils' @@ -558,8 +559,18 @@ async function processWebhook( executionTime: result.metadata?.duration, }) + // Build trace spans from execution logs + const { traceSpans, totalDuration } = buildTraceSpans(result) + + // Add trace spans to the execution result + const enrichedResult = { + ...result, + traceSpans, + totalDuration, + } + // Log each execution step and the final result - await persistExecutionLogs(foundWorkflow.id, executionId, result, 'webhook') + await persistExecutionLogs(foundWorkflow.id, executionId, enrichedResult, 'webhook') // Return the execution result return NextResponse.json(result, { status: 200 }) diff --git a/sim/app/api/workflow/[id]/deploy/route.ts b/sim/app/api/workflows/[id]/deploy/route.ts similarity index 100% rename from sim/app/api/workflow/[id]/deploy/route.ts rename to sim/app/api/workflows/[id]/deploy/route.ts diff --git a/sim/app/api/workflow/[id]/execute/route.ts b/sim/app/api/workflows/[id]/execute/route.ts similarity index 95% rename from sim/app/api/workflow/[id]/execute/route.ts rename to sim/app/api/workflows/[id]/execute/route.ts index dfe7f1b343..c18cb6cc1e 100644 --- a/sim/app/api/workflow/[id]/execute/route.ts +++ b/sim/app/api/workflows/[id]/execute/route.ts @@ -4,6 +4,7 @@ import { v4 as uuidv4 } from 'uuid' import { z } from 'zod' import { createLogger } from '@/lib/logs/console-logger' import { persistExecutionError, persistExecutionLogs } from '@/lib/logs/execution-logger' +import { buildTraceSpans } from '@/lib/logs/trace-spans' import { decryptSecret } from '@/lib/utils' import { mergeSubblockState } from '@/stores/workflows/utils' import { WorkflowState } from '@/stores/workflows/workflow/types' @@ -158,8 +159,18 @@ async function executeWorkflow(workflow: any, requestId: string, input?: any) { executionTime: result.metadata?.duration, }) + // Build trace spans from execution logs + const { traceSpans, totalDuration } = buildTraceSpans(result) + + // Add trace spans to the execution result + const enrichedResult = { + ...result, + traceSpans, + totalDuration, + } + // Log each execution step and the final result - await persistExecutionLogs(workflowId, executionId, result, 'api') + await persistExecutionLogs(workflowId, executionId, enrichedResult, 'api') return result } catch (error: any) { diff --git a/sim/app/api/workflow/[id]/log/route.ts b/sim/app/api/workflows/[id]/log/route.ts similarity index 100% rename from sim/app/api/workflow/[id]/log/route.ts rename to sim/app/api/workflows/[id]/log/route.ts diff --git a/sim/app/api/workflow/[id]/status/route.ts b/sim/app/api/workflows/[id]/status/route.ts similarity index 100% rename from sim/app/api/workflow/[id]/status/route.ts rename to sim/app/api/workflows/[id]/status/route.ts diff --git a/sim/app/api/workflow/middleware.ts b/sim/app/api/workflows/middleware.ts similarity index 100% rename from sim/app/api/workflow/middleware.ts rename to sim/app/api/workflows/middleware.ts diff --git a/sim/app/api/workflow/utils.ts b/sim/app/api/workflows/utils.ts similarity index 100% rename from sim/app/api/workflow/utils.ts rename to sim/app/api/workflows/utils.ts diff --git a/sim/app/w/[id]/components/control-bar/control-bar.tsx b/sim/app/w/[id]/components/control-bar/control-bar.tsx index 95e29ae5d8..65e708a78e 100644 --- a/sim/app/w/[id]/components/control-bar/control-bar.tsx +++ b/sim/app/w/[id]/components/control-bar/control-bar.tsx @@ -110,7 +110,7 @@ export function ControlBar() { } try { - const response = await fetch(`/api/workflow/${activeWorkflowId}/status`) + const response = await fetch(`/api/workflows/${activeWorkflowId}/status`) if (response.ok) { const data = await response.json() // Update the store with the status from the API @@ -198,11 +198,11 @@ export function ControlBar() { try { setIsDeploying(true) - const response = await fetch(`/api/workflow/${activeWorkflowId}/deploy/info`) + const response = await fetch(`/api/workflows/${activeWorkflowId}/deploy/info`) if (!response.ok) throw new Error('Failed to fetch deployment info') const { apiKey } = await response.json() - const endpoint = `${process.env.NEXT_PUBLIC_APP_URL}/api/workflow/${activeWorkflowId}/execute` + const endpoint = `${process.env.NEXT_PUBLIC_APP_URL}/api/workflows/${activeWorkflowId}/execute` // Create a new notification with the deployment info addNotification('api', 'Workflow deployment information', activeWorkflowId, { @@ -234,14 +234,14 @@ export function ControlBar() { try { setIsDeploying(true) - const response = await fetch(`/api/workflow/${activeWorkflowId}/deploy`, { + const response = await fetch(`/api/workflows/${activeWorkflowId}/deploy`, { method: 'POST', }) if (!response.ok) throw new Error('Failed to deploy workflow') const { apiKey, isDeployed: newDeployStatus, deployedAt } = await response.json() - const endpoint = `${process.env.NEXT_PUBLIC_APP_URL}/api/workflow/${activeWorkflowId}/execute` + const endpoint = `${process.env.NEXT_PUBLIC_APP_URL}/api/workflows/${activeWorkflowId}/execute` // Update the store with the deployment status setDeploymentStatus(newDeployStatus, deployedAt ? new Date(deployedAt) : undefined) diff --git a/sim/app/w/[id]/components/notifications/notifications.tsx b/sim/app/w/[id]/components/notifications/notifications.tsx index e01ced7bdb..15f66b1848 100644 --- a/sim/app/w/[id]/components/notifications/notifications.tsx +++ b/sim/app/w/[id]/components/notifications/notifications.tsx @@ -175,7 +175,7 @@ function NotificationAlert({ notification, isFading, onHide }: NotificationAlert if (!workflowId) return try { - const response = await fetch(`/api/workflow/${workflowId}/deploy`, { + const response = await fetch(`/api/workflows/${workflowId}/deploy`, { method: 'DELETE', }) diff --git a/sim/app/w/[id]/components/workflow-block/components/action-bar/schedule-status.tsx b/sim/app/w/[id]/components/workflow-block/components/action-bar/schedule-status.tsx index 8b005787a9..4d7760faa5 100644 --- a/sim/app/w/[id]/components/workflow-block/components/action-bar/schedule-status.tsx +++ b/sim/app/w/[id]/components/workflow-block/components/action-bar/schedule-status.tsx @@ -25,7 +25,7 @@ export function ScheduleStatus({ blockId }: ScheduleStatusProps) { setIsLoading(true) try { // Check if there's a schedule for this workflow - const response = await fetch(`/api/scheduled?workflowId=${workflowId}`) + const response = await fetch(`/api/schedules?workflowId=${workflowId}`) if (response.ok) { const data = await response.json() if (data.schedule) { diff --git a/sim/app/w/[id]/components/workflow-block/components/sub-block/components/schedule/schedule-config.tsx b/sim/app/w/[id]/components/workflow-block/components/sub-block/components/schedule/schedule-config.tsx index 3524d6ee1e..1467b39325 100644 --- a/sim/app/w/[id]/components/workflow-block/components/sub-block/components/schedule/schedule-config.tsx +++ b/sim/app/w/[id]/components/workflow-block/components/sub-block/components/schedule/schedule-config.tsx @@ -46,7 +46,7 @@ export function ScheduleConfig({ blockId, subBlockId, isConnecting }: ScheduleCo setIsLoading(true) try { // Check if there's a schedule for this workflow - const response = await fetch(`/api/scheduled?workflowId=${workflowId}`) + const response = await fetch(`/api/schedules?workflowId=${workflowId}`) if (response.ok) { const data = await response.json() if (data.schedule) { @@ -115,7 +115,7 @@ export function ScheduleConfig({ blockId, subBlockId, isConnecting }: ScheduleCo setIsSaving(true) try { // Send the complete workflow state to be saved/updated - const response = await fetch(`/api/scheduled/schedule`, { + const response = await fetch(`/api/schedules/schedule`, { method: 'POST', headers: { 'Content-Type': 'application/json', @@ -151,7 +151,7 @@ export function ScheduleConfig({ blockId, subBlockId, isConnecting }: ScheduleCo setIsDeleting(true) try { - const response = await fetch(`/api/scheduled/${scheduleId}`, { + const response = await fetch(`/api/schedules/${scheduleId}`, { method: 'DELETE', }) diff --git a/sim/app/w/[id]/hooks/use-workflow-execution.ts b/sim/app/w/[id]/hooks/use-workflow-execution.ts index 577df7e4ba..cc0909ce76 100644 --- a/sim/app/w/[id]/hooks/use-workflow-execution.ts +++ b/sim/app/w/[id]/hooks/use-workflow-execution.ts @@ -1,6 +1,7 @@ import { useCallback, useState } from 'react' import { v4 as uuidv4 } from 'uuid' import { createLogger } from '@/lib/logs/console-logger' +import { buildTraceSpans } from '@/lib/logs/trace-spans' import { useConsoleStore } from '@/stores/console/store' import { useExecutionStore } from '@/stores/execution/store' import { useNotificationStore } from '@/stores/notifications/store' @@ -8,250 +9,12 @@ import { useEnvironmentStore } from '@/stores/settings/environment/store' import { useWorkflowRegistry } from '@/stores/workflows/registry/store' import { mergeSubblockState } from '@/stores/workflows/utils' import { useWorkflowStore } from '@/stores/workflows/workflow/store' -import { TraceSpan } from '@/app/w/logs/stores/types' import { Executor } from '@/executor' import { ExecutionResult } from '@/executor/types' import { Serializer } from '@/serializer' const logger = createLogger('useWorkflowExecution') -// Helper function to build a tree of trace spans from execution logs -function buildTraceSpans(result: ExecutionResult): { - traceSpans: TraceSpan[] - totalDuration: number -} { - // If no logs, return empty spans - if (!result.logs || result.logs.length === 0) { - return { traceSpans: [], totalDuration: 0 } - } - - // Store all spans as a map for faster lookup - const spanMap = new Map() - - // First pass: Create spans for each block - result.logs.forEach((log) => { - // Skip logs that don't have block execution information - if (!log.blockId || !log.blockType) return - - // Create a unique ID for this span using blockId and timestamp - const spanId = `${log.blockId}-${new Date(log.startedAt).getTime()}` - - // Extract duration if available - const duration = log.durationMs || 0 - - // Create the span - const span: TraceSpan = { - id: spanId, - name: log.blockName || log.blockId, - type: log.blockType, - duration: duration, - startTime: log.startedAt, - endTime: log.endedAt, - status: log.error ? 'error' : 'success', - children: [], - } - - // Add provider timing data if it exists - if (log.output?.response?.providerTiming) { - const providerTiming = log.output.response.providerTiming - - // If we have time segments, use them to create a more detailed timeline - if (providerTiming.timeSegments && providerTiming.timeSegments.length > 0) { - const segmentStartTime = new Date(log.startedAt).getTime() - const children: TraceSpan[] = [] - - // Process segments in order - providerTiming.timeSegments.forEach( - ( - segment: { - type: string - name: string - startTime: number - endTime: number - duration: number - }, - index: number - ) => { - const relativeStart = segment.startTime - segmentStartTime - - // Enhance the segment name to include model information for model segments - let enhancedName = segment.name - if (segment.type === 'model') { - const modelName = log.output.response.model || '' - - if (segment.name === 'Initial response') { - enhancedName = `Initial response${modelName ? ` (${modelName})` : ''}` - } else if (segment.name.includes('iteration')) { - // Extract the iteration number - const iterationMatch = segment.name.match(/\(iteration (\d+)\)/) - const iterationNum = iterationMatch ? iterationMatch[1] : '' - - enhancedName = `Model response${iterationNum ? ` (iteration ${iterationNum})` : ''}${modelName ? ` (${modelName})` : ''}` - } - } - - const segmentSpan: TraceSpan = { - id: `${spanId}-segment-${index}`, - name: enhancedName, - // Make sure we handle model and tool types, and fallback to generic 'span' for anything else - type: segment.type === 'model' || segment.type === 'tool' ? segment.type : 'span', - duration: segment.duration, - startTime: new Date(segment.startTime).toISOString(), - endTime: new Date(segment.endTime).toISOString(), - status: 'success', - // Add relative timing display for segments after the first one - relativeStartMs: index === 0 ? undefined : relativeStart, - // For model segments, add token info if available - ...(segment.type === 'model' && { - tokens: index === 0 ? log.output.response.tokens?.completion : undefined, - }), - } - - children.push(segmentSpan) - } - ) - - // Add all segments as children - if (!span.children) span.children = [] - span.children.push(...children) - } - // If no segments but we have provider timing, create a provider span - else { - // Create a child span for the provider execution - const providerSpan: TraceSpan = { - id: `${spanId}-provider`, - name: log.output.response.model || 'AI Provider', - type: 'provider', - duration: providerTiming.duration || 0, - startTime: providerTiming.startTime || log.startedAt, - endTime: providerTiming.endTime || log.endedAt, - status: 'success', - tokens: log.output.response.tokens?.total, - } - - // If we have model time, create a child span for just the model processing - if (providerTiming.modelTime) { - const modelName = log.output.response.model || '' - const modelSpan: TraceSpan = { - id: `${spanId}-model`, - name: `Model Generation${modelName ? ` (${modelName})` : ''}`, - type: 'model', - duration: providerTiming.modelTime, - startTime: providerTiming.startTime, // Approximate - endTime: providerTiming.endTime, // Approximate - status: 'success', - tokens: log.output.response.tokens?.completion, - } - - if (!providerSpan.children) providerSpan.children = [] - providerSpan.children.push(modelSpan) - } - - if (!span.children) span.children = [] - span.children.push(providerSpan) - - // When using provider timing without segments, still add tool calls if they exist - if (log.output?.response?.toolCalls?.list) { - span.toolCalls = log.output.response.toolCalls.list.map((tc: any) => ({ - name: tc.name, - duration: tc.duration || 0, - startTime: tc.startTime || log.startedAt, - endTime: tc.endTime || log.endedAt, - status: tc.error ? 'error' : 'success', - input: tc.arguments || tc.input, - output: tc.result || tc.output, - error: tc.error, - })) - } - } - } else { - // When not using provider timing at all, add tool calls if they exist - if (log.output?.response?.toolCalls?.list) { - span.toolCalls = log.output.response.toolCalls.list.map((tc: any) => ({ - name: tc.name, - duration: tc.duration || 0, - startTime: tc.startTime || log.startedAt, - endTime: tc.endTime || log.endedAt, - status: tc.error ? 'error' : 'success', - input: tc.arguments || tc.input, - output: tc.result || tc.output, - error: tc.error, - })) - } - } - - // Store in map - spanMap.set(spanId, span) - }) - - // Second pass: Build the hierarchy - // We'll first need to sort logs chronologically - const sortedLogs = [...result.logs].sort((a, b) => { - const aTime = new Date(a.startedAt).getTime() - const bTime = new Date(b.startedAt).getTime() - return aTime - bTime - }) - - // Track parent spans using a stack - const spanStack: TraceSpan[] = [] - const rootSpans: TraceSpan[] = [] - - // Process logs to build the hierarchy - sortedLogs.forEach((log) => { - if (!log.blockId || !log.blockType) return - - const spanId = `${log.blockId}-${new Date(log.startedAt).getTime()}` - const span = spanMap.get(spanId) - if (!span) return - - // If we have a non-empty stack, check if this span should be a child - if (spanStack.length > 0) { - const potentialParent = spanStack[spanStack.length - 1] - const parentStartTime = new Date(potentialParent.startTime).getTime() - const parentEndTime = new Date(potentialParent.endTime).getTime() - const spanStartTime = new Date(span.startTime).getTime() - - // If this span starts after the parent starts and the parent is still on the stack, - // we'll assume it's a child span - if (spanStartTime >= parentStartTime && spanStartTime <= parentEndTime) { - if (!potentialParent.children) potentialParent.children = [] - potentialParent.children.push(span) - } else { - // This span doesn't belong to the current parent, pop from stack - while ( - spanStack.length > 0 && - new Date(spanStack[spanStack.length - 1].endTime).getTime() < spanStartTime - ) { - spanStack.pop() - } - - // Check if we still have a parent - if (spanStack.length > 0) { - const newParent = spanStack[spanStack.length - 1] - if (!newParent.children) newParent.children = [] - newParent.children.push(span) - } else { - // No parent, this is a root span - rootSpans.push(span) - } - } - } else { - // Empty stack, this is a root span - rootSpans.push(span) - } - - // Check if this span could be a parent to future spans - if (log.blockType === 'agent' || log.blockType === 'workflow') { - spanStack.push(span) - } - }) - - // Calculate total duration as the sum of root spans - const totalDuration = rootSpans.reduce((sum, span) => sum + span.duration, 0) - - return { traceSpans: rootSpans, totalDuration } -} - export function useWorkflowExecution() { const { blocks, edges, loops } = useWorkflowStore() const { activeWorkflowId } = useWorkflowRegistry() @@ -273,7 +36,7 @@ export function useWorkflowExecution() { totalDuration, } - const response = await fetch(`/api/workflow/${activeWorkflowId}/log`, { + const response = await fetch(`/api/workflows/${activeWorkflowId}/log`, { method: 'POST', headers: { 'Content-Type': 'application/json', diff --git a/sim/app/w/logs/logs.tsx b/sim/app/w/logs/logs.tsx index c850cad477..58ded0a743 100644 --- a/sim/app/w/logs/logs.tsx +++ b/sim/app/w/logs/logs.tsx @@ -26,9 +26,18 @@ const getLevelBadgeStyles = (level: string) => { // Helper function to get trigger badge styling const getTriggerBadgeStyles = (trigger: string) => { - return trigger.toLowerCase() === 'manual' - ? 'bg-secondary text-secondary-foreground' - : 'bg-blue-100 dark:bg-blue-950/40 text-blue-700 dark:text-blue-400' + switch (trigger.toLowerCase()) { + case 'manual': + return 'bg-secondary text-secondary-foreground' + case 'api': + return 'bg-blue-100 dark:bg-blue-950/40 text-blue-700 dark:text-blue-400' + case 'webhook': + return 'bg-purple-100 dark:bg-purple-950/40 text-purple-700 dark:text-purple-400' + case 'schedule': + return 'bg-green-100 dark:bg-green-950/40 text-green-700 dark:text-green-400' + default: + return 'bg-gray-100 dark:bg-gray-800 text-gray-700 dark:text-gray-400' + } } // Add a new CSS class for the selected row animation diff --git a/sim/lib/logs/trace-spans.ts b/sim/lib/logs/trace-spans.ts new file mode 100644 index 0000000000..3b9bf59e46 --- /dev/null +++ b/sim/lib/logs/trace-spans.ts @@ -0,0 +1,239 @@ +import { TraceSpan } from '@/app/w/logs/stores/types' +import { ExecutionResult } from '@/executor/types' + +// Helper function to build a tree of trace spans from execution logs +export function buildTraceSpans(result: ExecutionResult): { + traceSpans: TraceSpan[] + totalDuration: number +} { + // If no logs, return empty spans + if (!result.logs || result.logs.length === 0) { + return { traceSpans: [], totalDuration: 0 } + } + + // Store all spans as a map for faster lookup + const spanMap = new Map() + + // First pass: Create spans for each block + result.logs.forEach((log) => { + // Skip logs that don't have block execution information + if (!log.blockId || !log.blockType) return + + // Create a unique ID for this span using blockId and timestamp + const spanId = `${log.blockId}-${new Date(log.startedAt).getTime()}` + + // Extract duration if available + const duration = log.durationMs || 0 + + // Create the span + const span: TraceSpan = { + id: spanId, + name: log.blockName || log.blockId, + type: log.blockType, + duration: duration, + startTime: log.startedAt, + endTime: log.endedAt, + status: log.error ? 'error' : 'success', + children: [], + } + + // Add provider timing data if it exists + if (log.output?.response?.providerTiming) { + const providerTiming = log.output.response.providerTiming + + // If we have time segments, use them to create a more detailed timeline + if (providerTiming.timeSegments && providerTiming.timeSegments.length > 0) { + const segmentStartTime = new Date(log.startedAt).getTime() + const children: TraceSpan[] = [] + + // Process segments in order + providerTiming.timeSegments.forEach( + ( + segment: { + type: string + name: string + startTime: number + endTime: number + duration: number + }, + index: number + ) => { + const relativeStart = segment.startTime - segmentStartTime + + // Enhance the segment name to include model information for model segments + let enhancedName = segment.name + if (segment.type === 'model') { + const modelName = log.output.response.model || '' + + if (segment.name === 'Initial response') { + enhancedName = `Initial response${modelName ? ` (${modelName})` : ''}` + } else if (segment.name.includes('iteration')) { + // Extract the iteration number + const iterationMatch = segment.name.match(/\(iteration (\d+)\)/) + const iterationNum = iterationMatch ? iterationMatch[1] : '' + + enhancedName = `Model response${iterationNum ? ` (iteration ${iterationNum})` : ''}${modelName ? ` (${modelName})` : ''}` + } + } + + const segmentSpan: TraceSpan = { + id: `${spanId}-segment-${index}`, + name: enhancedName, + // Make sure we handle model and tool types, and fallback to generic 'span' for anything else + type: segment.type === 'model' || segment.type === 'tool' ? segment.type : 'span', + duration: segment.duration, + startTime: new Date(segment.startTime).toISOString(), + endTime: new Date(segment.endTime).toISOString(), + status: 'success', + // Add relative timing display for segments after the first one + relativeStartMs: index === 0 ? undefined : relativeStart, + // For model segments, add token info if available + ...(segment.type === 'model' && { + tokens: index === 0 ? log.output.response.tokens?.completion : undefined, + }), + } + + children.push(segmentSpan) + } + ) + + // Add all segments as children + if (!span.children) span.children = [] + span.children.push(...children) + } + // If no segments but we have provider timing, create a provider span + else { + // Create a child span for the provider execution + const providerSpan: TraceSpan = { + id: `${spanId}-provider`, + name: log.output.response.model || 'AI Provider', + type: 'provider', + duration: providerTiming.duration || 0, + startTime: providerTiming.startTime || log.startedAt, + endTime: providerTiming.endTime || log.endedAt, + status: 'success', + tokens: log.output.response.tokens?.total, + } + + // If we have model time, create a child span for just the model processing + if (providerTiming.modelTime) { + const modelName = log.output.response.model || '' + const modelSpan: TraceSpan = { + id: `${spanId}-model`, + name: `Model Generation${modelName ? ` (${modelName})` : ''}`, + type: 'model', + duration: providerTiming.modelTime, + startTime: providerTiming.startTime, // Approximate + endTime: providerTiming.endTime, // Approximate + status: 'success', + tokens: log.output.response.tokens?.completion, + } + + if (!providerSpan.children) providerSpan.children = [] + providerSpan.children.push(modelSpan) + } + + if (!span.children) span.children = [] + span.children.push(providerSpan) + + // When using provider timing without segments, still add tool calls if they exist + if (log.output?.response?.toolCalls?.list) { + span.toolCalls = log.output.response.toolCalls.list.map((tc: any) => ({ + name: tc.name, + duration: tc.duration || 0, + startTime: tc.startTime || log.startedAt, + endTime: tc.endTime || log.endedAt, + status: tc.error ? 'error' : 'success', + input: tc.arguments || tc.input, + output: tc.result || tc.output, + error: tc.error, + })) + } + } + } else { + // When not using provider timing at all, add tool calls if they exist + if (log.output?.response?.toolCalls?.list) { + span.toolCalls = log.output.response.toolCalls.list.map((tc: any) => ({ + name: tc.name, + duration: tc.duration || 0, + startTime: tc.startTime || log.startedAt, + endTime: tc.endTime || log.endedAt, + status: tc.error ? 'error' : 'success', + input: tc.arguments || tc.input, + output: tc.result || tc.output, + error: tc.error, + })) + } + } + + // Store in map + spanMap.set(spanId, span) + }) + + // Second pass: Build the hierarchy + // We'll first need to sort logs chronologically + const sortedLogs = [...result.logs].sort((a, b) => { + const aTime = new Date(a.startedAt).getTime() + const bTime = new Date(b.startedAt).getTime() + return aTime - bTime + }) + + // Track parent spans using a stack + const spanStack: TraceSpan[] = [] + const rootSpans: TraceSpan[] = [] + + // Process logs to build the hierarchy + sortedLogs.forEach((log) => { + if (!log.blockId || !log.blockType) return + + const spanId = `${log.blockId}-${new Date(log.startedAt).getTime()}` + const span = spanMap.get(spanId) + if (!span) return + + // If we have a non-empty stack, check if this span should be a child + if (spanStack.length > 0) { + const potentialParent = spanStack[spanStack.length - 1] + const parentStartTime = new Date(potentialParent.startTime).getTime() + const parentEndTime = new Date(potentialParent.endTime).getTime() + const spanStartTime = new Date(span.startTime).getTime() + + // If this span starts after the parent starts and the parent is still on the stack, + // we'll assume it's a child span + if (spanStartTime >= parentStartTime && spanStartTime <= parentEndTime) { + if (!potentialParent.children) potentialParent.children = [] + potentialParent.children.push(span) + } else { + // This span doesn't belong to the current parent, pop from stack + while ( + spanStack.length > 0 && + new Date(spanStack[spanStack.length - 1].endTime).getTime() < spanStartTime + ) { + spanStack.pop() + } + + // Check if we still have a parent + if (spanStack.length > 0) { + const newParent = spanStack[spanStack.length - 1] + if (!newParent.children) newParent.children = [] + newParent.children.push(span) + } else { + // No parent, this is a root span + rootSpans.push(span) + } + } + } else { + // Empty stack, this is a root span + rootSpans.push(span) + } + + // Check if this span could be a parent to future spans + if (log.blockType === 'agent' || log.blockType === 'workflow') { + spanStack.push(span) + } + }) + + // Calculate total duration as the sum of root spans + const totalDuration = rootSpans.reduce((sum, span) => sum + span.duration, 0) + + return { traceSpans: rootSpans, totalDuration } +} diff --git a/sim/stores/constants.ts b/sim/stores/constants.ts index fed163ce14..c56ee98f5d 100644 --- a/sim/stores/constants.ts +++ b/sim/stores/constants.ts @@ -7,7 +7,7 @@ export const STORAGE_KEYS = { export const API_ENDPOINTS = { WORKFLOW: '/api/db/workflow', ENVIRONMENT: '/api/db/environment', - SCHEDULE: '/api/scheduled/schedule', + SCHEDULE: '/api/schedules/schedule', } export const SYNC_INTERVALS = { diff --git a/sim/vercel.json b/sim/vercel.json index 588158865b..fe6ee80cfa 100644 --- a/sim/vercel.json +++ b/sim/vercel.json @@ -1,7 +1,7 @@ { "crons": [ { - "path": "/api/scheduled/execute", + "path": "/api/schedules/execute", "schedule": "*/1 * * * *" } ]