feat(logger): added new logger for remaining 1/2 of api endpoints

This commit is contained in:
Waleed Latif
2025-03-12 00:47:29 -07:00
parent 3e02c81e2a
commit e87e230416
15 changed files with 405 additions and 97 deletions
+16 -1
View File
@@ -1,10 +1,13 @@
import { NextRequest, NextResponse } from 'next/server'
import { Script, createContext } from 'vm'
import { createLogger } from '@/lib/logs/console-logger'
// Explicitly export allowed methods
export const dynamic = 'force-dynamic' // Disable static optimization
export const runtime = 'nodejs' // Use Node.js runtime
const logger = createLogger('FunctionExecuteAPI')
/**
* Resolves environment variables and tags in code
* @param code - Code with variables
@@ -40,6 +43,7 @@ function resolveCodeVariables(
}
export async function POST(req: NextRequest) {
const requestId = crypto.randomUUID().slice(0, 8)
const startTime = Date.now()
let stdout = ''
@@ -48,6 +52,12 @@ export async function POST(req: NextRequest) {
const { code, params = {}, timeout = 3000, envVars = {} } = body
logger.debug(`[${requestId}] Executing function with params`, {
hasParams: Object.keys(params).length > 0,
timeout,
hasEnvVars: Object.keys(envVars).length > 0,
})
// Resolve variables in the code with workflow environment variables
const resolvedCode = resolveCodeVariables(code, params, envVars)
@@ -68,7 +78,7 @@ export async function POST(req: NextRequest) {
args
.map((arg) => (typeof arg === 'object' ? JSON.stringify(arg, null, 2) : String(arg)))
.join(' ') + '\n'
console.error('❌ Code Console Error:', errorMessage.trim())
logger.error(`[${requestId}] Code Console Error:`, errorMessage.trim())
stdout += 'ERROR: ' + errorMessage
},
},
@@ -91,6 +101,7 @@ export async function POST(req: NextRequest) {
})
const executionTime = Date.now() - startTime
logger.info(`[${requestId}] Function executed successfully`, { executionTime })
const response = {
success: true,
@@ -104,6 +115,10 @@ export async function POST(req: NextRequest) {
return NextResponse.json(response)
} catch (error: any) {
const executionTime = Date.now() - startTime
logger.error(`[${requestId}] Function execution failed`, {
error: error.message || 'Unknown error',
executionTime,
})
const errorResponse = {
success: false,
+19 -4
View File
@@ -1,8 +1,10 @@
import { NextRequest, NextResponse } from 'next/server'
import { Resend } from 'resend'
import { z } from 'zod'
import { createLogger } from '@/lib/logs/console-logger'
const resend = new Resend(process.env.RESEND_API_KEY)
const logger = createLogger('HelpAPI')
// Define schema for validation
const helpFormSchema = z.object({
@@ -13,6 +15,8 @@ const helpFormSchema = z.object({
})
export async function POST(req: NextRequest) {
const requestId = crypto.randomUUID().slice(0, 8)
try {
// Handle multipart form data
const formData = await req.formData()
@@ -23,6 +27,11 @@ export async function POST(req: NextRequest) {
const message = formData.get('message') as string
const type = formData.get('type') as string
logger.info(`[${requestId}] Processing help request`, {
type,
email: email.substring(0, 3) + '***', // Log partial email for privacy
})
// Validate the form data
const result = helpFormSchema.safeParse({
email,
@@ -32,6 +41,9 @@ export async function POST(req: NextRequest) {
})
if (!result.success) {
logger.warn(`[${requestId}] Invalid help request data`, {
errors: result.error.format(),
})
return NextResponse.json(
{ error: 'Invalid request data', details: result.error.format() },
{ status: 400 }
@@ -54,6 +66,8 @@ export async function POST(req: NextRequest) {
}
}
logger.debug(`[${requestId}] Help request includes ${images.length} images`)
// Prepare email content
let emailText = `
Type: ${type}
@@ -82,10 +96,12 @@ ${message}
})
if (error) {
console.error('Error sending email:', error)
logger.error(`[${requestId}] Error sending help request email`, error)
return NextResponse.json({ error: 'Failed to send email' }, { status: 500 })
}
logger.info(`[${requestId}] Help request email sent successfully`)
// Send confirmation email to the user
await resend.emails
.send({
@@ -108,8 +124,7 @@ The Sim Studio Team
replyTo: 'help@simstudio.ai',
})
.catch((err) => {
// Log but don't fail if confirmation email fails
console.warn('Failed to send confirmation email:', err)
logger.warn(`[${requestId}] Failed to send confirmation email`, err)
})
return NextResponse.json(
@@ -117,7 +132,7 @@ The Sim Studio Team
{ status: 200 }
)
} catch (error) {
console.error('Error processing help request:', error)
logger.error(`[${requestId}] Error processing help request`, error)
return NextResponse.json({ error: 'Internal server error' }, { status: 500 })
}
}
+33 -1
View File
@@ -1,18 +1,40 @@
import { NextResponse } from 'next/server'
import { createLogger } from '@/lib/logs/console-logger'
import { executeTool, getTool } from '@/tools'
import { validateToolRequest } from '@/tools/utils'
const logger = createLogger('ProxyAPI')
export async function POST(request: Request) {
const requestId = crypto.randomUUID().slice(0, 8)
try {
const { toolId, params } = await request.json()
logger.debug(`[${requestId}] Proxy request for tool`, {
toolId,
hasParams: !!params && Object.keys(params).length > 0,
})
const tool = getTool(toolId)
// Validate the tool and its parameters
validateToolRequest(toolId, tool, params)
try {
validateToolRequest(toolId, tool, params)
} catch (error) {
logger.warn(`[${requestId}] Tool validation failed`, {
toolId,
error: error instanceof Error ? error.message : String(error),
})
return NextResponse.json({
success: false,
error: error instanceof Error ? error.message : String(error),
})
}
try {
if (!tool) {
logger.error(`[${requestId}] Tool not found`, { toolId })
throw new Error(`Tool not found: ${toolId}`)
}
@@ -20,6 +42,11 @@ export async function POST(request: Request) {
const result = await executeTool(toolId, params, true, true)
if (!result.success) {
logger.warn(`[${requestId}] Tool execution failed`, {
toolId,
error: result.error || 'Unknown error',
})
if (tool.transformError) {
try {
const errorResult = tool.transformError(result)
@@ -54,11 +81,16 @@ export async function POST(request: Request) {
}
}
logger.info(`[${requestId}] Tool executed successfully`, { toolId })
return NextResponse.json(result)
} catch (error: any) {
throw error
}
} catch (error: any) {
logger.error(`[${requestId}] Proxy request failed`, {
error: error instanceof Error ? error.message : String(error),
})
return NextResponse.json({
success: false,
error: error instanceof Error ? error.message : String(error),
+38 -13
View File
@@ -3,6 +3,7 @@ import { Cron } from 'croner'
import { eq, lte } from 'drizzle-orm'
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 { decryptSecret } from '@/lib/utils'
import { mergeSubblockState } from '@/stores/workflows/utils'
@@ -12,6 +13,8 @@ import { environment, workflow, workflowSchedule } from '@/db/schema'
import { Executor } from '@/executor'
import { Serializer } from '@/serializer'
const logger = createLogger('ScheduledExecuteAPI')
interface SubBlockValue {
value: string
}
@@ -151,6 +154,7 @@ const runningExecutions = new Set<string>()
// Add GET handler for cron job
export async function GET(req: NextRequest) {
const requestId = crypto.randomUUID().slice(0, 8)
const now = new Date()
let dueSchedules: (typeof workflowSchedule.$inferSelect)[] = []
@@ -164,16 +168,20 @@ export async function GET(req: NextRequest) {
// Limit to 10 workflows per minute to prevent overload
.limit(10)
logger.info(`[${requestId}] Processing ${dueSchedules.length} due scheduled workflows`)
for (const schedule of dueSchedules) {
const executionId = uuidv4()
try {
// Skip if this workflow is already running
if (runningExecutions.has(schedule.workflowId)) {
logger.debug(`[${requestId}] Skipping workflow ${schedule.workflowId} - already running`)
continue
}
runningExecutions.add(schedule.workflowId)
logger.debug(`[${requestId}] Starting execution of workflow ${schedule.workflowId}`)
// Retrieve the workflow record
const [workflowRecord] = await db
@@ -183,6 +191,7 @@ export async function GET(req: NextRequest) {
.limit(1)
if (!workflowRecord) {
logger.warn(`[${requestId}] Workflow ${schedule.workflowId} not found`)
runningExecutions.delete(schedule.workflowId)
continue
}
@@ -202,6 +211,9 @@ export async function GET(req: NextRequest) {
.limit(1)
if (!userEnv) {
logger.error(
`[${requestId}] No environment variables found for user ${workflowRecord.userId}`
)
throw new Error('No environment variables found for this user')
}
@@ -233,7 +245,10 @@ export async function GET(req: NextRequest) {
const { decrypted } = await decryptSecret(encryptedValue)
value = (value as string).replace(match, decrypted)
} catch (error: any) {
console.error('Error decrypting value:', error)
logger.error(
`[${requestId}] Error decrypting value for variable "${varName}"`,
error
)
throw new Error(
`Failed to decrypt environment variable "${varName}": ${error.message}`
)
@@ -259,7 +274,7 @@ export async function GET(req: NextRequest) {
const { decrypted } = await decryptSecret(encryptedValue)
decryptedEnvVars[key] = decrypted
} catch (error: any) {
console.error(`Failed to decrypt ${key}:`, error)
logger.error(`[${requestId}] Failed to decrypt environment variable "${key}"`, error)
throw new Error(`Failed to decrypt environment variable "${key}": ${error.message}`)
}
}
@@ -279,23 +294,19 @@ export async function GET(req: NextRequest) {
// Check if this block has a responseFormat that needs to be parsed
if (blockState.responseFormat && typeof blockState.responseFormat === 'string') {
try {
console.log(
`[Schedule Debug] Block ${blockId} has responseFormat as string:`,
blockState.responseFormat
)
logger.debug(`[${requestId}] Parsing responseFormat for block ${blockId}`)
// Attempt to parse the responseFormat if it's a string
const parsedResponseFormat = JSON.parse(blockState.responseFormat)
console.log(
`[Schedule Debug] Successfully parsed responseFormat for block ${blockId}:`,
parsedResponseFormat
)
acc[blockId] = {
...blockState,
responseFormat: parsedResponseFormat,
}
} catch (error) {
console.warn(`Failed to parse responseFormat for block ${blockId}:`, error)
logger.warn(
`[${requestId}] Failed to parse responseFormat for block ${blockId}`,
error
)
acc[blockId] = blockState
}
} else {
@@ -306,6 +317,7 @@ export async function GET(req: NextRequest) {
{} as Record<string, Record<string, any>>
)
logger.info(`[${requestId}] Executing workflow ${schedule.workflowId}`)
const executor = new Executor(
serializedWorkflow,
processedBlockStates, // Use the processed block states
@@ -319,6 +331,7 @@ export async function GET(req: NextRequest) {
// Only update next_run_at if execution was successful
if (result.success) {
logger.info(`[${requestId}] Workflow ${schedule.workflowId} executed successfully`)
// Calculate the next run time based on the schedule configuration
const nextRunAt = calculateNextRunTime(schedule, blocks)
@@ -331,7 +344,12 @@ export async function GET(req: NextRequest) {
nextRunAt,
})
.where(eq(workflowSchedule.id, schedule.id))
logger.debug(
`[${requestId}] Updated next run time for workflow ${schedule.workflowId} to ${nextRunAt.toISOString()}`
)
} else {
logger.warn(`[${requestId}] Workflow ${schedule.workflowId} execution failed`)
// If execution failed, increment next_run_at by a small delay to prevent immediate retries
const retryDelay = 1 * 60 * 1000 // 1 minute delay
const nextRetryAt = new Date(now.getTime() + retryDelay)
@@ -343,9 +361,16 @@ export async function GET(req: NextRequest) {
nextRunAt: nextRetryAt,
})
.where(eq(workflowSchedule.id, schedule.id))
logger.debug(
`[${requestId}] Scheduled retry for workflow ${schedule.workflowId} at ${nextRetryAt.toISOString()}`
)
}
} catch (error: any) {
console.error(`Error executing scheduled workflow ${schedule.workflowId}:`, error)
logger.error(
`[${requestId}] Error executing scheduled workflow ${schedule.workflowId}`,
error
)
// Log the error
await persistExecutionError(schedule.workflowId, executionId, error, 'schedule')
@@ -366,7 +391,7 @@ export async function GET(req: NextRequest) {
}
}
} catch (error: any) {
console.error('Error in scheduled execution:', error)
logger.error(`[${requestId}] Error in scheduled execution handler`, error)
return NextResponse.json({ error: error.message }, { status: 500 })
}
+22 -2
View File
@@ -2,9 +2,12 @@ import { NextRequest, NextResponse } from 'next/server'
import { eq } from 'drizzle-orm'
import { z } from 'zod'
import { getSession } from '@/lib/auth'
import { createLogger } from '@/lib/logs/console-logger'
import { BlockState } from '@/stores/workflows/workflow/types'
import { db } from '@/db'
import { workflow, workflowSchedule } from '@/db/schema'
import { workflowSchedule } from '@/db/schema'
const logger = createLogger('ScheduledScheduleAPI')
interface SubBlockValue {
value: string
@@ -26,21 +29,27 @@ const ScheduleRequestSchema = z.object({
})
export async function POST(req: NextRequest) {
const requestId = crypto.randomUUID().slice(0, 8)
try {
const session = await getSession()
if (!session?.user?.id) {
logger.warn(`[${requestId}] Unauthorized schedule update attempt`)
return NextResponse.json({ error: 'Unauthorized' }, { status: 401 })
}
const body = await req.json()
const { workflowId, state } = ScheduleRequestSchema.parse(body)
logger.info(`[${requestId}] Processing schedule update for workflow ${workflowId}`)
// Find the starter block to check if it's configured for scheduling
const starterBlock = Object.values(state.blocks).find(
(block: any) => block.type === 'starter'
) as BlockState | undefined
if (!starterBlock) {
logger.warn(`[${requestId}] No starter block found in workflow ${workflowId}`)
return NextResponse.json({ error: 'No starter block found in workflow' }, { status: 400 })
}
@@ -48,6 +57,7 @@ export async function POST(req: NextRequest) {
// If the workflow is not scheduled, delete any existing schedule
if (startWorkflow !== 'schedule') {
logger.info(`[${requestId}] Removing schedule for workflow ${workflowId}`)
await db.delete(workflowSchedule).where(eq(workflowSchedule.workflowId, workflowId))
return NextResponse.json({ message: 'Schedule removed' })
@@ -55,6 +65,7 @@ export async function POST(req: NextRequest) {
// Get schedule configuration from starter block
const scheduleType = getSubBlockValue(starterBlock, 'scheduleType')
logger.debug(`[${requestId}] Schedule type for workflow ${workflowId}: ${scheduleType}`)
// Calculate cron expression based on schedule type
let cronExpression: string | null = null
@@ -181,6 +192,7 @@ export async function POST(req: NextRequest) {
break
}
default:
logger.warn(`[${requestId}] Invalid schedule type: ${scheduleType}`)
return NextResponse.json({ error: 'Invalid schedule type' }, { status: 400 })
}
@@ -219,13 +231,21 @@ export async function POST(req: NextRequest) {
set: setValues,
})
logger.info(`[${requestId}] Schedule updated for workflow ${workflowId}`, {
nextRunAt: shouldUpdateNextRunAt
? nextRunAt?.toISOString()
: existingSchedule[0]?.nextRunAt?.toISOString(),
cronExpression,
})
return NextResponse.json({
message: 'Schedule updated',
nextRunAt: shouldUpdateNextRunAt ? nextRunAt : existingSchedule[0]?.nextRunAt,
cronExpression,
})
} catch (error) {
console.error('Error updating workflow schedule:', error)
logger.error(`[${requestId}] Error updating workflow schedule`, error)
if (error instanceof z.ZodError) {
return NextResponse.json(
{ error: 'Invalid request data', details: error.errors },
+6 -1
View File
@@ -2,14 +2,19 @@ import { NextResponse } from 'next/server'
import { eq } from 'drizzle-orm'
import { nanoid } from 'nanoid'
import { z } from 'zod'
import { createLogger } from '@/lib/logs/console-logger'
import { db } from '@/db'
import { waitlist } from '@/db/schema'
const logger = createLogger('WaitlistAPI')
const waitlistSchema = z.object({
email: z.string().email(),
})
export async function POST(request: Request) {
const requestId = crypto.randomUUID().slice(0, 8)
try {
const body = await request.json()
const { email } = waitlistSchema.parse(body)
@@ -35,7 +40,7 @@ export async function POST(request: Request) {
return NextResponse.json({ message: 'Successfully joined waitlist' }, { status: 200 })
} catch (error) {
console.error('Waitlist error:', error)
logger.error(`[${requestId}] Waitlist error`, error)
return NextResponse.json({ message: 'Failed to join waitlist' }, { status: 500 })
}
}
+33 -3
View File
@@ -1,18 +1,25 @@
import { NextRequest, NextResponse } from 'next/server'
import { and, eq } from 'drizzle-orm'
import { getSession } from '@/lib/auth'
import { createLogger } from '@/lib/logs/console-logger'
import { db } from '@/db'
import { webhook, workflow } from '@/db/schema'
const logger = createLogger('WebhookAPI')
export const dynamic = 'force-dynamic'
// Get a specific webhook
export async function GET(request: NextRequest, { params }: { params: Promise<{ id: string }> }) {
const requestId = crypto.randomUUID().slice(0, 8)
try {
const { id } = await params
logger.debug(`[${requestId}] Fetching webhook with ID: ${id}`)
const session = await getSession()
if (!session?.user?.id) {
logger.warn(`[${requestId}] Unauthorized webhook access attempt`)
return NextResponse.json({ error: 'Unauthorized' }, { status: 401 })
}
@@ -30,23 +37,29 @@ export async function GET(request: NextRequest, { params }: { params: Promise<{
.limit(1)
if (webhooks.length === 0) {
logger.warn(`[${requestId}] Webhook not found: ${id}`)
return NextResponse.json({ error: 'Webhook not found' }, { status: 404 })
}
logger.info(`[${requestId}] Successfully retrieved webhook: ${id}`)
return NextResponse.json({ webhook: webhooks[0] }, { status: 200 })
} catch (error) {
console.error('Error fetching webhook:', error)
logger.error(`[${requestId}] Error fetching webhook`, error)
return NextResponse.json({ error: 'Internal server error' }, { status: 500 })
}
}
// Update a webhook
export async function PATCH(request: NextRequest, { params }: { params: Promise<{ id: string }> }) {
const requestId = crypto.randomUUID().slice(0, 8)
try {
const { id } = await params
logger.debug(`[${requestId}] Updating webhook with ID: ${id}`)
const session = await getSession()
if (!session?.user?.id) {
logger.warn(`[${requestId}] Unauthorized webhook update attempt`)
return NextResponse.json({ error: 'Unauthorized' }, { status: 401 })
}
@@ -68,13 +81,22 @@ export async function PATCH(request: NextRequest, { params }: { params: Promise<
.limit(1)
if (webhooks.length === 0) {
logger.warn(`[${requestId}] Webhook not found: ${id}`)
return NextResponse.json({ error: 'Webhook not found' }, { status: 404 })
}
if (webhooks[0].workflow.userId !== session.user.id) {
logger.warn(`[${requestId}] Unauthorized webhook update attempt for webhook: ${id}`)
return NextResponse.json({ error: 'Unauthorized' }, { status: 403 })
}
logger.debug(`[${requestId}] Updating webhook properties`, {
hasPathUpdate: path !== undefined,
hasProviderUpdate: provider !== undefined,
hasConfigUpdate: providerConfig !== undefined,
hasActiveUpdate: isActive !== undefined,
})
// Update the webhook
const updatedWebhook = await db
.update(webhook)
@@ -89,9 +111,10 @@ export async function PATCH(request: NextRequest, { params }: { params: Promise<
.where(eq(webhook.id, id))
.returning()
logger.info(`[${requestId}] Successfully updated webhook: ${id}`)
return NextResponse.json({ webhook: updatedWebhook[0] }, { status: 200 })
} catch (error) {
console.error('Error updating webhook:', error)
logger.error(`[${requestId}] Error updating webhook`, error)
return NextResponse.json({ error: 'Internal server error' }, { status: 500 })
}
}
@@ -101,11 +124,15 @@ export async function DELETE(
request: NextRequest,
{ params }: { params: Promise<{ id: string }> }
) {
const requestId = crypto.randomUUID().slice(0, 8)
try {
const { id } = await params
logger.debug(`[${requestId}] Deleting webhook with ID: ${id}`)
const session = await getSession()
if (!session?.user?.id) {
logger.warn(`[${requestId}] Unauthorized webhook deletion attempt`)
return NextResponse.json({ error: 'Unauthorized' }, { status: 401 })
}
@@ -124,19 +151,22 @@ export async function DELETE(
.limit(1)
if (webhooks.length === 0) {
logger.warn(`[${requestId}] Webhook not found: ${id}`)
return NextResponse.json({ error: 'Webhook not found' }, { status: 404 })
}
if (webhooks[0].workflow.userId !== session.user.id) {
logger.warn(`[${requestId}] Unauthorized webhook deletion attempt for webhook: ${id}`)
return NextResponse.json({ error: 'Unauthorized' }, { status: 403 })
}
// Delete the webhook
await db.delete(webhook).where(eq(webhook.id, id))
logger.info(`[${requestId}] Successfully deleted webhook: ${id}`)
return NextResponse.json({ success: true }, { status: 200 })
} catch (error) {
console.error('Error deleting webhook:', error)
logger.error(`[${requestId}] Error deleting webhook`, error)
return NextResponse.json({ error: 'Internal server error' }, { status: 500 })
}
}
+40 -3
View File
@@ -2,16 +2,22 @@ import { NextRequest, NextResponse } from 'next/server'
import { and, eq } from 'drizzle-orm'
import { nanoid } from 'nanoid'
import { getSession } from '@/lib/auth'
import { createLogger } from '@/lib/logs/console-logger'
import { db } from '@/db'
import { webhook, workflow } from '@/db/schema'
const logger = createLogger('WebhooksAPI')
export const dynamic = 'force-dynamic'
// Get all webhooks for the current user
export async function GET(request: NextRequest) {
const requestId = crypto.randomUUID().slice(0, 8)
try {
const session = await getSession()
if (!session?.user?.id) {
logger.warn(`[${requestId}] Unauthorized webhooks access attempt`)
return NextResponse.json({ error: 'Unauthorized' }, { status: 401 })
}
@@ -19,6 +25,10 @@ export async function GET(request: NextRequest) {
const { searchParams } = new URL(request.url)
const workflowId = searchParams.get('workflowId')
logger.debug(`[${requestId}] Fetching webhooks for user ${session.user.id}`, {
filteredByWorkflow: !!workflowId,
})
// Create where condition
const whereCondition = workflowId
? and(eq(workflow.userId, session.user.id), eq(webhook.workflowId, workflowId))
@@ -36,18 +46,22 @@ export async function GET(request: NextRequest) {
.innerJoin(workflow, eq(webhook.workflowId, workflow.id))
.where(whereCondition)
logger.info(`[${requestId}] Retrieved ${webhooks.length} webhooks for user ${session.user.id}`)
return NextResponse.json({ webhooks }, { status: 200 })
} catch (error) {
console.error('Error fetching webhooks:', error)
logger.error(`[${requestId}] Error fetching webhooks`, error)
return NextResponse.json({ error: 'Internal server error' }, { status: 500 })
}
}
// Create a new webhook
export async function POST(request: NextRequest) {
const requestId = crypto.randomUUID().slice(0, 8)
try {
const session = await getSession()
if (!session?.user?.id) {
logger.warn(`[${requestId}] Unauthorized webhook creation attempt`)
return NextResponse.json({ error: 'Unauthorized' }, { status: 401 })
}
@@ -56,9 +70,18 @@ export async function POST(request: NextRequest) {
// Validate input
if (!workflowId || !path) {
logger.warn(`[${requestId}] Missing required fields for webhook creation`, {
hasWorkflowId: !!workflowId,
hasPath: !!path,
})
return NextResponse.json({ error: 'Missing required fields' }, { status: 400 })
}
logger.debug(`[${requestId}] Creating webhook for workflow ${workflowId}`, {
path,
provider: provider || 'generic',
})
// Check if the workflow belongs to the user
const workflows = await db
.select()
@@ -67,6 +90,7 @@ export async function POST(request: NextRequest) {
.limit(1)
if (workflows.length === 0) {
logger.warn(`[${requestId}] Workflow not found or not owned by user: ${workflowId}`)
return NextResponse.json({ error: 'Workflow not found' }, { status: 404 })
}
@@ -75,6 +99,10 @@ export async function POST(request: NextRequest) {
// If a webhook with the same path exists but belongs to a different workflow, return an error
if (existingWebhooks.length > 0 && existingWebhooks[0].workflowId !== workflowId) {
logger.warn(`[${requestId}] Webhook path conflict: ${path}`, {
existingWorkflowId: existingWebhooks[0].workflowId,
requestedWorkflowId: workflowId,
})
return NextResponse.json(
{ error: 'Webhook path already exists. Please use a different path.', code: 'PATH_EXISTS' },
{ status: 409 }
@@ -83,6 +111,8 @@ export async function POST(request: NextRequest) {
// If a webhook with the same path and workflowId exists, update it
if (existingWebhooks.length > 0 && existingWebhooks[0].workflowId === workflowId) {
logger.info(`[${requestId}] Updating existing webhook for path: ${path}`)
const updatedWebhook = await db
.update(webhook)
.set({
@@ -98,10 +128,17 @@ export async function POST(request: NextRequest) {
}
// Create a new webhook
const webhookId = nanoid()
logger.info(`[${requestId}] Creating new webhook with ID: ${webhookId}`, {
path,
workflowId,
provider: provider || 'generic',
})
const newWebhook = await db
.insert(webhook)
.values({
id: nanoid(),
id: webhookId,
workflowId,
path,
provider,
@@ -114,7 +151,7 @@ export async function POST(request: NextRequest) {
return NextResponse.json({ webhook: newWebhook[0] }, { status: 201 })
} catch (error) {
console.error('Error creating webhook:', error)
logger.error(`[${requestId}] Error creating webhook`, error)
return NextResponse.json({ error: 'Internal server error' }, { status: 500 })
}
}
+36 -2
View File
@@ -1,24 +1,33 @@
import { NextRequest, NextResponse } from 'next/server'
import { and, eq } from 'drizzle-orm'
import { eq } from 'drizzle-orm'
import { createLogger } from '@/lib/logs/console-logger'
import { db } from '@/db'
import { webhook } from '@/db/schema'
const logger = createLogger('WebhookTestAPI')
export const dynamic = 'force-dynamic'
export async function GET(request: NextRequest) {
const requestId = crypto.randomUUID().slice(0, 8)
try {
// Get the webhook ID and provider from the query parameters
const { searchParams } = new URL(request.url)
const webhookId = searchParams.get('id')
if (!webhookId) {
logger.warn(`[${requestId}] Missing webhook ID in test request`)
return NextResponse.json({ success: false, error: 'Webhook ID is required' }, { status: 400 })
}
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) {
logger.warn(`[${requestId}] Webhook not found: ${webhookId}`)
return NextResponse.json({ success: false, error: 'Webhook not found' }, { status: 404 })
}
@@ -30,12 +39,19 @@ export async function GET(request: NextRequest) {
const baseUrl = new URL(request.url).origin
const webhookUrl = `${baseUrl}/api/webhooks/trigger/${foundWebhook.path}`
logger.info(`[${requestId}] Testing webhook for provider: ${provider}`, {
webhookId,
path: foundWebhook.path,
isActive: foundWebhook.isActive,
})
// Provider-specific test logic
switch (provider) {
case 'whatsapp': {
const verificationToken = providerConfig.verificationToken
if (!verificationToken) {
logger.warn(`[${requestId}] WhatsApp webhook missing verification token: ${webhookId}`)
return NextResponse.json(
{ success: false, error: 'Webhook has no verification token' },
{ status: 400 }
@@ -48,6 +64,11 @@ export async function GET(request: NextRequest) {
// 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`, {
webhookId,
challenge,
})
// Make a request to the webhook endpoint
const response = await fetch(whatsappUrl, {
headers: {
@@ -63,6 +84,16 @@ export async function GET(request: NextRequest) {
// Check if the test was successful
const success = status === 200 && responseText === challenge
if (success) {
logger.info(`[${requestId}] WhatsApp webhook verification successful: ${webhookId}`)
} else {
logger.warn(`[${requestId}] WhatsApp webhook verification failed: ${webhookId}`, {
status,
contentType,
responseTextLength: responseText.length,
})
}
return NextResponse.json({
success,
webhook: {
@@ -99,6 +130,7 @@ export async function GET(request: NextRequest) {
case 'github': {
const contentType = providerConfig.contentType || 'application/json'
logger.info(`[${requestId}] GitHub webhook test successful: ${webhookId}`)
return NextResponse.json({
success: true,
webhook: {
@@ -118,6 +150,7 @@ export async function GET(request: NextRequest) {
}
case 'stripe': {
logger.info(`[${requestId}] Stripe webhook test successful: ${webhookId}`)
return NextResponse.json({
success: true,
webhook: {
@@ -139,6 +172,7 @@ export async function GET(request: NextRequest) {
default: {
// Generic webhook test
logger.info(`[${requestId}] Generic webhook test successful: ${webhookId}`)
return NextResponse.json({
success: true,
webhook: {
@@ -153,7 +187,7 @@ export async function GET(request: NextRequest) {
}
}
} catch (error: any) {
console.error('Error testing webhook:', error)
logger.error(`[${requestId}] Error testing webhook`, error)
return NextResponse.json(
{
success: false,
+64 -41
View File
@@ -1,15 +1,18 @@
import { NextRequest, NextResponse } from 'next/server'
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 { closeRedisConnection, hasProcessedMessage, markMessageAsProcessed } from '@/lib/redis'
import { decryptSecret } from '@/lib/utils'
import { mergeSubblockState, mergeSubblockStateAsync } from '@/stores/workflows/utils'
import { mergeSubblockStateAsync } from '@/stores/workflows/utils'
import { db } from '@/db'
import { environment, webhook, workflow } from '@/db/schema'
import { Executor } from '@/executor'
import { Serializer } from '@/serializer'
const logger = createLogger('WebhookTriggerAPI')
// Force dynamic rendering for webhook endpoints
export const dynamic = 'force-dynamic'
// Increase the response size limit for webhook payloads
@@ -20,6 +23,8 @@ export const maxDuration = 300 // 5 minutes max execution time for long-running
* Handles both WhatsApp verification and other webhook providers
*/
export async function GET(request: NextRequest, { params }: { params: Promise<{ path: string }> }) {
const requestId = crypto.randomUUID().slice(0, 8)
try {
const path = (await params).path
const url = new URL(request.url)
@@ -31,10 +36,10 @@ export async function GET(request: NextRequest, { params }: { params: Promise<{
if (mode && token && challenge) {
// This is a WhatsApp verification request
console.log('WhatsApp verification request received')
logger.info(`[${requestId}] WhatsApp verification request received for path: ${path}`)
if (mode !== 'subscribe') {
console.log('Invalid mode:', mode)
logger.warn(`[${requestId}] Invalid WhatsApp verification mode: ${mode}`)
return new NextResponse('Invalid mode', { status: 400 })
}
@@ -50,14 +55,12 @@ export async function GET(request: NextRequest, { params }: { params: Promise<{
const verificationToken = providerConfig.verificationToken
if (!verificationToken) {
console.log(`Webhook ${wh.id} has no verification token, skipping`)
logger.debug(`[${requestId}] Webhook ${wh.id} has no verification token, skipping`)
continue
}
if (token === verificationToken) {
console.log(
`Verification successful for webhook ${wh.id}, returning challenge: ${challenge}`
)
logger.info(`[${requestId}] WhatsApp verification successful for webhook ${wh.id}`)
// Return ONLY the challenge as plain text (exactly as WhatsApp expects)
return new NextResponse(challenge, {
status: 200,
@@ -68,12 +71,12 @@ export async function GET(request: NextRequest, { params }: { params: Promise<{
}
}
console.log('No matching verification token found')
logger.warn(`[${requestId}] No matching WhatsApp verification token found`)
return new NextResponse('Verification failed', { status: 403 })
}
// For non-WhatsApp verification requests
console.log('Looking for webhook with path:', path)
logger.debug(`[${requestId}] Looking for webhook with path: ${path}`)
// Find the webhook in the database
const webhooks = await db
@@ -85,13 +88,15 @@ export async function GET(request: NextRequest, { params }: { params: Promise<{
.limit(1)
if (webhooks.length === 0) {
logger.warn(`[${requestId}] No active webhook found for path: ${path}`)
return new NextResponse('Webhook not found', { status: 404 })
}
// For other providers, just return a 200 OK
logger.info(`[${requestId}] Webhook verification successful for path: ${path}`)
return new NextResponse('OK', { status: 200 })
} catch (error: any) {
console.error('Error processing webhook verification:', error)
logger.error(`[${requestId}] Error processing webhook verification`, error)
return new NextResponse(`Internal Server Error: ${error.message}`, { status: 500 })
} finally {
// Ensure Redis connection is properly closed in serverless environment
@@ -103,6 +108,7 @@ export async function POST(
request: NextRequest,
{ params }: { params: Promise<{ path: string }> }
) {
const requestId = crypto.randomUUID().slice(0, 8)
const executionId = uuidv4()
let foundWorkflow: any = null
@@ -111,16 +117,14 @@ export async function POST(
// Parse the request body
const body = await request.json().catch(() => ({}))
console.log(`Webhook POST request received for path: ${path}`)
logger.info(`[${requestId}] Webhook POST request received for path: ${path}`)
// Generate a unique request ID based on the request content
const requestHash = await generateRequestHash(path, body)
// Check if this exact request has been processed before
if (await hasProcessedMessage(requestHash)) {
console.log(
`Duplicate webhook request detected with hash: ${requestHash}. Skipping processing.`
)
logger.info(`[${requestId}] Duplicate webhook request detected with hash: ${requestHash}`)
// Return early for duplicate requests to prevent workflow execution
return new NextResponse('Duplicate request', { status: 200 })
}
@@ -137,12 +141,19 @@ export async function POST(
.limit(1)
if (webhooks.length === 0) {
logger.warn(`[${requestId}] No active webhook found for path: ${path}`)
return new NextResponse('Webhook not found', { status: 404 })
}
const { webhook: foundWebhook, workflow: workflowData } = webhooks[0]
foundWorkflow = workflowData
logger.info(`[${requestId}] Found webhook for path ${path}`, {
webhookId: foundWebhook.id,
provider: foundWebhook.provider,
workflowId: foundWorkflow.id,
})
// For WhatsApp, also check for duplicate messages using their message ID
if (foundWebhook.provider === 'whatsapp') {
const data = body?.entry?.[0]?.changes?.[0]?.value
@@ -154,9 +165,7 @@ export async function POST(
// Check if we've already processed this message using Redis
if (messageId && (await hasProcessedMessage(messageId))) {
console.log(
`Duplicate WhatsApp message detected with ID: ${messageId}. Skipping processing.`
)
logger.info(`[${requestId}] Duplicate WhatsApp message detected with ID: ${messageId}`)
// Return early for duplicate messages to prevent workflow execution
return new NextResponse('Duplicate message', { status: 200 })
}
@@ -171,14 +180,23 @@ export async function POST(
await markMessageAsProcessed(requestHash, 60 * 60 * 24)
// Process the webhook synchronously - complete the workflow before returning
const result = await processWebhook(foundWebhook, foundWorkflow, body, request, executionId)
const result = await processWebhook(
foundWebhook,
foundWorkflow,
body,
request,
executionId,
requestId
)
// After workflow execution is complete, return 200 OK
console.log(`Workflow execution complete for WhatsApp message ID: ${messageId}`)
logger.info(
`[${requestId}] Workflow execution complete for WhatsApp message ID: ${messageId}`
)
return result
} else {
// This might be a different type of notification (e.g., status update)
console.log('No messages in WhatsApp payload, might be a status update')
logger.debug(`[${requestId}] No messages in WhatsApp payload, might be a status update`)
return new NextResponse('OK', { status: 200 })
}
}
@@ -187,9 +205,9 @@ export async function POST(
await markMessageAsProcessed(requestHash, 60 * 60 * 24)
// For other providers, continue with synchronous processing
return await processWebhook(foundWebhook, foundWorkflow, body, request, executionId)
return await processWebhook(foundWebhook, foundWorkflow, body, request, executionId, requestId)
} catch (error: any) {
console.error('Error processing webhook:', error)
logger.error(`[${requestId}] Error processing webhook`, error)
// Log the error if we have a workflow ID
if (foundWorkflow?.id) {
@@ -225,7 +243,6 @@ async function generateRequestHash(path: string, body: any): Promise<string> {
return `request:${path}:${hash}`
} catch (error) {
// If hashing fails, use a UUID as fallback
console.error('Error generating request hash:', error)
return `request:${path}:${uuidv4()}`
}
}
@@ -267,7 +284,8 @@ async function processWebhook(
foundWorkflow: any,
body: any,
request: NextRequest,
executionId: string
executionId: string,
requestId: string
): Promise<NextResponse> {
try {
// Handle provider-specific verification and authentication
@@ -289,6 +307,7 @@ async function processWebhook(
if (providerConfig.token) {
const providedToken = authHeader?.startsWith('Bearer ') ? authHeader.substring(7) : null
if (!providedToken || providedToken !== providerConfig.token) {
logger.warn(`[${requestId}] Unauthorized webhook access attempt - invalid token`)
return new NextResponse('Unauthorized', { status: 401 })
}
}
@@ -311,9 +330,12 @@ async function processWebhook(
const timestamp = message.timestamp
const text = message.text?.body
console.log(
`Processing WhatsApp message: ${text ? text.substring(0, 50) : '[no text]'} from ${from}`
)
logger.info(`[${requestId}] Processing WhatsApp message from ${from}`, {
messageId,
textPreview: text
? `${text.substring(0, 30)}${text.length > 30 ? '...' : ''}`
: '[no text]',
})
input = {
whatsapp: {
@@ -358,20 +380,21 @@ async function processWebhook(
// Get the workflow state
if (!foundWorkflow.state) {
console.log(`Workflow ${foundWorkflow.id} has no state, skipping`)
logger.error(`[${requestId}] Workflow ${foundWorkflow.id} has no state`)
return new NextResponse('Workflow state not found', { status: 500 })
}
console.log(`Executing workflow ${foundWorkflow.id} for webhook ${foundWebhook.id}`)
logger.info(
`[${requestId}] Executing workflow ${foundWorkflow.id} for webhook ${foundWebhook.id}`
)
// Get the workflow state
const state = foundWorkflow.state as any
const { blocks, edges, loops } = state
// Use the async version of mergeSubblockState to ensure all values are properly resolved
console.log(`[Webhook Debug] Merging subblock states for workflow ${foundWorkflow.id}...`)
logger.debug(`[${requestId}] Merging subblock states for workflow ${foundWorkflow.id}`)
const mergedStates = await mergeSubblockStateAsync(blocks, foundWorkflow.id)
console.log(`[Webhook Debug] Subblock states merged successfully`)
// Retrieve environment variables for this user
const [userEnv] = await db
@@ -389,7 +412,7 @@ async function processWebhook(
const { decrypted } = await decryptSecret(encryptedValue)
return [key, decrypted] as const
} catch (error: any) {
console.error(`Failed to decrypt ${key}:`, error)
logger.error(`[${requestId}] Failed to decrypt environment variable "${key}"`, error)
throw new Error(`Failed to decrypt environment variable "${key}": ${error.message}`)
}
}
@@ -438,7 +461,6 @@ async function processWebhook(
}
// Ensure the responseFormat is properly structured for OpenAI
// This ensures compatibility with the provider's expected format
if (
processedState.responseFormat &&
typeof processedState.responseFormat === 'object'
@@ -456,7 +478,7 @@ async function processWebhook(
acc[blockId] = processedState
} catch (error) {
console.warn(`Failed to parse responseFormat for block ${blockId}:`, error)
logger.warn(`[${requestId}] Failed to parse responseFormat for block ${blockId}`, error)
acc[blockId] = blockState
}
} else {
@@ -467,10 +489,9 @@ async function processWebhook(
{} as Record<string, Record<string, any>>
)
console.log(`[Webhook Debug] Serialized workflow:`, serializedWorkflow)
console.log(`[Webhook Debug] Processed block states:`, processedBlockStates)
console.log(`[Webhook Debug] Decrypted env vars:`, decryptedEnvVars)
console.log(`[Webhook Debug] Enriched input:`, enrichedInput)
logger.debug(
`[${requestId}] Starting workflow execution with ${Object.keys(processedBlockStates).length} blocks`
)
const executor = new Executor(
serializedWorkflow,
@@ -479,9 +500,11 @@ async function processWebhook(
enrichedInput
)
const result = await executor.execute(foundWorkflow.id)
console.log(`[Webhook Debug] Final execution result:`, JSON.stringify(result))
console.log(`Successfully executed workflow ${foundWorkflow.id}`)
logger.info(`[${requestId}] Successfully executed workflow ${foundWorkflow.id}`, {
success: result.success,
executionTime: result.metadata?.duration,
})
// Log each execution step and the final result
await persistExecutionLogs(foundWorkflow.id, executionId, result, 'webhook')
@@ -489,7 +512,7 @@ async function processWebhook(
// Return the execution result
return NextResponse.json(result, { status: 200 })
} catch (error: any) {
console.error('Error processing webhook:', error)
logger.error(`[${requestId}] Error processing webhook`, error)
// Log the error if we have a workflow ID
if (foundWorkflow?.id) {
+13 -2
View File
@@ -1,21 +1,27 @@
import { NextRequest } from 'next/server'
import { eq } from 'drizzle-orm'
import { v4 as uuidv4 } from 'uuid'
import { createLogger } from '@/lib/logs/console-logger'
import { db } from '@/db'
import { workflow } from '@/db/schema'
import { validateWorkflowAccess } from '../../middleware'
import { createErrorResponse, createSuccessResponse } from '../../utils'
const logger = createLogger('WorkflowDeployAPI')
export const dynamic = 'force-dynamic'
export const runtime = 'nodejs'
export async function POST(request: NextRequest, { params }: { params: Promise<{ id: string }> }) {
const requestId = crypto.randomUUID().slice(0, 8)
const { id } = await params
try {
logger.debug(`[${requestId}] Deploying workflow: ${id}`)
const validation = await validateWorkflowAccess(request, id, false)
if (validation.error) {
logger.warn(`[${requestId}] Workflow deployment failed: ${validation.error.message}`)
return createErrorResponse(validation.error.message, validation.error.status)
}
@@ -33,9 +39,10 @@ export async function POST(request: NextRequest, { params }: { params: Promise<{
})
.where(eq(workflow.id, id))
logger.info(`[${requestId}] Workflow deployed successfully: ${id}`)
return createSuccessResponse({ apiKey, isDeployed: true, deployedAt })
} catch (error: any) {
console.error('Error deploying workflow:', error)
logger.error(`[${requestId}] Error deploying workflow: ${id}`, error)
return createErrorResponse(error.message || 'Failed to deploy workflow', 500)
}
}
@@ -44,12 +51,15 @@ export async function DELETE(
request: NextRequest,
{ params }: { params: Promise<{ id: string }> }
) {
const requestId = crypto.randomUUID().slice(0, 8)
const { id } = await params
try {
logger.debug(`[${requestId}] Undeploying workflow: ${id}`)
const validation = await validateWorkflowAccess(request, id, false)
if (validation.error) {
logger.warn(`[${requestId}] Workflow undeployment failed: ${validation.error.message}`)
return createErrorResponse(validation.error.message, validation.error.status)
}
@@ -63,9 +73,10 @@ export async function DELETE(
})
.where(eq(workflow.id, id))
logger.info(`[${requestId}] Workflow undeployed successfully: ${id}`)
return createSuccessResponse({ isDeployed: false, deployedAt: null, apiKey: null })
} catch (error: any) {
console.error('Error undeploying workflow:', error)
logger.error(`[${requestId}] Error undeploying workflow: ${id}`, error)
return createErrorResponse(error.message || 'Failed to undeploy workflow', 500)
}
}
+32 -17
View File
@@ -2,6 +2,7 @@ import { NextRequest } from 'next/server'
import { eq } from 'drizzle-orm'
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 { decryptSecret } from '@/lib/utils'
import { mergeSubblockState } from '@/stores/workflows/utils'
@@ -13,6 +14,8 @@ import { Serializer } from '@/serializer'
import { validateWorkflowAccess } from '../../middleware'
import { createErrorResponse, createSuccessResponse } from '../../utils'
const logger = createLogger('WorkflowExecuteAPI')
export const dynamic = 'force-dynamic'
export const runtime = 'nodejs'
@@ -22,17 +25,19 @@ const EnvVarsSchema = z.record(z.string())
// Keep track of running executions to prevent overlap
const runningExecutions = new Set<string>()
async function executeWorkflow(workflow: any, input?: any) {
async function executeWorkflow(workflow: any, requestId: string, input?: any) {
const workflowId = workflow.id
const executionId = uuidv4()
// Skip if this workflow is already running
if (runningExecutions.has(workflowId)) {
logger.warn(`[${requestId}] Workflow is already running: ${workflowId}`)
throw new Error('Workflow is already running')
}
try {
runningExecutions.add(workflowId)
logger.info(`[${requestId}] Starting workflow execution: ${workflowId}`)
// Get the workflow state
const state = workflow.state as WorkflowState
@@ -49,6 +54,7 @@ async function executeWorkflow(workflow: any, input?: any) {
.limit(1)
if (!userEnv) {
logger.error(`[${requestId}] No environment variables found for user: ${workflow.userId}`)
throw new Error('No environment variables found for this user')
}
@@ -80,7 +86,10 @@ async function executeWorkflow(workflow: any, input?: any) {
const { decrypted } = await decryptSecret(encryptedValue)
value = (value as string).replace(match, decrypted)
} catch (error: any) {
console.error('Error decrypting value:', error)
logger.error(
`[${requestId}] Error decrypting environment variable "${varName}"`,
error
)
throw new Error(
`Failed to decrypt environment variable "${varName}": ${error.message}`
)
@@ -106,35 +115,27 @@ async function executeWorkflow(workflow: any, input?: any) {
const { decrypted } = await decryptSecret(encryptedValue)
decryptedEnvVars[key] = decrypted
} catch (error: any) {
console.error(`Failed to decrypt ${key}:`, error)
logger.error(`[${requestId}] Failed to decrypt environment variable "${key}"`, error)
throw new Error(`Failed to decrypt environment variable "${key}": ${error.message}`)
}
}
// Process the block states to ensure response formats are properly parsed
// This is crucial for agent blocks with response format
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') {
try {
console.log(
`[API Debug] Block ${blockId} has responseFormat as string:`,
blockState.responseFormat
)
logger.debug(`[${requestId}] Parsing responseFormat for block ${blockId}`)
// Attempt to parse the responseFormat if it's a string
const parsedResponseFormat = JSON.parse(blockState.responseFormat)
console.log(
`[API Debug] Successfully parsed responseFormat for block ${blockId}:`,
parsedResponseFormat
)
acc[blockId] = {
...blockState,
responseFormat: parsedResponseFormat,
}
} catch (error) {
console.warn(`Failed to parse responseFormat for block ${blockId}:`, error)
logger.warn(`[${requestId}] Failed to parse responseFormat for block ${blockId}`, error)
acc[blockId] = blockState
}
} else {
@@ -146,15 +147,23 @@ async function executeWorkflow(workflow: any, input?: any) {
)
// Serialize and execute the workflow
logger.debug(`[${requestId}] Serializing workflow: ${workflowId}`)
const serializedWorkflow = new Serializer().serializeWorkflow(mergedStates, edges, loops)
const executor = new Executor(serializedWorkflow, processedBlockStates, decryptedEnvVars, input)
const result = await executor.execute(workflowId)
logger.info(`[${requestId}] Workflow execution completed: ${workflowId}`, {
success: result.success,
executionTime: result.metadata?.duration,
})
// Log each execution step and the final result
await persistExecutionLogs(workflowId, executionId, result, 'api')
return result
} catch (error: any) {
logger.error(`[${requestId}] Workflow execution failed: ${workflowId}`, error)
// Log the error
await persistExecutionError(workflowId, executionId, error, 'api')
throw error
@@ -164,18 +173,21 @@ async function executeWorkflow(workflow: any, input?: any) {
}
export async function GET(request: NextRequest, { params }: { params: Promise<{ id: string }> }) {
const requestId = crypto.randomUUID().slice(0, 8)
const { id } = await params
try {
logger.debug(`[${requestId}] GET execution request for workflow: ${id}`)
const validation = await validateWorkflowAccess(request, id)
if (validation.error) {
logger.warn(`[${requestId}] Workflow access validation failed: ${validation.error.message}`)
return createErrorResponse(validation.error.message, validation.error.status)
}
const result = await executeWorkflow(validation.workflow)
const result = await executeWorkflow(validation.workflow, requestId)
return createSuccessResponse(result)
} catch (error: any) {
console.error('Error executing workflow:', error)
logger.error(`[${requestId}] Error executing workflow: ${id}`, error)
return createErrorResponse(
error.message || 'Failed to execute workflow',
500,
@@ -185,19 +197,22 @@ export async function GET(request: NextRequest, { params }: { params: Promise<{
}
export async function POST(request: NextRequest, { params }: { params: Promise<{ id: string }> }) {
const requestId = crypto.randomUUID().slice(0, 8)
const { id } = await params
try {
logger.debug(`[${requestId}] POST execution request for workflow: ${id}`)
const validation = await validateWorkflowAccess(request, id)
if (validation.error) {
logger.warn(`[${requestId}] Workflow access validation failed: ${validation.error.message}`)
return createErrorResponse(validation.error.message, validation.error.status)
}
const body = await request.json().catch(() => ({}))
const result = await executeWorkflow(validation.workflow, body)
const result = await executeWorkflow(validation.workflow, requestId, body)
return createSuccessResponse(result)
} catch (error: any) {
console.error('Error executing workflow:', error)
logger.error(`[${requestId}] Error executing workflow: ${id}`, error)
return createErrorResponse(
error.message || 'Failed to execute workflow',
500,
+16 -1
View File
@@ -1,23 +1,38 @@
import { NextRequest } from 'next/server'
import { v4 as uuidv4 } from 'uuid'
import { createLogger } from '@/lib/logs/console-logger'
import { persistLog } from '@/lib/logs/execution-logger'
import { validateWorkflowAccess } from '../../middleware'
import { createErrorResponse, createSuccessResponse } from '../../utils'
const logger = createLogger('WorkflowLogAPI')
export const dynamic = 'force-dynamic'
export async function POST(request: NextRequest, { params }: { params: Promise<{ id: string }> }) {
const requestId = crypto.randomUUID().slice(0, 8)
const { id } = await params
try {
logger.debug(`[${requestId}] Persisting logs for workflow: ${id}`)
const validation = await validateWorkflowAccess(request, id, false)
if (validation.error) {
logger.warn(`[${requestId}] Workflow access validation failed: ${validation.error.message}`)
return createErrorResponse(validation.error.message, validation.error.status)
}
const body = await request.json()
const { logs, executionId } = body
if (!logs || !Array.isArray(logs) || logs.length === 0) {
logger.warn(`[${requestId}] No logs provided for workflow: ${id}`)
return createErrorResponse('No logs provided', 400)
}
logger.info(`[${requestId}] Persisting ${logs.length} logs for workflow: ${id}`, {
executionId,
})
// Persist each log
for (const log of logs) {
await persistLog({
@@ -34,7 +49,7 @@ export async function POST(request: NextRequest, { params }: { params: Promise<{
return createSuccessResponse({ message: 'Logs persisted successfully' })
} catch (error: any) {
console.error('Error persisting logs:', error)
logger.error(`[${requestId}] Error persisting logs for workflow: ${id}`, error)
return createErrorResponse(error.message || 'Failed to persist logs', 500)
}
}
+24 -6
View File
@@ -1,54 +1,72 @@
import { NextRequest } from 'next/server'
import { createLogger } from '@/lib/logs/console-logger'
import { Executor } from '@/executor'
import { SerializedWorkflow } from '@/serializer/types'
import { validateWorkflowAccess } from '../middleware'
import { createErrorResponse, createSuccessResponse } from '../utils'
const logger = createLogger('WorkflowAPI')
export const dynamic = 'force-dynamic'
async function executeWorkflow(workflow: any, input?: any) {
async function executeWorkflow(workflow: any, requestId: string, input?: any) {
try {
logger.info(`[${requestId}] Executing workflow: ${workflow.id}`)
const executor = new Executor(workflow.state as SerializedWorkflow, input)
const result = await executor.execute(workflow.id)
logger.info(`[${requestId}] Workflow execution completed: ${workflow.id}`, {
success: result.success,
})
return result
} catch (error: any) {
console.error('Workflow execution failed:', { workflowId: workflow.id, error })
logger.error(`[${requestId}] Workflow execution failed: ${workflow.id}`, error)
throw new Error(`Execution failed: ${error.message}`)
}
}
export async function GET(request: NextRequest, { params }: { params: Promise<{ id: string }> }) {
const requestId = crypto.randomUUID().slice(0, 8)
try {
const { id } = await params
logger.debug(`[${requestId}] GET request for workflow: ${id}`)
const validation = await validateWorkflowAccess(request, id)
if (validation.error) {
logger.warn(`[${requestId}] Workflow access validation failed: ${validation.error.message}`)
return createErrorResponse(validation.error.message, validation.error.status)
}
const result = await executeWorkflow(validation.workflow)
const result = await executeWorkflow(validation.workflow, requestId)
return createSuccessResponse(result)
} catch (error: any) {
console.error('Error executing workflow:', error)
logger.error(`[${requestId}] Error executing workflow: ${(await params).id}`, error)
return createErrorResponse('Failed to execute workflow', 500, 'EXECUTION_ERROR')
}
}
export async function POST(request: NextRequest, { params }: { params: Promise<{ id: string }> }) {
const requestId = crypto.randomUUID().slice(0, 8)
try {
const { id } = await params
logger.debug(`[${requestId}] POST request for workflow: ${id}`)
const validation = await validateWorkflowAccess(request, id)
if (validation.error) {
logger.warn(`[${requestId}] Workflow access validation failed: ${validation.error.message}`)
return createErrorResponse(validation.error.message, validation.error.status)
}
const body = await request.json().catch(() => ({}))
const result = await executeWorkflow(validation.workflow, body)
const result = await executeWorkflow(validation.workflow, requestId, body)
return createSuccessResponse(result)
} catch (error: any) {
console.error('Error executing workflow:', error)
logger.error(`[${requestId}] Error executing workflow: ${(await params).id}`, error)
return createErrorResponse('Failed to execute workflow', 500, 'EXECUTION_ERROR')
}
}
+13
View File
@@ -1,20 +1,33 @@
import { NextRequest } from 'next/server'
import { createLogger } from '@/lib/logs/console-logger'
import { validateWorkflowAccess } from '../../middleware'
import { createErrorResponse, createSuccessResponse } from '../../utils'
const logger = createLogger('WorkflowStatusAPI')
export async function GET(request: NextRequest, { params }: { params: Promise<{ id: string }> }) {
const requestId = crypto.randomUUID().slice(0, 8)
try {
const { id } = await params
logger.debug(`[${requestId}] Checking status for workflow: ${id}`)
const validation = await validateWorkflowAccess(request, id, false)
if (validation.error) {
logger.warn(`[${requestId}] Workflow access validation failed: ${validation.error.message}`)
return createErrorResponse(validation.error.message, validation.error.status)
}
logger.info(`[${requestId}] Retrieved status for workflow: ${id}`, {
isDeployed: validation.workflow.isDeployed,
})
return createSuccessResponse({
isDeployed: validation.workflow.isDeployed,
deployedAt: validation.workflow.deployedAt,
})
} catch (error) {
logger.error(`[${requestId}] Error getting status for workflow: ${(await params).id}`, error)
return createErrorResponse('Failed to get status', 500)
}
}