Skip to content

Commit b5ebbd9

Browse files
authored
4.x: Streamable + concat, using, test fixes (#8220)
1 parent ae52e2e commit b5ebbd9

8 files changed

Lines changed: 545 additions & 13 deletions

File tree

src/main/java/io/reactivex/rxjava4/core/Streamable.java

Lines changed: 71 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -79,10 +79,6 @@ public interface Streamable<@NonNull T> {
7979
@NonNull
8080
Streamer<T> stream(@NonNull DisposableContainer cancellation);
8181

82-
// oooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooo
83-
// HELPERS
84-
// oooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooo
85-
8682
// oooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooo
8783
// Data sources and wrappers
8884
// oooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooooo
@@ -111,6 +107,21 @@ public interface Streamable<@NonNull T> {
111107
}, executor);
112108
}
113109

110+
/**
111+
* Streams the {@code Streamable}s one after the other from the given iterable sequence
112+
* of {@code Streamable}s.
113+
* @param <T> the element type of the inner and resulting {@code Streamable}s
114+
* @param sources the iterable sequence of {@code Streamable}s.
115+
* @return the new {@code Streamable} source
116+
* @throws NullPointerException if {@code sources} is {@code null}
117+
*/
118+
@CheckReturnValue
119+
@NonNull
120+
static <@NonNull T> Streamable<T> concat(Iterable<? extends Streamable<? extends T>> sources) {
121+
Objects.requireNonNull(sources, "sources is null");
122+
return RxJavaPlugins.onAssembly(new StreamableConcatIterable<>(sources, ErrorMode.IMMEDIATE)); // TODO implement
123+
}
124+
114125
/**
115126
* Generate a sequence of values via a virtual generator callback (yielder)
116127
* which is free to block and is natively backpressured.
@@ -514,6 +525,8 @@ static Streamable<Long> rangeLong(long start, long count) {
514525
* @return the new {@code Streamable} instance
515526
* @throws NullPointerException if {@code unit} or {@code scheduler} is {@code null}
516527
*/
528+
@CheckReturnValue
529+
@NonNull
517530
static Streamable<Long> timer(long delay, TimeUnit unit, Scheduler scheduler) {
518531
Objects.requireNonNull(unit, "unit is null");
519532
Objects.requireNonNull(scheduler, "scheduler is null");
@@ -533,12 +546,39 @@ static Streamable<Long> timer(long delay, TimeUnit unit, Scheduler scheduler) {
533546
* @return the new {@code Streamable} instance
534547
* @throws NullPointerException if {@code unit} or {@code executor} is {@code null}
535548
*/
549+
@CheckReturnValue
550+
@NonNull
536551
static Streamable<Long> timer(long delay, TimeUnit unit, ExecutorService executor) {
537552
Objects.requireNonNull(unit, "unit is null");
538553
Objects.requireNonNull(executor, "executor is null");
539554
return RxJavaPlugins.onAssembly(new StreamableTimer(delay, unit, null, executor));
540555
}
541556

557+
/**
558+
* For each incoming streamer, this operator creates a resource, then
559+
* uses that resource to create the actual {@code Streamable} instance to
560+
* stream value of and then uses a cleaner callback to dissolve the resource
561+
* once the {@code Streamable} terminated.
562+
* @param <T> the element type of the sequence
563+
* @param <R> the resource type
564+
* @param resourceSupplier supplies a resource object per {@link #stream(DisposableContainer)} call
565+
* @param resourceMapper maps the supplied resource into a {@code Streamable} source
566+
* @param resourceCleaner cleans up the supplied resource
567+
* @return the new {@code Streamable} instance
568+
* @throws NullPointerException if {@code resourceSupplier} or {@code resourceMapper}
569+
* or {@code resourceCleaner} is {@code null}
570+
*/
571+
@CheckReturnValue
572+
@NonNull
573+
static <T, R> Streamable<T> using(Supplier<? extends R> resourceSupplier,
574+
Function<? super R, ? extends Streamable<? extends T>> resourceMapper,
575+
Consumer<? super R> resourceCleaner) {
576+
Objects.requireNonNull(resourceSupplier, "resourceSupplier is null");
577+
Objects.requireNonNull(resourceMapper, "resourceMapper is null");
578+
Objects.requireNonNull(resourceCleaner, "resourceCleaner is null");
579+
return RxJavaPlugins.onAssembly(new StreamableUsing<>(resourceSupplier, resourceMapper, resourceCleaner));
580+
}
581+
542582
/**
543583
* Takes the next element from each source {@code Streamable} and emits them a a single
544584
* row of {@link List}.
@@ -550,6 +590,8 @@ static Streamable<Long> timer(long delay, TimeUnit unit, ExecutorService executo
550590
* @return the new {@code Streamable} instance
551591
* @throws NullPointerException if {@code sources} is {@&ode null}
552592
*/
593+
@CheckReturnValue
594+
@NonNull
553595
static <T> Streamable<List<T>> zip(Iterable<? extends Streamable<? extends T>> sources) {
554596
Objects.requireNonNull(sources, "sources is null");
555597
return RxJavaPlugins.onAssembly(new StreamableZip<>(sources));
@@ -570,7 +612,9 @@ static <T> Streamable<List<T>> zip(Iterable<? extends Streamable<? extends T>> s
570612
* @return the new {@code Streamable} instance
571613
* @throws NullPointerException if {@code collector} is {@code null}
572614
*/
573-
default <A, R> Streamable<R> collect(Collector<T, A, R> collector) {
615+
@CheckReturnValue
616+
@NonNull
617+
default <A, R> Streamable<R> collect(@NonNull Collector<T, A, R> collector) {
574618
Objects.requireNonNull(collector, "collector is null");
575619
return RxJavaPlugins.onAssembly(new StreamableCollector<>(this, collector));
576620
}
@@ -583,7 +627,9 @@ default <A, R> Streamable<R> collect(Collector<T, A, R> collector) {
583627
* @return the new {@code Streamable} instance
584628
* @throws NullPointerException if {@code unit} or {@code scheduler} is {@code null}
585629
*/
586-
default Streamable<T> delay(long time, TimeUnit unit, Scheduler scheduler) {
630+
@CheckReturnValue
631+
@NonNull
632+
default Streamable<T> delay(long time, @NonNull TimeUnit unit, @NonNull Scheduler scheduler) {
587633
Objects.requireNonNull(unit, "unit is null");
588634
Objects.requireNonNull(scheduler, "scheduler is null");
589635
return RxJavaPlugins.onAssembly(new StreamableDelay<>(this, time, unit, scheduler));
@@ -595,7 +641,9 @@ default Streamable<T> delay(long time, TimeUnit unit, Scheduler scheduler) {
595641
* @return the new {@code Streamable} instance
596642
* @throws NullPointerException if {@code consumer} is {@code null}
597643
*/
598-
default Streamable<T> doOnError(Consumer<? super Throwable> consumer) {
644+
@CheckReturnValue
645+
@NonNull
646+
default Streamable<T> doOnError(@NonNull Consumer<? super Throwable> consumer) {
599647
Objects.requireNonNull(consumer, "consumer is null");
600648
return intercept(StreamableHelper.createOnError(consumer));
601649
}
@@ -606,7 +654,9 @@ default Streamable<T> doOnError(Consumer<? super Throwable> consumer) {
606654
* @return the new {@code Streamable} instance
607655
* @throws NullPointerException if {@code consumer} is {@code null}
608656
*/
609-
default Streamable<T> doOnNext(Consumer<? super T> consumer) {
657+
@CheckReturnValue
658+
@NonNull
659+
default Streamable<T> doOnNext(@NonNull Consumer<? super T> consumer) {
610660
Objects.requireNonNull(consumer, "consumer is null");
611661
return intercept(new StreamableInterceptConfig<>(v -> { consumer.accept(v); return v; } ));
612662
}
@@ -638,7 +688,9 @@ default <R> Streamable<R> flatMap(@NonNull Function<? super T, ? extends Streama
638688
* @return the new {@code Streamable} instance
639689
* @throws NullPointerException if {@code keySelector} is {@code null}
640690
*/
641-
default <@Nullable K> Streamable<GroupedStreamable<K, T>> groupBy(Function<? super T, ? extends K> keySelector) {
691+
@CheckReturnValue
692+
@NonNull
693+
default <@Nullable K> Streamable<GroupedStreamable<K, T>> groupBy(@NonNull Function<? super T, ? extends K> keySelector) {
642694
Objects.requireNonNull(keySelector, "keySelector is null");
643695
return RxJavaPlugins.onAssembly(new StreamableGroupBy<>(this, keySelector));
644696
}
@@ -732,7 +784,9 @@ default Streamable<T> intercept(StreamableInterceptConfig<T> config) {
732784
* @return the new {@code Streamable} instance
733785
* @throws NullPointerException if {@code fallbackMapper} is {@code null}
734786
*/
735-
default Streamable<T> onErrorResumeNext(Function<? super Throwable, ? extends Streamable<? extends T>> fallbackMapper) {
787+
@CheckReturnValue
788+
@NonNull
789+
default Streamable<T> onErrorResumeNext(@NonNull Function<? super Throwable, ? extends Streamable<? extends T>> fallbackMapper) {
736790
Objects.requireNonNull(fallbackMapper, "fallbackMapper is null");
737791
return RxJavaPlugins.onAssembly(new StreamableOnErrorResumeNext<>(this, fallbackMapper));
738792
}
@@ -798,7 +852,9 @@ default Streamable<T> takeWhile(@NonNull Predicate<? super T> predicate) {
798852
* @return the new {@code Streamable} instance
799853
* @throws NullPointerException if {@code unit} or {@code scheduler} or {@code fallback} is {@code null}
800854
*/
801-
default Streamable<T> timeout(long timeout, TimeUnit unit, Scheduler scheduler, Streamable<T> fallback) {
855+
@CheckReturnValue
856+
@NonNull
857+
default Streamable<T> timeout(long timeout, @NonNull TimeUnit unit, @NonNull Scheduler scheduler, @NonNull Streamable<T> fallback) {
802858
Objects.requireNonNull(unit, "unit is null");
803859
Objects.requireNonNull(scheduler, "scheduler is null");
804860
Objects.requireNonNull(fallback, "fallback is null");
@@ -855,6 +911,8 @@ default Flowable<T> toFlowable(@NonNull ExecutorService executor) {
855911
* or {@link ExecutorService} on its own.
856912
* @return the new {@code Observable} instance
857913
*/
914+
@CheckReturnValue
915+
@NonNull
858916
default Observable<T> toObservable() {
859917
return RxJavaPlugins.onAssembly(new StreamableToObservable<>(this));
860918
}
@@ -1003,6 +1061,7 @@ default void subscribe(@NonNull Flow.Subscriber<? super T> subscriber) {
10031061
* {@code Streamable} terminates
10041062
* @throws NullPointerException if {@code consumer} is {@code null}
10051063
*/
1064+
@NonNull
10061065
default CompletionStage<Void> subscribe(@NonNull StreamSink<? super T> consumer) {
10071066
return subscribe(consumer, Executors.newVirtualThreadPerTaskExecutor());
10081067
}
@@ -1016,6 +1075,7 @@ default CompletionStage<Void> subscribe(@NonNull StreamSink<? super T> consumer)
10161075
* {@code Streamable} terminates
10171076
* @throws NullPointerException if {@code consumer} or {@code executor} is {@code null}
10181077
*/
1078+
@NonNull
10191079
default CompletionStage<Void> subscribe(@NonNull StreamSink<? super T> consumer, ExecutorService executor) {
10201080
Objects.requireNonNull(consumer, "consumer is null");
10211081
Objects.requireNonNull(executor, "executor is null");
Lines changed: 124 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,124 @@
1+
/*
2+
* Copyright (c) 2016-present, RxJava Contributors.
3+
*
4+
* Licensed under the Apache License, Version 2.0 (the "License"); you may not use this file except in
5+
* compliance with the License. You may obtain a copy of the License at
6+
*
7+
* http://www.apache.org/licenses/LICENSE-2.0
8+
*
9+
* Unless required by applicable law or agreed to in writing, software distributed under the License is
10+
* distributed on an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See
11+
* the License for the specific language governing permissions and limitations under the License.
12+
*/
13+
14+
package io.reactivex.rxjava4.internal.operators.streamable;
15+
16+
import java.io.Serial;
17+
import java.util.Iterator;
18+
import java.util.concurrent.*;
19+
import java.util.concurrent.atomic.AtomicInteger;
20+
21+
import io.reactivex.rxjava4.annotations.NonNull;
22+
import io.reactivex.rxjava4.core.*;
23+
import io.reactivex.rxjava4.disposables.DisposableContainer;
24+
25+
public record StreamableConcatIterable<T>(
26+
Iterable<? extends Streamable<? extends T>> sources,
27+
ErrorMode errorMode
28+
) implements Streamable<T> {
29+
30+
@Override
31+
public @NonNull Streamer<@NonNull T> stream(@NonNull DisposableContainer cancellation) {
32+
return new ConcatIteratorStreamer<>(sources.iterator(), cancellation);
33+
}
34+
35+
static final class ConcatIteratorStreamer<T> extends AtomicInteger implements Streamer<T> {
36+
37+
@Serial
38+
private static final long serialVersionUID = -9136569444189652718L;
39+
40+
final Iterator<? extends Streamable<? extends T>> iterator;
41+
42+
final DisposableContainer cancellation;
43+
44+
DisposableContainer currentCancellation;
45+
46+
Streamer<? extends T> upstream;
47+
48+
CompletableFuture<Boolean> nextReady;
49+
50+
ConcatIteratorStreamer(Iterator<? extends Streamable<? extends T>> iterator,
51+
DisposableContainer cancellation) {
52+
this.iterator = iterator;
53+
this.cancellation = cancellation;
54+
}
55+
56+
@Override
57+
public @NonNull CompletionStage<Boolean> next() {
58+
nextReady = new CompletableFuture<Boolean>();
59+
drain();
60+
return nextReady;
61+
}
62+
63+
@Override
64+
public @NonNull T current() {
65+
return upstream.current();
66+
}
67+
68+
@Override
69+
public @NonNull CompletionStage<Void> finish() {
70+
var localUpstream = upstream;
71+
var localCurrentCancellation = currentCancellation;
72+
upstream = null;
73+
nextReady = null;
74+
currentCancellation = null;
75+
if (localUpstream != null) {
76+
cancellation.delete(localCurrentCancellation);
77+
return localUpstream.finish();
78+
}
79+
return FINISHED;
80+
}
81+
82+
void drain() {
83+
if (getAndIncrement() != 0) {
84+
return;
85+
}
86+
87+
do {
88+
if (upstream == null) {
89+
if (iterator.hasNext()) {
90+
currentCancellation = cancellation.derive();
91+
var nextStreamable = iterator.next();
92+
if (nextStreamable == null) {
93+
nextReady.completeExceptionally(new NullPointerException("The iterator returned a null Streamable"));
94+
} else {
95+
upstream = nextStreamable.stream(currentCancellation);
96+
drain();
97+
}
98+
} else {
99+
nextReady.complete(false);
100+
}
101+
} else {
102+
upstream.next().whenComplete((v, e) -> {
103+
if (e != null) {
104+
nextReady.completeExceptionally(e);
105+
} else
106+
if (v) {
107+
nextReady.complete(true);
108+
} else {
109+
cancellation.delete(currentCancellation);
110+
upstream.finish().whenComplete((_, u) -> {
111+
if (u != null) {
112+
nextReady.completeExceptionally(u);
113+
} else {
114+
upstream = null;
115+
drain();
116+
}
117+
});
118+
}
119+
});
120+
}
121+
} while (decrementAndGet() != 0);
122+
}
123+
}
124+
}

0 commit comments

Comments
 (0)