zanmato1984 commented on code in PR #51498:
URL: https://github.com/apache/arrow/pull/51498#discussion_r4118552933
##########
cpp/src/arrow/util/async_generator.h:
##########
@@ -1338,8 +1363,11 @@ class MergedGenerator {
// Now we have given up the lock and we can take all the actions we
decided we
// need to take.
+ if (signaled_error) {
+ RunErrorSignaledHookForTesting();
+ }
if (should_mark_final_error) {
- state->MarkFinalError(maybe_next->status(), std::move(sink));
+ state->DeliverFinalError(maybe_next->status(), std::move(sink));
Review Comment:
Thanks for the update. Registering the waiting error callback under the
state mutex closes the interleaving from my first comment, and the new
inner/outer tests cover that window. However, I don't think callback
registration order is sufficient to guarantee the required completion order.
`Future::AddCallback` explicitly does not guarantee callback execution
order; in particular, a callback added while the future is being marked
complete may run immediately, ahead of or concurrently with callbacks
registered earlier. `all_finished.MarkFinished()` marks `all_finished` complete
before it invokes the previously queued callback that calls
`sink.MarkFinished(err)`. A concurrent `operator()` in that interval sees
`broken`, calls `all_finished.Then(...)`, and its terminal continuation can run
synchronously while the older waiting future is still pending.
I reproduced this deterministically without sleeps:
1. Leave the first `merged()` future waiting on a pending inner future.
2. Use `waiting.TryAddCallback(...)` to hold that future's implementation
mutex, so the queued error callback blocks in `sink.MarkFinished(err)`.
3. Fail the inner future and wait for the error-state hook.
4. Pull again while `all_finished` is dispatching callbacks.
The later terminal future finishes while the earlier waiting error future is
still pending (`terminal_overtook_error == true`). This still violates the
AsyncGenerator rule that a terminal value must not complete while an earlier
returned future is outstanding. The current tests pass because the hook
performs the later pull before `all_finished.MarkFinished()`, so both callbacks
are already queued and happen to run in vector order; they do not cover a
callback added after `all_finished` has transitioned to finished.
Could we make terminal readiness depend on the waiting error actually being
delivered, rather than on callback registration order? For example, save the
error sink and complete it before calling `all_finished.MarkFinished()`, or use
a separate `error_delivered` gate that later terminal pulls wait on.
--
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]