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.