From 544ac63d70c5eed89a8913ee3bc8f1c2cfec5fb8 Mon Sep 17 00:00:00 2001 From: sachin Date: Sat, 30 Jun 2018 01:06:18 +0530 Subject: [PATCH 1/9] Added multiple disconnect methods to allow disconnection with proper closeCode, closeReason. Added methods to subscribe to all channels and unsubscribe at the same time. --- src/main/java/io/github/sac/Socket.java | 53 ++++++++++++++++++------- 1 file changed, 39 insertions(+), 14 deletions(-) diff --git a/src/main/java/io/github/sac/Socket.java b/src/main/java/io/github/sac/Socket.java index f8c9c0b..ab7aa7d 100644 --- a/src/main/java/io/github/sac/Socket.java +++ b/src/main/java/io/github/sac/Socket.java @@ -4,6 +4,7 @@ import com.neovisionaries.ws.client.StatusLine; import com.neovisionaries.ws.client.WebSocket; import com.neovisionaries.ws.client.WebSocketAdapter; +import com.neovisionaries.ws.client.WebSocketCloseCode; import com.neovisionaries.ws.client.WebSocketException; import com.neovisionaries.ws.client.WebSocketFactory; import com.neovisionaries.ws.client.WebSocketFrame; @@ -60,7 +61,11 @@ private void putDefaultHeaders() { } public Channel createChannel(String name) { - Channel channel = new Channel(name); + return createChannel(name, true); + } + + public Channel createChannel(String name, boolean autoSubscribe) { + Channel channel = new Channel(name, autoSubscribe); channels.add(channel); return channel; } @@ -69,6 +74,18 @@ public List getChannels() { return channels; } + public void subscribeAllChannels() { + for (Channel channel : channels) { + channel.subscribe(); + } + } + + public void unsubscribeAllChannels(){ + for (Channel channel: channels) { + channel.unsubscribe(); + } + } + public Channel getChannelByName(String name) { for (Channel channel : channels) { if (channel.getChannelName().equals(name)) @@ -175,7 +192,7 @@ public void onFrame(WebSocket websocket, WebSocketFrame frame) throws Exception case ISAUTHENTICATED: listener.onAuthentication(Socket.this, ((JSONObject) dataobject).getBoolean("isAuthenticated")); - subscribeChannels(); + subscribeAllChannels(); break; case PUBLISH: Socket.this.handlePublish(((JSONObject) dataobject).getString("channel"), ((JSONObject) dataobject).opt("data")); @@ -416,12 +433,6 @@ public void run() { } - private void subscribeChannels() { - for (Channel channel : channels) { - channel.subscribe(); - } - } - public void setExtraHeaders(Map extraHeaders, boolean overrideDefaultHeaders) { if (overrideDefaultHeaders) { headers.clear(); @@ -534,10 +545,23 @@ public void run() { } public void disconnect() { + disconnect(WebSocketCloseCode.NORMAL, null); + } + + public void disconnect(String closeReason) { + disconnect(WebSocketCloseCode.NORMAL, closeReason); + } + + public void disconnect(int closeCode, String closeReason) { + disconnect(closeCode, closeReason, -1); + } + + public void disconnect(int closeCode, String closeReason, long closeDelay) { + unsubscribeAllChannels(); + strategy.setAttemptsMade(strategy.maxAttempts); if (ws != null) { - ws.disconnect(); + ws.disconnect(closeCode, closeReason, closeDelay); } - strategy = null; } /** @@ -549,10 +573,11 @@ public void disconnect() { * OPEN */ - public WebSocketState getCurrentState() { + public WebSocketState getSocketStatus() { return ws.getState(); } + public Boolean isconnected() { return ws != null && ws.getState() == WebSocketState.OPEN; } @@ -569,13 +594,15 @@ public void disableLogging() { public class Channel { String channelName; + boolean autoSubscribe; public String getChannelName() { return channelName; } - public Channel(String channelName) { + public Channel(String channelName, boolean autoSubscribe) { this.channelName = channelName; + this.autoSubscribe = autoSubscribe; } public void subscribe() { @@ -600,12 +627,10 @@ public void publish(Object data, Ack ack) { public void unsubscribe() { Socket.this.unsubscribe(channelName); - channels.remove(this); } public void unsubscribe(Ack ack) { Socket.this.unsubscribe(channelName, ack); - channels.remove(this); } } From 6506aca90f92efd00c38ce01b03e0d9eefecdbfa Mon Sep 17 00:00:00 2001 From: sachin Date: Sat, 30 Jun 2018 01:33:15 +0530 Subject: [PATCH 2/9] Sachin | Created new enum for managing authState, Added method for getting authentication state of the socket --- src/main/java/io/github/sac/Socket.java | 15 ++++++++++++++- 1 file changed, 14 insertions(+), 1 deletion(-) diff --git a/src/main/java/io/github/sac/Socket.java b/src/main/java/io/github/sac/Socket.java index ab7aa7d..d1ef9fb 100644 --- a/src/main/java/io/github/sac/Socket.java +++ b/src/main/java/io/github/sac/Socket.java @@ -41,6 +41,7 @@ public class Socket extends Emitter { private List channels; private WebSocketAdapter adapter; private Map headers; + private AuthState authState; public Socket(String URL) { this.URL = URL; @@ -121,6 +122,11 @@ public void setAuthToken(String token) { AuthToken = token; } + public String getAuthToken() { + return AuthToken; + } + + public WebSocketAdapter getAdapter() { return new WebSocketAdapter() { @@ -191,7 +197,9 @@ public void onFrame(WebSocket websocket, WebSocketFrame frame) throws Exception switch (Parser.parse(dataobject, event)) { case ISAUTHENTICATED: - listener.onAuthentication(Socket.this, ((JSONObject) dataobject).getBoolean("isAuthenticated")); + boolean isAuthenticated = ((JSONObject) dataobject).getBoolean("isAuthenticated"); + authState = isAuthenticated ? AuthState.AUTHENTICATED : AuthState.UNAUTHENTICATED; + listener.onAuthentication(Socket.this, isAuthenticated); subscribeAllChannels(); break; case PUBLISH: @@ -634,6 +642,11 @@ public void unsubscribe(Ack ack) { } } + public enum AuthState { + AUTHENTICATED, + UNAUTHENTICATED + } + @Override protected void finalize() throws Throwable { ws.disconnect("Client socket garbage collected, closing connection"); From bb0d62d01d4f6db542be7246f7ee2eacb25e282e Mon Sep 17 00:00:00 2001 From: sachin Date: Tue, 3 Jul 2018 01:04:39 +0530 Subject: [PATCH 3/9] Added code to return current authState of socket. --- src/main/java/io/github/sac/Socket.java | 3 +++ 1 file changed, 3 insertions(+) diff --git a/src/main/java/io/github/sac/Socket.java b/src/main/java/io/github/sac/Socket.java index d1ef9fb..4078eba 100644 --- a/src/main/java/io/github/sac/Socket.java +++ b/src/main/java/io/github/sac/Socket.java @@ -126,6 +126,9 @@ public String getAuthToken() { return AuthToken; } + public AuthState getAuthState() { + return authState; + } public WebSocketAdapter getAdapter() { return new WebSocketAdapter() { From 176d1f39f276af84ecf070008cd6b118fc8ca1e8 Mon Sep 17 00:00:00 2001 From: sachin Date: Tue, 3 Jul 2018 01:24:33 +0530 Subject: [PATCH 4/9] Added ConnectionState to represent combined state of authentication and socket --- src/main/java/io/github/sac/Socket.java | 34 +++++++++++++++++++++++++ 1 file changed, 34 insertions(+) diff --git a/src/main/java/io/github/sac/Socket.java b/src/main/java/io/github/sac/Socket.java index 4078eba..3b071a7 100644 --- a/src/main/java/io/github/sac/Socket.java +++ b/src/main/java/io/github/sac/Socket.java @@ -588,6 +588,29 @@ public WebSocketState getSocketStatus() { return ws.getState(); } + public SocketState getConnectionState() { + switch (getSocketStatus()) { + case CREATED: + return SocketState.CREATED; + case CONNECTING: + return SocketState.CONNECTING; + case OPEN: + return SocketState.OPEN; + case CLOSING: + return SocketState.CLOSING; + case CLOSED: + return SocketState.CLOSED; + } + + switch (getAuthState()) { + case AUTHENTICATED: + return SocketState.AUTHENTICATED; + case UNAUTHENTICATED: + return SocketState.UNAUTHENTICATED; + } + return SocketState.NOTFOUND; + } + public Boolean isconnected() { return ws != null && ws.getState() == WebSocketState.OPEN; @@ -650,6 +673,17 @@ public enum AuthState { UNAUTHENTICATED } + public enum SocketState{ + CREATED, + CONNECTING, + OPEN, + CLOSING, + CLOSED, + AUTHENTICATED, + UNAUTHENTICATED, + NOTFOUND + } + @Override protected void finalize() throws Throwable { ws.disconnect("Client socket garbage collected, closing connection"); From 9190e615fac7200cb23bb3b9fa040dfcdbf1198d Mon Sep 17 00:00:00 2001 From: sachin Date: Tue, 3 Jul 2018 01:44:47 +0530 Subject: [PATCH 5/9] Added set of events that can be used as a lambda in java8, as well as can be easily inferred in any background service --- src/main/java/io/github/sac/Socket.java | 1 + .../io/github/sac/events/AuthenticationEvent.java | 11 +++++++++++ .../io/github/sac/events/ChannelKickoutEvent.java | 10 ++++++++++ src/main/java/io/github/sac/events/ConnectEvent.java | 12 ++++++++++++ .../java/io/github/sac/events/ConnectionAbort.java | 11 +++++++++++ .../java/io/github/sac/events/DisconnectEvent.java | 12 ++++++++++++ src/main/java/io/github/sac/events/ErrorEvent.java | 11 +++++++++++ src/main/java/io/github/sac/events/MessageEvent.java | 8 ++++++++ src/main/java/io/github/sac/events/RawEvent.java | 9 +++++++++ .../io/github/sac/events/SubscribeStateEvent.java | 8 ++++++++ 10 files changed, 93 insertions(+) create mode 100644 src/main/java/io/github/sac/events/AuthenticationEvent.java create mode 100644 src/main/java/io/github/sac/events/ChannelKickoutEvent.java create mode 100644 src/main/java/io/github/sac/events/ConnectEvent.java create mode 100644 src/main/java/io/github/sac/events/ConnectionAbort.java create mode 100644 src/main/java/io/github/sac/events/DisconnectEvent.java create mode 100644 src/main/java/io/github/sac/events/ErrorEvent.java create mode 100644 src/main/java/io/github/sac/events/MessageEvent.java create mode 100644 src/main/java/io/github/sac/events/RawEvent.java create mode 100644 src/main/java/io/github/sac/events/SubscribeStateEvent.java diff --git a/src/main/java/io/github/sac/Socket.java b/src/main/java/io/github/sac/Socket.java index 3b071a7..fd2bcdd 100644 --- a/src/main/java/io/github/sac/Socket.java +++ b/src/main/java/io/github/sac/Socket.java @@ -684,6 +684,7 @@ public enum SocketState{ NOTFOUND } + @Override protected void finalize() throws Throwable { ws.disconnect("Client socket garbage collected, closing connection"); diff --git a/src/main/java/io/github/sac/events/AuthenticationEvent.java b/src/main/java/io/github/sac/events/AuthenticationEvent.java new file mode 100644 index 0000000..459b811 --- /dev/null +++ b/src/main/java/io/github/sac/events/AuthenticationEvent.java @@ -0,0 +1,11 @@ +package io.github.sac.events; + +import io.github.sac.Socket; + +/** + * Created by sachin on 3/7/18. + */ +public interface AuthenticationEvent { + void onAuthenticated(Socket socket, String token); + void onDeauthentication(Socket socket); +} diff --git a/src/main/java/io/github/sac/events/ChannelKickoutEvent.java b/src/main/java/io/github/sac/events/ChannelKickoutEvent.java new file mode 100644 index 0000000..aa0749d --- /dev/null +++ b/src/main/java/io/github/sac/events/ChannelKickoutEvent.java @@ -0,0 +1,10 @@ +package io.github.sac.events; + +import io.github.sac.Socket; + +/** + * Created by sachin on 3/7/18. + */ +public interface ChannelKickoutEvent { + void onChannelKickout(Socket socket, String message, String channelName); +} diff --git a/src/main/java/io/github/sac/events/ConnectEvent.java b/src/main/java/io/github/sac/events/ConnectEvent.java new file mode 100644 index 0000000..c616a39 --- /dev/null +++ b/src/main/java/io/github/sac/events/ConnectEvent.java @@ -0,0 +1,12 @@ +package io.github.sac.events; + +import io.github.sac.Socket; +import java.util.List; +import java.util.Map; + +/** + * Created by sachin on 3/7/18. + */ +public interface ConnectEvent { + void onConnected(Socket socket, Map> headers); +} diff --git a/src/main/java/io/github/sac/events/ConnectionAbort.java b/src/main/java/io/github/sac/events/ConnectionAbort.java new file mode 100644 index 0000000..3e69739 --- /dev/null +++ b/src/main/java/io/github/sac/events/ConnectionAbort.java @@ -0,0 +1,11 @@ +package io.github.sac.events; + +import com.neovisionaries.ws.client.WebSocketException; +import io.github.sac.Socket; + +/** + * Created by sachin on 3/7/18. + */ +public interface ConnectionAbort { + void onConnectionAbort(Socket socket, WebSocketException exception); +} diff --git a/src/main/java/io/github/sac/events/DisconnectEvent.java b/src/main/java/io/github/sac/events/DisconnectEvent.java new file mode 100644 index 0000000..4738dc7 --- /dev/null +++ b/src/main/java/io/github/sac/events/DisconnectEvent.java @@ -0,0 +1,12 @@ +package io.github.sac.events; + +import com.neovisionaries.ws.client.WebSocketFrame; +import io.github.sac.Socket; + +/** + * Created by sachin on 3/7/18. + */ + +public interface DisconnectEvent { + void onDisconnected(Socket socket, WebSocketFrame serverCloseFrame, WebSocketFrame clientCloseFrame, boolean closedByServer); +} diff --git a/src/main/java/io/github/sac/events/ErrorEvent.java b/src/main/java/io/github/sac/events/ErrorEvent.java new file mode 100644 index 0000000..c8f7d9f --- /dev/null +++ b/src/main/java/io/github/sac/events/ErrorEvent.java @@ -0,0 +1,11 @@ +package io.github.sac.events; + +import io.github.sac.Socket; + +/** + * Created by sachin on 3/7/18. + */ + +public interface ErrorEvent { + void onError(Socket socket, Exception error); +} diff --git a/src/main/java/io/github/sac/events/MessageEvent.java b/src/main/java/io/github/sac/events/MessageEvent.java new file mode 100644 index 0000000..63d661e --- /dev/null +++ b/src/main/java/io/github/sac/events/MessageEvent.java @@ -0,0 +1,8 @@ +package io.github.sac.events; + +/** + * Created by sachin on 3/7/18. + */ +public interface MessageEvent { + void onMessage(String name, Object error, Object data); +} diff --git a/src/main/java/io/github/sac/events/RawEvent.java b/src/main/java/io/github/sac/events/RawEvent.java new file mode 100644 index 0000000..84426e6 --- /dev/null +++ b/src/main/java/io/github/sac/events/RawEvent.java @@ -0,0 +1,9 @@ +package io.github.sac.events; + +/** + * Created by sachin on 3/7/18. + */ + +public interface RawEvent { + void onRawEvent(Object error, Object data); +} diff --git a/src/main/java/io/github/sac/events/SubscribeStateEvent.java b/src/main/java/io/github/sac/events/SubscribeStateEvent.java new file mode 100644 index 0000000..650aa78 --- /dev/null +++ b/src/main/java/io/github/sac/events/SubscribeStateEvent.java @@ -0,0 +1,8 @@ +package io.github.sac.events; + +/** + * Created by sachin on 3/7/18. + */ +public interface SubscribeStateEvent { + void onSubscribeStateChange(); +} From c0ea98e457a9365833728fb94c0f463b0d32201d Mon Sep 17 00:00:00 2001 From: sachin Date: Tue, 3 Jul 2018 01:50:59 +0530 Subject: [PATCH 6/9] Added channelstate to represent internal state of channels, added method signature to onSubscribeStateChange --- src/main/java/io/github/sac/Socket.java | 5 +++++ src/main/java/io/github/sac/events/SubscribeStateEvent.java | 4 +++- 2 files changed, 8 insertions(+), 1 deletion(-) diff --git a/src/main/java/io/github/sac/Socket.java b/src/main/java/io/github/sac/Socket.java index fd2bcdd..f18c0ba 100644 --- a/src/main/java/io/github/sac/Socket.java +++ b/src/main/java/io/github/sac/Socket.java @@ -684,6 +684,11 @@ public enum SocketState{ NOTFOUND } + public enum ChannelState { + SUBSCRIBED, + PENDING, + UNSUBSCRIBED + } @Override protected void finalize() throws Throwable { diff --git a/src/main/java/io/github/sac/events/SubscribeStateEvent.java b/src/main/java/io/github/sac/events/SubscribeStateEvent.java index 650aa78..ee4437b 100644 --- a/src/main/java/io/github/sac/events/SubscribeStateEvent.java +++ b/src/main/java/io/github/sac/events/SubscribeStateEvent.java @@ -1,8 +1,10 @@ package io.github.sac.events; +import io.github.sac.Socket; + /** * Created by sachin on 3/7/18. */ public interface SubscribeStateEvent { - void onSubscribeStateChange(); + void onSubscribeStateChange(Socket.Channel channel, Socket.ChannelState oldState, Socket.ChannelState newState); } From 32871e64eca4dce50016e523408effdc4367afc1 Mon Sep 17 00:00:00 2001 From: sachin Date: Wed, 4 Jul 2018 01:55:18 +0530 Subject: [PATCH 7/9] Updated Socket Constructor to more meaningful format, Now most of the settings can be passed from constructor. Added default channel state to unsubscribed. Added code in emitter to support multipleListeners and multipleChannelWatchers. --- src/main/java/io/github/sac/Emitter.java | 153 ++++++++++++++--------- src/main/java/io/github/sac/Socket.java | 40 ++++-- 2 files changed, 126 insertions(+), 67 deletions(-) diff --git a/src/main/java/io/github/sac/Emitter.java b/src/main/java/io/github/sac/Emitter.java index cdf49b1..6eb5e90 100644 --- a/src/main/java/io/github/sac/Emitter.java +++ b/src/main/java/io/github/sac/Emitter.java @@ -4,15 +4,40 @@ * Created by sachin on 13/11/16. */ +import java.util.Iterator; import java.util.Map; import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ConcurrentLinkedQueue; public class Emitter { + private boolean multipleListenersEnabled; + private boolean multipleChannelWatchersEnabled; - private ConcurrentHashMap singlecallbacks = new ConcurrentHashMap<>(); - private ConcurrentHashMap singleackcallbacks = new ConcurrentHashMap<>(); - private ConcurrentHashMap publishcallbacks = new ConcurrentHashMap<>(); + public Emitter(boolean multipleListenersEnabled, boolean multipleChannelWatchersEnabled) { + this.multipleListenersEnabled = multipleListenersEnabled; + this.multipleChannelWatchersEnabled = multipleChannelWatchersEnabled; + } + + public void setMultipleListenersEnabled(boolean multipleListenersEnabled) { + this.multipleListenersEnabled = multipleListenersEnabled; + } + + public boolean isMultipleListenersEnabled() { + return multipleListenersEnabled; + } + + public boolean isMultipleChannelWatchersEnabled() { + return multipleChannelWatchersEnabled; + } + + public void setMultipleChannelWatchersEnabled(boolean multipleChannelWatchersEnabled) { + this.multipleChannelWatchersEnabled = multipleChannelWatchersEnabled; + } + + private ConcurrentHashMap> listeners = new ConcurrentHashMap<>(); + private ConcurrentHashMap> ackListeners = new ConcurrentHashMap<>(); + private ConcurrentHashMap> channelObservers = new ConcurrentHashMap<>(); /** * Listens on the event. @@ -21,97 +46,113 @@ public class Emitter { * @return a reference to this object. */ public Emitter on(String event, Listener fn) { + return on(event, fn, multipleListenersEnabled); + } - if (singlecallbacks.containsKey(event)) { - singlecallbacks.remove(event); - } - singlecallbacks.put(event, fn); + public Emitter on(String event, Listener fn, boolean multiListenersEnabled) { + registerEvent(listeners, event, fn, multiListenersEnabled); return this; } - public Emitter onSubscribe(String event, Listener fn) { + public Emitter on(String event, AckListener fn) { + return on(event, fn, multipleListenersEnabled); + } - if (publishcallbacks.containsKey(event)) { - publishcallbacks.remove(event); - } - publishcallbacks.put(event, fn); + public Emitter on(String event, AckListener fn, boolean multiListenersEnabled) { + registerEvent(ackListeners, event, fn, multiListenersEnabled); return this; } - public Emitter on(String event, AckListener fn) { - if (singleackcallbacks.containsKey(event)) { - singleackcallbacks.remove(event); - } - singleackcallbacks.put(event, fn); - return this; + public Emitter onSubscribe(String event, Listener fn) { + return onSubscribe(event, fn, multipleChannelWatchersEnabled); } + public Emitter onSubscribe(String event, Listener fn, boolean multipleChannelWatchersEnabled) { + registerEvent(channelObservers, event, fn, multipleChannelWatchersEnabled); + return this; + } - public Emitter handleEmit(String event, Object object) { - Listener listener = singlecallbacks.get(event); - if (listener != null) { - listener.call(event, object); + private static void registerEvent(ConcurrentHashMap> listeners, String event, T fn, boolean multiEnabled) { + if (listeners.containsKey(event)) { + if (!multiEnabled) { + listeners.get(event).clear(); + } + listeners.get(event).add(fn); + return; } - return this; + ConcurrentLinkedQueue linkedListeners = new ConcurrentLinkedQueue<>(); + linkedListeners.add(fn); + listeners.put(event, linkedListeners); } - public Emitter handlePublish(String event, Object object) { - - Listener listener = publishcallbacks.get(event); - if (listener != null) { - listener.call(event, object); - } + public Emitter handleEmit(String event, Object object) { + handleEvent(listeners, event, object, null); return this; } - public boolean hasEventAck(String event) { - return this.singleackcallbacks.get(event) != null; - } public Emitter handleEmitAck(String event, Object object, Ack ack) { + handleEvent(ackListeners, event, object, ack); + return this; + } - AckListener listener = singleackcallbacks.get(event); - if (listener != null) { - listener.call(event, object, ack); - } + public Emitter handlePublish(String event, Object object) { + handleEvent(channelObservers, event, object, null); return this; } - public interface Listener { - void call(String name, Object data); + public static void handleEvent(ConcurrentHashMap> listeners, String event, Object object, Ack ack) { + Iterator listenerIterator = listeners.get(event).iterator(); + while (listenerIterator.hasNext()) { + T listener = listenerIterator.next(); + if (listener instanceof Listener) { + ((Listener) listener).call(event, object); + } else { + ((AckListener) listener).call(event, object, ack); + } + } } - public interface AckListener { - void call(String name, Object data, Ack ack); + public void off(String event) { + listeners.remove(event); + ackListeners.remove(event); } - /** - * New methods ADDED - */ - - public void removeEmitCallback(String event) { - singlecallbacks.remove(event); - singleackcallbacks.remove(event); + public void off(String event, Listener listener) { + if (listeners.containsKey(event)) { + listeners.get(event).remove(listener); + } } - public void removeSubscribeCallback(String event) { - publishcallbacks.remove(event); + public void off(String event, AckListener ackListener) { + if (ackListeners.containsKey(event)) { + ackListeners.get(event).remove(ackListener); + } } - public void removeAllCallbacks() { - for (Map.Entry e : singlecallbacks.entrySet()) { - singlecallbacks.remove(e.getKey().toString()); + public void removeAllListeners() { + for (Map.Entry e : listeners.entrySet()) { + listeners.remove(e.getKey().toString()); } - for (Map.Entry e : singleackcallbacks.entrySet()) { - singleackcallbacks.remove(e.getKey().toString()); + for (Map.Entry e : ackListeners.entrySet()) { + ackListeners.remove(e.getKey().toString()); } - for (Map.Entry e : publishcallbacks.entrySet()) { - publishcallbacks.remove(e.getKey().toString()); + for (Map.Entry e : channelObservers.entrySet()) { + channelObservers.remove(e.getKey().toString()); } } + + public interface Listener { + void call(String name, Object data); + } + + public interface AckListener { + void call(String name, Object data, Ack ack); + } + } diff --git a/src/main/java/io/github/sac/Socket.java b/src/main/java/io/github/sac/Socket.java index f18c0ba..c31b5e3 100644 --- a/src/main/java/io/github/sac/Socket.java +++ b/src/main/java/io/github/sac/Socket.java @@ -44,7 +44,27 @@ public class Socket extends Emitter { private AuthState authState; public Socket(String URL) { + this(URL, null); + } + + public Socket(String URL, BasicListener listener) { + this(URL, listener, null); + } + + public Socket(String URL, BasicListener listener, String authToken) { + this(URL, listener, authToken, null); + } + + public Socket(String URL, BasicListener listener, String authToken, ReconnectStrategy reconnectStrategy) { + this(URL, listener, authToken, reconnectStrategy, false, false); + } + + public Socket(String URL, BasicListener listener, String AuthToken, ReconnectStrategy reconnectStrategy, boolean multipleListenersEnabled, boolean multipleChannelWatchersEnabled) { + super(multipleListenersEnabled, multipleChannelWatchersEnabled); this.URL = URL; + this.listener = listener; + this.AuthToken = AuthToken; + strategy = reconnectStrategy; factory = new WebSocketFactory().setConnectionTimeout(5000); counter = new AtomicInteger(1); acks = new HashMap<>(); @@ -81,8 +101,8 @@ public void subscribeAllChannels() { } } - public void unsubscribeAllChannels(){ - for (Channel channel: channels) { + public void unsubscribeAllChannels() { + for (Channel channel : channels) { channel.unsubscribe(); } } @@ -107,9 +127,10 @@ public void setListener(BasicListener listener) { this.listener = listener; } - public Logger getLogger(){ + public Logger getLogger() { return logger; } + /** * used to set up TLS/SSL connection to server for more details visit neovisionaries websocket client */ @@ -190,7 +211,6 @@ public void onFrame(WebSocket websocket, WebSocketFrame frame) throws Exception */ logger.info("Message :" + object.toString()); - try { Object dataobject = object.opt("data"); Integer rid = (Integer) object.opt("rid"); @@ -217,12 +237,8 @@ public void onFrame(WebSocket websocket, WebSocketFrame frame) throws Exception listener.onSetAuthToken(token, Socket.this); break; case EVENT: - if (hasEventAck(event)) { - handleEmitAck(event, dataobject, ack(Long.valueOf(cid))); - } else { - Socket.this.handleEmit(event, dataobject); - - } + handleEmitAck(event, dataobject, ack(Long.valueOf(cid))); + handleEmit(event, dataobject); break; case ACKRECEIVE: if (acks.containsKey((long) rid)) { @@ -629,6 +645,7 @@ public class Channel { String channelName; boolean autoSubscribe; + ChannelState channelState; public String getChannelName() { return channelName; @@ -637,6 +654,7 @@ public String getChannelName() { public Channel(String channelName, boolean autoSubscribe) { this.channelName = channelName; this.autoSubscribe = autoSubscribe; + this.channelState = ChannelState.UNSUBSCRIBED; } public void subscribe() { @@ -673,7 +691,7 @@ public enum AuthState { UNAUTHENTICATED } - public enum SocketState{ + public enum SocketState { CREATED, CONNECTING, OPEN, From c78f2e05c72a41af3367e79700b9cc6642fb2540 Mon Sep 17 00:00:00 2001 From: sachin Date: Fri, 6 Jul 2018 02:07:26 +0530 Subject: [PATCH 8/9] Added code to send and handle raw event with server. Added code to handle any type of message. Added error event handler to handle all types of errors --- src/main/java/io/github/sac/Emitter.java | 22 ++++- src/main/java/io/github/sac/Socket.java | 87 ++++++++++++++++--- .../io/github/sac/events/ConnectEvent.java | 12 --- .../io/github/sac/events/ConnectionAbort.java | 11 --- .../io/github/sac/events/DisconnectEvent.java | 12 --- .../io/github/sac/events/MessageEvent.java | 8 -- .../java/io/github/sac/events/RawEvent.java | 9 -- 7 files changed, 96 insertions(+), 65 deletions(-) delete mode 100644 src/main/java/io/github/sac/events/ConnectEvent.java delete mode 100644 src/main/java/io/github/sac/events/ConnectionAbort.java delete mode 100644 src/main/java/io/github/sac/events/DisconnectEvent.java delete mode 100644 src/main/java/io/github/sac/events/MessageEvent.java delete mode 100644 src/main/java/io/github/sac/events/RawEvent.java diff --git a/src/main/java/io/github/sac/Emitter.java b/src/main/java/io/github/sac/Emitter.java index 6eb5e90..bb585f8 100644 --- a/src/main/java/io/github/sac/Emitter.java +++ b/src/main/java/io/github/sac/Emitter.java @@ -11,6 +11,9 @@ public class Emitter { + public static String RAWEVENT = "raw"; + public static String MESSAGEEVENT = "message"; + private boolean multipleListenersEnabled; private boolean multipleChannelWatchersEnabled; @@ -45,6 +48,19 @@ public void setMultipleChannelWatchersEnabled(boolean multipleChannelWatchersEna * @param event event name. * @return a reference to this object. */ + + public Emitter onRawEvent(Listener fn) { + return on(RAWEVENT, fn); + } + + public Emitter onAnyMessage(Listener fn) { + return on(MESSAGEEVENT, fn); + } + + public Emitter onAnyMessage(AckListener fn) { + return on(MESSAGEEVENT, fn); + } + public Emitter on(String event, Listener fn) { return on(event, fn, multipleListenersEnabled); } @@ -105,7 +121,11 @@ public Emitter handlePublish(String event, Object object) { public static void handleEvent(ConcurrentHashMap> listeners, String event, Object object, Ack ack) { - Iterator listenerIterator = listeners.get(event).iterator(); + InvokeListeners(event, object, ack, listeners.get(event).iterator()); + InvokeListeners(event, object, ack, listeners.get(MESSAGEEVENT).iterator()); + } + + public static void InvokeListeners(String event, Object object, Ack ack, Iterator listenerIterator) { while (listenerIterator.hasNext()) { T listener = listenerIterator.next(); if (listener instanceof Listener) { diff --git a/src/main/java/io/github/sac/Socket.java b/src/main/java/io/github/sac/Socket.java index c31b5e3..0ffe7ef 100644 --- a/src/main/java/io/github/sac/Socket.java +++ b/src/main/java/io/github/sac/Socket.java @@ -9,6 +9,10 @@ import com.neovisionaries.ws.client.WebSocketFactory; import com.neovisionaries.ws.client.WebSocketFrame; import com.neovisionaries.ws.client.WebSocketState; +import io.github.sac.events.AuthenticationEvent; +import io.github.sac.events.ChannelKickoutEvent; +import io.github.sac.events.ErrorEvent; +import io.github.sac.events.SubscribeStateEvent; import java.io.IOException; import java.util.ArrayList; import java.util.HashMap; @@ -43,6 +47,13 @@ public class Socket extends Emitter { private Map headers; private AuthState authState; + // Extra definition of events + public AuthenticationEvent authenticationEventHandler; + public ChannelKickoutEvent channelKickoutEventHandler; + public ErrorEvent errorEventHandler; + public SubscribeStateEvent subscribeStateEventHandler; + + public Socket(String URL) { this(URL, null); } @@ -255,6 +266,9 @@ public void onFrame(WebSocket websocket, WebSocketFrame frame) throws Exception break; } } catch (Exception e) { + if (errorEventHandler != null) { + errorEventHandler.onError(Socket.this, e); + } logger.severe(e.toString()); } @@ -279,6 +293,10 @@ public void onSendError(WebSocket websocket, WebSocketException cause, WebSocket } + public Socket send(final Object object) { + return emit(RAWEVENT, object); + } + public Socket emit(final String event, final Object object) { EventThread.exec(new Runnable() { public void run() { @@ -288,6 +306,9 @@ public void run() { eventObject.put("data", object); } catch (JSONException e) { e.printStackTrace(); + if (errorEventHandler != null) { + errorEventHandler.onError(Socket.this, e); + } } ws.sendText(eventObject.toString()); } @@ -308,6 +329,9 @@ public void run() { eventObject.put("cid", counter.getAndIncrement()); } catch (JSONException e) { e.printStackTrace(); + if (errorEventHandler != null) { + errorEventHandler.onError(Socket.this, e); + } } ws.sendText(eventObject.toString()); } @@ -328,6 +352,9 @@ public void run() { subscribeObject.put("cid", counter.getAndIncrement()); } catch (JSONException e) { e.printStackTrace(); + if (errorEventHandler != null) { + errorEventHandler.onError(Socket.this, e); + } } ws.sendText(subscribeObject.toString()); } @@ -353,6 +380,9 @@ public void run() { subscribeObject.put("cid", counter.getAndIncrement()); } catch (JSONException e) { e.printStackTrace(); + if (errorEventHandler != null) { + errorEventHandler.onError(Socket.this, e); + } } ws.sendText(subscribeObject.toString()); } @@ -370,6 +400,9 @@ public void run() { subscribeObject.put("cid", counter.getAndIncrement()); } catch (JSONException e) { e.printStackTrace(); + if (errorEventHandler != null) { + errorEventHandler.onError(Socket.this, e); + } } ws.sendText(subscribeObject.toString()); } @@ -389,6 +422,9 @@ public void run() { subscribeObject.put("cid", counter.getAndIncrement()); } catch (JSONException e) { e.printStackTrace(); + if (errorEventHandler != null) { + errorEventHandler.onError(Socket.this, e); + } } ws.sendText(subscribeObject.toString()); } @@ -409,6 +445,9 @@ public void run() { publishObject.put("cid", counter.getAndIncrement()); } catch (JSONException e) { e.printStackTrace(); + if (errorEventHandler != null) { + errorEventHandler.onError(Socket.this, e); + } } ws.sendText(publishObject.toString()); } @@ -431,6 +470,9 @@ public void run() { publishObject.put("cid", counter.getAndIncrement()); } catch (JSONException e) { e.printStackTrace(); + if (errorEventHandler != null) { + errorEventHandler.onError(Socket.this, e); + } } ws.sendText(publishObject.toString()); } @@ -451,6 +493,9 @@ public void run() { object.put("rid", cid); } catch (JSONException e) { e.printStackTrace(); + if (errorEventHandler != null) { + errorEventHandler.onError(Socket.this, e); + } } ws.sendText(object.toString()); } @@ -473,17 +518,7 @@ public Map getHeaders() { } public void connect() { - - try { - ws = factory.createSocket(URL); - } catch (IOException e) { - logger.severe(e.toString()); - } - ws.addExtension("permessage-deflate; client_max_window_bits"); - for (Map.Entry entry : headers.entrySet()) { - ws.addHeader(entry.getKey(), entry.getValue()); - } - + CreateSocket(); ws.addListener(adapter); try { @@ -491,6 +526,9 @@ public void connect() { } catch (OpeningHandshakeException e) { // A violation against the WebSocket protocol was detected // during the opening handshake. + if (errorEventHandler != null) { + errorEventHandler.onError(Socket.this, e); + } logger.severe(e.toString()); // Status line. @@ -523,22 +561,31 @@ public void connect() { } } catch (WebSocketException e) { listener.onConnectError(Socket.this, e); + if (errorEventHandler != null) { + errorEventHandler.onError(Socket.this, e); + } reconnect(); } } - public void connectAsync() { + private void CreateSocket() { try { ws = factory.createSocket(URL); } catch (IOException e) { logger.severe(e.toString()); + if (errorEventHandler != null) { + errorEventHandler.onError(Socket.this, e); + } } ws.addExtension("permessage-deflate; client_max_window_bits"); for (Map.Entry entry : headers.entrySet()) { ws.addHeader(entry.getKey(), entry.getValue()); } + } + public void connectAsync() { + CreateSocket(); ws.addListener(adapter); ws.connectAsynchronously(); } @@ -708,6 +755,22 @@ public enum ChannelState { UNSUBSCRIBED } + public void setAuthenticationEventHandler(AuthenticationEvent authenticationEventHandler) { + this.authenticationEventHandler = authenticationEventHandler; + } + + public void setChannelKickoutEventHandler(ChannelKickoutEvent channelKickoutEventHandler) { + this.channelKickoutEventHandler = channelKickoutEventHandler; + } + + public void setErrorEventHandler(ErrorEvent errorEventHandler) { + this.errorEventHandler = errorEventHandler; + } + + public void setSubscribeStateEventHandler(SubscribeStateEvent subscribeStateEventHandler) { + this.subscribeStateEventHandler = subscribeStateEventHandler; + } + @Override protected void finalize() throws Throwable { ws.disconnect("Client socket garbage collected, closing connection"); diff --git a/src/main/java/io/github/sac/events/ConnectEvent.java b/src/main/java/io/github/sac/events/ConnectEvent.java deleted file mode 100644 index c616a39..0000000 --- a/src/main/java/io/github/sac/events/ConnectEvent.java +++ /dev/null @@ -1,12 +0,0 @@ -package io.github.sac.events; - -import io.github.sac.Socket; -import java.util.List; -import java.util.Map; - -/** - * Created by sachin on 3/7/18. - */ -public interface ConnectEvent { - void onConnected(Socket socket, Map> headers); -} diff --git a/src/main/java/io/github/sac/events/ConnectionAbort.java b/src/main/java/io/github/sac/events/ConnectionAbort.java deleted file mode 100644 index 3e69739..0000000 --- a/src/main/java/io/github/sac/events/ConnectionAbort.java +++ /dev/null @@ -1,11 +0,0 @@ -package io.github.sac.events; - -import com.neovisionaries.ws.client.WebSocketException; -import io.github.sac.Socket; - -/** - * Created by sachin on 3/7/18. - */ -public interface ConnectionAbort { - void onConnectionAbort(Socket socket, WebSocketException exception); -} diff --git a/src/main/java/io/github/sac/events/DisconnectEvent.java b/src/main/java/io/github/sac/events/DisconnectEvent.java deleted file mode 100644 index 4738dc7..0000000 --- a/src/main/java/io/github/sac/events/DisconnectEvent.java +++ /dev/null @@ -1,12 +0,0 @@ -package io.github.sac.events; - -import com.neovisionaries.ws.client.WebSocketFrame; -import io.github.sac.Socket; - -/** - * Created by sachin on 3/7/18. - */ - -public interface DisconnectEvent { - void onDisconnected(Socket socket, WebSocketFrame serverCloseFrame, WebSocketFrame clientCloseFrame, boolean closedByServer); -} diff --git a/src/main/java/io/github/sac/events/MessageEvent.java b/src/main/java/io/github/sac/events/MessageEvent.java deleted file mode 100644 index 63d661e..0000000 --- a/src/main/java/io/github/sac/events/MessageEvent.java +++ /dev/null @@ -1,8 +0,0 @@ -package io.github.sac.events; - -/** - * Created by sachin on 3/7/18. - */ -public interface MessageEvent { - void onMessage(String name, Object error, Object data); -} diff --git a/src/main/java/io/github/sac/events/RawEvent.java b/src/main/java/io/github/sac/events/RawEvent.java deleted file mode 100644 index 84426e6..0000000 --- a/src/main/java/io/github/sac/events/RawEvent.java +++ /dev/null @@ -1,9 +0,0 @@ -package io.github.sac.events; - -/** - * Created by sachin on 3/7/18. - */ - -public interface RawEvent { - void onRawEvent(Object error, Object data); -} From 851c588e91a667abc4f65d89a0dc0a63bc14f1ac Mon Sep 17 00:00:00 2001 From: sachin Date: Fri, 6 Jul 2018 02:13:57 +0530 Subject: [PATCH 9/9] Added todo for subscribe and unsubscribe on parser --- src/main/java/io/github/sac/Parser.java | 1 + 1 file changed, 1 insertion(+) diff --git a/src/main/java/io/github/sac/Parser.java b/src/main/java/io/github/sac/Parser.java index b0f8005..5877c78 100644 --- a/src/main/java/io/github/sac/Parser.java +++ b/src/main/java/io/github/sac/Parser.java @@ -9,6 +9,7 @@ */ public class Parser { + // todo : Probably need to add SUBSCRIBE AND UNSUBSCRIBE EVENTS FROM SERVER IN PARSERESULT public enum ParseResult { ISAUTHENTICATED, PUBLISH,