Lee-W commented on code in PR #72149:
URL: https://github.com/apache/airflow/pull/72149#discussion_r4090968653
##########
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:
Review Comment:
Dropped the branch.
`execute_complete` now reads `event["batch_id"]` the same way the success
path does and calls `_cancel_batch_quietly` unconditionally on a timeout.
--
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]