GH-51495: [C++] Fix race in MergedGenerator that could drop an error and end the stream early - #51498
Conversation
…error and end the stream early When an inner or outer subscription failed while no caller was waiting, MergedGenerator set `broken` under the mutex but stored the error in `final_error` only after releasing it. A pull in that window saw `broken` with an OK `final_error` and returned end-of-stream, so the error was lost. In the dataset scanner this made ToTable() return a truncated table instead of raising when a fragment could not be opened. Store `final_error` while the mutex is held. Add a test-only hook that runs right after the error state is entered, and use it for deterministic regression tests of the inner and outer error paths. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
|
|
zanmato1984
left a comment
There was a problem hiding this comment.
The no-waiter race identified by this PR is real, and storing final_error under the mutex fixes that path. However, the waiter-present path still has a related ordering race that allows a later end-of-stream future to complete before the earlier error future. I reproduced it deterministically using the newly added test hook, so I think this should be addressed before merge.
Non-blocking: please also consider keeping error_signaled_hook_for_testing private/internal instead of exposing a mutable static member from the public header.
| } | ||
| if (should_mark_final_error) { | ||
| state->MarkFinalError(maybe_next->status(), std::move(sink)); | ||
| state->DeliverFinalError(maybe_next->status(), std::move(sink)); |
There was a problem hiding this comment.
Could we also make the waiter-present path atomic with respect to later terminal pulls?
When sink is valid, SignalErrorUnlocked sets broken while holding the state mutex, but the error callback is not registered on all_finished until DeliverFinalError runs after the mutex is released. A concurrent operator() in that interval observes broken with an OK final_error and registers its IterationEnd continuation on all_finished first.
Since FutureImpl runs callbacks in registration order, the newer terminal future can complete before the older waiting future receives its error. This violates the AsyncGenerator contract that a terminal result must not complete while an earlier returned future is still outstanding.
The new hook makes this deterministic to reproduce:
- Use one inner generator whose first future is pending.
- Call
merged()once, leaving an existing waiter. - Call
merged()again fromerror_signaled_hook_for_testing. - Fail the pending future and record callback completion order.
Expected: error, terminal
Actual on this commit: terminal, error
Could we register or persist the waiting error as part of the locked state transition so that later terminal pulls cannot overtake it, and add a regression test for this waiter-present path? The same ordering should also be checked for the outer-error path.
There was a problem hiding this comment.
Thanks for the careful review and the repro steps.
Fixed in 9a48233:
SetFinalErrorUnlockednow attaches the waiting caller's error callback toall_finishedinside the same locked section that setsbroken, on both the inner and outer error paths. Any pull that observesbrokentherefore registers itsIterationEndcontinuation after the error callback.DeliverFinalErroris gone.all_finishedcan't be finished at that point, since the failed request is still counted inoutstanding_requests, so the callback never runs under the lock.operator()already attaches its terminalThen()toall_finishedwhile holding the same mutex.- Added
InnerErrorToWaiterNotOvertakenByLaterPullandOuterErrorToWaiterNotOvertakenByLaterPull, following your steps: one caller waiting, a second pull from the hook, then fail the pending future and check the completion order. Both fail on the previous commit with{"terminal", "error"}and pass now. - On the non-blocking point: the hook is now a private static member, and the
MergedGeneratorErrorHookTestfixture is a friend, the same pattern askey_hash_internal.handchunker_internal.h. It's no longer part of the public API.
All 133 tests in arrow-async-utility-test pass locally, and the MergedGenerator tests passed 200/200 with --gtest_repeat.
There was a problem hiding this 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:
- Leave the first
merged()future waiting on a pending inner future. - Use
waiting.TryAddCallback(...)to hold that future's implementation mutex, so the queued error callback blocks insink.MarkFinished(err). - Fail the inner future and wait for the error-state hook.
- Pull again while
all_finishedis 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.
There was a problem hiding this comment.
Pushed 364639c, which stops relying on callbacks for ordering altogether:
all_finishedis gone.MarkFinishedAndPurgecompletes the remaining futures itself, in the order they were handed out: the future receiving the error first (a waiter, or the first caller to ask after an error nobody was waiting for), then every waiting caller and every terminal item requested since, which getIterationEnd. Terminal pulls made after the generator is broken or exhausted are queued inwaiting_jobsinstead of being chained onall_finished. Pulls that arrive while the queue is being completed are queued behind it, and it is drained until it stays empty.- A caller that asks from the callbacks of the last pending future still gets an already-finished
IterationEnd, as it did before, so a callback that pulls and then blocks cannot deadlock. - While testing this I found that the same mechanism also let a later
IterationEndcomplete before earlier ones: all terminal pulls hung offall_finished, so one added mid-dispatch ran at once. That's fixed by the same change.
Tests:
{Inner,Outer}ErrorToWaiterNotOvertakenDuringCompletionandClaimedErrorNotOvertakenDuringCompletionuse your technique: hold the earlier future's mutex viaTryAddCallbackand pull once the generator has started completing. All three fail every time on the previous commit.PullFromLastFutureCallbackCompletesAtOnce.MergedGeneratorStressTest.TerminalNeverOvertakesEarlierFutureschecks the contract directly with inner and outer failures on the thread pool: no terminal completes while an earlier future is pending, and exactly one error is delivered.
All 138 tests in arrow-async-utility-test pass. The MergedGenerator tests pass 300/300 repeats in debug and 30/30 under TSAN with no reports.
…lls in MergedGenerator When a caller was already waiting for the item that failed, the error callback was attached to `all_finished` only after the mutex was released. A pull in that window saw `broken`, attached its end-of-stream continuation to `all_finished` first, and since callbacks run in registration order the later terminal item completed before the earlier waiting future got its error. Attach the waiting caller's error callback in the same locked section that sets `broken`, for both the inner and outer error paths. This removes `DeliverFinalError`. Add regression tests for the waiter-present path. Also make `error_signaled_hook_for_testing` private and friend the test fixture instead of exposing a mutable static in the public header. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…ad of relying on callback order Future callbacks have no guaranteed execution order: one added while the future is being marked finished may run immediately. MergedGenerator chained every terminal item, and the delivery of an error to a waiting caller, on `all_finished`, so a pull made while `all_finished` was dispatching could get its terminal item before an earlier future (the one receiving the error, or an earlier IterationEnd) had completed. Remove `all_finished`. Terminal pulls made once the generator is broken or exhausted are queued in `waiting_jobs`, and MarkFinishedAndPurge completes the remaining futures itself, in the order they were handed out: the future receiving the error first, then the waiting callers and terminal items. Pulls that arrive meanwhile are queued behind them, and the queue is drained until it stays empty. A caller that asks from the callbacks of the last pending future still gets a finished IterationEnd. Add deterministic tests for pulls made while the generator is completing (inner error, outer error, error nobody was waiting for), a test for pulling from the last future's callback, and a stress test that checks the AsyncGenerator ordering contract directly. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
zanmato1984
left a comment
There was a problem hiding this comment.
+1
Thanks for working on this.
Rationale for this change
MergedGeneratorcould turn an error from a subscription into a normal end-of-stream. When an inner or outer subscription failed while no caller was waiting, the callback setbroken = trueunder the mutex but only stored the error infinal_errorafter releasing it. A call tooperator()in that window sawbrokenwith an OKfinal_errorand returnedIterationEnd, and the error was never delivered.In the dataset scanner this shows up as
ToTable()returning a truncated table without raising when a fragment cannot be opened (#51495).What changes are included in this PR?
broken, viaSetFinalErrorUnlocked(replacingMarkFinalError), together with the caller that will receive it, if one is already waiting. Otherwise the error goes to the next caller that asks, so a concurrent pull that seesbrokenalso sees the error.all_finishedis removed, and the remaining futures are no longer completed through callbacks on it.Futurecallbacks have no guaranteed order (future.h): one added while a future is being marked finished may run immediately. With every terminal item chained onall_finished, a pull made during that dispatch could complete before an earlier future: the one receiving the error (pointed out in review), or an earlierIterationEnd.waiting_jobs.MarkFinishedAndPurgecompletes the remaining futures itself, in the order they were handed out: the future receiving the error first, then the waiting callers and terminal items, which getIterationEnd. Pulls that arrive meanwhile are queued behind them, and the queue is drained until it stays empty.IterationEnd, as before.MergedGenerator<T>::error_signaled_hook_for_testing, that runs right after the generator enters its error state and the mutex is released. It is empty by default and only reachable from theMergedGeneratorErrorHookTestfixture, which is a friend. Without it I couldn't find a way to test this deterministically, since no user code runs inside the window.MergedGeneratorErrorHookTest.{Inner,Outer}ErrorNotLostToConcurrentPull: nobody waiting; a pull from inside the hook must raise.MergedGeneratorErrorHookTest.{Inner,Outer}ErrorToWaiterNotOvertakenByLaterPull: a caller already waiting must get the error before a later pull gets end-of-stream.MergedGeneratorErrorHookTest.{Inner,Outer}ErrorToWaiterNotOvertakenDuringCompletionandClaimedErrorNotOvertakenDuringCompletion: the later pull is made while the generator is completing its futures, with the earlier future's completion held up (viaTryAddCallback, as in the review).MergedGeneratorTest.PullFromLastFutureCallbackCompletesAtOnce.MergedGeneratorStressTest.TerminalNeverOvertakesEarlierFutures: inner and outer failures on the thread pool; checks that no terminal item completes while an earlier future is pending and that exactly one error is delivered.Are these changes tested?
Yes.
arrow-async-utility-test.Latest commit (macOS arm64): all 138 tests pass. The
MergedGeneratortests pass 300/300 with--gtest_repeat=300(Debug) and 30/30 under TSAN with no reports. The three*DuringCompletiontests fail every time on the previous commit, and the stress test found theIterationEndordering issue before the fix.Original commit (Linux):
--gtest_repeat=1000.final_errorchange reverted (hook kept): both new tests fail every time (2000/2000 over--gtest_repeat=1000) withExpected '_fut.status()' to fail with Invalid, but got OK. The other 129 tests still pass.I also checked the fix end to end outside the test suite. I built
libarrow_datasetfrom the 25.0.1 tag with the same header change and swapped it into the 25.0.1 wheel. The Python reproducer from the issue (going through aPyFileSystemwrapper, which is what made it reproducible for me) went from 57/4500 truncated tables to 0 in about 11.6k runs.Are there any user-facing changes?
A scan that hits a fragment error now always raises. Before, it could sometimes return a partial result.
This PR contains a "Critical Fix". It fixes a bug where a dataset scan could return incorrect (truncated) results with no error.
Was AI used for this PR?
In accordance to the AI generation guidelines, please disclose below whether and how AI was used in this PR.
PR code and description written by:
Reviewed before submission by:
🤖 Generated with Claude Code