shahar1 opened a new pull request, #74124:
URL: https://github.com/apache/airflow/pull/74124

   `TestMessageQueueTrigger::test_trigger_run_good` failed in CI on an 
unrelated PR with:
   
   ```
   assert task.done() is True
   E   AssertionError: assert False is True
   E    +  where False = <Task pending ... wait_for=<Future pending 
cb=[shield.<locals>._outer_done_callback() ...]>>.done
   ```
   
   The Kafka trigger tests started `trigger.run().__anext__()` as a task, slept 
a fixed 1 second, and asserted the task was done. `AwaitMessageTrigger.run` 
crosses several `sync_to_async` hops (`get_consumer`, `poll`, `apply_function`, 
`commit`) on asgiref's single shared thread. On a loaded xdist worker those 
hops can take longer than 1 second, so the task is still pending on the 
`asyncio.shield` future and the assertion fails. Passes locally in isolation 
and in the full provider suite; the failure only shows under CI load.
   
   This replaces the fixed sleep with 
`asyncio.wait_for(trigger.run().__anext__(), timeout=10)` and asserts on the 
yielded `TriggerEvent` payload, following the pattern already used in 
`providers/amazon/tests/unit/amazon/aws/triggers/test_s3.py`. The tests now 
return as soon as the event is yielded, and the timeout only matters when the 
trigger is actually broken.
   
   The tombstone test is rewritten in the same way: `poll` yields a tombstone 
followed by a regular message, and the test asserts the yielded payload and 
both commit calls. The previous `task.done() is False` check after a fixed 
sleep could only pass by luck under load and could not reliably catch a trigger 
that yields an event for a tombstone.
   
   The negative `test_trigger_run_bad` tests keep their fixed sleep: a sleep 
followed by `done() is False` cannot fail wrongly under load.
   
   Test-only change in the Kafka provider, no changelog entry.
   
   Verified locally: `uv run --project providers/apache/kafka pytest 
providers/apache/kafka/tests/unit/apache/kafka/triggers/` (52 passed) and `prek 
run --from-ref upstream/main --stage pre-commit` (no failures). With 
`MockedConsumer.poll` slowed to 1.5 s, the old tests reproduce the CI failure 
and the new ones pass.
   
   ---
   
   ##### 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