From 65819d31183d2e9b3d7d442caa3ff79ae70ca262 Mon Sep 17 00:00:00 2001 From: myetcd Date: Wed, 30 Sep 2026 00:41:32 +0800 Subject: [PATCH] Fix finish error propagation in StreamableUsing Signed-off-by: myetcd --- .../operators/streamable/StreamableUsing.java | 12 ++++++--- .../streamable/StreamableUsingTest.java | 25 +++++++++++++++++++ 2 files changed, 33 insertions(+), 4 deletions(-) diff --git a/src/main/java/io/reactivex/rxjava4/internal/operators/streamable/StreamableUsing.java b/src/main/java/io/reactivex/rxjava4/internal/operators/streamable/StreamableUsing.java index c7675ea16f..fb1e025e0f 100644 --- a/src/main/java/io/reactivex/rxjava4/internal/operators/streamable/StreamableUsing.java +++ b/src/main/java/io/reactivex/rxjava4/internal/operators/streamable/StreamableUsing.java @@ -21,6 +21,7 @@ import io.reactivex.rxjava4.disposables.StreamerCancellation; import io.reactivex.rxjava4.exceptions.Exceptions; import io.reactivex.rxjava4.functions.*; +import io.reactivex.rxjava4.internal.util.ExceptionHelper; public record StreamableUsing( Supplier resourceSupplier, @@ -80,11 +81,14 @@ static final class UsingStreamer implements Streamer { cleanup = null; c.run(); } catch (Throwable ex) { - Exceptions.throwIfFatal(e); - cf.completeExceptionally(ex); - return; + Exceptions.throwIfFatal(ex); + e = ExceptionHelper.unwrapAndCombine(e, ex); + } + if (e != null) { + cf.completeExceptionally(e); + } else { + cf.complete(null); } - cf.complete(null); }); return cf; } diff --git a/src/test/java/io/reactivex/rxjava4/internal/operators/streamable/StreamableUsingTest.java b/src/test/java/io/reactivex/rxjava4/internal/operators/streamable/StreamableUsingTest.java index b4e22b636e..d34db6263c 100644 --- a/src/test/java/io/reactivex/rxjava4/internal/operators/streamable/StreamableUsingTest.java +++ b/src/test/java/io/reactivex/rxjava4/internal/operators/streamable/StreamableUsingTest.java @@ -87,4 +87,29 @@ public void resourceCleanerCrash() throws Throwable { .awaitDone(5, TimeUnit.SECONDS) .assertFailure(TestException.class, 5, 6, 7, 8, 9); } + + @Test + public void upstreamFinishCrash() { + var resource = new AtomicReference(); + + Streamable.using(() -> 1, _ -> StreamableFailingFinish.MAIN_COMPLETES, resource::set) + .test() + .awaitDone(5, TimeUnit.SECONDS) + .assertFailure(TestException.class) + .assertError(e -> e.getMessage().equals("StreamableFailingFinish.finish()")); + + assertEquals(1, resource.get(), "resource cleanup mismatch"); + } + + @Test + public void upstreamFinishAndResourceCleanerCrash() { + var cleanerError = new TestException("resourceCleaner"); + + Streamable.using(() -> 1, _ -> StreamableFailingFinish.MAIN_COMPLETES, _ -> { throw cleanerError; }) + .test() + .awaitDone(5, TimeUnit.SECONDS) + .assertFailure(TestException.class) + .assertError(e -> e.getMessage().equals("StreamableFailingFinish.finish()")) + .assertError(e -> e.getSuppressed().length == 1 && e.getSuppressed()[0] == cleanerError); + } }