From 5df0ee23863c67d85e53943ba22ba3cc77c0ce58 Mon Sep 17 00:00:00 2001 From: akarnokd Date: Wed, 5 Aug 2026 10:42:16 +0200 Subject: [PATCH] 4.x: Fix Streamable.fromStream not closing the Stream --- .../streamable/StreamableFromIterable.java | 17 +++++- .../streamable/StreamableFromStream.java | 10 +++- .../streamable/StreamableFromStreamTest.java | 56 +++++++++++++++++++ 3 files changed, 80 insertions(+), 3 deletions(-) diff --git a/src/main/java/io/reactivex/rxjava4/internal/operators/streamable/StreamableFromIterable.java b/src/main/java/io/reactivex/rxjava4/internal/operators/streamable/StreamableFromIterable.java index 82591d8d79..1339cefac8 100644 --- a/src/main/java/io/reactivex/rxjava4/internal/operators/streamable/StreamableFromIterable.java +++ b/src/main/java/io/reactivex/rxjava4/internal/operators/streamable/StreamableFromIterable.java @@ -38,19 +38,22 @@ public record StreamableFromIterable(@NonNull Iterable items) im Exceptions.throwIfFatal(ex); return StreamableError.createFailed(ex); } - return new IteratorStreamer<>(iterator); + return new IteratorStreamer<>(iterator, null); } static final class IteratorStreamer implements Streamer, EnumerableSource { Iterator iterator; + AutoCloseable toClose; + long index; T current; - IteratorStreamer(Iterator iterator) { + IteratorStreamer(Iterator iterator, AutoCloseable toClose) { this.iterator = iterator; + this.toClose = toClose; } @Override @@ -81,6 +84,16 @@ static NullPointerException createNullError(long index) { public @NonNull CompletionStage finish() { iterator = null; current = null; + var toClose = this.toClose; + this.toClose = null; + if (toClose != null) { + try { + toClose.close(); + } catch (Throwable ex) { + Exceptions.throwIfFatal(ex); + return CompletableFuture.failedFuture(ex); + } + } return FINISHED; } diff --git a/src/main/java/io/reactivex/rxjava4/internal/operators/streamable/StreamableFromStream.java b/src/main/java/io/reactivex/rxjava4/internal/operators/streamable/StreamableFromStream.java index c008bddbcc..58b8672d3b 100644 --- a/src/main/java/io/reactivex/rxjava4/internal/operators/streamable/StreamableFromStream.java +++ b/src/main/java/io/reactivex/rxjava4/internal/operators/streamable/StreamableFromStream.java @@ -19,6 +19,7 @@ import io.reactivex.rxjava4.annotations.NonNull; import io.reactivex.rxjava4.core.*; import io.reactivex.rxjava4.disposables.StreamerCancellation; +import io.reactivex.rxjava4.exceptions.Exceptions; import io.reactivex.rxjava4.internal.operators.streamable.StreamableEmpty.EmptyStreamer; public record StreamableFromStream(@NonNull Stream items) implements Streamable { @@ -28,8 +29,15 @@ public record StreamableFromStream(@NonNull Stream items) implem public @NonNull Streamer<@NonNull T> stream(@NonNull StreamerCancellation cancellation) { var iterator = Objects.requireNonNull(items.iterator(), "iterator is null"); if (!iterator.hasNext()) { + try { + items.close(); + } catch (Throwable ex) { + Exceptions.throwIfFatal(ex); + return StreamableError.createFailed(ex); + } + return (Streamer)EmptyStreamer.INSTANCE; } - return new StreamableFromIterable.IteratorStreamer<>(iterator); + return new StreamableFromIterable.IteratorStreamer<>(iterator, items); } } diff --git a/src/test/java/io/reactivex/rxjava4/internal/operators/streamable/StreamableFromStreamTest.java b/src/test/java/io/reactivex/rxjava4/internal/operators/streamable/StreamableFromStreamTest.java index cfbb1285af..d721c264ab 100644 --- a/src/test/java/io/reactivex/rxjava4/internal/operators/streamable/StreamableFromStreamTest.java +++ b/src/test/java/io/reactivex/rxjava4/internal/operators/streamable/StreamableFromStreamTest.java @@ -13,11 +13,17 @@ package io.reactivex.rxjava4.internal.operators.streamable; +import static org.junit.jupiter.api.Assertions.assertTrue; + import java.util.*; import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.stream.Stream; import org.junit.jupiter.api.Test; + import io.reactivex.rxjava4.core.Streamable; +import io.reactivex.rxjava4.exceptions.TestException; public class StreamableFromStreamTest extends StreamableBaseTest { @@ -64,4 +70,54 @@ public void hasNull2() throws Throwable { .assertError(t -> t.getMessage().equals("Item at index 0 is null.")); ; } + + @Test + public void ensureCloses() { + var hasClosed = new AtomicBoolean(); + Streamable.fromStream(Stream.empty().onClose(() -> hasClosed.set(true))) + .test() + .awaitDone(5, TimeUnit.SECONDS) + .assertResult(); + + assertTrue(hasClosed.get(), "Close not called"); + } + + @Test + public void ensureClosesCrash() { + var hasClosed = new AtomicBoolean(); + Streamable.fromStream(Stream.empty().onClose(() -> { + hasClosed.set(true); + throw new TestException(); + })) + .test() + .awaitDone(5, TimeUnit.SECONDS) + .assertFailure(TestException.class); + + assertTrue(hasClosed.get(), "Close not called"); + } + + @Test + public void ensureCloses2() { + var hasClosed = new AtomicBoolean(); + Streamable.fromStream(Stream.of(1, 2, 3, 4, 5).onClose(() -> hasClosed.set(true))) + .test() + .awaitDone(5, TimeUnit.SECONDS) + .assertResult(1, 2, 3, 4, 5); + + assertTrue(hasClosed.get(), "Close not called"); + } + + @Test + public void ensureCloses2Crash() { + var hasClosed = new AtomicBoolean(); + Streamable.fromStream(Stream.of(1, 2, 3, 4, 5).onClose(() -> { + hasClosed.set(true); + throw new TestException(); + })) + .test() + .awaitDone(5, TimeUnit.SECONDS) + .assertFailure(TestException.class, 1, 2, 3, 4, 5); + + assertTrue(hasClosed.get(), "Close not called"); + } }