LeaderInitiator

This commit is contained in:
Gytis Trikleris
2018-06-05 12:46:18 +02:00
committed by Ioannis Canellos
parent 91f21172fd
commit 47baca5cac
2 changed files with 134 additions and 141 deletions

View File

@@ -16,18 +16,9 @@
package org.springframework.cloud.kubernetes.leader;
import java.util.Map;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.locks.Lock;
import java.util.concurrent.locks.ReentrantLock;
import io.fabric8.kubernetes.api.model.ConfigMap;
import io.fabric8.kubernetes.api.model.ConfigMapBuilder;
import io.fabric8.kubernetes.api.model.ObjectMeta;
import io.fabric8.kubernetes.api.model.Pod;
import io.fabric8.kubernetes.client.KubernetesClient;
import org.springframework.context.SmartLifecycle;
import org.springframework.integration.leader.Candidate;
@@ -36,22 +27,22 @@ import org.springframework.integration.leader.Candidate;
*/
public class LeaderInitiator implements SmartLifecycle {
private final Lock lock = new ReentrantLock();
private final KubernetesClient kubernetesClient;
private final LeadershipController leadershipController;
private final Candidate candidate;
private final LeaderProperties leaderProperties;
private ScheduledExecutorService scheduledExecutorService;
private final ScheduledExecutorService scheduledExecutorService;
private String currentLeaderId; // TODO remove once events are implemented
private boolean isRunning;
public LeaderInitiator(KubernetesClient kubernetesClient, Candidate candidate, LeaderProperties leaderProperties) {
this.kubernetesClient = kubernetesClient;
public LeaderInitiator(LeadershipController leadershipController, Candidate candidate,
LeaderProperties leaderProperties, ScheduledExecutorService scheduledExecutorService) {
this.leadershipController = leadershipController;
this.candidate = candidate;
this.leaderProperties = leaderProperties;
this.scheduledExecutorService = scheduledExecutorService;
}
@Override
@@ -61,27 +52,17 @@ public class LeaderInitiator implements SmartLifecycle {
@Override
public void start() {
lock.lock();
try {
if (!isRunning()) {
scheduledExecutorService = Executors.newSingleThreadScheduledExecutor();
scheduledExecutorService.execute(this::update);
}
} finally {
lock.unlock();
if (!isRunning()) {
scheduledExecutorService.execute(this::update);
isRunning = true;
}
}
@Override
public void stop() {
lock.lock();
try {
if (isRunning()) {
scheduledExecutorService.shutdown();
scheduledExecutorService = null;
}
} finally {
lock.unlock();
if (isRunning()) {
scheduledExecutorService.shutdown();
isRunning = false;
}
}
@@ -93,7 +74,7 @@ public class LeaderInitiator implements SmartLifecycle {
@Override
public boolean isRunning() {
return scheduledExecutorService != null;
return isRunning;
}
@Override
@@ -101,104 +82,17 @@ public class LeaderInitiator implements SmartLifecycle {
return 0; // TODO implement
}
public String getCurrentLeaderId() {
return currentLeaderId;
}
private void update() {
ConfigMap configMap = getConfigMap();
Leader leader = getLeader(configMap);
if (leader == null) {
System.out.println("Currently there is no leader, trying to become one");
try {
takeLeadership(configMap);
currentLeaderId = candidate.getId();
scheduleUpdate(leaderProperties.getLeaseDuration());
} catch (Exception e) {
// Leadership takeover failed, try again later
System.out.println("Leadership takeover failed: " + e.getMessage());
scheduleUpdate(leaderProperties.getRetryPeriod());
}
} else if (!isValidLeader(leader)) {
System.out.println("Old leader is not valid any more, try to take over");
try {
takeLeadership(configMap);
currentLeaderId = candidate.getId();
scheduleUpdate(leaderProperties.getLeaseDuration()); // Is this needed?
} catch (Exception e) {
// Leadership takeover failed, try again later
System.out.println("Leadership takeover failed: " + e.getMessage());
scheduleUpdate(leaderProperties.getRetryPeriod());
}
} else if (!isCandidateALeader(leader)) {
currentLeaderId = leader.getId();
System.out.println(currentLeaderId + " is a leader, check in later");
if (leadershipController.acquire(candidate)) {
// We're a leader, check-in later
scheduleUpdate(leaderProperties.getLeaseDuration());
} else {
System.out.println("I am a leader, check in later");
currentLeaderId = candidate.getId();
scheduleUpdate(leaderProperties.getLeaseDuration()); // Is this needed?
// Couldn't become a leader, retry sooner.
// TODO maybe we should separate error and another leader scenarios?
scheduleUpdate(leaderProperties.getRetryPeriod());
}
}
private void takeLeadership(ConfigMap oldConfigMap) {
String leaderIdKey = leaderProperties.getLeaderIdPrefix() + candidate.getRole();
if (oldConfigMap == null) {
ConfigMap newConfigMap = new ConfigMapBuilder().withNewMetadata()
.withName(leaderProperties.getConfigMapName())
.addToLabels("provider", "spring-cloud-kubernetes")
.addToLabels("kind", "locks")
.endMetadata()
.addToData(leaderIdKey, candidate.getId())
.build();
kubernetesClient.configMaps()
.inNamespace(leaderProperties.getNamespace(kubernetesClient.getNamespace()))
.create(newConfigMap);
} else {
ConfigMap newConfigMap = new ConfigMapBuilder(oldConfigMap)
.addToData(leaderIdKey, candidate.getId())
.build();
kubernetesClient.configMaps()
.inNamespace(leaderProperties.getNamespace(kubernetesClient.getNamespace()))
.withName(leaderProperties.getConfigMapName())
.lockResourceVersion(oldConfigMap.getMetadata().getResourceVersion())
.replace(newConfigMap);
}
System.out.println(candidate.getId() + " is now a leader");
}
private ConfigMap getConfigMap() {
try {
return kubernetesClient.configMaps()
.inNamespace(leaderProperties.getNamespace(kubernetesClient.getNamespace()))
.withName(leaderProperties.getConfigMapName())
.get();
} catch (Exception e) {
System.out.println("Failed to get a ConfigMap: " + e.getMessage());
return null;
}
}
private Leader getLeader(ConfigMap configMap) {
if (configMap == null || configMap.getData() == null) {
return null;
}
Map<String, String> data = configMap.getData();
String leaderIdKey = leaderProperties.getLeaderIdPrefix() + candidate.getRole();
String leaderId = data.get(leaderIdKey);
if (leaderId == null) {
return null;
}
return new Leader(candidate.getRole(), leaderId);
}
private void scheduleUpdate(long waitPeriod) {
scheduledExecutorService.schedule(this::update, jitter(waitPeriod), TimeUnit.MILLISECONDS);
}
@@ -207,20 +101,4 @@ public class LeaderInitiator implements SmartLifecycle {
return (long) (num * (1 + Math.random() * (leaderProperties.getJitterFactor() - 1)));
}
private boolean isCandidateALeader(Leader leader) {
return candidate.getId().equals(leader.getId());
}
private boolean isValidLeader(Leader leader) {
return kubernetesClient.pods()
.inNamespace(leaderProperties.getNamespace(kubernetesClient.getNamespace()))
.withLabels(leaderProperties.getLabels())
.list()
.getItems()
.stream()
.map(Pod::getMetadata)
.map(ObjectMeta::getName)
.anyMatch(name -> name.equals(leader.getId()));
}
}

View File

@@ -0,0 +1,115 @@
package org.springframework.cloud.kubernetes.leader;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import org.junit.Before;
import org.junit.Test;
import org.junit.runner.RunWith;
import org.mockito.ArgumentCaptor;
import org.mockito.Mock;
import org.mockito.junit.MockitoJUnitRunner;
import org.springframework.integration.leader.Candidate;
import static org.assertj.core.api.Assertions.assertThat;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.BDDMockito.given;
import static org.mockito.Mockito.verify;
/**
* @author <a href="mailto:gytis@redhat.com">Gytis Trikleris</a>
*/
@RunWith(MockitoJUnitRunner.class)
public class LeaderInitiatorTest {
@Mock
private Candidate mockCandidate;
@Mock
private LeadershipController mockLeadershipController;
@Mock
private LeaderProperties mockLeaderProperties;
@Mock
private ScheduledExecutorService mockScheduledExecutorService;
@Mock
private Runnable mockRunnable;
private LeaderInitiator leaderInitiator;
@Before
public void before() {
leaderInitiator = new LeaderInitiator(mockLeadershipController, mockCandidate, mockLeaderProperties,
mockScheduledExecutorService);
}
@Test
public void testIsAutoStartup() {
assertThat(leaderInitiator.isAutoStartup()).isFalse();
verify(mockLeaderProperties).isAutoStartup();
}
@Test
public void shouldStartOnlyOnce() {
leaderInitiator.start();
leaderInitiator.start();
assertThat(leaderInitiator.isRunning()).isTrue();
verify(mockScheduledExecutorService).execute(any());
}
@Test
public void shouldStopOnlyOnce() {
leaderInitiator.start();
leaderInitiator.stop();
leaderInitiator.stop();
assertThat(leaderInitiator.isRunning()).isFalse();
verify(mockScheduledExecutorService).shutdown();
}
@Test
public void shouldStopAndExecuteCallback() {
leaderInitiator.start();
leaderInitiator.stop(mockRunnable);
assertThat(leaderInitiator.isRunning()).isFalse();
verify(mockScheduledExecutorService).shutdown();
verify(mockRunnable).run();
}
@Test
public void shouldScheduleUpdateIfLeadershipIsAcquired() {
given(mockLeadershipController.acquire(mockCandidate)).willReturn(true);
given(mockLeaderProperties.getLeaseDuration()).willReturn(1000L);
given(mockLeaderProperties.getJitterFactor()).willReturn(1.0);
leaderInitiator.start();
ArgumentCaptor<Runnable> updateRunnableCaptor = ArgumentCaptor.forClass(Runnable.class);
verify(mockScheduledExecutorService).execute(updateRunnableCaptor.capture());
Runnable updateRunnable = updateRunnableCaptor.getValue();
updateRunnable.run();
verify(mockLeadershipController).acquire(mockCandidate);
verify(mockScheduledExecutorService).schedule(any(Runnable.class), eq(1000L), eq(TimeUnit.MILLISECONDS));
}
@Test
public void shouldScheduleUpdateIfLeadershipIsNotAcquired() {
given(mockLeaderProperties.getRetryPeriod()).willReturn(1000L);
given(mockLeaderProperties.getJitterFactor()).willReturn(1.0);
leaderInitiator.start();
ArgumentCaptor<Runnable> updateRunnableCaptor = ArgumentCaptor.forClass(Runnable.class);
verify(mockScheduledExecutorService).execute(updateRunnableCaptor.capture());
Runnable updateRunnable = updateRunnableCaptor.getValue();
updateRunnable.run();
verify(mockLeadershipController).acquire(mockCandidate);
verify(mockScheduledExecutorService).schedule(any(Runnable.class), eq(1000L), eq(TimeUnit.MILLISECONDS));
}
}