Skip to content

GH-51495: [C++] Fix race in MergedGenerator that could drop an error and end the stream early - #51498

Merged
zanmato1984 merged 3 commits into
apache:mainfrom
kita-renji:GH-51495-merged-generator-final-error
Sep 28, 2026
Merged

zanmato1984 merged 3 commits into
apache:mainfrom
kita-renji:GH-51495-merged-generator-final-error

Conversation

@kita-renji

@kita-renji kita-renji commented Sep 26, 2026 •

Copy link
Copy Markdown
Contributor

Rationale for this change

MergedGenerator could 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 set broken = true under the mutex but only stored the error in final_error after releasing it. A call to operator() in that window saw broken with an OK final_error and returned IterationEnd, 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?

  • The first error is recorded in the same locked section that sets broken, via SetFinalErrorUnlocked (replacing MarkFinalError), 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 sees broken also sees the error.
  • all_finished is removed, and the remaining futures are no longer completed through callbacks on it. Future callbacks have no guaranteed order (future.h): one added while a future is being marked finished may run immediately. With every terminal item chained on all_finished, a pull made during that dispatch could complete before an earlier future: the one receiving the error (pointed out in review), or an earlier IterationEnd.
    • Terminal pulls made once the generator is broken or exhausted are queued in waiting_jobs.
    • 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, which get IterationEnd. 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 an already-finished IterationEnd, as before.
  • A private, test-only static hook, 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 the MergedGeneratorErrorHookTest fixture, which is a friend. Without it I couldn't find a way to test this deterministically, since no user code runs inside the window.
  • Tests, for both the inner and outer error paths:
    • 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}ErrorToWaiterNotOvertakenDuringCompletion and ClaimedErrorNotOvertakenDuringCompletion: the later pull is made while the generator is completing its futures, with the earlier future's completion held up (via TryAddCallback, 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 MergedGenerator tests pass 300/300 with --gtest_repeat=300 (Debug) and 30/30 under TSAN with no reports. The three *DuringCompletion tests fail every time on the previous commit, and the stress test found the IterationEnd ordering issue before the fix.

Original commit (Linux):

  • With this PR: all 131 tests pass. The two new tests also passed 1000/1000 with --gtest_repeat=1000.
  • With only the final_error change reverted (hook kept): both new tests fail every time (2000/2000 over --gtest_repeat=1000) with Expected '_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_dataset from 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 a PyFileSystem wrapper, 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:

  • Human
  • AI

Reviewed before submission by:

  • Human
  • AI
  • Not reviewed

🤖 Generated with Claude Code

…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>
@github-actions

Copy link
Copy Markdown

⚠️ GitHub issue #51495 has been automatically assigned in GitHub to PR creator.

@zanmato1984 zanmato1984 left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment thread cpp/src/arrow/util/async_generator.h Outdated
}
if (should_mark_final_error) {
state->MarkFinalError(maybe_next->status(), std::move(sink));
state->DeliverFinalError(maybe_next->status(), std::move(sink));

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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:

  1. Use one inner generator whose first future is pending.
  2. Call merged() once, leaving an existing waiter.
  3. Call merged() again from error_signaled_hook_for_testing.
  4. 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.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for the careful review and the repro steps.

Fixed in 9a48233:

  • SetFinalErrorUnlocked now attaches the waiting caller's error callback to all_finished inside the same locked section that sets broken, on both the inner and outer error paths. Any pull that observes broken therefore registers its IterationEnd continuation after the error callback. DeliverFinalError is gone. all_finished can't be finished at that point, since the failed request is still counted in outstanding_requests, so the callback never runs under the lock. operator() already attaches its terminal Then() to all_finished while holding the same mutex.
  • Added InnerErrorToWaiterNotOvertakenByLaterPull and OuterErrorToWaiterNotOvertakenByLaterPull, 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 MergedGeneratorErrorHookTest fixture is a friend, the same pattern as key_hash_internal.h and chunker_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.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Pushed 364639c, which stops relying on callbacks for ordering altogether:

  • all_finished is gone. MarkFinishedAndPurge completes 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 get IterationEnd. Terminal pulls made after the generator is broken or exhausted are queued in waiting_jobs instead of being chained on all_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 IterationEnd complete before earlier ones: all terminal pulls hung off all_finished, so one added mid-dispatch ran at once. That's fixed by the same change.

Tests:

  • {Inner,Outer}ErrorToWaiterNotOvertakenDuringCompletion and ClaimedErrorNotOvertakenDuringCompletion use your technique: hold the earlier future's mutex via TryAddCallback and pull once the generator has started completing. All three fail every time on the previous commit.
  • PullFromLastFutureCallbackCompletesAtOnce.
  • MergedGeneratorStressTest.TerminalNeverOvertakesEarlierFutures checks 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.

kita-renji and others added 2 commits September 28, 2026 11:11
…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 zanmato1984 left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

+1

Thanks for working on this.

@github-actions github-actions Bot added awaiting committer review Awaiting committer review and removed awaiting review Awaiting review labels Sep 28, 2026
@zanmato1984
zanmato1984 merged commit 065db5f into apache:main Sep 28, 2026
61 of 62 checks passed
@zanmato1984 zanmato1984 removed the awaiting committer review Awaiting committer review label Sep 28, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants