Skip to content

Instantly share code, notes, and snippets.

@FireZenk
Created December 21, 2016 12:51
Show Gist options
  • Star 2 You must be signed in to star a gist
  • Fork 1 You must be signed in to fork a gist
  • Save FireZenk/bbaf291bfcfe327c38aa581e37937b9b to your computer and use it in GitHub Desktop.
Save FireZenk/bbaf291bfcfe327c38aa581e37937b9b to your computer and use it in GitHub Desktop.
A working example of reactive implementation of a socket connection using okhttp3 socket api
import java.util.concurrent.TimeUnit;
import okhttp3.OkHttpClient;
import okhttp3.Request;
import okhttp3.Response;
import okhttp3.WebSocket;
import okhttp3.WebSocketListener;
import rx.Observable;
public class RxSocketConnection extends WebSocketListener {
private final BehaviorSubject<String> broadcaster = BehaviorSubject.create();
public Observable<String> subscribe() {
OkHttpClient client = new OkHttpClient.Builder()
.readTimeout(0, TimeUnit.MILLISECONDS)
.build();
Request request = new Request.Builder()
.url("ws://echo.websocket.org")
.build();
client.newWebSocket(request, this);
client.dispatcher().executorService().shutdown();
return broadcaster;
}
@Override public void onOpen(WebSocket webSocket, Response response) {
super.onOpen(webSocket, response);
broadcaster.onNext("Connection open");
}
@Override public void onClosed(WebSocket webSocket, int code, String reason) {
super.onClosed(webSocket, code, reason);
broadcaster.onNext("Connection closed");
}
@Override public void onMessage(WebSocket webSocket, String text) {
super.onMessage(webSocket, text);
broadcaster.onNext(text);
}
@Override public void onFailure(WebSocket webSocket, Throwable throwable, Response response) {
super.onFailure(webSocket, throwable, response);
broadcaster.onError(throwable);
}
}
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment