feat(api): added trace spans for asynchronous executions, modified api route naming convention (#154)

* improvement: added trace spans for asynchronous workflow executions, previously only appeared when workflow was executed manually

* improvement: updated api/workflow route naming convention

* improvement: updated api/schedules route naming convention

* improvement(logs/ui): add differentiator in UI for different trigger types
This commit is contained in:
Waleed Latif
2025-03-23 16:03:47 -07:00
committed by GitHub
parent c5d537f12d
commit 03e1699433
23 changed files with 340 additions and 291 deletions
+8 -5
View File
@@ -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)
}
}
}
+27 -25
View File
@@ -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)
}
}
}
@@ -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)
}
}
}
@@ -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) {
+12 -1
View File
@@ -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 })
@@ -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) {
@@ -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)
@@ -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',
})
@@ -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) {
@@ -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',
})
+2 -239
View File
@@ -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<string, TraceSpan>()
// 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',
+12 -3
View File
@@ -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
+239
View File
@@ -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<string, TraceSpan>()
// 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 }
}
+1 -1
View File
@@ -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 = {
+1 -1
View File
@@ -1,7 +1,7 @@
{
"crons": [
{
"path": "/api/scheduled/execute",
"path": "/api/schedules/execute",
"schedule": "*/1 * * * *"
}
]