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 f8b3a5bb467728836539e8f9dfa78526ff78810e
Author: Wei Lee <[email protected]>
AuthorDate: Sat Sep 19 17:04:14 2026 +0900

    Stop a failed cancellation from masking why the batch task stopped
    
    OpenAITriggerBatchOperator.on_kill and OpenAIHook.wait_for_batch's timeout
    branch both called cancel_batch bare, so an error from the cancellation 
request
    propagated in place of the reason the task was actually stopping. on_kill 
now
    goes through _cancel_batch_quietly, and the synchronous timeout path logs a
    warning and still raises OpenAIBatchTimeout. Both paths now behave the way 
the
    deferred timeout path already did.
---
 .../src/airflow/providers/openai/hooks/openai.py   |  7 ++++++-
 .../airflow/providers/openai/operators/openai.py   | 16 ++++++++-------
 .../openai/tests/unit/openai/hooks/test_openai.py  | 17 ++++++++++++++++
 .../tests/unit/openai/operators/test_openai.py     | 23 ++++++++++++++++++++++
 4 files changed, 55 insertions(+), 8 deletions(-)

diff --git a/providers/openai/src/airflow/providers/openai/hooks/openai.py 
b/providers/openai/src/airflow/providers/openai/hooks/openai.py
index 2ee588bf281..4e6562f27a6 100644
--- a/providers/openai/src/airflow/providers/openai/hooks/openai.py
+++ b/providers/openai/src/airflow/providers/openai/hooks/openai.py
@@ -709,7 +709,12 @@ class OpenAIHook(BaseHook):
         start = time.monotonic()
         while True:
             if start + timeout < time.monotonic():
-                self.cancel_batch(batch_id=batch_id)
+                try:
+                    self.cancel_batch(batch_id=batch_id)
+                except Exception as e:
+                    self.log.warning(
+                        "Failed to request cancellation of batch %s after 
timeout: %s", batch_id, e
+                    )
                 raise OpenAIBatchTimeout(f"Timeout: OpenAI Batch {batch_id} is 
not ready after {timeout}s")
             batch = self.get_batch(batch_id=batch_id)
 
diff --git a/providers/openai/src/airflow/providers/openai/operators/openai.py 
b/providers/openai/src/airflow/providers/openai/operators/openai.py
index cade98742eb..a0c936ba690 100644
--- a/providers/openai/src/airflow/providers/openai/operators/openai.py
+++ b/providers/openai/src/airflow/providers/openai/operators/openai.py
@@ -503,15 +503,17 @@ class OpenAITriggerBatchOperator(BaseOperator):
         """
         Best-effort request to cancel a batch; never raises.
 
-        Called from ``execute_complete`` after a deferred timeout, using the 
batch id carried
-        by the trigger event rather than ``self.batch_id`` — this method runs 
on a resumed task
-        instance, a fresh operator object on which ``execute``'s assignment to 
``self.batch_id``
-        never happened, so ``self.batch_id`` is ``None`` here.
+        Takes ``batch_id`` as a parameter rather than reading 
``self.batch_id`` because it has
+        two callers with different sources for it: ``execute_complete``, after 
a deferred
+        timeout, passes the batch id carried by the trigger event, since it 
runs on a resumed
+        task instance where ``execute``'s assignment to ``self.batch_id`` 
never happened;
+        ``on_kill`` passes ``self.batch_id`` directly, already set by 
``execute`` on this same
+        operator instance.
 
         Cancellation on OpenAI's side is asynchronous: the batch reports 
``cancelling`` for up
         to 10 minutes before it settles as ``cancelled``, so this only 
requests cancellation. A
-        failure to cancel is logged, not raised, so it never masks the timeout 
that is the
-        task's real failure reason.
+        failure to cancel is logged, not raised, so it never masks the real 
failure reason
+        (the timeout, or the kill).
         """
         try:
             self.hook.cancel_batch(batch_id)
@@ -522,4 +524,4 @@ class OpenAITriggerBatchOperator(BaseOperator):
         """Cancel the batch if task is cancelled."""
         if self.batch_id:
             self.log.info("on_kill: cancel the OpenAI Batch %s", self.batch_id)
-            self.hook.cancel_batch(self.batch_id)
+            self._cancel_batch_quietly(self.batch_id)
diff --git a/providers/openai/tests/unit/openai/hooks/test_openai.py 
b/providers/openai/tests/unit/openai/hooks/test_openai.py
index 6548d95654b..6b56df91f55 100644
--- a/providers/openai/tests/unit/openai/hooks/test_openai.py
+++ b/providers/openai/tests/unit/openai/hooks/test_openai.py
@@ -693,6 +693,23 @@ def 
test_wait_for_in_progress_batch_timeout(mock_openai_hook, mock_wip_batch):
     assert mock_openai_hook.conn.batches.cancel.call_count == 1
 
 
+def 
test_wait_for_in_progress_batch_timeout_cancel_failure_does_not_mask_timeout(
+    mock_openai_hook, mock_wip_batch, caplog
+):
+    """A cancellation failure inside the timeout branch must not replace 
``OpenAIBatchTimeout``
+    with the cancellation's own exception, and the failure must still be 
logged.
+    """
+    mock_openai_hook.conn.batches.retrieve.return_value = mock_wip_batch
+    mock_openai_hook.conn.batches.cancel.side_effect = RuntimeError("cancel 
failed")
+
+    with caplog.at_level("WARNING"):
+        with pytest.raises(OpenAIBatchTimeout, match="Timeout"):
+            mock_openai_hook.wait_for_batch(batch_id=BATCH_ID, 
wait_seconds=0.01, timeout=0.01)
+
+    assert mock_openai_hook.conn.batches.cancel.call_count == 1
+    assert any("Failed to request cancellation of batch" in message for 
message in caplog.messages)
+
+
 @pytest.mark.parametrize("status", ["cancelled", "cancelling"])
 def 
test_wait_for_cancelled_batch_raises_exact_cancelled_type(mock_openai_hook, 
status):
     """``OpenAIBatchCancelled`` is a subclass of ``OpenAIBatchJobException``, 
so asserting
diff --git a/providers/openai/tests/unit/openai/operators/test_openai.py 
b/providers/openai/tests/unit/openai/operators/test_openai.py
index 1c40f2951e5..3ff7581d328 100644
--- a/providers/openai/tests/unit/openai/operators/test_openai.py
+++ b/providers/openai/tests/unit/openai/operators/test_openai.py
@@ -903,6 +903,29 @@ def 
test_openai_trigger_batch_operator_deferred_logs_active_knob(mock_log, mock_
     )
 
 
+def test_openai_trigger_batch_operator_on_kill_cancels_batch_quietly(caplog):
+    """on_kill()'s cancellation failure is logged, not raised."""
+    operator = OpenAITriggerBatchOperator(
+        task_id=TASK_ID,
+        conn_id=CONN_ID,
+        file_id=FILE_ID,
+        endpoint=BATCH_ENDPOINT,
+    )
+    operator.batch_id = BATCH_ID
+    mock_hook_instance = Mock(spec=OpenAIHook)
+    mock_hook_instance.cancel_batch.side_effect = RuntimeError("cancel failed")
+    operator.hook = mock_hook_instance
+
+    with caplog.at_level("WARNING"):
+        try:
+            operator.on_kill()
+        except Exception as e:
+            pytest.fail(f"on_kill() should not raise: {e}")
+
+    mock_hook_instance.cancel_batch.assert_called_once_with(BATCH_ID)
+    assert any("Failed to request cancellation of batch" in message for 
message in caplog.messages)
+
+
 class TestOpenAITriggerBatchOperatorExecuteComplete:
     def _operator(self):
         return OpenAITriggerBatchOperator(

Reply via email to