dabla opened a new pull request, #74201: URL: https://github.com/apache/airflow/pull/74201
`@task.agent` cannot be used with an `async def` function today. The decorator calls the function without awaiting it, so the task fails in `validate_prompt` on the coroutine it gets back, while the operator already reports `is_async` through `DecoratedOperator`. An agent task therefore cannot await async hooks to build its prompt, and cannot share the event loop of an async task with other work, although the pydantic-ai agent behind it is async-native. This makes `AgentOperator` a `BaseAsyncOperator`, the way `PythonOperator` is one: - `@task.agent` on an `async def` function awaits it for the prompt and runs the agent on the task's event loop. A regular function keeps the synchronous path unchanged. `AgentOperator.is_async` is `False`; on the decorator it follows the function. - `AgentOperator.aexecute()` makes no blocking call to the supervisor on the path of a successful run: `PydanticAIHook.aget_hook()` fetches the connection once and primes the hook with it, the agent is built with `acreate_agent()` and run with `await agent.run()` (`CancellableAgentRunMixin.run_agent_async`, same cancellation token as the sync run), and the usage budget (`aload` / `asave` / `aclear`), the tool approval transcript and the XComs go through their async accessors. - The two rare branches reuse the synchronous code from a worker thread: the usage report of a failed run and the pause for a tool approval. Resuming after an approval stays synchronous. - The template fields are rendered in a worker thread, since a template can read a Variable or a Connection with a blocking call. - `durable=True` and `enable_hitl_review=True` do blocking I/O during the run (the durable storage inside the model and tool calls, the review poll loop) and are rejected with an `async def` function, at parse time on the decorator. They can follow in a later PR. - `axcom_push` exists since Airflow 3.3.0; on 3.2 the XComs are pushed with `xcom_push` from a worker thread. Depends on #74150 (`aget_conn` / `acreate_agent` on `PydanticAIHook`): this branch is built on top of it, so the diff shows its commits until it is merged. **What it changes for performance.** One agent run is not faster, its duration is the model's. The gain is in running many of them in one task process. Measured with pydantic-ai alone and a fake model that sleeps 1 s, one run per cell: a thread pool calling `run_sync` versus one event loop awaiting `run`. | Concurrent runs | Thread pool + `run_sync` | One loop + `run` | CPU, threads | CPU, loop | |---|---|---|---|---| | 8 | 1.06 s | 1.03 s | 0.07 s | 0.02 s | | 32 | 1.26 s | 1.10 s | 0.30 s | 0.09 s | | 64 | 1.41 s | 1.19 s | 0.51 s | 0.18 s | | 256 | 2.97 s | 1.82 s | 2.33 s | 0.78 s | **Tests.** New tests cover `run_agent_async` (token held and cleared, `on_kill` cancelling an awaited run), `aget_hook`, the async usage budget methods, `aexecute` with mocks and with a real `FunctionModel` run where every blocking state store call and `xcom_push` fails the test, the failed-run and tool-approval branches, the rejected features, the decorator with an `async def` function through both `aexecute` and `execute`, rendering off the event loop, and eight runs overlapping on one loop. Run locally: the whole `providers/common/ai/tests/unit` suite (3104 passed, 13 skipped), `mypy` on the five changed source files, and `prek` on the changed files. A backport of the same design to common.ai 0.9 also ran on a Celery worker with Airflow 3.3.2: one async `@task.agent` task, and eight agent runs sharing one event loop inside a single task, against an Azure OpenAI model, without a `DeadlockImminentError`. --- ##### Was generative AI tooling used to co-author this PR? - [X] Yes — Claude Code (Fable 5.1) Generated-by: Claude Code (Fable 5.1) following [the guidelines](https://github.com/apache/airflow/blob/main/contributing-docs/05_pull_requests.rst#gen-ai-assisted-contributions) -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
