blob: 099e17af840c57f92455f8d994d4d2ccfd0f9d61 [file] [log] [blame]
Jordan Halterman2bf177c2017-06-29 01:49:08 -07001/*
2 * Copyright 2017-present Open Networking Laboratory
3 *
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;
23import java.util.stream.Collectors;
24
25import io.atomix.protocols.raft.cluster.MemberId;
26import io.atomix.protocols.raft.protocol.CloseSessionRequest;
27import io.atomix.protocols.raft.protocol.CloseSessionResponse;
28import io.atomix.protocols.raft.protocol.CommandRequest;
29import io.atomix.protocols.raft.protocol.CommandResponse;
30import io.atomix.protocols.raft.protocol.KeepAliveRequest;
31import io.atomix.protocols.raft.protocol.KeepAliveResponse;
32import io.atomix.protocols.raft.protocol.MetadataRequest;
33import io.atomix.protocols.raft.protocol.MetadataResponse;
34import io.atomix.protocols.raft.protocol.OpenSessionRequest;
35import io.atomix.protocols.raft.protocol.OpenSessionResponse;
36import io.atomix.protocols.raft.protocol.PublishRequest;
37import io.atomix.protocols.raft.protocol.QueryRequest;
38import io.atomix.protocols.raft.protocol.QueryResponse;
39import io.atomix.protocols.raft.protocol.RaftClientProtocol;
40import io.atomix.protocols.raft.protocol.ResetRequest;
41import io.atomix.protocols.raft.session.SessionId;
42import org.onosproject.cluster.NodeId;
43import org.onosproject.cluster.PartitionId;
44import org.onosproject.store.cluster.messaging.ClusterCommunicationService;
45import org.onosproject.store.service.Serializer;
46
47/**
48 * Raft client protocol that uses a cluster communicator.
49 */
50public class RaftClientCommunicator extends RaftCommunicator implements RaftClientProtocol {
51
52 public RaftClientCommunicator(
53 PartitionId partitionId,
54 Serializer serializer,
55 ClusterCommunicationService clusterCommunicator) {
56 super(new RaftMessageContext(String.format("partition-%d", partitionId.id())), serializer, clusterCommunicator);
57 }
58
59 @Override
60 public CompletableFuture<OpenSessionResponse> openSession(MemberId memberId, OpenSessionRequest request) {
61 return sendAndReceive(context.openSessionSubject, request, memberId);
62 }
63
64 @Override
65 public CompletableFuture<CloseSessionResponse> closeSession(MemberId memberId, CloseSessionRequest request) {
66 return sendAndReceive(context.closeSessionSubject, request, memberId);
67 }
68
69 @Override
70 public CompletableFuture<KeepAliveResponse> keepAlive(MemberId memberId, KeepAliveRequest request) {
71 return sendAndReceive(context.keepAliveSubject, request, memberId);
72 }
73
74 @Override
75 public CompletableFuture<QueryResponse> query(MemberId memberId, QueryRequest request) {
76 return sendAndReceive(context.querySubject, request, memberId);
77 }
78
79 @Override
80 public CompletableFuture<CommandResponse> command(MemberId memberId, CommandRequest request) {
81 return sendAndReceive(context.commandSubject, request, memberId);
82 }
83
84 @Override
85 public CompletableFuture<MetadataResponse> metadata(MemberId memberId, MetadataRequest request) {
86 return sendAndReceive(context.metadataSubject, request, memberId);
87 }
88
89 @Override
90 public void reset(Collection<MemberId> members, ResetRequest request) {
91 Set<NodeId> nodes = members.stream().map(m -> NodeId.nodeId(m.id())).collect(Collectors.toSet());
92 clusterCommunicator.multicast(
93 request,
94 context.resetSubject(request.session()),
95 serializer::encode,
96 nodes);
97 }
98
99 @Override
100 public void registerPublishListener(SessionId sessionId, Consumer<PublishRequest> listener, Executor executor) {
101 clusterCommunicator.addSubscriber(
102 context.publishSubject(sessionId.id()),
103 serializer::decode,
104 listener,
105 executor);
106 }
107
108 @Override
109 public void unregisterPublishListener(SessionId sessionId) {
110 clusterCommunicator.removeSubscriber(context.publishSubject(sessionId.id()));
111 }
112}