This repository was archived by the owner on Oct 27, 2021. It is now read-only.
-
Notifications
You must be signed in to change notification settings - Fork 1
Expand file tree
/
Copy pathLazyConnection.java
More file actions
80 lines (67 loc) · 2.31 KB
/
LazyConnection.java
File metadata and controls
80 lines (67 loc) · 2.31 KB
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
package kr.jadekim.rxjava.websocket.connection;
import io.reactivex.rxjava3.core.Observable;
import io.reactivex.rxjava3.core.ObservableSource;
import io.reactivex.rxjava3.functions.Action;
import io.reactivex.rxjava3.functions.Supplier;
import kr.jadekim.rxjava.websocket.listener.WebSocketEventListener;
import java.util.concurrent.Callable;
public class LazyConnection implements Connection {
private ConnectionFactory connectionFactory;
private String url;
private boolean isErrorPropagation;
private WebSocketEventListener listener;
private volatile Connection connection;
private volatile Observable<String> stream;
public LazyConnection(ConnectionFactory connectionFactory, String url, boolean isErrorPropagation, WebSocketEventListener listener) {
this.connectionFactory = connectionFactory;
this.url = url;
this.isErrorPropagation = isErrorPropagation;
this.listener = listener;
}
@Override
public String getUrl() {
return url;
}
@Override
public Observable<String> getInboundStream() {
return Observable.defer(new Supplier<ObservableSource<? extends String>>() {
@Override
public ObservableSource<? extends String> get() throws Throwable {
if (connection == null) {
connect();
}
return stream;
}
});
}
@Override
public boolean sendMessage(String message) {
if (connection == null) {
connect();
}
return connection.sendMessage(message);
}
@Override
public void disconnect() {
synchronized (this) {
if (connection != null) {
connection.disconnect();
connection = null;
}
}
}
public synchronized Connection connect() {
if (connection == null) {
connection = connectionFactory.connect(url, isErrorPropagation, listener);
stream = connection.getInboundStream()
.doFinally(new Action() {
@Override
public void run() throws Exception {
disconnect();
}
})
.share();
}
return connection;
}
}