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]

Reply via email to