blob: 45f40ae25ec874e127b803dbf3d620fe9161b6a6 [file] [log] [blame]
Jordan Halterman2bf177c2017-06-29 01:49:08 -07001/*
Brian O'Connora09fe5b2017-08-03 21:12:30 -07002 * Copyright 2017-present Open Networking Foundation
Jordan Halterman2bf177c2017-06-29 01:49:08 -07003 *
4 * Licensed under the Apache License, Version 2.0 (the "License");
5 * you may not use this file except in compliance with the License.
6 * You may obtain a copy of the License at
7 *
8 * http://www.apache.org/licenses/LICENSE-2.0
9 *
10 * Unless required by applicable law or agreed to in writing, software
11 * distributed under the License is distributed on an "AS IS" BASIS,
12 * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
13 * See the License for the specific language governing permissions and
14 * limitations under the License.
15 */
16package org.onosproject.store.primitives.impl;
17
18import java.util.Collection;
19import java.util.Set;
20import java.util.concurrent.CompletableFuture;
21import java.util.concurrent.Executor;
22import java.util.function.Consumer;
Jordan Halterman19486e32017-11-02 15:00:06 -070023import java.util.function.Function;
Jordan Halterman2bf177c2017-06-29 01:49:08 -070024import java.util.stream.Collectors;
25
26import io.atomix.protocols.raft.cluster.MemberId;
27import io.atomix.protocols.raft.protocol.CloseSessionRequest;
28import io.atomix.protocols.raft.protocol.CloseSessionResponse;
29import io.atomix.protocols.raft.protocol.CommandRequest;
30import io.atomix.protocols.raft.protocol.CommandResponse;
Jordan Halterman19486e32017-11-02 15:00:06 -070031import io.atomix.protocols.raft.protocol.HeartbeatRequest;
32import io.atomix.protocols.raft.protocol.HeartbeatResponse;
Jordan Halterman2bf177c2017-06-29 01:49:08 -070033import io.atomix.protocols.raft.protocol.KeepAliveRequest;
34import io.atomix.protocols.raft.protocol.KeepAliveResponse;
35import io.atomix.protocols.raft.protocol.MetadataRequest;
36import io.atomix.protocols.raft.protocol.MetadataResponse;
37import io.atomix.protocols.raft.protocol.OpenSessionRequest;
38import io.atomix.protocols.raft.protocol.OpenSessionResponse;
39import io.atomix.protocols.raft.protocol.PublishRequest;
40import io.atomix.protocols.raft.protocol.QueryRequest;
41import io.atomix.protocols.raft.protocol.QueryResponse;
42import io.atomix.protocols.raft.protocol.RaftClientProtocol;
43import io.atomix.protocols.raft.protocol.ResetRequest;
44import io.atomix.protocols.raft.session.SessionId;
45import org.onosproject.cluster.NodeId;
Jordan Halterman28183ee2017-10-17 17:29:10 -070046import org.onosproject.store.cluster.messaging.ClusterCommunicationService;
Jordan Halterman2bf177c2017-06-29 01:49:08 -070047import org.onosproject.store.service.Serializer;
48
49/**
50 * Raft client protocol that uses a cluster communicator.
51 */
52public class RaftClientCommunicator extends RaftCommunicator implements RaftClientProtocol {
53
54 public RaftClientCommunicator(
Jordan Halterman980a8c12017-09-22 18:01:19 -070055 String prefix,
Jordan Halterman2bf177c2017-06-29 01:49:08 -070056 Serializer serializer,
Jordan Halterman28183ee2017-10-17 17:29:10 -070057 ClusterCommunicationService clusterCommunicator) {
Jordan Halterman980a8c12017-09-22 18:01:19 -070058 super(new RaftMessageContext(prefix), serializer, clusterCommunicator);
Jordan Halterman2bf177c2017-06-29 01:49:08 -070059 }
60
61 @Override
62 public CompletableFuture<OpenSessionResponse> openSession(MemberId memberId, OpenSessionRequest request) {
63 return sendAndReceive(context.openSessionSubject, request, memberId);
64 }
65
66 @Override
67 public CompletableFuture<CloseSessionResponse> closeSession(MemberId memberId, CloseSessionRequest request) {
68 return sendAndReceive(context.closeSessionSubject, request, memberId);
69 }
70
71 @Override
72 public CompletableFuture<KeepAliveResponse> keepAlive(MemberId memberId, KeepAliveRequest request) {
73 return sendAndReceive(context.keepAliveSubject, request, memberId);
74 }
75
76 @Override
77 public CompletableFuture<QueryResponse> query(MemberId memberId, QueryRequest request) {
78 return sendAndReceive(context.querySubject, request, memberId);
79 }
80
81 @Override
82 public CompletableFuture<CommandResponse> command(MemberId memberId, CommandRequest request) {
83 return sendAndReceive(context.commandSubject, request, memberId);
84 }
85
86 @Override
87 public CompletableFuture<MetadataResponse> metadata(MemberId memberId, MetadataRequest request) {
88 return sendAndReceive(context.metadataSubject, request, memberId);
89 }
90
91 @Override
Jordan Halterman19486e32017-11-02 15:00:06 -070092 public void registerHeartbeatHandler(Function<HeartbeatRequest, CompletableFuture<HeartbeatResponse>> function) {
93 clusterCommunicator.addSubscriber(context.heartbeatSubject, serializer::decode, function, serializer::encode);
94 }
95
96 @Override
97 public void unregisterHeartbeatHandler() {
98 clusterCommunicator.removeSubscriber(context.heartbeatSubject);
99 }
100
101 @Override
Jordan Halterman2bf177c2017-06-29 01:49:08 -0700102 public void reset(Collection<MemberId> members, ResetRequest request) {
103 Set<NodeId> nodes = members.stream().map(m -> NodeId.nodeId(m.id())).collect(Collectors.toSet());
104 clusterCommunicator.multicast(
105 request,
106 context.resetSubject(request.session()),
107 serializer::encode,
108 nodes);
109 }
110
111 @Override
112 public void registerPublishListener(SessionId sessionId, Consumer<PublishRequest> listener, Executor executor) {
113 clusterCommunicator.addSubscriber(
114 context.publishSubject(sessionId.id()),
115 serializer::decode,
116 listener,
117 executor);
118 }
119
120 @Override
121 public void unregisterPublishListener(SessionId sessionId) {
122 clusterCommunicator.removeSubscriber(context.publishSubject(sessionId.id()));
123 }
124}