[
https://issues.apache.org/jira/browse/FLINK-40604?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
ASF GitHub Bot updated FLINK-40604:
-----------------------------------
Labels: pull-request-available (was: )
> AsyncWaitOperator bypasses timeout handler when retry result mail runs after
> the deadline
> -----------------------------------------------------------------------------------------
>
> Key: FLINK-40604
> URL: https://issues.apache.org/jira/browse/FLINK-40604
> Project: Flink
> Issue Type: Bug
> Components: API / DataStream
> Reporter: David Peacock
> Priority: Major
> Labels: pull-request-available
>
> `AsyncWaitOperator` can forward a stale retry-attempt exception and fail the
> task after the overall async timeout has already elapsed, instead of invoking
> the configured `AsyncFunction#timeout` handler.
> The race is:
> 1. An async attempt completes exceptionally before the overall deadline.
> 2. `RetryableResultHandlerDelegator#completeExceptionally` determines that
> the result should be evaluated for retry and queues `processRetry` in the
> mailbox.
> 3. Processing time advances past the overall timeout before that retry mail
> is processed, while the timeout callback has not run yet.
> 4. `processRetry` evaluates `!isTimeout()` as false and falls into its
> generic terminal branch.
> 5. That branch forwards the stale attempt error to
> `ResultHandler#completeExceptionally`, which calls `failExternally()`.
> As a result, the outcome depends on mailbox/timer ordering. A user-provided
> timeout handler is bypassed even though the overall operation has timed out.
> The async I/O documentation defines the timeout as the maximum duration from
> the first invocation through final completion, including retries. It also
> says that users can override `AsyncFunction#timeout` to handle an expired
> request. The expired retry-result path should therefore use the same timeout
> handling path as the timer callback.
> h3. Expected behavior
> When retry-result mail is processed after the overall async deadline:
> - no further retry is scheduled;
> - `AsyncFunction#timeout` is invoked exactly once through the existing
> timeout path; and
> - the timeout handler decides whether to emit a fallback, emit nothing, or
> complete exceptionally.
> h3. Actual behavior
> The stale attempt result is completed normally or exceptionally. For an
> exceptional attempt, this can fail and restart the task without invoking
> `AsyncFunction#timeout`.
> h3. Reproducer
> A deterministic `AsyncWaitOperatorTest` can reproduce the race without sleeps:
> 1. Start an async operation with retries and a custom timeout handler that
> emits `-1`.
> 2. Complete the first attempt exceptionally so retry evaluation is queued in
> the mailbox.
> 3. Quiesce the processing-time service so the timeout callback does not run
> first.
> 4. Advance processing time to `timeout + 1`.
> 5. Process the queued retry mail.
> On unmodified master at `4e9d5412d3dcd4776b750e8a94ba6e86d2b409c6`, the test
> errors with:
> {code:java}
> java.lang.UnsupportedOperationException: MockEnvironment does not support
> external task failure. at MockEnvironment.failExternally(...) at
> AsyncWaitOperator$ResultHandler.completeExceptionally(...) at
> AsyncWaitOperator$RetryableResultHandlerDelegator.processRetry(...) {code}
> With the proposed fix, the timeout handler is called once, `-1` is emitted,
> and the task is not failed externally.
> h3. Proposed fix
> At the beginning of `processRetry`, after its existing duplicate-call guard,
> check whether the overall timeout has elapsed. If it has, invoke the existing
> `timerTriggered()` method and return before evaluating retry predicates or
> forwarding the attempt result.
> Because `timerTriggered()` may propagate an exception from a user timeout
> handler, narrow the `CollectionSupplier` overload's existing `try/catch` to
> cover only `supplier.get()`. This prevents a timeout-handler exception from
> being reclassified as a supplier failure and then ignored by the
> duplicate-call guard.
> This is deliberately narrow:
> - retry behavior before the deadline is unchanged;
> - Flink's existing strict `>` timeout boundary is unchanged;
> - no public API or serializer changes are involved; and
> - the timeout state transition remains centralized in `timerTriggered()`.
> I have a fix ready and can submit the PR if the approach has consensus, could
> a committer please assign FLINK-40604 to me?
> h3. Related issues
> - FLINK-30477 prevents retries from continuing after a timeout.
> - FLINK-40456 fixes the opposite ordering, where the timeout callback runs
> while retry mail is already queued and the timeout result could be dropped.
> - FLINK-34501 concerns behavior after retry exhaustion; this issue concerns
> retry-result mail processed after the overall timeout.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)