This is an automated email from the ASF dual-hosted git repository. Lee-W pushed a commit to branch openai-batch-termination-reason in repository https://gitbox.apache.org/repos/asf/airflow.git
commit 461d0392d666824fc04a8cfccc191a8062add58a Author: Wei Lee <[email protected]> AuthorDate: Sat Sep 19 16:36:40 2026 +0900 Drop the unreachable missing-batch_id branch in execute_complete Every TriggerEvent OpenAIBatchTrigger yields sets batch_id from a required constructor argument, and OpenAITriggerBatchOperator only ever defers to that trigger, so the branch that logged "carried no batch_id" could not be reached and had no test covering it. The batch id is now read the way the success path in the same method already reads it, with event["batch_id"], so the two paths are consistently defensive. --- .../airflow/providers/openai/operators/openai.py | 21 +++++++-------------- 1 file changed, 7 insertions(+), 14 deletions(-) diff --git a/providers/openai/src/airflow/providers/openai/operators/openai.py b/providers/openai/src/airflow/providers/openai/operators/openai.py index f58258639f4..cade98742eb 100644 --- a/providers/openai/src/airflow/providers/openai/operators/openai.py +++ b/providers/openai/src/airflow/providers/openai/operators/openai.py @@ -487,20 +487,13 @@ class OpenAITriggerBatchOperator(BaseOperator): event = validate_execute_complete_event(event) if event["status"] != "success": if event.get("termination_reason") == "timeout": - batch_id = event.get("batch_id") - if batch_id: - self.log.warning( - "%s timed out waiting for batch %s; requesting cancellation.", - self.task_id, - batch_id, - ) - self._cancel_batch_quietly(batch_id) - else: - self.log.warning( - "%s timed out but the trigger event carried no batch_id; " - "skipping cancellation request.", - self.task_id, - ) + batch_id = event["batch_id"] + self.log.warning( + "%s timed out waiting for batch %s; requesting cancellation.", + self.task_id, + batch_id, + ) + self._cancel_batch_quietly(batch_id) raise build_batch_error(event["message"], event.get("termination_reason")) self.log.info("%s completed successfully.", self.task_id)
