Added loops to serialized blob sent to executor, loop between evaluator -> agent works

This commit is contained in:
Waleed Latif
2025-02-11 14:53:32 -08:00
parent 5100e2e2fc
commit 3d52dec731
4 changed files with 179 additions and 42 deletions
+2 -2
View File
@@ -11,7 +11,7 @@ import { Serializer } from '@/serializer'
export function useWorkflowExecution() {
const [isExecuting, setIsExecuting] = useState(false)
const [executionResult, setExecutionResult] = useState<ExecutionResult | null>(null)
const { blocks, edges } = useWorkflowStore()
const { blocks, edges, loops } = useWorkflowStore()
const { activeWorkflowId } = useWorkflowRegistry()
const { addNotification } = useNotificationStore()
const { addConsole, toggleConsole, isOpen } = useConsoleStore()
@@ -50,7 +50,7 @@ export function useWorkflowExecution() {
)
// Execute workflow
const workflow = new Serializer().serializeWorkflow(blocks, edges)
const workflow = new Serializer().serializeWorkflow(blocks, edges, loops)
const executor = new Executor(workflow, currentBlockStates, envVarValues)
const result = await executor.execute('my-run-id')
+164 -38
View File
@@ -92,6 +92,13 @@ export class Executor {
*/
private async executeInParallel(context: ExecutionContext): Promise<BlockOutput> {
const { blocks, connections } = this.workflow
const MAX_ITERATIONS = 2 // Add safety limit for loops
// Track iterations per loop
const loopIterations = new Map<string, number>()
for (const [loopId, loop] of Object.entries(this.workflow.loops || {})) {
loopIterations.set(loopId, 0)
}
// Build dependency graphs: inDegree (number of incoming edges) and adjacency (outgoing connections)
const inDegree = new Map<string, number>()
@@ -102,35 +109,34 @@ export class Executor {
adjacency.set(block.id, [])
}
// Set to track which connections are counted in inDegree.
// Set to track which connections are counted in inDegree
const countedEdges = new Set<(typeof connections)[number]>()
// Populate inDegree and adjacency.
// Populate inDegree and adjacency
for (const conn of connections) {
const sourceBlock = blocks.find((b) => b.id === conn.source)
let countEdge = true
if (conn.condition) {
countEdge = false
} else if (sourceBlock && sourceBlock.metadata?.type === 'evaluator') {
// For evaluator edges, count the dependency only if the target block's config references the evaluator output.
// For evaluator edges, count the dependency only if the target block's config references the evaluator output
const targetBlock = blocks.find((b) => b.id === conn.target)
if (targetBlock) {
const paramsStr = JSON.stringify(targetBlock.config.params || {})
// Look for the evaluator block's id or normalized title (lowercase, no spaces) in the template.
const evaluatorIdRef = `<${sourceBlock.id}`
const evaluatorTitleRef = sourceBlock.metadata?.title
? `<${sourceBlock.metadata.title.toLowerCase().replace(/\s+/g, '')}`
: ''
if (
!(
paramsStr.includes(evaluatorIdRef) ||
(evaluatorTitleRef && paramsStr.includes(evaluatorTitleRef))
)
) {
const evaluatorRef = `<${sourceBlock.metadata?.title?.toLowerCase().replace(/\s+/g, '')}`
const altEvaluatorRef = `<${sourceBlock.id}`
// If target block references evaluator output, count the edge
if (paramsStr.includes(evaluatorRef) || paramsStr.includes(altEvaluatorRef)) {
countEdge = true
} else {
// For paths that don't use evaluator output, handle via decisions
countEdge = false
}
}
}
if (countEdge) {
inDegree.set(conn.target, (inDegree.get(conn.target) || 0) + 1)
countedEdges.add(conn)
@@ -138,11 +144,29 @@ export class Executor {
adjacency.get(conn.source)?.push(conn.target)
}
// Maps for router and conditional decisions.
// Function to reset inDegree for blocks in a loop
const resetLoopBlocksDegrees = (loopId: string) => {
const loop = this.workflow.loops?.[loopId]
if (!loop) return
for (const blockId of loop.nodes) {
// For each block in the loop, recalculate its initial inDegree
let degree = 0
for (const conn of connections) {
if (conn.target === blockId && loop.nodes.includes(conn.source)) {
degree++
}
}
inDegree.set(blockId, degree)
}
}
// Maps for decisions
const routerDecisions = new Map<string, string>()
const evaluatorDecisions = new Map<string, string>()
const activeConditionalPaths = new Map<string, string>()
// Initial queue: all blocks with zero inDegree.
// Initial queue: all blocks with zero inDegree
const queue: string[] = []
for (const [blockId, degree] of inDegree) {
if (degree === 0) {
@@ -156,15 +180,22 @@ export class Executor {
const currentLayer = [...queue]
queue.length = 0
// Filtering: only execute blocks that satisfy router and condition decisions.
// Filter executable blocks
const executableBlocks = currentLayer.filter((blockId) => {
const block = blocks.find((b) => b.id === blockId)
if (!block || block.enabled === false) return false
// Check router decisions
for (const [routerId, chosenPath] of routerDecisions) {
if (!this.isInChosenPath(blockId, chosenPath, routerId)) return false
}
// Check evaluator decisions
for (const [evaluatorId, chosenPath] of evaluatorDecisions) {
if (!this.isInChosenPath(blockId, chosenPath, evaluatorId)) return false
}
// Check conditional paths
for (const [conditionBlockId, selectedConditionId] of activeConditionalPaths) {
const connection = connections.find(
(conn) =>
@@ -180,7 +211,7 @@ export class Executor {
return true
})
// Execute all blocks in the current layer in parallel.
// Execute all blocks in the current layer in parallel
const layerResults = await Promise.all(
executableBlocks.map(async (blockId) => {
const block = blocks.find((b) => b.id === blockId)
@@ -201,6 +232,16 @@ export class Executor {
}
}
routerDecisions.set(block.id, routerResult.response.selectedPath.blockId)
} else if (block.metadata?.type === 'evaluator') {
const evaluatorResult = result as {
response: {
content: string
model: string
tokens: { prompt: number; completion: number; total: number }
selectedPath: { blockId: string }
}
}
evaluatorDecisions.set(block.id, evaluatorResult.response.selectedPath.blockId)
} else if (block.metadata?.type === 'condition') {
const conditionResult = result as {
response: {
@@ -219,17 +260,48 @@ export class Executor {
})
)
// After executing a layer, update inDegree for all adjacent blocks using only counted edges.
// Process outgoing connections and update queue
for (const finishedBlockId of layerResults) {
const outgoingConns = connections.filter(
(conn) => conn.source === finishedBlockId && countedEdges.has(conn)
)
const outgoingConns = connections.filter((conn) => conn.source === finishedBlockId)
for (const conn of outgoingConns) {
if (!conn.sourceHandle || !conn.sourceHandle.startsWith('condition-')) {
const sourceBlock = blocks.find((b) => b.id === conn.source)
if (sourceBlock?.metadata?.type === 'evaluator') {
// Only add to queue if this is the chosen path
const chosenPath = evaluatorDecisions.get(sourceBlock.id)
if (conn.target === chosenPath) {
// CHANGED: Don't rely on inDegree for loop targets
const targetBlock = blocks.find((b) => b.id === conn.target)
const isInLoop = Object.values(this.workflow.loops || {}).some((loop) =>
loop.nodes.includes(conn.target)
)
if (isInLoop) {
// If target is in a loop, queue it directly
queue.push(conn.target)
} else {
// For non-loop targets, use normal inDegree logic
const newDegree = (inDegree.get(conn.target) || 0) - 1
inDegree.set(conn.target, newDegree)
if (newDegree === 0) queue.push(conn.target)
}
}
} else if (sourceBlock?.metadata?.type === 'router') {
// Only add to queue if this is the chosen path
const chosenPath = routerDecisions.get(sourceBlock.id)
if (conn.target === chosenPath) {
const newDegree = (inDegree.get(conn.target) || 0) - 1
inDegree.set(conn.target, newDegree)
if (newDegree === 0) queue.push(conn.target)
}
} else if (!conn.sourceHandle?.startsWith('condition-')) {
// Normal connection
const newDegree = (inDegree.get(conn.target) || 0) - 1
inDegree.set(conn.target, newDegree)
if (newDegree === 0) queue.push(conn.target)
} else {
// Condition connection
const conditionId = conn.sourceHandle.replace('condition-', '')
if (activeConditionalPaths.get(finishedBlockId) === conditionId) {
const newDegree = (inDegree.get(conn.target) || 0) - 1
@@ -239,6 +311,37 @@ export class Executor {
}
}
}
// Check if we need to reset any loops
for (const [loopId, loop] of Object.entries(this.workflow.loops || {})) {
const loopBlocks = new Set(loop.nodes)
const executedLoopBlocks = layerResults.filter((blockId) => loopBlocks.has(blockId))
if (executedLoopBlocks.length > 0) {
const iterations = loopIterations.get(loopId) || 0
if (iterations < MAX_ITERATIONS) {
// Check if the evaluator chose a block within the loop
const evaluatorInLoop = executedLoopBlocks.find((blockId) => {
const block = blocks.find((b) => b.id === blockId)
return block?.metadata?.type === 'evaluator'
})
if (evaluatorInLoop) {
const chosenPath = evaluatorDecisions.get(evaluatorInLoop)
if (chosenPath && loopBlocks.has(chosenPath)) {
// Reset the loop blocks' inDegrees and add them back to queue if needed
resetLoopBlocksDegrees(loopId)
for (const blockId of loop.nodes) {
if (inDegree.get(blockId) === 0) {
queue.push(blockId)
}
}
loopIterations.set(loopId, iterations + 1)
}
}
}
}
}
}
return lastOutput
@@ -691,11 +794,10 @@ export class Executor {
const resolvedInputs = this.resolveInputs(block, context)
console.log('Evaluator: Resolved inputs:', resolvedInputs)
console.log('Evaluator: Filtering outgoing connections for the block.')
// Get all possible target blocks from outgoing connections
const outgoingConnections = this.workflow.connections.filter((conn) => conn.source === block.id)
console.log('Evaluator: Outgoing connections:', outgoingConnections)
console.log('Evaluator: Mapping target blocks from outgoing connections.')
const targetBlocks = outgoingConnections.map((conn) => {
const targetBlock = this.workflow.blocks.find((b) => b.id === conn.target)
if (!targetBlock) {
@@ -724,10 +826,8 @@ export class Executor {
const model = evaluatorConfig.model || 'gpt-4o'
const providerId = getProviderFromModel(model)
// Add logging before making the request
console.log('Evaluator: Sending request with prompt:', evaluatorConfig.prompt)
console.log('Evaluator: Available target blocks:', targetBlocks)
// Generate and execute the evaluator prompt
console.log('Evaluator: Sending request with config:', evaluatorConfig)
const response = await executeProviderRequest(providerId, {
model: evaluatorConfig.model,
systemPrompt: generateEvaluatorPrompt(
@@ -740,19 +840,19 @@ export class Executor {
apiKey: evaluatorConfig.apiKey,
})
// Add logging after getting the response
console.log('Evaluator: Raw response:', response)
console.log('Evaluator: Chosen block ID:', response.content.trim().toLowerCase())
const chosenBlockId = response.content.trim().toLowerCase()
console.log('Evaluator: Chosen block ID:', chosenBlockId)
const chosenBlock = targetBlocks.find((b) => b.id === chosenBlockId)
if (!chosenBlock) {
throw new Error(`Invalid routing decision: ${chosenBlockId}`)
throw new Error(`Invalid evaluation decision: ${chosenBlockId}`)
}
// Store the evaluation result in the context
const tokens = response.tokens || { prompt: 0, completion: 0, total: 0 }
return {
content: resolvedInputs.prompt,
const result = {
content: evaluatorConfig.prompt,
model: response.model,
tokens: {
prompt: tokens.prompt || 0,
@@ -765,6 +865,13 @@ export class Executor {
blockTitle: chosenBlock.title || 'Untitled Block',
},
}
// ADDED: Explicitly store the evaluation decision in the context
context.blockStates.set(block.id, {
response: result,
})
return result
}
/**
@@ -772,21 +879,40 @@ export class Executor {
*
* This uses a breadth-first search starting from the chosen block id.
*/
private isInChosenPath(blockId: string, chosenBlockId: string, routerId: string): boolean {
private isInChosenPath(blockId: string, chosenBlockId: string, decisionBlockId: string): boolean {
const visited = new Set<string>()
const queue = [chosenBlockId]
// Add the decision block (router/evaluator) itself as valid
if (blockId === decisionBlockId) {
return true
}
while (queue.length > 0) {
const currentId = queue.shift()!
if (visited.has(currentId)) continue
visited.add(currentId)
// If we found the block we're looking for
if (currentId === blockId) {
return true
}
// Get all outgoing connections from current block
const connections = this.workflow.connections.filter((conn) => conn.source === currentId)
for (const conn of connections) {
queue.push(conn.target)
// Don't follow connections from other routers/evaluators
const sourceBlock = this.workflow.blocks.find((b) => b.id === conn.source)
if (
sourceBlock?.metadata?.type !== 'router' &&
sourceBlock?.metadata?.type !== 'evaluator'
) {
queue.push(conn.target)
}
}
}
return blockId === routerId || visited.has(blockId)
return false
}
/**
+7 -2
View File
@@ -1,10 +1,14 @@
import { Edge } from 'reactflow'
import { BlockState, SubBlockState } from '@/stores/workflow/types'
import { BlockState, Loop, SubBlockState } from '@/stores/workflow/types'
import { getBlock } from '@/blocks'
import { SerializedBlock, SerializedConnection, SerializedWorkflow } from './types'
export class Serializer {
serializeWorkflow(blocks: Record<string, BlockState>, edges: Edge[]): SerializedWorkflow {
serializeWorkflow(
blocks: Record<string, BlockState>,
edges: Edge[],
loops: Record<string, Loop>
): SerializedWorkflow {
return {
version: '1.0',
blocks: Object.values(blocks).map((block) => this.serializeBlock(block)),
@@ -14,6 +18,7 @@ export class Serializer {
sourceHandle: edge.sourceHandle || undefined,
targetHandle: edge.targetHandle || undefined,
})),
loops,
}
}
+6
View File
@@ -5,6 +5,7 @@ export interface SerializedWorkflow {
version: string
blocks: SerializedBlock[]
connections: SerializedConnection[]
loops: Record<string, SerializedLoop>
}
export interface SerializedConnection {
@@ -37,3 +38,8 @@ export interface SerializedBlock {
}
enabled: boolean
}
export interface SerializedLoop {
id: string
nodes: string[]
}