Feature/api (#82)

* my test changes for branch protection

* feat(api): introduced 'deploy as an API' button and updated workflows db to include status of deployment

* feat(api): added 'trigger' column for logs table to indicate source of workflow run, persist logs from API executions, removed session validation in favor of API key

* fix(bug): cleanup old reference to JSX element in favor of ReactElement

* feat(api): added persistent notification for one-click deployment with copy boxes for url, keys, & ex curl

* fix(ui/notifications): cleaned up deploy with one-click button ui
This commit is contained in:
waleedlatif1
2025-02-23 13:46:50 -08:00
committed by GitHub
parent 48b6095d53
commit f52de5d1d6
27 changed files with 2184 additions and 1480 deletions
+40
View File
@@ -0,0 +1,40 @@
import { NextRequest } from 'next/server'
import { eq } from 'drizzle-orm'
import { v4 as uuidv4 } from 'uuid'
import { db } from '@/db'
import { workflow } from '@/db/schema'
import { validateWorkflowAccess } from '../../middleware'
import { createErrorResponse, createSuccessResponse } from '../../utils'
export const dynamic = 'force-dynamic'
export const runtime = 'nodejs'
export async function POST(request: NextRequest, { params }: { params: Promise<{ id: string }> }) {
const { id } = await params
try {
const validation = await validateWorkflowAccess(request, id, false)
if (validation.error) {
return createErrorResponse(validation.error.message, validation.error.status)
}
// Generate a new API key
const apiKey = `wf_${uuidv4().replace(/-/g, '')}`
// Update the workflow with the API key and deployment status
await db
.update(workflow)
.set({
apiKey,
isDeployed: true,
deployedAt: new Date(),
})
.where(eq(workflow.id, id))
return createSuccessResponse({ apiKey })
} catch (error: any) {
console.error('Error deploying workflow:', error)
return createErrorResponse(error.message || 'Failed to deploy workflow', 500)
}
}
+214
View File
@@ -0,0 +1,214 @@
import { NextRequest } from 'next/server'
import { eq } from 'drizzle-orm'
import { v4 as uuidv4 } from 'uuid'
import { z } from 'zod'
import { persistLog } from '@/lib/logging'
import { decryptSecret } from '@/lib/utils'
import { WorkflowState } from '@/stores/workflow/types'
import { mergeSubblockState } from '@/stores/workflow/utils'
import { db } from '@/db'
import { environment } from '@/db/schema'
import { Executor } from '@/executor'
import { Serializer } from '@/serializer'
import { validateWorkflowAccess } from '../../middleware'
import { createErrorResponse, createSuccessResponse } from '../../utils'
export const dynamic = 'force-dynamic'
export const runtime = 'nodejs'
// Define the schema for environment variables
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) {
const workflowId = workflow.id
const executionId = uuidv4()
// Skip if this workflow is already running
if (runningExecutions.has(workflowId)) {
throw new Error('Workflow is already running')
}
try {
runningExecutions.add(workflowId)
// Get the workflow state
const state = workflow.state as WorkflowState
const { blocks, edges, loops } = state
// Use the same execution flow as in scheduled executions
const mergedStates = mergeSubblockState(blocks)
// Retrieve environment variables for this user
const [userEnv] = await db
.select()
.from(environment)
.where(eq(environment.userId, workflow.userId))
.limit(1)
if (!userEnv) {
throw new Error('No environment variables found for this user')
}
// Parse and validate environment variables
const variables = EnvVarsSchema.parse(userEnv.variables)
// Replace environment variables in the block states
const currentBlockStates = await Object.entries(mergedStates).reduce(
async (accPromise, [id, block]) => {
const acc = await accPromise
acc[id] = await Object.entries(block.subBlocks).reduce(
async (subAccPromise, [key, subBlock]) => {
const subAcc = await subAccPromise
let value = subBlock.value
// If the value is a string and contains environment variable syntax
if (typeof value === 'string' && value.includes('{{') && value.includes('}}')) {
const matches = value.match(/{{([^}]+)}}/g)
if (matches) {
// Process all matches sequentially
for (const match of matches) {
const varName = match.slice(2, -2) // Remove {{ and }}
const encryptedValue = variables[varName]
if (!encryptedValue) {
throw new Error(`Environment variable "${varName}" was not found`)
}
try {
const { decrypted } = await decryptSecret(encryptedValue)
value = (value as string).replace(match, decrypted)
} catch (error: any) {
console.error('Error decrypting value:', error)
throw new Error(
`Failed to decrypt environment variable "${varName}": ${error.message}`
)
}
}
}
}
subAcc[key] = value
return subAcc
},
Promise.resolve({} as Record<string, any>)
)
return acc
},
Promise.resolve({} as Record<string, Record<string, any>>)
)
// Create a map of decrypted environment variables
const decryptedEnvVars: Record<string, string> = {}
for (const [key, encryptedValue] of Object.entries(variables)) {
try {
const { decrypted } = await decryptSecret(encryptedValue)
decryptedEnvVars[key] = decrypted
} catch (error: any) {
console.error(`Failed to decrypt ${key}:`, error)
throw new Error(`Failed to decrypt environment variable "${key}": ${error.message}`)
}
}
// Serialize and execute the workflow
const serializedWorkflow = new Serializer().serializeWorkflow(mergedStates, edges, loops)
const executor = new Executor(serializedWorkflow, currentBlockStates, decryptedEnvVars)
const result = await executor.execute(workflowId)
// Log each execution step
for (const log of result.logs || []) {
await persistLog({
id: uuidv4(),
workflowId,
executionId,
level: log.success ? 'info' : 'error',
message: `Block ${log.blockName || log.blockId} (${log.blockType}): ${
log.error || 'Completed successfully'
}`,
duration: log.success ? `${log.durationMs}ms` : 'NA',
trigger: 'api',
createdAt: new Date(log.endedAt || log.startedAt),
})
}
// Calculate total duration from successful block logs
const totalDuration = (result.logs || [])
.filter((log) => log.success)
.reduce((sum, log) => sum + log.durationMs, 0)
// Log the final execution result
await persistLog({
id: uuidv4(),
workflowId,
executionId,
level: result.success ? 'info' : 'error',
message: result.success
? 'API workflow executed successfully'
: `API workflow execution failed: ${result.error}`,
duration: result.success ? `${totalDuration}ms` : 'NA',
trigger: 'api',
createdAt: new Date(),
})
return result
} catch (error: any) {
// Log the error
await persistLog({
id: uuidv4(),
workflowId,
executionId,
level: 'error',
message: `API workflow execution failed: ${error.message}`,
duration: 'NA',
trigger: 'api',
createdAt: new Date(),
})
throw error
} finally {
runningExecutions.delete(workflowId)
}
}
export async function GET(request: NextRequest, { params }: { params: Promise<{ id: string }> }) {
const { id } = await params
try {
const validation = await validateWorkflowAccess(request, id)
if (validation.error) {
return createErrorResponse(validation.error.message, validation.error.status)
}
const result = await executeWorkflow(validation.workflow)
return createSuccessResponse(result)
} catch (error: any) {
console.error('Error executing workflow:', error)
return createErrorResponse(
error.message || 'Failed to execute workflow',
500,
'EXECUTION_ERROR'
)
}
}
export async function POST(request: NextRequest, { params }: { params: Promise<{ id: string }> }) {
const { id } = await params
try {
const validation = await validateWorkflowAccess(request, id)
if (validation.error) {
return createErrorResponse(validation.error.message, validation.error.status)
}
const body = await request.json().catch(() => ({}))
const result = await executeWorkflow(validation.workflow, body)
return createSuccessResponse(result)
} catch (error: any) {
console.error('Error executing workflow:', error)
return createErrorResponse(
error.message || 'Failed to execute workflow',
500,
'EXECUTION_ERROR'
)
}
}
+54
View File
@@ -0,0 +1,54 @@
import { NextRequest } from 'next/server'
import { Executor } from '@/executor'
import { SerializedWorkflow } from '@/serializer/types'
import { validateWorkflowAccess } from '../middleware'
import { createErrorResponse, createSuccessResponse } from '../utils'
export const dynamic = 'force-dynamic'
async function executeWorkflow(workflow: any, input?: any) {
try {
const executor = new Executor(workflow.state as SerializedWorkflow, input)
const result = await executor.execute(workflow.id)
return result
} catch (error: any) {
console.error('Workflow execution failed:', { workflowId: workflow.id, error })
throw new Error(`Execution failed: ${error.message}`)
}
}
export async function GET(request: NextRequest, { params }: { params: Promise<{ id: string }> }) {
try {
const { id } = await params
const validation = await validateWorkflowAccess(request, id)
if (validation.error) {
return createErrorResponse(validation.error.message, validation.error.status)
}
const result = await executeWorkflow(validation.workflow)
return createSuccessResponse(result)
} catch (error: any) {
console.error('Error executing workflow:', error)
return createErrorResponse('Failed to execute workflow', 500, 'EXECUTION_ERROR')
}
}
export async function POST(request: NextRequest, { params }: { params: Promise<{ id: string }> }) {
try {
const { id } = await params
const validation = await validateWorkflowAccess(request, id)
if (validation.error) {
return createErrorResponse(validation.error.message, validation.error.status)
}
const body = await request.json().catch(() => ({}))
const result = await executeWorkflow(validation.workflow, body)
return createSuccessResponse(result)
} catch (error: any) {
console.error('Error executing workflow:', error)
return createErrorResponse('Failed to execute workflow', 500, 'EXECUTION_ERROR')
}
}
+20
View File
@@ -0,0 +1,20 @@
import { NextRequest } from 'next/server'
import { validateWorkflowAccess } from '../../middleware'
import { createErrorResponse, createSuccessResponse } from '../../utils'
export async function GET(request: NextRequest, { params }: { params: Promise<{ id: string }> }) {
try {
const { id } = await params
const validation = await validateWorkflowAccess(request, id, false)
if (validation.error) {
return createErrorResponse(validation.error.message, validation.error.status)
}
return createSuccessResponse({
isDeployed: validation.workflow.isDeployed,
deployedAt: validation.workflow.deployedAt,
})
} catch (error) {
return createErrorResponse('Failed to get status', 500)
}
}
+56
View File
@@ -0,0 +1,56 @@
import { NextRequest } from 'next/server'
import { getWorkflowById } from '@/lib/workflows'
export interface ValidationResult {
error?: { message: string; status: number }
workflow?: any
}
export async function validateWorkflowAccess(
request: NextRequest,
workflowId: string,
requireDeployment = true
): Promise<ValidationResult> {
try {
const workflow = await getWorkflowById(workflowId)
if (!workflow) {
return {
error: {
message: 'Workflow not found',
status: 404,
},
}
}
if (requireDeployment) {
if (!workflow.isDeployed) {
return {
error: {
message: 'Workflow is not deployed',
status: 403,
},
}
}
// API key authentication
const apiKey = request.headers.get('x-api-key')
if (!apiKey || !workflow.apiKey || apiKey !== workflow.apiKey) {
return {
error: {
message: 'Unauthorized',
status: 401,
},
}
}
}
return { workflow }
} catch (error) {
console.error('Validation error:', error)
return {
error: {
message: 'Internal server error',
status: 500,
},
}
}
}
+15
View File
@@ -0,0 +1,15 @@
import { NextResponse } from 'next/server'
export function createErrorResponse(error: string, status: number, code?: string) {
return NextResponse.json(
{
error,
code: code || error.toUpperCase().replace(/\s+/g, '_'),
},
{ status }
)
}
export function createSuccessResponse(data: any) {
return NextResponse.json(data)
}