Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -25,7 +25,6 @@
import java.time.Duration;
import java.util.Map;
import java.util.Map.Entry;
import java.util.Optional;
import java.util.SortedSet;
import java.util.TreeSet;

Expand All @@ -38,7 +37,6 @@
import org.apache.accumulo.core.data.Key;
import org.apache.accumulo.core.data.Value;
import org.apache.accumulo.core.lock.ServiceLock;
import org.apache.accumulo.core.lock.ServiceLockData;
import org.apache.accumulo.core.metadata.SystemTables;
import org.apache.accumulo.core.metadata.schema.MetadataSchema.DeletesSection;
import org.apache.accumulo.core.security.Authorizations;
Expand Down Expand Up @@ -98,28 +96,15 @@ public void alterConfig() throws Exception {
ClientContext context = (ClientContext) client) {
ZooCache zcache = context.getZooCache();
var path = context.getServerPaths().createGarbageCollectorPath();
Optional<ServiceLockData> gcLockData;
do {
gcLockData = ServiceLock.getLockData(zcache, path, null);
if (gcLockData.isPresent()) {
log.info("Waiting for GC ZooKeeper lock to expire");
Thread.sleep(2000);
}
} while (gcLockData.isPresent());

Wait.waitFor(() -> ServiceLock.getLockData(zcache, path, null).isEmpty(), 30_000, 250,
"GC ZooKeeper lock did not expire");
log.info("GC lock was lost");

getCluster().getClusterControl().startAllServers(ServerType.GARBAGE_COLLECTOR);
log.info("Garbage collector was restarted");

do {
gcLockData = ServiceLock.getLockData(zcache, path, null);
if (gcLockData.isEmpty()) {
log.info("Waiting for GC ZooKeeper lock to be acquired");
Thread.sleep(2000);
}
} while (gcLockData.isEmpty());

Wait.waitFor(() -> ServiceLock.getLockData(zcache, path, null).isPresent(), 30_000, 250,
"GC ZooKeeper lock was not acquired");
log.info("GC lock was acquired");
}
}
Expand Down
12 changes: 6 additions & 6 deletions test/src/main/java/org/apache/accumulo/test/ExistingMacIT.java
Original file line number Diff line number Diff line change
Expand Up @@ -51,6 +51,7 @@
import org.apache.accumulo.miniclusterImpl.ProcessReference;
import org.apache.accumulo.server.util.AccumuloStatus;
import org.apache.accumulo.test.functional.ConfigurableMacBase;
import org.apache.accumulo.test.util.Wait;
import org.apache.commons.io.FileUtils;
import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.fs.RawLocalFileSystem;
Expand Down Expand Up @@ -116,10 +117,8 @@ public void testExistingInstance() throws Exception {
}
}

while (!AccumuloStatus.isAccumuloOffline((ClientContext) client)) {
log.debug("Accumulo services still have their ZK locks held");
Thread.sleep(1000);
}
Wait.waitFor(() -> AccumuloStatus.isAccumuloOffline((ClientContext) client), 30_000, 250,
"Accumulo services still have ZooKeeper locks after processes were killed");

File hadoopConfDir = createTestDir(ExistingMacIT.class.getSimpleName() + "_hadoop_conf");
FileUtils.deleteQuietly(hadoopConfDir);
Expand All @@ -138,9 +137,10 @@ public void testExistingInstance() throws Exception {
MiniAccumuloClusterImpl accumulo2 = new MiniAccumuloClusterImpl(macConfig2);
accumulo2.start();

client = accumulo2.createAccumuloClient(rootUser, new PasswordToken(ROOT_PASSWORD));
AccumuloClient client2 =
accumulo2.createAccumuloClient(rootUser, new PasswordToken(ROOT_PASSWORD));

try (Scanner scanner = client.createScanner(table, Authorizations.EMPTY)) {
try (Scanner scanner = client2.createScanner(table, Authorizations.EMPTY)) {
int sum = 0;
for (Entry<Key,Value> entry : scanner) {
sum += Integer.parseInt(entry.getValue().toString());
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -98,13 +98,10 @@ public void testListRunningCompactions() throws Exception {
compact(client, tableName, 2, GROUP7, false);

// wait for the compaction to start
var expected =
Wait.waitFor(() -> !ExternalCompactionTestUtils
.getRunningCompactions(getCluster().getServerContext()).isEmpty(), 10_000, 250);
final var expected =
ExternalCompactionTestUtils.getRunningCompactions(getCluster().getServerContext());
while (expected.isEmpty()) {
Thread.sleep(1000);
expected =
ExternalCompactionTestUtils.getRunningCompactions(getCluster().getServerContext());
}

final List<RunningCompactionSummary> running = new ArrayList<>();
final Map<String,RunningCompactionSummary> compactionsByEcid = new HashMap<>();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -53,7 +53,6 @@
import org.apache.accumulo.core.lock.ServiceLockPaths;
import org.apache.accumulo.core.lock.ServiceLockPaths.ServiceLockPath;
import org.apache.accumulo.core.metadata.SystemTables;
import org.apache.accumulo.core.util.UtilWaitThread;
import org.apache.accumulo.manager.Manager;
import org.apache.accumulo.manager.tableOps.FateEnv;
import org.apache.accumulo.minicluster.ServerType;
Expand All @@ -62,6 +61,7 @@
import org.apache.accumulo.server.ServerContext;
import org.apache.accumulo.test.fate.FastFate;
import org.apache.accumulo.test.functional.ConfigurableMacBase;
import org.apache.accumulo.test.util.Wait;
import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.io.Text;
import org.junit.jupiter.api.Test;
Expand Down Expand Up @@ -261,15 +261,13 @@ private static void waitToSeeManagers(ClientContext context, int expectedManager
.stream().map(FateStore.FateReservation::getReservationUUID).collect(toSet());
log.debug("existingReservationUUIDs {}", existingReservationUUIDs);

// Wait for there to be the expected number of managers in zookeeper. After manager processes
// are killed these entries in zookeeper may persist for a bit.
Wait.waitFor(() -> context.getServerPaths()
.getAssistantManagers(ServiceLockPaths.AddressSelector.all(), true).size()
== expectedManagers, 120_000, 250, "Expected manager count was not reached");
var assistants =
context.getServerPaths().getAssistantManagers(ServiceLockPaths.AddressSelector.all(), true);
// Wait for there to be the expected number of managers in zookeeper. After manager processes
// are kill these entries in zookeeper may persist for a bit.
while (assistants.size() != expectedManagers) {
UtilWaitThread.sleep(1);
assistants = context.getServerPaths()
.getAssistantManagers(ServiceLockPaths.AddressSelector.all(), true);
}

var expectedServers = assistants.stream().map(ServiceLockPath::getServer)
.map(HostAndPort::fromString).collect(toSet());
Expand All @@ -285,7 +283,7 @@ private static void waitToSeeManagers(ClientContext context, int expectedManager
// reassigned.
Set<Character> seenPrefixes = new HashSet<>();

while (reservationsSeen.size() < expectedManagers || !seenPrefixes.equals(expectedPrefixes)) {
Wait.waitFor(() -> {
var reservations =
store.getActiveReservations(Set.of(FatePartition.all(FateInstanceType.USER)));
reservations.forEach((fateId, reservation) -> {
Expand All @@ -308,8 +306,8 @@ private static void waitToSeeManagers(ClientContext context, int expectedManager
}
}
});
UtilWaitThread.sleep(1);
}
return reservationsSeen.size() >= expectedManagers && seenPrefixes.equals(expectedPrefixes);
}, 120_000, 250, "Expected manager FATE reservations and UUID prefixes were not seen");

log.debug("managers seen in fate reservations :{}", reservationsSeen);
if (managersKilled) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -81,6 +81,7 @@
import org.apache.accumulo.core.util.tables.TableNameUtil;
import org.apache.accumulo.test.constraints.NumericValueConstraint;
import org.apache.accumulo.test.harness.SharedMiniClusterBase;
import org.apache.accumulo.test.util.Wait;
import org.apache.commons.lang3.StringUtils;
import org.apache.hadoop.io.Text;
import org.junit.jupiter.api.AfterAll;
Expand Down Expand Up @@ -560,21 +561,21 @@ public void verifyConstraintInheritance() throws Exception {

c.namespaceOperations().addConstraint(namespace, constraintClassName);
// loop until constraint is seen (or until test timeout)
while (!c.namespaceOperations().listConstraints(namespace).containsKey(constraintClassName)
|| !c.tableOperations().listConstraints(t1).containsKey(constraintClassName)) {
Thread.sleep(500);
}
Wait.waitFor(
() -> c.namespaceOperations().listConstraints(namespace).containsKey(constraintClassName)
&& c.tableOperations().listConstraints(t1).containsKey(constraintClassName),
30_000, 250, "Constraint was not propagated to the namespace and table");

Integer namespaceNum = null;
Integer tableNum = null;
Integer[] constraintNums = new Integer[2];
// loop until constraint is seen in namespace and table (or until test times out)
while (namespaceNum == null || tableNum == null) {
namespaceNum = c.namespaceOperations().listConstraints(namespace).get(constraintClassName);
tableNum = c.tableOperations().listConstraints(t1).get(constraintClassName);
if (namespaceNum == null || tableNum == null) {
Thread.sleep(500);
}
}
Wait.waitFor(() -> {
constraintNums[0] =
c.namespaceOperations().listConstraints(namespace).get(constraintClassName);
constraintNums[1] = c.tableOperations().listConstraints(t1).get(constraintClassName);
return constraintNums[0] != null && constraintNums[1] != null;
}, 30_000, 250, "Constraint IDs were not propagated to the namespace and table");
Integer namespaceNum = constraintNums[0];
Integer tableNum = constraintNums[1];
assertEquals(namespaceNum, tableNum);

Mutation m1 = new Mutation("r1");
Expand All @@ -585,41 +586,39 @@ public void verifyConstraintInheritance() throws Exception {
m3.put("c", "d", new Value("zyxwv"));

// loop until constraint is activated and rejects mutations (or until test timeout)
boolean mutationsRejected = false;
while (!mutationsRejected) {
Wait.waitFor(() -> {
BatchWriter bw = c.createBatchWriter(t1);
bw.addMutations(Arrays.asList(m1, m2, m3));
try {
bw.close();
Thread.sleep(500);
return false;
} catch (MutationsRejectedException e) {
mutationsRejected = true;
assertEquals(1, e.getConstraintViolationSummaries().size());
assertEquals(2, e.getConstraintViolationSummaries().get(0).getNumberOfViolatingMutations());
return true;
}
}
}, 30_000, 250, "Constraint did not become active and reject invalid mutations");

assertNotNull(namespaceNum, "Namespace constraint ID should not be null");
c.namespaceOperations().removeConstraint(namespace, namespaceNum);

// loop until constraint is removed from config (or until test timeout)
while (c.namespaceOperations().listConstraints(namespace).containsKey(constraintClassName)
|| c.tableOperations().listConstraints(t1).containsKey(constraintClassName)) {
Thread.sleep(500);
}
Wait.waitFor(
() -> !c.namespaceOperations().listConstraints(namespace).containsKey(constraintClassName)
&& !c.tableOperations().listConstraints(t1).containsKey(constraintClassName),
30_000, 250, "Constraint was not removed from the namespace and table");

// loop until constraint is removed and stops rejecting (or until test timeout)
boolean mutationsAccepted = false;
while (!mutationsAccepted) {
Wait.waitFor(() -> {
BatchWriter bw = c.createBatchWriter(t1);
try {
bw.addMutations(Arrays.asList(m1, m2, m3));
bw.close();
mutationsAccepted = true;
return true;
} catch (MutationsRejectedException e) {
Thread.sleep(500);
return false;
}
}
}, 30_000, 250, "Mutations continued to be rejected after constraint removal");
}

@Test
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -162,10 +162,9 @@ public void testScanServerMetadataEntries() throws Exception {

}

// close happens asynchronously. Let the test fail by timeout
while (ctx.getAmple().scanServerRefs().list().findAny().isPresent()) {
Thread.sleep(1000);
}
// close happens asynchronously; poll until the references are removed
Wait.waitFor(() -> ctx.getAmple().scanServerRefs().list().findAny().isEmpty(), 30_000, 250,
"Scan server references were not removed after scanner close");
}
}

Expand Down Expand Up @@ -196,10 +195,9 @@ public void testBatchScanServerMetadataEntries() throws Exception {

}

// close happens asynchronously. Let the test fail by timeout
while (ctx.getAmple().scanServerRefs().list().findAny().isPresent()) {
Thread.sleep(1000);
}
// close happens asynchronously; poll until the references are removed
Wait.waitFor(() -> ctx.getAmple().scanServerRefs().list().findAny().isEmpty(), 30_000, 250,
"Scan server references were not removed after scanner close");
}
}

Expand Down Expand Up @@ -265,10 +263,9 @@ public void testGcRunScanServerReferences() throws Exception {
assertEquals(fileCount, deduplicatedReferences.size());
client.tableOperations().delete(tableName);
}
// close happens asynchronously. Let the test fail by timeout
while (ctx.getAmple().scanServerRefs().list().findAny().isPresent()) {
Thread.sleep(1000);
}
// close happens asynchronously; poll until the references are removed
Wait.waitFor(() -> ctx.getAmple().scanServerRefs().list().findAny().isEmpty(), 30_000, 250,
"Scan server references were not removed after table deletion");

}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,8 +18,6 @@
*/
package org.apache.accumulo.test;

import static java.util.concurrent.TimeUnit.SECONDS;

import java.time.Duration;
import java.util.concurrent.atomic.AtomicReference;

Expand Down Expand Up @@ -90,9 +88,9 @@ public void test() throws Exception {
});
backgroundWriter.start();
// wait for the tserver to give up on writing to the WAL
while (client.instanceOperations().getServers(ServerId.Type.TABLET_SERVER).size() == 1) {
Thread.sleep(SECONDS.toMillis(1));
}
Wait.waitFor(
() -> client.instanceOperations().getServers(ServerId.Type.TABLET_SERVER).size() != 1,
30_000, 250, "Tablet server did not leave the live list");
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,7 @@
import org.apache.accumulo.core.security.Authorizations;
import org.apache.accumulo.miniclusterImpl.MiniAccumuloConfigImpl;
import org.apache.accumulo.test.functional.ConfigurableMacBase;
import org.apache.accumulo.test.util.Wait;
import org.apache.hadoop.conf.Configuration;
import org.junit.jupiter.api.Test;

Expand All @@ -56,9 +57,9 @@ public void configure(MiniAccumuloConfigImpl cfg, Configuration hadoopCoreSite)
public void test() throws Exception {
try (AccumuloClient client = Accumulo.newClient().from(getClientProperties()).build()) {
// wait until a tablet server is up
while (client.instanceOperations().getServers(ServerId.Type.TABLET_SERVER).isEmpty()) {
Thread.sleep(50);
}
Wait.waitFor(
() -> !client.instanceOperations().getServers(ServerId.Type.TABLET_SERVER).isEmpty(),
30_000, 250, "Tablet server did not start");
final String tableName = getUniqueNames(1)[0];
client.tableOperations().create(tableName);
try (BatchWriter bw = client.createBatchWriter(tableName)) {
Expand Down
Loading
Loading