-
Notifications
You must be signed in to change notification settings - Fork 215
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
Issue #559: Introduced factory for creating shard region proxy actors.
Signed-off-by: Juergen Fickel <juergen.fickel@bosch.io>
- Loading branch information
Juergen Fickel
committed
Oct 29, 2021
1 parent
7a68338
commit 32c5b1f
Showing
5 changed files
with
245 additions
and
36 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
75 changes: 75 additions & 0 deletions
75
.../src/main/java/org/eclipse/ditto/internal/utils/cluster/ShardRegionProxyActorFactory.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,75 @@ | ||
/* | ||
* Copyright (c) 2021 Contributors to the Eclipse Foundation | ||
* | ||
* See the NOTICE file(s) distributed with this work for additional | ||
* information regarding copyright ownership. | ||
* | ||
* This program and the accompanying materials are made available under the | ||
* terms of the Eclipse Public License 2.0 which is available at | ||
* http://www.eclipse.org/legal/epl-2.0 | ||
* | ||
* SPDX-License-Identifier: EPL-2.0 | ||
*/ | ||
package org.eclipse.ditto.internal.utils.cluster; | ||
|
||
import java.util.Optional; | ||
|
||
import javax.annotation.concurrent.NotThreadSafe; | ||
|
||
import org.eclipse.ditto.base.model.common.ConditionChecker; | ||
import org.eclipse.ditto.internal.utils.cluster.config.ClusterConfig; | ||
|
||
import akka.actor.ActorRef; | ||
import akka.actor.ActorSystem; | ||
import akka.cluster.sharding.ClusterSharding; | ||
|
||
/** | ||
* Factory for creating shard region proxy actors. | ||
* | ||
* @since 2.2.0 | ||
*/ | ||
@NotThreadSafe | ||
public final class ShardRegionProxyActorFactory { | ||
|
||
private final ShardRegionExtractor extractor; | ||
private final ClusterSharding clusterSharding; | ||
|
||
private ShardRegionProxyActorFactory(final ShardRegionExtractor extractor, final ClusterSharding clusterSharding) { | ||
this.extractor = extractor; | ||
this.clusterSharding = clusterSharding; | ||
} | ||
|
||
/** | ||
* Returns a new instance of {@code ShardRegionFactory}. | ||
* | ||
* @return the instance. | ||
* @throws NullPointerException if any argument is {@code null}. | ||
*/ | ||
public static ShardRegionProxyActorFactory newInstance(final ActorSystem actorSystem, | ||
final ClusterConfig clusterConfig) { | ||
|
||
ConditionChecker.checkNotNull(actorSystem, "actorSystem"); | ||
ConditionChecker.checkNotNull(clusterConfig, "clusterConfig"); | ||
|
||
return new ShardRegionProxyActorFactory(ShardRegionExtractor.of(clusterConfig.getNumberOfShards(), actorSystem), | ||
ClusterSharding.get(actorSystem)); | ||
} | ||
|
||
/** | ||
* Starts a proxy of a shard region specified by the cluster role and shard region name arguments. | ||
* The actor reference is not being cached. | ||
* | ||
* @param shardRegionName name of the shard region. | ||
* @param clusterRole role of cluster members where the shard region resides. | ||
* @return reference of the shard region proxy. | ||
* @throws NullPointerException if any argument is {@code null}. | ||
* @throws IllegalArgumentException if any argument is empty. | ||
*/ | ||
public ActorRef getShardRegionProxyActor(final CharSequence clusterRole, final CharSequence shardRegionName) { | ||
ConditionChecker.argumentNotEmpty(clusterRole, "clusterRole"); | ||
ConditionChecker.argumentNotEmpty(shardRegionName, "shardRegionName"); | ||
|
||
return clusterSharding.startProxy(shardRegionName.toString(), Optional.of(clusterRole.toString()), extractor); | ||
} | ||
|
||
} |
145 changes: 145 additions & 0 deletions
145
.../test/java/org/eclipse/ditto/internal/utils/cluster/ShardRegionProxyActorFactoryTest.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,145 @@ | ||
/* | ||
* Copyright (c) 2021 Contributors to the Eclipse Foundation | ||
* | ||
* See the NOTICE file(s) distributed with this work for additional | ||
* information regarding copyright ownership. | ||
* | ||
* This program and the accompanying materials are made available under the | ||
* terms of the Eclipse Public License 2.0 which is available at | ||
* http://www.eclipse.org/legal/epl-2.0 | ||
* | ||
* SPDX-License-Identifier: EPL-2.0 | ||
*/ | ||
package org.eclipse.ditto.internal.utils.cluster; | ||
|
||
import static org.assertj.core.api.Assertions.assertThat; | ||
import static org.assertj.core.api.Assertions.assertThatIllegalArgumentException; | ||
import static org.assertj.core.api.Assertions.assertThatNullPointerException; | ||
|
||
import java.util.Map; | ||
|
||
import org.eclipse.ditto.internal.utils.akka.ActorSystemResource; | ||
import org.eclipse.ditto.internal.utils.cluster.config.ClusterConfig; | ||
import org.junit.Before; | ||
import org.junit.ClassRule; | ||
import org.junit.Test; | ||
import org.junit.runner.RunWith; | ||
import org.mockito.Mock; | ||
import org.mockito.Mockito; | ||
import org.mockito.junit.MockitoJUnitRunner; | ||
|
||
import com.typesafe.config.ConfigFactory; | ||
|
||
/** | ||
* Unit test for {@link ShardRegionProxyActorFactoryTest}. | ||
*/ | ||
@RunWith(MockitoJUnitRunner.class) | ||
public final class ShardRegionProxyActorFactoryTest { | ||
|
||
@ClassRule | ||
public static final ActorSystemResource ACTOR_SYSTEM_RESOURCE = ActorSystemResource.newInstance( | ||
ConfigFactory.parseMap(Map.ofEntries( | ||
Map.entry(MappingStrategies.CONFIG_KEY_DITTO_MAPPING_STRATEGY_IMPLEMENTATION, | ||
TestMappingStrategies.class.getName()), | ||
Map.entry("akka.actor.provider", "cluster") | ||
)) | ||
); | ||
|
||
private static final int NUMBER_OF_SHARDS = 3; | ||
|
||
@Mock | ||
private ClusterConfig clusterConfig; | ||
|
||
@Before | ||
public void before() { | ||
Mockito.when(clusterConfig.getNumberOfShards()).thenReturn(NUMBER_OF_SHARDS); | ||
} | ||
|
||
@Test | ||
public void newInstanceWithNullActorSystemThrowsException() { | ||
assertThatNullPointerException() | ||
.isThrownBy(() -> ShardRegionProxyActorFactory.newInstance(null, clusterConfig)) | ||
.withMessage("The actorSystem must not be null!") | ||
.withNoCause(); | ||
} | ||
|
||
@Test | ||
public void newInstanceWithNullClusterConfigThrowsException() { | ||
assertThatNullPointerException() | ||
.isThrownBy(() -> ShardRegionProxyActorFactory.newInstance(ACTOR_SYSTEM_RESOURCE.getActorSystem(), | ||
null)) | ||
.withMessage("The clusterConfig must not be null!") | ||
.withNoCause(); | ||
} | ||
|
||
@Test | ||
public void newInstanceReturnsNotNull() { | ||
final var underTest = | ||
ShardRegionProxyActorFactory.newInstance(ACTOR_SYSTEM_RESOURCE.getActorSystem(), clusterConfig); | ||
|
||
assertThat(underTest).isNotNull(); | ||
} | ||
|
||
@Test | ||
public void getShardRegionProxyActorWithNullClusterRoleThrowsException() { | ||
final var underTest = | ||
ShardRegionProxyActorFactory.newInstance(ACTOR_SYSTEM_RESOURCE.getActorSystem(), clusterConfig); | ||
|
||
assertThatNullPointerException() | ||
.isThrownBy(() -> underTest.getShardRegionProxyActor(null, "myShardRegion")) | ||
.withMessage("The clusterRole must not be null!") | ||
.withNoCause(); | ||
} | ||
|
||
@Test | ||
public void getShardRegionProxyActorWithEmptyClusterRoleThrowsException() { | ||
final var underTest = | ||
ShardRegionProxyActorFactory.newInstance(ACTOR_SYSTEM_RESOURCE.getActorSystem(), clusterConfig); | ||
|
||
assertThatIllegalArgumentException() | ||
.isThrownBy(() -> underTest.getShardRegionProxyActor("", "myShardRegion")) | ||
.withMessage("The argument 'clusterRole' must not be empty!") | ||
.withNoCause(); | ||
} | ||
|
||
@Test | ||
public void getShardRegionProxyActorWithNullShardRegionNameThrowsException() { | ||
final var underTest = | ||
ShardRegionProxyActorFactory.newInstance(ACTOR_SYSTEM_RESOURCE.getActorSystem(), clusterConfig); | ||
|
||
assertThatNullPointerException() | ||
.isThrownBy(() -> underTest.getShardRegionProxyActor("myClusterRole", null)) | ||
.withMessage("The shardRegionName must not be null!") | ||
.withNoCause(); | ||
} | ||
|
||
@Test | ||
public void getShardRegionProxyActorWithEmptyShardRegionNameThrowsException() { | ||
final var underTest = | ||
ShardRegionProxyActorFactory.newInstance(ACTOR_SYSTEM_RESOURCE.getActorSystem(), clusterConfig); | ||
|
||
assertThatIllegalArgumentException() | ||
.isThrownBy(() -> underTest.getShardRegionProxyActor("myClusterRole", "")) | ||
.withMessage("The argument 'shardRegionName' must not be empty!") | ||
.withNoCause(); | ||
} | ||
|
||
@Test | ||
public void getShardRegionProxyActorReturnsNotNull() { | ||
final var underTest = | ||
ShardRegionProxyActorFactory.newInstance(ACTOR_SYSTEM_RESOURCE.getActorSystem(), clusterConfig); | ||
|
||
final var shardRegionProxyActor = underTest.getShardRegionProxyActor("myClusterRole", "myShardRegionName"); | ||
|
||
assertThat(shardRegionProxyActor).isNotNull(); | ||
} | ||
|
||
public static final class TestMappingStrategies extends MappingStrategies { | ||
|
||
public TestMappingStrategies() { | ||
super(Map.of()); | ||
} | ||
|
||
} | ||
|
||
} |