Added a messaging service implementation on top of IOLoop. Added ability to easily switch between netty and io loop (default is netty)
Change-Id: Id9af0756bf0a542f832f3611b486b2ac680b91e4
diff --git a/core/store/dist/src/main/java/org/onosproject/store/cluster/impl/DistributedClusterStore.java b/core/store/dist/src/main/java/org/onosproject/store/cluster/impl/DistributedClusterStore.java
index f458bfe..c81079d 100644
--- a/core/store/dist/src/main/java/org/onosproject/store/cluster/impl/DistributedClusterStore.java
+++ b/core/store/dist/src/main/java/org/onosproject/store/cluster/impl/DistributedClusterStore.java
@@ -19,14 +19,12 @@
import com.google.common.collect.Maps;
import com.google.common.collect.Sets;
import com.hazelcast.util.AddressUtil;
+
import org.apache.felix.scr.annotations.Activate;
import org.apache.felix.scr.annotations.Component;
import org.apache.felix.scr.annotations.Deactivate;
import org.apache.felix.scr.annotations.Service;
import org.joda.time.DateTime;
-import org.onlab.netty.Endpoint;
-import org.onlab.netty.Message;
-import org.onlab.netty.MessageHandler;
import org.onlab.netty.NettyMessagingService;
import org.onlab.packet.IpAddress;
import org.onlab.util.KryoNamespace;
@@ -38,6 +36,7 @@
import org.onosproject.cluster.DefaultControllerNode;
import org.onosproject.cluster.NodeId;
import org.onosproject.store.AbstractStore;
+import org.onosproject.store.cluster.messaging.Endpoint;
import org.onosproject.store.consistent.impl.DatabaseDefinition;
import org.onosproject.store.consistent.impl.DatabaseDefinitionStore;
import org.onosproject.store.serializers.KryoNamespaces;
@@ -56,6 +55,7 @@
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
+import java.util.function.Consumer;
import java.util.stream.Collectors;
import static com.google.common.base.Preconditions.checkNotNull;
@@ -108,7 +108,7 @@
private final Map<NodeId, ControllerNode> allNodes = Maps.newConcurrentMap();
private final Map<NodeId, State> nodeStates = Maps.newConcurrentMap();
private final Map<NodeId, DateTime> nodeStateLastUpdatedTimes = Maps.newConcurrentMap();
- private NettyMessagingService messagingService = new NettyMessagingService();
+ private NettyMessagingService messagingService;
private ScheduledExecutorService heartBeatSender = Executors.newSingleThreadScheduledExecutor(
groupedThreads("onos/cluster/membership", "heartbeat-sender"));
private ExecutorService heartBeatMessageHandler = Executors.newSingleThreadExecutor(
@@ -149,7 +149,6 @@
establishSelfIdentity();
messagingService = new NettyMessagingService(HEARTBEAT_FD_PORT);
-
try {
messagingService.activate();
} catch (InterruptedException e) {
@@ -376,10 +375,10 @@
throw new IllegalStateException("Unable to determine local ip");
}
- private class HeartbeatMessageHandler implements MessageHandler {
+ private class HeartbeatMessageHandler implements Consumer<byte[]> {
@Override
- public void handle(Message message) throws IOException {
- HeartbeatMessage hb = SERIALIZER.decode(message.payload());
+ public void accept(byte[] message) {
+ HeartbeatMessage hb = SERIALIZER.decode(message);
failureDetector.report(hb.source().id());
hb.knownPeers().forEach(node -> {
allNodes.put(node.id(), node);
diff --git a/core/store/dist/src/main/java/org/onosproject/store/cluster/messaging/impl/ClusterCommunicationManager.java b/core/store/dist/src/main/java/org/onosproject/store/cluster/messaging/impl/ClusterCommunicationManager.java
index 042bda2..6f47b48 100644
--- a/core/store/dist/src/main/java/org/onosproject/store/cluster/messaging/impl/ClusterCommunicationManager.java
+++ b/core/store/dist/src/main/java/org/onosproject/store/cluster/messaging/impl/ClusterCommunicationManager.java
@@ -21,18 +21,17 @@
import org.apache.felix.scr.annotations.Reference;
import org.apache.felix.scr.annotations.ReferenceCardinality;
import org.apache.felix.scr.annotations.Service;
-import org.onlab.netty.Endpoint;
-import org.onlab.netty.Message;
-import org.onlab.netty.MessageHandler;
-import org.onlab.netty.MessagingService;
import org.onlab.netty.NettyMessagingService;
+import org.onlab.nio.service.IOLoopMessagingService;
import org.onosproject.cluster.ClusterService;
import org.onosproject.cluster.ControllerNode;
import org.onosproject.cluster.NodeId;
import org.onosproject.store.cluster.messaging.ClusterCommunicationService;
import org.onosproject.store.cluster.messaging.ClusterMessage;
import org.onosproject.store.cluster.messaging.ClusterMessageHandler;
+import org.onosproject.store.cluster.messaging.Endpoint;
import org.onosproject.store.cluster.messaging.MessageSubject;
+import org.onosproject.store.cluster.messaging.MessagingService;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -64,17 +63,28 @@
// TODO: This probably should not be a OSGi service.
private MessagingService messagingService;
+ private final boolean useNetty = true;
+
@Activate
public void activate() {
ControllerNode localNode = clusterService.getLocalNode();
- NettyMessagingService netty = new NettyMessagingService(localNode.ip(), localNode.tcpPort());
- // FIXME: workaround until it becomes a service.
- try {
- netty.activate();
- } catch (Exception e) {
- log.error("NettyMessagingService#activate", e);
+ if (useNetty) {
+ NettyMessagingService netty = new NettyMessagingService(localNode.ip(), localNode.tcpPort());
+ try {
+ netty.activate();
+ messagingService = netty;
+ } catch (Exception e) {
+ log.error("NettyMessagingService#activate", e);
+ }
+ } else {
+ IOLoopMessagingService ioLoop = new IOLoopMessagingService(localNode.ip(), localNode.tcpPort());
+ try {
+ ioLoop.activate();
+ messagingService = ioLoop;
+ } catch (Exception e) {
+ log.error("IOLoopMessagingService#activate", e);
+ }
}
- messagingService = netty;
log.info("Started on {}:{}", localNode.ip(), localNode.tcpPort());
}
@@ -83,9 +93,13 @@
// TODO: cleanup messageingService if needed.
// FIXME: workaround until it becomes a service.
try {
- ((NettyMessagingService) messagingService).deactivate();
+ if (useNetty) {
+ ((NettyMessagingService) messagingService).deactivate();
+ } else {
+ ((IOLoopMessagingService) messagingService).deactivate();
+ }
} catch (Exception e) {
- log.error("NettyMessagingService#deactivate", e);
+ log.error("MessagingService#deactivate", e);
}
log.info("Stopped");
}
@@ -232,7 +246,9 @@
public void addSubscriber(MessageSubject subject,
ClusterMessageHandler subscriber,
ExecutorService executor) {
- messagingService.registerHandler(subject.value(), new InternalClusterMessageHandler(subscriber), executor);
+ messagingService.registerHandler(subject.value(),
+ new InternalClusterMessageHandler(subscriber),
+ executor);
}
@Override
@@ -240,31 +256,6 @@
messagingService.unregisterHandler(subject.value());
}
- private final class InternalClusterMessageHandler implements MessageHandler {
-
- private final ClusterMessageHandler handler;
-
- public InternalClusterMessageHandler(ClusterMessageHandler handler) {
- this.handler = handler;
- }
-
- @Override
- public void handle(Message message) {
- final ClusterMessage clusterMessage;
- try {
- clusterMessage = ClusterMessage.fromBytes(message.payload());
- } catch (Exception e) {
- log.error("Failed decoding {}", message, e);
- throw e;
- }
- try {
- handler.handle(new InternalClusterMessage(clusterMessage, message));
- } catch (Exception e) {
- log.trace("Failed handling {}", clusterMessage, e);
- throw e;
- }
- }
- }
@Override
public <M, R> void addSubscriber(MessageSubject subject,
@@ -287,7 +278,22 @@
executor);
}
- private class InternalMessageResponder<M, R> implements MessageHandler {
+ private class InternalClusterMessageHandler implements Function<byte[], byte[]> {
+ private ClusterMessageHandler handler;
+
+ public InternalClusterMessageHandler(ClusterMessageHandler handler) {
+ this.handler = handler;
+ }
+
+ @Override
+ public byte[] apply(byte[] bytes) {
+ ClusterMessage message = ClusterMessage.fromBytes(bytes);
+ handler.handle(message);
+ return message.response();
+ }
+ }
+
+ private class InternalMessageResponder<M, R> implements Function<byte[], byte[]> {
private final Function<byte[], M> decoder;
private final Function<R, byte[]> encoder;
private final Function<M, R> handler;
@@ -299,14 +305,15 @@
this.encoder = encoder;
this.handler = handler;
}
+
@Override
- public void handle(Message message) throws IOException {
- R response = handler.apply(decoder.apply(ClusterMessage.fromBytes(message.payload()).payload()));
- message.respond(encoder.apply(response));
+ public byte[] apply(byte[] bytes) {
+ R reply = handler.apply(decoder.apply(ClusterMessage.fromBytes(bytes).payload()));
+ return encoder.apply(reply);
}
}
- private class InternalMessageConsumer<M> implements MessageHandler {
+ private class InternalMessageConsumer<M> implements Consumer<byte[]> {
private final Function<byte[], M> decoder;
private final Consumer<M> consumer;
@@ -314,24 +321,10 @@
this.decoder = decoder;
this.consumer = consumer;
}
- @Override
- public void handle(Message message) throws IOException {
- consumer.accept(decoder.apply(ClusterMessage.fromBytes(message.payload()).payload()));
- }
- }
-
- public static final class InternalClusterMessage extends ClusterMessage {
-
- private final Message rawMessage;
-
- public InternalClusterMessage(ClusterMessage clusterMessage, Message rawMessage) {
- super(clusterMessage.sender(), clusterMessage.subject(), clusterMessage.payload());
- this.rawMessage = rawMessage;
- }
@Override
- public void respond(byte[] response) throws IOException {
- rawMessage.respond(response);
+ public void accept(byte[] bytes) {
+ consumer.accept(decoder.apply(ClusterMessage.fromBytes(bytes).payload()));
}
}
}
diff --git a/core/store/dist/src/main/java/org/onosproject/store/flow/impl/DistributedFlowRuleStore.java b/core/store/dist/src/main/java/org/onosproject/store/flow/impl/DistributedFlowRuleStore.java
index 8ae7d2a..7f8769b 100644
--- a/core/store/dist/src/main/java/org/onosproject/store/flow/impl/DistributedFlowRuleStore.java
+++ b/core/store/dist/src/main/java/org/onosproject/store/flow/impl/DistributedFlowRuleStore.java
@@ -77,7 +77,6 @@
import org.osgi.service.component.ComponentContext;
import org.slf4j.Logger;
-import java.io.IOException;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collections;
@@ -278,11 +277,7 @@
FlowRule rule = SERIALIZER.decode(message.payload());
log.trace("received get flow entry request for {}", rule);
FlowEntry flowEntry = flowTable.getFlowEntry(rule); //getFlowEntryInternal(rule);
- try {
- message.respond(SERIALIZER.encode(flowEntry));
- } catch (IOException e) {
- log.error("Failed to respond back", e);
- }
+ message.respond(SERIALIZER.encode(flowEntry));
}
}, executor);
@@ -293,11 +288,7 @@
DeviceId deviceId = SERIALIZER.decode(message.payload());
log.trace("Received get flow entries request for {} from {}", deviceId, message.sender());
Set<FlowEntry> flowEntries = flowTable.getFlowEntries(deviceId);
- try {
- message.respond(SERIALIZER.encode(flowEntries));
- } catch (IOException e) {
- log.error("Failed to respond to peer's getFlowEntries request", e);
- }
+ message.respond(SERIALIZER.encode(flowEntries));
}
}, executor);
@@ -308,11 +299,7 @@
FlowEntry rule = SERIALIZER.decode(message.payload());
log.trace("received get flow entry request for {}", rule);
FlowRuleEvent event = removeFlowRuleInternal(rule);
- try {
- message.respond(SERIALIZER.encode(event));
- } catch (IOException e) {
- log.error("Failed to respond back", e);
- }
+ message.respond(SERIALIZER.encode(event));
}
}, executor);
}
@@ -691,11 +678,7 @@
// TODO: we might want to wrap response in envelope
// to distinguish sw programming failure and hand over
// it make sense in the latter case to retry immediately.
- try {
- message.respond(SERIALIZER.encode(allFailed));
- } catch (IOException e) {
- log.error("Failed to respond back", e);
- }
+ message.respond(SERIALIZER.encode(allFailed));
return;
}
diff --git a/core/store/dist/src/main/java/org/onosproject/store/flowext/impl/DefaultFlowRuleExtRouter.java b/core/store/dist/src/main/java/org/onosproject/store/flowext/impl/DefaultFlowRuleExtRouter.java
index b609a3c..b7d69b2 100644
--- a/core/store/dist/src/main/java/org/onosproject/store/flowext/impl/DefaultFlowRuleExtRouter.java
+++ b/core/store/dist/src/main/java/org/onosproject/store/flowext/impl/DefaultFlowRuleExtRouter.java
@@ -51,7 +51,6 @@
import org.onosproject.store.serializers.impl.DistributedStoreSerializers;
import org.slf4j.Logger;
-import java.io.IOException;
import java.util.Collection;
import java.util.Collections;
import java.util.HashSet;
@@ -143,11 +142,7 @@
@Override
public void run() {
FlowExtCompletedOperation result = Futures.getUnchecked(f);
- try {
- message.respond(SERIALIZER.encode(result));
- } catch (IOException e) {
- log.error("Failed to respond back", e);
- }
+ message.respond(SERIALIZER.encode(result));
}
}, futureListeners);
}
diff --git a/core/store/dist/src/main/java/org/onosproject/store/mastership/impl/ConsistentDeviceMastershipStore.java b/core/store/dist/src/main/java/org/onosproject/store/mastership/impl/ConsistentDeviceMastershipStore.java
index 27a33b3..bc31e08 100644
--- a/core/store/dist/src/main/java/org/onosproject/store/mastership/impl/ConsistentDeviceMastershipStore.java
+++ b/core/store/dist/src/main/java/org/onosproject/store/mastership/impl/ConsistentDeviceMastershipStore.java
@@ -22,7 +22,6 @@
import static org.slf4j.LoggerFactory.getLogger;
import static com.google.common.base.Preconditions.checkArgument;
-import java.io.IOException;
import java.util.List;
import java.util.Map;
import java.util.Set;
@@ -300,11 +299,7 @@
@Override
public void handle(ClusterMessage message) {
DeviceId deviceId = SERIALIZER.decode(message.payload());
- try {
- message.respond(SERIALIZER.encode(getRole(localNodeId, deviceId)));
- } catch (IOException e) {
- log.error("Failed to responsd to role query", e);
- }
+ message.respond(SERIALIZER.encode(getRole(localNodeId, deviceId)));
}
}
@@ -318,11 +313,7 @@
@Override
public void handle(ClusterMessage message) {
DeviceId deviceId = SERIALIZER.decode(message.payload());
- try {
- message.respond(SERIALIZER.encode(relinquishRole(localNodeId, deviceId)));
- } catch (IOException e) {
- log.error("Failed to relinquish role.", e);
- }
+ message.respond(SERIALIZER.encode(relinquishRole(localNodeId, deviceId)));
}
}
@@ -371,4 +362,4 @@
return m.matches();
}
-}
\ No newline at end of file
+}
diff --git a/core/store/dist/src/main/java/org/onosproject/store/statistic/impl/DistributedStatisticStore.java b/core/store/dist/src/main/java/org/onosproject/store/statistic/impl/DistributedStatisticStore.java
index 907631d..cfbcec8 100644
--- a/core/store/dist/src/main/java/org/onosproject/store/statistic/impl/DistributedStatisticStore.java
+++ b/core/store/dist/src/main/java/org/onosproject/store/statistic/impl/DistributedStatisticStore.java
@@ -43,7 +43,6 @@
import org.onosproject.store.serializers.KryoSerializer;
import org.slf4j.Logger;
-import java.io.IOException;
import java.util.Collections;
import java.util.HashSet;
import java.util.Map;
@@ -118,11 +117,7 @@
@Override
public void handle(ClusterMessage message) {
ConnectPoint cp = SERIALIZER.decode(message.payload());
- try {
- message.respond(SERIALIZER.encode(getCurrentStatisticInternal(cp)));
- } catch (IOException e) {
- log.error("Failed to respond back", e);
- }
+ message.respond(SERIALIZER.encode(getCurrentStatisticInternal(cp)));
}
}, messageHandlingExecutor);
@@ -131,11 +126,7 @@
@Override
public void handle(ClusterMessage message) {
ConnectPoint cp = SERIALIZER.decode(message.payload());
- try {
- message.respond(SERIALIZER.encode(getPreviousStatisticInternal(cp)));
- } catch (IOException e) {
- log.error("Failed to respond back", e);
- }
+ message.respond(SERIALIZER.encode(getPreviousStatisticInternal(cp)));
}
}, messageHandlingExecutor);
log.info("Started");