Add MembershipListenerAdaper class implementing and adapting the o.a.g.distributed.internal.MembershipListener interface.
The internal Apache Geode API MembershipListener interface is useful for receiving notification in the event a peer cache member is disconnected (forced disconnected) and then auto-reconnected to the cluster.
This commit is contained in:
@@ -0,0 +1,198 @@
|
||||
/*
|
||||
* Copyright 2020 the original author or authors.
|
||||
*
|
||||
* Licensed under the Apache License, Version 2.0 (the "License");
|
||||
* you may not use this file except in compliance with the License.
|
||||
* You may obtain a copy of the License at
|
||||
*
|
||||
* https://www.apache.org/licenses/LICENSE-2.0
|
||||
*
|
||||
* Unless required by applicable law or agreed to in writing, software
|
||||
* distributed under the License is distributed on an "AS IS" BASIS,
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express
|
||||
* or implied. See the License for the specific language governing
|
||||
* permissions and limitations under the License.
|
||||
*/
|
||||
package org.springframework.geode.distributed;
|
||||
|
||||
import java.util.List;
|
||||
import java.util.Optional;
|
||||
import java.util.Set;
|
||||
import java.util.function.BiConsumer;
|
||||
import java.util.function.Consumer;
|
||||
|
||||
import org.apache.geode.cache.Cache;
|
||||
import org.apache.geode.distributed.internal.DistributionManager;
|
||||
import org.apache.geode.distributed.internal.InternalDistributedSystem;
|
||||
import org.apache.geode.distributed.internal.MembershipListener;
|
||||
import org.apache.geode.distributed.internal.membership.InternalDistributedMember;
|
||||
|
||||
import org.springframework.geode.util.function.QuadConsumer;
|
||||
import org.springframework.geode.util.function.TriConsumer;
|
||||
|
||||
/**
|
||||
* A {@link MembershipListener} implementation using the
|
||||
* <a href="https://en.wikipedia.org/wiki/Adapter_pattern">Adapter Software Design Pattern</a>
|
||||
* to delegate membership event callbacks to a {@link Consumer} of those membership events.
|
||||
*
|
||||
* @author John Blum
|
||||
* @see java.util.function.BiConsumer
|
||||
* @see java.util.function.Consumer
|
||||
* @see org.apache.geode.distributed.internal.DistributionManager
|
||||
* @see org.apache.geode.distributed.internal.MembershipListener
|
||||
* @see org.apache.geode.distributed.internal.membership.InternalDistributedMember
|
||||
* @since 1.3.0
|
||||
*/
|
||||
@SuppressWarnings("unused")
|
||||
public class MembershipListenerAdapter implements MembershipListener {
|
||||
|
||||
/**
|
||||
* Factory method used to construct a new instance of the {@link MembershipListenerAdapter}.
|
||||
*
|
||||
* @return a new instance of {@link MembershipListenerAdapter}.
|
||||
*/
|
||||
public static MembershipListenerAdapter create() {
|
||||
return new MembershipListenerAdapter();
|
||||
}
|
||||
|
||||
private TriConsumer<DistributionManager, InternalDistributedMember, Boolean> memberDepartedConsumer =
|
||||
(manager, member, crashed) -> {};
|
||||
|
||||
private BiConsumer<DistributionManager, InternalDistributedMember> memberJoinedConsumer = (manager, member) -> {};
|
||||
|
||||
private QuadConsumer<DistributionManager, InternalDistributedMember, InternalDistributedMember, String> memberSuspectConsumer =
|
||||
(manage, member, suspect, reason) -> {};
|
||||
|
||||
private TriConsumer<DistributionManager, Set<InternalDistributedMember>, List<InternalDistributedMember>> quorumLostConsumer =
|
||||
(manager, failures, remaining) -> {};
|
||||
|
||||
/**
|
||||
* @inheritDoc
|
||||
*/
|
||||
@Override
|
||||
public void memberDeparted(DistributionManager distributionManager, InternalDistributedMember distributedMember,
|
||||
boolean crashed) {
|
||||
|
||||
this.memberDepartedConsumer.accept(distributionManager, distributedMember, crashed);
|
||||
}
|
||||
|
||||
/**
|
||||
* @inheritDoc
|
||||
*/
|
||||
@Override
|
||||
public void memberJoined(DistributionManager distributionManager, InternalDistributedMember distributedMember) {
|
||||
this.memberJoinedConsumer.accept(distributionManager, distributedMember);
|
||||
}
|
||||
|
||||
/**
|
||||
* @inheritDoc
|
||||
*/
|
||||
@Override
|
||||
public void memberSuspect(DistributionManager distributionManager, InternalDistributedMember distributedMember,
|
||||
InternalDistributedMember whoSuspected, String reason) {
|
||||
|
||||
this.memberSuspectConsumer.accept(distributionManager, distributedMember, whoSuspected, reason);
|
||||
}
|
||||
|
||||
/**
|
||||
* @inheritDoc
|
||||
*/
|
||||
@Override
|
||||
public void quorumLost(DistributionManager distributionManager, Set<InternalDistributedMember> failures,
|
||||
List<InternalDistributedMember> remaining) {
|
||||
|
||||
this.quorumLostConsumer.accept(distributionManager, failures, remaining);
|
||||
}
|
||||
|
||||
/**
|
||||
* Registers this {@link MembershipListener} with the given {@literal peer} {@link Cache}.
|
||||
*
|
||||
* @param peerCache {@literal peer} {@link Cache} on which to register this {@link MembershipListener}.
|
||||
* @return this {@link MembershipListenerAdapter}.
|
||||
* @see org.apache.geode.cache.Cache
|
||||
*/
|
||||
public MembershipListenerAdapter register(Cache peerCache) {
|
||||
|
||||
Optional.ofNullable(peerCache)
|
||||
.map(Cache::getDistributedSystem)
|
||||
.filter(InternalDistributedSystem.class::isInstance)
|
||||
.map(InternalDistributedSystem.class::cast)
|
||||
.map(InternalDistributedSystem::getDistributionManager)
|
||||
.ifPresent(distributionManager -> distributionManager
|
||||
.addMembershipListener(this));
|
||||
|
||||
return this;
|
||||
}
|
||||
|
||||
/**
|
||||
* Null-safe builder method used to add a {@link #memberDeparted(DistributionManager, InternalDistributedMember, boolean)}
|
||||
* {@link TriConsumer} event handler.
|
||||
*
|
||||
* @param memberDepartedConsumer {@link TriConsumer} handling {@literal memberDeparted} events.
|
||||
* @return this {@link MembershipListenerAdapter}.
|
||||
* @see org.springframework.geode.util.function.TriConsumer
|
||||
*/
|
||||
public MembershipListenerAdapter withMemberDepartedConsumer(
|
||||
TriConsumer<DistributionManager, InternalDistributedMember, Boolean> memberDepartedConsumer) {
|
||||
|
||||
if (memberDepartedConsumer != null) {
|
||||
this.memberDepartedConsumer = memberDepartedConsumer;
|
||||
}
|
||||
|
||||
return this;
|
||||
}
|
||||
|
||||
/**
|
||||
* Null-safe builder method used to add a {@link #memberJoined(DistributionManager, InternalDistributedMember)}
|
||||
* {@link BiConsumer} event handler.
|
||||
*
|
||||
* @param memberJoinedConsumer {@link BiConsumer} handling {@literal memberJoined} events.
|
||||
* @return this {@link MembershipListenerAdapter}.
|
||||
* @see java.util.function.BiConsumer
|
||||
*/
|
||||
public MembershipListenerAdapter withMemberJoinedConsumer(
|
||||
BiConsumer<DistributionManager, InternalDistributedMember> memberJoinedConsumer) {
|
||||
|
||||
if (memberJoinedConsumer != null) {
|
||||
this.memberJoinedConsumer = memberJoinedConsumer;
|
||||
}
|
||||
|
||||
return this;
|
||||
}
|
||||
|
||||
/**
|
||||
* Null-safe builder method used to add a {@link #memberSuspect(DistributionManager, InternalDistributedMember, InternalDistributedMember, String)}
|
||||
* {@link QuadConsumer} event handler.
|
||||
*
|
||||
* @param memberSuspectConsumer {@link QuadConsumer} handling {@literal memberSuspect} events.
|
||||
* @return this {@link MembershipListenerAdapter}.
|
||||
* @see org.springframework.geode.util.function.QuadConsumer
|
||||
*/
|
||||
public MembershipListenerAdapter withMemberSuspectConsumer(
|
||||
QuadConsumer<DistributionManager, InternalDistributedMember, InternalDistributedMember, String> memberSuspectConsumer) {
|
||||
|
||||
if (memberSuspectConsumer != null) {
|
||||
this.memberSuspectConsumer = memberSuspectConsumer;
|
||||
}
|
||||
|
||||
return this;
|
||||
}
|
||||
|
||||
/**
|
||||
* Null-safe build method used to add a {@link #quorumLost(DistributionManager, Set, List)} {@link TriConsumer}
|
||||
* event handler.
|
||||
*
|
||||
* @param quorumLostConsumer {@link TriConsumer} handling {@literal quorumLost} events.
|
||||
* @return this {@link MembershipListenerAdapter}.
|
||||
* @see org.springframework.geode.util.function.TriConsumer
|
||||
*/
|
||||
public MembershipListenerAdapter withQuorumLostConsumer(
|
||||
TriConsumer<DistributionManager, Set<InternalDistributedMember>, List<InternalDistributedMember>> quorumLostConsumer) {
|
||||
|
||||
if (quorumLostConsumer != null) {
|
||||
this.quorumLostConsumer = quorumLostConsumer;
|
||||
}
|
||||
|
||||
return this;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,235 @@
|
||||
/*
|
||||
* Copyright 2020 the original author or authors.
|
||||
*
|
||||
* Licensed under the Apache License, Version 2.0 (the "License");
|
||||
* you may not use this file except in compliance with the License.
|
||||
* You may obtain a copy of the License at
|
||||
*
|
||||
* https://www.apache.org/licenses/LICENSE-2.0
|
||||
*
|
||||
* Unless required by applicable law or agreed to in writing, software
|
||||
* distributed under the License is distributed on an "AS IS" BASIS,
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express
|
||||
* or implied. See the License for the specific language governing
|
||||
* permissions and limitations under the License.
|
||||
*/
|
||||
package org.springframework.geode.distributed;
|
||||
|
||||
import static org.assertj.core.api.Assertions.assertThat;
|
||||
import static org.mockito.ArgumentMatchers.eq;
|
||||
import static org.mockito.Mockito.doReturn;
|
||||
import static org.mockito.Mockito.mock;
|
||||
import static org.mockito.Mockito.times;
|
||||
import static org.mockito.Mockito.verify;
|
||||
import static org.mockito.Mockito.verifyNoInteractions;
|
||||
|
||||
import java.util.Collections;
|
||||
import java.util.List;
|
||||
import java.util.Set;
|
||||
import java.util.function.BiConsumer;
|
||||
|
||||
import org.junit.Test;
|
||||
import org.junit.runner.RunWith;
|
||||
import org.mockito.Mock;
|
||||
import org.mockito.junit.MockitoJUnitRunner;
|
||||
|
||||
import org.apache.geode.cache.Cache;
|
||||
import org.apache.geode.distributed.DistributedSystem;
|
||||
import org.apache.geode.distributed.internal.DistributionManager;
|
||||
import org.apache.geode.distributed.internal.InternalDistributedSystem;
|
||||
import org.apache.geode.distributed.internal.membership.InternalDistributedMember;
|
||||
|
||||
import org.springframework.geode.util.function.QuadConsumer;
|
||||
import org.springframework.geode.util.function.TriConsumer;
|
||||
|
||||
/**
|
||||
* Unit Tests for {@link MembershipListenerAdapter}.
|
||||
*
|
||||
* @author John Blum
|
||||
* @see org.junit.Test
|
||||
* @see org.mockito.Mockito
|
||||
* @see org.mockito.junit.MockitoJUnitRunner
|
||||
* @see org.apache.geode.distributed.internal.DistributionManager
|
||||
* @see org.apache.geode.distributed.internal.membership.InternalDistributedMember
|
||||
* @see org.springframework.geode.distributed.MembershipListenerAdapter
|
||||
* @since 1.3.0
|
||||
*/
|
||||
@RunWith(MockitoJUnitRunner.class)
|
||||
public class MembershipListenerAdapterUnitTests {
|
||||
|
||||
@Mock
|
||||
private DistributionManager mockDistributionManager;
|
||||
|
||||
@Mock
|
||||
private InternalDistributedMember mockDistributedMember;
|
||||
|
||||
@Test
|
||||
public void createConstructsNewMembershipListenerAdapterWithNoOpEventHandlers() {
|
||||
|
||||
InternalDistributedMember mockSuspectMember = mock(InternalDistributedMember.class);
|
||||
|
||||
MembershipListenerAdapter membershipListener = MembershipListenerAdapter.create();
|
||||
|
||||
assertThat(membershipListener).isNotNull();
|
||||
|
||||
membershipListener.memberDeparted(this.mockDistributionManager, this.mockDistributedMember, true);
|
||||
membershipListener.memberJoined(this.mockDistributionManager, this.mockDistributedMember);
|
||||
membershipListener.memberSuspect(this.mockDistributionManager, this.mockDistributedMember, mockSuspectMember,
|
||||
"System is unstable!!");
|
||||
membershipListener.quorumLost(this.mockDistributionManager, Collections.singleton(mockSuspectMember),
|
||||
Collections.singletonList(this.mockDistributedMember));
|
||||
|
||||
verifyNoInteractions(this.mockDistributedMember);
|
||||
verifyNoInteractions(this.mockDistributedMember);
|
||||
verifyNoInteractions(mockSuspectMember);
|
||||
}
|
||||
|
||||
@Test
|
||||
public void registersListenerWithPeerCache() {
|
||||
|
||||
Cache mockCache = mock(Cache.class);
|
||||
|
||||
DistributionManager mockDistributionManager = mock(DistributionManager.class);
|
||||
|
||||
InternalDistributedSystem mockDistributedSystem = mock(InternalDistributedSystem.class);
|
||||
|
||||
doReturn(mockDistributedSystem).when(mockCache).getDistributedSystem();
|
||||
doReturn(mockDistributionManager).when(mockDistributedSystem).getDistributionManager();
|
||||
|
||||
MembershipListenerAdapter membershipListener = new MembershipListenerAdapter();
|
||||
|
||||
assertThat(membershipListener.register(mockCache)).isSameAs(membershipListener);
|
||||
|
||||
verify(mockCache, times(1)).getDistributedSystem();
|
||||
verify(mockDistributedSystem, times(1)).getDistributionManager();
|
||||
verify(mockDistributionManager, times(1)).addMembershipListener(eq(membershipListener));
|
||||
}
|
||||
|
||||
@Test
|
||||
public void registerListenerWithNullCacheIsNullSafe() {
|
||||
|
||||
MembershipListenerAdapter membershipListener = new MembershipListenerAdapter();
|
||||
|
||||
assertThat(membershipListener.register(null)).isSameAs(membershipListener);
|
||||
}
|
||||
|
||||
@Test
|
||||
public void registerListenerWithNullDistributedSystemIsNullSafe() {
|
||||
|
||||
Cache mockCache = mock(Cache.class);
|
||||
|
||||
MembershipListenerAdapter membershipListener = new MembershipListenerAdapter();
|
||||
|
||||
assertThat(membershipListener.register(mockCache)).isSameAs(membershipListener);
|
||||
|
||||
verify(mockCache, times(1)).getDistributedSystem();
|
||||
}
|
||||
|
||||
@Test
|
||||
public void registerListenerWithNonInternalDistributedSystemIsSafe() {
|
||||
|
||||
Cache mockCache = mock(Cache.class);
|
||||
|
||||
DistributedSystem mockDistributedSystem = mock(DistributedSystem.class);
|
||||
|
||||
MembershipListenerAdapter membershipListener = new MembershipListenerAdapter();
|
||||
|
||||
assertThat(membershipListener.register(mockCache)).isSameAs(membershipListener);
|
||||
|
||||
verify(mockCache, times(1)).getDistributedSystem();
|
||||
verifyNoInteractions(mockDistributedSystem);
|
||||
}
|
||||
|
||||
@Test
|
||||
public void registerListenerWithNullDistributionManagerIsNullSafe() {
|
||||
|
||||
Cache mockCache = mock(Cache.class);
|
||||
|
||||
InternalDistributedSystem mockDistributedSystem = mock(InternalDistributedSystem.class);
|
||||
|
||||
doReturn(mockDistributedSystem).when(mockCache).getDistributedSystem();
|
||||
|
||||
MembershipListenerAdapter membershipListener = new MembershipListenerAdapter();
|
||||
|
||||
assertThat(membershipListener.register(mockCache)).isSameAs(membershipListener);
|
||||
|
||||
verify(mockCache, times(1)).getDistributedSystem();
|
||||
verify(mockDistributedSystem, times(1)).getDistributionManager();
|
||||
}
|
||||
|
||||
@Test
|
||||
@SuppressWarnings("unchecked")
|
||||
public void withMemberDepartedEventHandler() {
|
||||
|
||||
TriConsumer<DistributionManager, InternalDistributedMember, Boolean> mockConsumer = mock(TriConsumer.class);
|
||||
|
||||
MembershipListenerAdapter membershipListener = MembershipListenerAdapter.create()
|
||||
.withMemberDepartedConsumer(mockConsumer);
|
||||
|
||||
assertThat(membershipListener).isNotNull();
|
||||
|
||||
membershipListener.memberDeparted(this.mockDistributionManager, this.mockDistributedMember, true);
|
||||
|
||||
verify(mockConsumer, times(1))
|
||||
.accept(eq(this.mockDistributionManager), eq(this.mockDistributedMember), eq(true));
|
||||
}
|
||||
|
||||
@Test
|
||||
@SuppressWarnings("unchecked")
|
||||
public void withMemberJoinedEventHandler() {
|
||||
|
||||
BiConsumer<DistributionManager, InternalDistributedMember> mockConsumer = mock(BiConsumer.class);
|
||||
|
||||
MembershipListenerAdapter membershipListener = MembershipListenerAdapter.create()
|
||||
.withMemberJoinedConsumer(mockConsumer);
|
||||
|
||||
assertThat(membershipListener).isNotNull();
|
||||
|
||||
membershipListener.memberJoined(this.mockDistributionManager, this.mockDistributedMember);
|
||||
|
||||
verify(mockConsumer, times(1))
|
||||
.accept(eq(this.mockDistributionManager), eq(this.mockDistributedMember));
|
||||
}
|
||||
|
||||
@Test
|
||||
@SuppressWarnings("unchecked")
|
||||
public void withMemberSuspectEventHandler() {
|
||||
|
||||
QuadConsumer<DistributionManager, InternalDistributedMember, InternalDistributedMember, String> mockConsumer =
|
||||
mock(QuadConsumer.class);
|
||||
|
||||
InternalDistributedMember mockSuspectMember = mock(InternalDistributedMember.class);
|
||||
|
||||
MembershipListenerAdapter membershipListener = MembershipListenerAdapter.create()
|
||||
.withMemberSuspectConsumer(mockConsumer);
|
||||
|
||||
assertThat(membershipListener).isNotNull();
|
||||
|
||||
membershipListener.memberSuspect(this.mockDistributionManager, this.mockDistributedMember, mockSuspectMember,
|
||||
"System is a lost cause!!");
|
||||
|
||||
verify(mockConsumer, times(1))
|
||||
.accept(eq(this.mockDistributionManager), eq(this.mockDistributedMember), eq(mockSuspectMember),
|
||||
eq("System is a lost cause!!"));
|
||||
}
|
||||
|
||||
@Test
|
||||
@SuppressWarnings("unchecked")
|
||||
public void withQuorumLostEventHandler() {
|
||||
|
||||
TriConsumer<DistributionManager, Set<InternalDistributedMember>, List<InternalDistributedMember>> mockConsumer
|
||||
= mock(TriConsumer.class);
|
||||
|
||||
MembershipListenerAdapter membershipListener = MembershipListenerAdapter.create()
|
||||
.withQuorumLostConsumer(mockConsumer);
|
||||
|
||||
assertThat(membershipListener).isNotNull();
|
||||
|
||||
membershipListener.quorumLost(this.mockDistributionManager, Collections.singleton(this.mockDistributedMember),
|
||||
Collections.singletonList(this.mockDistributedMember));
|
||||
|
||||
verify(mockConsumer, times(1)).accept(eq(this.mockDistributionManager),
|
||||
eq(Collections.singleton(this.mockDistributedMember)),
|
||||
eq(Collections.singletonList(this.mockDistributedMember)));
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user