From 268d40d675010bc0b877880e0210635bdad5ca15 Mon Sep 17 00:00:00 2001 From: Christopher Tubbs Date: Thu, 1 Oct 2026 19:30:52 -0400 Subject: [PATCH 1/3] Improve integration test wait efficiency Replace fixed-sleep polling loops with condition-based Wait.waitFor calls across several ITs. Co-authored-by: Dom G. Co-authored-by: OpenAI GPT-5.6 Luna --- .../test/BadDeleteMarkersCreatedIT.java | 23 ++----- .../apache/accumulo/test/ExistingMacIT.java | 12 ++-- .../accumulo/test/ListCompactionsIT.java | 9 +-- .../test/ScanServerMetadataEntriesIT.java | 21 +++---- .../test/TabletServerHdfsRestartIT.java | 7 ++- .../ZooKeeperPropertiesIT_SimpleSuite.java | 60 ++++++------------- .../test/compaction/CompactionExecutorIT.java | 16 +++-- .../compaction/ExternalCompaction2ITBase.java | 24 ++++---- .../test/fate/FateExecutionOrderITBase.java | 6 +- .../apache/accumulo/test/fate/FateITBase.java | 7 +-- .../test/functional/BatchScanSplitIT.java | 13 ++-- .../test/functional/ConstraintIT.java | 21 +++---- .../test/functional/LastLocationIT.java | 27 +++++---- .../test/functional/ManyWriteAheadLogsIT.java | 9 +-- .../accumulo/test/functional/MetadataIT.java | 6 +- .../test/functional/MonitorSslIT.java | 18 +++--- .../accumulo/test/functional/ReadWriteIT.java | 32 ++++------ .../accumulo/test/functional/RestartIT.java | 33 +++------- .../accumulo/test/functional/ScanIdIT.java | 7 +-- .../test/functional/TabletMetadataIT.java | 19 ++---- .../accumulo/test/lock/ServiceLockIT.java | 10 ++-- .../accumulo/test/metrics/MetricsIT.java | 6 +- .../org/apache/accumulo/test/rpc/Mocket.java | 1 - .../accumulo/test/tracing/ScanTracingIT.java | 1 - 24 files changed, 143 insertions(+), 245 deletions(-) diff --git a/test/src/main/java/org/apache/accumulo/test/BadDeleteMarkersCreatedIT.java b/test/src/main/java/org/apache/accumulo/test/BadDeleteMarkersCreatedIT.java index 84bf25dbe13..2324b611772 100644 --- a/test/src/main/java/org/apache/accumulo/test/BadDeleteMarkersCreatedIT.java +++ b/test/src/main/java/org/apache/accumulo/test/BadDeleteMarkersCreatedIT.java @@ -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; @@ -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; @@ -98,28 +96,15 @@ public void alterConfig() throws Exception { ClientContext context = (ClientContext) client) { ZooCache zcache = context.getZooCache(); var path = context.getServerPaths().createGarbageCollectorPath(); - Optional 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"); } } diff --git a/test/src/main/java/org/apache/accumulo/test/ExistingMacIT.java b/test/src/main/java/org/apache/accumulo/test/ExistingMacIT.java index 9f70102d447..c7f4ec6c34e 100644 --- a/test/src/main/java/org/apache/accumulo/test/ExistingMacIT.java +++ b/test/src/main/java/org/apache/accumulo/test/ExistingMacIT.java @@ -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; @@ -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); @@ -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 entry : scanner) { sum += Integer.parseInt(entry.getValue().toString()); diff --git a/test/src/main/java/org/apache/accumulo/test/ListCompactionsIT.java b/test/src/main/java/org/apache/accumulo/test/ListCompactionsIT.java index be14de80a36..1101f82701d 100644 --- a/test/src/main/java/org/apache/accumulo/test/ListCompactionsIT.java +++ b/test/src/main/java/org/apache/accumulo/test/ListCompactionsIT.java @@ -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 running = new ArrayList<>(); final Map compactionsByEcid = new HashMap<>(); diff --git a/test/src/main/java/org/apache/accumulo/test/ScanServerMetadataEntriesIT.java b/test/src/main/java/org/apache/accumulo/test/ScanServerMetadataEntriesIT.java index 4bf1aa54ab1..5612f95d9bf 100644 --- a/test/src/main/java/org/apache/accumulo/test/ScanServerMetadataEntriesIT.java +++ b/test/src/main/java/org/apache/accumulo/test/ScanServerMetadataEntriesIT.java @@ -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"); } } @@ -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"); } } @@ -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"); } diff --git a/test/src/main/java/org/apache/accumulo/test/TabletServerHdfsRestartIT.java b/test/src/main/java/org/apache/accumulo/test/TabletServerHdfsRestartIT.java index de7fccefc51..90cbd0c2620 100644 --- a/test/src/main/java/org/apache/accumulo/test/TabletServerHdfsRestartIT.java +++ b/test/src/main/java/org/apache/accumulo/test/TabletServerHdfsRestartIT.java @@ -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; @@ -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)) { diff --git a/test/src/main/java/org/apache/accumulo/test/ZooKeeperPropertiesIT_SimpleSuite.java b/test/src/main/java/org/apache/accumulo/test/ZooKeeperPropertiesIT_SimpleSuite.java index 5fa83c094c3..a7a7f39fdb0 100644 --- a/test/src/main/java/org/apache/accumulo/test/ZooKeeperPropertiesIT_SimpleSuite.java +++ b/test/src/main/java/org/apache/accumulo/test/ZooKeeperPropertiesIT_SimpleSuite.java @@ -20,7 +20,6 @@ import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertThrows; -import static org.junit.jupiter.api.Assertions.fail; import java.util.List; import java.util.Map; @@ -41,6 +40,7 @@ import org.apache.accumulo.server.conf.store.TablePropKey; import org.apache.accumulo.server.util.PropUtil; import org.apache.accumulo.test.harness.SharedMiniClusterBase; +import org.apache.accumulo.test.util.Wait; import org.junit.jupiter.api.AfterAll; import org.junit.jupiter.api.BeforeAll; import org.junit.jupiter.api.Test; @@ -87,31 +87,18 @@ public void testTablePropUtils() throws AccumuloException, TableExistsException, PropUtil.setProperties(context, tablePropKey, Map.of(Property.TABLE_BLOOM_ENABLED.getKey(), "true")); - // add a sleep to give the property change time to propagate - properties = client.tableOperations().getConfiguration(tableName); - while (properties.get(Property.TABLE_BLOOM_ENABLED.getKey()).equals("false")) { - try { - Thread.sleep(250); - } catch (InterruptedException e) { - Thread.currentThread().interrupt(); - fail("Thread interrupted while waiting for tablePropUtil update"); - } - properties = client.tableOperations().getConfiguration(tableName); - } + Wait.waitFor( + () -> !client.tableOperations().getConfiguration(tableName) + .get(Property.TABLE_BLOOM_ENABLED.getKey()).equals("false"), + 30_000, 250, "Table property update did not propagate"); PropUtil.removeProperties(context, tablePropKey, List.of(Property.TABLE_BLOOM_ENABLED.getKey())); - properties = client.tableOperations().getConfiguration(tableName); - while (properties.get(Property.TABLE_BLOOM_ENABLED.getKey()).equals("true")) { - try { - Thread.sleep(250); - } catch (InterruptedException e) { - Thread.currentThread().interrupt(); - fail("Thread interrupted while waiting for tablePropUtil update"); - } - properties = client.tableOperations().getConfiguration(tableName); - } + Wait.waitFor( + () -> !client.tableOperations().getConfiguration(tableName) + .get(Property.TABLE_BLOOM_ENABLED.getKey()).equals("true"), + 30_000, 250, "Table property removal did not propagate"); // Add invalid property assertThrows(IllegalArgumentException.class, @@ -142,31 +129,18 @@ public void testNamespacePropUtils() throws AccumuloException, AccumuloSecurityE PropUtil.setProperties(context, namespacePropKey, Map.of(Property.TABLE_FILE_MAX.getKey(), "31")); - // add a sleep to give the property change time to propagate - properties = client.namespaceOperations().getConfiguration(namespace); - while (!properties.get(Property.TABLE_FILE_MAX.getKey()).equals("31")) { - try { - Thread.sleep(250); - } catch (InterruptedException e) { - Thread.currentThread().interrupt(); - fail("Thread interrupted while waiting for namespacePropUtil update"); - } - properties = client.namespaceOperations().getConfiguration(namespace); - } + Wait.waitFor( + () -> client.namespaceOperations().getConfiguration(namespace) + .get(Property.TABLE_FILE_MAX.getKey()).equals("31"), + 30_000, 250, "Namespace property update did not propagate"); PropUtil.removeProperties(context, namespacePropKey, List.of(Property.TABLE_FILE_MAX.getKey())); - properties = client.namespaceOperations().getConfiguration(namespace); - while (!properties.get(Property.TABLE_FILE_MAX.getKey()).equals("15")) { - try { - Thread.sleep(250); - } catch (InterruptedException e) { - Thread.currentThread().interrupt(); - fail("Thread interrupted while waiting for namespacePropUtil update"); - } - properties = client.namespaceOperations().getConfiguration(namespace); - } + Wait.waitFor( + () -> client.namespaceOperations().getConfiguration(namespace) + .get(Property.TABLE_FILE_MAX.getKey()).equals("15"), + 30_000, 250, "Namespace property removal did not propagate"); // Add invalid property assertThrows(IllegalArgumentException.class, diff --git a/test/src/main/java/org/apache/accumulo/test/compaction/CompactionExecutorIT.java b/test/src/main/java/org/apache/accumulo/test/compaction/CompactionExecutorIT.java index fc64cee8995..404e5e315bd 100644 --- a/test/src/main/java/org/apache/accumulo/test/compaction/CompactionExecutorIT.java +++ b/test/src/main/java/org/apache/accumulo/test/compaction/CompactionExecutorIT.java @@ -73,6 +73,7 @@ import org.apache.accumulo.miniclusterImpl.MiniAccumuloConfigImpl; import org.apache.accumulo.test.harness.MiniClusterConfigurationCallback; import org.apache.accumulo.test.harness.SharedMiniClusterBase; +import org.apache.accumulo.test.util.Wait; import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.fs.FileSystem; import org.apache.hadoop.fs.RawLocalFileSystem; @@ -307,9 +308,8 @@ public void testReconfigureCompactionService() throws Exception { addFiles(client, "rctt", 22); - while (getFiles(client, "rctt").size() > 2) { - Thread.sleep(100); - } + Wait.waitFor(() -> getFiles(client, "rctt").size() <= 2, 30_000, 100, + "Compaction did not reduce rctt to two files"); assertEquals(2, getFiles(client, "rctt").size()); @@ -326,9 +326,8 @@ public void testReconfigureCompactionService() throws Exception { addFiles(client, "rctt", 10); - while (getFiles(client, "rctt").size() > 4) { - Thread.sleep(100); - } + Wait.waitFor(() -> getFiles(client, "rctt").size() <= 4, 30_000, 100, + "Compaction did not reduce rctt to four files"); assertEquals(4, getFiles(client, "rctt").size()); } @@ -359,9 +358,8 @@ public void testAddCompactionService() throws Exception { addFiles(client, "acst", 42); - while (getFiles(client, "acst").size() > 6) { - Thread.sleep(100); - } + Wait.waitFor(() -> getFiles(client, "acst").size() <= 6, 30_000, 100, + "Compaction did not reduce acst to six files"); assertEquals(6, getFiles(client, "acst").size()); } diff --git a/test/src/main/java/org/apache/accumulo/test/compaction/ExternalCompaction2ITBase.java b/test/src/main/java/org/apache/accumulo/test/compaction/ExternalCompaction2ITBase.java index 64f89cf6c38..80d0621926e 100644 --- a/test/src/main/java/org/apache/accumulo/test/compaction/ExternalCompaction2ITBase.java +++ b/test/src/main/java/org/apache/accumulo/test/compaction/ExternalCompaction2ITBase.java @@ -178,10 +178,10 @@ public void testUserCompactionCancellation() throws Exception { // when the compaction starts it will create a selected files column in the tablet, wait for // that to happen - while (countTablets(getCluster().getServerContext(), table1, - tm -> tm.getSelectedFiles() != null) == 0) { - Thread.sleep(1000); - } + Wait.waitFor( + () -> countTablets(getCluster().getServerContext(), table1, + tm -> tm.getSelectedFiles() != null) > 0, + 30_000, 250, "Compaction did not select files"); client.tableOperations().cancelCompaction(table1); @@ -193,10 +193,10 @@ public void testUserCompactionCancellation() throws Exception { confirmCompactionsNoLongerRunning(getCluster().getServerContext(), ecids); // ensure the canceled compaction deletes any tablet metadata related to the compaction - while (countTablets(getCluster().getServerContext(), table1, - tm -> tm.getSelectedFiles() != null || !tm.getCompacted().isEmpty()) > 0) { - Thread.sleep(1000); - } + Wait.waitFor( + () -> countTablets(getCluster().getServerContext(), table1, + tm -> tm.getSelectedFiles() != null || !tm.getCompacted().isEmpty()) == 0, + 30_000, 250, "Canceled compaction metadata was not removed"); // Verify that the tmp file are cleaned up Wait.waitFor(() -> FindCompactionTmpFiles @@ -226,10 +226,10 @@ public void testDeleteTableCancelsUserExternalCompaction() throws Exception { // when the compaction starts it will create a selected files column in the tablet, wait for // that to happen - while (countTablets(getCluster().getServerContext(), table1, - tm -> tm.getSelectedFiles() != null) == 0) { - Thread.sleep(1000); - } + Wait.waitFor( + () -> countTablets(getCluster().getServerContext(), table1, + tm -> tm.getSelectedFiles() != null) > 0, + 30_000, 250, "Compaction did not select files"); client.tableOperations().delete(table1); diff --git a/test/src/main/java/org/apache/accumulo/test/fate/FateExecutionOrderITBase.java b/test/src/main/java/org/apache/accumulo/test/fate/FateExecutionOrderITBase.java index 60badd3025a..b10b9636cc3 100644 --- a/test/src/main/java/org/apache/accumulo/test/fate/FateExecutionOrderITBase.java +++ b/test/src/main/java/org/apache/accumulo/test/fate/FateExecutionOrderITBase.java @@ -59,6 +59,7 @@ import org.apache.accumulo.core.fate.Repo; import org.apache.accumulo.server.ServerContext; import org.apache.accumulo.test.harness.SharedMiniClusterBase; +import org.apache.accumulo.test.util.Wait; import org.apache.hadoop.io.Text; import org.junit.jupiter.api.AfterAll; import org.junit.jupiter.api.BeforeAll; @@ -194,9 +195,8 @@ public void before() throws Exception { } private void waitFor(FateStore store, FateId txid) throws Exception { - while (store.read(txid).getStatus() != SUCCESSFUL) { - Thread.sleep(50); - } + Wait.waitFor(() -> store.read(txid).getStatus() == SUCCESSFUL, 30_000, 50, + "FATE transaction did not complete successfully"); } protected Fate initializeFate(AccumuloClient client, FateStore store) { diff --git a/test/src/main/java/org/apache/accumulo/test/fate/FateITBase.java b/test/src/main/java/org/apache/accumulo/test/fate/FateITBase.java index c044b9dd879..b42533a3d35 100644 --- a/test/src/main/java/org/apache/accumulo/test/fate/FateITBase.java +++ b/test/src/main/java/org/apache/accumulo/test/fate/FateITBase.java @@ -652,11 +652,8 @@ protected void testShutdownDoesNotFailTx(FateStore store, ServerContext finishCall.countDown(); // This should complete normally, cleaning up the tx and deleting it from ZK - TStatus status = getTxStatus(sctx, txid); - while (status != TStatus.UNKNOWN) { - Thread.sleep(100); - status = getTxStatus(sctx, txid); - } + Wait.waitFor(() -> getTxStatus(sctx, txid) == TStatus.UNKNOWN, 30_000, 100, + "FATE transaction was not cleaned up"); assertNull(interruptedException.get()); } finally { fate.shutdown(10, TimeUnit.MINUTES); diff --git a/test/src/main/java/org/apache/accumulo/test/functional/BatchScanSplitIT.java b/test/src/main/java/org/apache/accumulo/test/functional/BatchScanSplitIT.java index ebe2a23175e..a3d0840626f 100644 --- a/test/src/main/java/org/apache/accumulo/test/functional/BatchScanSplitIT.java +++ b/test/src/main/java/org/apache/accumulo/test/functional/BatchScanSplitIT.java @@ -22,7 +22,6 @@ import java.time.Duration; import java.util.ArrayList; -import java.util.Collection; import java.util.HashMap; import java.util.Map.Entry; @@ -36,6 +35,7 @@ import org.apache.accumulo.core.data.Range; import org.apache.accumulo.core.data.Value; import org.apache.accumulo.test.harness.AccumuloClusterHarness; +import org.apache.accumulo.test.util.Wait; import org.apache.hadoop.io.Text; import org.junit.jupiter.api.Test; import org.slf4j.Logger; @@ -69,13 +69,8 @@ public void test() throws Exception { c.tableOperations().setProperty(tableName, Property.TABLE_SPLIT_THRESHOLD.getKey(), "4K"); - Collection splits = c.tableOperations().listSplits(tableName); - while (splits.size() < 2) { - Thread.sleep(1); - splits = c.tableOperations().listSplits(tableName); - } - - System.out.println("splits : " + splits); + Wait.waitFor(() -> c.tableOperations().listSplits(tableName).size() >= 2, 30_000, 100, + "Expected table splits were not created"); HashMap expected = new HashMap<>(); ArrayList ranges = new ArrayList<>(); @@ -113,7 +108,7 @@ public void test() throws Exception { } } - splits = c.tableOperations().listSplits(tableName); + var splits = c.tableOperations().listSplits(tableName); log.info("splits : {}", splits); } } diff --git a/test/src/main/java/org/apache/accumulo/test/functional/ConstraintIT.java b/test/src/main/java/org/apache/accumulo/test/functional/ConstraintIT.java index c4811605111..b84f31fb3b3 100644 --- a/test/src/main/java/org/apache/accumulo/test/functional/ConstraintIT.java +++ b/test/src/main/java/org/apache/accumulo/test/functional/ConstraintIT.java @@ -43,6 +43,7 @@ import org.apache.accumulo.test.constraints.AlphaNumKeyConstraint; import org.apache.accumulo.test.constraints.NumericValueConstraint; import org.apache.accumulo.test.harness.AccumuloClusterHarness; +import org.apache.accumulo.test.util.Wait; import org.apache.hadoop.io.Text; import org.junit.jupiter.api.Test; import org.slf4j.Logger; @@ -66,20 +67,14 @@ public void run() throws Exception { c.tableOperations().addConstraint(table, AlphaNumKeyConstraint.class.getName()); } - // A static sleep to just let ZK do its thing - Thread.sleep(10_000); - - // Then check that the client has at least gotten the updates + // Wait for both constraints to become visible for (String table : tableNames) { - log.debug("Checking constraints on {}", table); - Map constraints = c.tableOperations().listConstraints(table); - while (!constraints.containsKey(NumericValueConstraint.class.getName()) - || !constraints.containsKey(AlphaNumKeyConstraint.class.getName())) { - log.debug("Failed to verify constraints. Sleeping and retrying"); - Thread.sleep(2000); - constraints = c.tableOperations().listConstraints(table); - } - log.debug("Verified all constraints on {}", table); + Wait.waitFor(() -> { + log.debug("Checking constraints on {}", table); + Map constraints = c.tableOperations().listConstraints(table); + return constraints.containsKey(NumericValueConstraint.class.getName()) + && constraints.containsKey(AlphaNumKeyConstraint.class.getName()); + }, 30_000, 250, "Constraints did not propagate for " + table); } log.debug("Verified constraints on all tables. Running tests"); diff --git a/test/src/main/java/org/apache/accumulo/test/functional/LastLocationIT.java b/test/src/main/java/org/apache/accumulo/test/functional/LastLocationIT.java index 27aab35177b..cec37017b3b 100644 --- a/test/src/main/java/org/apache/accumulo/test/functional/LastLocationIT.java +++ b/test/src/main/java/org/apache/accumulo/test/functional/LastLocationIT.java @@ -33,7 +33,7 @@ import org.apache.accumulo.core.client.security.tokens.PasswordToken; import org.apache.accumulo.core.data.Mutation; import org.apache.accumulo.core.metadata.schema.TabletMetadata; -import org.apache.accumulo.core.util.UtilWaitThread; +import org.apache.accumulo.test.util.Wait; import org.junit.jupiter.api.Test; public class LastLocationIT extends ConfigurableMacBase { @@ -53,13 +53,15 @@ public void test() throws Exception { c.tableOperations().create(tableName, ntc); String tableId = c.tableOperations().tableIdMap().get(tableName); // wait for the table to be online - TabletMetadata newTablet; - do { - UtilWaitThread.sleep(250); - newTablet = ManagerAssignmentIT.getTabletMetadata(c, tableId, null); - } while (!newTablet.hasCurrent()); + TabletMetadata[] newTablet = {null}; + Wait.waitFor(() -> { + newTablet[0] = ManagerAssignmentIT.getTabletMetadata(c, tableId, null); + return newTablet[0].hasCurrent(); + }, 30_000, 250, "Tablet was not assigned"); + TabletMetadata assignedTablet = newTablet[0]; // this would be null if the mode was not "assign" - assertEquals(newTablet.getLocation().getHostPort(), newTablet.getLast().getHostPort()); + assertEquals(assignedTablet.getLocation().getHostPort(), + assignedTablet.getLast().getHostPort()); // put something in it try (BatchWriter bw = c.createBatchWriter(tableName)) { @@ -70,23 +72,24 @@ public void test() throws Exception { // last location should not be set yet TabletMetadata unflushed = ManagerAssignmentIT.getTabletMetadata(c, tableId, null); - assertEquals(newTablet.getLocation().getHostPort(), unflushed.getLocation().getHostPort()); - assertEquals(newTablet.getLocation().getHostPort(), unflushed.getLast().getHostPort()); - assertTrue(newTablet.hasCurrent()); + assertEquals(assignedTablet.getLocation().getHostPort(), + unflushed.getLocation().getHostPort()); + assertEquals(assignedTablet.getLocation().getHostPort(), unflushed.getLast().getHostPort()); + assertTrue(assignedTablet.hasCurrent()); // take the tablet offline c.tableOperations().offline(tableName, true); TabletMetadata offline = ManagerAssignmentIT.getTabletMetadata(c, tableId, null); assertNull(offline.getLocation()); assertFalse(offline.hasCurrent()); - assertEquals(newTablet.getLocation().getHostPort(), offline.getLast().getHostPort()); + assertEquals(assignedTablet.getLocation().getHostPort(), offline.getLast().getHostPort()); // put it back online, should have the same last location c.tableOperations().online(tableName, true); TabletMetadata online = ManagerAssignmentIT.getTabletMetadata(c, tableId, null); assertTrue(online.hasCurrent()); assertNotNull(online.getLocation()); - assertEquals(newTablet.getLast().getHostPort(), online.getLast().getHostPort()); + assertEquals(assignedTablet.getLast().getHostPort(), online.getLast().getHostPort()); } } diff --git a/test/src/main/java/org/apache/accumulo/test/functional/ManyWriteAheadLogsIT.java b/test/src/main/java/org/apache/accumulo/test/functional/ManyWriteAheadLogsIT.java index 06841fa7d09..6c2ad5193f4 100644 --- a/test/src/main/java/org/apache/accumulo/test/functional/ManyWriteAheadLogsIT.java +++ b/test/src/main/java/org/apache/accumulo/test/functional/ManyWriteAheadLogsIT.java @@ -41,6 +41,7 @@ import org.apache.accumulo.server.ServerContext; import org.apache.accumulo.server.log.WalStateManager.WalState; import org.apache.accumulo.test.harness.AccumuloClusterHarness; +import org.apache.accumulo.test.util.Wait; import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.fs.RawLocalFileSystem; import org.apache.hadoop.io.Text; @@ -162,12 +163,8 @@ public void testMany() throws Exception { "Number of WALs seen was less than expected " + allWalsSeen.size()); // the total number of closed write ahead logs should get small - int closedLogs = countClosedWals(context); - while (closedLogs > 3) { - log.debug("Waiting for wals to shrink " + closedLogs); - Thread.sleep(250); - closedLogs = countClosedWals(context); - } + Wait.waitFor(() -> countClosedWals(context) <= 3, 30_000, 250, + "Closed WAL count did not decrease"); } } diff --git a/test/src/main/java/org/apache/accumulo/test/functional/MetadataIT.java b/test/src/main/java/org/apache/accumulo/test/functional/MetadataIT.java index d51d1cf6db6..ee6dee2e4fa 100644 --- a/test/src/main/java/org/apache/accumulo/test/functional/MetadataIT.java +++ b/test/src/main/java/org/apache/accumulo/test/functional/MetadataIT.java @@ -61,6 +61,7 @@ import org.apache.accumulo.core.security.TablePermission; import org.apache.accumulo.core.util.time.SteadyTime; import org.apache.accumulo.test.harness.SharedMiniClusterBase; +import org.apache.accumulo.test.util.Wait; import org.apache.hadoop.io.Text; import org.junit.jupiter.api.AfterAll; import org.junit.jupiter.api.BeforeAll; @@ -142,9 +143,8 @@ public void mergeMeta() throws Exception { c.tableOperations().merge(SystemTables.METADATA.tableName(), null, null); try (Scanner s = c.createScanner(SystemTables.ROOT.tableName(), Authorizations.EMPTY)) { s.setRange(DeletesSection.getRange()); - while (s.stream().findAny().isEmpty()) { - Thread.sleep(100); - } + Wait.waitFor(() -> s.stream().findAny().isPresent(), 30_000, 100, + "Metadata delete marker did not appear"); assertEquals(0, c.tableOperations().listSplits(SystemTables.METADATA.tableName()).size()); } } diff --git a/test/src/main/java/org/apache/accumulo/test/functional/MonitorSslIT.java b/test/src/main/java/org/apache/accumulo/test/functional/MonitorSslIT.java index 2d0d319a576..c215b2fbfdd 100644 --- a/test/src/main/java/org/apache/accumulo/test/functional/MonitorSslIT.java +++ b/test/src/main/java/org/apache/accumulo/test/functional/MonitorSslIT.java @@ -44,6 +44,7 @@ import org.apache.accumulo.core.util.MonitorUtil; import org.apache.accumulo.minicluster.ServerType; import org.apache.accumulo.miniclusterImpl.MiniAccumuloConfigImpl; +import org.apache.accumulo.test.util.Wait; import org.apache.hadoop.conf.Configuration; import org.junit.jupiter.api.BeforeAll; import org.junit.jupiter.api.Test; @@ -129,21 +130,18 @@ public void configure(MiniAccumuloConfigImpl cfg, Configuration hadoopCoreSite) public void test() throws Exception { log.debug("Starting Monitor"); cluster.getClusterControl().startAllServers(ServerType.MONITOR); - String monitorLocation = null; + String[] monitorLocation = {null}; try (AccumuloClient client = Accumulo.newClient().from(getClientProperties()).build()) { - while (monitorLocation == null) { + Wait.waitFor(() -> { try { - monitorLocation = MonitorUtil.getLocation((ClientContext) client); + monitorLocation[0] = MonitorUtil.getLocation((ClientContext) client); } catch (Exception e) { - // ignored + // The monitor may not have published its location yet. } - if (monitorLocation == null) { - log.debug("Could not fetch monitor HTTP address from zookeeper"); - Thread.sleep(2000); - } - } + return monitorLocation[0] != null; + }, 30_000, 250, "Monitor location was not published to ZooKeeper"); } - var url = new URI(monitorLocation).toURL(); + var url = new URI(monitorLocation[0]).toURL(); log.debug("Fetching web page {}", url); String result = FunctionalTestUtils.readWebPage(url).body(); assertTrue(result.length() > 100); diff --git a/test/src/main/java/org/apache/accumulo/test/functional/ReadWriteIT.java b/test/src/main/java/org/apache/accumulo/test/functional/ReadWriteIT.java index a4d4e1f5cb8..cd0cd6745f8 100644 --- a/test/src/main/java/org/apache/accumulo/test/functional/ReadWriteIT.java +++ b/test/src/main/java/org/apache/accumulo/test/functional/ReadWriteIT.java @@ -27,7 +27,6 @@ import java.net.URI; import java.security.cert.X509Certificate; import java.time.Duration; -import java.util.Optional; import java.util.concurrent.atomic.AtomicBoolean; import javax.net.ssl.HostnameVerifier; @@ -47,7 +46,6 @@ import org.apache.accumulo.core.data.Mutation; 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.util.MonitorUtil; import org.apache.accumulo.minicluster.ServerType; import org.apache.accumulo.miniclusterImpl.MiniAccumuloConfigImpl; @@ -56,6 +54,7 @@ import org.apache.accumulo.test.VerifyIngest; import org.apache.accumulo.test.VerifyIngest.VerifyParams; import org.apache.accumulo.test.harness.AccumuloClusterHarness; +import org.apache.accumulo.test.util.Wait; import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.io.Text; import org.junit.jupiter.api.Tag; @@ -97,14 +96,12 @@ public void sunnyDay() throws Exception { String tableName = getUniqueNames(1)[0]; ingest(accumuloClient, ROWS, COLS, 50, 0, tableName); verify(accumuloClient, ROWS, COLS, 50, 0, tableName); - String monitorLocation = null; - while (monitorLocation == null) { - monitorLocation = MonitorUtil.getLocation((ClientContext) accumuloClient); - if (monitorLocation == null) { - log.debug("Could not fetch monitor HTTP address from zookeeper"); - Thread.sleep(2000); - } - } + String[] monitorLocation = {null}; + Wait.waitFor(() -> { + monitorLocation[0] = MonitorUtil.getLocation((ClientContext) accumuloClient); + return monitorLocation[0] != null; + }, 30_000, 250, "Monitor location was not published to ZooKeeper"); + String monitorAddress = monitorLocation[0]; if (getCluster() instanceof StandaloneAccumuloCluster) { String monitorSslKeystore = getCluster().getSiteConfiguration().get(Property.MONITOR_SSL_KEYSTORE); @@ -119,22 +116,17 @@ public void sunnyDay() throws Exception { HttpsURLConnection.setDefaultHostnameVerifier(new TestHostnameVerifier()); } } - var url = new URI(monitorLocation).toURL(); + var url = new URI(monitorAddress).toURL(); log.debug("Fetching web page {}", url); String result = FunctionalTestUtils.readWebPage(url).body(); assertTrue(result.length() > 100); log.debug("Stopping accumulo cluster"); ClusterControl control = cluster.getClusterControl(); control.adminStopAll(); - Optional managerLockData; - do { - managerLockData = ServiceLock.getLockData(cluster.getServerContext().getZooCache(), - getServerContext().getServerPaths().createManagerPath(), null); - if (managerLockData.isPresent()) { - log.info("Manager lock is still held"); - Thread.sleep(1000); - } - } while (managerLockData.isPresent()); + Wait.waitFor( + () -> ServiceLock.getLockData(cluster.getServerContext().getZooCache(), + getServerContext().getServerPaths().createManagerPath(), null).isEmpty(), + 30_000, 250, "Manager lock was not released during shutdown"); control.stopAllServers(ServerType.MANAGER); control.stopAllServers(ServerType.TABLET_SERVER); control.stopAllServers(ServerType.GARBAGE_COLLECTOR); diff --git a/test/src/main/java/org/apache/accumulo/test/functional/RestartIT.java b/test/src/main/java/org/apache/accumulo/test/functional/RestartIT.java index 0360aa430a2..72c338cd5d7 100644 --- a/test/src/main/java/org/apache/accumulo/test/functional/RestartIT.java +++ b/test/src/main/java/org/apache/accumulo/test/functional/RestartIT.java @@ -24,7 +24,6 @@ import java.io.IOException; import java.util.Map.Entry; -import java.util.Optional; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.Future; @@ -36,7 +35,6 @@ import org.apache.accumulo.core.conf.ClientProperty; import org.apache.accumulo.core.conf.Property; 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.zookeeper.ZooCache; import org.apache.accumulo.minicluster.ServerType; @@ -46,6 +44,7 @@ import org.apache.accumulo.test.VerifyIngest; import org.apache.accumulo.test.VerifyIngest.VerifyParams; import org.apache.accumulo.test.harness.AccumuloClusterHarness; +import org.apache.accumulo.test.util.Wait; import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.fs.RawLocalFileSystem; import org.junit.jupiter.api.AfterEach; @@ -132,27 +131,15 @@ public void restartManagerRecovery() throws Exception { ZooCache zcache = ((ClientContext) c).getZooCache(); var zLockPath = getServerContext().getServerPaths().createManagerPath(); - Optional managerLockData; - do { - managerLockData = ServiceLock.getLockData(zcache, zLockPath, null); - if (managerLockData.isPresent()) { - log.info("Manager lock is still held"); - Thread.sleep(1000); - } - } while (managerLockData.isPresent()); + Wait.waitFor(() -> ServiceLock.getLockData(zcache, zLockPath, null).isEmpty(), 30_000, 250, + "Manager lock was not released before restart"); cluster.start(); Thread.sleep(5); control.stopAllServers(ServerType.MANAGER); - managerLockData = null; - do { - managerLockData = ServiceLock.getLockData(zcache, zLockPath, null); - if (managerLockData.isPresent()) { - log.info("Manager lock is still held"); - Thread.sleep(1000); - } - } while (managerLockData.isPresent()); + Wait.waitFor(() -> ServiceLock.getLockData(zcache, zLockPath, null).isEmpty(), 30_000, 250, + "Manager lock was not released before second restart"); cluster.start(); VerifyIngest.verifyIngest(c, params); } @@ -182,14 +169,8 @@ public void restartManagerSplit() throws Exception { ZooCache zcache = ((ClientContext) c).getZooCache(); var zLockPath = getServerContext().getServerPaths().createManagerPath(); - Optional managerLockData; - do { - managerLockData = ServiceLock.getLockData(zcache, zLockPath, null); - if (managerLockData.isPresent()) { - log.info("Manager lock is still held"); - Thread.sleep(1000); - } - } while (managerLockData.isPresent()); + Wait.waitFor(() -> ServiceLock.getLockData(zcache, zLockPath, null).isEmpty(), 30_000, 250, + "Manager lock was not released before restart"); cluster.start(); assertEquals(0, ret.get().intValue()); diff --git a/test/src/main/java/org/apache/accumulo/test/functional/ScanIdIT.java b/test/src/main/java/org/apache/accumulo/test/functional/ScanIdIT.java index 837a1202b9b..198c94e2940 100644 --- a/test/src/main/java/org/apache/accumulo/test/functional/ScanIdIT.java +++ b/test/src/main/java/org/apache/accumulo/test/functional/ScanIdIT.java @@ -62,6 +62,7 @@ import org.apache.accumulo.core.security.Authorizations; import org.apache.accumulo.core.security.ColumnVisibility; import org.apache.accumulo.test.harness.AccumuloClusterHarness; +import org.apache.accumulo.test.util.Wait; import org.apache.hadoop.io.Text; import org.junit.jupiter.api.AfterAll; import org.junit.jupiter.api.BeforeAll; @@ -180,10 +181,8 @@ public void testScanId() throws Exception { scanThreadsToClose.forEach(st -> st.scanner.close()); batchScanThreadsToClose.forEach(bst -> bst.bs.close()); - while (!getScanIds(client).isEmpty()) { - log.debug("Waiting for active scans to stop..."); - Thread.sleep(200); - } + Wait.waitFor(() -> getScanIds(client).isEmpty(), 30_000, 100, + "Scan IDs remained active after scanners closed"); } } diff --git a/test/src/main/java/org/apache/accumulo/test/functional/TabletMetadataIT.java b/test/src/main/java/org/apache/accumulo/test/functional/TabletMetadataIT.java index 4e9f9a3c357..7f0dc5c5a07 100644 --- a/test/src/main/java/org/apache/accumulo/test/functional/TabletMetadataIT.java +++ b/test/src/main/java/org/apache/accumulo/test/functional/TabletMetadataIT.java @@ -18,7 +18,6 @@ */ package org.apache.accumulo.test.functional; -import static java.util.concurrent.TimeUnit.SECONDS; import static org.apache.accumulo.minicluster.ServerType.TABLET_SERVER; import static org.junit.jupiter.api.Assertions.assertEquals; @@ -32,16 +31,14 @@ import org.apache.accumulo.core.metadata.TServerInstance; import org.apache.accumulo.core.metadata.schema.TabletMetadata; import org.apache.accumulo.miniclusterImpl.MiniAccumuloConfigImpl; +import org.apache.accumulo.test.util.Wait; import org.apache.hadoop.conf.Configuration; import org.junit.jupiter.api.Test; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; /** * Tests features of the Ample TabletMetadata class that can't be tested in TabletMetadataTest */ public class TabletMetadataIT extends ConfigurableMacBase { - private static final Logger log = LoggerFactory.getLogger(TabletMetadataIT.class); private static final int NUM_TSERVERS = 3; @Override @@ -57,11 +54,8 @@ public void configure(MiniAccumuloConfigImpl cfg, Configuration conf) { @Test public void getLiveTServersTest() throws Exception { try (AccumuloClient c = Accumulo.newClient().from(getClientProperties()).build()) { - while (c.instanceOperations().getServers(ServerId.Type.TABLET_SERVER).size() - != NUM_TSERVERS) { - log.info("Waiting for tservers to start up..."); - Thread.sleep(SECONDS.toMillis(5)); - } + Wait.waitFor(() -> c.instanceOperations().getServers(ServerId.Type.TABLET_SERVER).size() + == NUM_TSERVERS, 30_000, 250, "Tablet servers did not start up"); Set servers = TabletMetadata.getLiveTServers((ClientContext) c); assertEquals(NUM_TSERVERS, servers.size()); @@ -69,11 +63,8 @@ public void getLiveTServersTest() throws Exception { getCluster().killProcess(TABLET_SERVER, getCluster().getProcesses().get(TABLET_SERVER).iterator().next()); - while (c.instanceOperations().getServers(ServerId.Type.TABLET_SERVER).size() - == NUM_TSERVERS) { - log.info("Waiting for a tserver to die..."); - Thread.sleep(SECONDS.toMillis(5)); - } + Wait.waitFor(() -> c.instanceOperations().getServers(ServerId.Type.TABLET_SERVER).size() + != NUM_TSERVERS, 30_000, 250, "Tablet server did not leave the live list"); servers = TabletMetadata.getLiveTServers((ClientContext) c); assertEquals(NUM_TSERVERS - 1, servers.size()); } diff --git a/test/src/main/java/org/apache/accumulo/test/lock/ServiceLockIT.java b/test/src/main/java/org/apache/accumulo/test/lock/ServiceLockIT.java index 8324e84eb40..1fa3ee0478f 100644 --- a/test/src/main/java/org/apache/accumulo/test/lock/ServiceLockIT.java +++ b/test/src/main/java/org/apache/accumulo/test/lock/ServiceLockIT.java @@ -46,6 +46,7 @@ import org.apache.accumulo.core.lock.ServiceLockData.ThriftService; import org.apache.accumulo.core.lock.ServiceLockPaths.ServiceLockPath; import org.apache.accumulo.core.zookeeper.ZooSession; +import org.apache.accumulo.test.util.Wait; import org.apache.accumulo.test.zookeeper.ZooKeeperTestingServer; import org.apache.zookeeper.CreateMode; import org.apache.zookeeper.KeeperException; @@ -596,13 +597,12 @@ public void testLockParallel() throws Exception { workers.forEach(w -> assertNull(w.getException())); for (int i = 4; i > 0; i--) { + final int expectedChildren = i; + Wait.waitFor(() -> zk.getChildren(parent.toString(), null).size() == expectedChildren, + 30_000, 100, "Unexpected number of lock children"); List children = ServiceLock.validateAndSort(parent, zk.getChildren(parent.toString(), null)); - while (children.size() != i) { - Thread.sleep(100); - children = zk.getChildren(parent.toString(), null); - } - assertEquals(i, children.size()); + assertEquals(expectedChildren, children.size()); String first = children.get(0); int workerWithLock = parseLockWorkerName(first); LockWorker worker = workers.get(workerWithLock); diff --git a/test/src/main/java/org/apache/accumulo/test/metrics/MetricsIT.java b/test/src/main/java/org/apache/accumulo/test/metrics/MetricsIT.java index 598edaf0377..89d4f593cbd 100644 --- a/test/src/main/java/org/apache/accumulo/test/metrics/MetricsIT.java +++ b/test/src/main/java/org/apache/accumulo/test/metrics/MetricsIT.java @@ -85,6 +85,7 @@ import org.apache.accumulo.test.fate.SlowFateSplitManager; import org.apache.accumulo.test.functional.ConfigurableMacBase; import org.apache.accumulo.test.functional.SlowIterator; +import org.apache.accumulo.test.util.Wait; import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.io.Text; import org.junit.jupiter.api.AfterAll; @@ -506,9 +507,8 @@ static void doWorkToGenerateMetrics(AccumuloClient client, Class testClass) t cc.setWait(false); client.tableOperations().compact(tableName, new CompactionConfig().setWait(true)); client.tableOperations().delete(tableName); - while (client.tableOperations().exists(tableName)) { - Thread.sleep(1000); - } + Wait.waitFor(() -> !client.tableOperations().exists(tableName), 30_000, 250, + "Deleted table remained visible"); } @Override diff --git a/test/src/main/java/org/apache/accumulo/test/rpc/Mocket.java b/test/src/main/java/org/apache/accumulo/test/rpc/Mocket.java index e641fcf77f3..8b581707bc0 100644 --- a/test/src/main/java/org/apache/accumulo/test/rpc/Mocket.java +++ b/test/src/main/java/org/apache/accumulo/test/rpc/Mocket.java @@ -39,7 +39,6 @@ public class Mocket { private final TTransport clientTransport; private final TServerTransport serverTransport; private final AtomicBoolean closed = new AtomicBoolean(false); - private final TConfiguration config = new TConfiguration(); public Mocket() { Buffer serverQueue = new Buffer(); diff --git a/test/src/main/java/org/apache/accumulo/test/tracing/ScanTracingIT.java b/test/src/main/java/org/apache/accumulo/test/tracing/ScanTracingIT.java index 71359a408ff..b0771e85a99 100644 --- a/test/src/main/java/org/apache/accumulo/test/tracing/ScanTracingIT.java +++ b/test/src/main/java/org/apache/accumulo/test/tracing/ScanTracingIT.java @@ -20,7 +20,6 @@ import static org.apache.accumulo.core.trace.TraceAttributes.EXECUTOR_KEY; import static org.apache.accumulo.core.trace.TraceAttributes.EXTENT_KEY; -import static org.apache.accumulo.tserver.tablet.KVEntry.MEMORY_OVERHEAD; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertTrue; From c0b7c6f751c5a1df573df52939c8b30513637f2f Mon Sep 17 00:00:00 2001 From: Christopher Tubbs Date: Thu, 1 Oct 2026 20:03:12 -0400 Subject: [PATCH 2/3] Update additional tests Co-authored-by: OpenAI GPT-6 Luna --- .../accumulo/test/MultipleManagerFateIT.java | 24 +-- .../test/NamespacesIT_SimpleSuite.java | 53 +++--- .../accumulo/test/TabletServerGivesUpIT.java | 8 +- ...hriftServerBindsBeforeZooKeeperLockIT.java | 168 ++++++++---------- .../apache/accumulo/test/VolumeITBase.java | 13 +- .../test/compaction/CompactionExecutorIT.java | 12 +- .../CompactionPriorityQueueMetricsIT.java | 30 ++-- .../compaction/ExternalCompaction_3_IT.java | 89 +++++----- .../test/conf/store/ZooBasedConfigIT.java | 7 +- .../functional/AccumuloConfigurationIT.java | 22 ++- .../BalanceAfterCommsFailureIT.java | 18 +- .../BalanceInPresenceOfOfflineTableIT.java | 41 ++--- .../test/functional/CompactionIT.java | 40 ++--- .../test/functional/FateConcurrencyIT.java | 39 ++-- .../test/functional/GarbageCollectorIT.java | 49 ++--- .../test/functional/HalfDeadTServerIT.java | 9 +- .../test/functional/ManyWriteAheadLogsIT.java | 27 ++- .../test/functional/MemoryStarvedMajCIT.java | 10 +- .../test/functional/MemoryStarvedMinCIT.java | 5 +- .../test/functional/MemoryStarvedScanIT.java | 29 ++- .../test/functional/MetadataMaxFilesIT.java | 10 +- .../test/functional/PerTableCryptoIT.java | 23 +-- .../test/functional/RegexGroupBalanceIT.java | 43 ++--- .../accumulo/test/functional/ScanIdIT.java | 22 ++- .../accumulo/test/functional/SplitIT.java | 20 +-- .../accumulo/test/functional/SummaryIT.java | 20 +-- .../TabletManagementIteratorIT.java | 17 +- .../TabletResourceGroupBalanceIT.java | 33 ++-- .../test/functional/WALSunnyDayITBase.java | 34 ++-- .../accumulo/test/lock/ServiceLockIT.java | 5 +- .../test/manager/SuspendedTabletsIT.java | 63 ++++--- .../accumulo/test/shell/ShellServerIT.java | 20 +-- 32 files changed, 480 insertions(+), 523 deletions(-) diff --git a/test/src/main/java/org/apache/accumulo/test/MultipleManagerFateIT.java b/test/src/main/java/org/apache/accumulo/test/MultipleManagerFateIT.java index a15b1777f61..f0abb25e361 100644 --- a/test/src/main/java/org/apache/accumulo/test/MultipleManagerFateIT.java +++ b/test/src/main/java/org/apache/accumulo/test/MultipleManagerFateIT.java @@ -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; @@ -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; @@ -261,15 +261,15 @@ private static void waitToSeeManagers(ClientContext context, int expectedManager .stream().map(FateStore.FateReservation::getReservationUUID).collect(toSet()); log.debug("existingReservationUUIDs {}", existingReservationUUIDs); - 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); - } + // 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); var expectedServers = assistants.stream().map(ServiceLockPath::getServer) .map(HostAndPort::fromString).collect(toSet()); @@ -285,7 +285,7 @@ private static void waitToSeeManagers(ClientContext context, int expectedManager // reassigned. Set 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) -> { @@ -308,8 +308,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) { diff --git a/test/src/main/java/org/apache/accumulo/test/NamespacesIT_SimpleSuite.java b/test/src/main/java/org/apache/accumulo/test/NamespacesIT_SimpleSuite.java index f91124245ee..712d93ab800 100644 --- a/test/src/main/java/org/apache/accumulo/test/NamespacesIT_SimpleSuite.java +++ b/test/src/main/java/org/apache/accumulo/test/NamespacesIT_SimpleSuite.java @@ -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; @@ -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"); @@ -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 diff --git a/test/src/main/java/org/apache/accumulo/test/TabletServerGivesUpIT.java b/test/src/main/java/org/apache/accumulo/test/TabletServerGivesUpIT.java index faa1584d063..ed3918d4e44 100644 --- a/test/src/main/java/org/apache/accumulo/test/TabletServerGivesUpIT.java +++ b/test/src/main/java/org/apache/accumulo/test/TabletServerGivesUpIT.java @@ -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; @@ -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"); } } } diff --git a/test/src/main/java/org/apache/accumulo/test/ThriftServerBindsBeforeZooKeeperLockIT.java b/test/src/main/java/org/apache/accumulo/test/ThriftServerBindsBeforeZooKeeperLockIT.java index d8cd8e31982..ae4dbd6ec67 100644 --- a/test/src/main/java/org/apache/accumulo/test/ThriftServerBindsBeforeZooKeeperLockIT.java +++ b/test/src/main/java/org/apache/accumulo/test/ThriftServerBindsBeforeZooKeeperLockIT.java @@ -22,6 +22,7 @@ import java.io.IOException; import java.net.HttpURLConnection; +import java.net.InetSocketAddress; import java.net.Socket; import java.net.URI; import java.util.Collection; @@ -41,6 +42,7 @@ import org.apache.accumulo.server.util.PortUtils; import org.apache.accumulo.test.functional.FunctionalTestUtils; import org.apache.accumulo.test.harness.AccumuloClusterHarness; +import org.apache.accumulo.test.util.Wait; import org.junit.jupiter.api.Tag; import org.junit.jupiter.api.Test; import org.slf4j.Logger; @@ -72,57 +74,59 @@ public void testMonitorService() throws Exception { getClusterControl().start(ServerType.MONITOR, "localhost"); } - while (true) { + String[] monitorLocation = {null}; + Wait.waitFor(() -> { try { - MonitorUtil.getLocation(getServerContext()); - break; + monitorLocation[0] = MonitorUtil.getLocation(getServerContext()); } catch (Exception e) { LOG.debug("Failed to find active monitor location, retrying", e); - Thread.sleep(1000); } - } + return monitorLocation[0] != null; + }, 30_000, 250, "Active monitor location was not published to ZooKeeper"); LOG.debug("Found active monitor"); - int freePort = PortUtils.getRandomFreePort(); - String monitorUrl = "http://localhost:" + freePort; - Process monitor = null; + int[] freePort = {PortUtils.getRandomFreePort()}; + Process[] monitor = {null}; try { - LOG.debug("Starting standby monitor on {}", freePort); - monitor = startProcess(cluster, ServerType.MONITOR, freePort); + LOG.debug("Starting standby monitor on {}", freePort[0]); + monitor[0] = startProcess(cluster, ServerType.MONITOR, freePort[0]); - while (true) { - var url = new URI(monitorUrl).toURL(); + Wait.waitFor(() -> { + var url = new URI("http://localhost:" + freePort[0]).toURL(); try { HttpURLConnection cnxn = (HttpURLConnection) url.openConnection(); - final int responseCode = cnxn.getResponseCode(); - String errorText; - // This is our "assertion", but we want to re-check it if it's not what we expect - if (responseCode == HttpURLConnection.HTTP_OK) { - return; - } else { - errorText = FunctionalTestUtils.readAll(cnxn.getErrorStream()); + cnxn.setConnectTimeout(1000); + cnxn.setReadTimeout(1000); + try { + final int responseCode = cnxn.getResponseCode(); + String errorText; + // This is our "assertion", but we want to re-check it if it's not what we expect + if (responseCode == HttpURLConnection.HTTP_OK) { + return true; + } else { + errorText = FunctionalTestUtils.readAll(cnxn.getErrorStream()); + } + LOG.debug("Unexpected responseCode and/or error text, will retry: '{}' '{}'", + responseCode, errorText); + } finally { + cnxn.disconnect(); } - LOG.debug("Unexpected responseCode and/or error text, will retry: '{}' '{}'", - responseCode, errorText); } catch (Exception e) { LOG.debug("Caught exception trying to fetch monitor info", e); } - // Wait before trying again - Thread.sleep(1000); // Make sure the process is still up. Possible the "randomFreePort" we got wasn't actually - // free and the process - // died trying to bind it. Pick a new port and restart it in that case. - if (!monitor.isAlive()) { - freePort = PortUtils.getRandomFreePort(); - monitorUrl = "http://localhost:" + freePort; - LOG.debug("Monitor died, restarting it listening on {}", freePort); - monitor = startProcess(cluster, ServerType.MONITOR, freePort); + // free and the process died trying to bind it. Pick a new port and restart it in that case. + if (!monitor[0].isAlive()) { + freePort[0] = PortUtils.getRandomFreePort(); + LOG.debug("Monitor died, restarting it listening on {}", freePort[0]); + monitor[0] = startProcess(cluster, ServerType.MONITOR, freePort[0]); } - } + return false; + }, 30_000, 250, "Standby monitor did not serve requests"); } finally { - if (monitor != null) { - monitor.destroyForcibly(); + if (monitor[0] != null) { + monitor[0].destroyForcibly(); } } } @@ -135,50 +139,43 @@ public void testManagerService() throws Exception { try (AccumuloClient client = Accumulo.newClient().from(getClientProps()).build()) { // Wait for the Manager to grab its lock - while (true) { + Wait.waitFor(() -> { try { ServiceLockPath managerLockPath = getServerContext().getServerPaths().getManager(true); - if (managerLockPath != null) { - break; - } + return managerLockPath != null; } catch (Exception e) { LOG.debug("Failed to find active manager location, retrying", e); - Thread.sleep(1000); + return false; } - } + }, 30_000, 250, "Active manager lock was not acquired"); LOG.debug("Found active manager"); - int freePort = PortUtils.getRandomFreePort(); - Process manager = null; + int[] freePort = {PortUtils.getRandomFreePort()}; + Process[] manager = {null}; try { - LOG.debug("Starting standby manager on {}", freePort); - manager = startProcess(cluster, ServerType.MANAGER, freePort); + LOG.debug("Starting standby manager on {}", freePort[0]); + manager[0] = startProcess(cluster, ServerType.MANAGER, freePort[0]); - while (true) { - try (Socket s = new Socket("localhost", freePort)) { - if (s.isConnected()) { - // Pass - return; - } + Wait.waitFor(() -> { + try (Socket s = new Socket()) { + s.connect(new InetSocketAddress("localhost", freePort[0]), 1000); + return s.isConnected(); } catch (Exception e) { LOG.debug("Caught exception trying to connect to Manager", e); } - // Wait before trying again - Thread.sleep(1000); // Make sure the process is still up. Possible the "randomFreePort" we got wasn't - // actually - // free and the process - // died trying to bind it. Pick a new port and restart it in that case. - if (!manager.isAlive()) { - freePort = PortUtils.getRandomFreePort(); - LOG.debug("Manager died, restarting it listening on {}", freePort); - manager = startProcess(cluster, ServerType.MANAGER, freePort); + // actually free and the process died trying to bind it. Pick a new port and restart it. + if (!manager[0].isAlive()) { + freePort[0] = PortUtils.getRandomFreePort(); + LOG.debug("Manager died, restarting it listening on {}", freePort[0]); + manager[0] = startProcess(cluster, ServerType.MANAGER, freePort[0]); } - } + return false; + }, 30_000, 250, "Standby manager did not accept connections"); } finally { - if (manager != null) { - manager.destroyForcibly(); + if (manager[0] != null) { + manager[0].destroyForcibly(); } } } @@ -192,50 +189,43 @@ public void testGarbageCollectorPorts() throws Exception { try (AccumuloClient client = Accumulo.newClient().from(getClientProps()).build()) { // Wait for the Manager to grab its lock - while (true) { + Wait.waitFor(() -> { try { ServiceLockPath slp = getServerContext().getServerPaths().getGarbageCollector(true); - if (slp != null) { - break; - } + return slp != null; } catch (Exception e) { LOG.debug("Failed to find active gc location, retrying", e); - Thread.sleep(1000); + return false; } - } + }, 30_000, 250, "Active garbage collector lock was not acquired"); LOG.debug("Found active gc"); - int freePort = PortUtils.getRandomFreePort(); - Process manager = null; + int[] freePort = {PortUtils.getRandomFreePort()}; + Process[] manager = {null}; try { - LOG.debug("Starting standby gc on {}", freePort); - manager = startProcess(cluster, ServerType.GARBAGE_COLLECTOR, freePort); + LOG.debug("Starting standby gc on {}", freePort[0]); + manager[0] = startProcess(cluster, ServerType.GARBAGE_COLLECTOR, freePort[0]); - while (true) { - try (Socket s = new Socket("localhost", freePort)) { - if (s.isConnected()) { - // Pass - return; - } + Wait.waitFor(() -> { + try (Socket s = new Socket()) { + s.connect(new InetSocketAddress("localhost", freePort[0]), 1000); + return s.isConnected(); } catch (Exception e) { LOG.debug("Caught exception trying to connect to GC", e); } - // Wait before trying again - Thread.sleep(1000); // Make sure the process is still up. Possible the "randomFreePort" we got wasn't - // actually - // free and the process - // died trying to bind it. Pick a new port and restart it in that case. - if (!manager.isAlive()) { - freePort = PortUtils.getRandomFreePort(); - LOG.debug("GC died, restarting it listening on {}", freePort); - manager = startProcess(cluster, ServerType.GARBAGE_COLLECTOR, freePort); + // actually free and the process died trying to bind it. Pick a new port and restart it. + if (!manager[0].isAlive()) { + freePort[0] = PortUtils.getRandomFreePort(); + LOG.debug("GC died, restarting it listening on {}", freePort[0]); + manager[0] = startProcess(cluster, ServerType.GARBAGE_COLLECTOR, freePort[0]); } - } + return false; + }, 30_000, 250, "Standby garbage collector did not accept connections"); } finally { - if (manager != null) { - manager.destroyForcibly(); + if (manager[0] != null) { + manager[0].destroyForcibly(); } } } diff --git a/test/src/main/java/org/apache/accumulo/test/VolumeITBase.java b/test/src/main/java/org/apache/accumulo/test/VolumeITBase.java index 3c1b301744b..d9c94b69ab4 100644 --- a/test/src/main/java/org/apache/accumulo/test/VolumeITBase.java +++ b/test/src/main/java/org/apache/accumulo/test/VolumeITBase.java @@ -64,7 +64,6 @@ import org.apache.accumulo.core.metadata.schema.MetadataSchema; import org.apache.accumulo.core.security.Authorizations; import org.apache.accumulo.core.security.TablePermission; -import org.apache.accumulo.core.util.UtilWaitThread; import org.apache.accumulo.miniclusterImpl.MiniAccumuloConfigImpl; import org.apache.accumulo.server.ServerContext; import org.apache.accumulo.server.log.WalStateManager; @@ -72,6 +71,7 @@ import org.apache.accumulo.server.util.adminCommand.StopAll; import org.apache.accumulo.test.functional.ConfigurableMacBase; import org.apache.accumulo.test.util.FileMetadataUtil; +import org.apache.accumulo.test.util.Wait; import org.apache.commons.configuration2.PropertiesConfiguration; import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.fs.Path; @@ -204,7 +204,7 @@ private void verifyVolumesUsed(AccumuloClient client, String tableName, boolean } // keep retrying until WAL state information in ZooKeeper stabilizes or until test times out - retry: while (true) { + Wait.waitFor(() -> { WalStateManager wals = new WalStateManager(getServerContext()); try { outer: for (var wal : wals.getAllState()) { @@ -214,19 +214,18 @@ private void verifyVolumesUsed(AccumuloClient client, String tableName, boolean } } log.warn("Unexpected volume " + wal.path() + " (" + wal.state() + ")"); - UtilWaitThread.sleep(100); - continue retry; + return false; } + return true; } catch (WalStateManager.WalMarkerException e) { Throwable cause = e.getCause(); if (cause instanceof KeeperException.NoNodeException) { // ignore WALs being cleaned up - continue retry; + return false; } throw e; } - break; - } + }, 30_000, 100, "WAL state information did not stabilize"); // if a volume is chosen randomly for each tablet, then the probability that a volume will not // be chosen for any tablet is ((num_volumes - diff --git a/test/src/main/java/org/apache/accumulo/test/compaction/CompactionExecutorIT.java b/test/src/main/java/org/apache/accumulo/test/compaction/CompactionExecutorIT.java index 404e5e315bd..fb1a1e15e9b 100644 --- a/test/src/main/java/org/apache/accumulo/test/compaction/CompactionExecutorIT.java +++ b/test/src/main/java/org/apache/accumulo/test/compaction/CompactionExecutorIT.java @@ -391,9 +391,9 @@ public void testDispatchSystem() throws Exception { addFiles(client, "dst1", 1); addFiles(client, "dst2", 1); - while (getFiles(client, "dst1").size() > 3 || getFiles(client, "dst2").size() > 2) { - Thread.sleep(100); - } + Wait.waitFor( + () -> getFiles(client, "dst1").size() <= 3 && getFiles(client, "dst2").size() <= 2, 30_000, + 100, "Compaction did not reduce dst1 and dst2 to the expected file counts"); assertEquals(3, getFiles(client, "dst1").size()); assertEquals(2, getFiles(client, "dst2").size()); @@ -425,9 +425,9 @@ public void testDispatchUser() throws Exception { client.tableOperations().compact("dut2", new CompactionConfig().setWait(false) .setExecutionHints(Map.of("compaction_type", "special"))); - while (getFiles(client, "dut1").size() > 2 || getFiles(client, "dut2").size() > 3) { - Thread.sleep(100); - } + Wait.waitFor( + () -> getFiles(client, "dut1").size() <= 2 && getFiles(client, "dut2").size() <= 3, 30_000, + 100, "Compaction did not reduce dut1 and dut2 to the expected file counts"); assertEquals(2, getFiles(client, "dut1").size()); assertEquals(3, getFiles(client, "dut2").size()); diff --git a/test/src/main/java/org/apache/accumulo/test/compaction/CompactionPriorityQueueMetricsIT.java b/test/src/main/java/org/apache/accumulo/test/compaction/CompactionPriorityQueueMetricsIT.java index c8df5040c3a..073abb6c04c 100644 --- a/test/src/main/java/org/apache/accumulo/test/compaction/CompactionPriorityQueueMetricsIT.java +++ b/test/src/main/java/org/apache/accumulo/test/compaction/CompactionPriorityQueueMetricsIT.java @@ -338,29 +338,28 @@ public void testQueueMetrics() throws Exception { verifyData(c, tableName, 0, 100 * 100 - 1, false); } - boolean sawMetricsQ1 = false; - boolean sawMetricsQ1Size = false; + AtomicBoolean sawMetricsQ1 = new AtomicBoolean(false); + AtomicBoolean sawMetricsQ1Size = new AtomicBoolean(false); - while (!sawMetricsQ1 || !sawMetricsQ1Size) { + Wait.waitFor(() -> { while (!queueMetrics.isEmpty()) { var qm = queueMetrics.take(); if (qm.getName().contains(COMPACTOR_JOB_PRIORITY_QUEUE_JOBS_QUEUED.getName()) && qm.getTags().containsValue(QUEUE1_METRIC_LABEL)) { if (Integer.parseInt(qm.getValue()) > 0) { - sawMetricsQ1 = true; + sawMetricsQ1.set(true); } } if (qm.getName().contains(COMPACTOR_JOB_PRIORITY_QUEUE_JOBS_SIZE.getName()) && qm.getTags().containsValue(QUEUE1_METRIC_LABEL)) { if (Integer.parseInt(qm.getValue()) > 0) { - sawMetricsQ1Size = true; + sawMetricsQ1Size.set(true); } } } - - // If metrics are not found in the queue, sleep until the next poll. - UtilWaitThread.sleep(TestStatsDRegistryFactory.pollingFrequency.toMillis()); - } + return sawMetricsQ1.get() && sawMetricsQ1Size.get(); + }, 60_000, TestStatsDRegistryFactory.pollingFrequency.toMillis(), + "Expected queue 1 jobs and size metrics"); // Set lowest priority to the lowest possible system compaction priority long lowestPriority = Short.MIN_VALUE; @@ -404,18 +403,18 @@ public void testQueueMetrics() throws Exception { getCluster().getConfig().getClusterServerConfiguration().addCompactorResourceGroup(QUEUE1, 1); getCluster().getClusterControl().start(ServerType.COMPACTOR); - boolean emptyQueue = false; + AtomicBoolean emptyQueue = new AtomicBoolean(false); // Make sure that metrics added to the queue are recent UtilWaitThread.sleep(TestStatsDRegistryFactory.pollingFrequency.toMillis()); - while (!emptyQueue) { + Wait.waitFor(() -> { while (!queueMetrics.isEmpty()) { var metric = queueMetrics.take(); if (metric.getName().contains(COMPACTOR_JOB_PRIORITY_QUEUE_JOBS_QUEUED.getName()) && metric.getTags().containsValue(QUEUE1_METRIC_LABEL)) { if (Integer.parseInt(metric.getValue()) == 0) { - emptyQueue = true; + emptyQueue.set(true); } } @@ -423,12 +422,13 @@ public void testQueueMetrics() throws Exception { // above queue. if (metric.getName().equals(COMPACTOR_JOB_PRIORITY_QUEUES.getName())) { if (Integer.parseInt(metric.getValue()) == 0) { - emptyQueue = true; + emptyQueue.set(true); } } } - UtilWaitThread.sleep(TestStatsDRegistryFactory.pollingFrequency.toMillis()); - } + return emptyQueue.get(); + }, 60_000, TestStatsDRegistryFactory.pollingFrequency.toMillis(), + "Expected the compaction queue to empty"); } /** diff --git a/test/src/main/java/org/apache/accumulo/test/compaction/ExternalCompaction_3_IT.java b/test/src/main/java/org/apache/accumulo/test/compaction/ExternalCompaction_3_IT.java index f1edec9bc76..54c3190aed8 100644 --- a/test/src/main/java/org/apache/accumulo/test/compaction/ExternalCompaction_3_IT.java +++ b/test/src/main/java/org/apache/accumulo/test/compaction/ExternalCompaction_3_IT.java @@ -37,6 +37,7 @@ import java.util.Optional; import java.util.Set; import java.util.TreeMap; +import java.util.concurrent.atomic.AtomicReference; import java.util.stream.Collectors; import org.apache.accumulo.core.client.Accumulo; @@ -52,7 +53,6 @@ import org.apache.accumulo.core.metadata.schema.TabletMetadata; import org.apache.accumulo.core.metadata.schema.TabletMetadata.ColumnType; import org.apache.accumulo.core.metadata.schema.TabletsMetadata; -import org.apache.accumulo.core.util.UtilWaitThread; import org.apache.accumulo.core.util.compaction.ExternalCompactionUtil; import org.apache.accumulo.core.util.compaction.RunningCompactionInfo; import org.apache.accumulo.manager.compaction.coordinator.CompactionCoordinator; @@ -139,17 +139,15 @@ public void testMergeCancelsExternalCompaction() throws Exception { confirmCompactionsNoLongerRunning(getCluster().getServerContext(), ecids); // ensure compaction ids were deleted by merge operation from metadata table - try (TabletsMetadata tm = getCluster().getServerContext().getAmple().readTablets() - .forTable(tid).fetch(ColumnType.ECOMP).build()) { - Set ecids2 = tm.stream() - .flatMap(t -> t.getExternalCompactions().keySet().stream()).collect(Collectors.toSet()); - // keep checking until test times out - while (!Collections.disjoint(ecids, ecids2)) { - UtilWaitThread.sleep(25); - ecids2 = tm.stream().flatMap(t -> t.getExternalCompactions().keySet().stream()) + Wait.waitFor(() -> { + try (TabletsMetadata tm = getCluster().getServerContext().getAmple().readTablets() + .forTable(tid).fetch(ColumnType.ECOMP).build()) { + Set ecids2 = tm.stream() + .flatMap(t -> t.getExternalCompactions().keySet().stream()) .collect(Collectors.toSet()); + return Collections.disjoint(ecids, ecids2); } - } + }); // Verify that the tmp file are cleaned up Wait.waitFor(() -> FindCompactionTmpFiles @@ -193,11 +191,16 @@ public void testCoordinatorRestartsDuringCompaction() throws Exception { ServerContext ctx = getCluster().getServerContext(); // Wait for all compactions to start - Map originalRunningInfo = null; - do { - originalRunningInfo = getRunningCompactionInformation(ctx, ecids); - } while (originalRunningInfo == null - || originalRunningInfo.values().stream().allMatch(rci -> rci.duration == 0)); + AtomicReference> originalRunningInfoRef = + new AtomicReference<>(); + Wait.waitFor(() -> { + Map runningInfo = + getRunningCompactionInformation(ctx, ecids); + originalRunningInfoRef.set(runningInfo); + return runningInfo.values().stream().anyMatch(rci -> rci.duration > 0); + }, 30_000, 250, "Compaction did not start"); + Map originalRunningInfo = + originalRunningInfoRef.get(); // Stop the Manager (Coordinator) getCluster().getClusterControl().stop(ServerType.MANAGER); @@ -227,41 +230,41 @@ public void testCoordinatorRestartsDuringCompaction() throws Exception { } private Map getRunningCompactionInformation( - ServerContext ctx, Set ecids) throws InterruptedException { + ServerContext ctx, Set ecids) { final Map results = new HashMap<>(); - while (results.isEmpty()) { - Map running = null; - while (running == null || running.isEmpty()) { - try { - Optional coordinatorHost = - ExternalCompactionUtil.findCompactionCoordinator(ctx); - if (coordinatorHost.isEmpty()) { - throw new TTransportException( - "Unable to get CompactionCoordinator address from ZooKeeper"); - } - running = getRunningCompactions(ctx); - } catch (TException t) { - running = null; - Thread.sleep(2000); + Wait.waitFor(() -> { + try { + Optional coordinatorHost = + ExternalCompactionUtil.findCompactionCoordinator(ctx); + if (coordinatorHost.isEmpty()) { + throw new TTransportException( + "Unable to get CompactionCoordinator address from ZooKeeper"); } - } - for (ExternalCompactionId ecid : ecids) { - final TExternalCompaction tec = running.get(ecid.canonical()); - if (tec != null && tec.getUpdatesSize() > 0) { - // When the coordinator restarts it inserts a message into the updates. If this - // is the last message, then don't insert this into the results. We want to get - // an actual update from the Compactor. - TreeMap sorted = new TreeMap<>(tec.getUpdates()); - var lastEntry = sorted.lastEntry(); - if (lastEntry.getValue().getMessage().equals(CompactionCoordinator.RESTART_UPDATE_MSG)) { - continue; + Map running = getRunningCompactions(ctx); + if (running.isEmpty()) { + return false; + } + for (ExternalCompactionId ecid : ecids) { + final TExternalCompaction tec = running.get(ecid.canonical()); + if (tec != null && tec.getUpdatesSize() > 0) { + // When the coordinator restarts it inserts a message into the updates. If this + // is the last message, then don't insert this into the results. We want to get + // an actual update from the Compactor. + TreeMap sorted = new TreeMap<>(tec.getUpdates()); + var lastEntry = sorted.lastEntry(); + if (lastEntry.getValue().getMessage().equals(CompactionCoordinator.RESTART_UPDATE_MSG)) { + continue; + } + results.put(ecid, new RunningCompactionInfo(tec)); } - results.put(ecid, new RunningCompactionInfo(tec)); } + return !results.isEmpty(); + } catch (TException t) { + return false; } - } + }, 30_000, 250, "No running compaction information was reported"); return results; } diff --git a/test/src/main/java/org/apache/accumulo/test/conf/store/ZooBasedConfigIT.java b/test/src/main/java/org/apache/accumulo/test/conf/store/ZooBasedConfigIT.java index 8677ac8ccd2..48edee4fd25 100644 --- a/test/src/main/java/org/apache/accumulo/test/conf/store/ZooBasedConfigIT.java +++ b/test/src/main/java/org/apache/accumulo/test/conf/store/ZooBasedConfigIT.java @@ -59,6 +59,7 @@ import org.apache.accumulo.server.conf.store.SystemPropKey; import org.apache.accumulo.server.conf.store.TablePropKey; import org.apache.accumulo.server.conf.store.impl.ZooPropStore; +import org.apache.accumulo.test.util.Wait; import org.apache.accumulo.test.zookeeper.ZooKeeperTestingServer; import org.apache.zookeeper.CreateMode; import org.apache.zookeeper.ZooDefs; @@ -278,10 +279,8 @@ public void expireTest() throws Exception { // allow ZooKeeper notification time to propagate - int retries = 5; - do { - Thread.sleep(25); - } while (changeCount >= testListener.getZkChangeCount() && --retries > 0); + Wait.waitFor(() -> changeCount < testListener.getZkChangeCount(), 30_000, 250, + "ZooKeeper change notification was not received"); assertTrue(changeCount < testListener.getZkChangeCount()); diff --git a/test/src/main/java/org/apache/accumulo/test/functional/AccumuloConfigurationIT.java b/test/src/main/java/org/apache/accumulo/test/functional/AccumuloConfigurationIT.java index 9ea4c9f38f8..f3eeb354fa3 100644 --- a/test/src/main/java/org/apache/accumulo/test/functional/AccumuloConfigurationIT.java +++ b/test/src/main/java/org/apache/accumulo/test/functional/AccumuloConfigurationIT.java @@ -29,6 +29,7 @@ import org.apache.accumulo.server.ServerContext; import org.apache.accumulo.test.harness.MiniClusterConfigurationCallback; import org.apache.accumulo.test.harness.SharedMiniClusterBase; +import org.apache.accumulo.test.util.Wait; import org.apache.hadoop.conf.Configuration; import org.junit.jupiter.api.AfterAll; import org.junit.jupiter.api.BeforeAll; @@ -71,16 +72,19 @@ public void testInvalidation() throws Exception { ctx.getConfiguration().invalidateCache(); - int oldValueReturned = 0; - while (ctx.getConfiguration().get(fakeProperty).equals(initialThreads)) { - oldValueReturned++; - Thread.sleep(25); - } - System.out.println("Configuration returned old value " + oldValueReturned + " times and took " - + timer.elapsed(TimeUnit.MILLISECONDS) + "ms"); - + int[] oldValueReturned = {0}; + Wait.waitFor(() -> { + String value = ctx.getConfiguration().get(fakeProperty); + if (initialThreads.equals(value)) { + oldValueReturned[0]++; + } + return !initialThreads.equals(value); + }, 30_000, 250, "Configuration did not return updated value after cache invalidation"); + System.out.println("Configuration returned old value " + oldValueReturned[0] + + " times and took " + timer.elapsed(TimeUnit.MILLISECONDS) + "ms"); + + assertEquals(0, oldValueReturned[0]); assertEquals("4", ctx.getConfiguration().get(fakeProperty)); - assertEquals(0, oldValueReturned); } diff --git a/test/src/main/java/org/apache/accumulo/test/functional/BalanceAfterCommsFailureIT.java b/test/src/main/java/org/apache/accumulo/test/functional/BalanceAfterCommsFailureIT.java index 30c03500b9c..f5cc437424e 100644 --- a/test/src/main/java/org/apache/accumulo/test/functional/BalanceAfterCommsFailureIT.java +++ b/test/src/main/java/org/apache/accumulo/test/functional/BalanceAfterCommsFailureIT.java @@ -39,6 +39,7 @@ import org.apache.accumulo.miniclusterImpl.MiniAccumuloConfigImpl; import org.apache.accumulo.miniclusterImpl.ProcessReference; import org.apache.accumulo.test.BalanceIT; +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; @@ -97,19 +98,20 @@ public void test() throws Exception { private void checkBalance(AccumuloClient c) throws Exception { - Map tableLocations = null; - long unassignedTablets = 1; - for (int i = 0; unassignedTablets > 0 && i < 10; i++) { - tableLocations = BalanceIT.countLocations(c, "test"); - unassignedTablets = - tableLocations.entrySet().stream().filter(e -> e.getKey().equals("none")).count(); + Wait.waitFor(() -> { + Map tableLocations = BalanceIT.countLocations(c, "test"); + long unassignedTablets = tableLocations.entrySet().stream() + .filter(e -> e.getKey().equals("none")).count(); if (unassignedTablets > 0) { log.info("Found {} unassigned tablets, sleeping 3 seconds for tablet assignment", unassignedTablets); - Thread.sleep(3000); } - } + return unassignedTablets == 0; + }, 30_000, 3_000, "Unassigned tablets were not assigned within 30 seconds"); + Map tableLocations = BalanceIT.countLocations(c, "test"); + long unassignedTablets = tableLocations.entrySet().stream() + .filter(e -> e.getKey().equals("none")).count(); assertEquals(0, unassignedTablets, "Unassigned tablets were not assigned within 30 seconds"); assertNotNull(tableLocations); assertTrue(tableLocations.size() > 1, "Expected to have at least two TabletServers"); diff --git a/test/src/main/java/org/apache/accumulo/test/functional/BalanceInPresenceOfOfflineTableIT.java b/test/src/main/java/org/apache/accumulo/test/functional/BalanceInPresenceOfOfflineTableIT.java index 42d512f8917..690c1f8b592 100644 --- a/test/src/main/java/org/apache/accumulo/test/functional/BalanceInPresenceOfOfflineTableIT.java +++ b/test/src/main/java/org/apache/accumulo/test/functional/BalanceInPresenceOfOfflineTableIT.java @@ -18,7 +18,6 @@ */ package org.apache.accumulo.test.functional; -import static org.junit.jupiter.api.Assertions.assertTrue; import static org.junit.jupiter.api.Assumptions.assumeTrue; import java.util.HashMap; @@ -138,19 +137,12 @@ public void test() throws Exception { VerifyIngest.verifyIngest(accumuloClient, params); log.debug("waiting for balancing, up to ~5 minutes to allow for migration cleanup."); - final long startTime = System.currentTimeMillis(); - long currentWait = 1_000; - boolean balancingWorked = false; - - while (!balancingWorked && (System.currentTimeMillis() - startTime) < ((5 * 60 + 15) * 1000)) { - Thread.sleep(currentWait); - currentWait = Math.min(currentWait * 2, 30_000); - + Wait.waitFor(() -> { log.debug("fetch the list of tablets assigned to each tserver."); if (accumuloClient.instanceOperations().getServers(ServerId.Type.TABLET_SERVER).size() < 2) { - log.debug("we need >= 2 servers. sleeping for {}ms", currentWait); - continue; + log.debug("we need >= 2 servers"); + return false; } Map tableLocations = BalanceIT.countLocations(accumuloClient, TEST_TABLE); @@ -158,13 +150,13 @@ public void test() throws Exception { long unassignedTablets = tableLocations.entrySet().stream().filter(e -> e.getKey().equals("none")).count(); if (unassignedTablets != 0) { - log.debug("We shouldn't have unassigned tablets. sleeping for {}ms", currentWait); - continue; + log.debug("We shouldn't have unassigned tablets."); + return false; } var numSplits = accumuloClient.tableOperations().listSplits(TEST_TABLE).size(); if (numSplits < 20) { - log.debug("Waiting for 20 splits, saw {} split. sleeping for {}ms", numSplits, currentWait); - continue; + log.debug("Waiting for 20 splits, saw {} split.", numSplits); + return false; } HashMap serversWithTablets = tableLocations.entrySet().stream() @@ -176,27 +168,22 @@ public void test() throws Exception { }); if (serversWithTablets.isEmpty()) { - continue; + return false; } if (serversWithTablets.values().iterator().next() <= 10) { - log.debug("We should have > 10 tablets. sleeping for {}ms tabletsPerServer:{}", currentWait, - List.of(serversWithTablets)); - continue; + log.debug("We should have > 10 tablets. tabletsPerServer:{}", List.of(serversWithTablets)); + return false; } Integer min = serversWithTablets.values().stream().min(Long::compare).orElseThrow(); Integer max = serversWithTablets.values().stream().max(Long::compare).orElseThrow(); log.debug("Min={}, Max={}", min, max); if ((min / ((double) max)) < 0.5) { - log.debug( - "ratio of min to max tablets per server should be roughly even. sleeping for {}ms", - currentWait); - continue; + log.debug("ratio of min to max tablets per server should be roughly even"); + return false; } - balancingWorked = true; - } - - assertTrue(balancingWorked, "did not properly balance"); + return true; + }, (5 * 60 + 15) * 1000, 1_000, "did not properly balance"); } } diff --git a/test/src/main/java/org/apache/accumulo/test/functional/CompactionIT.java b/test/src/main/java/org/apache/accumulo/test/functional/CompactionIT.java index 000ffd26e07..96753127994 100644 --- a/test/src/main/java/org/apache/accumulo/test/functional/CompactionIT.java +++ b/test/src/main/java/org/apache/accumulo/test/functional/CompactionIT.java @@ -453,18 +453,15 @@ public void testConcurrentWithIterators() throws Exception { }; // wait until the filtering done by all three compactions is seen - while (!expected.equals(actualSupplier.get())) { - Thread.sleep(250); - } + Wait.waitFor(() -> expected.equals(actualSupplier.get()), 30_000, 250, + "Concurrent compactions did not filter the expected data"); // eventually the compactions should clean up all of their metadata, wait for this to happen - while (countTablets(tableName, + Wait.waitFor(() -> countTablets(tableName, tabletMetadata -> !tabletMetadata.getCompacted().isEmpty() || tabletMetadata.getSelectedFiles() != null - || !tabletMetadata.getExternalCompactions().isEmpty()) - > 0) { - Thread.sleep(250); - } + || !tabletMetadata.getExternalCompactions().isEmpty()) == 0, + 30_000, 250, "Concurrent compactions did not clean up their metadata"); } } @@ -555,18 +552,15 @@ && countTablets(tableName, tabletMetadata -> !tabletMetadata.getCompacted().isEm }; // wait until the filtering done by all three compactions is seen - while (!expected.equals(actualSupplier.get())) { - Thread.sleep(250); - } + Wait.waitFor(() -> expected.equals(actualSupplier.get()), 30_000, 250, + "Concurrent compactions did not filter the expected data"); // eventually the compactions should clean up all of their metadata, wait for this to happen - while (countTablets(tableName, + Wait.waitFor(() -> countTablets(tableName, tabletMetadata -> !tabletMetadata.getCompacted().isEmpty() || tabletMetadata.getSelectedFiles() != null - || !tabletMetadata.getExternalCompactions().isEmpty()) - > 0) { - Thread.sleep(250); - } + || !tabletMetadata.getExternalCompactions().isEmpty()) == 0, + 30_000, 250, "Concurrent compactions did not clean up their metadata"); } } @@ -1174,7 +1168,8 @@ public void testGetActiveCompactions() throws Exception { started.await(); List compactions = new ArrayList<>(); - do { + Wait.waitFor(() -> { + compactions.clear(); getActiveCompactions(client.instanceOperations()).forEach((ac) -> { try { if (ac.getTable().equals(table1)) { @@ -1184,15 +1179,16 @@ public void testGetActiveCompactions() throws Exception { fail("Table was deleted during test, should not happen"); } }); - Thread.sleep(1000); - } while (compactions.isEmpty()); + return !compactions.isEmpty(); + }, 30_000, 1_000, "Compaction did not become active"); ActiveCompaction running1 = compactions.get(0); ServerId host = running1.getServerId(); assertTrue(host.getType() == ServerId.Type.COMPACTOR); compactions.clear(); - do { + Wait.waitFor(() -> { + compactions.clear(); client.instanceOperations().getActiveCompactions(List.of(host)).forEach((ac) -> { try { if (ac.getTable().equals(table1)) { @@ -1202,8 +1198,8 @@ public void testGetActiveCompactions() throws Exception { fail("Table was deleted during test, should not happen"); } }); - Thread.sleep(1000); - } while (compactions.isEmpty()); + return !compactions.isEmpty(); + }, 30_000, 1_000, "Compaction did not resume on the compactor"); ActiveCompaction running2 = compactions.get(0); assertEquals(running1.getInputFiles(), running2.getInputFiles()); diff --git a/test/src/main/java/org/apache/accumulo/test/functional/FateConcurrencyIT.java b/test/src/main/java/org/apache/accumulo/test/functional/FateConcurrencyIT.java index 5eec44fb2b7..f615ea055d8 100644 --- a/test/src/main/java/org/apache/accumulo/test/functional/FateConcurrencyIT.java +++ b/test/src/main/java/org/apache/accumulo/test/functional/FateConcurrencyIT.java @@ -195,26 +195,29 @@ public void changeTableStateTest() throws Exception { private boolean findFate(String aTableName) { boolean isMeta = aTableName.startsWith(Namespace.ACCUMULO.name()); log.debug("Look for fate {}", aTableName); - for (int retry = 0; retry < 5; retry++) { - try { - boolean found = - isMeta ? lookupFateInZookeeper(aTableName) : lookupFateInAccumulo(aTableName); - log.trace("Try {}: Fate in {} for table {} : {}", retry, isMeta ? "zk" : "accumulo", - aTableName, found); - if (found) { - log.debug("Found fate {}", aTableName); - return true; - } else { - Thread.sleep(150); + AtomicInteger retries = new AtomicInteger(); + try { + Wait.waitFor(() -> { + int retry = retries.getAndIncrement(); + try { + boolean found = + isMeta ? lookupFateInZookeeper(aTableName) : lookupFateInAccumulo(aTableName); + log.trace("Try {}: Fate in {} for table {} : {}", retry, isMeta ? "zk" : "accumulo", + aTableName, found); + if (found) { + log.debug("Found fate {}", aTableName); + } + return found; + } catch (Exception ex) { + log.debug("Find fate failed for table name {} with exception, will retry", aTableName, ex); + return false; } - } catch (InterruptedException ex) { - Thread.currentThread().interrupt(); - return false; - } catch (Exception ex) { - log.debug("Find fate failed for table name {} with exception, will retry", aTableName, ex); - } + // Keep the original five lookup attempts at 150 ms intervals. + }, 600, 150, "FATE operation was not found for table " + aTableName); + return true; + } catch (IllegalStateException ex) { + return false; } - return false; } /** diff --git a/test/src/main/java/org/apache/accumulo/test/functional/GarbageCollectorIT.java b/test/src/main/java/org/apache/accumulo/test/functional/GarbageCollectorIT.java index 69f1a482e7f..90765fe4723 100644 --- a/test/src/main/java/org/apache/accumulo/test/functional/GarbageCollectorIT.java +++ b/test/src/main/java/org/apache/accumulo/test/functional/GarbageCollectorIT.java @@ -73,6 +73,7 @@ import org.apache.accumulo.test.TestIngest; import org.apache.accumulo.test.VerifyIngest; import org.apache.accumulo.test.VerifyIngest.VerifyParams; +import org.apache.accumulo.test.util.Wait; import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.fs.Path; import org.apache.hadoop.fs.RawLocalFileSystem; @@ -148,14 +149,16 @@ public void gcTest() throws Exception { int before = countFiles(pathString); log.info("Counted {} files in path: {}", before, pathString); - while (true) { - Thread.sleep(SECONDS.toMillis(1)); + int[] previous = {before}; + Wait.waitFor(() -> { int more = countFiles(pathString); - if (more <= before) { - break; + if (more > previous[0]) { + previous[0] = more; + return false; } - before = more; - } + return true; + }, 30_000, SECONDS.toMillis(1), "File count did not stabilize in " + pathString); + before = previous[0]; // restart GC log.info("Restarting GC..."); @@ -181,17 +184,21 @@ public void gcLotsOfCandidatesIT() throws Exception { cluster.getConfig().setDefaultMemory(32, MemoryUnit.MEGABYTE); ProcessInfo gc = cluster.exec(SimpleGarbageCollector.class); Thread.sleep(SECONDS.toMillis(20)); - String output = ""; - while (!output.contains("has exceeded the threshold")) { - try { - output = gc.readStdOut(); - } catch (UncheckedIOException ex) { - log.error("IO error reading the IT's accumulo-gc STDOUT", ex); - break; - } + String[] output = {""}; + try { + Wait.waitFor(() -> { + try { + output[0] = gc.readStdOut(); + } catch (UncheckedIOException ex) { + log.error("IO error reading the IT's accumulo-gc STDOUT", ex); + return true; + } + return output[0].contains("has exceeded the threshold"); + }, 30_000, 100, "GC output did not contain the threshold message"); + } finally { + gc.getProcess().destroy(); } - gc.getProcess().destroy(); - assertTrue(output.contains("has exceeded the threshold")); + assertTrue(output[0].contains("has exceeded the threshold")); } } @@ -256,15 +263,15 @@ public void testInvalidDelete() throws Exception { ProcessInfo gc = cluster.exec(SimpleGarbageCollector.class); try { - String output = ""; - while (!output.contains("Ignoring invalid deletion candidate")) { - Thread.sleep(250); + String[] output = {""}; + Wait.waitFor(() -> { try { - output = gc.readStdOut(); + output[0] = gc.readStdOut(); } catch (UncheckedIOException ioe) { log.error("Could not read all from cluster.", ioe); } - } + return output[0].contains("Ignoring invalid deletion candidate"); + }, 30_000, 250, "GC output did not contain the invalid deletion candidate message"); } finally { gc.getProcess().destroy(); } diff --git a/test/src/main/java/org/apache/accumulo/test/functional/HalfDeadTServerIT.java b/test/src/main/java/org/apache/accumulo/test/functional/HalfDeadTServerIT.java index 8bb368dfc9a..fe8084fc57c 100644 --- a/test/src/main/java/org/apache/accumulo/test/functional/HalfDeadTServerIT.java +++ b/test/src/main/java/org/apache/accumulo/test/functional/HalfDeadTServerIT.java @@ -44,6 +44,7 @@ import org.apache.accumulo.start.Main; import org.apache.accumulo.test.TestIngest; import org.apache.accumulo.test.VerifyIngest; +import org.apache.accumulo.test.util.Wait; import org.apache.accumulo.tserver.TabletServer; import org.apache.hadoop.conf.Configuration; import org.junit.jupiter.api.BeforeAll; @@ -159,10 +160,10 @@ public void testTimeout() throws Exception { public String test(int seconds, boolean expectTserverDied) throws Exception { assumeTrue(sharedLibBuilt.get(), "Shared library did not build"); try (AccumuloClient client = Accumulo.newClient().from(getClientProperties()).build()) { - while (client.instanceOperations().getServers(ServerId.Type.TABLET_SERVER).isEmpty()) { - // wait until the tserver that we need to kill is running - Thread.sleep(50); - } + // wait until the tserver that we need to kill is running + Wait.waitFor( + () -> !client.instanceOperations().getServers(ServerId.Type.TABLET_SERVER).isEmpty(), + 30_000, 250, "Tablet server did not start"); // create our own tablet server with the special test library Path confDirPath = cluster.getConfig().getDir().toPath(); diff --git a/test/src/main/java/org/apache/accumulo/test/functional/ManyWriteAheadLogsIT.java b/test/src/main/java/org/apache/accumulo/test/functional/ManyWriteAheadLogsIT.java index 6c2ad5193f4..e9ca380c047 100644 --- a/test/src/main/java/org/apache/accumulo/test/functional/ManyWriteAheadLogsIT.java +++ b/test/src/main/java/org/apache/accumulo/test/functional/ManyWriteAheadLogsIT.java @@ -170,35 +170,32 @@ public void testMany() throws Exception { private void addOpenWals(ServerContext c, Set allWalsSeen) throws Exception { - int open = 0; - int attempts = 0; - boolean foundWal = false; + int[] open = {0}; + int[] attempts = {0}; - while (open == 0) { - attempts++; + Wait.waitFor(() -> { + attempts[0]++; + open[0] = 0; Map wals = WALSunnyDayIT._getWals(c); Set> es = wals.entrySet(); for (Entry entry : es) { if (entry.getValue() == WalState.OPEN) { - open++; + open[0]++; allWalsSeen.add(entry.getKey()); - foundWal = true; } else { // log CLOSED or UNREFERENCED to help debug this test log.debug("The WalState for {} is {}", entry.getKey(), entry.getValue()); } } - if (!foundWal) { - Thread.sleep(50); - if (attempts % 50 == 0) { - log.debug("No open WALs found in {} attempts.", attempts); - } + if (open[0] == 0 && attempts[0] % 50 == 0) { + log.debug("No open WALs found in {} attempts.", attempts[0]); } - } + return open[0] > 0; + }, 30_000, 250, "No open WALs found"); - log.debug("It took {} attempt(s) to find {} open WALs", attempts, open); - assertTrue(open > 0 && open < 4, "Open WALs not in expected range " + open); + log.debug("It took {} attempt(s) to find {} open WALs", attempts[0], open[0]); + assertTrue(open[0] > 0 && open[0] < 4, "Open WALs not in expected range " + open[0]); } private int countClosedWals(ServerContext c) throws Exception { diff --git a/test/src/main/java/org/apache/accumulo/test/functional/MemoryStarvedMajCIT.java b/test/src/main/java/org/apache/accumulo/test/functional/MemoryStarvedMajCIT.java index 6821fd1f539..385b0df1645 100644 --- a/test/src/main/java/org/apache/accumulo/test/functional/MemoryStarvedMajCIT.java +++ b/test/src/main/java/org/apache/accumulo/test/functional/MemoryStarvedMajCIT.java @@ -39,7 +39,6 @@ import org.apache.accumulo.core.conf.Property; import org.apache.accumulo.core.lock.ServiceLockPaths.AddressSelector; import org.apache.accumulo.core.lock.ServiceLockPaths.ResourceGroupPredicate; -import org.apache.accumulo.core.util.UtilWaitThread; import org.apache.accumulo.core.util.compaction.ExternalCompactionUtil; import org.apache.accumulo.minicluster.ServerType; import org.apache.accumulo.miniclusterImpl.MiniAccumuloConfigImpl; @@ -165,15 +164,14 @@ public void testMajCPauses() throws Exception { // Calling getRunningCompaction on the MemoryConsumingCompactor // will consume the free memory LOG.info("Calling getRunningCompaction on {}", compactorAddr); - boolean success = false; - while (!success) { + waitFor(() -> { try { ExternalCompactionUtil.getRunningCompaction(compactorAddr, ctx); - success = true; + return true; } catch (Exception e) { - UtilWaitThread.sleep(3000); + return false; } - } + }, 30_000, 250, "getRunningCompaction RPC did not succeed"); ReadWriteIT.ingest(client, 100, 100, 100, 0, table); compactionThread.start(); diff --git a/test/src/main/java/org/apache/accumulo/test/functional/MemoryStarvedMinCIT.java b/test/src/main/java/org/apache/accumulo/test/functional/MemoryStarvedMinCIT.java index c7b14e173c9..aae61fef8a4 100644 --- a/test/src/main/java/org/apache/accumulo/test/functional/MemoryStarvedMinCIT.java +++ b/test/src/main/java/org/apache/accumulo/test/functional/MemoryStarvedMinCIT.java @@ -163,10 +163,7 @@ public void testMinCPauses() throws Exception { ingestThread.start(); - while (paused <= 0) { - Thread.sleep(1000); - paused = MINC_PAUSED_COUNT.intValue(); - } + waitFor(() -> MINC_PAUSED_COUNT.intValue() > 0); assertTrue(getActiveCompactions(client.instanceOperations()).stream() .anyMatch(ac -> ac.getPausedCount() > 0)); diff --git a/test/src/main/java/org/apache/accumulo/test/functional/MemoryStarvedScanIT.java b/test/src/main/java/org/apache/accumulo/test/functional/MemoryStarvedScanIT.java index 14b9ac16a44..496fdea62e4 100644 --- a/test/src/main/java/org/apache/accumulo/test/functional/MemoryStarvedScanIT.java +++ b/test/src/main/java/org/apache/accumulo/test/functional/MemoryStarvedScanIT.java @@ -256,11 +256,8 @@ public void testScanPauses() throws Exception { LOG.info("Waiting for memory to be consumed"); // Wait until the dataConsumingScanner has started fetching data - int currentCount = fetched.get(); - while (currentCount == 0) { - Thread.sleep(500); - currentCount = fetched.get(); - } + waitFor(() -> fetched.get() > 0, MINUTES.toMillis(5), 200, + "Data-consuming scan did not fetch any rows"); // This should block until the LowMemoryDetector runs and notices that the // VM is low on memory. @@ -268,7 +265,7 @@ public void testScanPauses() throws Exception { assertTrue(consumingIter.hasNext()); // Confirm that some data was fetched by the memoryConsumingScanner - currentCount = fetched.get(); + int currentCount = fetched.get(); assertTrue(currentCount > 0 && currentCount < 100); LOG.info("Memory consumed after reading {} rows", currentCount); @@ -382,10 +379,13 @@ public void testBatchScanPauses() throws Exception { waitFor(() -> fetched.get() > 0, MINUTES.toMillis(5), 200); // Make sure memory is free before trying to consume memory - while (LOW_MEM_DETECTED.get() == 1) { - freeServerMemory(client); - Thread.sleep(5000); - } + waitFor(() -> { + if (LOW_MEM_DETECTED.get() == 1) { + freeServerMemory(client); + } + return LOW_MEM_DETECTED.get() == 0; + }, MINUTES.toMillis(2), SECONDS.toMillis(5), + "Tablet server low-memory condition did not clear"); // This should block until the LowMemoryDetector runs and notices that the // VM is low on memory. @@ -491,11 +491,8 @@ public void testLowMemoryFlapping() throws Exception { t.start(); // Wait until the dataConsumingScanner has started fetching data - int currentCount = fetched.get(); - while (currentCount == 0) { - Thread.sleep(500); - currentCount = fetched.get(); - } + waitFor(() -> fetched.get() > 0, MINUTES.toMillis(5), 200, + "Data-consuming scan did not fetch any rows"); // This should block until the LowMemoryDetector runs and notices that the // VM is low on memory. @@ -503,7 +500,7 @@ public void testLowMemoryFlapping() throws Exception { assertTrue(consumingIter.hasNext()); // Confirm that some data was fetched by the dataConsumingScanner - currentCount = fetched.get(); + int currentCount = fetched.get(); assertTrue(currentCount > 0 && currentCount < 100); // Grab the current paused count, wait two seconds and then confirm that diff --git a/test/src/main/java/org/apache/accumulo/test/functional/MetadataMaxFilesIT.java b/test/src/main/java/org/apache/accumulo/test/functional/MetadataMaxFilesIT.java index d1b7afd2495..c31e1895f07 100644 --- a/test/src/main/java/org/apache/accumulo/test/functional/MetadataMaxFilesIT.java +++ b/test/src/main/java/org/apache/accumulo/test/functional/MetadataMaxFilesIT.java @@ -32,6 +32,7 @@ import org.apache.accumulo.core.data.TableId; import org.apache.accumulo.core.metadata.SystemTables; import org.apache.accumulo.miniclusterImpl.MiniAccumuloConfigImpl; +import org.apache.accumulo.test.util.Wait; import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.fs.RawLocalFileSystem; import org.apache.hadoop.io.Text; @@ -76,15 +77,12 @@ public void test() throws Exception { TableId tid0 = TableId.of(c.tableOperations().tableIdMap().get("table0")); TableId tid1 = TableId.of(c.tableOperations().tableIdMap().get("table1")); - while (true) { + Wait.waitFor(() -> { long hostedTabletCount = ManagerAssignmentIT.countTabletsWithLocation(c, tid0) + ManagerAssignmentIT.countTabletsWithLocation(c, tid1); log.info("Online tablets " + hostedTabletCount); - if (hostedTabletCount == 2002) { - break; - } - Thread.sleep(SECONDS.toMillis(1)); - } + return hostedTabletCount == 2002; + }, 30_000, SECONDS.toMillis(1), "Hosted tablet count did not reach 2002"); } } } diff --git a/test/src/main/java/org/apache/accumulo/test/functional/PerTableCryptoIT.java b/test/src/main/java/org/apache/accumulo/test/functional/PerTableCryptoIT.java index 7a4085680d8..e6ccd549af6 100644 --- a/test/src/main/java/org/apache/accumulo/test/functional/PerTableCryptoIT.java +++ b/test/src/main/java/org/apache/accumulo/test/functional/PerTableCryptoIT.java @@ -54,6 +54,7 @@ import org.apache.accumulo.test.TestIngest; import org.apache.accumulo.test.VerifyIngest; import org.apache.accumulo.test.harness.AccumuloClusterHarness; +import org.apache.accumulo.test.util.Wait; import org.apache.accumulo.tserver.logger.LogReader; import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.fs.FSDataOutputStream; @@ -218,30 +219,24 @@ private void checkTableEncryption(TableId tableId, boolean expectEncrypt) throws private void checkWALEncryption() throws Exception { Set walsSeen = new HashSet<>(); - int open = 0; - int attempts = 0; - boolean foundWal = false; - while (open == 0) { - attempts++; + int[] attempts = {0}; + Wait.waitFor(() -> { + attempts[0]++; + walsSeen.clear(); Map wals = WALSunnyDayIT._getWals(getServerContext()); for (var entry : wals.entrySet()) { if (entry.getValue() == WalStateManager.WalState.OPEN) { - open++; walsSeen.add(entry.getKey()); - foundWal = true; } else { // log CLOSED or UNREFERENCED to help debug this test log.debug("The WalState for {} is {}", entry.getKey(), entry.getValue()); } } - - if (!foundWal) { - Thread.sleep(50); - if (attempts % 50 == 0) { - log.debug("No open WALs found in {} attempts.", attempts); - } + if (walsSeen.isEmpty() && attempts[0] % 50 == 0) { + log.debug("No open WALs found in {} attempts.", attempts[0]); } - } + return !walsSeen.isEmpty(); + }, 30_000, 50, "No open WALs found"); assertFalse(walsSeen.isEmpty(), "Did not see any WALs"); diff --git a/test/src/main/java/org/apache/accumulo/test/functional/RegexGroupBalanceIT.java b/test/src/main/java/org/apache/accumulo/test/functional/RegexGroupBalanceIT.java index e87450211ac..e3705b1fe81 100644 --- a/test/src/main/java/org/apache/accumulo/test/functional/RegexGroupBalanceIT.java +++ b/test/src/main/java/org/apache/accumulo/test/functional/RegexGroupBalanceIT.java @@ -44,6 +44,7 @@ import org.apache.accumulo.core.security.Authorizations; import org.apache.accumulo.core.spi.balancer.RegexGroupBalancer; import org.apache.accumulo.miniclusterImpl.MiniAccumuloConfigImpl; +import org.apache.accumulo.test.util.Wait; import org.apache.commons.lang3.mutable.MutableInt; import org.apache.hadoop.io.Text; import org.junit.jupiter.api.Test; @@ -92,9 +93,7 @@ public void testBalancing() throws Exception { client.tableOperations().create(tablename, new NewTableConfiguration().setProperties(props) .withSplits(splits).withInitialTabletAvailability(TabletAvailability.HOSTED)); - while (true) { - Thread.sleep(250); - + Wait.waitFor(() -> { Table groupLocationCounts = getCounts(client, tablename); boolean allGood = true; @@ -102,11 +101,11 @@ public void testBalancing() throws Exception { allGood &= checkGroup(groupLocationCounts, "02", 1, 1, 4); allGood &= checkGroup(groupLocationCounts, "03", 1, 2, 4); allGood &= checkTabletsPerTserver(groupLocationCounts, 3, 3, 4); - - if (allGood) { - break; - } - } + return allGood; + }, 30_000, 250, + "Expected initial distribution was not reached (group 01: 1 tablet on 3 servers, " + + "group 02: 1 tablet on 4 servers, group 03: 1-2 tablets on 4 servers, " + + "3 tablets per server)"); splits.clear(); splits.add(new Text("01b")); @@ -115,9 +114,7 @@ public void testBalancing() throws Exception { splits.add(new Text("01r")); client.tableOperations().addSplits(tablename, splits); - while (true) { - Thread.sleep(250); - + Wait.waitFor(() -> { Table groupLocationCounts = getCounts(client, tablename); boolean allGood = true; @@ -125,18 +122,16 @@ public void testBalancing() throws Exception { allGood &= checkGroup(groupLocationCounts, "02", 1, 1, 4); allGood &= checkGroup(groupLocationCounts, "03", 1, 2, 4); allGood &= checkTabletsPerTserver(groupLocationCounts, 4, 4, 4); - - if (allGood) { - break; - } - } + return allGood; + }, 30_000, 250, + "Expected distribution after adding splits was not reached (group 01: 1-2 tablets on " + + "4 servers, group 02: 1 tablet on 4 servers, group 03: 1-2 tablets on 4 servers, " + + "4 tablets per server)"); // merge group 01 down to one tablet client.tableOperations().merge(tablename, null, new Text("01z")); - while (true) { - Thread.sleep(250); - + Wait.waitFor(() -> { Table groupLocationCounts = getCounts(client, tablename); boolean allGood = true; @@ -144,11 +139,11 @@ public void testBalancing() throws Exception { allGood &= checkGroup(groupLocationCounts, "02", 1, 1, 4); allGood &= checkGroup(groupLocationCounts, "03", 1, 2, 4); allGood &= checkTabletsPerTserver(groupLocationCounts, 2, 3, 4); - - if (allGood) { - break; - } - } + return allGood; + }, 30_000, 250, + "Expected distribution after merging group 01 was not reached (group 01: 1 tablet on " + + "1 server, group 02: 1 tablet on 4 servers, group 03: 1-2 tablets on 4 servers, " + + "2-3 tablets per server)"); } } diff --git a/test/src/main/java/org/apache/accumulo/test/functional/ScanIdIT.java b/test/src/main/java/org/apache/accumulo/test/functional/ScanIdIT.java index 198c94e2940..27b342ae81a 100644 --- a/test/src/main/java/org/apache/accumulo/test/functional/ScanIdIT.java +++ b/test/src/main/java/org/apache/accumulo/test/functional/ScanIdIT.java @@ -154,22 +154,20 @@ public void testScanId() throws Exception { } // wait for scanners to report a result. - while (testInProgress.get()) { - + Wait.waitFor(() -> { if (resultsByWorker.size() < NUM_TOTAL_SCANNERS) { log.trace("Results reported {}", resultsByWorker.size()); - Thread.sleep(750); - } else { - // each worker has reported at least one result. - testInProgress.set(false); - - log.debug("Final result count {}", resultsByWorker.size()); - - // delay to allow scanners to react to end of test and cleanly close. - Thread.sleep(SECONDS.toMillis(1)); + return false; } + return true; + }, 30_000, 100, "Not all scanner workers reported a result"); + // each worker has reported at least one result. + testInProgress.set(false); - } + log.debug("Final result count {}", resultsByWorker.size()); + + // delay to allow scanners to react to end of test and cleanly close. + Thread.sleep(SECONDS.toMillis(1)); Set scanIds = getScanIds(client); assertTrue(scanIds.size() >= NUM_TOTAL_SCANNERS, diff --git a/test/src/main/java/org/apache/accumulo/test/functional/SplitIT.java b/test/src/main/java/org/apache/accumulo/test/functional/SplitIT.java index ac6065ff705..8daa4655685 100644 --- a/test/src/main/java/org/apache/accumulo/test/functional/SplitIT.java +++ b/test/src/main/java/org/apache/accumulo/test/functional/SplitIT.java @@ -218,9 +218,8 @@ public void tabletShouldSplit() throws Exception { VerifyParams params = new VerifyParams(getClientProps(), table, 100_000); TestIngest.ingest(c, params); VerifyIngest.verifyIngest(c, params); - while (c.tableOperations().listSplits(table).size() < 10) { - Thread.sleep(SECONDS.toMillis(15)); - } + Wait.waitFor(() -> c.tableOperations().listSplits(table).size() >= 10, 30_000, 100, + "Expected table splits were not created"); TableId id = TableId.of(c.tableOperations().tableIdMap().get(table)); try (Scanner s = c.createScanner(SystemTables.METADATA.tableName(), Authorizations.EMPTY)) { KeyExtent extent = new KeyExtent(id, null, null); @@ -266,12 +265,9 @@ public void interleaveSplit() throws Exception { Thread.sleep(SECONDS.toMillis(5)); ReadWriteIT.interleaveTest(c, tableName); Thread.sleep(SECONDS.toMillis(5)); + Wait.waitFor(() -> c.tableOperations().listSplits(tableName).size() > 20, 30_000, 100, + "Expected at least 20 splits"); int numSplits = c.tableOperations().listSplits(tableName).size(); - while (numSplits <= 20) { - log.info("Waiting for splits to happen"); - Thread.sleep(2000); - numSplits = c.tableOperations().listSplits(tableName).size(); - } assertTrue(numSplits > 20, "Expected at least 20 splits, saw " + numSplits); } } @@ -285,12 +281,8 @@ public void deleteSplit() throws Exception { "10K", Property.TABLE_FILE_COMPRESSED_BLOCK_SIZE.getKey(), "1K"))); DeleteIT.deleteTest(c, getCluster(), tableName); c.tableOperations().flush(tableName, null, null, true); - for (int i = 0; i < 5; i++) { - Thread.sleep(SECONDS.toMillis(10)); - if (c.tableOperations().listSplits(tableName).size() > 20) { - break; - } - } + Wait.waitFor(() -> c.tableOperations().listSplits(tableName).size() > 20, 50_000, 250, + "Expected more than 20 splits"); assertTrue(c.tableOperations().listSplits(tableName).size() > 20); } } diff --git a/test/src/main/java/org/apache/accumulo/test/functional/SummaryIT.java b/test/src/main/java/org/apache/accumulo/test/functional/SummaryIT.java index e37fdd61b49..f25b2ea615c 100644 --- a/test/src/main/java/org/apache/accumulo/test/functional/SummaryIT.java +++ b/test/src/main/java/org/apache/accumulo/test/functional/SummaryIT.java @@ -88,6 +88,7 @@ import org.apache.accumulo.test.harness.MiniClusterConfigurationCallback; import org.apache.accumulo.test.harness.SharedMiniClusterBase; import org.apache.accumulo.test.util.FileMetadataUtil; +import org.apache.accumulo.test.util.Wait; import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.io.Text; import org.junit.jupiter.api.AfterAll; @@ -605,19 +606,18 @@ public void testPermissions() throws Exception { c.securityOperations().grantTablePermission("user1", table, TablePermission.GET_SUMMARIES); - int tries = 0; - while (tries < 10) { + Summary[] summary = new Summary[1]; + Wait.waitFor(() -> { try { - Summary summary = c2.tableOperations().summaries(table).retrieve().get(0); - assertEquals(2, summary.getStatistics().size()); - assertEquals(2L, (long) summary.getStatistics().getOrDefault("bars", 0L)); - assertEquals(1L, (long) summary.getStatistics().getOrDefault("foos", 0L)); - break; + summary[0] = c2.tableOperations().summaries(table).retrieve().get(0); + return true; } catch (AccumuloSecurityException ase) { - Thread.sleep(500); - tries++; + return false; } - } + }, 5_000, 500, "Timed out waiting for GET_SUMMARIES permission to take effect"); + assertEquals(2, summary[0].getStatistics().size()); + assertEquals(2L, (long) summary[0].getStatistics().getOrDefault("bars", 0L)); + assertEquals(1L, (long) summary[0].getStatistics().getOrDefault("foos", 0L)); } } } diff --git a/test/src/main/java/org/apache/accumulo/test/functional/TabletManagementIteratorIT.java b/test/src/main/java/org/apache/accumulo/test/functional/TabletManagementIteratorIT.java index 8936f6e80bd..1f9a1145b20 100644 --- a/test/src/main/java/org/apache/accumulo/test/functional/TabletManagementIteratorIT.java +++ b/test/src/main/java/org/apache/accumulo/test/functional/TabletManagementIteratorIT.java @@ -96,7 +96,6 @@ import org.apache.accumulo.core.metadata.schema.TabletsMetadata; import org.apache.accumulo.core.security.Authorizations; import org.apache.accumulo.core.tabletserver.log.LogEntry; -import org.apache.accumulo.core.util.UtilWaitThread; import org.apache.accumulo.core.util.time.SteadyTime; import org.apache.accumulo.miniclusterImpl.MiniAccumuloConfigImpl; import org.apache.accumulo.server.manager.LiveTServerSet; @@ -104,6 +103,7 @@ import org.apache.accumulo.server.manager.state.TabletManagementParameters; import org.apache.accumulo.server.metadata.TabletsMutatorImpl; import org.apache.accumulo.test.harness.AccumuloClusterHarness; +import org.apache.accumulo.test.util.Wait; import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.fs.Path; import org.apache.hadoop.io.Text; @@ -191,14 +191,15 @@ public void test(String compressionType) throws Exception { Map> expected; TabletManagementParameters tabletMgmtParams = createParameters(client); - Map> tabletsInFlux = - findTabletsNeedingAttention(client, metaCopy1, tabletMgmtParams, false); - while (!tabletsInFlux.isEmpty()) { - log.debug("Waiting for {} tablets for {}", tabletsInFlux, metaCopy1); - UtilWaitThread.sleep(500); + Wait.waitFor(() -> { copyTable(client, SystemTables.METADATA.tableName(), metaCopy1); - tabletsInFlux = findTabletsNeedingAttention(client, metaCopy1, tabletMgmtParams, false); - } + Map> tabletsInFlux = + findTabletsNeedingAttention(client, metaCopy1, tabletMgmtParams, false); + if (!tabletsInFlux.isEmpty()) { + log.debug("Waiting for {} tablets for {}", tabletsInFlux, metaCopy1); + } + return tabletsInFlux.isEmpty(); + }, 30_000, 250, "Tablets remained in flux"); expected = Map.of(); assertEquals(expected, findTabletsNeedingAttention(client, metaCopy1, tabletMgmtParams, false), diff --git a/test/src/main/java/org/apache/accumulo/test/functional/TabletResourceGroupBalanceIT.java b/test/src/main/java/org/apache/accumulo/test/functional/TabletResourceGroupBalanceIT.java index 7dca8195189..ca3470ff5d2 100644 --- a/test/src/main/java/org/apache/accumulo/test/functional/TabletResourceGroupBalanceIT.java +++ b/test/src/main/java/org/apache/accumulo/test/functional/TabletResourceGroupBalanceIT.java @@ -301,14 +301,20 @@ public void testResourceGroupPropertyChange(AccumuloClient client, String tableN locations = getLocations(ample, tableId); // wait for GROUP1 to show up in the list of locations as the current location - while ((locations == null || locations.isEmpty() || locations.size() != numExpectedSplits - || locations.get(0).getLocation() == null - || locations.get(0).getLocation().getType() == LocationType.FUTURE) - || (locations.get(0).getLocation().getType() == LocationType.CURRENT - && !tserverGroups.get(locations.get(0).getLocation().getHostAndPort().toString()) - .equals(ResourceGroupId.of("GROUP1")))) { - locations = getLocations(ample, tableId); - } + AtomicReference> locationsRef = new AtomicReference<>(locations); + Wait.waitFor(() -> { + List currentLocations = getLocations(ample, tableId); + locationsRef.set(currentLocations); + return !((currentLocations == null || currentLocations.isEmpty() + || currentLocations.size() != numExpectedSplits + || currentLocations.get(0).getLocation() == null + || currentLocations.get(0).getLocation().getType() == LocationType.FUTURE) + || (currentLocations.get(0).getLocation().getType() == LocationType.CURRENT + && !tserverGroups + .get(currentLocations.get(0).getLocation().getHostAndPort().toString()) + .equals(ResourceGroupId.of("GROUP1")))); + }, 30_000, 250, "Tablets were not assigned to GROUP1"); + locations = locationsRef.get(); Location group1Location = locations.get(0).getLocation(); assertTrue(tserverGroups.get(group1Location.getHostAndPort().toString()) .equals(ResourceGroupId.of("GROUP1"))); @@ -317,9 +323,14 @@ public void testResourceGroupPropertyChange(AccumuloClient client, String tableN // validate that all tablets have the same location as the first tablet locations = getLocations(ample, tableId); - while (locations == null || locations.isEmpty() || locations.size() != numExpectedSplits) { - locations = getLocations(ample, tableId); - } + locationsRef.set(locations); + Wait.waitFor(() -> { + List currentLocations = getLocations(ample, tableId); + locationsRef.set(currentLocations); + return currentLocations != null && !currentLocations.isEmpty() + && currentLocations.size() == numExpectedSplits; + }, 30_000, 250, "Tablets did not reach the expected count"); + locations = locationsRef.get(); if (locations.stream().map(TabletMetadata::getLocation) .allMatch((l) -> group1Location.equals(l))) { LOG.info("Group1 location: {} matches all tablet locations: {}", group1Location, diff --git a/test/src/main/java/org/apache/accumulo/test/functional/WALSunnyDayITBase.java b/test/src/main/java/org/apache/accumulo/test/functional/WALSunnyDayITBase.java index e197c49002d..1f1692bfdf3 100644 --- a/test/src/main/java/org/apache/accumulo/test/functional/WALSunnyDayITBase.java +++ b/test/src/main/java/org/apache/accumulo/test/functional/WALSunnyDayITBase.java @@ -30,7 +30,6 @@ import static org.apache.accumulo.minicluster.ServerType.TABLET_SERVER; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertTrue; -import static org.junit.jupiter.api.Assertions.fail; import java.util.ArrayList; import java.util.Collection; @@ -39,6 +38,7 @@ import java.util.List; import java.util.Map; import java.util.Map.Entry; +import java.util.concurrent.atomic.AtomicReference; import org.apache.accumulo.core.client.Accumulo; import org.apache.accumulo.core.client.AccumuloClient; @@ -239,27 +239,17 @@ private Map> getRecoveryMarkers(AccumuloClient c) throws private Map getWALsAndAssertCount(ServerContext c, int expectedCount) throws Exception { // see https://issues.apache.org/jira/browse/ACCUMULO-4110. Sometimes this test counts the logs - // before - // the new standby log is actually ready. So let's try a few times before failing, returning the - // last - // wals variable with the the correct count. - Map wals = _getWals(c); - if (wals.size() == expectedCount) { - return wals; - } - - int waitLonger = Wait.getTimeoutFactor(e -> 1); // default to 1 - for (int i = 1; i <= TIMES_TO_COUNT; i++) { - Thread.sleep(i * PAUSE_BETWEEN_COUNTS * waitLonger); - wals = _getWals(c); - if (wals.size() == expectedCount) { - return wals; - } - } - - fail( - "Unable to get the correct number of WALs, expected " + expectedCount + " but got " + wals); - return new HashMap<>(); + // before the new standby log is actually ready, so wait for the expected count. The timeout is + // the sum of the original progressive retry delays; Wait applies the configured timeout factor. + long waitMillis = (long) TIMES_TO_COUNT * (TIMES_TO_COUNT + 1) / 2 * PAUSE_BETWEEN_COUNTS; + long retryMillis = (long) (TIMES_TO_COUNT + 1) * PAUSE_BETWEEN_COUNTS / 2; + AtomicReference> result = new AtomicReference<>(); + Wait.waitFor(() -> { + result.set(_getWals(c)); + return result.get().size() == expectedCount; + }, waitMillis, retryMillis, + "Unable to get the correct number of WALs, expected " + expectedCount); + return result.get(); } static Map _getWals(ServerContext c) throws Exception { diff --git a/test/src/main/java/org/apache/accumulo/test/lock/ServiceLockIT.java b/test/src/main/java/org/apache/accumulo/test/lock/ServiceLockIT.java index 1fa3ee0478f..d636b02c466 100644 --- a/test/src/main/java/org/apache/accumulo/test/lock/ServiceLockIT.java +++ b/test/src/main/java/org/apache/accumulo/test/lock/ServiceLockIT.java @@ -36,7 +36,6 @@ import java.util.UUID; import java.util.concurrent.CountDownLatch; import java.util.concurrent.atomic.AtomicInteger; -import java.util.concurrent.locks.LockSupport; import org.apache.accumulo.core.data.ResourceGroupId; import org.apache.accumulo.core.lock.ServiceLock; @@ -483,9 +482,7 @@ public void testLockSerial() throws Exception { assertFalse(zl1.verifyLockAtSource()); zk1.close(); - while (!zlw2.isLockHeld()) { - LockSupport.parkNanos(50); - } + Wait.waitFor(zlw2::isLockHeld, 5_000, 50, "Second lock worker did not acquire the lock"); assertTrue(zlw2.isLockHeld()); assertTrue(zl2.verifyLockAtSource()); diff --git a/test/src/main/java/org/apache/accumulo/test/manager/SuspendedTabletsIT.java b/test/src/main/java/org/apache/accumulo/test/manager/SuspendedTabletsIT.java index f7f6c529168..ca3e1a31144 100644 --- a/test/src/main/java/org/apache/accumulo/test/manager/SuspendedTabletsIT.java +++ b/test/src/main/java/org/apache/accumulo/test/manager/SuspendedTabletsIT.java @@ -207,21 +207,23 @@ private void suspensionTestBody(TServerKiller serverStopper, AfterSuspendAction // Wait for all of the tablets to hosted ... log.info("Waiting on hosting and balance"); - TabletLocations ds; - for (ds = TabletLocations.retrieve(ctx, tableName); ds.hostedCount != TABLETS; - ds = TabletLocations.retrieve(ctx, tableName)) { - Thread.sleep(1000); - } + TabletLocations[] observed = {null}; + Wait.waitFor(() -> { + observed[0] = TabletLocations.retrieve(ctx, tableName); + return observed[0].hostedCount == TABLETS; + }, 30_000, 250, "Tablets were not hosted"); + TabletLocations ds = observed[0]; log.info("Tablets hosted"); // ... and balanced. ctx.instanceOperations().waitForBalance(); log.info("Tablets balanced."); - do { - // Keep checking until all tablets are hosted and spread out across the tablet servers - Thread.sleep(1000); - ds = TabletLocations.retrieve(ctx, tableName); - } while (ds.hostedCount != TABLETS || ds.hosted.keySet().size() != (TSERVERS - 1)); + Wait.waitFor(() -> { + observed[0] = TabletLocations.retrieve(ctx, tableName); + return observed[0].hostedCount == TABLETS + && observed[0].hosted.keySet().size() == (TSERVERS - 1); + }, 30_000, 250, "Tablets were not balanced across the tablet servers"); + ds = observed[0]; // Given the loop exit condition above, at this point we're sure that all tablets are hosted // and some are hosted on each of the tablet servers other than the one reserved for hosting @@ -237,11 +239,12 @@ private void suspensionTestBody(TServerKiller serverStopper, AfterSuspendAction // All tablets should be either hosted or suspended. log.info("Waiting on suspended tablets"); - do { - Thread.sleep(1000); - ds = TabletLocations.retrieve(ctx, tableName); - } while (ds.suspended.keySet().size() != (TSERVERS - 1) - || (ds.suspendedCount + ds.hostedCount) != TABLETS); + Wait.waitFor(() -> { + observed[0] = TabletLocations.retrieve(ctx, tableName); + return observed[0].suspended.keySet().size() == (TSERVERS - 1) + && (observed[0].suspendedCount + observed[0].hostedCount) == TABLETS; + }, 30_000, 250, "Tablets were not suspended after tablet servers stopped"); + ds = observed[0]; SetMultimap deadTabletsByServer = ds.suspended; @@ -258,11 +261,12 @@ private void suspensionTestBody(TServerKiller serverStopper, AfterSuspendAction if (action == AfterSuspendAction.OFFLINE) { client.tableOperations().offline(tableName, true); - while (ds.suspendedCount > 0) { - Thread.sleep(1000); - ds = TabletLocations.retrieve(ctx, tableName); - log.info("Waiting for suspended {}", ds.suspended); - } + Wait.waitFor(() -> { + observed[0] = TabletLocations.retrieve(ctx, tableName); + return observed[0].suspendedCount == 0; + }, 30_000, 250, "Suspended tablets were not cleared when taking the table offline"); + ds = observed[0]; + log.info("Suspended tablets cleared after taking the table offline"); } else if (action == AfterSuspendAction.RESUME) { // Restart the first tablet server, making sure it ends up on the same port HostAndPort restartedServer = deadTabletsByServer.keySet().iterator().next(); @@ -273,18 +277,21 @@ private void suspensionTestBody(TServerKiller serverStopper, AfterSuspendAction // Eventually, the suspended tablets should be reassigned to the newly alive tserver. log.info("Awaiting tablet unsuspension for tablets belonging to " + restartedServer); - while (ds.suspended.containsKey(restartedServer) || ds.assignedCount != 0) { - Thread.sleep(1000); - ds = TabletLocations.retrieve(ctx, tableName); - } + Wait.waitFor(() -> { + observed[0] = TabletLocations.retrieve(ctx, tableName); + return !observed[0].suspended.containsKey(restartedServer) + && observed[0].assignedCount == 0; + }, 30_000, 250, "Tablets were not unsuspended after the tablet server restarted"); + ds = observed[0]; assertEquals(deadTabletsByServer.get(restartedServer), ds.hosted.get(restartedServer)); // Finally, after much longer, remaining suspended tablets should be reassigned. log.info("Awaiting tablet reassignment for remaining tablets (suspension timeout)"); - while (ds.hostedCount != TABLETS) { - Thread.sleep(1000); - ds = TabletLocations.retrieve(ctx, tableName); - } + Wait.waitFor(() -> { + observed[0] = TabletLocations.retrieve(ctx, tableName); + return observed[0].hostedCount == TABLETS; + }, 120_000, 250, "Remaining suspended tablets were not reassigned after suspension timeout"); + ds = observed[0]; // Ensure all suspension markers in the metadata table were cleared. assertTrue(ds.suspended.isEmpty()); diff --git a/test/src/main/java/org/apache/accumulo/test/shell/ShellServerIT.java b/test/src/main/java/org/apache/accumulo/test/shell/ShellServerIT.java index 4283c8d837d..7e71a32a525 100644 --- a/test/src/main/java/org/apache/accumulo/test/shell/ShellServerIT.java +++ b/test/src/main/java/org/apache/accumulo/test/shell/ShellServerIT.java @@ -518,30 +518,26 @@ public void addauths() throws Exception { final String table = getUniqueNames(1)[0]; // addauths ts.exec("createtable " + table + " -evc"); - boolean success = false; // Rely on the timeout rule in AccumuloIT - while (!success) { + Wait.waitFor(() -> { try { ts.exec("insert a b c d -l foo", false, "does not have authorization", true, new ErrorMessageCallback(getClientProps())); - success = true; + return true; } catch (AssertionError e) { - Thread.sleep(500); + return false; } - } + }, 30_000, 250, "Insert did not succeed after granting authorization"); ts.exec("addauths -s foo,bar", true); - boolean passed = false; - // Rely on the timeout rule in AccumuloIT - while (!passed) { + Wait.waitFor(() -> { try { ts.exec("getauths", true, "foo", true); ts.exec("getauths", true, "bar", true); - passed = true; + return true; } catch (AssertionError | Exception e) { - Thread.sleep(500); + return false; } - } - assertTrue(passed, "Could not successfully see updated authoriations"); + }, 30_000, 250, "Could not successfully see updated authoriations"); ts.exec("insert a b c d -l foo"); ts.exec("scan", true, "[foo]"); ts.exec("scan -s bar", true, "[foo]", false); From baa936833de9804f905b6b56ba264e33a791c6b9 Mon Sep 17 00:00:00 2001 From: Christopher Tubbs Date: Thu, 1 Oct 2026 21:03:37 -0400 Subject: [PATCH 3/3] Fix formatting (still has compile bug) --- .../accumulo/test/MultipleManagerFateIT.java | 12 +++++------ .../test/NamespacesIT_SimpleSuite.java | 4 ++-- .../test/compaction/CompactionExecutorIT.java | 8 ++++---- .../compaction/ExternalCompaction_3_IT.java | 13 ++++++------ .../BalanceAfterCommsFailureIT.java | 8 ++++---- .../test/functional/CompactionIT.java | 20 +++++++++++-------- .../test/functional/FateConcurrencyIT.java | 5 +++-- .../test/manager/SuspendedTabletsIT.java | 3 ++- 8 files changed, 39 insertions(+), 34 deletions(-) diff --git a/test/src/main/java/org/apache/accumulo/test/MultipleManagerFateIT.java b/test/src/main/java/org/apache/accumulo/test/MultipleManagerFateIT.java index f0abb25e361..3bfb65f14dc 100644 --- a/test/src/main/java/org/apache/accumulo/test/MultipleManagerFateIT.java +++ b/test/src/main/java/org/apache/accumulo/test/MultipleManagerFateIT.java @@ -263,13 +263,11 @@ private static void waitToSeeManagers(ClientContext context, int expectedManager // 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.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); var expectedServers = assistants.stream().map(ServiceLockPath::getServer) .map(HostAndPort::fromString).collect(toSet()); diff --git a/test/src/main/java/org/apache/accumulo/test/NamespacesIT_SimpleSuite.java b/test/src/main/java/org/apache/accumulo/test/NamespacesIT_SimpleSuite.java index 712d93ab800..47c663bfbdd 100644 --- a/test/src/main/java/org/apache/accumulo/test/NamespacesIT_SimpleSuite.java +++ b/test/src/main/java/org/apache/accumulo/test/NamespacesIT_SimpleSuite.java @@ -569,8 +569,8 @@ public void verifyConstraintInheritance() throws Exception { Integer[] constraintNums = new Integer[2]; // loop until constraint is seen in namespace and table (or until test times out) Wait.waitFor(() -> { - constraintNums[0] = c.namespaceOperations().listConstraints(namespace) - .get(constraintClassName); + 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"); diff --git a/test/src/main/java/org/apache/accumulo/test/compaction/CompactionExecutorIT.java b/test/src/main/java/org/apache/accumulo/test/compaction/CompactionExecutorIT.java index fb1a1e15e9b..8e9d19d532c 100644 --- a/test/src/main/java/org/apache/accumulo/test/compaction/CompactionExecutorIT.java +++ b/test/src/main/java/org/apache/accumulo/test/compaction/CompactionExecutorIT.java @@ -392,8 +392,8 @@ public void testDispatchSystem() throws Exception { addFiles(client, "dst2", 1); Wait.waitFor( - () -> getFiles(client, "dst1").size() <= 3 && getFiles(client, "dst2").size() <= 2, 30_000, - 100, "Compaction did not reduce dst1 and dst2 to the expected file counts"); + () -> getFiles(client, "dst1").size() <= 3 && getFiles(client, "dst2").size() <= 2, + 30_000, 100, "Compaction did not reduce dst1 and dst2 to the expected file counts"); assertEquals(3, getFiles(client, "dst1").size()); assertEquals(2, getFiles(client, "dst2").size()); @@ -426,8 +426,8 @@ public void testDispatchUser() throws Exception { .setExecutionHints(Map.of("compaction_type", "special"))); Wait.waitFor( - () -> getFiles(client, "dut1").size() <= 2 && getFiles(client, "dut2").size() <= 3, 30_000, - 100, "Compaction did not reduce dut1 and dut2 to the expected file counts"); + () -> getFiles(client, "dut1").size() <= 2 && getFiles(client, "dut2").size() <= 3, + 30_000, 100, "Compaction did not reduce dut1 and dut2 to the expected file counts"); assertEquals(2, getFiles(client, "dut1").size()); assertEquals(3, getFiles(client, "dut2").size()); diff --git a/test/src/main/java/org/apache/accumulo/test/compaction/ExternalCompaction_3_IT.java b/test/src/main/java/org/apache/accumulo/test/compaction/ExternalCompaction_3_IT.java index 54c3190aed8..d61c8f9663b 100644 --- a/test/src/main/java/org/apache/accumulo/test/compaction/ExternalCompaction_3_IT.java +++ b/test/src/main/java/org/apache/accumulo/test/compaction/ExternalCompaction_3_IT.java @@ -142,9 +142,9 @@ public void testMergeCancelsExternalCompaction() throws Exception { Wait.waitFor(() -> { try (TabletsMetadata tm = getCluster().getServerContext().getAmple().readTablets() .forTable(tid).fetch(ColumnType.ECOMP).build()) { - Set ecids2 = tm.stream() - .flatMap(t -> t.getExternalCompactions().keySet().stream()) - .collect(Collectors.toSet()); + Set ecids2 = + tm.stream().flatMap(t -> t.getExternalCompactions().keySet().stream()) + .collect(Collectors.toSet()); return Collections.disjoint(ecids, ecids2); } }); @@ -229,8 +229,8 @@ public void testCoordinatorRestartsDuringCompaction() throws Exception { } } - private Map getRunningCompactionInformation( - ServerContext ctx, Set ecids) { + private Map + getRunningCompactionInformation(ServerContext ctx, Set ecids) { final Map results = new HashMap<>(); @@ -254,7 +254,8 @@ private Map getRunningCompactionInfo // an actual update from the Compactor. TreeMap sorted = new TreeMap<>(tec.getUpdates()); var lastEntry = sorted.lastEntry(); - if (lastEntry.getValue().getMessage().equals(CompactionCoordinator.RESTART_UPDATE_MSG)) { + if (lastEntry.getValue().getMessage() + .equals(CompactionCoordinator.RESTART_UPDATE_MSG)) { continue; } results.put(ecid, new RunningCompactionInfo(tec)); diff --git a/test/src/main/java/org/apache/accumulo/test/functional/BalanceAfterCommsFailureIT.java b/test/src/main/java/org/apache/accumulo/test/functional/BalanceAfterCommsFailureIT.java index f5cc437424e..b857f2800c8 100644 --- a/test/src/main/java/org/apache/accumulo/test/functional/BalanceAfterCommsFailureIT.java +++ b/test/src/main/java/org/apache/accumulo/test/functional/BalanceAfterCommsFailureIT.java @@ -100,8 +100,8 @@ private void checkBalance(AccumuloClient c) throws Exception { Wait.waitFor(() -> { Map tableLocations = BalanceIT.countLocations(c, "test"); - long unassignedTablets = tableLocations.entrySet().stream() - .filter(e -> e.getKey().equals("none")).count(); + long unassignedTablets = + tableLocations.entrySet().stream().filter(e -> e.getKey().equals("none")).count(); if (unassignedTablets > 0) { log.info("Found {} unassigned tablets, sleeping 3 seconds for tablet assignment", unassignedTablets); @@ -110,8 +110,8 @@ private void checkBalance(AccumuloClient c) throws Exception { }, 30_000, 3_000, "Unassigned tablets were not assigned within 30 seconds"); Map tableLocations = BalanceIT.countLocations(c, "test"); - long unassignedTablets = tableLocations.entrySet().stream() - .filter(e -> e.getKey().equals("none")).count(); + long unassignedTablets = + tableLocations.entrySet().stream().filter(e -> e.getKey().equals("none")).count(); assertEquals(0, unassignedTablets, "Unassigned tablets were not assigned within 30 seconds"); assertNotNull(tableLocations); assertTrue(tableLocations.size() > 1, "Expected to have at least two TabletServers"); diff --git a/test/src/main/java/org/apache/accumulo/test/functional/CompactionIT.java b/test/src/main/java/org/apache/accumulo/test/functional/CompactionIT.java index 96753127994..c6fd29d03d0 100644 --- a/test/src/main/java/org/apache/accumulo/test/functional/CompactionIT.java +++ b/test/src/main/java/org/apache/accumulo/test/functional/CompactionIT.java @@ -457,10 +457,12 @@ public void testConcurrentWithIterators() throws Exception { "Concurrent compactions did not filter the expected data"); // eventually the compactions should clean up all of their metadata, wait for this to happen - Wait.waitFor(() -> countTablets(tableName, - tabletMetadata -> !tabletMetadata.getCompacted().isEmpty() - || tabletMetadata.getSelectedFiles() != null - || !tabletMetadata.getExternalCompactions().isEmpty()) == 0, + Wait.waitFor( + () -> countTablets(tableName, + tabletMetadata -> !tabletMetadata.getCompacted().isEmpty() + || tabletMetadata.getSelectedFiles() != null + || !tabletMetadata.getExternalCompactions().isEmpty()) + == 0, 30_000, 250, "Concurrent compactions did not clean up their metadata"); } } @@ -556,10 +558,12 @@ && countTablets(tableName, tabletMetadata -> !tabletMetadata.getCompacted().isEm "Concurrent compactions did not filter the expected data"); // eventually the compactions should clean up all of their metadata, wait for this to happen - Wait.waitFor(() -> countTablets(tableName, - tabletMetadata -> !tabletMetadata.getCompacted().isEmpty() - || tabletMetadata.getSelectedFiles() != null - || !tabletMetadata.getExternalCompactions().isEmpty()) == 0, + Wait.waitFor( + () -> countTablets(tableName, + tabletMetadata -> !tabletMetadata.getCompacted().isEmpty() + || tabletMetadata.getSelectedFiles() != null + || !tabletMetadata.getExternalCompactions().isEmpty()) + == 0, 30_000, 250, "Concurrent compactions did not clean up their metadata"); } } diff --git a/test/src/main/java/org/apache/accumulo/test/functional/FateConcurrencyIT.java b/test/src/main/java/org/apache/accumulo/test/functional/FateConcurrencyIT.java index f615ea055d8..8c7fa3e7619 100644 --- a/test/src/main/java/org/apache/accumulo/test/functional/FateConcurrencyIT.java +++ b/test/src/main/java/org/apache/accumulo/test/functional/FateConcurrencyIT.java @@ -209,10 +209,11 @@ private boolean findFate(String aTableName) { } return found; } catch (Exception ex) { - log.debug("Find fate failed for table name {} with exception, will retry", aTableName, ex); + log.debug("Find fate failed for table name {} with exception, will retry", aTableName, + ex); return false; } - // Keep the original five lookup attempts at 150 ms intervals. + // Keep the original five lookup attempts at 150 ms intervals. }, 600, 150, "FATE operation was not found for table " + aTableName); return true; } catch (IllegalStateException ex) { diff --git a/test/src/main/java/org/apache/accumulo/test/manager/SuspendedTabletsIT.java b/test/src/main/java/org/apache/accumulo/test/manager/SuspendedTabletsIT.java index ca3e1a31144..73fd6c8a802 100644 --- a/test/src/main/java/org/apache/accumulo/test/manager/SuspendedTabletsIT.java +++ b/test/src/main/java/org/apache/accumulo/test/manager/SuspendedTabletsIT.java @@ -290,7 +290,8 @@ private void suspensionTestBody(TServerKiller serverStopper, AfterSuspendAction Wait.waitFor(() -> { observed[0] = TabletLocations.retrieve(ctx, tableName); return observed[0].hostedCount == TABLETS; - }, 120_000, 250, "Remaining suspended tablets were not reassigned after suspension timeout"); + }, 120_000, 250, + "Remaining suspended tablets were not reassigned after suspension timeout"); ds = observed[0]; // Ensure all suspension markers in the metadata table were cleared.