diff --git a/README.md b/README.md
index 9f0eaf24..d04bfca0 100644
--- a/README.md
+++ b/README.md
@@ -57,6 +57,9 @@ Tests are run using the following scripts in `bin/`:
* `agitator` - Runs agitator
* `gcs` - Runs garbage collection simulation
* `monitor` - Runs availability monitor probe
+ * `manager-stress/managerstress.sh` - Stresses Manager and Compaction Coordinator operations
+ * `manager-stress/stresscompactor.sh` - Starts in-JVM Compactor simulators in a separate JVM
+ * `manager-stress/configure-mgr-stress-test.sh` - Prompts for settings and prints both commands
Run the scripts without arguments to view usage.
@@ -176,6 +179,121 @@ verification test.
See [gcs.md](docs/gcs.md).
+## Manager and Compaction Coordinator stress test
+
+The `managerstress` command starts local client JVMs that issue weighted create-table, delete-table,
+split, merge, tablet-availability, compaction, and bulk-import tasks against an existing Accumulo
+instance. It creates the requested number of uniquely named tables before starting workers and
+replenishes the table pool after deletes. A selected create task creates a new table generation
+before deleting the old generation so the pool returns to its configured size. Weights apply to
+selected workload tasks; the follow-up table creations after deletes are unweighted pool maintenance.
+
+All created table generations start with `TabletAvailability.UNHOSTED` and 10–999 evenly spaced
+initial splits, leaving at most 1,000 tablets. Split tasks add a random number of unique split rows
+without exceeding that tablet limit. Merge tasks sample a target from 1–100, clamp it to a valid
+2–100 tablet range, and use the table's actual split boundaries for the merge. Split rows and the
+bulk-import seed key use the same fixed-width hexadecimal row-key scheme.
+
+Availability tasks sample a random number of unique tablets from one through the table's full
+tablet count. Each selected tablet is assigned a random availability state different from its
+current state. Nonadjacent selections remain separate; adjacent selections with the same target
+state are grouped into one row-range request.
+
+Bulk import requires a shared HDFS directory accessible to all local workers. Each worker creates
+one one-entry RFile there at startup and stages a copy of it for every bulk import. Compaction tasks
+submit non-blocking major compactions. The coordinator removes the run-specific HDFS staging
+directory during shutdown, including any files left by workers that did not exit cleanly.
+
+`managerstress` temporarily configures the `mgrstress` compaction service in the system
+configuration. It uses `org.apache.accumulo.core.spi.compaction.RatioBasedCompactionPlanner`, with
+`compaction.service.mgrstress.planner.opts.groups` set to the group specified by
+the required `--compactor-resource-group` option. The stress namespace is created with
+`table.compaction.dispatcher.opts.service=mgrstress`. Previous ZooKeeper system-property overrides
+are restored when the test finishes; pre-existing values do not prevent startup.
+
+The required `--namespace` must name a namespace that does not already exist. `managerstress` creates
+it before creating tables, places all stress tables in it, and removes its tables and namespace when
+the run finishes. The configured Accumulo user needs permission to create and drop namespaces and
+tables, alter tables and namespaces, modify system configuration, compact, and bulk import. It also
+needs access to the configured HDFS directory.
+
+For example, the command below runs four client processes for 30 minutes, maintaining eight tables.
+Operation weights are relative; setting a weight to zero disables that operation.
+
+```bash
+./bin/manager-stress/managerstress.sh \
+ --namespace mgrstress_example \
+ --compactor-resource-group default \
+ --duration 30m \
+ --clients 4 \
+ --tables 8 \
+ --hdfs-dir hdfs://namenode:8020/tmp/accumulo-managerstress \
+ --create-weight 1 \
+ --delete-weight 1 \
+ --split-weight 2 \
+ --merge-weight 2 \
+ --availability-weight 1 \
+ --compact-weight 4 \
+ --bulk-import-weight 1
+```
+
+Use `--table-prefix` to choose a prefix for the generated tables; if omitted, a unique prefix is
+created for the run. Choose a namespace name that is unused on the Accumulo instance; the test fails
+if it already exists. The test removes all tables in its namespace and deletes the namespace after
+workers stop. Per-worker logs include operation attempts, completions, skips, races, errors, and
+average latency.
+
+When workers stop, `managerstress` logs totals aggregated from each worker's atomic metrics snapshot.
+For each operation, `submitted` counts workload operations that reached their Manager API call,
+`completed` counts successful synchronous operations, and `asyncAccepted` counts compaction and
+tablet-availability requests accepted by Accumulo. Availability and compaction requests do not wait
+for tablets to reach their target state. `failed` counts operation errors; `outcomeUnknown` is the
+subset with a transport failure or interruption after submission. Snapshot files are refreshed
+periodically in the run's temporary control directory and removed after the summary is logged.
+
+### Compactor simulator
+
+`stresscompactor` is designed to be used with the `managerstress` test framework. Run it alongside
+`managerstress` to provide Compactor simulators for external compaction requests generated by the
+stress tables. Each invocation starts the configured number of independently registered Compactor
+simulators as threads in that JVM. Start additional JVMs by invoking the command again (or on other
+hosts); set `JAVA_OPTS` separately for each process to tune its heap and garbage collector.
+
+Use the same group for `managerstress --compactor-resource-group` and
+`stresscompactor --resource-group`. Run `./bin/manager-stress/configure-mgr-stress-test.sh` to
+interactively enter settings for both processes and print the two commands. Press Enter to accept
+the displayed defaults. The client config prompt defaults to `../../conf/accumulo-client.properties`
+relative to the helper script; credentials should be stored in that file. The helper prints commands
+to stdout without starting the processes.
+
+The Manager must be able to reach the advertised host and callback ports. Configure the resource
+group to match the Compactor resource group used by the external compaction service for the stress
+tables. The configured Accumulo user must be permitted to perform system actions. When using SASL,
+the advertised host should be the host's canonical name.
+
+```bash
+JAVA_OPTS="-Xms2g -Xmx2g -XX:+UseG1GC" ./bin/manager-stress/stresscompactor.sh \
+ --compactors-per-jvm 16 \
+ --duration 30m \
+ --resource-group default \
+ --host compactor-client01.example.net \
+ --port-range 9600-9699 \
+ --success-weight 1 \
+ --failure-weight 1 \
+ --cancellation-weight 1
+```
+
+Each simulator applies `--resource-group` to its own `compactor.group` configuration property,
+registers its Compactor address, and serves the Manager's status, wake, and cancel RPCs. On receiving
+a job it immediately reports a result according to the configured relative weights. A successful
+result reports zero output entries, which commits as a zero-output compaction and removes the
+input-file references from the stress table. Simulators unregister and close their connections when
+the JVM's duration expires or it is stopped.
+
+Each simulator logs its job outcomes on shutdown, followed by a per-JVM aggregate. These counts are
+external compaction jobs, which may be multiple for one table-level compaction request; they are
+reported separately from `managerstress`'s accepted asynchronous request count.
+
## Agitator
The agitator will periodically kill the Accumulo manager, tablet server, and Hadoop data node
diff --git a/bin/manager-stress/configure-mgr-stress-test.sh b/bin/manager-stress/configure-mgr-stress-test.sh
new file mode 100755
index 00000000..176df54b
--- /dev/null
+++ b/bin/manager-stress/configure-mgr-stress-test.sh
@@ -0,0 +1,215 @@
+#! /usr/bin/env bash
+#
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements. See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership. The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License. You may obtain a copy of the License at
+#
+# https://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied. See the License for the
+# specific language governing permissions and limitations
+# under the License.
+#
+
+set -euo pipefail
+
+script_dir=$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)
+default_client_config="${script_dir}/../../conf/accumulo-client.properties"
+default_namespace="mgrstress_$(date +%Y%m%d_%H%M%S)_${RANDOM}"
+default_host=$(hostname -f 2>/dev/null || hostname)
+
+prompt_default() {
+ local label=$1
+ local default=$2
+ local answer
+ printf '%s [%s]: ' "$label" "$default" >&2
+ IFS= read -r answer
+ REPLY=${answer:-$default}
+}
+
+prompt_optional() {
+ local label=$1
+ local default_description=$2
+ printf '%s [%s]: ' "$label" "$default_description" >&2
+ IFS= read -r REPLY
+}
+
+shell_quote() {
+ local value=$1
+ value=${value//\'/\'\\\'\'}
+ printf "'%s'" "$value"
+}
+
+printf '%s\n' \
+ 'Enter the ManagerStress settings. Press Enter to accept the displayed defaults.' \
+ 'Client passwords should be stored in the selected client properties file.' >&2
+
+prompt_default 'ManagerStress duration' '5m'
+manager_duration=$REPLY
+prompt_default 'ManagerStress client JVMs' '2'
+manager_clients=$REPLY
+prompt_default 'ManagerStress table count' '4'
+manager_tables=$REPLY
+prompt_default 'Namespace (generated names must not already exist)' "$default_namespace"
+manager_namespace=$REPLY
+prompt_default 'Compactor resource group (shared with StressCompactor)' 'default'
+compactor_resource_group=$REPLY
+prompt_optional 'Table prefix (blank generates one)' 'generated automatically'
+manager_table_prefix=$REPLY
+prompt_default 'Shared HDFS directory for bulk imports' 'hdfs:///tmp/accumulo-managerstress'
+manager_hdfs_dir=$REPLY
+prompt_default 'Create-table weight' '1'
+manager_create_weight=$REPLY
+prompt_default 'Delete-table weight' '1'
+manager_delete_weight=$REPLY
+prompt_default 'Split weight' '1'
+manager_split_weight=$REPLY
+prompt_default 'Merge weight' '1'
+manager_merge_weight=$REPLY
+prompt_default 'Tablet-availability weight' '1'
+manager_availability_weight=$REPLY
+prompt_default 'Compaction weight' '1'
+manager_compact_weight=$REPLY
+prompt_default 'Bulk-import weight' '1'
+manager_bulk_import_weight=$REPLY
+prompt_optional 'ManagerStress random seed (blank uses a random seed)' 'random per run'
+manager_seed=$REPLY
+prompt_optional 'ManagerStress JAVA_OPTS (blank uses JVM defaults)' 'empty'
+manager_java_opts=$REPLY
+
+printf '\n%s\n' 'Enter the StressCompactor settings.' >&2
+prompt_default 'StressCompactor JVM duration' '5m'
+compactor_duration=$REPLY
+prompt_default 'Compactors per JVM' '1'
+compactors_per_jvm=$REPLY
+prompt_default 'Compactor callback host' "$default_host"
+compactor_host=$REPLY
+prompt_default 'Compactor callback port range' '0'
+compactor_port_range=$REPLY
+prompt_default 'Compaction success weight' '1'
+compactor_success_weight=$REPLY
+prompt_default 'Compaction failure weight' '1'
+compactor_failure_weight=$REPLY
+prompt_default 'Compaction cancellation weight' '1'
+compactor_cancellation_weight=$REPLY
+prompt_optional 'StressCompactor random seed (blank uses a random seed)' 'random per run'
+compactor_seed=$REPLY
+prompt_optional 'StressCompactor JAVA_OPTS (blank uses JVM defaults)' 'empty'
+compactor_java_opts=$REPLY
+
+printf '\n%s\n' 'Enter shared Accumulo client settings.' >&2
+prompt_default 'Accumulo client config file' "$default_client_config"
+client_config=$REPLY
+prompt_optional 'Accumulo user override (blank uses client config)' 'from client config'
+client_user=$REPLY
+prompt_optional 'Authorizations (blank uses none)' 'empty'
+client_auths=$REPLY
+prompt_default 'Enable distributed tracing? (yes/no)' 'no'
+trace_choice=${REPLY,,}
+case "$trace_choice" in
+ yes | y) enable_trace=true ;;
+ no | n) enable_trace=false ;;
+ *)
+ printf 'Please answer yes or no.\n' >&2
+ exit 2
+ ;;
+esac
+
+client_args=(--config-file "$client_config")
+if [[ -n "$client_user" ]]; then
+ client_args+=(--user "$client_user")
+fi
+if [[ -n "$client_auths" ]]; then
+ client_args+=(--auths "$client_auths")
+fi
+if [[ "$enable_trace" == true ]]; then
+ client_args+=(--trace)
+fi
+while true; do
+ prompt_optional 'Additional client property override (key=value; blank to finish)' 'none'
+ if [[ -z "$REPLY" ]]; then
+ break
+ fi
+ client_args+=(-o "$REPLY")
+done
+
+manager_args=(
+ --namespace "$manager_namespace"
+ --compactor-resource-group "$compactor_resource_group"
+ --duration "$manager_duration"
+ --clients "$manager_clients"
+ --tables "$manager_tables"
+ --hdfs-dir "$manager_hdfs_dir"
+ --create-weight "$manager_create_weight"
+ --delete-weight "$manager_delete_weight"
+ --split-weight "$manager_split_weight"
+ --merge-weight "$manager_merge_weight"
+ --availability-weight "$manager_availability_weight"
+ --compact-weight "$manager_compact_weight"
+ --bulk-import-weight "$manager_bulk_import_weight"
+)
+if [[ -n "$manager_table_prefix" ]]; then
+ manager_args+=(--table-prefix "$manager_table_prefix")
+fi
+if [[ -n "$manager_seed" ]]; then
+ manager_args+=(--seed "$manager_seed")
+fi
+manager_args+=("${client_args[@]}")
+
+compactor_args=(
+ --resource-group "$compactor_resource_group"
+ --duration "$compactor_duration"
+ --compactors-per-jvm "$compactors_per_jvm"
+ --host "$compactor_host"
+ --port-range "$compactor_port_range"
+ --success-weight "$compactor_success_weight"
+ --failure-weight "$compactor_failure_weight"
+ --cancellation-weight "$compactor_cancellation_weight"
+)
+if [[ -n "$compactor_seed" ]]; then
+ compactor_args+=(--seed "$compactor_seed")
+fi
+compactor_args+=("${client_args[@]}")
+
+print_command() {
+ local label=$1
+ local java_opts=$2
+ local script=$3
+ local argument
+ local index=0
+ shift 3
+ local -a arguments=("$@")
+
+ printf '# %s\n' "$label"
+ if [[ -n "$java_opts" ]]; then
+ printf 'JAVA_OPTS=%s ' "$(shell_quote "$java_opts")"
+ fi
+ printf '%q' "$script"
+ while ((index < ${#arguments[@]})); do
+ argument=${arguments[index]}
+ printf ' \\\n %q' "$argument"
+ if [[ "$argument" == --trace ]]; then
+ index=$((index + 1))
+ continue
+ fi
+ if ((index + 1 >= ${#arguments[@]})); then
+ printf 'Option %s is missing its value.\n' "$argument" >&2
+ return 2
+ fi
+ printf ' %q' "${arguments[index + 1]}"
+ index=$((index + 2))
+ done
+ printf '\n\n'
+}
+
+print_command 'ManagerStress' "$manager_java_opts" "${script_dir}/managerstress.sh" \
+ "${manager_args[@]}"
+print_command 'StressCompactor' "$compactor_java_opts" "${script_dir}/stresscompactor.sh" \
+ "${compactor_args[@]}"
diff --git a/bin/manager-stress/managerstress.sh b/bin/manager-stress/managerstress.sh
new file mode 100755
index 00000000..4218c5c0
--- /dev/null
+++ b/bin/manager-stress/managerstress.sh
@@ -0,0 +1,40 @@
+#! /usr/bin/env bash
+#
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements. See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership. The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License. You may obtain a copy of the License at
+#
+# https://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied. See the License for the
+# specific language governing permissions and limitations
+# under the License.
+#
+
+bin_dir=$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)
+source "${bin_dir}/../build"
+
+export CLASSPATH="$TEST_JAR_PATH:$HADOOP_API_JAR:$HADOOP_RUNTIME_JAR:$CLASSPATH"
+
+client_config_set=false
+for arg in "$@"; do
+ case "$arg" in
+ -c | --config-file | --config-file=*)
+ client_config_set=true
+ break
+ ;;
+ esac
+done
+if [[ "$client_config_set" == false ]]; then
+ set -- "$@" -c "$ACCUMULO_CLIENT_PROPS"
+fi
+
+java $JAVA_OPTS -Dlog4j.configurationFile="file:$TEST_LOG4J" \
+ org.apache.accumulo.testing.manager.stress.ManagerStress "$@"
diff --git a/bin/manager-stress/stresscompactor.sh b/bin/manager-stress/stresscompactor.sh
new file mode 100755
index 00000000..705a1f62
--- /dev/null
+++ b/bin/manager-stress/stresscompactor.sh
@@ -0,0 +1,44 @@
+#! /usr/bin/env bash
+#
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements. See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership. The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License. You may obtain a copy of the License at
+#
+# https://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied. See the License for the
+# specific language governing permissions and limitations
+# under the License.
+#
+
+bin_dir=$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)
+source "${bin_dir}/../build"
+
+export CLASSPATH="$ACCUMULO_HOME/conf:$ACCUMULO_HOME/lib/*:$TEST_JAR_PATH:$HADOOP_API_JAR:$HADOOP_RUNTIME_JAR:$CLASSPATH"
+
+if [[ "$#" -eq 0 ]]; then
+ set -- --help
+fi
+
+client_config_set=false
+for arg in "$@"; do
+ case "$arg" in
+ -c | --config-file | --config-file=*)
+ client_config_set=true
+ break
+ ;;
+ esac
+done
+if [[ "$client_config_set" == false ]]; then
+ set -- "$@" -c "$ACCUMULO_CLIENT_PROPS"
+fi
+
+java $JAVA_OPTS -Dlog4j.configurationFile="file:$TEST_LOG4J" \
+ org.apache.accumulo.testing.manager.stress.compactor.StressCompactor "$@"
diff --git a/pom.xml b/pom.xml
index c27351cd..ece487b5 100644
--- a/pom.xml
+++ b/pom.xml
@@ -105,6 +105,11 @@
org.slf4j
slf4j-api
+
+ org.apache.accumulo
+ accumulo-server-base
+ provided
+
org.apache.accumulo
accumulo-minicluster
@@ -156,6 +161,8 @@
SCRIPT_STYLE
SCRIPT_STYLE
SCRIPT_STYLE
+ SCRIPT_STYLE
+ SCRIPT_STYLE
SCRIPT_STYLE
SCRIPT_STYLE
SCRIPT_STYLE
@@ -165,6 +172,7 @@
SCRIPT_STYLE
SCRIPT_STYLE
SLASHSTAR_STYLE
+ SCRIPT_STYLE
SCRIPT_STYLE
SCRIPT_STYLE
diff --git a/src/build/checkstyle/import-control.xml b/src/build/checkstyle/import-control.xml
index 43008f5c..b602a2a3 100644
--- a/src/build/checkstyle/import-control.xml
+++ b/src/build/checkstyle/import-control.xml
@@ -30,6 +30,7 @@
+
@@ -41,6 +42,24 @@
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
diff --git a/src/main/java/org/apache/accumulo/testing/manager/stress/ManagerStress.java b/src/main/java/org/apache/accumulo/testing/manager/stress/ManagerStress.java
new file mode 100644
index 00000000..57be1341
--- /dev/null
+++ b/src/main/java/org/apache/accumulo/testing/manager/stress/ManagerStress.java
@@ -0,0 +1,457 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * https://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+package org.apache.accumulo.testing.manager.stress;
+
+import static java.nio.charset.StandardCharsets.UTF_8;
+import static java.nio.file.StandardCopyOption.ATOMIC_MOVE;
+
+import java.io.IOException;
+import java.nio.file.Files;
+import java.nio.file.NoSuchFileException;
+import java.nio.file.Path;
+import java.time.Duration;
+import java.util.ArrayList;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.Objects;
+import java.util.UUID;
+import java.util.concurrent.CopyOnWriteArrayList;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
+
+import org.apache.accumulo.core.client.Accumulo;
+import org.apache.accumulo.core.client.AccumuloClient;
+import org.apache.accumulo.core.conf.Property;
+import org.apache.accumulo.core.data.ResourceGroupId;
+import org.apache.accumulo.core.spi.compaction.RatioBasedCompactionPlanner;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+/** Starts local client processes which issue Manager and Compaction Coordinator work. */
+public class ManagerStress {
+
+ private static final Logger log = LoggerFactory.getLogger(ManagerStress.class);
+ private static final Duration WORKER_STARTUP_TIMEOUT = Duration.ofMinutes(2);
+ private static final Duration WORKER_SHUTDOWN_TIMEOUT = Duration.ofSeconds(30);
+ private static final String COMPACTION_SERVICE_NAME = "mgrstress";
+
+ private record PreviousSystemProperty(boolean wasSet, String value) {
+ }
+
+ public static void main(String[] args) throws Exception {
+ ManagerStressOptions options = new ManagerStressOptions();
+ options.parseArgs(ManagerStress.class.getName(), args);
+ if (options.workerId >= 0) {
+ options.validate(false);
+ new ManagerStressWorker(options).run();
+ } else {
+ options.validate(true);
+ new ManagerStress(options, args).run();
+ }
+ }
+
+ private final ManagerStressOptions options;
+ private final String[] originalArgs;
+ private final List workers = new CopyOnWriteArrayList<>();
+ private final AtomicBoolean stopping = new AtomicBoolean();
+ private final AtomicBoolean metricsReported = new AtomicBoolean();
+ private volatile boolean namespaceCreated;
+ private volatile boolean compactionServicePropertiesConfigured;
+ private volatile Map previousCompactionServiceProperties =
+ Map.of();
+ private volatile Map desiredCompactionServiceProperties = Map.of();
+ private Path stateDir;
+ private String runId;
+
+ private ManagerStress(ManagerStressOptions options, String[] originalArgs) {
+ this.options = options;
+ this.originalArgs = originalArgs.clone();
+ }
+
+ private void run() throws Exception {
+ String tablePrefix = options.tablePrefix;
+ if (tablePrefix == null || tablePrefix.isBlank()) {
+ tablePrefix = "mgrstress_" + UUID.randomUUID().toString().replace("-", "").substring(0, 12);
+ options.tablePrefix = tablePrefix;
+ }
+ runId = UUID.randomUUID().toString().replace("-", "");
+ stateDir = Files.createTempDirectory("accumulo-manager-stress-");
+
+ Runtime.getRuntime().addShutdownHook(new Thread(() -> {
+ stopWorkers();
+ try {
+ reportRunMetrics();
+ } catch (Exception e) {
+ log.error("Could not aggregate manager-stress metrics", e);
+ }
+ try {
+ cleanupNamespace();
+ } catch (Exception e) {
+ log.error("Could not clean up manager-stress namespace {}", options.namespace, e);
+ } finally {
+ try {
+ cleanupBulkImportDirectory();
+ } finally {
+ try {
+ restoreCompactionServiceProperties();
+ } catch (Exception e) {
+ log.error("Could not restore compaction service system properties", e);
+ } finally {
+ deleteControlDirectory();
+ }
+ }
+ }
+ }, "manager-stress-shutdown"));
+
+ try {
+ try (AccumuloClient client = Accumulo.newClient().from(options.getClientProps()).build()) {
+ createNamespace(client);
+ configureCompactionService(client);
+ new TablePool(stateDir, options.namespace, tablePrefix, options.tables, options.seed)
+ .initialize(client);
+ }
+
+ log.info("Created {} stress tables with prefix {}", options.tables, tablePrefix);
+ startWorkers();
+ waitForWorkersReady();
+
+ long deadline = Math.addExact(System.currentTimeMillis(), options.duration);
+ Path startFile = stateDir.resolve("start");
+ Path tempStartFile = stateDir.resolve("start.tmp");
+ Files.writeString(tempStartFile, Long.toString(deadline), UTF_8);
+ Files.move(tempStartFile, startFile, ATOMIC_MOVE);
+ log.info("Started {} client processes for {} ms", options.clients, options.duration);
+ awaitWorkers(deadline);
+ } finally {
+ stopWorkers();
+ try {
+ reportRunMetrics();
+ } finally {
+ try {
+ cleanupNamespace();
+ } finally {
+ try {
+ cleanupBulkImportDirectory();
+ } finally {
+ try {
+ restoreCompactionServiceProperties();
+ } finally {
+ deleteControlDirectory();
+ }
+ }
+ }
+ }
+ }
+ }
+
+ private synchronized void reportRunMetrics() {
+ if (!metricsReported.compareAndSet(false, true)) {
+ return;
+ }
+ List snapshots = new ArrayList<>();
+ for (int workerId = 0; workerId < options.clients; workerId++) {
+ Path snapshotPath = ManagerStressMetrics.snapshotPath(stateDir, workerId);
+ if (Files.exists(snapshotPath)) {
+ try {
+ snapshots.add(ManagerStressMetrics.readSnapshot(snapshotPath));
+ } catch (Exception e) {
+ log.warn("Could not read metrics snapshot for worker {}", workerId, e);
+ }
+ } else {
+ log.warn("No final metrics snapshot was available for worker {}", workerId);
+ }
+ }
+
+ ManagerStressMetrics.RunSnapshot totals = ManagerStressMetrics.aggregate(snapshots);
+ log.info("ManagerStress totals from {}/{} worker snapshots:", totals.workers(),
+ options.clients);
+ for (Operation operation : Operation.values()) {
+ ManagerStressMetrics.OperationSnapshot counts = totals.operations().get(operation);
+ double averageMillis =
+ counts.attempts() == 0 ? 0 : (double) counts.totalNanos() / counts.attempts() / 1_000_000;
+ log.info(
+ "ManagerStress {}: attempts={}, submitted={}, completed={}, asyncAccepted={}, "
+ + "skipped={}, races={}, failed={}, outcomeUnknown={}, avgMs={}",
+ operation, counts.attempts(), counts.submitted(), counts.completed(),
+ counts.asyncAccepted(), counts.skipped(), counts.races(), counts.failed(),
+ counts.outcomeUnknown(), String.format("%.3f", averageMillis));
+ }
+ log.info("ManagerStress replacement-table maintenance creates={}", totals.maintenanceCreates());
+ }
+
+ private void createNamespace(AccumuloClient client) throws Exception {
+ if (client.namespaceOperations().exists(options.namespace)) {
+ throw new IllegalStateException("Refusing to use existing namespace " + options.namespace);
+ }
+ client.namespaceOperations().create(options.namespace);
+ namespaceCreated = true;
+ client.namespaceOperations().setProperty(options.namespace,
+ Property.TABLE_COMPACTION_DISPATCHER_OPTS.getKey() + "service", COMPACTION_SERVICE_NAME);
+ log.info("Created namespace {} for manager-stress run", options.namespace);
+ }
+
+ static Map compactionServiceProperties(String compactorResourceGroup) {
+ String servicePrefix = Property.COMPACTION_SERVICE_PREFIX.getKey() + COMPACTION_SERVICE_NAME;
+ String group = ResourceGroupId.of(compactorResourceGroup).canonical();
+ return Map.of(servicePrefix + ".planner", RatioBasedCompactionPlanner.class.getName(),
+ servicePrefix + ".planner.opts.groups", "[{\"group\":\"" + group + "\"}]");
+ }
+
+ private synchronized void configureCompactionService(AccumuloClient client) throws Exception {
+ Map desired = compactionServiceProperties(options.compactorResourceGroup);
+ Map previous = new HashMap<>();
+ desiredCompactionServiceProperties = desired;
+ client.instanceOperations().modifyProperties(systemProperties -> {
+ previous.clear();
+ desired.forEach((key, value) -> {
+ previous.put(key, new PreviousSystemProperty(systemProperties.containsKey(key),
+ systemProperties.get(key)));
+ systemProperties.put(key, value);
+ });
+ previousCompactionServiceProperties = Map.copyOf(previous);
+ compactionServicePropertiesConfigured = true;
+ });
+ log.info("Configured compaction service {} for resource group {}", COMPACTION_SERVICE_NAME,
+ options.compactorResourceGroup);
+ }
+
+ private synchronized void restoreCompactionServiceProperties() throws Exception {
+ if (!compactionServicePropertiesConfigured) {
+ return;
+ }
+ try (AccumuloClient client = Accumulo.newClient().from(options.getClientProps()).build()) {
+ client.instanceOperations().modifyProperties(systemProperties -> {
+ previousCompactionServiceProperties.forEach((key, previous) -> {
+ if (!Objects.equals(systemProperties.get(key),
+ desiredCompactionServiceProperties.get(key))) {
+ log.warn(
+ "System property {} changed during the ManagerStress run; preserving its value",
+ key);
+ return;
+ }
+ if (previous.wasSet()) {
+ systemProperties.put(key, previous.value());
+ } else {
+ systemProperties.remove(key);
+ }
+ });
+ });
+ compactionServicePropertiesConfigured = false;
+ log.info("Restored compaction service system properties after ManagerStress run");
+ }
+ }
+
+ private synchronized void cleanupNamespace() throws Exception {
+ if (!namespaceCreated) {
+ return;
+ }
+ try (AccumuloClient client = Accumulo.newClient().from(options.getClientProps()).build()) {
+ var tableOps = client.tableOperations();
+ String namespacePrefix = options.namespace + ".";
+ List tables = tableOps.list().stream()
+ .filter(tableName -> tableName.startsWith(namespacePrefix)).toList();
+ Exception cleanupError = null;
+ for (String table : tables) {
+ try {
+ tableOps.delete(table);
+ } catch (Exception e) {
+ if (cleanupError == null) {
+ cleanupError = new IllegalStateException(
+ "Could not remove all tables from manager-stress namespace " + options.namespace);
+ }
+ cleanupError.addSuppressed(e);
+ }
+ }
+ if (cleanupError != null) {
+ throw cleanupError;
+ }
+ client.namespaceOperations().delete(options.namespace);
+ namespaceCreated = false;
+ log.info("Deleted namespace {} after manager-stress run", options.namespace);
+ }
+ }
+
+ private void cleanupBulkImportDirectory() {
+ if (options.hdfsDir == null || options.hdfsDir.isBlank() || runId == null) {
+ return;
+ }
+ try {
+ deleteBulkImportRunDirectory(options.hdfsDir, runId);
+ log.info("Deleted HDFS bulk-import directory for run {}", runId);
+ } catch (IOException e) {
+ log.warn("Could not remove HDFS bulk-import directory for run {}", runId, e);
+ }
+ }
+
+ static void deleteBulkImportRunDirectory(String configuredDirectory, String runId)
+ throws IOException {
+ var runDirectory =
+ new org.apache.hadoop.fs.Path(new org.apache.hadoop.fs.Path(configuredDirectory), runId);
+ var configuration = new org.apache.hadoop.conf.Configuration();
+ try (var fileSystem =
+ org.apache.hadoop.fs.FileSystem.newInstance(runDirectory.toUri(), configuration)) {
+ boolean deleted = fileSystem.delete(runDirectory, true);
+ if (!deleted && fileSystem.exists(runDirectory)) {
+ throw new IOException("Could not delete HDFS bulk-import directory " + runDirectory);
+ }
+ }
+ }
+
+ private void startWorkers() throws IOException {
+ for (int i = 0; i < options.clients; i++) {
+ List command = new ArrayList<>();
+ command.add(Path.of(System.getProperty("java.home"), "bin", "java").toString());
+ String logConfig = System.getProperty("log4j.configurationFile");
+ if (logConfig != null) {
+ command.add("-Dlog4j.configurationFile=" + logConfig);
+ }
+ command.add("-cp");
+ command.add(System.getProperty("java.class.path"));
+ command.add(ManagerStress.class.getName());
+ command.addAll(workerArguments());
+ command.add("--table-prefix");
+ command.add(options.tablePrefix);
+ command.add("--internal-worker-id");
+ command.add(Integer.toString(i));
+ command.add("--internal-state-dir");
+ command.add(stateDir.toString());
+ command.add("--internal-run-id");
+ command.add(runId);
+
+ Process process = new ProcessBuilder(command).inheritIO().start();
+ workers.add(process);
+ log.info("Started worker {} with pid {}", i, process.pid());
+ }
+ }
+
+ private void waitForWorkersReady() throws Exception {
+ long deadline = System.nanoTime() + WORKER_STARTUP_TIMEOUT.toNanos();
+ while (System.nanoTime() < deadline) {
+ int ready = 0;
+ for (int i = 0; i < workers.size(); i++) {
+ if (Files.exists(stateDir.resolve("ready-" + i))) {
+ ready++;
+ } else if (!workers.get(i).isAlive()) {
+ throw new IllegalStateException("Worker " + i + " exited before startup completed");
+ }
+ }
+ if (ready == workers.size()) {
+ return;
+ }
+ Thread.sleep(100);
+ }
+ throw new IllegalStateException("Timed out waiting for client workers to initialize");
+ }
+
+ private void awaitWorkers(long deadline) throws InterruptedException {
+ long untilDeadline = deadline - System.currentTimeMillis();
+ if (untilDeadline > 0) {
+ Thread.sleep(untilDeadline);
+ }
+ for (int i = 0; i < workers.size(); i++) {
+ Process process = workers.get(i);
+ if (!process.waitFor(WORKER_SHUTDOWN_TIMEOUT.toMillis(), TimeUnit.MILLISECONDS)) {
+ log.warn("Worker {} did not stop before the shutdown timeout", i);
+ process.destroy();
+ try {
+ if (!process.waitFor(5, TimeUnit.SECONDS)) {
+ process.destroyForcibly();
+ }
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ process.destroyForcibly();
+ throw e;
+ }
+ } else if (process.exitValue() != 0) {
+ log.warn("Worker {} exited with status {}", i, process.exitValue());
+ }
+ }
+ }
+
+ private List workerArguments() {
+ List result = new ArrayList<>();
+ for (int i = 0; i < originalArgs.length; i++) {
+ String arg = originalArgs[i];
+ if (arg.equals("--table-prefix") || arg.equals("--internal-worker-id")
+ || arg.equals("--internal-state-dir") || arg.equals("--internal-run-id")
+ || arg.startsWith("--table-prefix=") || arg.startsWith("--internal-worker-id=")
+ || arg.startsWith("--internal-state-dir=") || arg.startsWith("--internal-run-id=")) {
+ if (arg.contains("=")) {
+ continue;
+ }
+ i++;
+ } else {
+ result.add(arg);
+ }
+ }
+ return result;
+ }
+
+ private synchronized void stopWorkers() {
+ if (!stopping.compareAndSet(false, true)) {
+ return;
+ }
+ for (Process process : workers) {
+ if (process.isAlive()) {
+ process.destroy();
+ }
+ }
+ for (Process process : workers) {
+ if (process.isAlive()) {
+ try {
+ if (!process.waitFor(WORKER_SHUTDOWN_TIMEOUT.toMillis(), TimeUnit.MILLISECONDS)) {
+ process.destroyForcibly();
+ process.waitFor();
+ }
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ process.destroyForcibly();
+ }
+ }
+ }
+ }
+
+ private synchronized void deleteControlDirectory() {
+ Path directory = stateDir;
+ stateDir = null;
+ if (directory == null) {
+ return;
+ }
+ try {
+ deleteRecursivelyIfExists(directory);
+ } catch (IOException e) {
+ log.warn("Could not remove temporary manager-stress directory {}", directory, e);
+ }
+ }
+
+ static void deleteRecursivelyIfExists(Path directory) throws IOException {
+ try (var paths = Files.walk(directory)) {
+ for (Path path : paths.sorted((a, b) -> b.compareTo(a)).toList()) {
+ try {
+ Files.deleteIfExists(path);
+ } catch (NoSuchFileException ignored) {
+ // The entry may have been removed by another cleanup invocation.
+ }
+ }
+ } catch (NoSuchFileException ignored) {
+ // The root directory may already have been removed.
+ }
+ }
+}
diff --git a/src/main/java/org/apache/accumulo/testing/manager/stress/ManagerStressMetrics.java b/src/main/java/org/apache/accumulo/testing/manager/stress/ManagerStressMetrics.java
new file mode 100644
index 00000000..a4aee5a0
--- /dev/null
+++ b/src/main/java/org/apache/accumulo/testing/manager/stress/ManagerStressMetrics.java
@@ -0,0 +1,181 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * https://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+package org.apache.accumulo.testing.manager.stress;
+
+import static java.nio.file.StandardCopyOption.ATOMIC_MOVE;
+import static java.nio.file.StandardCopyOption.REPLACE_EXISTING;
+
+import java.io.IOException;
+import java.io.OutputStream;
+import java.nio.file.AtomicMoveNotSupportedException;
+import java.nio.file.Files;
+import java.nio.file.Path;
+import java.util.Collection;
+import java.util.EnumMap;
+import java.util.Map;
+import java.util.Properties;
+
+/** Per-worker workload metrics and their atomic on-disk snapshots. */
+final class ManagerStressMetrics {
+
+ static final class Counts {
+ long attempts;
+ long submitted;
+ long completed;
+ long asyncAccepted;
+ long skipped;
+ long races;
+ long failed;
+ long outcomeUnknown;
+ long totalNanos;
+
+ OperationSnapshot snapshot() {
+ return new OperationSnapshot(attempts, submitted, completed, asyncAccepted, skipped, races,
+ failed, outcomeUnknown, totalNanos);
+ }
+ }
+
+ record OperationSnapshot(long attempts, long submitted, long completed, long asyncAccepted,
+ long skipped, long races, long failed, long outcomeUnknown, long totalNanos) {
+ }
+
+ record Snapshot(int workerId, long maintenanceCreates,
+ Map operations) {
+ }
+
+ record RunSnapshot(int workers, long maintenanceCreates,
+ Map operations) {
+ }
+
+ private final EnumMap counts = new EnumMap<>(Operation.class);
+ private long maintenanceCreates;
+
+ ManagerStressMetrics() {
+ for (Operation operation : Operation.values()) {
+ counts.put(operation, new Counts());
+ }
+ }
+
+ Counts counts(Operation operation) {
+ return counts.get(operation);
+ }
+
+ void addMaintenanceCreates(long creates) {
+ maintenanceCreates += creates;
+ }
+
+ long maintenanceCreates() {
+ return maintenanceCreates;
+ }
+
+ void writeSnapshot(Path stateDir, int workerId) throws IOException {
+ Properties properties = new Properties();
+ properties.setProperty("workerId", Integer.toString(workerId));
+ properties.setProperty("maintenanceCreates", Long.toString(maintenanceCreates));
+ for (Operation operation : Operation.values()) {
+ Counts operationCounts = counts.get(operation);
+ String prefix = operation.name() + ".";
+ properties.setProperty(prefix + "attempts", Long.toString(operationCounts.attempts));
+ properties.setProperty(prefix + "submitted", Long.toString(operationCounts.submitted));
+ properties.setProperty(prefix + "completed", Long.toString(operationCounts.completed));
+ properties.setProperty(prefix + "asyncAccepted",
+ Long.toString(operationCounts.asyncAccepted));
+ properties.setProperty(prefix + "skipped", Long.toString(operationCounts.skipped));
+ properties.setProperty(prefix + "races", Long.toString(operationCounts.races));
+ properties.setProperty(prefix + "failed", Long.toString(operationCounts.failed));
+ properties.setProperty(prefix + "outcomeUnknown",
+ Long.toString(operationCounts.outcomeUnknown));
+ properties.setProperty(prefix + "totalNanos", Long.toString(operationCounts.totalNanos));
+ }
+
+ Path target = snapshotPath(stateDir, workerId);
+ Path temporary = Files.createTempFile(stateDir, "worker-metrics-", ".tmp");
+ try {
+ try (OutputStream output = Files.newOutputStream(temporary)) {
+ properties.store(output, "ManagerStress worker metrics");
+ }
+ try {
+ Files.move(temporary, target, ATOMIC_MOVE, REPLACE_EXISTING);
+ } catch (AtomicMoveNotSupportedException e) {
+ Files.move(temporary, target, REPLACE_EXISTING);
+ }
+ } finally {
+ Files.deleteIfExists(temporary);
+ }
+ }
+
+ static Path snapshotPath(Path stateDir, int workerId) {
+ return stateDir.resolve("worker-metrics-" + workerId + ".properties");
+ }
+
+ static Snapshot readSnapshot(Path snapshotPath) throws IOException {
+ Properties properties = new Properties();
+ try (var input = Files.newInputStream(snapshotPath)) {
+ properties.load(input);
+ }
+ int workerId = Integer.parseInt(properties.getProperty("workerId"));
+ long maintenanceCreates = Long.parseLong(properties.getProperty("maintenanceCreates"));
+ EnumMap operations = new EnumMap<>(Operation.class);
+ for (Operation operation : Operation.values()) {
+ String prefix = operation.name() + ".";
+ operations.put(operation, new OperationSnapshot(readLong(properties, prefix + "attempts"),
+ readLong(properties, prefix + "submitted"), readLong(properties, prefix + "completed"),
+ readLong(properties, prefix + "asyncAccepted"), readLong(properties, prefix + "skipped"),
+ readLong(properties, prefix + "races"), readLong(properties, prefix + "failed"),
+ readLong(properties, prefix + "outcomeUnknown"),
+ readLong(properties, prefix + "totalNanos")));
+ }
+ return new Snapshot(workerId, maintenanceCreates, Map.copyOf(operations));
+ }
+
+ static RunSnapshot aggregate(Collection snapshots) {
+ EnumMap totals = new EnumMap<>(Operation.class);
+ for (Operation operation : Operation.values()) {
+ totals.put(operation, new Counts());
+ }
+ long maintenanceCreates = 0;
+ int workers = 0;
+ for (Snapshot snapshot : snapshots) {
+ workers++;
+ maintenanceCreates += snapshot.maintenanceCreates();
+ for (Operation operation : Operation.values()) {
+ OperationSnapshot current = snapshot.operations().get(operation);
+ Counts total = totals.get(operation);
+ total.attempts += current.attempts();
+ total.submitted += current.submitted();
+ total.completed += current.completed();
+ total.asyncAccepted += current.asyncAccepted();
+ total.skipped += current.skipped();
+ total.races += current.races();
+ total.failed += current.failed();
+ total.outcomeUnknown += current.outcomeUnknown();
+ total.totalNanos += current.totalNanos();
+ }
+ }
+ EnumMap operations = new EnumMap<>(Operation.class);
+ for (Operation operation : Operation.values()) {
+ operations.put(operation, totals.get(operation).snapshot());
+ }
+ return new RunSnapshot(workers, maintenanceCreates, Map.copyOf(operations));
+ }
+
+ private static long readLong(Properties properties, String key) {
+ return Long.parseLong(properties.getProperty(key, "0"));
+ }
+}
diff --git a/src/main/java/org/apache/accumulo/testing/manager/stress/ManagerStressOptions.java b/src/main/java/org/apache/accumulo/testing/manager/stress/ManagerStressOptions.java
new file mode 100644
index 00000000..bd90ba94
--- /dev/null
+++ b/src/main/java/org/apache/accumulo/testing/manager/stress/ManagerStressOptions.java
@@ -0,0 +1,134 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * https://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+package org.apache.accumulo.testing.manager.stress;
+
+import org.apache.accumulo.core.data.ResourceGroupId;
+import org.apache.accumulo.testing.cli.ClientOpts;
+
+import com.beust.jcommander.Parameter;
+
+class ManagerStressOptions extends ClientOpts {
+
+ @Parameter(names = "--duration", converter = TimeConverter.class,
+ description = "duration of the stress run (for example 10m or 1h)")
+ long duration = 5 * 60 * 1000L;
+
+ @Parameter(names = "--clients", description = "number of local worker JVMs")
+ int clients = 2;
+
+ @Parameter(names = "--tables", description = "target number of managed tables")
+ int tables = 4;
+
+ @Parameter(names = "--namespace", required = true,
+ description = "new namespace used for all tables created by this run")
+ String namespace;
+
+ @Parameter(names = "--compactor-resource-group", required = true,
+ description = "resource group used by the external Compactor simulators")
+ String compactorResourceGroup;
+
+ @Parameter(names = "--table-prefix", description = "prefix for tables created by this run")
+ String tablePrefix;
+
+ @Parameter(names = "--hdfs-dir", description = "shared HDFS directory for bulk import files")
+ String hdfsDir;
+
+ @Parameter(names = "--create-weight", description = "relative weight for create-table tasks")
+ double createWeight = 1;
+
+ @Parameter(names = "--delete-weight", description = "relative weight for delete-table tasks")
+ double deleteWeight = 1;
+
+ @Parameter(names = "--split-weight", description = "relative weight for split tasks")
+ double splitWeight = 1;
+
+ @Parameter(names = "--merge-weight", description = "relative weight for merge tasks")
+ double mergeWeight = 1;
+
+ @Parameter(names = "--availability-weight",
+ description = "relative weight for tablet availability changes")
+ double availabilityWeight = 1;
+
+ @Parameter(names = "--compact-weight", description = "relative weight for compaction tasks")
+ double compactWeight = 1;
+
+ @Parameter(names = "--bulk-import-weight", description = "relative weight for bulk import tasks")
+ double bulkImportWeight = 1;
+
+ @Parameter(names = "--seed", description = "random seed (defaults to a per-run random seed)")
+ long seed = System.nanoTime();
+
+ @Parameter(names = "--internal-worker-id", hidden = true)
+ int workerId = -1;
+
+ @Parameter(names = "--internal-state-dir", hidden = true)
+ String stateDir;
+
+ @Parameter(names = "--internal-run-id", hidden = true)
+ String runId;
+
+ void validate(boolean coordinator) {
+ if (duration <= 0) {
+ throw new IllegalArgumentException("--duration must be positive");
+ }
+ if (clients <= 0) {
+ throw new IllegalArgumentException("--clients must be positive");
+ }
+ if (tables <= 0) {
+ throw new IllegalArgumentException("--tables must be positive");
+ }
+ if (namespace == null || namespace.isBlank()) {
+ throw new IllegalArgumentException("--namespace is required");
+ }
+ if (!namespace.matches("[A-Za-z0-9_]+")) {
+ throw new IllegalArgumentException("--namespace may contain only letters, digits, and _");
+ }
+ if (compactorResourceGroup == null || compactorResourceGroup.isBlank()) {
+ throw new IllegalArgumentException("--compactor-resource-group is required");
+ }
+ ResourceGroupId.of(compactorResourceGroup);
+ if (coordinator && workerId >= 0) {
+ throw new IllegalArgumentException("--internal-worker-id is reserved for worker JVMs");
+ }
+ if (!coordinator && (workerId < 0 || stateDir == null || runId == null)) {
+ throw new IllegalArgumentException("Missing internal worker configuration");
+ }
+ if (tablePrefix == null || tablePrefix.isBlank()) {
+ if (!coordinator) {
+ throw new IllegalArgumentException("Worker table prefix was not supplied");
+ }
+ } else if (!tablePrefix.matches("[A-Za-z0-9_]+")) {
+ throw new IllegalArgumentException("--table-prefix may contain only letters, digits, and _");
+ }
+ double totalWeight = createWeight + deleteWeight + splitWeight + mergeWeight
+ + availabilityWeight + compactWeight + bulkImportWeight;
+ for (double weight : new double[] {createWeight, deleteWeight, splitWeight, mergeWeight,
+ availabilityWeight, compactWeight, bulkImportWeight}) {
+ if (!Double.isFinite(weight) || weight < 0) {
+ throw new IllegalArgumentException("Operation weights must be finite and non-negative");
+ }
+ }
+ if (!Double.isFinite(totalWeight) || totalWeight <= 0) {
+ throw new IllegalArgumentException("At least one operation weight must be positive");
+ }
+ if (bulkImportWeight > 0 && (hdfsDir == null || hdfsDir.isBlank())) {
+ throw new IllegalArgumentException("--hdfs-dir is required when bulk import is enabled");
+ }
+ }
+}
diff --git a/src/main/java/org/apache/accumulo/testing/manager/stress/ManagerStressRows.java b/src/main/java/org/apache/accumulo/testing/manager/stress/ManagerStressRows.java
new file mode 100644
index 00000000..03024b27
--- /dev/null
+++ b/src/main/java/org/apache/accumulo/testing/manager/stress/ManagerStressRows.java
@@ -0,0 +1,201 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * https://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+package org.apache.accumulo.testing.manager.stress;
+
+import java.math.BigInteger;
+import java.util.ArrayList;
+import java.util.Collection;
+import java.util.Collections;
+import java.util.List;
+import java.util.Random;
+import java.util.SortedSet;
+import java.util.TreeSet;
+
+import org.apache.accumulo.core.client.admin.NewTableConfiguration;
+import org.apache.accumulo.core.client.admin.TabletAvailability;
+import org.apache.accumulo.core.data.RowRange;
+import org.apache.hadoop.io.Text;
+
+/** Row keys and tablet ranges used by the Manager stress workload. */
+final class ManagerStressRows {
+
+ static final int MAX_TABLETS = 1000;
+ static final int MIN_INITIAL_SPLITS = 10;
+ static final int MAX_INITIAL_SPLITS = MAX_TABLETS - 1;
+
+ private static final BigInteger ROW_SPACE = BigInteger.ONE.shiftLeft(64);
+
+ record MergeRange(Text start, Text end, int tabletCount) {
+ }
+
+ record AvailabilityChange(int firstTablet, int tabletCount, TabletAvailability availability,
+ RowRange rowRange) {
+ }
+
+ private ManagerStressRows() {}
+
+ static NewTableConfiguration newTableConfiguration(Random random) {
+ int splitRange = MAX_INITIAL_SPLITS - MIN_INITIAL_SPLITS + 1;
+ int splitCount = MIN_INITIAL_SPLITS + random.nextInt(splitRange);
+ return new NewTableConfiguration().withInitialTabletAvailability(TabletAvailability.UNHOSTED)
+ .withSplits(initialSplits(splitCount));
+ }
+
+ static SortedSet initialSplits(int splitCount) {
+ if (splitCount < MIN_INITIAL_SPLITS || splitCount > MAX_INITIAL_SPLITS) {
+ throw new IllegalArgumentException("Initial split count must be between " + MIN_INITIAL_SPLITS
+ + " and " + MAX_INITIAL_SPLITS);
+ }
+ TreeSet splits = new TreeSet<>();
+ BigInteger denominator = BigInteger.valueOf((long) splitCount + 1);
+ for (int i = 1; i <= splitCount; i++) {
+ BigInteger row = ROW_SPACE.multiply(BigInteger.valueOf(i)).divide(denominator);
+ splits.add(row(row));
+ }
+ return splits;
+ }
+
+ static Text randomRow(Random random) {
+ String hex = Long.toUnsignedString(random.nextLong(), 16);
+ return new Text("0".repeat(16 - hex.length()) + hex);
+ }
+
+ static SortedSet additionalSplits(Collection currentSplits, Random random) {
+ TreeSet existing = new TreeSet<>(currentSplits);
+ int currentTablets = existing.size() + 1;
+ int available = MAX_TABLETS - currentTablets;
+ TreeSet additions = new TreeSet<>();
+ if (available <= 0) {
+ return additions;
+ }
+
+ int requested = 1 + random.nextInt(available);
+ while (additions.size() < requested) {
+ Text candidate = randomRow(random);
+ if (!existing.contains(candidate)) {
+ additions.add(candidate);
+ }
+ }
+ return additions;
+ }
+
+ static MergeRange selectMergeRange(Collection currentSplits, int targetTabletCount,
+ Random random) {
+ List splits = new ArrayList<>(new TreeSet<>(currentSplits));
+ int totalTablets = splits.size() + 1;
+ if (totalTablets < 2) {
+ return null;
+ }
+ int mergeTablets = Math.max(2, Math.min(Math.min(targetTabletCount, 100), totalTablets));
+ int firstTablet = random.nextInt(totalTablets - mergeTablets + 1);
+ return mergeRange(splits, firstTablet, mergeTablets);
+ }
+
+ static MergeRange mergeRange(Collection currentSplits, int firstTablet, int tabletCount) {
+ return mergeRange(new ArrayList<>(new TreeSet<>(currentSplits)), firstTablet, tabletCount);
+ }
+
+ static MergeRange mergeRange(List sortedSplits, int firstTablet, int tabletCount) {
+ int totalTablets = sortedSplits.size() + 1;
+ if (tabletCount < 2 || firstTablet < 0 || firstTablet + tabletCount > totalTablets) {
+ throw new IllegalArgumentException("Invalid tablet range");
+ }
+ Text start = firstTablet == 0 ? null : sortedSplits.get(firstTablet - 1);
+ int lastTablet = firstTablet + tabletCount - 1;
+ Text end = lastTablet == sortedSplits.size() ? null : sortedSplits.get(lastTablet);
+ return new MergeRange(start, end, tabletCount);
+ }
+
+ static List randomAvailabilityChanges(List sortedSplits,
+ List currentAvailability, Random random) {
+ int totalTablets = sortedSplits.size() + 1;
+ if (currentAvailability.size() != totalTablets) {
+ throw new IllegalArgumentException("Availability must be supplied for every tablet");
+ }
+
+ List tabletIndices = new ArrayList<>(totalTablets);
+ for (int i = 0; i < totalTablets; i++) {
+ tabletIndices.add(i);
+ }
+ Collections.shuffle(tabletIndices, random);
+ int selectedCount = 1 + random.nextInt(totalTablets);
+ tabletIndices = new ArrayList<>(tabletIndices.subList(0, selectedCount));
+ tabletIndices.sort(Integer::compareTo);
+
+ List changes = new ArrayList<>();
+ int runStart = tabletIndices.get(0);
+ int runEnd = runStart;
+ TabletAvailability runAvailability =
+ chooseDifferentAvailability(currentAvailability.get(runStart), random);
+ for (int i = 1; i < tabletIndices.size(); i++) {
+ int tablet = tabletIndices.get(i);
+ TabletAvailability nextAvailability =
+ chooseDifferentAvailability(currentAvailability.get(tablet), random);
+ if (tablet == runEnd + 1 && nextAvailability == runAvailability) {
+ runEnd = tablet;
+ } else {
+ changes.add(availabilityChange(sortedSplits, runStart, runEnd, runAvailability));
+ runStart = tablet;
+ runEnd = tablet;
+ runAvailability = nextAvailability;
+ }
+ }
+ changes.add(availabilityChange(sortedSplits, runStart, runEnd, runAvailability));
+ return List.copyOf(changes);
+ }
+
+ static RowRange tabletRowRange(List sortedSplits, int firstTablet, int lastTablet) {
+ int totalTablets = sortedSplits.size() + 1;
+ if (firstTablet < 0 || lastTablet < firstTablet || lastTablet >= totalTablets) {
+ throw new IllegalArgumentException("Invalid tablet range");
+ }
+ if (firstTablet == 0 && lastTablet == totalTablets - 1) {
+ return RowRange.all();
+ }
+ if (lastTablet == totalTablets - 1) {
+ return RowRange.greaterThan(sortedSplits.get(firstTablet - 1));
+ }
+ if (firstTablet == 0) {
+ return RowRange.atMost(sortedSplits.get(lastTablet));
+ }
+ return RowRange.openClosed(sortedSplits.get(firstTablet - 1), sortedSplits.get(lastTablet));
+ }
+
+ private static AvailabilityChange availabilityChange(List sortedSplits, int firstTablet,
+ int lastTablet, TabletAvailability availability) {
+ return new AvailabilityChange(firstTablet, lastTablet - firstTablet + 1, availability,
+ tabletRowRange(sortedSplits, firstTablet, lastTablet));
+ }
+
+ private static TabletAvailability chooseDifferentAvailability(TabletAvailability current,
+ Random random) {
+ List candidates = new ArrayList<>();
+ for (TabletAvailability availability : TabletAvailability.values()) {
+ if (availability != current) {
+ candidates.add(availability);
+ }
+ }
+ return candidates.get(random.nextInt(candidates.size()));
+ }
+
+ private static Text row(BigInteger value) {
+ String hex = value.toString(16);
+ return new Text("0".repeat(16 - hex.length()) + hex);
+ }
+}
diff --git a/src/main/java/org/apache/accumulo/testing/manager/stress/ManagerStressWorker.java b/src/main/java/org/apache/accumulo/testing/manager/stress/ManagerStressWorker.java
new file mode 100644
index 00000000..942e1922
--- /dev/null
+++ b/src/main/java/org/apache/accumulo/testing/manager/stress/ManagerStressWorker.java
@@ -0,0 +1,420 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * https://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+package org.apache.accumulo.testing.manager.stress;
+
+import static java.nio.charset.StandardCharsets.UTF_8;
+
+import java.io.IOException;
+import java.nio.file.Files;
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.List;
+import java.util.Random;
+import java.util.TreeSet;
+import java.util.UUID;
+
+import org.apache.accumulo.core.client.Accumulo;
+import org.apache.accumulo.core.client.AccumuloClient;
+import org.apache.accumulo.core.client.TableExistsException;
+import org.apache.accumulo.core.client.TableNotFoundException;
+import org.apache.accumulo.core.client.TableOfflineException;
+import org.apache.accumulo.core.client.admin.CompactionConfig;
+import org.apache.accumulo.core.client.admin.TabletAvailability;
+import org.apache.accumulo.core.client.admin.TabletInformation;
+import org.apache.accumulo.core.client.rfile.RFile;
+import org.apache.accumulo.core.client.rfile.RFileWriter;
+import org.apache.accumulo.core.data.Key;
+import org.apache.accumulo.core.data.RowRange;
+import org.apache.accumulo.core.data.Value;
+import org.apache.hadoop.conf.Configuration;
+import org.apache.hadoop.fs.FSDataInputStream;
+import org.apache.hadoop.fs.FSDataOutputStream;
+import org.apache.hadoop.fs.FileSystem;
+import org.apache.hadoop.fs.Path;
+import org.apache.hadoop.io.IOUtils;
+import org.apache.hadoop.io.Text;
+import org.apache.thrift.transport.TTransportException;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+final class ManagerStressWorker {
+
+ private static final Logger log = LoggerFactory.getLogger(ManagerStressWorker.class);
+
+ private static class SubmissionTracker {
+ private final ManagerStressMetrics.Counts counts;
+ private boolean submitted;
+
+ SubmissionTracker(ManagerStressMetrics.Counts counts) {
+ this.counts = counts;
+ }
+
+ void markSubmitted() {
+ if (!submitted) {
+ counts.submitted++;
+ submitted = true;
+ }
+ }
+ }
+
+ private final ManagerStressOptions options;
+ private final ManagerStressMetrics metrics = new ManagerStressMetrics();
+ private java.nio.file.Path stateDirectory;
+ private long lastMetricsSnapshotNanos;
+ private Path workerBulkDir;
+ private Path seedFile;
+ private FileSystem fileSystem;
+
+ ManagerStressWorker(ManagerStressOptions options) {
+ this.options = options;
+ }
+
+ void run() throws Exception {
+ stateDirectory = java.nio.file.Path.of(options.stateDir);
+ TablePool tablePool = new TablePool(stateDirectory, options.namespace, options.tablePrefix,
+ options.tables, options.seed);
+ Random random = new Random(options.seed + (long) options.workerId * 0x9e3779b97f4a7c15L);
+
+ try (AccumuloClient client = Accumulo.newClient().from(options.getClientProps()).build()) {
+ if (options.bulkImportWeight > 0) {
+ prepareBulkFile(random);
+ }
+ writeMetricsSnapshot();
+ Files.createFile(stateDirectory.resolve("ready-" + options.workerId));
+ long deadline = awaitStart(stateDirectory);
+ log.info("Worker {} started workload", options.workerId);
+
+ while (System.currentTimeMillis() < deadline) {
+ try {
+ metrics.addMaintenanceCreates(replenish(tablePool, client));
+ } catch (Exception e) {
+ log.warn("Worker {} could not replenish the table pool", options.workerId, e);
+ }
+
+ List activeSlots =
+ tablePool.snapshot().stream().filter(slot -> slot.tableName() != null).toList();
+ if (activeSlots.isEmpty()) {
+ Thread.sleep(100);
+ continue;
+ }
+
+ Operation operation = Operation.choose(options, random);
+ ManagerStressMetrics.Counts operationCounts = metrics.counts(operation);
+ operationCounts.attempts++;
+ SubmissionTracker submission = new SubmissionTracker(operationCounts);
+ long start = System.nanoTime();
+ try {
+ if (execute(operation, activeSlots, tablePool, client, random, submission)) {
+ if (operation.isAsynchronous()) {
+ operationCounts.asyncAccepted++;
+ } else {
+ operationCounts.completed++;
+ }
+ } else {
+ operationCounts.skipped++;
+ }
+ } catch (Exception e) {
+ if (isExpectedRace(e)) {
+ operationCounts.races++;
+ log.debug("Worker {} observed a table lifecycle race during {}", options.workerId,
+ operation, e);
+ } else {
+ operationCounts.failed++;
+ if (submission.submitted && isOutcomeUnknown(e)) {
+ operationCounts.outcomeUnknown++;
+ }
+ log.warn("Worker {} failed {} operation", options.workerId, operation, e);
+ }
+ } finally {
+ operationCounts.totalNanos += System.nanoTime() - start;
+ writeMetricsSnapshotIfDue();
+ }
+ }
+ try {
+ metrics.addMaintenanceCreates(replenish(tablePool, client));
+ } catch (Exception e) {
+ log.warn("Worker {} could not restore the table pool at shutdown", options.workerId, e);
+ }
+ } finally {
+ try {
+ if (fileSystem != null) {
+ try {
+ if (workerBulkDir != null) {
+ fileSystem.delete(workerBulkDir, true);
+ }
+ } catch (IOException e) {
+ log.warn("Worker {} could not remove its HDFS work directory {}", options.workerId,
+ workerBulkDir, e);
+ } finally {
+ fileSystem.close();
+ }
+ }
+ } finally {
+ report();
+ writeMetricsSnapshot();
+ }
+ }
+ }
+
+ private int replenish(TablePool pool, AccumuloClient client) throws Exception {
+ int created = pool.replenish(client);
+ if (created > 0) {
+ log.info("Worker {} replenished {} deleted table slot(s)", options.workerId, created);
+ }
+ return created;
+ }
+
+ private boolean execute(Operation operation, List activeSlots, TablePool pool,
+ AccumuloClient client, Random random, SubmissionTracker submission) throws Exception {
+ TablePool.Slot slot = activeSlots.get(random.nextInt(activeSlots.size()));
+ String tableName = slot.tableName();
+ return switch (operation) {
+ case CREATE -> {
+ if (!pool.replace(client, slot.index(), submission::markSubmitted)) {
+ throw new TableNotFoundException(null, tableName,
+ "all table slots are being updated by other workers");
+ }
+ log.debug("Worker {} created a replacement table for slot {}", options.workerId,
+ slot.index());
+ yield true;
+ }
+ case DELETE -> {
+ boolean deleted = pool.delete(client, slot.index(), submission::markSubmitted);
+ if (!deleted) {
+ throw new TableNotFoundException(null, tableName, "table was deleted by another worker");
+ }
+ log.debug("Worker {} deleted table {}", options.workerId, tableName);
+ yield true;
+ }
+ case SPLIT -> split(client, pool, slot, random, submission);
+ case MERGE -> merge(client, pool, slot, random, submission);
+ case AVAILABILITY -> changeAvailability(client, pool, slot, random, submission);
+ case COMPACT -> {
+ submission.markSubmitted();
+ client.tableOperations().compact(tableName,
+ new CompactionConfig().setFlush(true).setWait(false));
+ log.debug("Worker {} submitted a compaction for table {}", options.workerId, tableName);
+ yield true;
+ }
+ case BULK_IMPORT -> {
+ bulkImport(client, tableName, submission);
+ yield true;
+ }
+ };
+ }
+
+ private boolean split(AccumuloClient client, TablePool pool, TablePool.Slot selected,
+ Random random, SubmissionTracker submission) throws Exception {
+ TablePool.Slot reservation = pool.lockTable(selected.index());
+ if (reservation == null) {
+ return false;
+ }
+ try {
+ var tableOps = client.tableOperations();
+ var currentSplits = new TreeSet<>(tableOps.listSplits(reservation.tableName()));
+ var newSplits = ManagerStressRows.additionalSplits(currentSplits, random);
+ if (newSplits.isEmpty()) {
+ log.debug("Worker {} skipped split for {} at the {}-tablet limit", options.workerId,
+ reservation.tableName(), ManagerStressRows.MAX_TABLETS);
+ return false;
+ }
+ submission.markSubmitted();
+ tableOps.addSplits(reservation.tableName(), newSplits);
+ log.debug("Worker {} added {} split(s) to {}", options.workerId, newSplits.size(),
+ reservation.tableName());
+ return true;
+ } finally {
+ pool.unlockTable(reservation);
+ }
+ }
+
+ private boolean merge(AccumuloClient client, TablePool pool, TablePool.Slot selected,
+ Random random, SubmissionTracker submission) throws Exception {
+ TablePool.Slot reservation = pool.lockTable(selected.index());
+ if (reservation == null) {
+ return false;
+ }
+ try {
+ var tableOps = client.tableOperations();
+ var splits = tableOps.listSplits(reservation.tableName());
+ int targetTablets = random.nextInt(100) + 1;
+ ManagerStressRows.MergeRange range =
+ ManagerStressRows.selectMergeRange(splits, targetTablets, random);
+ if (range == null) {
+ log.debug("Worker {} skipped merge for {} because it has fewer than two tablets",
+ options.workerId, reservation.tableName());
+ return false;
+ }
+ submission.markSubmitted();
+ tableOps.merge(reservation.tableName(), range.start(), range.end());
+ log.debug("Worker {} merged {} tablets in {} over ({}, {}]", options.workerId,
+ range.tabletCount(), reservation.tableName(), range.start(), range.end());
+ return true;
+ } finally {
+ pool.unlockTable(reservation);
+ }
+ }
+
+ private boolean changeAvailability(AccumuloClient client, TablePool pool, TablePool.Slot selected,
+ Random random, SubmissionTracker submission) throws Exception {
+ TablePool.Slot reservation = pool.lockTable(selected.index());
+ if (reservation == null) {
+ return false;
+ }
+ try {
+ var tableOps = client.tableOperations();
+ List splits = new ArrayList<>(tableOps.listSplits(reservation.tableName()));
+ List currentAvailability =
+ new ArrayList<>(Collections.nCopies(splits.size() + 1, null));
+ try (var tabletInformation = tableOps.getTabletInformation(reservation.tableName(),
+ List.of(RowRange.all()), TabletInformation.Field.AVAILABILITY)) {
+ var iterator = tabletInformation.iterator();
+ while (iterator.hasNext()) {
+ var tablet = iterator.next();
+ Text endRow = tablet.getTabletId().getEndRow();
+ int tabletIndex =
+ endRow == null ? splits.size() : Collections.binarySearch(splits, endRow);
+ if (tabletIndex < 0 || tabletIndex >= currentAvailability.size()
+ || currentAvailability.set(tabletIndex, tablet.getTabletAvailability()) != null) {
+ throw new IllegalStateException(
+ "Could not map tablet availability to split boundaries for "
+ + reservation.tableName());
+ }
+ }
+ }
+ if (currentAvailability.contains(null)) {
+ throw new IllegalStateException(
+ "Could not obtain availability for every tablet in " + reservation.tableName());
+ }
+
+ List changes =
+ ManagerStressRows.randomAvailabilityChanges(splits, currentAvailability, random);
+ int changedTablets = 0;
+ for (ManagerStressRows.AvailabilityChange change : changes) {
+ submission.markSubmitted();
+ tableOps.setTabletAvailability(reservation.tableName(), change.rowRange(),
+ change.availability());
+ changedTablets += change.tabletCount();
+ }
+ log.debug("Worker {} submitted {} availability range change(s) for {} tablet(s) in {}",
+ options.workerId, changes.size(), changedTablets, reservation.tableName());
+ return !changes.isEmpty();
+ } finally {
+ pool.unlockTable(reservation);
+ }
+ }
+
+ private void bulkImport(AccumuloClient client, String tableName, SubmissionTracker submission)
+ throws Exception {
+ Path importDir =
+ new Path(workerBulkDir, "import-" + UUID.randomUUID().toString().replace("-", ""));
+ Path stagedFile = new Path(importDir, "one-entry.rf");
+ fileSystem.mkdirs(importDir);
+ try (FSDataInputStream in = fileSystem.open(seedFile);
+ FSDataOutputStream out = fileSystem.create(stagedFile, false)) {
+ IOUtils.copyBytes(in, out, new Configuration(), false);
+ }
+ try {
+ submission.markSubmitted();
+ client.tableOperations().importDirectory(importDir.toString()).to(tableName).tableTime(true)
+ .load();
+ log.debug("Worker {} bulk imported data into table {}", options.workerId, tableName);
+ } finally {
+ fileSystem.delete(importDir, true);
+ }
+ }
+
+ private void prepareBulkFile(Random random) throws IOException {
+ Configuration configuration = new Configuration();
+ Path configuredDir = new Path(options.hdfsDir);
+ fileSystem = configuredDir.getFileSystem(configuration);
+ workerBulkDir = new Path(new Path(configuredDir, options.runId), "worker-" + options.workerId);
+ if (!fileSystem.mkdirs(workerBulkDir) && !fileSystem.exists(workerBulkDir)) {
+ throw new IOException("Could not create HDFS work directory " + workerBulkDir);
+ }
+ seedFile = new Path(workerBulkDir, "seed.rf");
+ Key key = new Key(ManagerStressRows.randomRow(random), new Text("cf"), new Text("cq"));
+ Value value = new Value(("manager-stress-" + options.runId).getBytes(UTF_8));
+ try (RFileWriter writer =
+ RFile.newWriter().to(seedFile.toString()).withFileSystem(fileSystem).build()) {
+ writer.startDefaultLocalityGroup();
+ writer.append(key, value);
+ }
+ log.info("Worker {} created one-entry RFile at {}", options.workerId, seedFile);
+ }
+
+ private long awaitStart(java.nio.file.Path stateDir) throws Exception {
+ java.nio.file.Path startFile = stateDir.resolve("start");
+ while (!Files.exists(startFile)) {
+ Thread.sleep(100);
+ }
+ return Long.parseLong(Files.readString(startFile, UTF_8));
+ }
+
+ private boolean isExpectedRace(Throwable error) {
+ for (Throwable current = error; current != null; current = current.getCause()) {
+ if (current instanceof TableNotFoundException || current instanceof TableExistsException
+ || current instanceof TableOfflineException) {
+ return true;
+ }
+ }
+ return false;
+ }
+
+ private boolean isOutcomeUnknown(Throwable error) {
+ for (Throwable current = error; current != null; current = current.getCause()) {
+ if (current instanceof TTransportException || current instanceof InterruptedException) {
+ return true;
+ }
+ }
+ return false;
+ }
+
+ private void writeMetricsSnapshotIfDue() {
+ if (System.nanoTime() - lastMetricsSnapshotNanos >= 1_000_000_000L) {
+ writeMetricsSnapshot();
+ }
+ }
+
+ private void writeMetricsSnapshot() {
+ try {
+ metrics.writeSnapshot(stateDirectory, options.workerId);
+ lastMetricsSnapshotNanos = System.nanoTime();
+ } catch (IOException e) {
+ log.warn("Worker {} could not write its metrics snapshot", options.workerId, e);
+ }
+ }
+
+ private void report() {
+ for (Operation operation : Operation.values()) {
+ ManagerStressMetrics.Counts result = metrics.counts(operation);
+ if (result.attempts > 0) {
+ double averageMillis = (double) result.totalNanos / result.attempts / 1_000_000;
+ log.info(
+ "Worker {} {}: attempts={}, submitted={}, completed={}, asyncAccepted={}, skipped={}, "
+ + "races={}, failed={}, outcomeUnknown={}, avgMs={}",
+ options.workerId, operation, result.attempts, result.submitted, result.completed,
+ result.asyncAccepted, result.skipped, result.races, result.failed,
+ result.outcomeUnknown, String.format("%.3f", averageMillis));
+ }
+ }
+ log.info("Worker {} replacement-table maintenance creates={}", options.workerId,
+ metrics.maintenanceCreates());
+ }
+}
diff --git a/src/main/java/org/apache/accumulo/testing/manager/stress/Operation.java b/src/main/java/org/apache/accumulo/testing/manager/stress/Operation.java
new file mode 100644
index 00000000..433a3e8f
--- /dev/null
+++ b/src/main/java/org/apache/accumulo/testing/manager/stress/Operation.java
@@ -0,0 +1,77 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * https://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+package org.apache.accumulo.testing.manager.stress;
+
+import java.util.random.RandomGenerator;
+
+enum Operation {
+ CREATE, DELETE, SPLIT, MERGE, AVAILABILITY, COMPACT, BULK_IMPORT;
+
+ boolean isAsynchronous() {
+ return this == AVAILABILITY || this == COMPACT;
+ }
+
+ static Operation choose(ManagerStressOptions options, RandomGenerator random) {
+ double total =
+ options.createWeight + options.deleteWeight + options.splitWeight + options.mergeWeight
+ + options.availabilityWeight + options.compactWeight + options.bulkImportWeight;
+ return choose(options, random.nextDouble() * total);
+ }
+
+ static Operation choose(ManagerStressOptions options, double selection) {
+ if ((selection -= options.createWeight) < 0) {
+ return CREATE;
+ }
+ if ((selection -= options.deleteWeight) < 0) {
+ return DELETE;
+ }
+ if ((selection -= options.splitWeight) < 0) {
+ return SPLIT;
+ }
+ if ((selection -= options.mergeWeight) < 0) {
+ return MERGE;
+ }
+ if ((selection -= options.availabilityWeight) < 0) {
+ return AVAILABILITY;
+ }
+ if ((selection -= options.compactWeight) < 0) {
+ return COMPACT;
+ }
+ if (options.bulkImportWeight > 0) {
+ return BULK_IMPORT;
+ }
+ // Guard floating point rounding at the final positive-weight boundary.
+ if (options.availabilityWeight > 0) {
+ return AVAILABILITY;
+ }
+ if (options.compactWeight > 0) {
+ return COMPACT;
+ }
+ if (options.mergeWeight > 0) {
+ return MERGE;
+ }
+ if (options.splitWeight > 0) {
+ return SPLIT;
+ }
+ if (options.deleteWeight > 0) {
+ return DELETE;
+ }
+ return CREATE;
+ }
+}
diff --git a/src/main/java/org/apache/accumulo/testing/manager/stress/TablePool.java b/src/main/java/org/apache/accumulo/testing/manager/stress/TablePool.java
new file mode 100644
index 00000000..5f9d8e12
--- /dev/null
+++ b/src/main/java/org/apache/accumulo/testing/manager/stress/TablePool.java
@@ -0,0 +1,294 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * https://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+package org.apache.accumulo.testing.manager.stress;
+
+import static java.nio.charset.StandardCharsets.UTF_8;
+import static java.nio.file.StandardCopyOption.ATOMIC_MOVE;
+import static java.nio.file.StandardCopyOption.REPLACE_EXISTING;
+import static java.nio.file.StandardOpenOption.CREATE;
+import static java.nio.file.StandardOpenOption.WRITE;
+
+import java.io.IOException;
+import java.nio.channels.FileChannel;
+import java.nio.channels.FileLock;
+import java.nio.file.Files;
+import java.nio.file.Path;
+import java.util.ArrayList;
+import java.util.List;
+import java.util.Random;
+
+import org.apache.accumulo.core.client.AccumuloClient;
+import org.apache.accumulo.core.client.TableNotFoundException;
+import org.apache.accumulo.core.client.admin.NewTableConfiguration;
+
+/** Shared, locally locked table-pool state for coordinator and worker JVMs. */
+final class TablePool {
+
+ record Slot(int index, long generation, String tableName, boolean busy) {
+ }
+
+ @FunctionalInterface
+ private interface LockedFunction {
+ T apply(List slots) throws Exception;
+ }
+
+ private final Path stateDir;
+ private final Path stateFile;
+ private final Path lockFile;
+ private final String namespace;
+ private final String tablePrefix;
+ private final int tableCount;
+ private final long splitSeed;
+
+ TablePool(Path stateDir, String namespace, String tablePrefix, int tableCount, long splitSeed) {
+ this.stateDir = stateDir;
+ this.stateFile = stateDir.resolve("tables.state");
+ this.lockFile = stateDir.resolve("tables.lock");
+ this.namespace = namespace;
+ this.tablePrefix = tablePrefix;
+ this.tableCount = tableCount;
+ this.splitSeed = splitSeed;
+ }
+
+ void initialize(AccumuloClient client) throws Exception {
+ Files.createDirectories(stateDir);
+ List slots = new ArrayList<>(tableCount);
+ try {
+ for (int i = 0; i < tableCount; i++) {
+ String name = tableName(i, 1);
+ if (client.tableOperations().exists(name)) {
+ throw new IllegalStateException("Refusing to use existing table " + name);
+ }
+ createTable(client, name, i, 1);
+ slots.add(new Slot(i, 1, name, false));
+ }
+ writeSlots(slots);
+ } catch (Exception e) {
+ for (Slot slot : slots) {
+ try {
+ client.tableOperations().delete(slot.tableName());
+ } catch (Exception cleanupError) {
+ e.addSuppressed(cleanupError);
+ }
+ }
+ throw e;
+ }
+ }
+
+ List snapshot() throws Exception {
+ return withLock(slots -> List.copyOf(slots));
+ }
+
+ int replenish(AccumuloClient client) throws Exception {
+ int created = 0;
+ for (int i = 0; i < tableCount; i++) {
+ Slot reserved = reserve(i, true, true);
+ if (reserved != null) {
+ createInVacantSlot(client, reserved);
+ created++;
+ }
+ }
+ return created;
+ }
+
+ boolean delete(AccumuloClient client, int slotIndex, Runnable markSubmitted) throws Exception {
+ Slot reserved = reserve(slotIndex, false, false);
+ if (reserved == null) {
+ return false;
+ }
+ try {
+ markSubmitted.run();
+ client.tableOperations().delete(reserved.tableName());
+ } catch (TableNotFoundException e) {
+ updateSlot(new Slot(reserved.index(), reserved.generation(), null, false));
+ return false;
+ } catch (Exception e) {
+ updateSlot(new Slot(reserved.index(), reserved.generation(), reserved.tableName(), false));
+ throw e;
+ }
+ updateSlot(new Slot(reserved.index(), reserved.generation(), null, false));
+ return true;
+ }
+
+ boolean replace(AccumuloClient client, int slotIndex, Runnable markSubmitted) throws Exception {
+ Slot reserved = reserve(slotIndex, false, true);
+ if (reserved == null) {
+ return false;
+ }
+ String replacementName = tableName(reserved.index(), reserved.generation());
+ boolean replacementCreated = false;
+ boolean oldTableDeleted = false;
+ try {
+ markSubmitted.run();
+ createTable(client, replacementName, reserved.index(), reserved.generation());
+ replacementCreated = true;
+ if (reserved.tableName() != null) {
+ try {
+ client.tableOperations().delete(reserved.tableName());
+ oldTableDeleted = true;
+ } catch (TableNotFoundException e) {
+ // A competing client completed the deletion first.
+ oldTableDeleted = true;
+ }
+ }
+ updateSlot(new Slot(reserved.index(), reserved.generation(), replacementName, false));
+ return true;
+ } catch (Exception e) {
+ if (replacementCreated) {
+ try {
+ client.tableOperations().delete(replacementName);
+ } catch (Exception cleanupError) {
+ e.addSuppressed(cleanupError);
+ }
+ }
+ String previousName = oldTableDeleted ? null : reserved.tableName();
+ try {
+ updateSlot(new Slot(reserved.index(), reserved.generation(), previousName, false));
+ } catch (Exception updateError) {
+ e.addSuppressed(updateError);
+ }
+ throw e;
+ }
+ }
+
+ private void createInVacantSlot(AccumuloClient client, Slot reserved) throws Exception {
+ String replacementName = tableName(reserved.index(), reserved.generation());
+ boolean created = false;
+ try {
+ createTable(client, replacementName, reserved.index(), reserved.generation());
+ created = true;
+ updateSlot(new Slot(reserved.index(), reserved.generation(), replacementName, false));
+ } catch (Exception e) {
+ if (created) {
+ try {
+ client.tableOperations().delete(replacementName);
+ } catch (Exception cleanupError) {
+ e.addSuppressed(cleanupError);
+ }
+ }
+ try {
+ updateSlot(new Slot(reserved.index(), reserved.generation(), null, false));
+ } catch (Exception updateError) {
+ e.addSuppressed(updateError);
+ }
+ throw e;
+ }
+ }
+
+ Slot lockTable(int preferredIndex) throws Exception {
+ return reserve(preferredIndex, false, false);
+ }
+
+ void unlockTable(Slot reservation) throws Exception {
+ withLock(slots -> {
+ Slot current = slots.get(reservation.index());
+ if (!current.busy() || current.generation() != reservation.generation()
+ || !java.util.Objects.equals(current.tableName(), reservation.tableName())) {
+ throw new IllegalStateException(
+ "Table slot reservation changed while in use: " + reservation.index());
+ }
+ slots.set(reservation.index(),
+ new Slot(reservation.index(), reservation.generation(), reservation.tableName(), false));
+ writeSlots(slots);
+ return null;
+ });
+ }
+
+ private void createTable(AccumuloClient client, String tableName, int slot, long generation)
+ throws Exception {
+ long seed = splitSeed ^ ((long) slot * 0x9e3779b97f4a7c15L) ^ Long.rotateLeft(generation, 32);
+ NewTableConfiguration config = ManagerStressRows.newTableConfiguration(new Random(seed));
+ client.tableOperations().create(tableName, config);
+ }
+
+ private Slot reserve(int preferredIndex, boolean vacant, boolean incrementGeneration)
+ throws Exception {
+ return withLock(slots -> {
+ for (int offset = 0; offset < slots.size(); offset++) {
+ int index = (preferredIndex + offset) % slots.size();
+ Slot current = slots.get(index);
+ if (!current.busy() && (vacant == (current.tableName() == null))) {
+ long generation = current.generation() + (incrementGeneration ? 1 : 0);
+ Slot reserved = new Slot(index, generation, current.tableName(), true);
+ slots.set(index, reserved);
+ writeSlots(slots);
+ return reserved;
+ }
+ }
+ return null;
+ });
+ }
+
+ private void updateSlot(Slot updated) throws Exception {
+ withLock(slots -> {
+ slots.set(updated.index(), updated);
+ writeSlots(slots);
+ return null;
+ });
+ }
+
+ String tableName(int slot, long generation) {
+ return namespace + "." + tablePrefix + "_" + slot + "_" + generation;
+ }
+
+ private T withLock(LockedFunction function) throws Exception {
+ try (FileChannel channel = FileChannel.open(lockFile, CREATE, WRITE);
+ FileLock ignored = channel.lock()) {
+ return function.apply(readSlots());
+ }
+ }
+
+ private List readSlots() throws IOException {
+ List lines = Files.readAllLines(stateFile, UTF_8);
+ if (lines.size() != tableCount) {
+ throw new IOException("Invalid table-pool state in " + stateFile);
+ }
+ List slots = new ArrayList<>(tableCount);
+ for (int i = 0; i < lines.size(); i++) {
+ String[] values = lines.get(i).split("\\t", -1);
+ if (values.length != 3) {
+ throw new IOException("Invalid table-pool state line: " + lines.get(i));
+ }
+ long generation = Long.parseLong(values[0]);
+ String name = values[1].isEmpty() ? null : values[1];
+ slots.add(new Slot(i, generation, name, Boolean.parseBoolean(values[2])));
+ }
+ return slots;
+ }
+
+ private void writeSlots(List slots) throws IOException {
+ Path tempFile = Files.createTempFile(stateDir, "tables", ".tmp");
+ List<
+ String> lines =
+ slots.stream()
+ .map(slot -> slot.generation() + "\t"
+ + (slot.tableName() == null ? "" : slot.tableName()) + "\t" + slot.busy())
+ .toList();
+ try {
+ Files.write(tempFile, lines, UTF_8);
+ try {
+ Files.move(tempFile, stateFile, ATOMIC_MOVE, REPLACE_EXISTING);
+ } catch (UnsupportedOperationException | IOException e) {
+ Files.move(tempFile, stateFile, REPLACE_EXISTING);
+ }
+ } finally {
+ Files.deleteIfExists(tempFile);
+ }
+ }
+}
diff --git a/src/main/java/org/apache/accumulo/testing/manager/stress/compactor/CompactionOutcome.java b/src/main/java/org/apache/accumulo/testing/manager/stress/compactor/CompactionOutcome.java
new file mode 100644
index 00000000..da6b05f1
--- /dev/null
+++ b/src/main/java/org/apache/accumulo/testing/manager/stress/compactor/CompactionOutcome.java
@@ -0,0 +1,46 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * https://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+package org.apache.accumulo.testing.manager.stress.compactor;
+
+import java.util.random.RandomGenerator;
+
+enum CompactionOutcome {
+ SUCCESS, FAILURE, CANCELLATION;
+
+ static CompactionOutcome choose(StressCompactorOptions options, RandomGenerator random) {
+ double total = options.successWeight + options.failureWeight + options.cancellationWeight;
+ return choose(options, random.nextDouble() * total);
+ }
+
+ static CompactionOutcome choose(StressCompactorOptions options, double selection) {
+ if ((selection -= options.successWeight) < 0) {
+ return SUCCESS;
+ }
+ if ((selection -= options.failureWeight) < 0) {
+ return FAILURE;
+ }
+ if (options.cancellationWeight > 0) {
+ return CANCELLATION;
+ }
+ if (options.failureWeight > 0) {
+ return FAILURE;
+ }
+ return SUCCESS;
+ }
+}
diff --git a/src/main/java/org/apache/accumulo/testing/manager/stress/compactor/SimulatedCompactor.java b/src/main/java/org/apache/accumulo/testing/manager/stress/compactor/SimulatedCompactor.java
new file mode 100644
index 00000000..e190abab
--- /dev/null
+++ b/src/main/java/org/apache/accumulo/testing/manager/stress/compactor/SimulatedCompactor.java
@@ -0,0 +1,485 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * https://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+package org.apache.accumulo.testing.manager.stress.compactor;
+
+import static java.util.concurrent.TimeUnit.MILLISECONDS;
+
+import java.util.List;
+import java.util.UUID;
+import java.util.concurrent.Semaphore;
+import java.util.concurrent.ThreadLocalRandom;
+import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.concurrent.atomic.AtomicReference;
+import java.util.concurrent.atomic.LongAdder;
+import java.util.random.RandomGenerator;
+
+import org.apache.accumulo.core.cli.ServerOpts;
+import org.apache.accumulo.core.client.Accumulo;
+import org.apache.accumulo.core.client.AccumuloClient;
+import org.apache.accumulo.core.client.AccumuloSecurityException;
+import org.apache.accumulo.core.client.admin.servers.ServerId.Type;
+import org.apache.accumulo.core.clientImpl.ClientContext;
+import org.apache.accumulo.core.clientImpl.thrift.SecurityErrorCode;
+import org.apache.accumulo.core.clientImpl.thrift.TInfo;
+import org.apache.accumulo.core.clientImpl.thrift.ThriftSecurityException;
+import org.apache.accumulo.core.compaction.thrift.CompactorService;
+import org.apache.accumulo.core.compaction.thrift.TCompactionState;
+import org.apache.accumulo.core.compaction.thrift.TCompactionStatusUpdate;
+import org.apache.accumulo.core.compaction.thrift.TExternalCompaction;
+import org.apache.accumulo.core.compaction.thrift.TNextCompactionJob;
+import org.apache.accumulo.core.compaction.thrift.UnknownCompactionIdException;
+import org.apache.accumulo.core.conf.Property;
+import org.apache.accumulo.core.conf.SiteConfiguration;
+import org.apache.accumulo.core.data.ResourceGroupId;
+import org.apache.accumulo.core.lock.ServiceLock;
+import org.apache.accumulo.core.lock.ServiceLock.LockLossReason;
+import org.apache.accumulo.core.lock.ServiceLock.LockWatcher;
+import org.apache.accumulo.core.lock.ServiceLockData;
+import org.apache.accumulo.core.lock.ServiceLockData.ServiceDescriptor;
+import org.apache.accumulo.core.lock.ServiceLockData.ServiceDescriptors;
+import org.apache.accumulo.core.lock.ServiceLockData.ThriftService;
+import org.apache.accumulo.core.lock.ServiceLockPaths.ServiceLockPath;
+import org.apache.accumulo.core.lock.ServiceLockSupport;
+import org.apache.accumulo.core.metadata.schema.ExternalCompactionId;
+import org.apache.accumulo.core.metrics.MetricsInfo;
+import org.apache.accumulo.core.rpc.ThriftUtil;
+import org.apache.accumulo.core.securityImpl.thrift.TCredentials;
+import org.apache.accumulo.core.tabletserver.thrift.ActiveCompaction;
+import org.apache.accumulo.core.tabletserver.thrift.TCompactionStats;
+import org.apache.accumulo.core.tabletserver.thrift.TExternalCompactionJob;
+import org.apache.accumulo.core.trace.TraceUtil;
+import org.apache.accumulo.core.util.compaction.ExternalCompactionUtil;
+import org.apache.accumulo.server.AbstractServer;
+import org.apache.accumulo.server.ServerContext;
+import org.apache.accumulo.server.client.ClientServiceHandler;
+import org.apache.accumulo.server.rpc.ServerAddress;
+import org.apache.accumulo.server.rpc.TServerUtils;
+import org.apache.accumulo.server.rpc.ThriftProcessorTypes;
+import org.apache.thrift.TException;
+import org.apache.zookeeper.KeeperException;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import com.google.common.net.HostAndPort;
+
+/** A registered Compactor identity which polls the Manager without doing file compaction. */
+final class SimulatedCompactor extends AbstractServer implements CompactorService.Iface {
+
+ private static final Logger log = LoggerFactory.getLogger(SimulatedCompactor.class);
+
+ static final class Counters {
+ final LongAdder jobRequests = new LongAdder();
+ final LongAdder idleResponses = new LongAdder();
+ final LongAdder jobsReceived = new LongAdder();
+ final LongAdder successes = new LongAdder();
+ final LongAdder failures = new LongAdder();
+ final LongAdder cancellations = new LongAdder();
+ final LongAdder rpcFailures = new LongAdder();
+ final LongAdder outcomeReportFailures = new LongAdder();
+ }
+
+ private final StressCompactorOptions options;
+ private final ServerContext serverContext;
+ private final int instanceNumber;
+ private final UUID serverUuid = UUID.randomUUID();
+ private final ResourceGroupId group;
+ private final ClientContext clientContext;
+ private final RandomGenerator random;
+ private final AtomicBoolean stopping = new AtomicBoolean();
+ private final AtomicBoolean closed = new AtomicBoolean();
+ private final AtomicBoolean cancelRequested = new AtomicBoolean();
+ private final AtomicReference currentId = new AtomicReference<>();
+ private final AtomicReference currentCompaction = new AtomicReference<>();
+ private final Semaphore wakeups = new Semaphore(0);
+ private final Counters counters = new Counters();
+
+ private ServerAddress serviceAddress;
+ private ServiceLock serviceLock;
+ private Thread worker;
+ private volatile String advertisedAddress;
+
+ SimulatedCompactor(StressCompactorOptions options, int instanceNumber) throws Exception {
+ super(Type.COMPACTOR, new ServerOpts(), ServerContext::new, serverArgs(options));
+ this.options = options;
+ this.instanceNumber = instanceNumber;
+ this.group = getResourceGroup();
+ this.serverContext = getContext();
+ ClientContext newClientContext;
+ try {
+ newClientContext = createClientContext(options.getClientProps());
+ } catch (Exception e) {
+ try {
+ super.close();
+ } catch (Exception cleanupError) {
+ e.addSuppressed(cleanupError);
+ }
+ throw e;
+ }
+ this.clientContext = newClientContext;
+ String configuredGroup =
+ serverContext.getSiteConfiguration().get(Property.COMPACTOR_GROUP_NAME);
+ if (!group.canonical().equals(configuredGroup)
+ || !clientContext.getInstanceID().equals(serverContext.getInstanceID())) {
+ try {
+ clientContext.close();
+ } finally {
+ super.close();
+ }
+ throw new IllegalStateException("Simulated Compactor configuration mismatch: group="
+ + configuredGroup + ", expectedGroup=" + group.canonical()
+ + "; client and server contexts must also use the same Accumulo instance");
+ }
+ this.random = new java.util.Random(options.seed + instanceNumber * 0x9e3779b97f4a7c15L);
+ }
+
+ private static String[] serverArgs(StressCompactorOptions options) {
+ String compactorGroup = ResourceGroupId.of(options.resourceGroup).canonical();
+ return new String[] {"-o", Property.COMPACTOR_GROUP_NAME.getKey() + "=" + compactorGroup};
+ }
+
+ @Override
+ protected String getResourceGroupPropertyValue(SiteConfiguration siteConfig) {
+ return siteConfig.get(Property.COMPACTOR_GROUP_NAME);
+ }
+
+ @Override
+ public ServiceLock getLock() {
+ return serviceLock;
+ }
+
+ Counters getCounters() {
+ return counters;
+ }
+
+ void start() throws Exception {
+ var processor = ThriftProcessorTypes.getCompactorTProcessor(this,
+ new ClientServiceHandler(serverContext), this, serverContext);
+ serviceAddress = createThriftServer(processor);
+ serviceAddress.startThriftServer("SimulatedCompactorService-" + instanceNumber);
+ updateAdvertiseAddress(serviceAddress.address);
+ advertisedAddress = serviceAddress.address.toString();
+
+ ServiceLockPath lockPath =
+ serverContext.getServerPaths().createCompactorPath(group, serviceAddress.address);
+ ServiceLockSupport.createNonHaServiceLockPath(Type.COMPACTOR,
+ serverContext.getZooSession().asReaderWriter(), lockPath);
+ serviceLock = new ServiceLock(serverContext.getZooSession(), lockPath, serverUuid);
+ ServiceDescriptors descriptors = new ServiceDescriptors();
+ descriptors.addService(
+ new ServiceDescriptor(serverUuid, ThriftService.CLIENT, advertisedAddress, group));
+ descriptors.addService(
+ new ServiceDescriptor(serverUuid, ThriftService.COMPACTOR, advertisedAddress, group));
+ LockWatcher lockWatcher = new LockWatcher() {
+ @Override
+ public void lostLock(LockLossReason reason) {
+ log.warn("Simulated Compactor {} lost its service registration: {}", instanceNumber,
+ reason);
+ stopping.set(true);
+ wakeups.release();
+ }
+
+ @Override
+ public void unableToMonitorLockNode(Exception e) {
+ log.warn("Unable to monitor simulated Compactor {} service registration", instanceNumber,
+ e);
+ stopping.set(true);
+ wakeups.release();
+ }
+ };
+ if (!serviceLock.tryLock(lockWatcher, new ServiceLockData(descriptors))) {
+ serviceAddress.server.stop();
+ throw new IllegalStateException(
+ "Could not register simulated Compactor at " + advertisedAddress);
+ }
+
+ log.info("Started simulated Compactor {} at {} in resource group {}", instanceNumber,
+ advertisedAddress, group);
+ }
+
+ void startPolling() {
+ worker = new Thread(this, "CompactorSimulator-" + instanceNumber);
+ worker.setDaemon(false);
+ worker.start();
+ }
+
+ @Override
+ public void run() {
+ runLoop();
+ }
+
+ private ServerAddress createThriftServer(org.apache.thrift.TProcessor processor) {
+ var conf = serverContext.getConfiguration();
+ var address = options.addresses();
+ long threadTimeout =
+ Math.max(10_000, conf.getTimeInMillis(Property.COMPACTOR_MINTHREADS_TIMEOUT));
+ return TServerUtils.createThriftServer(conf, serverContext.getThriftServerType(), processor,
+ serverContext.getInstanceID(), "SimulatedCompactor-" + instanceNumber, 1, threadTimeout,
+ conf.getTimeInMillis(Property.COMPACTOR_THREADCHECK),
+ conf.getAsBytes(Property.RPC_MAX_MESSAGE_SIZE), serverContext.getServerSslParams(),
+ serverContext.getSaslParams(), serverContext.getClientTimeoutInMillis(),
+ conf.getCount(Property.RPC_BACKLOG), serverContext.getMetricsInfo(),
+ address.toArray(HostAndPort[]::new));
+ }
+
+ private static ClientContext createClientContext(java.util.Properties clientProperties) {
+ AccumuloClient client = Accumulo.newClient().from(clientProperties).build();
+ if (client instanceof ClientContext context) {
+ return context;
+ }
+ try {
+ client.close();
+ } catch (Exception e) {
+ log.warn("Could not close unsupported client implementation", e);
+ }
+ throw new IllegalStateException("Expected the Accumulo client builder to return ClientContext");
+ }
+
+ private void runLoop() {
+ var metricsInfo = serverContext.getMetricsInfo();
+ metricsInfo.addMetricsProducers(this);
+ try {
+ metricsInfo.init(MetricsInfo.serviceTags(serverContext.getInstanceName(),
+ getApplicationName(), serviceAddress.address, group));
+ } catch (IllegalStateException e) {
+ // Catch this because it is thrown due to many Compactor instances in the same VM.
+ }
+
+ long minWait =
+ serverContext.getConfiguration().getTimeInMillis(Property.COMPACTOR_MIN_JOB_WAIT_TIME);
+ long maxWait =
+ serverContext.getConfiguration().getTimeInMillis(Property.COMPACTOR_MAX_JOB_WAIT_TIME);
+ while (!stopping.get()) {
+ String externalId = ExternalCompactionId.generate(UUID.randomUUID()).canonical();
+ boolean assignedJob = false;
+ cancelRequested.set(false);
+ currentId.set(externalId);
+ counters.jobRequests.increment();
+ try {
+ TNextCompactionJob next = requestJob(externalId);
+ TExternalCompactionJob job = next.getJob();
+ if (!job.isSetExternalCompactionId()) {
+ currentId.compareAndSet(externalId, null);
+ counters.idleResponses.increment();
+ waitForWork(calculateWait(next.getCompactorCount(), minWait, maxWait));
+ continue;
+ }
+
+ counters.jobsReceived.increment();
+ assignedJob = true;
+ if (!externalId.equals(job.getExternalCompactionId())) {
+ throw new IllegalStateException("Manager returned compaction ID "
+ + job.getExternalCompactionId() + " for request " + externalId);
+ }
+ TExternalCompaction running = new TExternalCompaction();
+ running.setCompactor(advertisedAddress);
+ running.setGroupName(group.canonical());
+ long startTime = System.currentTimeMillis();
+ running.setStartTime(startTime);
+ running.putToUpdates(startTime, new TCompactionStatusUpdate(TCompactionState.STARTED,
+ "Simulated compaction started", -1, -1, -1, 0));
+ running.setJob(job);
+ currentCompaction.set(running);
+
+ CompactionOutcome outcome = CompactionOutcome.choose(options, random);
+ if (cancelRequested.get()) {
+ outcome = CompactionOutcome.CANCELLATION;
+ }
+ reportOutcome(job, outcome);
+ } catch (Exception e) {
+ counters.rpcFailures.increment();
+ if (assignedJob) {
+ counters.outcomeReportFailures.increment();
+ }
+ log.warn("Simulated Compactor {} failed a Manager RPC", instanceNumber, e);
+ waitForWork(minWait);
+ } finally {
+ currentCompaction.set(null);
+ currentId.set(null);
+ cancelRequested.set(false);
+ }
+ }
+ }
+
+ private TNextCompactionJob requestJob(String externalId) throws Exception {
+ var coordinator = ExternalCompactionUtil.getCoordinatorClient(clientContext);
+ try {
+ return coordinator.getCompactionJob(TraceUtil.traceInfo(), clientContext.rpcCreds(),
+ group.canonical(), advertisedAddress, externalId);
+ } finally {
+ ThriftUtil.returnClient(coordinator, clientContext);
+ }
+ }
+
+ private void reportOutcome(TExternalCompactionJob job, CompactionOutcome outcome)
+ throws Exception {
+ var coordinator = ExternalCompactionUtil.getCoordinatorClient(clientContext);
+ try {
+ switch (outcome) {
+ case SUCCESS -> {
+ TCompactionStats stats = new TCompactionStats();
+ stats.setEntriesRead(0);
+ stats.setEntriesWritten(0);
+ stats.setFileSize(0);
+ coordinator.compactionCompleted(TraceUtil.traceInfo(), clientContext.rpcCreds(),
+ job.getExternalCompactionId(), job.getExtent(), stats, group.canonical(),
+ advertisedAddress);
+ counters.successes.increment();
+ }
+ case FAILURE -> {
+ coordinator.compactionFailed(TraceUtil.traceInfo(), clientContext.rpcCreds(),
+ job.getExternalCompactionId(), job.getExtent(), "simulated compaction failure",
+ TCompactionState.FAILED, group.canonical(), advertisedAddress);
+ counters.failures.increment();
+ }
+ case CANCELLATION -> {
+ coordinator.compactionFailed(TraceUtil.traceInfo(), clientContext.rpcCreds(),
+ job.getExternalCompactionId(), job.getExtent(), "simulated compaction cancellation",
+ TCompactionState.CANCELLED, group.canonical(), advertisedAddress);
+ counters.cancellations.increment();
+ }
+ }
+ TCompactionState state = switch (outcome) {
+ case SUCCESS -> TCompactionState.SUCCEEDED;
+ case FAILURE -> TCompactionState.FAILED;
+ case CANCELLATION -> TCompactionState.CANCELLED;
+ };
+ log.info("Simulated compaction {} finished with state {}", job.getExternalCompactionId(),
+ state);
+ } finally {
+ ThriftUtil.returnClient(coordinator, clientContext);
+ }
+ }
+
+ private static long calculateWait(int compactorCount, long minWait, long maxWait) {
+ long wait = Math.max(minWait, (long) compactorCount * minWait / 3);
+ wait = Math.min(maxWait, wait);
+ return (long) (wait * (0.9 + 0.2 * ThreadLocalRandom.current().nextDouble()));
+ }
+
+ private void waitForWork(long millis) {
+ if (stopping.get()) {
+ return;
+ }
+ try {
+ wakeups.tryAcquire(Math.max(1, millis), MILLISECONDS);
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ stopping.set(true);
+ }
+ }
+
+ private void requireSystemAction(TCredentials credentials) throws ThriftSecurityException {
+ if (!serverContext.getSecurityOperation().canPerformSystemActions(credentials)) {
+ throw new AccumuloSecurityException(credentials.getPrincipal(),
+ SecurityErrorCode.PERMISSION_DENIED).asThriftException();
+ }
+ }
+
+ @Override
+ public void gracefulShutdown(TCredentials credentials) {
+ super.gracefulShutdown(credentials);
+ if (isShutdownRequested()) {
+ stopping.set(true);
+ wakeups.release();
+ }
+ }
+
+ @Override
+ public void close() {
+ if (!closed.compareAndSet(false, true)) {
+ return;
+ }
+ stopping.set(true);
+ wakeups.release();
+ if (worker != null) {
+ worker.interrupt();
+ try {
+ worker.join(10_000);
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ }
+ }
+ if (serviceLock != null) {
+ try {
+ serviceLock.unlock();
+ } catch (KeeperException | InterruptedException e) {
+ if (e instanceof InterruptedException) {
+ Thread.currentThread().interrupt();
+ }
+ log.warn("Could not unregister simulated Compactor {}", advertisedAddress, e);
+ }
+ }
+ if (serviceAddress != null) {
+ serviceAddress.server.stop();
+ }
+ try {
+ clientContext.close();
+ } catch (Exception e) {
+ log.warn("Could not close client context for simulated Compactor {}", instanceNumber, e);
+ }
+ super.close();
+ log.info(
+ "Compactor {} totals: polls={}, idle={}, jobs={}, success={}, failure={}, cancelled={}, "
+ + "rpcFailures={}, outcomeReportFailures={}",
+ instanceNumber, counters.jobRequests.sum(), counters.idleResponses.sum(),
+ counters.jobsReceived.sum(), counters.successes.sum(), counters.failures.sum(),
+ counters.cancellations.sum(), counters.rpcFailures.sum(),
+ counters.outcomeReportFailures.sum());
+ }
+
+ @Override
+ public TExternalCompaction getRunningCompaction(TInfo tinfo, TCredentials credentials)
+ throws ThriftSecurityException {
+ requireSystemAction(credentials);
+ TExternalCompaction current = currentCompaction.get();
+ return current == null ? new TExternalCompaction() : new TExternalCompaction(current);
+ }
+
+ @Override
+ public String getRunningCompactionId(TInfo tinfo, TCredentials credentials)
+ throws ThriftSecurityException {
+ requireSystemAction(credentials);
+ String id = currentId.get();
+ return id == null ? "" : id;
+ }
+
+ @Override
+ public List getActiveCompactions(TInfo tinfo, TCredentials credentials)
+ throws ThriftSecurityException {
+ requireSystemAction(credentials);
+ return List.of();
+ }
+
+ public void wake(TInfo tinfo, TCredentials credentials) throws ThriftSecurityException {
+ requireSystemAction(credentials);
+ wakeups.release();
+ }
+
+ @Override
+ public void cancel(TInfo tinfo, TCredentials credentials, String externalCompactionId)
+ throws TException {
+ requireSystemAction(credentials);
+ if (externalCompactionId.equals(currentId.get())) {
+ cancelRequested.set(true);
+ } else {
+ throw new UnknownCompactionIdException();
+ }
+ }
+
+}
diff --git a/src/main/java/org/apache/accumulo/testing/manager/stress/compactor/StressCompactor.java b/src/main/java/org/apache/accumulo/testing/manager/stress/compactor/StressCompactor.java
new file mode 100644
index 00000000..b0435f94
--- /dev/null
+++ b/src/main/java/org/apache/accumulo/testing/manager/stress/compactor/StressCompactor.java
@@ -0,0 +1,105 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * https://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+package org.apache.accumulo.testing.manager.stress.compactor;
+
+import java.net.InetAddress;
+import java.time.Duration;
+import java.util.ArrayList;
+import java.util.List;
+import java.util.concurrent.atomic.AtomicBoolean;
+
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+/** Starts multiple registered Compactor simulators as threads in this JVM. */
+public class StressCompactor {
+
+ private static final Logger log = LoggerFactory.getLogger(StressCompactor.class);
+
+ public static void main(String[] args) throws Exception {
+ if (args.length == 0) {
+ args = new String[] {"--help"};
+ }
+ StressCompactorOptions options = new StressCompactorOptions();
+ options.parseArgs(StressCompactor.class.getName(), args);
+ if (options.host == null || options.host.isBlank()) {
+ options.host = InetAddress.getLocalHost().getCanonicalHostName();
+ }
+ options.validate();
+
+ List compactors = new ArrayList<>(options.compactorsPerJvm);
+ AtomicBoolean stopping = new AtomicBoolean();
+ Runtime.getRuntime().addShutdownHook(new Thread(() -> {
+ if (stopping.compareAndSet(false, true)) {
+ stopCompactors(compactors);
+ }
+ }, "compactor-stress-shutdown"));
+
+ try {
+ for (int i = 0; i < options.compactorsPerJvm; i++) {
+ SimulatedCompactor compactor = new SimulatedCompactor(options, i);
+ compactors.add(compactor);
+ compactor.start();
+ }
+
+ long deadline = Math.addExact(System.currentTimeMillis(), options.duration);
+ compactors.forEach(SimulatedCompactor::startPolling);
+ log.info("Started {} Compactor simulator(s) in JVM {} for {} ms", compactors.size(),
+ ProcessHandle.current().pid(), options.duration);
+ while (System.currentTimeMillis() < deadline) {
+ long remaining = deadline - System.currentTimeMillis();
+ Thread.sleep(Math.min(remaining, Duration.ofSeconds(1).toMillis()));
+ }
+ } finally {
+ if (stopping.compareAndSet(false, true)) {
+ stopCompactors(compactors);
+ }
+ }
+ }
+
+ private static void stopCompactors(List compactors) {
+ for (int i = compactors.size() - 1; i >= 0; i--) {
+ compactors.get(i).close();
+ }
+ long polls = 0;
+ long idle = 0;
+ long jobs = 0;
+ long successes = 0;
+ long failures = 0;
+ long cancellations = 0;
+ long rpcFailures = 0;
+ long outcomeReportFailures = 0;
+ for (SimulatedCompactor compactor : compactors) {
+ SimulatedCompactor.Counters counters = compactor.getCounters();
+ polls += counters.jobRequests.sum();
+ idle += counters.idleResponses.sum();
+ jobs += counters.jobsReceived.sum();
+ successes += counters.successes.sum();
+ failures += counters.failures.sum();
+ cancellations += counters.cancellations.sum();
+ rpcFailures += counters.rpcFailures.sum();
+ outcomeReportFailures += counters.outcomeReportFailures.sum();
+ }
+ log.info(
+ "Compactor JVM totals: simulators={}, polls={}, idle={}, jobs={}, success={}, failure={}, "
+ + "cancelled={}, rpcFailures={}, outcomeReportFailures={}",
+ compactors.size(), polls, idle, jobs, successes, failures, cancellations, rpcFailures,
+ outcomeReportFailures);
+ }
+}
diff --git a/src/main/java/org/apache/accumulo/testing/manager/stress/compactor/StressCompactorOptions.java b/src/main/java/org/apache/accumulo/testing/manager/stress/compactor/StressCompactorOptions.java
new file mode 100644
index 00000000..609444b6
--- /dev/null
+++ b/src/main/java/org/apache/accumulo/testing/manager/stress/compactor/StressCompactorOptions.java
@@ -0,0 +1,112 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * https://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+package org.apache.accumulo.testing.manager.stress.compactor;
+
+import java.util.ArrayList;
+import java.util.List;
+
+import org.apache.accumulo.testing.cli.ClientOpts;
+
+import com.beust.jcommander.Parameter;
+import com.google.common.net.HostAndPort;
+
+class StressCompactorOptions extends ClientOpts {
+
+ @Parameter(names = "--duration", converter = TimeConverter.class,
+ description = "duration of this simulator JVM (for example 10m or 1h)")
+ long duration = 5 * 60 * 1000L;
+
+ @Parameter(names = "--compactors-per-jvm",
+ description = "number of simulated Compactor instances in this JVM")
+ int compactorsPerJvm = 1;
+
+ @Parameter(names = "--resource-group",
+ description = "Compactor resource group to register and request jobs from")
+ String resourceGroup = "default";
+
+ @Parameter(names = "--host",
+ description = "local host/interface reachable by the Manager for callbacks")
+ String host;
+
+ @Parameter(names = "--port-range",
+ description = "Compactor callback port or inclusive range; 0 selects an ephemeral port")
+ String portRange = "0";
+
+ @Parameter(names = "--success-weight",
+ description = "relative weight for successful zero-output completions")
+ double successWeight = 1;
+
+ @Parameter(names = "--failure-weight", description = "relative weight for failed jobs")
+ double failureWeight = 1;
+
+ @Parameter(names = "--cancellation-weight", description = "relative weight for cancelled jobs")
+ double cancellationWeight = 1;
+
+ @Parameter(names = "--seed", description = "random seed for outcome selection")
+ long seed = System.nanoTime();
+
+ void validate() {
+ if (duration <= 0) {
+ throw new IllegalArgumentException("--duration must be positive");
+ }
+ if (compactorsPerJvm <= 0) {
+ throw new IllegalArgumentException("--compactors-per-jvm must be positive");
+ }
+ if (resourceGroup == null || resourceGroup.isBlank()) {
+ throw new IllegalArgumentException("--resource-group must not be blank");
+ }
+ if (host == null || host.isBlank()) {
+ throw new IllegalArgumentException(
+ "--host must be a host/interface reachable by the Manager");
+ }
+ double totalWeight = successWeight + failureWeight + cancellationWeight;
+ for (double weight : new double[] {successWeight, failureWeight, cancellationWeight}) {
+ if (!Double.isFinite(weight) || weight < 0) {
+ throw new IllegalArgumentException("Outcome weights must be finite and non-negative");
+ }
+ }
+ if (!Double.isFinite(totalWeight) || totalWeight <= 0) {
+ throw new IllegalArgumentException("At least one outcome weight must be positive");
+ }
+ if (portRange == null || !portRange.matches("(?:0|[1-9][0-9]{0,4}(?:-[1-9][0-9]{0,4})?)")) {
+ throw new IllegalArgumentException("--port-range must be a port, 0, or inclusive port range");
+ }
+ String[] ports = portRange.split("-");
+ int first = Integer.parseInt(ports[0]);
+ int last = ports.length == 1 ? first : Integer.parseInt(ports[1]);
+ if (first > 65535 || last > 65535 || last < first || first == 0 && last != 0) {
+ throw new IllegalArgumentException("Invalid --port-range: " + portRange);
+ }
+ if (last - first > 4095) {
+ throw new IllegalArgumentException("--port-range cannot contain more than 4096 ports");
+ }
+ }
+
+ List addresses() {
+ String[] ports = portRange.split("-");
+ int first = Integer.parseInt(ports[0]);
+ int last = ports.length == 1 ? first : Integer.parseInt(ports[1]);
+ List addresses = new ArrayList<>(last - first + 1);
+ for (int port = first; port <= last; port++) {
+ addresses.add(HostAndPort.fromParts(host, port));
+ }
+ return addresses;
+ }
+
+}
diff --git a/src/test/java/org/apache/accumulo/testing/manager/stress/ManagerStressCleanupTest.java b/src/test/java/org/apache/accumulo/testing/manager/stress/ManagerStressCleanupTest.java
new file mode 100644
index 00000000..6f80b8d2
--- /dev/null
+++ b/src/test/java/org/apache/accumulo/testing/manager/stress/ManagerStressCleanupTest.java
@@ -0,0 +1,62 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * https://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+package org.apache.accumulo.testing.manager.stress;
+
+import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+
+import java.io.IOException;
+import java.nio.file.Files;
+import java.nio.file.Path;
+
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.io.TempDir;
+
+class ManagerStressCleanupTest {
+
+ @TempDir
+ Path temporaryDirectory;
+
+ @Test
+ void controlDirectoryCleanupIsIdempotent() throws IOException {
+ Path controlDirectory = temporaryDirectory.resolve("manager-stress-control");
+ Path nestedFile = controlDirectory.resolve("nested").resolve("snapshot.properties");
+ Files.createDirectories(nestedFile.getParent());
+ Files.writeString(nestedFile, "metrics=true");
+
+ ManagerStress.deleteRecursivelyIfExists(controlDirectory);
+
+ assertFalse(Files.exists(controlDirectory));
+ assertDoesNotThrow(() -> ManagerStress.deleteRecursivelyIfExists(controlDirectory));
+ }
+
+ @Test
+ void hdfsBulkImportRunDirectoryCleanupIsRecursiveAndIdempotent() throws IOException {
+ Path runDirectory = temporaryDirectory.resolve("run-1");
+ Path stagedImport = runDirectory.resolve("worker-0/import-1/data.rf");
+ Files.createDirectories(stagedImport.getParent());
+ Files.writeString(stagedImport, "bulk import data");
+
+ ManagerStress.deleteBulkImportRunDirectory(temporaryDirectory.toUri().toString(), "run-1");
+
+ assertFalse(Files.exists(runDirectory));
+ assertDoesNotThrow(() -> ManagerStress
+ .deleteBulkImportRunDirectory(temporaryDirectory.toUri().toString(), "run-1"));
+ }
+}
diff --git a/src/test/java/org/apache/accumulo/testing/manager/stress/ManagerStressMetricsTest.java b/src/test/java/org/apache/accumulo/testing/manager/stress/ManagerStressMetricsTest.java
new file mode 100644
index 00000000..c7e743ce
--- /dev/null
+++ b/src/test/java/org/apache/accumulo/testing/manager/stress/ManagerStressMetricsTest.java
@@ -0,0 +1,72 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * https://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+package org.apache.accumulo.testing.manager.stress;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+
+import java.io.IOException;
+import java.nio.file.Path;
+import java.util.List;
+
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.io.TempDir;
+
+class ManagerStressMetricsTest {
+
+ @TempDir
+ Path temporaryDirectory;
+
+ @Test
+ void snapshotsAreReadAndAggregatedAcrossWorkers() throws IOException {
+ ManagerStressMetrics workerOne = new ManagerStressMetrics();
+ ManagerStressMetrics.Counts create = workerOne.counts(Operation.CREATE);
+ create.attempts = 4;
+ create.submitted = 3;
+ create.completed = 2;
+ create.skipped = 1;
+ create.failed = 1;
+ create.outcomeUnknown = 1;
+ create.totalNanos = 40;
+ workerOne.addMaintenanceCreates(5);
+ workerOne.writeSnapshot(temporaryDirectory, 0);
+
+ ManagerStressMetrics workerTwo = new ManagerStressMetrics();
+ ManagerStressMetrics.Counts compact = workerTwo.counts(Operation.COMPACT);
+ compact.attempts = 2;
+ compact.submitted = 2;
+ compact.asyncAccepted = 2;
+ compact.totalNanos = 60;
+ workerTwo.addMaintenanceCreates(7);
+ workerTwo.writeSnapshot(temporaryDirectory, 1);
+
+ var first =
+ ManagerStressMetrics.readSnapshot(ManagerStressMetrics.snapshotPath(temporaryDirectory, 0));
+ var second =
+ ManagerStressMetrics.readSnapshot(ManagerStressMetrics.snapshotPath(temporaryDirectory, 1));
+ var totals = ManagerStressMetrics.aggregate(List.of(first, second));
+
+ assertEquals(2, totals.workers());
+ assertEquals(12, totals.maintenanceCreates());
+ assertEquals(4, totals.operations().get(Operation.CREATE).attempts());
+ assertEquals(3, totals.operations().get(Operation.CREATE).submitted());
+ assertEquals(2, totals.operations().get(Operation.CREATE).completed());
+ assertEquals(1, totals.operations().get(Operation.CREATE).outcomeUnknown());
+ assertEquals(2, totals.operations().get(Operation.COMPACT).asyncAccepted());
+ }
+}
diff --git a/src/test/java/org/apache/accumulo/testing/manager/stress/ManagerStressOptionsTest.java b/src/test/java/org/apache/accumulo/testing/manager/stress/ManagerStressOptionsTest.java
new file mode 100644
index 00000000..335534b6
--- /dev/null
+++ b/src/test/java/org/apache/accumulo/testing/manager/stress/ManagerStressOptionsTest.java
@@ -0,0 +1,73 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * https://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+package org.apache.accumulo.testing.manager.stress;
+
+import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+
+class ManagerStressOptionsTest {
+
+ @org.junit.jupiter.api.Test
+ void requiresNamespace() {
+ ManagerStressOptions options = validOptions();
+ options.namespace = null;
+
+ assertThrows(IllegalArgumentException.class, () -> options.validate(true));
+ }
+
+ @org.junit.jupiter.api.Test
+ void validatesNamespaceCharacters() {
+ ManagerStressOptions options = validOptions();
+ options.namespace = "invalid.namespace";
+
+ assertThrows(IllegalArgumentException.class, () -> options.validate(true));
+ }
+
+ @org.junit.jupiter.api.Test
+ void acceptsValidNamespace() {
+ assertDoesNotThrow(() -> validOptions().validate(true));
+ }
+
+ @org.junit.jupiter.api.Test
+ void requiresCompactorResourceGroup() {
+ ManagerStressOptions options = validOptions();
+ options.compactorResourceGroup = null;
+
+ assertThrows(IllegalArgumentException.class, () -> options.validate(true));
+ }
+
+ @org.junit.jupiter.api.Test
+ void configuresMgrstressCompactionPlannerForSelectedGroup() {
+ var properties = ManagerStress.compactionServiceProperties("stressGroup");
+
+ assertEquals("org.apache.accumulo.core.spi.compaction.RatioBasedCompactionPlanner",
+ properties.get("compaction.service.mgrstress.planner"));
+ assertEquals("[{\"group\":\"stressGroup\"}]",
+ properties.get("compaction.service.mgrstress.planner.opts.groups"));
+ }
+
+ private static ManagerStressOptions validOptions() {
+ ManagerStressOptions options = new ManagerStressOptions();
+ options.namespace = "stress_namespace_1";
+ options.compactorResourceGroup = "default";
+ options.hdfsDir = "/tmp/manager-stress-test";
+ return options;
+ }
+}
diff --git a/src/test/java/org/apache/accumulo/testing/manager/stress/ManagerStressRowsTest.java b/src/test/java/org/apache/accumulo/testing/manager/stress/ManagerStressRowsTest.java
new file mode 100644
index 00000000..4d55e1bd
--- /dev/null
+++ b/src/test/java/org/apache/accumulo/testing/manager/stress/ManagerStressRowsTest.java
@@ -0,0 +1,119 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * https://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+package org.apache.accumulo.testing.manager.stress;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+import java.util.List;
+import java.util.Random;
+import java.util.TreeSet;
+
+import org.apache.accumulo.core.client.admin.TabletAvailability;
+import org.apache.accumulo.core.data.RowRange;
+import org.apache.hadoop.io.Text;
+
+class ManagerStressRowsTest {
+
+ @org.junit.jupiter.api.Test
+ void createsUnhostedTablesWithBoundedInitialSplitCounts() {
+ var config = ManagerStressRows.newTableConfiguration(new Random(7));
+
+ assertEquals(TabletAvailability.UNHOSTED, config.getInitialTabletAvailability());
+ assertTrue(config.getSplits().size() >= ManagerStressRows.MIN_INITIAL_SPLITS);
+ assertTrue(config.getSplits().size() <= ManagerStressRows.MAX_INITIAL_SPLITS);
+ assertEquals(config.getSplits().size(), new TreeSet<>(config.getSplits()).size());
+ config.getSplits().forEach(split -> assertTrue(split.toString().matches("[0-9a-f]{16}")));
+ }
+
+ @org.junit.jupiter.api.Test
+ void addedSplitsRespectTheThousandTabletLimit() {
+ var additions = ManagerStressRows.additionalSplits(List.of(), new Random(5));
+ assertFalse(additions.isEmpty());
+ assertTrue(additions.size() + 1 <= ManagerStressRows.MAX_TABLETS);
+
+ var nearlyFull = ManagerStressRows.initialSplits(998);
+ var finalSplit = ManagerStressRows.additionalSplits(nearlyFull, new Random(11));
+ assertEquals(1, finalSplit.size());
+ assertTrue(finalSplit.stream().noneMatch(nearlyFull::contains));
+
+ var full = ManagerStressRows.initialSplits(ManagerStressRows.MAX_INITIAL_SPLITS);
+ assertTrue(ManagerStressRows.additionalSplits(full, new Random(13)).isEmpty());
+ }
+
+ @org.junit.jupiter.api.Test
+ void mergeRangesUseTheAdjacentSplitBoundaries() {
+ List splits = List.of(new Text("10"), new Text("20"), new Text("30"), new Text("40"));
+
+ var first = ManagerStressRows.mergeRange(splits, 0, 2);
+ assertNull(first.start());
+ assertEquals(new Text("20"), first.end());
+
+ var middle = ManagerStressRows.mergeRange(splits, 1, 2);
+ assertEquals(new Text("10"), middle.start());
+ assertEquals(new Text("30"), middle.end());
+
+ var last = ManagerStressRows.mergeRange(splits, 3, 2);
+ assertEquals(new Text("30"), last.start());
+ assertNull(last.end());
+
+ assertEquals(2, ManagerStressRows.selectMergeRange(splits, 1, new Random(17)).tabletCount());
+ assertEquals(5, ManagerStressRows.selectMergeRange(splits, 100, new Random(19)).tabletCount());
+ }
+
+ @org.junit.jupiter.api.Test
+ void availabilityRangesMapExactlyToTabletBoundaries() {
+ List splits = List.of(new Text("10"), new Text("20"), new Text("30"), new Text("40"));
+
+ assertEquals(RowRange.atMost(new Text("10")), ManagerStressRows.tabletRowRange(splits, 0, 0));
+ assertEquals(RowRange.atMost(new Text("30")), ManagerStressRows.tabletRowRange(splits, 0, 2));
+ assertEquals(RowRange.openClosed(new Text("10"), new Text("30")),
+ ManagerStressRows.tabletRowRange(splits, 1, 2));
+ assertEquals(RowRange.greaterThan(new Text("30")),
+ ManagerStressRows.tabletRowRange(splits, 3, 4));
+ assertEquals(RowRange.all(), ManagerStressRows.tabletRowRange(splits, 0, 4));
+ }
+
+ @org.junit.jupiter.api.Test
+ void availabilityChangesSelectUniqueTabletsAndUseDifferentStates() {
+ List splits = List.of(new Text("10"), new Text("20"), new Text("30"), new Text("40"));
+ List current =
+ List.of(TabletAvailability.ONDEMAND, TabletAvailability.ONDEMAND,
+ TabletAvailability.ONDEMAND, TabletAvailability.ONDEMAND, TabletAvailability.ONDEMAND);
+
+ var changes = ManagerStressRows.randomAvailabilityChanges(splits, current, new Random(43));
+ int selectedCount =
+ changes.stream().mapToInt(ManagerStressRows.AvailabilityChange::tabletCount).sum();
+
+ assertTrue(selectedCount >= 1);
+ assertTrue(selectedCount <= current.size());
+ for (int i = 0; i < changes.size(); i++) {
+ var change = changes.get(i);
+ assertTrue(change.availability() != TabletAvailability.ONDEMAND);
+ assertEquals(ManagerStressRows.tabletRowRange(splits, change.firstTablet(),
+ change.firstTablet() + change.tabletCount() - 1), change.rowRange());
+ if (i > 0) {
+ var previous = changes.get(i - 1);
+ assertTrue(previous.firstTablet() + previous.tabletCount() <= change.firstTablet());
+ }
+ }
+ }
+}
diff --git a/src/test/java/org/apache/accumulo/testing/manager/stress/OperationTest.java b/src/test/java/org/apache/accumulo/testing/manager/stress/OperationTest.java
new file mode 100644
index 00000000..0521a293
--- /dev/null
+++ b/src/test/java/org/apache/accumulo/testing/manager/stress/OperationTest.java
@@ -0,0 +1,56 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * https://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+package org.apache.accumulo.testing.manager.stress;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+
+class OperationTest {
+
+ @org.junit.jupiter.api.Test
+ void choosesOperationByRelativeWeight() {
+ ManagerStressOptions options = new ManagerStressOptions();
+ options.createWeight = 0;
+ options.deleteWeight = 2;
+ options.splitWeight = 0;
+ options.mergeWeight = 3;
+ options.availabilityWeight = 0;
+ options.compactWeight = 0;
+ options.bulkImportWeight = 0;
+
+ assertEquals(Operation.DELETE, Operation.choose(options, 0));
+ assertEquals(Operation.DELETE, Operation.choose(options, 1.99));
+ assertEquals(Operation.MERGE, Operation.choose(options, 2));
+ assertEquals(Operation.MERGE, Operation.choose(options, 4.99));
+ }
+
+ @org.junit.jupiter.api.Test
+ void choosesAvailabilityByRelativeWeight() {
+ ManagerStressOptions options = new ManagerStressOptions();
+ options.createWeight = 0;
+ options.deleteWeight = 0;
+ options.splitWeight = 0;
+ options.mergeWeight = 0;
+ options.availabilityWeight = 4;
+ options.compactWeight = 0;
+ options.bulkImportWeight = 0;
+
+ assertEquals(Operation.AVAILABILITY, Operation.choose(options, 0));
+ assertEquals(Operation.AVAILABILITY, Operation.choose(options, 3.99));
+ }
+}
diff --git a/src/test/java/org/apache/accumulo/testing/manager/stress/TablePoolTest.java b/src/test/java/org/apache/accumulo/testing/manager/stress/TablePoolTest.java
new file mode 100644
index 00000000..3144220a
--- /dev/null
+++ b/src/test/java/org/apache/accumulo/testing/manager/stress/TablePoolTest.java
@@ -0,0 +1,33 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * https://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+package org.apache.accumulo.testing.manager.stress;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+
+import java.nio.file.Path;
+
+class TablePoolTest {
+
+ @org.junit.jupiter.api.Test
+ void qualifiesGeneratedTableNamesWithTheNamespace() {
+ TablePool pool = new TablePool(Path.of("unused"), "stress_namespace", "stress", 2, 17);
+
+ assertEquals("stress_namespace.stress_3_7", pool.tableName(3, 7));
+ }
+}
diff --git a/src/test/java/org/apache/accumulo/testing/manager/stress/compactor/CompactionOutcomeTest.java b/src/test/java/org/apache/accumulo/testing/manager/stress/compactor/CompactionOutcomeTest.java
new file mode 100644
index 00000000..dab5495c
--- /dev/null
+++ b/src/test/java/org/apache/accumulo/testing/manager/stress/compactor/CompactionOutcomeTest.java
@@ -0,0 +1,37 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * https://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+package org.apache.accumulo.testing.manager.stress.compactor;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+
+class CompactionOutcomeTest {
+
+ @org.junit.jupiter.api.Test
+ void choosesConfiguredRelativeWeights() {
+ StressCompactorOptions options = new StressCompactorOptions();
+ options.successWeight = 0;
+ options.failureWeight = 2;
+ options.cancellationWeight = 3;
+
+ assertEquals(CompactionOutcome.FAILURE, CompactionOutcome.choose(options, 0));
+ assertEquals(CompactionOutcome.FAILURE, CompactionOutcome.choose(options, 1.99));
+ assertEquals(CompactionOutcome.CANCELLATION, CompactionOutcome.choose(options, 2));
+ assertEquals(CompactionOutcome.CANCELLATION, CompactionOutcome.choose(options, 4.99));
+ }
+}