core.async_executor
Async DAG executor for topological-parallel prompt execution.
Runs prompts level-by-level using asyncio.gather per level. Prompts
on the same level execute concurrently; levels execute sequentially.
Async DAG executor for topological-parallel prompt execution.
Runs prompts level-by-level using asyncio.gather per level. Prompts
on the same level execute concurrently; levels execute sequentially.
- class GraphResult(results=<factory>, success_count=0, failed_count=0, skipped_count=0, aborted=False, aborted_count=0)[source]
Bases:
objectAggregate result from executing a prompt dependency graph.
- Parameters:
- results
Mapping of prompt_name to
ResponseResult.
- results: dict[str, ResponseResult]
- class AsyncGraphExecutor(executor_fn, max_concurrency=10, prompt_resolver=None)[source]
Bases:
objectExecute a prompt DAG with topological-parallel async calls.
- Parameters:
executor_fn (Callable[..., Awaitable[ResponseResult]]) – Async callable that takes prompt kwargs and returns a
ResponseResult.max_concurrency (int) – Maximum number of concurrent API calls, enforced by an
asyncio.Semaphore.prompt_resolver (Callable[[dict[str, Any], dict[str, dict[str, Any]]], tuple[str, set[str]]] | None) – Optional callback
(prompt_spec, results_by_name) -> (resolved_prompt, interpolated_names). When provided, each node’s prompt is resolved before execution. WhenNone, prompts are sent as-is.
- async execute(prompts)[source]
Build and execute a prompt dependency graph.
- Parameters:
prompts (list[dict[str, Any]]) – List of prompt dicts with keys: prompt_name, prompt, history, condition, abort_condition.
- Returns:
GraphResultwith per-prompt results and aggregate counts.- Raises:
ValueError – If a dependency cycle is detected.
- Return type:
Classes
- class AsyncGraphExecutor(executor_fn, max_concurrency=10, prompt_resolver=None)[source]
Bases:
objectExecute a prompt DAG with topological-parallel async calls.
- Parameters:
executor_fn (Callable[..., Awaitable[ResponseResult]]) – Async callable that takes prompt kwargs and returns a
ResponseResult.max_concurrency (int) – Maximum number of concurrent API calls, enforced by an
asyncio.Semaphore.prompt_resolver (Callable[[dict[str, Any], dict[str, dict[str, Any]]], tuple[str, set[str]]] | None) – Optional callback
(prompt_spec, results_by_name) -> (resolved_prompt, interpolated_names). When provided, each node’s prompt is resolved before execution. WhenNone, prompts are sent as-is.
- async execute(prompts)[source]
Build and execute a prompt dependency graph.
- Parameters:
prompts (list[dict[str, Any]]) – List of prompt dicts with keys: prompt_name, prompt, history, condition, abort_condition.
- Returns:
GraphResultwith per-prompt results and aggregate counts.- Raises:
ValueError – If a dependency cycle is detected.
- Return type: