Lee-W commented on code in PR #72149:
URL: https://github.com/apache/airflow/pull/72149#discussion_r4090996186


##########
providers/openai/src/airflow/providers/openai/operators/openai.py:
##########
@@ -446,15 +464,58 @@ def execute_complete(self, context: Context, event: Any = 
None) -> str:
         Invoke this callback when the trigger fires; return immediately.
 
         Relies on trigger to throw an exception, otherwise it assumes 
execution was
-        successful.
+        successful. The exception raised depends on the event's 
``termination_reason``:
+        ``OpenAIBatchTimeout`` for a timeout, ``OpenAIBatchCancelled`` for a 
cancellation,
+        and ``OpenAIBatchJobException`` for any other failure (including 
events from a
+        trigger serialized before ``termination_reason`` existed).
+
+        On a timeout, cancellation of the batch is requested before the 
timeout is raised
+        (see :meth:`_cancel_batch_quietly`). No other termination reason 
triggers
+        cancellation: a ``polling_error`` may be a transient, Airflow-side 
failure rather than
+        a real batch problem, and cancellation is irreversible, so it is left 
alone to run to
+        its own 24-hour completion window instead.
         """
         event = validate_execute_complete_event(event)
         if event["status"] != "success":
-            raise OpenAIBatchJobException(event["message"])
+            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,
+                    )
+            raise build_batch_error(event["message"], 
event.get("termination_reason"))
 
         self.log.info("%s completed successfully.", self.task_id)
         return event["batch_id"]
 
+    def _cancel_batch_quietly(self, batch_id: str) -> None:

Review Comment:
   `on_kill` now goes through `_cancel_batch_quietly`, and the timeout branch 
in `OpenAIHook.wait_for_batch` logs a warning on a failed cancel and still 
raises `OpenAIBatchTimeout`.
   
   Both have tests for the failing-cancel case. With that, the sync path does 
have the same protection, so I left the "mirroring the synchronous path" 
wording as is.



-- 
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