Skip to content

latent.flows.agent_inference_flow

agent_inference_flow — Run any BaseAgent on a dataset and collect results.

Uses Prefect @task + .map() for concurrent execution.

Functions

agent_inference_flow

agent_inference_flow(eval_data: pd.DataFrame, agent_factory: Callable[[], Any] | None = None, agent_factory_for_row: Callable[[dict[str, Any]], Any] | None = None, question_column: str = 'question', context_columns: list[str] | None = None, context_formatter: Callable[[dict[str, Any]], str] | None = None, id_column: str | None = None, concurrency: int = 1, include_events: bool = False, per_row_timeout_s: float | None = None) -> dict[str, Any]

Run an agent on each row of eval_data and collect inference results.

Runs each row concurrently via task.map() (asyncio.gather bounded by an asyncio.Semaphore of size concurrency).

Args: eval_data: DataFrame with at least a question column. agent_factory: Callable returning a BaseAgent instance. Use this when every row gets the same kind of agent. agent_factory_for_row: Callable(row_dict) -> BaseAgent, for when the agent depends on the row — a per-channel context, a persona, a conversation history. Mutually exclusive with agent_factory. Two parameters rather than one that sniffs arity: a factory with an optional parameter is genuinely ambiguous, and guessing wrong silently feeds it a row dict. question_column: Column containing the question/prompt. context_columns: Optional columns to include as context. context_formatter: Callable(row_dict) -> str for context formatting. id_column: Column to use as row ID. If None, uses DataFrame index. concurrency: Number of parallel workers via Prefect task.map(). include_events: If True, include raw agent events in results.

Returns: Dict with inference_results and summary.