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(
