Implementing Parallel vs Sequential Execution Strategies in ChatDev Workflows

ChatDev implements parallel vs sequential execution through a strategy pattern where DagExecutionStrategy, CycleExecutionStrategy, and MajorityVoteStrategy delegate to a centralized ParallelExecutor that automatically threads non-blocking nodes while serializing blocking operations like human interactions.

The ChatDev framework by OpenBMB provides a sophisticated workflow execution engine that dynamically selects concurrency models based on graph topology. Understanding how the runtime switches between parallel and sequential modes is essential for optimizing multi-agent workflows and ensuring correct handling of interactive steps.

Execution Strategy Selection

The entry point for execution mode decisions resides in workflow/runtime/execution_strategy.py, which defines three concrete strategies. Each strategy instantiates a specialized executor that orchestrates how nodes traverse the graph while delegating actual concurrency decisions to the ParallelExecutor.

When the workflow graph is a Directed Acyclic Graph (DAG), the runtime selects DagExecutionStrategy. This strategy instantiates a DAGExecutor from workflow/executor/dag_executor.py that walks the graph layer-by-layer, handing each topological layer to the parallel executor for concurrent processing where possible.

If the graph contains cycles, CycleExecutionStrategy takes over. This strategy creates a CycleExecutor from workflow/executor/cycle_executor.py that builds "super-node" layers to handle nested cycles. It repeatedly executes scoped sub-graphs until exit conditions or maximum iteration limits are reached, while still leveraging parallel execution within each iteration.

For majority-vote configurations (graphs with no edges), MajorityVoteStrategy executes every node in parallel via the ParallelExecutor, then aggregates results to select the most common output.

All three strategies receive identical runtime objects: a structured log_manager, an immutable {node_id: Node} map defined in entity/configs.py, and an execute_node_func callable that implements the actual business logic (LLM calls, tool invocations, or human steps).

How the ParallelExecutor Decides Concurrency

The ParallelExecutor class in workflow/executor/parallel_executor.py serves as the single source of truth for parallel vs sequential decisions. This design centralizes concurrency logic, preventing race conditions while maximizing throughput for I/O-bound operations.

The executor distinguishes between two execution paths using an optional has_blocking_func parameter:

  • Parallel path: When has_blocking_func returns False, the executor creates a ThreadPoolExecutor with max_workers=len(items) and submits each executor_func(item) as an independent thread. This maximizes CPU utilization for LLM calls and tool invocations.

  • Sequential path: Items classified as blocking (typically Human nodes requiring user interaction) are collected in blocking_items and processed one-by-one via _execute_sequential_batch. This prevents multiple interactive prompts from interleaving or creating race conditions on shared input streams.

The classification function typically checks node types defined in entity/configs.py:

def has_blocking_func(node_id: str) -> bool:
    # Human nodes must run sequentially

    return self.nodes_dict[node_id].node_type == "Human"

Both DAGExecutor and CycleExecutor invoke ParallelExecutor.execute_items_parallel (or the convenience wrapper execute_nodes_parallel) with an appropriate has_blocking_func. This architecture allows each layer to contain a mixture of parallel and serialized work without additional orchestration code.

Layer-by-Layer Execution Patterns

Understanding the specific traversal mechanics reveals why the hybrid execution model scales effectively for complex workflow topologies.

DAG Execution

In workflow/executor/dag_executor.py, the DAGExecutor._execute_layer method iterates over node IDs in a topological layer. It first checks node.is_triggered() to determine if a node should run, then hands the filtered list to ParallelExecutor.execute_nodes_parallel. This ensures that within a single layer, independent nodes execute concurrently while maintaining dependency order between layers.

Cycle Execution

The CycleExecutor in workflow/executor/cycle_executor.py handles more complex topologies through _execute_super_layer_parallel. This method constructs lists of "super-items" (individual nodes or nested cycles) and processes them concurrently.

After each iteration, _is_initial_node_retriggered checks whether the entry node was re-triggered by an internal edge. If re-triggered, the executor rebuilds the scoped sub-graph and launches another parallel execution round. Within each cycle iteration, the blocking predicate still applies, ensuring Human nodes never run in parallel even when the surrounding cycle logic executes concurrently.

Practical Implementation Examples

Running a DAG Workflow with Parallel Layers

To execute a DAG workflow where each topological layer runs in parallel (except blocking nodes), instantiate the strategy with your node definitions and execution function:

from workflow.runtime.execution_strategy import DagExecutionStrategy
from utils.log_manager import LogManager

log = LogManager(name="my_workflow")
nodes = load_nodes()          # {node_id: Node}

layers = topological_layers() # List[List[node_id]]

def execute_node(node):
    # Your node-specific logic here (LLM calls, etc.)

    node.run()

dag_strategy = DagExecutionStrategy(
    log_manager=log,
    nodes=nodes,
    layers=layers,
    execute_node_func=execute_node,
)

dag_strategy.run()   # Executes each layer in parallel where safe

This strategy automatically utilizes DAGExecutor which internally calls ParallelExecutor.execute_nodes_parallel according to the implementation in workflow/executor/parallel_executor.py.

Executing Cyclic Workflows with Nested Cycles

For workflows containing loops or iterative refinement cycles, use the cycle-aware strategy:

from workflow.runtime.execution_strategy import CycleExecutionStrategy
from workflow.cycle_manager import CycleManager
from utils.log_manager import LogManager

log = LogManager(name="cycle_demo")
nodes = load_nodes()
cycle_order = build_cycle_execution_order()  # From GraphTopologyBuilder

cycle_mgr = CycleManager(nodes)

def execute_node(node):
    node.run()

cycle_strategy = CycleExecutionStrategy(
    log_manager=log,
    nodes=nodes,
    cycle_execution_order=cycle_order,
    cycle_manager=cycle_mgr,
    execute_node_func=execute_node,
)

cycle_strategy.run()   # Handles parallel execution inside cycles

The CycleExecutor._execute_super_layer_parallel method delegates to ParallelExecutor.execute_items_parallel, allowing concurrent execution within cycle iterations while respecting entry-node retriggering rules defined in workflow/executor/cycle_executor.py.

Forcing Sequential Execution for Specific Node Types

To ensure certain node types (like Human interaction nodes) always run sequentially even when grouped with parallelizable tasks, implement a custom blocking function:

from workflow.executor.parallel_executor import ParallelExecutor
from utils.log_manager import LogManager

log = LogManager(name="forced_seq")
nodes_dict = load_nodes()

def exec_func(node_id):
    # Execute the node logic

    nodes_dict[node_id].run()

def has_blocking(node_id: str) -> bool:
    # Force sequential execution for Human nodes

    return nodes_dict[node_id].node_type == "Human"

par_exec = ParallelExecutor(log_manager=log, nodes_dict=nodes_dict)
node_ids = ["agent_1", "agent_2", "human_review", "agent_3"]

par_exec.execute_items_parallel(
    items=node_ids,
    executor_func=exec_func,
    item_desc_func=lambda nid: f"node {nid}",
    has_blocking_func=has_blocking,
)

In this configuration, agent_1 and agent_2 execute in parallel, but human_review waits for both to complete and blocks agent_3 until the human interaction finishes.

Summary

  • Strategy Pattern: ChatDev uses DagExecutionStrategy, CycleExecutionStrategy, and MajorityVoteStrategy in workflow/runtime/execution_strategy.py to match execution algorithms to graph topology.
  • Centralized Concurrency: The ParallelExecutor class in workflow/executor/parallel_executor.py centralizes all threading decisions, using ThreadPoolExecutor for parallel batches and sequential loops for blocking items.
  • Blocking Detection: The optional has_blocking_func parameter allows fine-grained control over which nodes must serialize (typically Human nodes) while permitting LLM and tool nodes to run concurrently.
  • Layer-wise Processing: DAG executors process topological layers in parallel where possible, while cycle executors handle iterative workflows through super-node layers that respect the same concurrency rules.
  • Extensibility: New execution modes require only a new strategy class utilizing the existing ParallelExecutor API, maintaining separation between graph traversal logic and concurrency implementation.

Frequently Asked Questions

How does ChatDev handle parallel execution without race conditions?

ChatDev prevents race conditions by isolating shared state in the nodes map (immutable node objects from entity/configs.py) and executing the actual business logic (execute_node_func) within thread-local contexts. The ParallelExecutor manages the ThreadPoolExecutor lifecycle and ensures that blocking operations—which typically involve shared resources like stdin/stdout or human attention—are automatically serialized into the sequential batch processing queue.

Can I mix parallel and sequential nodes in the same workflow layer?

Yes. The ParallelExecutor.execute_items_parallel method accepts a has_blocking_func predicate that filters items into parallel and sequential batches within the same layer. When processing a mixed layer, non-blocking items execute concurrently via threading, while blocking items wait for the parallel batch to complete before running sequentially. This allows Human nodes to coexist with Agent nodes in the same topological layer without manual orchestration.

What is the difference between CycleExecutionStrategy and DagExecutionStrategy?

DagExecutionStrategy handles Directed Acyclic Graphs using DAGExecutor, which processes nodes in strict topological order—each layer completes before the next begins, with parallelization occurring horizontally within layers. CycleExecutionStrategy handles graphs containing loops using CycleExecutor, which detects cycles, creates "super-nodes" for nested cycles, and executes iterative rounds until exit conditions are met. Both strategies delegate the actual threading decisions to ParallelExecutor, but the cycle strategy adds iteration management and scoped sub-graph rebuilding logic.

Where are node types defined that determine blocking behavior?

Node types and their configurations are defined in entity/configs.py, which contains the Node class specifying attributes like node_type, predecessors, successors, and trigger conditions. The has_blocking_func check typically references nodes_dict[node_id].node_type to identify "Human" nodes or other types requiring sequential execution, allowing the ParallelExecutor to automatically route these to the sequential processing path regardless of the high-level strategy in use.

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:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →