kaxil commented on code in PR #72149:
URL: https://github.com/apache/airflow/pull/72149#discussion_r4123396711
##########
providers/openai/docs/changelog.rst:
##########
@@ -51,6 +51,13 @@ Changelog
metadata DB on every run, regardless of the operator's ``do_xcom_push``
setting.
+.. note::
+ A deferred ``OpenAITriggerBatchOperator`` that times out now requests
cancellation of the
Review Comment:
The earlier revision's note about deferred timeouts now raising
`OpenAIBatchTimeout` seems to have been lost in the rebase, and that's the one
users most need: `OpenAIBatchTimeout` isn't a subclass of
`OpenAIBatchJobException`, and 1.8.2 raised the latter for a deferred timeout,
so an `on_failure_callback` or retry rule keyed on `OpenAIBatchJobException`
silently stops matching. Also, this sits under 2.0.0, but
`providers-openai/2.0.0rc1` was cut from a commit without this change, so it
probably belongs above that heading.
##########
providers/openai/tests/unit/openai/operators/test_openai.py:
##########
@@ -935,3 +963,100 @@ def test_failed_event_raises(self, event):
def test_invalid_event_raises_instead_of_succeeding(self, event):
with pytest.raises(OpenAITriggerEventError):
self._operator().execute_complete(Context(), event)
+
+ @pytest.mark.parametrize(
+ ("termination_reason", "expected_exc"),
+ [
+ pytest.param("timeout", OpenAIBatchTimeout, id="timeout"),
Review Comment:
The `timeout` case never mocks the hook, so `_cancel_batch_quietly` goes
through a real `OpenAIHook` looking up `test_conn_id`, and the test passes only
because that failure is swallowed. Folding this and
`test_non_timeout_termination_reasons_do_not_cancel` into one table of (reason,
expected exception, cancel expected) with a `Mock(spec=OpenAIHook)` would cover
both and keep the test off the connection path.
##########
providers/openai/src/airflow/providers/openai/operators/openai.py:
##########
@@ -361,8 +365,16 @@ class OpenAITriggerBatchOperator(BaseOperator):
:param wait_seconds: Optional. Number of seconds between checks. Only used
when ``deferrable`` is False.
Defaults to 3 seconds.
:param timeout: Optional. The amount of time, in seconds, to wait for the
request to complete.
- Applies in both deferrable and non-deferrable mode. Defaults to 24
hours, which is the SLA for
- OpenAI Batch API.
+ Applies in both deferrable and non-deferrable mode: in the synchronous
path it bounds
+ ``wait_for_batch``; in the deferrable path it bounds the trigger's
poll loop. When the
+ deferrable path times out, the operator requests cancellation of the
batch using the
+ batch id carried by the trigger event, mirroring the synchronous path.
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, it does not
+ wait for it. If ``execution_timeout`` is set shorter than ``timeout``,
the scheduler's
+ deferral timeout fires first: the task is failed with
``TaskDeferralTimeout`` before the
+ trigger ever times out, ``execute_complete`` is never called, and this
cancellation path
+ does not run. Defaults to 24 hours, which is the SLA for OpenAI Batch
API.
Review Comment:
This param is now 11 lines, and the reason-to-exception mapping appears both
in the class docstring and in `execute_complete`'s, with the "cancelling for up
to 10 minutes" caveat in two places too. Could `:param timeout:` keep its
one-line meaning plus the `execution_timeout` caveat, with the mapping written
down once?
--
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]