feat(copilot): JSON sanitization logic + operations sequence diff correctness (#1521)

* add state sending capability

* progress

* add ability to add title and description to workflow state

* progress in language

* fix

* cleanup code

* fix type issue

* fix subflow deletion case

* Workflow console tool

* fix lint

---------

Co-authored-by: Siddharth Ganesan <siddharthganesan@gmail.com>
This commit is contained in:
Vikhyath Mondreti
2025-10-02 15:11:03 -07:00
committed by GitHub
co-authored by Siddharth Ganesan
parent 15138629cb
commit 4bc37db547
11 changed files with 1138 additions and 625 deletions
@@ -0,0 +1,59 @@
import { type NextRequest, NextResponse } from 'next/server'
import { env } from '@/lib/env'
import { createLogger } from '@/lib/logs/console/logger'
const logger = createLogger('CopilotTrainingExamplesAPI')
export const runtime = 'nodejs'
export const dynamic = 'force-dynamic'
export async function POST(request: NextRequest) {
const baseUrl = env.AGENT_INDEXER_URL
if (!baseUrl) {
logger.error('Missing AGENT_INDEXER_URL environment variable')
return NextResponse.json({ error: 'Missing AGENT_INDEXER_URL env' }, { status: 500 })
}
const apiKey = env.AGENT_INDEXER_API_KEY
if (!apiKey) {
logger.error('Missing AGENT_INDEXER_API_KEY environment variable')
return NextResponse.json({ error: 'Missing AGENT_INDEXER_API_KEY env' }, { status: 500 })
}
try {
const body = await request.json()
logger.info('Sending workflow example to agent indexer', {
hasJsonField: typeof body?.json === 'string',
})
const upstream = await fetch(`${baseUrl}/examples/add`, {
method: 'POST',
headers: {
'Content-Type': 'application/json',
'x-api-key': apiKey,
},
body: JSON.stringify(body),
})
if (!upstream.ok) {
const errorText = await upstream.text()
logger.error('Agent indexer rejected the example', {
status: upstream.status,
error: errorText,
})
return NextResponse.json({ error: errorText }, { status: upstream.status })
}
const data = await upstream.json()
logger.info('Successfully sent workflow example to agent indexer')
return NextResponse.json(data, {
headers: { 'content-type': 'application/json' },
})
} catch (err) {
const errorMessage = err instanceof Error ? err.message : 'Failed to add example'
logger.error('Failed to send workflow example', { error: err })
return NextResponse.json({ error: errorMessage }, { status: 502 })
}
}
@@ -76,6 +76,8 @@ export async function GET() {
telemetryEnabled: userSettings.telemetryEnabled,
emailPreferences: userSettings.emailPreferences ?? {},
billingUsageNotificationsEnabled: userSettings.billingUsageNotificationsEnabled ?? true,
showFloatingControls: userSettings.showFloatingControls ?? true,
showTrainingControls: userSettings.showTrainingControls ?? false,
},
},
{ status: 200 }
@@ -30,6 +30,7 @@ import { Textarea } from '@/components/ui/textarea'
import { cn } from '@/lib/utils'
import { sanitizeForCopilot } from '@/lib/workflows/json-sanitizer'
import { formatEditSequence } from '@/lib/workflows/training/compute-edit-sequence'
import { useCurrentWorkflow } from '@/app/workspace/[workspaceId]/w/[workflowId]/hooks/use-current-workflow'
import { useCopilotTrainingStore } from '@/stores/copilot-training/store'
/**
@@ -52,6 +53,8 @@ export function TrainingModal() {
markDatasetSent,
} = useCopilotTrainingStore()
const currentWorkflow = useCurrentWorkflow()
const [localPrompt, setLocalPrompt] = useState(currentPrompt)
const [localTitle, setLocalTitle] = useState(currentTitle)
const [copiedId, setCopiedId] = useState<string | null>(null)
@@ -63,6 +66,11 @@ export function TrainingModal() {
const [sendingSelected, setSendingSelected] = useState(false)
const [sentDatasets, setSentDatasets] = useState<Set<string>>(new Set())
const [failedDatasets, setFailedDatasets] = useState<Set<string>>(new Set())
const [sendingLiveWorkflow, setSendingLiveWorkflow] = useState(false)
const [liveWorkflowSent, setLiveWorkflowSent] = useState(false)
const [liveWorkflowFailed, setLiveWorkflowFailed] = useState(false)
const [liveWorkflowTitle, setLiveWorkflowTitle] = useState('')
const [liveWorkflowDescription, setLiveWorkflowDescription] = useState('')
const handleStart = () => {
if (localTitle.trim() && localPrompt.trim()) {
@@ -285,6 +293,46 @@ export function TrainingModal() {
}
}
const handleSendLiveWorkflow = async () => {
if (!liveWorkflowTitle.trim() || !liveWorkflowDescription.trim()) {
return
}
setLiveWorkflowSent(false)
setLiveWorkflowFailed(false)
setSendingLiveWorkflow(true)
try {
const sanitizedWorkflow = sanitizeForCopilot(currentWorkflow.workflowState)
const response = await fetch('/api/copilot/training/examples', {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({
json: JSON.stringify(sanitizedWorkflow),
source_path: liveWorkflowTitle,
summary: liveWorkflowDescription,
}),
})
if (!response.ok) {
const error = await response.json()
throw new Error(error.error || 'Failed to send live workflow')
}
setLiveWorkflowSent(true)
setLiveWorkflowTitle('')
setLiveWorkflowDescription('')
setTimeout(() => setLiveWorkflowSent(false), 5000)
} catch (error) {
console.error('Failed to send live workflow:', error)
setLiveWorkflowFailed(true)
setTimeout(() => setLiveWorkflowFailed(false), 5000)
} finally {
setSendingLiveWorkflow(false)
}
}
return (
<Dialog open={showModal} onOpenChange={toggleModal}>
<DialogContent className='max-w-3xl'>
@@ -335,24 +383,24 @@ export function TrainingModal() {
)}
<Tabs defaultValue={isTraining ? 'datasets' : 'new'} className='mt-4'>
<TabsList className='grid w-full grid-cols-2'>
<TabsList className='grid w-full grid-cols-3'>
<TabsTrigger value='new' disabled={isTraining}>
New Session
</TabsTrigger>
<TabsTrigger value='datasets'>Datasets ({datasets.length})</TabsTrigger>
<TabsTrigger value='live'>Send Live State</TabsTrigger>
</TabsList>
{/* New Training Session Tab */}
<TabsContent value='new' className='space-y-4'>
{startSnapshot && (
<div className='rounded-lg border bg-muted/50 p-3'>
<p className='font-medium text-muted-foreground text-sm'>Current Workflow State</p>
<p className='text-sm'>
{Object.keys(startSnapshot.blocks).length} blocks, {startSnapshot.edges.length}{' '}
edges
</p>
</div>
)}
<div className='rounded-lg border bg-muted/50 p-3'>
<p className='mb-2 font-medium text-muted-foreground text-sm'>
Current Workflow State
</p>
<p className='text-sm'>
{currentWorkflow.getBlockCount()} blocks, {currentWorkflow.getEdgeCount()} edges
</p>
</div>
<div className='space-y-2'>
<Label htmlFor='title'>Title</Label>
@@ -628,6 +676,94 @@ export function TrainingModal() {
</>
)}
</TabsContent>
{/* Send Live State Tab */}
<TabsContent value='live' className='space-y-4'>
<div className='rounded-lg border bg-muted/50 p-3'>
<p className='mb-2 font-medium text-muted-foreground text-sm'>
Current Workflow State
</p>
<p className='text-sm'>
{currentWorkflow.getBlockCount()} blocks, {currentWorkflow.getEdgeCount()} edges
</p>
</div>
<div className='space-y-2'>
<Label htmlFor='live-title'>Title</Label>
<Input
id='live-title'
placeholder='e.g., Customer Onboarding Workflow'
value={liveWorkflowTitle}
onChange={(e) => setLiveWorkflowTitle(e.target.value)}
/>
<p className='text-muted-foreground text-xs'>
A short title identifying this workflow
</p>
</div>
<div className='space-y-2'>
<Label htmlFor='live-description'>Description</Label>
<Textarea
id='live-description'
placeholder='Describe what this workflow does...'
value={liveWorkflowDescription}
onChange={(e) => setLiveWorkflowDescription(e.target.value)}
rows={3}
/>
<p className='text-muted-foreground text-xs'>
Explain the purpose and functionality of this workflow
</p>
</div>
<Button
onClick={handleSendLiveWorkflow}
disabled={
!liveWorkflowTitle.trim() ||
!liveWorkflowDescription.trim() ||
sendingLiveWorkflow ||
currentWorkflow.getBlockCount() === 0
}
className='w-full'
>
{sendingLiveWorkflow ? (
<>
<div className='mr-2 h-4 w-4 animate-spin rounded-full border-2 border-current border-t-transparent' />
Sending...
</>
) : liveWorkflowSent ? (
<>
<CheckCircle2 className='mr-2 h-4 w-4' />
Sent Successfully
</>
) : liveWorkflowFailed ? (
<>
<XCircle className='mr-2 h-4 w-4' />
Failed - Try Again
</>
) : (
<>
<Send className='mr-2 h-4 w-4' />
Send Live Workflow State
</>
)}
</Button>
{liveWorkflowSent && (
<div className='rounded-lg border bg-green-50 p-3 dark:bg-green-950/30'>
<p className='text-green-700 text-sm dark:text-green-300'>
Workflow state sent successfully!
</p>
</div>
)}
{liveWorkflowFailed && (
<div className='rounded-lg border bg-red-50 p-3 dark:bg-red-950/30'>
<p className='text-red-700 text-sm dark:text-red-300'>
Failed to send workflow state. Please try again.
</p>
</div>
)}
</TabsContent>
</Tabs>
</DialogContent>
</Dialog>
-1
View File
@@ -268,7 +268,6 @@ async function processWorkflowFromDb(
logger.info('Processed sanitized workflow context', {
workflowId,
blocks: Object.keys(sanitizedState.blocks || {}).length,
edges: sanitizedState.edges.length,
})
// Use the provided kind for the type
return { type: kind, tag, content }
+8
View File
@@ -262,6 +262,14 @@ const ExecutionEntry = z.object({
totalTokens: z.number().nullable(),
blockExecutions: z.array(z.any()), // can be detailed per need
output: z.any().optional(),
errorMessage: z.string().optional(),
errorBlock: z
.object({
blockId: z.string().optional(),
blockName: z.string().optional(),
blockType: z.string().optional(),
})
.optional(),
})
export const ToolResultSchemas = {
@@ -98,7 +98,35 @@ export class EditWorkflowClientTool extends BaseClientTool {
// Prepare currentUserWorkflow JSON from stores to preserve block IDs
let currentUserWorkflow = args?.currentUserWorkflow
if (!currentUserWorkflow) {
const diffStoreState = useWorkflowDiffStore.getState()
let usedDiffWorkflow = false
if (!currentUserWorkflow && diffStoreState.isDiffReady && diffStoreState.diffWorkflow) {
try {
const diffWorkflow = diffStoreState.diffWorkflow
const normalizedDiffWorkflow = {
...diffWorkflow,
blocks: diffWorkflow.blocks || {},
edges: diffWorkflow.edges || [],
loops: diffWorkflow.loops || {},
parallels: diffWorkflow.parallels || {},
}
currentUserWorkflow = JSON.stringify(normalizedDiffWorkflow)
usedDiffWorkflow = true
logger.info('Using diff workflow state as base for edit_workflow operations', {
toolCallId: this.toolCallId,
blocksCount: Object.keys(normalizedDiffWorkflow.blocks).length,
edgesCount: normalizedDiffWorkflow.edges.length,
})
} catch (e) {
logger.warn(
'Failed to serialize diff workflow state; falling back to active workflow',
e as any
)
}
}
if (!currentUserWorkflow && !usedDiffWorkflow) {
try {
const workflowStore = useWorkflowStore.getState()
const fullState = workflowStore.getWorkflowState()
@@ -77,13 +77,13 @@ export interface CopilotBlockMetadata {
name: string
description: string
bestPractices?: string
commonParameters: CopilotSubblockMetadata[]
inputs?: Record<string, any>
inputSchema: CopilotSubblockMetadata[]
inputDefinitions?: Record<string, any>
triggerAllowed?: boolean
authType?: 'OAuth' | 'API Key' | 'Bot Token'
tools: CopilotToolMetadata[]
triggers: CopilotTriggerMetadata[]
operationParameters: Record<string, CopilotSubblockMetadata[]>
operationInputSchema: Record<string, CopilotSubblockMetadata[]>
operations?: Record<
string,
{
@@ -92,7 +92,7 @@ export interface CopilotBlockMetadata {
description?: string
inputs?: Record<string, any>
outputs?: Record<string, any>
parameters?: CopilotSubblockMetadata[]
inputSchema?: CopilotSubblockMetadata[]
}
>
yamlDocumentation?: string
@@ -125,11 +125,11 @@ export const getBlocksMetadataServerTool: BaseServerTool<
id: specialBlock.id,
name: specialBlock.name,
description: specialBlock.description || '',
commonParameters: commonParameters,
inputs: specialBlock.inputs || {},
inputSchema: commonParameters,
inputDefinitions: specialBlock.inputs || {},
tools: [],
triggers: [],
operationParameters,
operationInputSchema: operationParameters,
}
;(metadata as any).subBlocks = undefined
} else {
@@ -192,7 +192,7 @@ export const getBlocksMetadataServerTool: BaseServerTool<
description: toolCfg?.description || undefined,
inputs: { ...filteredToolParams, ...(operationInputs[opId] || {}) },
outputs: toolOutputs,
parameters: operationParameters[opId] || [],
inputSchema: operationParameters[opId] || [],
}
}
@@ -201,13 +201,13 @@ export const getBlocksMetadataServerTool: BaseServerTool<
name: blockConfig.name || blockId,
description: blockConfig.longDescription || blockConfig.description || '',
bestPractices: blockConfig.bestPractices,
commonParameters: commonParameters,
inputs: blockInputs,
inputSchema: commonParameters,
inputDefinitions: blockInputs,
triggerAllowed: !!blockConfig.triggerAllowed,
authType: resolveAuthType(blockConfig.authMode),
tools,
triggers,
operationParameters,
operationInputSchema: operationParameters,
operations,
}
}
@@ -420,7 +420,7 @@ function splitParametersByOperation(
operationParameters[key].push(processed)
}
} else {
// Override description from blockInputs if available (by id or canonicalParamId)
// Override description from inputDefinitions if available (by id or canonicalParamId)
if (blockInputsForDescriptions) {
const candidates = [sb.id, sb.canonicalParamId].filter(Boolean)
for (const key of candidates) {
@@ -11,7 +11,7 @@ import { resolveOutputType } from '@/blocks/utils'
import { generateLoopBlocks, generateParallelBlocks } from '@/stores/workflows/workflow/utils'
interface EditWorkflowOperation {
operation_type: 'add' | 'edit' | 'delete'
operation_type: 'add' | 'edit' | 'delete' | 'insert_into_subflow' | 'extract_from_subflow'
block_id: string
params?: Record<string, any>
}
@@ -22,6 +22,78 @@ interface EditWorkflowParams {
currentUserWorkflow?: string
}
/**
* Helper to create a block state from operation params
*/
function createBlockFromParams(blockId: string, params: any, parentId?: string): any {
const blockConfig = getAllBlocks().find((b) => b.type === params.type)
const blockState: any = {
id: blockId,
type: params.type,
name: params.name,
position: { x: 0, y: 0 },
enabled: params.enabled !== undefined ? params.enabled : true,
horizontalHandles: true,
isWide: false,
advancedMode: params.advancedMode || false,
height: 0,
triggerMode: params.triggerMode || false,
subBlocks: {},
outputs: params.outputs || (blockConfig ? resolveOutputType(blockConfig.outputs) : {}),
data: parentId ? { parentId, extent: 'parent' as const } : {},
}
// Add inputs as subBlocks
if (params.inputs) {
Object.entries(params.inputs).forEach(([key, value]) => {
blockState.subBlocks[key] = {
id: key,
type: 'short-input',
value: value,
}
})
}
// Set up subBlocks from block configuration
if (blockConfig) {
blockConfig.subBlocks.forEach((subBlock) => {
if (!blockState.subBlocks[subBlock.id]) {
blockState.subBlocks[subBlock.id] = {
id: subBlock.id,
type: subBlock.type,
value: null,
}
}
})
}
return blockState
}
/**
* Helper to add connections as edges for a block
*/
function addConnectionsAsEdges(
modifiedState: any,
blockId: string,
connections: Record<string, any>
): void {
Object.entries(connections).forEach(([sourceHandle, targets]) => {
const targetArray = Array.isArray(targets) ? targets : [targets]
targetArray.forEach((targetId: string) => {
modifiedState.edges.push({
id: crypto.randomUUID(),
source: blockId,
sourceHandle,
target: targetId,
targetHandle: 'target',
type: 'default',
})
})
})
}
/**
* Apply operations directly to the workflow JSON state
*/
@@ -43,11 +115,19 @@ function applyOperationsToWorkflowState(
})),
})
// Reorder operations: delete -> add -> edit to ensure consistent application semantics
// Reorder operations: delete -> extract -> add -> insert -> edit
const deletes = operations.filter((op) => op.operation_type === 'delete')
const extracts = operations.filter((op) => op.operation_type === 'extract_from_subflow')
const adds = operations.filter((op) => op.operation_type === 'add')
const inserts = operations.filter((op) => op.operation_type === 'insert_into_subflow')
const edits = operations.filter((op) => op.operation_type === 'edit')
const orderedOperations: EditWorkflowOperation[] = [...deletes, ...adds, ...edits]
const orderedOperations: EditWorkflowOperation[] = [
...deletes,
...extracts,
...adds,
...inserts,
...edits,
]
for (const operation of orderedOperations) {
const { operation_type, block_id, params } = operation
@@ -105,6 +185,23 @@ function applyOperationsToWorkflowState(
block.subBlocks[key].value = value
}
})
// Update loop/parallel configuration in block.data
if (block.type === 'loop') {
block.data = block.data || {}
if (params.inputs.loopType !== undefined) block.data.loopType = params.inputs.loopType
if (params.inputs.iterations !== undefined)
block.data.count = params.inputs.iterations
if (params.inputs.collection !== undefined)
block.data.collection = params.inputs.collection
} else if (block.type === 'parallel') {
block.data = block.data || {}
if (params.inputs.parallelType !== undefined)
block.data.parallelType = params.inputs.parallelType
if (params.inputs.count !== undefined) block.data.count = params.inputs.count
if (params.inputs.collection !== undefined)
block.data.collection = params.inputs.collection
}
}
// Update basic properties
@@ -123,6 +220,50 @@ function applyOperationsToWorkflowState(
}
}
// Handle advanced mode toggle
if (typeof params?.advancedMode === 'boolean') {
block.advancedMode = params.advancedMode
}
// Handle nested nodes update (for loops/parallels)
if (params?.nestedNodes) {
// Remove all existing child blocks
const existingChildren = Object.keys(modifiedState.blocks).filter(
(id) => modifiedState.blocks[id].data?.parentId === block_id
)
existingChildren.forEach((childId) => delete modifiedState.blocks[childId])
// Remove edges to/from removed children
modifiedState.edges = modifiedState.edges.filter(
(edge: any) =>
!existingChildren.includes(edge.source) && !existingChildren.includes(edge.target)
)
// Add new nested blocks
Object.entries(params.nestedNodes).forEach(([childId, childBlock]: [string, any]) => {
const childBlockState = createBlockFromParams(childId, childBlock, block_id)
modifiedState.blocks[childId] = childBlockState
// Add connections for child block
if (childBlock.connections) {
addConnectionsAsEdges(modifiedState, childId, childBlock.connections)
}
})
// Update loop/parallel configuration based on type
if (block.type === 'loop') {
block.data = block.data || {}
if (params.inputs?.loopType) block.data.loopType = params.inputs.loopType
if (params.inputs?.iterations) block.data.count = params.inputs.iterations
if (params.inputs?.collection) block.data.collection = params.inputs.collection
} else if (block.type === 'parallel') {
block.data = block.data || {}
if (params.inputs?.parallelType) block.data.parallelType = params.inputs.parallelType
if (params.inputs?.count) block.data.count = params.inputs.count
if (params.inputs?.collection) block.data.collection = params.inputs.collection
}
}
// Handle connections update (convert to edges)
if (params?.connections) {
// Remove existing edges from this block
@@ -191,82 +332,135 @@ function applyOperationsToWorkflowState(
case 'add': {
if (params?.type && params?.name) {
// Get block configuration
const blockConfig = getAllBlocks().find((block) => block.type === params.type)
// Create new block with proper structure
const newBlock: any = {
id: block_id,
type: params.type,
name: params.name,
position: { x: 0, y: 0 }, // Default position
enabled: true,
horizontalHandles: true,
isWide: false,
advancedMode: false,
height: 0,
triggerMode: false,
subBlocks: {},
outputs: blockConfig ? resolveOutputType(blockConfig.outputs) : {},
data: {},
}
const newBlock = createBlockFromParams(block_id, params)
// Add inputs as subBlocks
if (params.inputs) {
Object.entries(params.inputs).forEach(([key, value]) => {
newBlock.subBlocks[key] = {
id: key,
type: 'short-input',
value: value,
// Handle nested nodes (for loops/parallels created from scratch)
if (params.nestedNodes) {
Object.entries(params.nestedNodes).forEach(([childId, childBlock]: [string, any]) => {
const childBlockState = createBlockFromParams(childId, childBlock, block_id)
modifiedState.blocks[childId] = childBlockState
if (childBlock.connections) {
addConnectionsAsEdges(modifiedState, childId, childBlock.connections)
}
})
}
// Set up subBlocks from block configuration
if (blockConfig) {
blockConfig.subBlocks.forEach((subBlock) => {
if (!newBlock.subBlocks[subBlock.id]) {
newBlock.subBlocks[subBlock.id] = {
id: subBlock.id,
type: subBlock.type,
value: null,
}
// Set loop/parallel data on parent block
if (params.type === 'loop') {
newBlock.data = {
...newBlock.data,
loopType: params.inputs?.loopType || 'for',
...(params.inputs?.collection && { collection: params.inputs.collection }),
...(params.inputs?.iterations && { count: params.inputs.iterations }),
}
})
} else if (params.type === 'parallel') {
newBlock.data = {
...newBlock.data,
parallelType: params.inputs?.parallelType || 'count',
...(params.inputs?.collection && { collection: params.inputs.collection }),
...(params.inputs?.count && { count: params.inputs.count }),
}
}
}
modifiedState.blocks[block_id] = newBlock
// Add connections as edges
if (params.connections) {
Object.entries(params.connections).forEach(([sourceHandle, targets]) => {
const addEdge = (targetBlock: string, targetHandle?: string) => {
modifiedState.edges.push({
id: crypto.randomUUID(),
source: block_id,
sourceHandle: sourceHandle,
target: targetBlock,
targetHandle: targetHandle || 'target',
type: 'default',
})
}
addConnectionsAsEdges(modifiedState, block_id, params.connections)
}
}
break
}
if (typeof targets === 'string') {
addEdge(targets)
} else if (Array.isArray(targets)) {
targets.forEach((target: any) => {
if (typeof target === 'string') {
addEdge(target)
} else if (target?.block) {
addEdge(target.block, target.handle)
}
})
} else if (typeof targets === 'object' && (targets as any)?.block) {
addEdge((targets as any).block, (targets as any).handle)
case 'insert_into_subflow': {
const subflowId = params?.subflowId
if (!subflowId || !params?.type || !params?.name) {
logger.warn('Missing required params for insert_into_subflow', { block_id, params })
break
}
const subflowBlock = modifiedState.blocks[subflowId]
if (!subflowBlock || (subflowBlock.type !== 'loop' && subflowBlock.type !== 'parallel')) {
logger.warn('Subflow block not found or invalid type', {
subflowId,
type: subflowBlock?.type,
})
break
}
// Get block configuration
const blockConfig = getAllBlocks().find((block) => block.type === params.type)
// Check if block already exists (moving into subflow) or is new
const existingBlock = modifiedState.blocks[block_id]
if (existingBlock) {
// Moving existing block into subflow - just update parent
existingBlock.data = {
...existingBlock.data,
parentId: subflowId,
extent: 'parent' as const,
}
// Update inputs if provided
if (params.inputs) {
Object.entries(params.inputs).forEach(([key, value]) => {
if (!existingBlock.subBlocks[key]) {
existingBlock.subBlocks[key] = { id: key, type: 'short-input', value }
} else {
existingBlock.subBlocks[key].value = value
}
})
}
} else {
// Create new block as child of subflow
const newBlock = createBlockFromParams(block_id, params, subflowId)
modifiedState.blocks[block_id] = newBlock
}
// Add/update connections as edges
if (params.connections) {
// Remove existing edges from this block
modifiedState.edges = modifiedState.edges.filter((edge: any) => edge.source !== block_id)
// Add new connections
addConnectionsAsEdges(modifiedState, block_id, params.connections)
}
break
}
case 'extract_from_subflow': {
const subflowId = params?.subflowId
if (!subflowId) {
logger.warn('Missing subflowId for extract_from_subflow', { block_id })
break
}
const block = modifiedState.blocks[block_id]
if (!block) {
logger.warn('Block not found for extraction', { block_id })
break
}
// Verify it's actually a child of this subflow
if (block.data?.parentId !== subflowId) {
logger.warn('Block is not a child of specified subflow', {
block_id,
actualParent: block.data?.parentId,
specifiedParent: subflowId,
})
}
// Remove parent relationship
if (block.data) {
block.data.parentId = undefined
block.data.extent = undefined
}
// Note: We keep the block and its edges, just remove parent relationship
// The block becomes a root-level block
break
}
}
@@ -43,6 +43,12 @@ interface ExecutionEntry {
totalTokens: number | null
blockExecutions: BlockExecution[]
output?: any
errorMessage?: string
errorBlock?: {
blockId?: string
blockName?: string
blockType?: string
}
}
function extractBlockExecutionsFromTraceSpans(traceSpans: any[]): BlockExecution[] {
@@ -74,6 +80,140 @@ function extractBlockExecutionsFromTraceSpans(traceSpans: any[]): BlockExecution
return blockExecutions
}
function normalizeErrorMessage(errorValue: unknown): string | undefined {
if (!errorValue) return undefined
if (typeof errorValue === 'string') return errorValue
if (errorValue instanceof Error) return errorValue.message
if (typeof errorValue === 'object') {
try {
return JSON.stringify(errorValue)
} catch {}
}
try {
return String(errorValue)
} catch {
return undefined
}
}
function extractErrorFromExecutionData(executionData: any): ExecutionEntry['errorBlock'] & {
message?: string
} {
if (!executionData) return {}
const errorDetails = executionData.errorDetails
if (errorDetails) {
const message = normalizeErrorMessage(errorDetails.error || errorDetails.message)
if (message) {
return {
message,
blockId: errorDetails.blockId,
blockName: errorDetails.blockName,
blockType: errorDetails.blockType,
}
}
}
const finalOutputError = normalizeErrorMessage(executionData.finalOutput?.error)
if (finalOutputError) {
return {
message: finalOutputError,
blockName: 'Workflow',
}
}
const genericError = normalizeErrorMessage(executionData.error)
if (genericError) {
return {
message: genericError,
blockName: 'Workflow',
}
}
return {}
}
function extractErrorFromTraceSpans(traceSpans: any[]): ExecutionEntry['errorBlock'] & {
message?: string
} {
if (!Array.isArray(traceSpans) || traceSpans.length === 0) return {}
const queue = [...traceSpans]
while (queue.length > 0) {
const span = queue.shift()
if (!span || typeof span !== 'object') continue
const message =
normalizeErrorMessage(span.output?.error) ||
normalizeErrorMessage(span.error) ||
normalizeErrorMessage(span.output?.message) ||
normalizeErrorMessage(span.message)
const status = span.status
if (status === 'error' || message) {
return {
message,
blockId: span.blockId,
blockName: span.blockName || span.name || (span.blockId ? undefined : 'Workflow'),
blockType: span.blockType || span.type,
}
}
if (Array.isArray(span.children)) {
queue.push(...span.children)
}
}
return {}
}
function deriveExecutionErrorSummary(params: {
blockExecutions: BlockExecution[]
traceSpans: any[]
executionData: any
}): { message?: string; block?: ExecutionEntry['errorBlock'] } {
const { blockExecutions, traceSpans, executionData } = params
const blockError = blockExecutions.find((block) => block.status === 'error' && block.errorMessage)
if (blockError) {
return {
message: blockError.errorMessage,
block: {
blockId: blockError.blockId,
blockName: blockError.blockName,
blockType: blockError.blockType,
},
}
}
const executionDataError = extractErrorFromExecutionData(executionData)
if (executionDataError.message) {
return {
message: executionDataError.message,
block: {
blockId: executionDataError.blockId,
blockName:
executionDataError.blockName || (executionDataError.blockId ? undefined : 'Workflow'),
blockType: executionDataError.blockType,
},
}
}
const traceError = extractErrorFromTraceSpans(traceSpans)
if (traceError.message) {
return {
message: traceError.message,
block: {
blockId: traceError.blockId,
blockName: traceError.blockName || (traceError.blockId ? undefined : 'Workflow'),
blockType: traceError.blockType,
},
}
}
return {}
}
export const getWorkflowConsoleServerTool: BaseServerTool<GetWorkflowConsoleArgs, any> = {
name: 'get_workflow_console',
async execute(rawArgs: GetWorkflowConsoleArgs): Promise<any> {
@@ -108,7 +248,8 @@ export const getWorkflowConsoleServerTool: BaseServerTool<GetWorkflowConsoleArgs
.limit(limit)
const formattedEntries: ExecutionEntry[] = executionLogs.map((log) => {
const traceSpans = (log.executionData as any)?.traceSpans || []
const executionData = log.executionData as any
const traceSpans = executionData?.traceSpans || []
const blockExecutions = includeDetails ? extractBlockExecutionsFromTraceSpans(traceSpans) : []
let finalOutput: any
@@ -125,6 +266,12 @@ export const getWorkflowConsoleServerTool: BaseServerTool<GetWorkflowConsoleArgs
if (outputBlock) finalOutput = outputBlock.outputData
}
const { message: errorMessage, block: errorBlock } = deriveExecutionErrorSummary({
blockExecutions,
traceSpans,
executionData,
})
return {
id: log.id,
executionId: log.executionId,
@@ -137,6 +284,8 @@ export const getWorkflowConsoleServerTool: BaseServerTool<GetWorkflowConsoleArgs
totalTokens: (log.cost as any)?.tokens?.total ?? null,
blockExecutions,
output: finalOutput,
errorMessage: errorMessage,
errorBlock: errorBlock,
}
})
+202 -262
View File
@@ -1,43 +1,30 @@
import type { Edge } from 'reactflow'
import type {
BlockState,
Loop,
Parallel,
Position,
WorkflowState,
} from '@/stores/workflows/workflow/types'
import type { BlockState, Loop, Parallel, WorkflowState } from '@/stores/workflows/workflow/types'
/**
* Sanitized workflow state for copilot (removes all UI-specific data)
* Connections are embedded in blocks for consistency with operations format
* Loops and parallels use nested structure - no separate loops/parallels objects
*/
export interface CopilotWorkflowState {
blocks: Record<string, CopilotBlockState>
edges: CopilotEdge[]
loops: Record<string, Loop>
parallels: Record<string, Parallel>
}
/**
* Block state for copilot (no positions, no UI dimensions)
* Block state for copilot (no positions, no UI dimensions, no redundant IDs)
* Connections are embedded here instead of separate edges array
* Loops and parallels have nested structure for clarity
*/
export interface CopilotBlockState {
id: string
type: string
name: string
subBlocks: BlockState['subBlocks']
inputs?: Record<string, string | number | string[][]>
outputs: BlockState['outputs']
connections?: Record<string, string | string[]>
nestedNodes?: Record<string, CopilotBlockState>
enabled: boolean
advancedMode?: boolean
triggerMode?: boolean
// Keep semantic data only (no width/height)
data?: {
parentId?: string
extent?: 'parent'
loopType?: 'for' | 'forEach'
parallelType?: 'collection' | 'count'
collection?: any
count?: number
}
}
/**
@@ -66,55 +53,208 @@ export interface ExportWorkflowState {
}
/**
* Sanitize workflow state for copilot by removing all UI-specific data
* Copilot doesn't need to see positions, dimensions, or visual styling
* Check if a subblock contains sensitive/secret data
*/
export function sanitizeForCopilot(state: WorkflowState): CopilotWorkflowState {
const sanitizedBlocks: Record<string, CopilotBlockState> = {}
function isSensitiveSubBlock(key: string, subBlock: BlockState['subBlocks'][string]): boolean {
// Check if it's an OAuth input type
if (subBlock.type === 'oauth-input') {
return true
}
// Sanitize blocks - remove position and UI-only fields
Object.entries(state.blocks).forEach(([blockId, block]) => {
const sanitizedData: CopilotBlockState['data'] = block.data
? {
// Keep semantic fields only
...(block.data.parentId !== undefined && { parentId: block.data.parentId }),
...(block.data.extent !== undefined && { extent: block.data.extent }),
...(block.data.loopType !== undefined && { loopType: block.data.loopType }),
...(block.data.parallelType !== undefined && { parallelType: block.data.parallelType }),
...(block.data.collection !== undefined && { collection: block.data.collection }),
...(block.data.count !== undefined && { count: block.data.count }),
}
: undefined
// Check if the field name suggests it contains sensitive data
const sensitivePattern = /credential|oauth|api[_-]?key|token|secret|auth|password|bearer/i
if (sensitivePattern.test(key)) {
return true
}
sanitizedBlocks[blockId] = {
id: block.id,
type: block.type,
name: block.name,
subBlocks: block.subBlocks,
outputs: block.outputs,
enabled: block.enabled,
...(block.advancedMode !== undefined && { advancedMode: block.advancedMode }),
...(block.triggerMode !== undefined && { triggerMode: block.triggerMode }),
...(sanitizedData && Object.keys(sanitizedData).length > 0 && { data: sanitizedData }),
// Check if the value itself looks like a secret (but not environment variable references)
if (typeof subBlock.value === 'string' && subBlock.value.length > 0) {
// Don't sanitize environment variable references like {{VAR_NAME}}
if (subBlock.value.startsWith('{{') && subBlock.value.endsWith('}}')) {
return false
}
// If it matches sensitive patterns in the value, it's likely a hardcoded secret
if (sensitivePattern.test(subBlock.value)) {
return true
}
}
return false
}
/**
* Sanitize subblocks by removing null values, secrets, and simplifying structure
* Maps each subblock key directly to its value instead of the full object
*/
function sanitizeSubBlocks(
subBlocks: BlockState['subBlocks']
): Record<string, string | number | string[][]> {
const sanitized: Record<string, string | number | string[][]> = {}
Object.entries(subBlocks).forEach(([key, subBlock]) => {
// Skip null/undefined values
if (subBlock.value === null || subBlock.value === undefined) {
return
}
// For sensitive fields, either omit or replace with placeholder
if (isSensitiveSubBlock(key, subBlock)) {
// If it's an environment variable reference, keep it
if (
typeof subBlock.value === 'string' &&
subBlock.value.startsWith('{{') &&
subBlock.value.endsWith('}}')
) {
sanitized[key] = subBlock.value
}
// Otherwise omit the sensitive value entirely
return
}
// For non-sensitive, non-null values, include them
sanitized[key] = subBlock.value
})
return sanitized
}
/**
* Reconstruct full subBlock structure from simplified copilot format
* Uses existing block structure as template for id and type fields
*/
function reconstructSubBlocks(
simplifiedSubBlocks: Record<string, string | number | string[][]>,
existingSubBlocks?: BlockState['subBlocks']
): BlockState['subBlocks'] {
const reconstructed: BlockState['subBlocks'] = {}
Object.entries(simplifiedSubBlocks).forEach(([key, value]) => {
const existingSubBlock = existingSubBlocks?.[key]
reconstructed[key] = {
id: existingSubBlock?.id || key,
type: existingSubBlock?.type || 'short-input',
value,
}
})
// Sanitize edges - keep only semantic connection data
const sanitizedEdges: CopilotEdge[] = state.edges.map((edge) => ({
id: edge.id,
source: edge.source,
target: edge.target,
...(edge.sourceHandle !== undefined &&
edge.sourceHandle !== null && { sourceHandle: edge.sourceHandle }),
...(edge.targetHandle !== undefined &&
edge.targetHandle !== null && { targetHandle: edge.targetHandle }),
}))
return reconstructed
}
/**
* Extract connections for a block from edges and format as operations-style connections
*/
function extractConnectionsForBlock(
blockId: string,
edges: WorkflowState['edges']
): Record<string, string | string[]> | undefined {
const connections: Record<string, string[]> = {}
// Find all outgoing edges from this block
const outgoingEdges = edges.filter((edge) => edge.source === blockId)
if (outgoingEdges.length === 0) {
return undefined
}
// Group by source handle
for (const edge of outgoingEdges) {
const handle = edge.sourceHandle || 'source'
if (!connections[handle]) {
connections[handle] = []
}
connections[handle].push(edge.target)
}
// Simplify single-element arrays to just the string
const simplified: Record<string, string | string[]> = {}
for (const [handle, targets] of Object.entries(connections)) {
simplified[handle] = targets.length === 1 ? targets[0] : targets
}
return simplified
}
/**
* Sanitize workflow state for copilot by removing all UI-specific data
* Creates nested structure for loops/parallels with their child blocks inside
*/
export function sanitizeForCopilot(state: WorkflowState): CopilotWorkflowState {
const sanitizedBlocks: Record<string, CopilotBlockState> = {}
const processedBlocks = new Set<string>()
// Helper to find child blocks of a parent (loop/parallel container)
const findChildBlocks = (parentId: string): string[] => {
return Object.keys(state.blocks).filter(
(blockId) => state.blocks[blockId].data?.parentId === parentId
)
}
// Helper to recursively sanitize a block and its children
const sanitizeBlock = (blockId: string, block: BlockState): CopilotBlockState => {
const connections = extractConnectionsForBlock(blockId, state.edges)
// For loop/parallel blocks, extract config from block.data instead of subBlocks
let inputs: Record<string, string | number | string[][]> = {}
if (block.type === 'loop' || block.type === 'parallel') {
// Extract configuration from block.data
if (block.data?.loopType) inputs.loopType = block.data.loopType
if (block.data?.count !== undefined) inputs.iterations = block.data.count
if (block.data?.collection !== undefined) inputs.collection = block.data.collection
if (block.data?.parallelType) inputs.parallelType = block.data.parallelType
} else {
// For regular blocks, sanitize subBlocks
inputs = sanitizeSubBlocks(block.subBlocks)
}
// Check if this is a loop or parallel (has children)
const childBlockIds = findChildBlocks(blockId)
const nestedNodes: Record<string, CopilotBlockState> = {}
if (childBlockIds.length > 0) {
// Recursively sanitize child blocks
childBlockIds.forEach((childId) => {
const childBlock = state.blocks[childId]
if (childBlock) {
nestedNodes[childId] = sanitizeBlock(childId, childBlock)
processedBlocks.add(childId)
}
})
}
const result: CopilotBlockState = {
type: block.type,
name: block.name,
outputs: block.outputs,
enabled: block.enabled,
}
if (Object.keys(inputs).length > 0) result.inputs = inputs
if (connections) result.connections = connections
if (Object.keys(nestedNodes).length > 0) result.nestedNodes = nestedNodes
if (block.advancedMode !== undefined) result.advancedMode = block.advancedMode
if (block.triggerMode !== undefined) result.triggerMode = block.triggerMode
return result
}
// Process only root-level blocks (those without a parent)
Object.entries(state.blocks).forEach(([blockId, block]) => {
// Skip if already processed as a child
if (processedBlocks.has(blockId)) return
// Skip if it has a parent (it will be processed as nested)
if (block.data?.parentId) return
sanitizedBlocks[blockId] = sanitizeBlock(blockId, block)
})
return {
blocks: sanitizedBlocks,
edges: sanitizedEdges,
loops: state.loops || {},
parallels: state.parallels || {},
}
}
@@ -167,203 +307,3 @@ export function sanitizeForExport(state: WorkflowState): ExportWorkflowState {
state: clonedState,
}
}
/**
* Validate that edges reference existing blocks
*/
export function validateEdges(
blocks: Record<string, any>,
edges: CopilotEdge[]
): {
valid: boolean
errors: string[]
} {
const errors: string[] = []
const blockIds = new Set(Object.keys(blocks))
edges.forEach((edge, index) => {
if (!blockIds.has(edge.source)) {
errors.push(`Edge ${index} references non-existent source block: ${edge.source}`)
}
if (!blockIds.has(edge.target)) {
errors.push(`Edge ${index} references non-existent target block: ${edge.target}`)
}
})
return {
valid: errors.length === 0,
errors,
}
}
/**
* Generate position for a new block based on its connections
* Uses compact horizontal spacing and intelligent positioning
*/
export function generatePositionForNewBlock(
blockId: string,
edges: CopilotEdge[],
existingBlocks: Record<string, BlockState>
): Position {
const HORIZONTAL_SPACING = 550
const VERTICAL_SPACING = 200
const incomingEdges = edges.filter((e) => e.target === blockId)
if (incomingEdges.length > 0) {
const sourceBlocks = incomingEdges
.map((e) => existingBlocks[e.source])
.filter((b) => b !== undefined)
if (sourceBlocks.length > 0) {
const rightmostX = Math.max(...sourceBlocks.map((b) => b.position.x))
const avgY = sourceBlocks.reduce((sum, b) => sum + b.position.y, 0) / sourceBlocks.length
return {
x: rightmostX + HORIZONTAL_SPACING,
y: avgY,
}
}
}
const outgoingEdges = edges.filter((e) => e.source === blockId)
if (outgoingEdges.length > 0) {
const targetBlocks = outgoingEdges
.map((e) => existingBlocks[e.target])
.filter((b) => b !== undefined)
if (targetBlocks.length > 0) {
const leftmostX = Math.min(...targetBlocks.map((b) => b.position.x))
const avgY = targetBlocks.reduce((sum, b) => sum + b.position.y, 0) / targetBlocks.length
return {
x: Math.max(150, leftmostX - HORIZONTAL_SPACING),
y: avgY,
}
}
}
const existingPositions = Object.values(existingBlocks).map((b) => b.position)
if (existingPositions.length > 0) {
const maxY = Math.max(...existingPositions.map((p) => p.y))
return {
x: 150,
y: maxY + VERTICAL_SPACING,
}
}
return { x: 150, y: 300 }
}
/**
* Merge sanitized copilot state with full UI state
* Preserves positions for existing blocks, generates positions for new blocks
*/
export function mergeWithUIState(
sanitized: CopilotWorkflowState,
fullState: WorkflowState
): WorkflowState {
const mergedBlocks: Record<string, BlockState> = {}
const existingBlocks = fullState.blocks
// Convert sanitized edges to full edges for position generation
const sanitizedEdges = sanitized.edges
// Process each block from sanitized state
Object.entries(sanitized.blocks).forEach(([blockId, sanitizedBlock]) => {
const existingBlock = existingBlocks[blockId]
if (existingBlock) {
// Existing block - preserve position and UI fields, update semantic fields
mergedBlocks[blockId] = {
...existingBlock,
// Update semantic fields from sanitized
type: sanitizedBlock.type,
name: sanitizedBlock.name,
subBlocks: sanitizedBlock.subBlocks,
outputs: sanitizedBlock.outputs,
enabled: sanitizedBlock.enabled,
advancedMode: sanitizedBlock.advancedMode,
triggerMode: sanitizedBlock.triggerMode,
// Merge data carefully
data: sanitizedBlock.data
? {
...existingBlock.data,
...sanitizedBlock.data,
}
: existingBlock.data,
}
} else {
// New block - generate position
const position = generatePositionForNewBlock(blockId, sanitizedEdges, existingBlocks)
mergedBlocks[blockId] = {
id: sanitizedBlock.id,
type: sanitizedBlock.type,
name: sanitizedBlock.name,
position,
subBlocks: sanitizedBlock.subBlocks,
outputs: sanitizedBlock.outputs,
enabled: sanitizedBlock.enabled,
horizontalHandles: true,
isWide: false,
height: 0,
advancedMode: sanitizedBlock.advancedMode,
triggerMode: sanitizedBlock.triggerMode,
data: sanitizedBlock.data
? {
...sanitizedBlock.data,
// Add UI dimensions if it's a container
...(sanitizedBlock.type === 'loop' || sanitizedBlock.type === 'parallel'
? {
width: 500,
height: 300,
type: 'subflowNode',
}
: {}),
}
: undefined,
}
}
})
// Convert sanitized edges to full edges
const mergedEdges: Edge[] = sanitized.edges.map((edge) => {
// Try to find existing edge to preserve styling
const existingEdge = fullState.edges.find(
(e) =>
e.source === edge.source &&
e.target === edge.target &&
e.sourceHandle === edge.sourceHandle &&
e.targetHandle === edge.targetHandle
)
if (existingEdge) {
return existingEdge
}
// New edge - create with defaults
return {
id: edge.id,
source: edge.source,
target: edge.target,
sourceHandle: edge.sourceHandle,
targetHandle: edge.targetHandle,
type: 'default',
data: {},
} as Edge
})
return {
blocks: mergedBlocks,
edges: mergedEdges,
loops: sanitized.loops,
parallels: sanitized.parallels,
lastSaved: Date.now(),
// Preserve deployment info
isDeployed: fullState.isDeployed,
deployedAt: fullState.deployedAt,
deploymentStatuses: fullState.deploymentStatuses,
}
}
@@ -1,29 +1,19 @@
import type { CopilotWorkflowState } from '@/lib/workflows/json-sanitizer'
export interface EditOperation {
operation_type: 'add' | 'edit' | 'delete'
operation_type: 'add' | 'edit' | 'delete' | 'insert_into_subflow' | 'extract_from_subflow'
block_id: string
params?: {
type?: string
name?: string
outputs?: Record<string, any>
enabled?: boolean
triggerMode?: boolean
advancedMode?: boolean
inputs?: Record<string, any>
connections?: Record<string, any>
removeEdges?: Array<{ targetBlockId: string; sourceHandle?: string }>
loopConfig?: {
nodes?: string[]
iterations?: number
loopType?: 'for' | 'forEach'
forEachItems?: any
}
parallelConfig?: {
nodes?: string[]
distribution?: any
count?: number
parallelType?: 'count' | 'collection'
}
parentId?: string
extent?: 'parent'
nestedNodes?: Record<string, any>
subflowId?: string
}
}
@@ -38,6 +28,79 @@ export interface WorkflowDiff {
}
}
/**
* Flatten nested blocks into a single-level map for comparison
* Returns map of blockId -> {block, parentId}
*/
function flattenBlocks(
blocks: Record<string, any>
): Record<string, { block: any; parentId?: string }> {
const flattened: Record<string, { block: any; parentId?: string }> = {}
const processBlock = (blockId: string, block: any, parentId?: string) => {
flattened[blockId] = { block, parentId }
// Recursively process nested nodes
if (block.nestedNodes) {
Object.entries(block.nestedNodes).forEach(([nestedId, nestedBlock]) => {
processBlock(nestedId, nestedBlock, blockId)
})
}
}
Object.entries(blocks).forEach(([blockId, block]) => {
processBlock(blockId, block)
})
return flattened
}
/**
* Extract all edges from blocks with embedded connections (including nested)
*/
function extractAllEdgesFromBlocks(blocks: Record<string, any>): Array<{
source: string
target: string
sourceHandle?: string | null
targetHandle?: string | null
}> {
const edges: Array<{
source: string
target: string
sourceHandle?: string | null
targetHandle?: string | null
}> = []
const processBlockConnections = (block: any, blockId: string) => {
if (block.connections) {
Object.entries(block.connections).forEach(([sourceHandle, targets]) => {
const targetArray = Array.isArray(targets) ? targets : [targets]
targetArray.forEach((target: string) => {
edges.push({
source: blockId,
target,
sourceHandle,
targetHandle: 'target',
})
})
})
}
// Process nested nodes
if (block.nestedNodes) {
Object.entries(block.nestedNodes).forEach(([nestedId, nestedBlock]) => {
processBlockConnections(nestedBlock, nestedId)
})
}
}
Object.entries(blocks).forEach(([blockId, block]) => {
processBlockConnections(block, blockId)
})
return edges
}
/**
* Compute the edit sequence (operations) needed to transform startState into endState
* This analyzes the differences and generates operations that can recreate the changes
@@ -51,12 +114,14 @@ export function computeEditSequence(
const startBlocks = startState.blocks || {}
const endBlocks = endState.blocks || {}
const startEdges = startState.edges || []
const endEdges = endState.edges || []
const startLoops = startState.loops || {}
const endLoops = endState.loops || {}
const startParallels = startState.parallels || {}
const endParallels = endState.parallels || {}
// Flatten nested blocks for comparison (includes nested nodes at top level)
const startFlattened = flattenBlocks(startBlocks)
const endFlattened = flattenBlocks(endBlocks)
// Extract edges from connections for tracking
const startEdges = extractAllEdgesFromBlocks(startBlocks)
const endEdges = extractAllEdgesFromBlocks(endBlocks)
// Track statistics
let blocksAdded = 0
@@ -65,74 +130,171 @@ export function computeEditSequence(
let edgesChanged = 0
let subflowsChanged = 0
// Track which blocks are being deleted (including subflows)
const deletedBlocks = new Set<string>()
for (const blockId in startFlattened) {
if (!(blockId in endFlattened)) {
deletedBlocks.add(blockId)
}
}
// 1. Find deleted blocks (exist in start but not in end)
for (const blockId in startBlocks) {
if (!(blockId in endBlocks)) {
operations.push({
operation_type: 'delete',
block_id: blockId,
})
blocksDeleted++
for (const blockId in startFlattened) {
if (!(blockId in endFlattened)) {
const { parentId } = startFlattened[blockId]
// Skip if parent is also being deleted (cascade delete is implicit)
if (parentId && deletedBlocks.has(parentId)) {
continue
}
if (parentId) {
// Block was inside a subflow and was removed (but subflow still exists)
operations.push({
operation_type: 'extract_from_subflow',
block_id: blockId,
params: {
subflowId: parentId,
},
})
subflowsChanged++
} else {
// Regular block deletion
operations.push({
operation_type: 'delete',
block_id: blockId,
})
blocksDeleted++
}
}
}
// 2. Find added blocks (exist in end but not in start)
for (const blockId in endBlocks) {
if (!(blockId in startBlocks)) {
const block = endBlocks[blockId]
const addParams: EditOperation['params'] = {
type: block.type,
name: block.name,
inputs: extractInputValues(block),
connections: extractConnections(blockId, endEdges),
triggerMode: Boolean(block?.triggerMode),
}
for (const blockId in endFlattened) {
if (!(blockId in startFlattened)) {
const { block, parentId } = endFlattened[blockId]
if (parentId) {
// Block was added inside a subflow - include full block state
const addParams: EditOperation['params'] = {
subflowId: parentId,
type: block.type,
name: block.name,
outputs: block.outputs,
enabled: block.enabled !== undefined ? block.enabled : true,
...(block?.triggerMode !== undefined && { triggerMode: Boolean(block.triggerMode) }),
...(block?.advancedMode !== undefined && { advancedMode: Boolean(block.advancedMode) }),
}
// Add loop/parallel configuration if this block is in a subflow
const loopConfig = findLoopConfigForBlock(blockId, endLoops)
if (loopConfig) {
;(addParams as any).loopConfig = loopConfig
// Add inputs if present
const inputs = extractInputValues(block)
if (Object.keys(inputs).length > 0) {
addParams.inputs = inputs
}
// Add connections if present
const connections = extractConnections(blockId, endEdges)
if (connections && Object.keys(connections).length > 0) {
addParams.connections = connections
}
operations.push({
operation_type: 'insert_into_subflow',
block_id: blockId,
params: addParams,
})
subflowsChanged++
}
} else {
// Regular block addition at root level
const addParams: EditOperation['params'] = {
type: block.type,
name: block.name,
...(block?.triggerMode !== undefined && { triggerMode: Boolean(block.triggerMode) }),
...(block?.advancedMode !== undefined && { advancedMode: Boolean(block.advancedMode) }),
}
const parallelConfig = findParallelConfigForBlock(blockId, endParallels)
if (parallelConfig) {
;(addParams as any).parallelConfig = parallelConfig
subflowsChanged++
}
// Add inputs if present
const inputs = extractInputValues(block)
if (Object.keys(inputs).length > 0) {
addParams.inputs = inputs
}
// Add parent-child relationship if present
if (block.data?.parentId) {
addParams.parentId = block.data.parentId
addParams.extent = block.data.extent
}
// Add connections if present
const connections = extractConnections(blockId, endEdges)
if (connections && Object.keys(connections).length > 0) {
addParams.connections = connections
}
operations.push({
operation_type: 'add',
block_id: blockId,
params: addParams,
})
blocksAdded++
// Add nested nodes if present (for loops/parallels created from scratch)
if (block.nestedNodes && Object.keys(block.nestedNodes).length > 0) {
addParams.nestedNodes = block.nestedNodes
subflowsChanged++
}
operations.push({
operation_type: 'add',
block_id: blockId,
params: addParams,
})
blocksAdded++
}
}
}
// 3. Find modified blocks (exist in both but have changes)
for (const blockId in endBlocks) {
if (blockId in startBlocks) {
const startBlock = startBlocks[blockId]
const endBlock = endBlocks[blockId]
const changes = computeBlockChanges(
startBlock,
endBlock,
blockId,
startEdges,
endEdges,
startLoops,
endLoops,
startParallels,
endParallels
)
for (const blockId in endFlattened) {
if (blockId in startFlattened) {
const { block: startBlock, parentId: startParentId } = startFlattened[blockId]
const { block: endBlock, parentId: endParentId } = endFlattened[blockId]
// Check if parent changed (moved in/out of subflow)
if (startParentId !== endParentId) {
// Extract from old parent if it had one
if (startParentId) {
operations.push({
operation_type: 'extract_from_subflow',
block_id: blockId,
params: { subflowId: startParentId },
})
subflowsChanged++
}
// Insert into new parent if it has one - include full block state
if (endParentId) {
const addParams: EditOperation['params'] = {
subflowId: endParentId,
type: endBlock.type,
name: endBlock.name,
outputs: endBlock.outputs,
enabled: endBlock.enabled !== undefined ? endBlock.enabled : true,
...(endBlock?.triggerMode !== undefined && {
triggerMode: Boolean(endBlock.triggerMode),
}),
...(endBlock?.advancedMode !== undefined && {
advancedMode: Boolean(endBlock.advancedMode),
}),
}
const inputs = extractInputValues(endBlock)
if (Object.keys(inputs).length > 0) {
addParams.inputs = inputs
}
const connections = extractConnections(blockId, endEdges)
if (connections && Object.keys(connections).length > 0) {
addParams.connections = connections
}
operations.push({
operation_type: 'insert_into_subflow',
block_id: blockId,
params: addParams,
})
subflowsChanged++
}
}
// Check for other changes (only if parent didn't change)
const changes = computeBlockChanges(startBlock, endBlock, blockId, startEdges, endEdges)
if (changes) {
operations.push({
operation_type: 'edit',
@@ -140,24 +302,13 @@ export function computeEditSequence(
params: changes,
})
blocksModified++
if (changes.connections || changes.removeEdges) {
if (changes.connections) {
edgesChanged++
}
if (changes.loopConfig || changes.parallelConfig) {
subflowsChanged++
}
}
}
}
// 4. Check for standalone loop/parallel changes (not tied to specific blocks)
const loopChanges = detectSubflowChanges(startLoops, endLoops, 'loop')
const parallelChanges = detectSubflowChanges(startParallels, endParallels, 'parallel')
if (loopChanges > 0 || parallelChanges > 0) {
subflowsChanged += loopChanges + parallelChanges
}
return {
operations,
summary: {
@@ -171,20 +322,21 @@ export function computeEditSequence(
}
/**
* Extract input values from a block's subBlocks
* Extract input values from a block
* Works with sanitized format where inputs is Record<string, value>
*/
function extractInputValues(block: any): Record<string, any> {
const inputs: Record<string, any> = {}
if (block.subBlocks) {
for (const [subBlockId, subBlock] of Object.entries(block.subBlocks)) {
if ((subBlock as any).value !== undefined && (subBlock as any).value !== null) {
inputs[subBlockId] = (subBlock as any).value
}
}
// New sanitized format uses 'inputs' field
if (block.inputs) {
return { ...block.inputs }
}
return inputs
// Fallback for any legacy data
if (block.subBlocks) {
return { ...block.subBlocks }
}
return {}
}
/**
@@ -233,101 +385,6 @@ function extractConnections(
return connections
}
/**
* Find loop configuration for a block
*/
function findLoopConfigForBlock(
blockId: string,
loops: Record<string, any>
):
| {
nodes?: string[]
iterations?: number
loopType?: 'for' | 'forEach'
forEachItems?: any
}
| undefined {
for (const loop of Object.values(loops)) {
if (loop.id === blockId || loop.nodes?.includes(blockId)) {
return {
nodes: loop.nodes,
iterations: loop.iterations,
loopType: loop.loopType,
forEachItems: loop.forEachItems,
}
}
}
return undefined
}
/**
* Find parallel configuration for a block
*/
function findParallelConfigForBlock(
blockId: string,
parallels: Record<string, any>
):
| {
nodes?: string[]
distribution?: any
count?: number
parallelType?: 'count' | 'collection'
}
| undefined {
for (const parallel of Object.values(parallels)) {
if (parallel.id === blockId || parallel.nodes?.includes(blockId)) {
return {
nodes: parallel.nodes,
distribution: parallel.distribution,
count: parallel.count,
parallelType: parallel.parallelType,
}
}
}
return undefined
}
/**
* Detect changes in subflow configurations
*/
function detectSubflowChanges(
startSubflows: Record<string, any>,
endSubflows: Record<string, any>,
type: 'loop' | 'parallel'
): number {
let changes = 0
// Check for added/removed subflows
const startIds = new Set(Object.keys(startSubflows))
const endIds = new Set(Object.keys(endSubflows))
for (const id of endIds) {
if (!startIds.has(id)) {
changes++ // New subflow
}
}
for (const id of startIds) {
if (!endIds.has(id)) {
changes++ // Removed subflow
}
}
// Check for modified subflows
for (const id of endIds) {
if (startIds.has(id)) {
const startSubflow = startSubflows[id]
const endSubflow = endSubflows[id]
if (JSON.stringify(startSubflow) !== JSON.stringify(endSubflow)) {
changes++ // Modified subflow
}
}
}
return changes
}
/**
* Compute what changed in a block between two states
*/
@@ -346,11 +403,7 @@ function computeBlockChanges(
target: string
sourceHandle?: string | null
targetHandle?: string | null
}>,
startLoops: Record<string, any>,
endLoops: Record<string, any>,
startParallels: Record<string, any>,
endParallels: Record<string, any>
}>
): Record<string, any> | null {
const changes: Record<string, any> = {}
let hasChanges = false
@@ -375,6 +428,14 @@ function computeBlockChanges(
hasChanges = true
}
// Check advanced mode change
const startAdvanced = Boolean(startBlock?.advancedMode)
const endAdvanced = Boolean(endBlock?.advancedMode)
if (startAdvanced !== endAdvanced) {
changes.advancedMode = endAdvanced
hasChanges = true
}
// Check input value changes
const startInputs = extractInputValues(startBlock)
const endInputs = extractInputValues(endBlock)
@@ -389,79 +450,7 @@ function computeBlockChanges(
const endConnections = extractConnections(blockId, endEdges)
if (JSON.stringify(startConnections) !== JSON.stringify(endConnections)) {
// Compute which edges were removed
const removedEdges: Array<{ targetBlockId: string; sourceHandle?: string }> = []
for (const handle in startConnections) {
const startTargets = Array.isArray(startConnections[handle])
? startConnections[handle]
: [startConnections[handle]]
const endTargets = endConnections[handle]
? Array.isArray(endConnections[handle])
? endConnections[handle]
: [endConnections[handle]]
: []
for (const target of startTargets) {
const targetId = typeof target === 'object' ? target.block : target
const isPresent = endTargets.some(
(t: any) => (typeof t === 'object' ? t.block : t) === targetId
)
if (!isPresent) {
removedEdges.push({
targetBlockId: targetId,
sourceHandle: handle !== 'default' ? handle : undefined,
})
}
}
}
if (removedEdges.length > 0) {
changes.removeEdges = removedEdges
}
// Add new connections
if (Object.keys(endConnections).length > 0) {
changes.connections = endConnections
}
hasChanges = true
}
// Check loop membership changes
const startLoopConfig = findLoopConfigForBlock(blockId, startLoops)
const endLoopConfig = findLoopConfigForBlock(blockId, endLoops)
if (JSON.stringify(startLoopConfig) !== JSON.stringify(endLoopConfig)) {
if (endLoopConfig) {
;(changes as any).loopConfig = endLoopConfig
}
hasChanges = true
}
// Check parallel membership changes
const startParallelConfig = findParallelConfigForBlock(blockId, startParallels)
const endParallelConfig = findParallelConfigForBlock(blockId, endParallels)
if (JSON.stringify(startParallelConfig) !== JSON.stringify(endParallelConfig)) {
if (endParallelConfig) {
;(changes as any).parallelConfig = endParallelConfig
}
hasChanges = true
}
// Check parent-child relationship changes
const startParentId = startBlock.data?.parentId
const endParentId = endBlock.data?.parentId
const startExtent = startBlock.data?.extent
const endExtent = endBlock.data?.extent
if (startParentId !== endParentId || startExtent !== endExtent) {
if (endParentId) {
changes.parentId = endParentId
changes.extent = endExtent
}
changes.connections = endConnections
hasChanges = true
}
@@ -478,20 +467,29 @@ export function formatEditSequence(operations: EditOperation[]): string[] {
return `Add block "${op.params?.name || op.block_id}" (${op.params?.type || 'unknown'})`
case 'delete':
return `Delete block "${op.block_id}"`
case 'insert_into_subflow':
return `Insert "${op.params?.name || op.block_id}" into subflow "${op.params?.subflowId}"`
case 'extract_from_subflow':
return `Extract "${op.block_id}" from subflow "${op.params?.subflowId}"`
case 'edit': {
const changes: string[] = []
if (op.params?.type) changes.push(`type to ${op.params.type}`)
if (op.params?.name) changes.push(`name to "${op.params.name}"`)
if (op.params?.inputs) changes.push('inputs')
if (op.params?.triggerMode !== undefined)
changes.push(`trigger mode to ${op.params.triggerMode}`)
if (op.params?.advancedMode !== undefined)
changes.push(`advanced mode to ${op.params.advancedMode}`)
if (op.params?.inputs) {
const inputKeys = Object.keys(op.params.inputs)
if (inputKeys.length > 0) {
changes.push(`inputs (${inputKeys.join(', ')})`)
}
}
if (op.params?.connections) changes.push('connections')
if (op.params?.removeEdges) changes.push(`remove ${op.params.removeEdges.length} edge(s)`)
if ((op.params as any)?.loopConfig) changes.push('loop configuration')
if ((op.params as any)?.parallelConfig) changes.push('parallel configuration')
if (op.params?.parentId) changes.push('parent-child relationship')
return `Edit block "${op.block_id}": ${changes.join(', ')}`
}
default:
return `Unknown operation on block "${op.block_id}"`
return `Unknown operation: ${op.operation_type}`
}
})
}