Madan Jampani | 9b19a82 | 2014-11-04 21:37:13 -0800 | [diff] [blame] | 1 | package org.onlab.onos.store.service.impl; |
| 2 | |
| 3 | import static org.slf4j.LoggerFactory.getLogger; |
| 4 | |
| 5 | import java.util.concurrent.CompletableFuture; |
Madan Jampani | 778f7ad | 2014-11-05 22:46:15 -0800 | [diff] [blame] | 6 | import java.util.function.BiConsumer; |
Madan Jampani | 9b19a82 | 2014-11-04 21:37:13 -0800 | [diff] [blame] | 7 | |
| 8 | import net.kuujo.copycat.protocol.PingRequest; |
Madan Jampani | 778f7ad | 2014-11-05 22:46:15 -0800 | [diff] [blame] | 9 | import net.kuujo.copycat.protocol.PingResponse; |
Madan Jampani | 9b19a82 | 2014-11-04 21:37:13 -0800 | [diff] [blame] | 10 | import net.kuujo.copycat.protocol.PollRequest; |
Madan Jampani | 778f7ad | 2014-11-05 22:46:15 -0800 | [diff] [blame] | 11 | import net.kuujo.copycat.protocol.PollResponse; |
Madan Jampani | 9b19a82 | 2014-11-04 21:37:13 -0800 | [diff] [blame] | 12 | import net.kuujo.copycat.protocol.RequestHandler; |
| 13 | import net.kuujo.copycat.protocol.SubmitRequest; |
Madan Jampani | 778f7ad | 2014-11-05 22:46:15 -0800 | [diff] [blame] | 14 | import net.kuujo.copycat.protocol.SubmitResponse; |
Madan Jampani | 9b19a82 | 2014-11-04 21:37:13 -0800 | [diff] [blame] | 15 | import net.kuujo.copycat.protocol.SyncRequest; |
Madan Jampani | 778f7ad | 2014-11-05 22:46:15 -0800 | [diff] [blame] | 16 | import net.kuujo.copycat.protocol.SyncResponse; |
Madan Jampani | 9b19a82 | 2014-11-04 21:37:13 -0800 | [diff] [blame] | 17 | import net.kuujo.copycat.spi.protocol.ProtocolServer; |
| 18 | |
| 19 | import org.onlab.onos.store.cluster.messaging.ClusterCommunicationService; |
| 20 | import org.onlab.onos.store.cluster.messaging.ClusterMessage; |
| 21 | import org.onlab.onos.store.cluster.messaging.ClusterMessageHandler; |
| 22 | import org.slf4j.Logger; |
| 23 | |
| 24 | /** |
Madan Jampani | dfbfa18 | 2014-11-04 22:06:41 -0800 | [diff] [blame] | 25 | * ONOS Cluster messaging based Copycat protocol server. |
Madan Jampani | 9b19a82 | 2014-11-04 21:37:13 -0800 | [diff] [blame] | 26 | */ |
Madan Jampani | 9b19a82 | 2014-11-04 21:37:13 -0800 | [diff] [blame] | 27 | public class ClusterMessagingProtocolServer implements ProtocolServer { |
| 28 | |
| 29 | private final Logger log = getLogger(getClass()); |
Yuta HIGUCHI | 76b54bf | 2014-11-07 01:56:55 -0800 | [diff] [blame] | 30 | private volatile RequestHandler handler; |
| 31 | private ClusterCommunicationService clusterCommunicator; |
Madan Jampani | 9b19a82 | 2014-11-04 21:37:13 -0800 | [diff] [blame] | 32 | |
| 33 | public ClusterMessagingProtocolServer(ClusterCommunicationService clusterCommunicator) { |
Yuta HIGUCHI | 76b54bf | 2014-11-07 01:56:55 -0800 | [diff] [blame] | 34 | this.clusterCommunicator = clusterCommunicator; |
Madan Jampani | 9b19a82 | 2014-11-04 21:37:13 -0800 | [diff] [blame] | 35 | |
Madan Jampani | 9b19a82 | 2014-11-04 21:37:13 -0800 | [diff] [blame] | 36 | } |
| 37 | |
| 38 | @Override |
| 39 | public void requestHandler(RequestHandler handler) { |
| 40 | this.handler = handler; |
| 41 | } |
| 42 | |
| 43 | @Override |
| 44 | public CompletableFuture<Void> listen() { |
Yuta HIGUCHI | 76b54bf | 2014-11-07 01:56:55 -0800 | [diff] [blame] | 45 | clusterCommunicator.addSubscriber(ClusterMessagingProtocol.COPYCAT_PING, |
| 46 | new CopycatMessageHandler<PingRequest>()); |
| 47 | clusterCommunicator.addSubscriber(ClusterMessagingProtocol.COPYCAT_SYNC, |
| 48 | new CopycatMessageHandler<SyncRequest>()); |
| 49 | clusterCommunicator.addSubscriber(ClusterMessagingProtocol.COPYCAT_POLL, |
| 50 | new CopycatMessageHandler<PollRequest>()); |
| 51 | clusterCommunicator.addSubscriber(ClusterMessagingProtocol.COPYCAT_SUBMIT, |
| 52 | new CopycatMessageHandler<SubmitRequest>()); |
Madan Jampani | 9b19a82 | 2014-11-04 21:37:13 -0800 | [diff] [blame] | 53 | return CompletableFuture.completedFuture(null); |
| 54 | } |
| 55 | |
| 56 | @Override |
| 57 | public CompletableFuture<Void> close() { |
Yuta HIGUCHI | 76b54bf | 2014-11-07 01:56:55 -0800 | [diff] [blame] | 58 | clusterCommunicator.removeSubscriber(ClusterMessagingProtocol.COPYCAT_PING); |
| 59 | clusterCommunicator.removeSubscriber(ClusterMessagingProtocol.COPYCAT_SYNC); |
| 60 | clusterCommunicator.removeSubscriber(ClusterMessagingProtocol.COPYCAT_POLL); |
| 61 | clusterCommunicator.removeSubscriber(ClusterMessagingProtocol.COPYCAT_SUBMIT); |
Madan Jampani | 9b19a82 | 2014-11-04 21:37:13 -0800 | [diff] [blame] | 62 | return CompletableFuture.completedFuture(null); |
| 63 | } |
| 64 | |
| 65 | private class CopycatMessageHandler<T> implements ClusterMessageHandler { |
| 66 | |
| 67 | @Override |
| 68 | public void handle(ClusterMessage message) { |
| 69 | T request = ClusterMessagingProtocol.SERIALIZER.decode(message.payload()); |
Yuta HIGUCHI | 5fb6c96 | 2014-11-07 13:07:40 -0800 | [diff] [blame] | 70 | if (handler == null) { |
| 71 | // there is a slight window of time during state transition, |
| 72 | // where handler becomes null |
| 73 | for (int i = 0; i < 10; ++i) { |
| 74 | if (handler != null) { |
| 75 | break; |
| 76 | } |
| 77 | try { |
| 78 | Thread.sleep(1); |
| 79 | } catch (InterruptedException e) { |
| 80 | log.trace("Exception", e); |
| 81 | } |
| 82 | } |
| 83 | if (handler == null) { |
| 84 | log.error("There was no handler for registered!"); |
| 85 | return; |
| 86 | } |
| 87 | } |
Madan Jampani | 9b19a82 | 2014-11-04 21:37:13 -0800 | [diff] [blame] | 88 | if (request.getClass().equals(PingRequest.class)) { |
Yuta HIGUCHI | 535c915 | 2014-11-07 15:18:57 -0800 | [diff] [blame] | 89 | handler.ping((PingRequest) request) |
| 90 | .whenComplete(new PostExecutionTask<PingResponse>(message)); |
Madan Jampani | 9b19a82 | 2014-11-04 21:37:13 -0800 | [diff] [blame] | 91 | } else if (request.getClass().equals(PollRequest.class)) { |
Yuta HIGUCHI | 535c915 | 2014-11-07 15:18:57 -0800 | [diff] [blame] | 92 | handler.poll((PollRequest) request) |
| 93 | .whenComplete(new PostExecutionTask<PollResponse>(message)); |
Madan Jampani | 9b19a82 | 2014-11-04 21:37:13 -0800 | [diff] [blame] | 94 | } else if (request.getClass().equals(SyncRequest.class)) { |
Yuta HIGUCHI | 535c915 | 2014-11-07 15:18:57 -0800 | [diff] [blame] | 95 | handler.sync((SyncRequest) request) |
| 96 | .whenComplete(new PostExecutionTask<SyncResponse>(message)); |
Madan Jampani | 9b19a82 | 2014-11-04 21:37:13 -0800 | [diff] [blame] | 97 | } else if (request.getClass().equals(SubmitRequest.class)) { |
Yuta HIGUCHI | 535c915 | 2014-11-07 15:18:57 -0800 | [diff] [blame] | 98 | handler.submit((SubmitRequest) request) |
| 99 | .whenComplete(new PostExecutionTask<SubmitResponse>(message)); |
Madan Jampani | 778f7ad | 2014-11-05 22:46:15 -0800 | [diff] [blame] | 100 | } else { |
| 101 | throw new IllegalStateException("Unknown request type: " + request.getClass().getName()); |
| 102 | } |
| 103 | } |
| 104 | |
| 105 | private class PostExecutionTask<R> implements BiConsumer<R, Throwable> { |
| 106 | |
| 107 | private final ClusterMessage message; |
| 108 | |
| 109 | public PostExecutionTask(ClusterMessage message) { |
| 110 | this.message = message; |
| 111 | } |
| 112 | |
| 113 | @Override |
| 114 | public void accept(R response, Throwable t) { |
| 115 | if (t != null) { |
| 116 | log.error("Processing for " + message.subject() + " failed.", t); |
| 117 | } else { |
Madan Jampani | 9b19a82 | 2014-11-04 21:37:13 -0800 | [diff] [blame] | 118 | try { |
Yuta HIGUCHI | 5fb6c96 | 2014-11-07 13:07:40 -0800 | [diff] [blame] | 119 | log.trace("responding to {}", message.subject()); |
Madan Jampani | 9b19a82 | 2014-11-04 21:37:13 -0800 | [diff] [blame] | 120 | message.respond(ClusterMessagingProtocol.SERIALIZER.encode(response)); |
| 121 | } catch (Exception e) { |
Madan Jampani | 778f7ad | 2014-11-05 22:46:15 -0800 | [diff] [blame] | 122 | log.error("Failed to respond to " + response.getClass().getName(), e); |
Madan Jampani | 9b19a82 | 2014-11-04 21:37:13 -0800 | [diff] [blame] | 123 | } |
Madan Jampani | 778f7ad | 2014-11-05 22:46:15 -0800 | [diff] [blame] | 124 | } |
Madan Jampani | 9b19a82 | 2014-11-04 21:37:13 -0800 | [diff] [blame] | 125 | } |
| 126 | } |
| 127 | } |
| 128 | } |