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: object

Aggregate result from executing a prompt dependency graph.

Parameters:
results

Mapping of prompt_name to ResponseResult.

Type:

dict[str, ffai.core.response_result.ResponseResult]

success_count

Number of prompts that succeeded.

Type:

int

failed_count

Number of prompts that failed.

Type:

int

skipped_count

Number of prompts skipped by conditions.

Type:

int

aborted

Whether execution was aborted.

Type:

bool

aborted_count

Number of prompts skipped due to abort.

Type:

int

results: dict[str, ResponseResult]
success_count: int = 0
failed_count: int = 0
skipped_count: int = 0
aborted: bool = False
aborted_count: int = 0
class AsyncGraphExecutor(executor_fn, max_concurrency=10, prompt_resolver=None)[source]

Bases: object

Execute 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. When None, 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:

GraphResult with per-prompt results and aggregate counts.

Raises:

ValueError – If a dependency cycle is detected.

Return type:

GraphResult

Classes

class AsyncGraphExecutor(executor_fn, max_concurrency=10, prompt_resolver=None)[source]

Bases: object

Execute 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. When None, 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:

GraphResult with per-prompt results and aggregate counts.

Raises:

ValueError – If a dependency cycle is detected.

Return type:

GraphResult

class GraphResult(results=<factory>, success_count=0, failed_count=0, skipped_count=0, aborted=False, aborted_count=0)[source]

Bases: object

Aggregate result from executing a prompt dependency graph.

Parameters:
results

Mapping of prompt_name to ResponseResult.

Type:

dict[str, ffai.core.response_result.ResponseResult]

success_count

Number of prompts that succeeded.

Type:

int

failed_count

Number of prompts that failed.

Type:

int

skipped_count

Number of prompts skipped by conditions.

Type:

int

aborted

Whether execution was aborted.

Type:

bool

aborted_count

Number of prompts skipped due to abort.

Type:

int

results: dict[str, ResponseResult]
success_count: int = 0
failed_count: int = 0
skipped_count: int = 0
aborted: bool = False
aborted_count: int = 0