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]

Reply via email to