/
BufferingFlowableWrapper.java
100 lines (89 loc) · 3.22 KB
/
BufferingFlowableWrapper.java
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
/*
* Copyright (c) 2023 Contributors to the Eclipse Foundation
*
* See the NOTICE file(s) distributed with this work for additional
* information regarding copyright ownership.
*
* This program and the accompanying materials are made available under the
* terms of the Eclipse Public License 2.0 which is available at
* http://www.eclipse.org/legal/epl-2.0
*
* SPDX-License-Identifier: EPL-2.0
*/
package org.eclipse.ditto.connectivity.service.messaging.mqtt.hivemq.client;
import io.reactivex.BackpressureStrategy;
import io.reactivex.Flowable;
import io.reactivex.disposables.Disposable;
import io.reactivex.subjects.PublishSubject;
/**
* Wrapper around flowable that buffers the items until it is told to stop.
* When buffering is enabled, all subscribers get all missed items, when the
* buffering is disabled, subscribers will get only new items.
*
* @param <T> type of items
*/
public class BufferingFlowableWrapper<T> implements Disposable {
private final Flowable<T> originalFlowable;
private final PublishSubject<T> buffered;
private final PublishSubject<T> unbuffered;
private final Disposable originalSubscription;
private final Flowable<T> flowable;
private final Disposable subscription;
private boolean isBuffering = true;
private BufferingFlowableWrapper(final Flowable<T> flowable) {
this.originalFlowable = flowable;
this.buffered = PublishSubject.<T>create();
this.unbuffered = PublishSubject.<T>create();
this.flowable = buffered
.replay()
.autoConnect()
.mergeWith(unbuffered)
.toFlowable(BackpressureStrategy.BUFFER);
this.originalSubscription = flowable.subscribe(
x -> (isBuffering ? buffered : unbuffered).onNext(x),
e -> (isBuffering ? buffered : unbuffered).onError(e),
() -> {
buffered.onComplete();
unbuffered.onComplete();
isBuffering = false;
});
this.subscription = this.flowable.subscribe(
ignored -> {},
ignored -> {});
}
/**
* Creates new wrapper around provided {@code Flowable}.
*
* @param flowable flowable to wrap.
* @return wrapper around the provided flowable.
* @param <T> type of items of the flowable.
*/
public static <T> BufferingFlowableWrapper<T> of(final Flowable<T> flowable) {
return new BufferingFlowableWrapper<>(flowable);
}
/**
* @return the {@code Flowable} which can be used to consume messages from original flowable.
*/
public Flowable<T> toFlowable() {
return isBuffering ? this.flowable : this.originalFlowable;
}
/**
* Stops buffering items. All new subscribers using {@code toFlowable} method will get
* only new items.
*/
public void stopBuffering() {
isBuffering = false;
buffered.onComplete();
}
private boolean isDisposed = false;
@Override
public void dispose() {
this.originalSubscription.dispose();
this.subscription.dispose();
isDisposed = true;
}
@Override
public boolean isDisposed() {
return isDisposed;
}
}