Skip to content

Commit d36352f

Browse files
authored
4.x: Fix Streamable.fromStream not closing the Stream (#8267)
1 parent 3895db1 commit d36352f

3 files changed

Lines changed: 80 additions & 3 deletions

File tree

src/main/java/io/reactivex/rxjava4/internal/operators/streamable/StreamableFromIterable.java

Lines changed: 15 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -38,19 +38,22 @@ public record StreamableFromIterable<T>(@NonNull Iterable<? extends T> items) im
3838
Exceptions.throwIfFatal(ex);
3939
return StreamableError.createFailed(ex);
4040
}
41-
return new IteratorStreamer<>(iterator);
41+
return new IteratorStreamer<>(iterator, null);
4242
}
4343

4444
static final class IteratorStreamer<T> implements Streamer<T>, EnumerableSource<T> {
4545

4646
Iterator<? extends T> iterator;
4747

48+
AutoCloseable toClose;
49+
4850
long index;
4951

5052
T current;
5153

52-
IteratorStreamer(Iterator<? extends T> iterator) {
54+
IteratorStreamer(Iterator<? extends T> iterator, AutoCloseable toClose) {
5355
this.iterator = iterator;
56+
this.toClose = toClose;
5457
}
5558

5659
@Override
@@ -81,6 +84,16 @@ static NullPointerException createNullError(long index) {
8184
public @NonNull CompletionStage<Void> finish() {
8285
iterator = null;
8386
current = null;
87+
var toClose = this.toClose;
88+
this.toClose = null;
89+
if (toClose != null) {
90+
try {
91+
toClose.close();
92+
} catch (Throwable ex) {
93+
Exceptions.throwIfFatal(ex);
94+
return CompletableFuture.failedFuture(ex);
95+
}
96+
}
8497
return FINISHED;
8598
}
8699

src/main/java/io/reactivex/rxjava4/internal/operators/streamable/StreamableFromStream.java

Lines changed: 9 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,7 @@
1919
import io.reactivex.rxjava4.annotations.NonNull;
2020
import io.reactivex.rxjava4.core.*;
2121
import io.reactivex.rxjava4.disposables.StreamerCancellation;
22+
import io.reactivex.rxjava4.exceptions.Exceptions;
2223
import io.reactivex.rxjava4.internal.operators.streamable.StreamableEmpty.EmptyStreamer;
2324

2425
public record StreamableFromStream<T>(@NonNull Stream<? extends T> items) implements Streamable<T> {
@@ -28,8 +29,15 @@ public record StreamableFromStream<T>(@NonNull Stream<? extends T> items) implem
2829
public @NonNull Streamer<@NonNull T> stream(@NonNull StreamerCancellation cancellation) {
2930
var iterator = Objects.requireNonNull(items.iterator(), "iterator is null");
3031
if (!iterator.hasNext()) {
32+
try {
33+
items.close();
34+
} catch (Throwable ex) {
35+
Exceptions.throwIfFatal(ex);
36+
return StreamableError.createFailed(ex);
37+
}
38+
3139
return (Streamer<T>)EmptyStreamer.INSTANCE;
3240
}
33-
return new StreamableFromIterable.IteratorStreamer<>(iterator);
41+
return new StreamableFromIterable.IteratorStreamer<>(iterator, items);
3442
}
3543
}

src/test/java/io/reactivex/rxjava4/internal/operators/streamable/StreamableFromStreamTest.java

Lines changed: 56 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -13,11 +13,17 @@
1313

1414
package io.reactivex.rxjava4.internal.operators.streamable;
1515

16+
import static org.junit.jupiter.api.Assertions.assertTrue;
17+
1618
import java.util.*;
1719
import java.util.concurrent.TimeUnit;
20+
import java.util.concurrent.atomic.AtomicBoolean;
21+
import java.util.stream.Stream;
1822

1923
import org.junit.jupiter.api.Test;
24+
2025
import io.reactivex.rxjava4.core.Streamable;
26+
import io.reactivex.rxjava4.exceptions.TestException;
2127

2228
public class StreamableFromStreamTest extends StreamableBaseTest {
2329

@@ -64,4 +70,54 @@ public void hasNull2() throws Throwable {
6470
.assertError(t -> t.getMessage().equals("Item at index 0 is null."));
6571
;
6672
}
73+
74+
@Test
75+
public void ensureCloses() {
76+
var hasClosed = new AtomicBoolean();
77+
Streamable.fromStream(Stream.empty().onClose(() -> hasClosed.set(true)))
78+
.test()
79+
.awaitDone(5, TimeUnit.SECONDS)
80+
.assertResult();
81+
82+
assertTrue(hasClosed.get(), "Close not called");
83+
}
84+
85+
@Test
86+
public void ensureClosesCrash() {
87+
var hasClosed = new AtomicBoolean();
88+
Streamable.fromStream(Stream.empty().onClose(() -> {
89+
hasClosed.set(true);
90+
throw new TestException();
91+
}))
92+
.test()
93+
.awaitDone(5, TimeUnit.SECONDS)
94+
.assertFailure(TestException.class);
95+
96+
assertTrue(hasClosed.get(), "Close not called");
97+
}
98+
99+
@Test
100+
public void ensureCloses2() {
101+
var hasClosed = new AtomicBoolean();
102+
Streamable.fromStream(Stream.of(1, 2, 3, 4, 5).onClose(() -> hasClosed.set(true)))
103+
.test()
104+
.awaitDone(5, TimeUnit.SECONDS)
105+
.assertResult(1, 2, 3, 4, 5);
106+
107+
assertTrue(hasClosed.get(), "Close not called");
108+
}
109+
110+
@Test
111+
public void ensureCloses2Crash() {
112+
var hasClosed = new AtomicBoolean();
113+
Streamable.fromStream(Stream.of(1, 2, 3, 4, 5).onClose(() -> {
114+
hasClosed.set(true);
115+
throw new TestException();
116+
}))
117+
.test()
118+
.awaitDone(5, TimeUnit.SECONDS)
119+
.assertFailure(TestException.class, 1, 2, 3, 4, 5);
120+
121+
assertTrue(hasClosed.get(), "Close not called");
122+
}
67123
}

0 commit comments

Comments
 (0)