How the Sim Studio Workflow Executor Handles Parallel Blocks and Concurrent Execution
The Sim Studio workflow executor expands parallel blocks into independent cloned sub-flows, schedules branch entry nodes via a dynamic pending queue for true concurrent execution, and aggregates results using an event-driven completion tracker.
The simstudioai/sim repository implements a sophisticated concurrency model that transforms high-level parallel blocks into independently executing branches. When a workflow definition contains a type: "parallel" block, the executor clones the entire sub-flow for each distribution item and manages concurrent execution through dynamic node scheduling. Understanding how this system handles parallel blocks and concurrent execution is essential for building high-performance workflows that maximize resource utilization while maintaining deterministic output aggregation.
DAG Construction and Parallel Expansion
Registering Parallel Configurations
During the DAG construction phase, the builder in apps/sim/executor/dag/construction/parallels.ts identifies parallel blocks and registers a SerializedParallel configuration in dag.parallelConfigs. This registration captures the distribution items and sub-flow structure before execution begins, establishing the foundation for branch cloning.
Cloning Sub-flows with ParallelExpander
The heavy lifting of parallel expansion happens in apps/sim/executor/utils/parallel-expansion.ts. When ParallelOrchestrator.initializeParallelScope invokes the ParallelExpander, it clones the original sub-flow for every branch in the distribution array.
// Inside ParallelOrchestrator.initializeParallelScope
const { entryNodes, clonedSubflows } = this.expander.expandParallel(
this.dag,
parallelId,
branchCount,
items
);
// entryNodes now contain IDs like "log-item__branch-1", "log-item__branch-2", …
Each cloned node receives a unique ID suffixed with __branch-N, and the system records parent-child relationships in ctx.subflowParentMap to ensure proper variable resolution across branches.
Concurrent Branch Execution
The ParallelScope State Manager
The orchestrator creates a ParallelScope object that tracks execution state across all branches. According to the source code in apps/sim/executor/orchestrators/parallel.ts, this scope maintains a Map<number, NormalizedBlockOutput[]> for branch outputs, the total branch count, and distribution items. This structure isolates branch data while providing a unified interface for completion tracking.
Dynamic Node Scheduling
True concurrent execution emerges through the executor's main loop in apps/sim/executor/execution/executor.ts. Entry nodes belonging to branches—all except __branch-0 entry points—are added to ctx.pendingDynamicNodes. The executor processes this queue in FIFO order, allowing multiple branches to run simultaneously without implicit ordering constraints.
Event-Driven Completion Tracking
As nodes complete, ParallelOrchestrator.handleParallelBranchCompletion records outputs in the scope's branchOutputs map. This system uses event-driven completion rather than a simple counter, extracting the branch index from the node ID suffix (extractBranchIndex) to handle conditional paths and early exits without breaking completion detection.
// Called by the block executor when a node finishes
orchestrator.handleParallelBranchCompletion(ctx, parallelId, nodeId, output);
Result Aggregation and Completion
Sentinel Node Triggering
A sentinel end-node—generated during DAG construction—triggers the final aggregation phase. When executed, it calls ParallelOrchestrator.aggregateParallelResults, which walks through scope.branchOutputs and assembles a two-dimensional array results[branch][node].
const agg = await orchestrator.aggregateParallelResults(ctx, parallelId);
// agg.results => [
// [{ result: 'apple' }], // branch 0
// [{ result: 'banana' }], // branch 1
// [{ result: 'cherry' }] // branch 2
// ]
The final output is stored via state.setBlockOutput(parallelId, { results }) and emitted to downstream blocks.
Validation and Error Handling
The orchestrator validates distribution arrays before execution. If a parallel block has an empty distribution or exceeds DEFAULTS.MAX_PARALLEL_BRANCHES, the system logs an error, creates a minimal ParallelScope with a validationError property, and short-circuits execution to prevent resource exhaustion.
Example Parallel Block Configuration
Here's how a parallel block appears in workflow definitions:
{
"type": "workflow",
"blocks": [
{
"id": "parallel-1",
"type": "parallel",
"distribution": ["apple", "banana", "cherry"],
"subBlocks": [
{
"id": "log-item",
"type": "tool",
"tool": "log",
"params": {
"message": "<{{item}}>"
}
}
]
}
]
}
The executor reads the distribution array, creates three branches, and clones the log-item block for each fruit, executing them concurrently.
Summary
- The SerializedParallel configuration in
dag.parallelConfigsregisters parallel blocks during DAG construction inapps/sim/executor/dag/construction/parallels.ts. - ParallelExpander (
apps/sim/executor/utils/parallel-expansion.ts) clones sub-flows for each branch, assigning__branch-Nsuffixes to node IDs and recording parent mappings inctx.subflowParentMap. - Branch entry nodes enter
ctx.pendingDynamicNodes, enabling the main executor loop inapps/sim/executor/execution/executor.tsto process them concurrently without ordering constraints. - Event-driven completion tracking via
handleParallelBranchCompletionstores outputs inParallelScope.branchOutputs, supporting conditional paths and early exits without race conditions. - A sentinel node triggers
aggregateParallelResultsto merge branch outputs into a unified two-dimensional array stored viastate.setBlockOutput. - Built-in validation prevents execution of empty distributions or branches exceeding
DEFAULTS.MAX_PARALLEL_BRANCHES.
Frequently Asked Questions
How does the executor prevent race conditions between parallel branches?
Each branch operates on cloned nodes with unique IDs suffixed by __branch-N, and outputs are isolated in a Map<number, NormalizedBlockOutput[]> within the ParallelScope. The event-driven completion handler extracts branch indices from node IDs to ensure outputs are recorded in the correct bucket without cross-branch contamination.
What happens if one branch fails while others are still running?
The event-driven completion tracking system records each branch's output independently as nodes finish. If a branch encounters an error, it stores the error state in its specific branch output array. The sentinel node still triggers aggregation once all branches complete, and downstream blocks receive the complete results matrix including error states for proper handling.
Is there a limit to how many branches can run concurrently?
Yes. The executor checks against DEFAULTS.MAX_PARALLEL_BRANCHES during initialization. If the distribution array exceeds this limit, ParallelOrchestrator.initializeParallelScope creates a ParallelScope with a validationError and short-circuits execution before any branches spawn, preventing resource exhaustion.
How does variable resolution work across cloned branches?
The ParallelExpander records parent-child relationships in ctx.subflowParentMap during the cloning process. When resolving variables, the system traverses this mapping to locate the correct parent context, ensuring that each branch accesses the appropriate execution state despite having cloned node IDs.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →