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)

Reply via email to