import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertFalse;
import static org.junit.Assert.assertTrue;
-import static org.mockito.Matchers.any;
-import static org.mockito.Matchers.anyInt;
-import static org.mockito.Matchers.anyObject;
-import static org.mockito.Mockito.doAnswer;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyInt;
import static org.mockito.Mockito.doNothing;
import static org.mockito.Mockito.doReturn;
+import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.times;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.verifyNoMoreInteractions;
+import akka.actor.ActorRef;
+import akka.actor.ActorSystem;
+import akka.actor.Props;
import akka.persistence.RecoveryCompleted;
import akka.persistence.SnapshotMetadata;
import akka.persistence.SnapshotOffer;
import com.google.common.collect.Sets;
-import java.io.Serializable;
+import com.google.common.util.concurrent.MoreExecutors;
+import java.io.OutputStream;
import java.util.Arrays;
import java.util.Collections;
-import java.util.List;
-import org.apache.commons.lang3.SerializationUtils;
-import org.hamcrest.Description;
+import java.util.Optional;
+import java.util.concurrent.ExecutionException;
+import java.util.concurrent.Executors;
+import java.util.concurrent.ScheduledExecutorService;
+import java.util.concurrent.ScheduledFuture;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicInteger;
+import java.util.function.Consumer;
+import org.junit.Assert;
import org.junit.Before;
import org.junit.Test;
-import org.mockito.ArgumentMatcher;
+import org.mockito.ArgumentMatchers;
import org.mockito.InOrder;
-import org.mockito.Matchers;
import org.mockito.Mock;
import org.mockito.Mockito;
import org.mockito.MockitoAnnotations;
import org.opendaylight.controller.cluster.raft.persisted.Snapshot;
import org.opendaylight.controller.cluster.raft.persisted.UpdateElectionTerm;
import org.opendaylight.controller.cluster.raft.protobuff.client.messages.Payload;
+import org.opendaylight.controller.cluster.raft.utils.DoNothingActor;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@Mock
private DataPersistenceProvider mockPersistence;
-
@Mock
private RaftActorRecoveryCohort mockCohort;
- @Mock
- private RaftActorSnapshotCohort mockSnapshotCohort;
-
@Mock
PersistentDataProvider mockPersistentProvider;
+ ActorRef mockActorRef;
+
+ ActorSystem mockActorSystem;
+
private RaftActorRecoverySupport support;
private RaftActorContext context;
private final DefaultConfigParamsImpl configParams = new DefaultConfigParamsImpl();
private final String localId = "leader";
-
@Before
public void setup() {
MockitoAnnotations.initMocks(this);
-
- context = new RaftActorContextImpl(null, null, localId, new ElectionTermImpl(mockPersistentProvider, "test",
- LOG), -1, -1, Collections.<String,String>emptyMap(), configParams,
- mockPersistence, applyState -> { }, LOG);
+ mockActorSystem = ActorSystem.create();
+ mockActorRef = mockActorSystem.actorOf(Props.create(DoNothingActor.class));
+ context = new RaftActorContextImpl(mockActorRef, null, localId,
+ new ElectionTermImpl(mockPersistentProvider, "test", LOG), -1, -1,
+ Collections.<String, String>emptyMap(), configParams, mockPersistence, applyState -> {
+ }, LOG, MoreExecutors.directExecutor());
support = new RaftActorRecoverySupport(context, mockCohort);
context.setReplicatedLog(ReplicatedLogImpl.newInstance(context));
}
- private void sendMessageToSupport(Object message) {
+ private void sendMessageToSupport(final Object message) {
sendMessageToSupport(message, false);
}
- private void sendMessageToSupport(Object message, boolean expComplete) {
+ private void sendMessageToSupport(final Object message, final boolean expComplete) {
boolean complete = support.handleRecoveryMessage(message, mockPersistentProvider);
assertEquals("complete", expComplete, complete);
}
inOrder.verifyNoMoreInteractions();
}
+ @Test
+ public void testIncrementalRecovery() {
+ int recoverySnapshotInterval = 3;
+ int numberOfEntries = 5;
+ configParams.setRecoverySnapshotIntervalSeconds(recoverySnapshotInterval);
+ Consumer<Optional<OutputStream>> mockSnapshotConsumer = mock(Consumer.class);
+ context.getSnapshotManager().setCreateSnapshotConsumer(mockSnapshotConsumer);
+
+ ScheduledExecutorService applyEntriesExecutor = Executors.newSingleThreadScheduledExecutor();
+ ReplicatedLog replicatedLog = context.getReplicatedLog();
+
+ for (int i = 0; i <= numberOfEntries; i++) {
+ replicatedLog.append(new SimpleReplicatedLogEntry(i, 1,
+ new MockRaftActorContext.MockPayload(String.valueOf(i))));
+ }
+
+ AtomicInteger entryCount = new AtomicInteger();
+ ScheduledFuture<?> applyEntriesFuture = applyEntriesExecutor.scheduleAtFixedRate(() -> {
+ int run = entryCount.getAndIncrement();
+ LOG.info("Sending entry number {}", run);
+ sendMessageToSupport(new ApplyJournalEntries(run));
+ }, 0, 1, TimeUnit.SECONDS);
+
+ ScheduledFuture<Boolean> canceller = applyEntriesExecutor.schedule(() -> applyEntriesFuture.cancel(false),
+ numberOfEntries, TimeUnit.SECONDS);
+ try {
+ canceller.get();
+ verify(mockSnapshotConsumer, times(1)).accept(any());
+ applyEntriesExecutor.shutdown();
+ } catch (InterruptedException | ExecutionException e) {
+ Assert.fail();
+ }
+ }
+
@Test
public void testOnSnapshotOffer() {
lastAppliedDuringSnapshotCapture, 1, electionTerm, electionVotedFor, null);
SnapshotMetadata metadata = new SnapshotMetadata("test", 6, 12345);
- SnapshotOffer snapshotOffer = new SnapshotOffer(metadata , snapshot);
-
- sendMessageToSupport(snapshotOffer);
-
- assertEquals("Journal log size", 2, context.getReplicatedLog().size());
- assertEquals("Journal data size", 9, context.getReplicatedLog().dataSize());
- assertEquals("Last index", lastIndexDuringSnapshotCapture, context.getReplicatedLog().lastIndex());
- assertEquals("Last applied", lastAppliedDuringSnapshotCapture, context.getLastApplied());
- assertEquals("Commit index", lastAppliedDuringSnapshotCapture, context.getCommitIndex());
- assertEquals("Snapshot term", 1, context.getReplicatedLog().getSnapshotTerm());
- assertEquals("Snapshot index", lastAppliedDuringSnapshotCapture, context.getReplicatedLog().getSnapshotIndex());
- assertEquals("Election term", electionTerm, context.getTermInformation().getCurrentTerm());
- assertEquals("Election votedFor", electionVotedFor, context.getTermInformation().getVotedFor());
- assertFalse("Dynamic server configuration", context.isDynamicServerConfigurationInUse());
-
- verify(mockCohort).applyRecoverySnapshot(snapshotState);
- }
-
- @Deprecated
- @Test
- public void testOnSnapshotOfferWithPreCarbonSnapshot() {
-
- ReplicatedLogEntry unAppliedEntry1 = new SimpleReplicatedLogEntry(4, 1,
- new MockRaftActorContext.MockPayload("4", 4));
-
- ReplicatedLogEntry unAppliedEntry2 = new SimpleReplicatedLogEntry(5, 1,
- new MockRaftActorContext.MockPayload("5", 5));
-
- long lastAppliedDuringSnapshotCapture = 3;
- long lastIndexDuringSnapshotCapture = 5;
- long electionTerm = 2;
- String electionVotedFor = "member-2";
-
- List<Object> snapshotData = Arrays.asList(new MockPayload("1"));
- final MockSnapshotState snapshotState = new MockSnapshotState(snapshotData);
-
- org.opendaylight.controller.cluster.raft.Snapshot snapshot = org.opendaylight.controller.cluster.raft.Snapshot
- .create(SerializationUtils.serialize((Serializable) snapshotData),
- Arrays.asList(unAppliedEntry1, unAppliedEntry2), lastIndexDuringSnapshotCapture, 1,
- lastAppliedDuringSnapshotCapture, 1, electionTerm, electionVotedFor, null);
-
- SnapshotMetadata metadata = new SnapshotMetadata("test", 6, 12345);
- SnapshotOffer snapshotOffer = new SnapshotOffer(metadata , snapshot);
-
- doAnswer(invocation -> new MockSnapshotState(SerializationUtils.deserialize(
- invocation.getArgumentAt(0, byte[].class))))
- .when(mockCohort).deserializePreCarbonSnapshot(any(byte[].class));
+ SnapshotOffer snapshotOffer = new SnapshotOffer(metadata, snapshot);
sendMessageToSupport(snapshotOffer);
@Test
public void testDataRecoveredWithPersistenceDisabled() {
- doNothing().when(mockCohort).applyRecoverySnapshot(anyObject());
+ doNothing().when(mockCohort).applyRecoverySnapshot(any());
doReturn(false).when(mockPersistence).isRecoveryApplicable();
doReturn(10L).when(mockPersistentProvider).getLastSequenceNumber();
sendMessageToSupport(RecoveryCompleted.getInstance(), true);
- verify(mockCohort, never()).applyRecoverySnapshot(anyObject());
+ verify(mockCohort, never()).applyRecoverySnapshot(any());
verify(mockCohort, never()).getRestoreFromSnapshot();
verifyNoMoreInteractions(mockCohort);
}
static UpdateElectionTerm updateElectionTerm(final long term, final String votedFor) {
- return Matchers.argThat(new ArgumentMatcher<UpdateElectionTerm>() {
- @Override
- public boolean matches(Object argument) {
- UpdateElectionTerm other = (UpdateElectionTerm) argument;
- return term == other.getCurrentTerm() && votedFor.equals(other.getVotedFor());
- }
-
- @Override
- public void describeTo(Description description) {
- description.appendValue(new UpdateElectionTerm(term, votedFor));
- }
- });
+ return ArgumentMatchers.argThat(other ->
+ term == other.getCurrentTerm() && votedFor.equals(other.getVotedFor()));
}
@Test
long electionTerm = 2;
String electionVotedFor = "member-2";
ServerConfigurationPayload serverPayload = new ServerConfigurationPayload(Arrays.asList(
- new ServerInfo(localId, true),
- new ServerInfo("follower1", true),
- new ServerInfo("follower2", true)));
+ new ServerInfo(localId, true),
+ new ServerInfo("follower1", true),
+ new ServerInfo("follower2", true)));
MockSnapshotState snapshotState = new MockSnapshotState(Arrays.asList(new MockPayload("1")));
Snapshot snapshot = Snapshot.create(snapshotState, Collections.<ReplicatedLogEntry>emptyList(),
-1, -1, -1, -1, electionTerm, electionVotedFor, serverPayload);
SnapshotMetadata metadata = new SnapshotMetadata("test", 6, 12345);
- SnapshotOffer snapshotOffer = new SnapshotOffer(metadata , snapshot);
+ SnapshotOffer snapshotOffer = new SnapshotOffer(metadata, snapshot);
sendMessageToSupport(snapshotOffer);
assertEquals("Election votedFor", electionVotedFor, context.getTermInformation().getVotedFor());
assertTrue("Dynamic server configuration", context.isDynamicServerConfigurationInUse());
assertEquals("Peer List", Sets.newHashSet("follower1", "follower2"),
- Sets.newHashSet(context.getPeerIds()));
+ Sets.newHashSet(context.getPeerIds()));
}
-}
+}
\ No newline at end of file