import akka.actor.AddressFromURIString;
import akka.cluster.Cluster;
import akka.dispatch.Futures;
-import akka.pattern.AskTimeoutException;
import akka.pattern.Patterns;
import akka.testkit.JavaTestKit;
import com.google.common.base.Optional;
@Parameters(name = "{0}")
public static Collection<Object[]> data() {
return Arrays.asList(new Object[][] {
- { DistributedDataStore.class }, { ClientBackedDataStore.class }
+ { DistributedDataStore.class, 7}, { ClientBackedDataStore.class, 60 }
});
}
- @Parameter
+ @Parameter(0)
public Class<? extends AbstractDataStore> testParameter;
+ @Parameter(1)
+ public int commitTimeout;
private static final String[] CARS_AND_PEOPLE = {"cars", "people"};
private static final String[] CARS = {"cars"};
private void initDatastores(final String type, final String moduleShardsConfig, final String[] shards)
throws Exception {
- leaderTestKit = new IntegrationTestKit(leaderSystem, leaderDatastoreContextBuilder);
+ leaderTestKit = new IntegrationTestKit(leaderSystem, leaderDatastoreContextBuilder, commitTimeout);
leaderDistributedDataStore = leaderTestKit.setupAbstractDataStore(
testParameter, type, moduleShardsConfig, false, shards);
- followerTestKit = new IntegrationTestKit(followerSystem, followerDatastoreContextBuilder);
+ followerTestKit = new IntegrationTestKit(followerSystem, followerDatastoreContextBuilder, commitTimeout);
followerDistributedDataStore = followerTestKit.setupAbstractDataStore(
testParameter, type, moduleShardsConfig, false, shards);
final ActorSystem newSystem = newActorSystem("reinstated-member2", "Member2");
- try (final AbstractDataStore member2Datastore = new IntegrationTestKit(newSystem, leaderDatastoreContextBuilder)
+ try (AbstractDataStore member2Datastore = new IntegrationTestKit(newSystem, leaderDatastoreContextBuilder,
+ commitTimeout)
.setupAbstractDataStore(testParameter, testName, "module-shards-member2", true, CARS)) {
verifyCars(member2Datastore.newReadOnlyTransaction(), car2);
}
@Test
public void testSingleShardTransactionsWithLeaderChanges() throws Exception {
- //TODO remove when test passes also for ClientBackedDataStore
- Assume.assumeTrue(testParameter.equals(DistributedDataStore.class));
final String testName = "testSingleShardTransactionsWithLeaderChanges";
initDatastoresWithCars(testName);
final DatastoreContext.Builder newMember1Builder = DatastoreContext.newBuilder()
.shardHeartbeatIntervalInMillis(100).shardElectionTimeoutFactor(5);
- IntegrationTestKit newMember1TestKit = new IntegrationTestKit(leaderSystem, newMember1Builder);
+ IntegrationTestKit newMember1TestKit = new IntegrationTestKit(leaderSystem, newMember1Builder, commitTimeout);
try (final AbstractDataStore ds =
newMember1TestKit.setupAbstractDataStore(
initDatastores(testName, MODULE_SHARDS_CARS_PEOPLE_1_2_3, CARS_AND_PEOPLE);
final IntegrationTestKit follower2TestKit = new IntegrationTestKit(follower2System,
- DatastoreContext.newBuilderFrom(followerDatastoreContextBuilder.build()).operationTimeoutInMillis(100));
+ DatastoreContext.newBuilderFrom(followerDatastoreContextBuilder.build()).operationTimeoutInMillis(100),
+ commitTimeout);
try (final AbstractDataStore follower2DistributedDataStore = follower2TestKit.setupAbstractDataStore(
testParameter, testName, MODULE_SHARDS_CARS_PEOPLE_1_2_3, false)) {
leaderTestKit.doCommit(successTxCohort);
}
- @Test(expected = AskTimeoutException.class)
+ @Test
public void testTransactionWithShardLeaderNotResponding() throws Exception {
- //TODO remove when test passes also for ClientBackedDataStore
- Assume.assumeTrue(testParameter.equals(DistributedDataStore.class));
followerDatastoreContextBuilder.shardElectionTimeoutFactor(50);
initDatastoresWithCars("testTransactionWithShardLeaderNotResponding");
try {
followerTestKit.doCommit(rwTx.ready());
+ fail("Exception expected");
} catch (final ExecutionException e) {
- assertTrue("Expected ShardLeaderNotRespondingException cause. Actual: " + e.getCause(),
- e.getCause() instanceof ShardLeaderNotRespondingException);
- assertNotNull("Expected a nested cause", e.getCause().getCause());
- Throwables.propagateIfInstanceOf(e.getCause().getCause(), Exception.class);
- Throwables.propagate(e.getCause().getCause());
+ final String msg = "Unexpected exception: " + Throwables.getStackTraceAsString(e.getCause());
+ assertTrue(msg, Throwables.getRootCause(e) instanceof NoShardLeaderException
+ || e.getCause() instanceof ShardLeaderNotRespondingException);
}
}
- @Test(expected = NoShardLeaderException.class)
+ @Test
public void testTransactionWithCreateTxFailureDueToNoLeader() throws Exception {
- //TODO remove when test passes also for ClientBackedDataStore
- Assume.assumeTrue(testParameter.equals(DistributedDataStore.class));
initDatastoresWithCars("testTransactionWithCreateTxFailureDueToNoLeader");
// Do an initial read to get the primary shard info cached.
try {
followerTestKit.doCommit(rwTx.ready());
+ fail("Exception expected");
} catch (final ExecutionException e) {
- Throwables.propagateIfInstanceOf(e.getCause(), Exception.class);
- Throwables.propagate(e.getCause());
+ final String msg = "Expected instance of NoShardLeaderException, actual: \n"
+ + Throwables.getStackTraceAsString(e.getCause());
+ assertTrue(msg, Throwables.getRootCause(e) instanceof NoShardLeaderException);
}
}
@Test
public void testTransactionRetryWithInitialAskTimeoutExOnCreateTx() throws Exception {
- //TODO remove when test passes also for ClientBackedDataStore
- Assume.assumeTrue(testParameter.equals(DistributedDataStore.class));
String testName = "testTransactionRetryWithInitialAskTimeoutExOnCreateTx";
initDatastores(testName, MODULE_SHARDS_CARS_PEOPLE_1_2_3, CARS);
final DatastoreContext.Builder follower2DatastoreContextBuilder = DatastoreContext.newBuilder()
.shardHeartbeatIntervalInMillis(100).shardElectionTimeoutFactor(5);
final IntegrationTestKit follower2TestKit = new IntegrationTestKit(
- follower2System, follower2DatastoreContextBuilder);
+ follower2System, follower2DatastoreContextBuilder, commitTimeout);
try (final AbstractDataStore ds =
follower2TestKit.setupAbstractDataStore(
protected DatastoreContext.Builder datastoreContextBuilder;
protected DatastoreSnapshot restoreFromSnapshot;
+ private final int commitTimeout;
public IntegrationTestKit(final ActorSystem actorSystem, final Builder datastoreContextBuilder) {
+ this(actorSystem, datastoreContextBuilder, 7);
+ }
+
+ public IntegrationTestKit(final ActorSystem actorSystem, final Builder datastoreContextBuilder, int commitTimeout) {
super(actorSystem);
this.datastoreContextBuilder = datastoreContextBuilder;
+ this.commitTimeout = commitTimeout;
}
public DatastoreContext.Builder getDatastoreContextBuilder() {
}
public void doCommit(final DOMStoreThreePhaseCommitCohort cohort) throws Exception {
- Boolean canCommit = cohort.canCommit().get(7, TimeUnit.SECONDS);
+ Boolean canCommit = cohort.canCommit().get(commitTimeout, TimeUnit.SECONDS);
assertEquals("canCommit", true, canCommit);
cohort.preCommit().get(5, TimeUnit.SECONDS);
cohort.commit().get(5, TimeUnit.SECONDS);
void doCommit(final ListenableFuture<Boolean> canCommitFuture, final DOMStoreThreePhaseCommitCohort cohort)
throws Exception {
- Boolean canCommit = canCommitFuture.get(7, TimeUnit.SECONDS);
+ Boolean canCommit = canCommitFuture.get(commitTimeout, TimeUnit.SECONDS);
assertEquals("canCommit", true, canCommit);
cohort.preCommit().get(5, TimeUnit.SECONDS);
cohort.commit().get(5, TimeUnit.SECONDS);