blob: 2c2f47d28f556b65d3f6686af3707b93e4421cbc [file] [log] [blame]
Madan Jampani5e5b3d62016-02-01 16:03:33 -08001/*
Brian O'Connor5ab426f2016-04-09 01:19:45 -07002 * Copyright 2016-present Open Networking Laboratory
Madan Jampani5e5b3d62016-02-01 16:03:33 -08003 *
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.resources.impl;
17
Madan Jampanid5b200f2016-06-06 17:15:25 -070018import io.atomix.copycat.Operation;
Madan Jampani5e5b3d62016-02-01 16:03:33 -080019import io.atomix.copycat.client.CopycatClient;
Madan Jampani65f24bb2016-03-15 15:16:18 -070020import io.atomix.resource.AbstractResource;
Madan Jampani5e5b3d62016-02-01 16:03:33 -080021import io.atomix.resource.ResourceTypeInfo;
22
Madan Jampanid7ff34d2016-05-26 18:20:55 -070023import java.util.Collection;
Madan Jampani0c0cdc62016-02-22 16:54:06 -080024import java.util.List;
Madan Jampani5e5b3d62016-02-01 16:03:33 -080025import java.util.Map;
Madan Jampani65f24bb2016-03-15 15:16:18 -070026import java.util.Properties;
Madan Jampani5e5b3d62016-02-01 16:03:33 -080027import java.util.Set;
28import java.util.concurrent.CompletableFuture;
29import java.util.function.Consumer;
30
Madan Jampanid5b200f2016-06-06 17:15:25 -070031import org.onlab.util.Tools;
Madan Jampani5e5b3d62016-02-01 16:03:33 -080032import org.onosproject.cluster.Leadership;
33import org.onosproject.cluster.NodeId;
34import org.onosproject.event.Change;
Madan Jampani27a2a3d2016-03-03 08:46:24 -080035import org.onosproject.store.primitives.resources.impl.AtomixLeaderElectorCommands.Anoint;
36import org.onosproject.store.primitives.resources.impl.AtomixLeaderElectorCommands.GetAllLeaderships;
37import org.onosproject.store.primitives.resources.impl.AtomixLeaderElectorCommands.GetElectedTopics;
38import org.onosproject.store.primitives.resources.impl.AtomixLeaderElectorCommands.GetLeadership;
39import org.onosproject.store.primitives.resources.impl.AtomixLeaderElectorCommands.Listen;
40import org.onosproject.store.primitives.resources.impl.AtomixLeaderElectorCommands.Promote;
41import org.onosproject.store.primitives.resources.impl.AtomixLeaderElectorCommands.Run;
42import org.onosproject.store.primitives.resources.impl.AtomixLeaderElectorCommands.Unlisten;
43import org.onosproject.store.primitives.resources.impl.AtomixLeaderElectorCommands.Withdraw;
Madan Jampani5e5b3d62016-02-01 16:03:33 -080044import org.onosproject.store.service.AsyncLeaderElector;
Madan Jampanid5b200f2016-06-06 17:15:25 -070045import org.onosproject.store.service.StorageException;
Madan Jampani5e5b3d62016-02-01 16:03:33 -080046
Madan Jampanid7ff34d2016-05-26 18:20:55 -070047import com.google.common.collect.ImmutableSet;
Madan Jampani630e7ac2016-05-31 11:34:05 -070048import com.google.common.cache.CacheBuilder;
49import com.google.common.cache.CacheLoader;
50import com.google.common.cache.LoadingCache;
Madan Jampani5e5b3d62016-02-01 16:03:33 -080051import com.google.common.collect.Sets;
52
53/**
54 * Distributed resource providing the {@link AsyncLeaderElector} primitive.
55 */
Madan Jampani966a5852016-04-01 10:39:57 -070056@ResourceTypeInfo(id = -152, factory = AtomixLeaderElectorFactory.class)
Madan Jampani65f24bb2016-03-15 15:16:18 -070057public class AtomixLeaderElector extends AbstractResource<AtomixLeaderElector>
Madan Jampanifc981772016-02-16 09:46:42 -080058 implements AsyncLeaderElector {
Madan Jampanid7ff34d2016-05-26 18:20:55 -070059 private final Set<Consumer<Status>> statusChangeListeners =
60 Sets.newCopyOnWriteArraySet();
Madan Jampani5e5b3d62016-02-01 16:03:33 -080061 private final Set<Consumer<Change<Leadership>>> leadershipChangeListeners =
Madan Jampanid7ff34d2016-05-26 18:20:55 -070062 Sets.newCopyOnWriteArraySet();
Madan Jampani630e7ac2016-05-31 11:34:05 -070063 private final Consumer<Change<Leadership>> cacheUpdater;
64 private final Consumer<Status> statusListener;
Madan Jampani5e5b3d62016-02-01 16:03:33 -080065
Madan Jampanic94b4852016-02-23 18:18:37 -080066 public static final String CHANGE_SUBJECT = "leadershipChangeEvents";
Madan Jampani630e7ac2016-05-31 11:34:05 -070067 private final LoadingCache<String, CompletableFuture<Leadership>> cache;
Madan Jampani5e5b3d62016-02-01 16:03:33 -080068
Madan Jampani65f24bb2016-03-15 15:16:18 -070069 public AtomixLeaderElector(CopycatClient client, Properties properties) {
70 super(client, properties);
Madan Jampani630e7ac2016-05-31 11:34:05 -070071 cache = CacheBuilder.newBuilder()
72 .maximumSize(1000)
Madan Jampanid5b200f2016-06-06 17:15:25 -070073 .build(CacheLoader.from(topic -> submit(new GetLeadership(topic))));
Madan Jampani630e7ac2016-05-31 11:34:05 -070074
75 cacheUpdater = change -> {
76 Leadership leadership = change.newValue();
77 cache.put(leadership.topic(), CompletableFuture.completedFuture(leadership));
78 };
79 statusListener = status -> {
80 if (status == Status.SUSPENDED || status == Status.INACTIVE) {
81 cache.invalidateAll();
82 }
83 };
84 addStatusChangeListener(statusListener);
85 }
86
87 @Override
88 public CompletableFuture<Void> destroy() {
89 removeStatusChangeListener(statusListener);
90 return removeChangeListener(cacheUpdater);
Madan Jampani5e5b3d62016-02-01 16:03:33 -080091 }
92
93 @Override
94 public String name() {
95 return null;
96 }
97
98 @Override
99 public CompletableFuture<AtomixLeaderElector> open() {
100 return super.open().thenApply(result -> {
Madan Jampani0c0cdc62016-02-22 16:54:06 -0800101 client.onEvent(CHANGE_SUBJECT, this::handleEvent);
Madan Jampani5e5b3d62016-02-01 16:03:33 -0800102 return result;
103 });
104 }
105
Madan Jampani630e7ac2016-05-31 11:34:05 -0700106 public CompletableFuture<AtomixLeaderElector> setupCache() {
107 return addChangeListener(cacheUpdater).thenApply(v -> this);
108 }
109
Madan Jampani0c0cdc62016-02-22 16:54:06 -0800110 private void handleEvent(List<Change<Leadership>> changes) {
111 changes.forEach(change -> leadershipChangeListeners.forEach(l -> l.accept(change)));
Madan Jampani5e5b3d62016-02-01 16:03:33 -0800112 }
113
114 @Override
115 public CompletableFuture<Leadership> run(String topic, NodeId nodeId) {
Madan Jampanid5b200f2016-06-06 17:15:25 -0700116 return submit(new Run(topic, nodeId)).whenComplete((r, e) -> cache.invalidate(topic));
Madan Jampani5e5b3d62016-02-01 16:03:33 -0800117 }
118
119 @Override
120 public CompletableFuture<Void> withdraw(String topic) {
Madan Jampanid5b200f2016-06-06 17:15:25 -0700121 return submit(new Withdraw(topic)).whenComplete((r, e) -> cache.invalidate(topic));
Madan Jampani5e5b3d62016-02-01 16:03:33 -0800122 }
123
124 @Override
125 public CompletableFuture<Boolean> anoint(String topic, NodeId nodeId) {
Madan Jampanid5b200f2016-06-06 17:15:25 -0700126 return submit(new Anoint(topic, nodeId)).whenComplete((r, e) -> cache.invalidate(topic));
Madan Jampani5e5b3d62016-02-01 16:03:33 -0800127 }
128
129 @Override
Madan Jampani0c0cdc62016-02-22 16:54:06 -0800130 public CompletableFuture<Boolean> promote(String topic, NodeId nodeId) {
Madan Jampanid5b200f2016-06-06 17:15:25 -0700131 return submit(new Promote(topic, nodeId)).whenComplete((r, e) -> cache.invalidate(topic));
Madan Jampani0c0cdc62016-02-22 16:54:06 -0800132 }
133
134 @Override
135 public CompletableFuture<Void> evict(NodeId nodeId) {
Madan Jampanid5b200f2016-06-06 17:15:25 -0700136 return submit(new AtomixLeaderElectorCommands.Evict(nodeId));
Madan Jampani0c0cdc62016-02-22 16:54:06 -0800137 }
138
139 @Override
Madan Jampani5e5b3d62016-02-01 16:03:33 -0800140 public CompletableFuture<Leadership> getLeadership(String topic) {
Madan Jampani77012442016-06-02 07:47:42 -0700141 return cache.getUnchecked(topic)
142 .whenComplete((r, e) -> {
143 if (e != null) {
144 cache.invalidate(topic);
145 }
146 });
Madan Jampani5e5b3d62016-02-01 16:03:33 -0800147 }
148
149 @Override
150 public CompletableFuture<Map<String, Leadership>> getLeaderships() {
Madan Jampanid5b200f2016-06-06 17:15:25 -0700151 return submit(new GetAllLeaderships());
Madan Jampani5e5b3d62016-02-01 16:03:33 -0800152 }
153
154 public CompletableFuture<Set<String>> getElectedTopics(NodeId nodeId) {
Madan Jampanid5b200f2016-06-06 17:15:25 -0700155 return submit(new GetElectedTopics(nodeId));
Madan Jampani5e5b3d62016-02-01 16:03:33 -0800156 }
157
Madan Jampani27a2a3d2016-03-03 08:46:24 -0800158 @Override
159 public synchronized CompletableFuture<Void> addChangeListener(Consumer<Change<Leadership>> consumer) {
160 if (leadershipChangeListeners.isEmpty()) {
Madan Jampanid5b200f2016-06-06 17:15:25 -0700161 return submit(new Listen()).thenRun(() -> leadershipChangeListeners.add(consumer));
Madan Jampani27a2a3d2016-03-03 08:46:24 -0800162 } else {
163 leadershipChangeListeners.add(consumer);
164 return CompletableFuture.completedFuture(null);
Madan Jampani5e5b3d62016-02-01 16:03:33 -0800165 }
166 }
167
168 @Override
Madan Jampani27a2a3d2016-03-03 08:46:24 -0800169 public synchronized CompletableFuture<Void> removeChangeListener(Consumer<Change<Leadership>> consumer) {
Madan Jampani68b1f5a2016-04-05 19:17:58 -0700170 if (leadershipChangeListeners.remove(consumer) && leadershipChangeListeners.isEmpty()) {
Madan Jampanid5b200f2016-06-06 17:15:25 -0700171 return submit(new Unlisten()).thenApply(v -> null);
Madan Jampani5e5b3d62016-02-01 16:03:33 -0800172 }
173 return CompletableFuture.completedFuture(null);
174 }
Madan Jampanid7ff34d2016-05-26 18:20:55 -0700175
176 @Override
177 public void addStatusChangeListener(Consumer<Status> listener) {
178 statusChangeListeners.add(listener);
179 }
180
181 @Override
182 public void removeStatusChangeListener(Consumer<Status> listener) {
183 statusChangeListeners.remove(listener);
184 }
185
186 @Override
187 public Collection<Consumer<Status>> statusChangeListeners() {
188 return ImmutableSet.copyOf(statusChangeListeners);
189 }
Madan Jampanid5b200f2016-06-06 17:15:25 -0700190
191 <T> CompletableFuture<T> submit(Operation<T> command) {
192 if (client.state() == CopycatClient.State.SUSPENDED || client.state() == CopycatClient.State.CLOSED) {
193 return Tools.exceptionalFuture(new StorageException.Unavailable());
194 }
195 return client.submit(command);
196 }
Madan Jampani5e5b3d62016-02-01 16:03:33 -0800197}