This is an automated email from the ASF dual-hosted git repository.

dsmiley pushed a commit to branch branch_10x
in repository https://gitbox.apache.org/repos/asf/solr.git

commit 4c9603c9a6d9b971afeeedb067b4e049557aa7f3
Author: Eric Pugh <[email protected]>
AuthorDate: Fri Aug 28 14:15:02 2026 -0400

    TIdy cross DC code.  The application and the solr module. (#4744)
    
    (cherry picked from commit ae262ba33e4ebfbe92e0ba591b4d1ecee26e6890)
---
 .../crossdc/manager/consumer/ConsumerMetrics.java  |  2 +-
 .../crossdc/manager/consumer/PartitionManager.java | 16 ---------
 .../solr/crossdc/manager/consumer/ThreadDump.java  |  6 ++--
 .../manager/consumer/ThreadDumpServlet.java        |  2 +-
 .../messageprocessor/SolrMessageProcessor.java     | 41 +---------------------
 .../crossdc/manager/DeleteByQueryToIdTest.java     |  1 -
 .../crossdc/manager/RetryQueueIntegrationTest.java | 10 +++---
 .../manager/SolrAndKafkaIntegrationTest.java       |  6 ++--
 ...SolrAndKafkaMultiCollectionIntegrationTest.java |  1 -
 .../crossdc/manager/SolrAndKafkaReindexTest.java   |  5 ++-
 .../crossdc/manager/ZkConfigIntegrationTest.java   |  1 -
 .../manager/consumer/KafkaCrossDcConsumerTest.java | 27 ++------------
 .../messageprocessor/TestMessageProcessor.java     | 11 +++++-
 .../apache/solr/crossdc/common/ConfigProperty.java | 10 ------
 .../apache/solr/crossdc/common/IQueueHandler.java  |  2 +-
 .../solr/crossdc/common/KafkaCrossDcConf.java      |  6 ++--
 .../solr/crossdc/common/KafkaMirroringSink.java    |  5 ++-
 .../solr/crossdc/common/MirroredSolrRequest.java   | 21 +++++------
 .../update/processor/MirroringException.java       | 36 -------------------
 .../update/processor/MirroringUpdateProcessor.java | 10 ++----
 .../MirroringUpdateRequestProcessorFactory.java    | 15 +-------
 .../apache/solr/crossdc/common/ConfUtilTest.java   |  2 +-
 .../common/MirroredSolrRequestSerializerTest.java  |  4 +--
 .../handler/MirroringCollectionsHandlerTest.java   |  3 +-
 .../handler/MirroringConfigSetsHandlerTest.java    |  6 ++--
 .../processor/MirroringUpdateProcessorTest.java    | 14 ++++----
 26 files changed, 58 insertions(+), 205 deletions(-)

diff --git 
a/solr/cross-dc-manager/src/java/org/apache/solr/crossdc/manager/consumer/ConsumerMetrics.java
 
b/solr/cross-dc-manager/src/java/org/apache/solr/crossdc/manager/consumer/ConsumerMetrics.java
index 92d96fae866..7d5f11a427e 100644
--- 
a/solr/cross-dc-manager/src/java/org/apache/solr/crossdc/manager/consumer/ConsumerMetrics.java
+++ 
b/solr/cross-dc-manager/src/java/org/apache/solr/crossdc/manager/consumer/ConsumerMetrics.java
@@ -131,7 +131,7 @@ public interface ConsumerMetrics {
 
   /**
    * Records the batch size of the output request. Batch size is defined as 
the number of operations
-   * in an output {@link SolrRequest} (which may be different than the input 
size due to
+   * in an output {@link SolrRequest} (which may be different from the input 
size due to
    * collapsing).
    *
    * @param type the type of the request, corresponding to one of the {@link
diff --git 
a/solr/cross-dc-manager/src/java/org/apache/solr/crossdc/manager/consumer/PartitionManager.java
 
b/solr/cross-dc-manager/src/java/org/apache/solr/crossdc/manager/consumer/PartitionManager.java
index c93740f25ea..cef9bfefb98 100644
--- 
a/solr/cross-dc-manager/src/java/org/apache/solr/crossdc/manager/consumer/PartitionManager.java
+++ 
b/solr/cross-dc-manager/src/java/org/apache/solr/crossdc/manager/consumer/PartitionManager.java
@@ -119,22 +119,6 @@ public class PartitionManager {
     }
   }
 
-  /**
-   * Reset the local offset so that the consumer reads the records from Kafka 
again.
-   *
-   * @param partition The TopicPartition to reset the offset for
-   * @param partitionRecords PartitionRecords for the specified partition
-   */
-  private void resetOffsetForPartition(
-      TopicPartition partition,
-      List<ConsumerRecord<String, MirroredSolrRequest<?>>> partitionRecords) {
-    if (log.isTraceEnabled()) {
-      log.trace("Resetting offset to: {}", partitionRecords.get(0).offset());
-    }
-    long resetOffset = partitionRecords.get(0).offset();
-    consumer.seek(partition, resetOffset);
-  }
-
   /**
    * Logs and updates the commit point for the partition that has been 
processed.
    *
diff --git 
a/solr/cross-dc-manager/src/java/org/apache/solr/crossdc/manager/consumer/ThreadDump.java
 
b/solr/cross-dc-manager/src/java/org/apache/solr/crossdc/manager/consumer/ThreadDump.java
index 2f0fb116bf9..aed28f2c952 100644
--- 
a/solr/cross-dc-manager/src/java/org/apache/solr/crossdc/manager/consumer/ThreadDump.java
+++ 
b/solr/cross-dc-manager/src/java/org/apache/solr/crossdc/manager/consumer/ThreadDump.java
@@ -42,7 +42,7 @@ public class ThreadDump {
   }
 
   /**
-   * Dumps all of the threads' current information, including synchronization, 
to an output stream.
+   * Dumps all the threads' current information, including synchronization, to 
an output stream.
    *
    * @param out an output stream
    */
@@ -51,8 +51,8 @@ public class ThreadDump {
   }
 
   /**
-   * Dumps all of the threads' current information, optionally including 
synchronization, to an
-   * output stream.
+   * Dumps all the threads' current information, optionally including 
synchronization, to an output
+   * stream.
    *
    * <p>Having control over including synchronization info allows using this 
method (and its
    * wrappers, i.e. ThreadDumpServlet) in environments where getting object 
monitor and/or ownable
diff --git 
a/solr/cross-dc-manager/src/java/org/apache/solr/crossdc/manager/consumer/ThreadDumpServlet.java
 
b/solr/cross-dc-manager/src/java/org/apache/solr/crossdc/manager/consumer/ThreadDumpServlet.java
index cea2bb33988..7df3e16b0e2 100644
--- 
a/solr/cross-dc-manager/src/java/org/apache/solr/crossdc/manager/consumer/ThreadDumpServlet.java
+++ 
b/solr/cross-dc-manager/src/java/org/apache/solr/crossdc/manager/consumer/ThreadDumpServlet.java
@@ -45,7 +45,7 @@ public class ThreadDumpServlet extends HttpServlet {
   @Override
   public void init() throws ServletException {
     try {
-      // Some PaaS like Google App Engine blacklist java.lang.managament
+      // Some PaaS like Google App Engine blacklist java.lang.management
       this.threadDump = new ThreadDump(ManagementFactory.getThreadMXBean());
     } catch (NoClassDefFoundError ncdfe) {
       this.threadDump = null; // we won't be able to provide thread dump
diff --git 
a/solr/cross-dc-manager/src/java/org/apache/solr/crossdc/manager/messageprocessor/SolrMessageProcessor.java
 
b/solr/cross-dc-manager/src/java/org/apache/solr/crossdc/manager/messageprocessor/SolrMessageProcessor.java
index 98f0bc5e8d4..ad9033fe1b3 100644
--- 
a/solr/cross-dc-manager/src/java/org/apache/solr/crossdc/manager/messageprocessor/SolrMessageProcessor.java
+++ 
b/solr/cross-dc-manager/src/java/org/apache/solr/crossdc/manager/messageprocessor/SolrMessageProcessor.java
@@ -31,10 +31,8 @@ import org.apache.solr.common.SolrInputField;
 import org.apache.solr.common.cloud.ClusterState;
 import org.apache.solr.common.cloud.DocCollection;
 import org.apache.solr.common.cloud.ZkStateReader;
-import org.apache.solr.common.params.ModifiableSolrParams;
 import org.apache.solr.common.params.SolrParams;
 import org.apache.solr.common.util.TimeSource;
-import org.apache.solr.crossdc.common.CrossDcConstants;
 import org.apache.solr.crossdc.common.IQueueHandler;
 import org.apache.solr.crossdc.common.MirroredSolrRequest;
 import org.apache.solr.crossdc.common.ResubmitBackoffPolicy;
@@ -165,7 +163,7 @@ public class SolrMessageProcessor extends MessageProcessor
   }
 
   private void logIf4xxException(SolrException solrException) {
-    // This shouldn't really happen but if it doesn, it most likely requires 
fixing in the return
+    // This shouldn't really happen but if it does, it most likely requires 
fixing in the return
     // code from Solr.
     if (solrException != null && 400 <= solrException.code() && 
solrException.code() < 500) {
       log.error("Exception occurred with 4xx response. {}", 
solrException.code(), solrException);
@@ -335,43 +333,6 @@ public class SolrMessageProcessor extends MessageProcessor
     }
   }
 
-  /**
-   * Adds {@link CrossDcConstants#SHOULD_MIRROR}=false to the params if it's 
not already specified.
-   * Logs a warning if it is specified and NOT set to false. (i.e. circular 
mirror may occur)
-   *
-   * @param mirroredSolrRequest MirroredSolrRequest object that is being 
processed.
-   */
-  void preventCircularMirroring(MirroredSolrRequest<?> mirroredSolrRequest) {
-    if (mirroredSolrRequest.getSolrRequest() instanceof UpdateRequest 
updateRequest) {
-      ModifiableSolrParams params = updateRequest.getParams();
-      String shouldMirror = (params == null ? null : 
params.get(CrossDcConstants.SHOULD_MIRROR));
-      if (shouldMirror == null) {
-        log.warn(
-            "{} param is missing - setting to false. Request={}",
-            CrossDcConstants.SHOULD_MIRROR,
-            mirroredSolrRequest);
-        updateRequest.setParam(CrossDcConstants.SHOULD_MIRROR, "false");
-      } else if (!"false".equalsIgnoreCase(shouldMirror)) {
-        log.warn("{} param equal to {}", CrossDcConstants.SHOULD_MIRROR, 
shouldMirror);
-      }
-    } else {
-      SolrParams params = mirroredSolrRequest.getSolrRequest().getParams();
-      assert params != null;
-      String shouldMirror = params.get(CrossDcConstants.SHOULD_MIRROR);
-      if (shouldMirror == null) {
-        if (params instanceof ModifiableSolrParams) {
-          log.warn("{} param is missing - setting to false", 
CrossDcConstants.SHOULD_MIRROR);
-          ((ModifiableSolrParams) params).set(CrossDcConstants.SHOULD_MIRROR, 
"false");
-        } else {
-          log.warn(
-              "{} param is missing and params are not modifiable", 
CrossDcConstants.SHOULD_MIRROR);
-        }
-      } else if (!"false".equalsIgnoreCase(shouldMirror)) {
-        log.warn("{} param is present and set to {}", 
CrossDcConstants.SHOULD_MIRROR, shouldMirror);
-      }
-    }
-  }
-
   private void connectToSolrIfNeeded() {
     // Don't try to consume anything if we can't connect to the solr server
     boolean connected = false;
diff --git 
a/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/DeleteByQueryToIdTest.java
 
b/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/DeleteByQueryToIdTest.java
index 47a9a7fb024..216435999dc 100644
--- 
a/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/DeleteByQueryToIdTest.java
+++ 
b/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/DeleteByQueryToIdTest.java
@@ -60,7 +60,6 @@ import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
 @ThreadLeakFilters(
-    defaultFilters = true,
     filters = {
       SolrIgnoredThreadsFilter.class,
       QuickPatchThreadsFilter.class,
diff --git 
a/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/RetryQueueIntegrationTest.java
 
b/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/RetryQueueIntegrationTest.java
index e69171a01cd..621080e920f 100644
--- 
a/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/RetryQueueIntegrationTest.java
+++ 
b/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/RetryQueueIntegrationTest.java
@@ -125,10 +125,8 @@ public class RetryQueueIntegrationTest extends 
SolrTestCaseJ4 {
       zkTestServer2.run();
     }
 
-    solrCluster1 = startCluster(solrCluster1, zkTestServer1, baseDir1);
-    solrCluster2 = startCluster(solrCluster2, zkTestServer2, baseDir2);
-
-    CloudSolrClient client = solrCluster1.getSolrClient(COLLECTION);
+    solrCluster1 = startCluster(zkTestServer1, baseDir1);
+    solrCluster2 = startCluster(zkTestServer2, baseDir2);
 
     log.info("bootstrapServers={}", bootstrapServers);
 
@@ -141,8 +139,8 @@ public class RetryQueueIntegrationTest extends 
SolrTestCaseJ4 {
     consumer.start(properties);
   }
 
-  private static MiniSolrCloudCluster startCluster(
-      MiniSolrCloudCluster solrCluster, ZkTestServer zkTestServer, Path 
baseDir) throws Exception {
+  private static MiniSolrCloudCluster startCluster(ZkTestServer zkTestServer, 
Path baseDir)
+      throws Exception {
     MiniSolrCloudCluster cluster =
         new MiniSolrCloudCluster(
             1,
diff --git 
a/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/SolrAndKafkaIntegrationTest.java
 
b/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/SolrAndKafkaIntegrationTest.java
index 5206def632a..f3c107bea97 100644
--- 
a/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/SolrAndKafkaIntegrationTest.java
+++ 
b/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/SolrAndKafkaIntegrationTest.java
@@ -582,14 +582,14 @@ public class SolrAndKafkaIntegrationTest extends 
SolrCloudTestCase {
               (InputStream) rsp.get(InputStreamResponseParser.STREAM_KEY), 
StandardCharsets.UTF_8);
       assertTrue(content, 
content.contains("solr_crossdc_consumer_output_total"));
 
-      // test the healtcheck endpoint
+      // test the healthcheck endpoint
       req = new GenericSolrRequest(SolrRequest.METHOD.GET, "/health");
       req.setResponseParser(new InputStreamResponseParser(null));
       rsp = httpJettySolrClient.request(req);
       content =
           IOUtils.toString(
               (InputStream) rsp.get(InputStreamResponseParser.STREAM_KEY), 
StandardCharsets.UTF_8);
-      assertEquals(Integer.valueOf(200), rsp.get("responseStatus"));
+      assertEquals(200, rsp.get("responseStatus"));
       Map<String, Object> map = (Map<String, Object>) 
ObjectBuilder.fromJSON(content);
       assertEquals(Boolean.TRUE, map.get("kafka"));
       assertEquals(Boolean.TRUE, map.get("solr"));
@@ -603,7 +603,7 @@ public class SolrAndKafkaIntegrationTest extends 
SolrCloudTestCase {
       content =
           IOUtils.toString(
               (InputStream) rsp.get(InputStreamResponseParser.STREAM_KEY), 
StandardCharsets.UTF_8);
-      assertEquals(Integer.valueOf(503), rsp.get("responseStatus"));
+      assertEquals(503, rsp.get("responseStatus"));
       map = (Map<String, Object>) ObjectBuilder.fromJSON(content);
       assertEquals(Boolean.TRUE, map.get("kafka"));
       assertEquals(Boolean.FALSE, map.get("solr"));
diff --git 
a/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/SolrAndKafkaMultiCollectionIntegrationTest.java
 
b/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/SolrAndKafkaMultiCollectionIntegrationTest.java
index 5e0795d4ff8..233694c7402 100644
--- 
a/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/SolrAndKafkaMultiCollectionIntegrationTest.java
+++ 
b/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/SolrAndKafkaMultiCollectionIntegrationTest.java
@@ -61,7 +61,6 @@ import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
 @ThreadLeakFilters(
-    defaultFilters = true,
     filters = {
       SolrIgnoredThreadsFilter.class,
       QuickPatchThreadsFilter.class,
diff --git 
a/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/SolrAndKafkaReindexTest.java
 
b/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/SolrAndKafkaReindexTest.java
index 7e1cc540b6d..38dddf5090f 100644
--- 
a/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/SolrAndKafkaReindexTest.java
+++ 
b/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/SolrAndKafkaReindexTest.java
@@ -50,7 +50,6 @@ import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
 @ThreadLeakFilters(
-    defaultFilters = true,
     filters = {
       SolrIgnoredThreadsFilter.class,
       QuickPatchThreadsFilter.class,
@@ -122,7 +121,7 @@ public class SolrAndKafkaReindexTest extends 
SolrCloudTestCase {
   }
 
   @AfterClass
-  public static void afterSolrAndKafkaIntegrationTest() throws Exception {
+  public static void afterSolrAndKafkaIntegrationTest() {
     ObjectReleaseTracker.clear();
 
     if (solrCluster1 != null) {
@@ -254,7 +253,7 @@ public class SolrAndKafkaReindexTest extends 
SolrCloudTestCase {
     doc2.addField("id", id2);
     doc2.addField("text", "some test two " + tag);
 
-    List<SolrInputDocument> docs = new ArrayList<SolrInputDocument>(2);
+    List<SolrInputDocument> docs = new ArrayList<>(2);
     docs.add(doc1);
     docs.add(doc2);
 
diff --git 
a/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/ZkConfigIntegrationTest.java
 
b/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/ZkConfigIntegrationTest.java
index d4daa6f8040..da589afa482 100644
--- 
a/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/ZkConfigIntegrationTest.java
+++ 
b/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/ZkConfigIntegrationTest.java
@@ -49,7 +49,6 @@ import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
 @ThreadLeakFilters(
-    defaultFilters = true,
     filters = {
       SolrIgnoredThreadsFilter.class,
       QuickPatchThreadsFilter.class,
diff --git 
a/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/consumer/KafkaCrossDcConsumerTest.java
 
b/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/consumer/KafkaCrossDcConsumerTest.java
index fcbd2edd84d..2bc50db48fd 100644
--- 
a/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/consumer/KafkaCrossDcConsumerTest.java
+++ 
b/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/consumer/KafkaCrossDcConsumerTest.java
@@ -151,29 +151,6 @@ public class KafkaCrossDcConsumerTest {
     kafkaCrossDcConsumer.shutdown();
   }
 
-  private ConsumerRecord<String, MirroredSolrRequest<?>> 
createSampleConsumerRecord() {
-    return new ConsumerRecord<>("sample-topic", 0, 0, "key", 
createSampleMirroredSolrRequest());
-  }
-
-  private ConsumerRecords<String, MirroredSolrRequest<?>> 
createSampleConsumerRecords() {
-    TopicPartition topicPartition = new TopicPartition("sample-topic", 0);
-    List<ConsumerRecord<String, MirroredSolrRequest<?>>> recordsList = new 
ArrayList<>();
-    recordsList.add(
-        new ConsumerRecord<>("sample-topic", 0, 0, "key", 
createSampleMirroredSolrRequest()));
-    return new ConsumerRecords<>(Map.of(topicPartition, recordsList));
-  }
-
-  private MirroredSolrRequest<?> createSampleMirroredSolrRequest() {
-    // Create a sample MirroredSolrRequest for testing
-    SolrInputDocument solrInputDocument = new SolrInputDocument();
-    solrInputDocument.addField("id", "1");
-    solrInputDocument.addField("title", "Sample title");
-    solrInputDocument.addField("content", "Sample content");
-    UpdateRequest updateRequest = new UpdateRequest();
-    updateRequest.add(solrInputDocument);
-    return new MirroredSolrRequest<>(updateRequest);
-  }
-
   /** Should create a KafkaCrossDcConsumer with the given configuration and 
startLatch */
   @Test
   public void kafkaCrossDcConsumerCreationWithConfigurationAndStartLatch() {
@@ -204,7 +181,7 @@ public class KafkaCrossDcConsumerTest {
   }
 
   @Test
-  public void testSolrClientSupplier() throws Exception {
+  public void testSolrClientSupplier() {
     supplier.get();
     assertEquals(1, solrClientCounter.get());
     clusterStateProviderIsClosed = true;
@@ -331,7 +308,7 @@ public class KafkaCrossDcConsumerTest {
   }
 
   @Test
-  public void testHandleValidAdminRequest() throws Exception {
+  public void testHandleValidAdminRequest() {
     KafkaConsumer<String, MirroredSolrRequest<?>> mockConsumer = 
mock(KafkaConsumer.class);
     KafkaCrossDcConsumer spyConsumer = createCrossDcConsumerSpy(mockConsumer);
     doReturn(new IQueueHandler.Result<>(IQueueHandler.ResultStatus.HANDLED, 
null))
diff --git 
a/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/messageprocessor/TestMessageProcessor.java
 
b/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/messageprocessor/TestMessageProcessor.java
index 5303d8f2323..f7a82d2cd5b 100644
--- 
a/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/messageprocessor/TestMessageProcessor.java
+++ 
b/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/messageprocessor/TestMessageProcessor.java
@@ -41,6 +41,7 @@ import org.apache.solr.crossdc.common.MirroredSolrRequest;
 import org.apache.solr.crossdc.common.ResubmitBackoffPolicy;
 import org.apache.solr.crossdc.manager.consumer.ConsumerMetrics;
 import org.apache.solr.crossdc.manager.consumer.OtelMetrics;
+import org.junit.After;
 import org.junit.Before;
 import org.junit.BeforeClass;
 import org.junit.Ignore;
@@ -54,6 +55,7 @@ public class TestMessageProcessor {
 
   @Mock private CloudSolrClient solrClient;
   private SolrMessageProcessor processor;
+  private AutoCloseable mocks;
 
   private final ResubmitBackoffPolicy backoffPolicy =
       spy(
@@ -71,7 +73,7 @@ public class TestMessageProcessor {
 
   @Before
   public void setUp() {
-    MockitoAnnotations.initMocks(this);
+    mocks = MockitoAnnotations.openMocks(this);
 
     // handleItem() probes the cluster through the state provider, so the mock 
must supply one
     ClusterStateProvider clusterStateProvider = 
Mockito.mock(ClusterStateProvider.class);
@@ -82,6 +84,13 @@ public class TestMessageProcessor {
     Mockito.doNothing().when(processor).uncheckedSleep(anyLong());
   }
 
+  @After
+  public void tearDown() throws Exception {
+    if (mocks != null) {
+      mocks.close();
+    }
+  }
+
   @Test
   public void testDocumentSanitization() {
     UpdateRequest request = spy(new UpdateRequest());
diff --git 
a/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/common/ConfigProperty.java
 
b/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/common/ConfigProperty.java
index fc2bd416ec6..660f8ce2d94 100644
--- 
a/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/common/ConfigProperty.java
+++ 
b/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/common/ConfigProperty.java
@@ -25,12 +25,6 @@ public class ConfigProperty {
 
   private boolean required = false;
 
-  public ConfigProperty(String key, String defaultValue, boolean required) {
-    this.key = key;
-    this.defaultValue = defaultValue;
-    this.required = required;
-  }
-
   public ConfigProperty(String key, String defaultValue) {
     this.key = key;
     this.defaultValue = defaultValue;
@@ -49,10 +43,6 @@ public class ConfigProperty {
     return required;
   }
 
-  public String getDefaultValue() {
-    return defaultValue;
-  }
-
   public String getValue(Map<?, ?> properties) {
     String val = (String) properties.get(key);
     if (val == null) {
diff --git 
a/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/common/IQueueHandler.java
 
b/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/common/IQueueHandler.java
index 145630deee7..c7f852e3716 100644
--- 
a/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/common/IQueueHandler.java
+++ 
b/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/common/IQueueHandler.java
@@ -21,7 +21,7 @@ public interface IQueueHandler<T> {
     /** Item was successfully processed */
     HANDLED,
 
-    /** Item was not processed, and the consumer should shutdown */
+    /** Item was not processed, and the consumer should shut down */
     NOT_HANDLED_SHUTDOWN,
 
     /** Item processing failed, and the item should be retried immediately */
diff --git 
a/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/common/KafkaCrossDcConf.java
 
b/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/common/KafkaCrossDcConf.java
index 02526216b36..88db372683e 100644
--- 
a/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/common/KafkaCrossDcConf.java
+++ 
b/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/common/KafkaCrossDcConf.java
@@ -225,7 +225,7 @@ public class KafkaCrossDcConf extends CrossDcConf {
   private final Map<String, Object> properties;
 
   public KafkaCrossDcConf(Map<String, Object> properties) {
-    List<String> nullValueKeys = new ArrayList<String>();
+    List<String> nullValueKeys = new ArrayList<>();
     properties.forEach(
         (k, v) -> {
           if (v == null) {
@@ -275,7 +275,7 @@ public class KafkaCrossDcConf extends CrossDcConf {
         (key, v) -> {
           try {
             int intVal = Integer.parseInt((String) v);
-            integerProperties.put(key.toString(), intVal);
+            integerProperties.put(key, intVal);
           } catch (NumberFormatException ignored) {
 
           }
@@ -312,7 +312,7 @@ public class KafkaCrossDcConf extends CrossDcConf {
         
sb.append(configProperty.getKey()).append("=").append(printablePropertyValue).append(",");
       }
     }
-    if (sb.length() > 0) {
+    if (!sb.isEmpty()) {
       sb.setLength(sb.length() - 1);
     }
 
diff --git 
a/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/common/KafkaMirroringSink.java
 
b/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/common/KafkaMirroringSink.java
index 70f48457814..c348f032c21 100644
--- 
a/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/common/KafkaMirroringSink.java
+++ 
b/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/common/KafkaMirroringSink.java
@@ -110,13 +110,12 @@ public class KafkaMirroringSink implements 
RequestMirroringSink, Closeable {
             }
           });
 
-      long lastSuccessfulEnqueueNanos = System.nanoTime();
       // Record time since last successful enqueue as 0
       long elapsedTimeMillis = TimeUnit.NANOSECONDS.toMillis(System.nanoTime() 
- enqueueStartNanos);
       // Update elapsed time
 
       if (elapsedTimeMillis > conf.getInt(SLOW_SUBMIT_THRESHOLD_MS)) {
-        slowSubmitAction(request, elapsedTimeMillis);
+        slowSubmitAction(elapsedTimeMillis);
       }
     } catch (Exception e) {
       // We are intentionally catching all exceptions, the expected exception 
form this function is
@@ -220,7 +219,7 @@ public class KafkaMirroringSink implements 
RequestMirroringSink, Closeable {
         kafkaConsumerProperties, new StringDeserializer(), new 
MirroredSolrRequestSerializer());
   }
 
-  private void slowSubmitAction(Object request, long elapsedTimeMillis) {
+  private void slowSubmitAction(long elapsedTimeMillis) {
     log.warn(
         "Enqueuing the request to Kafka took more than {} millis. 
enqueueElapsedTime={}",
         conf.get(KafkaCrossDcConf.SLOW_SUBMIT_THRESHOLD_MS),
diff --git 
a/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/common/MirroredSolrRequest.java
 
b/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/common/MirroredSolrRequest.java
index 7d27878f89b..4d29239adbe 100644
--- 
a/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/common/MirroredSolrRequest.java
+++ 
b/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/common/MirroredSolrRequest.java
@@ -49,7 +49,7 @@ public class MirroredSolrRequest<T extends SolrResponse> {
     CONFIGSET,
     UNKNOWN;
 
-    public static final Type get(String s) {
+    public static Type get(String s) {
       if (s == null) {
         return UNKNOWN;
       } else {
@@ -227,10 +227,6 @@ public class MirroredSolrRequest<T extends SolrResponse> {
     return submitTimeNanos;
   }
 
-  public void setSubmitTimeNanos(final long submitTimeNanos) {
-    this.submitTimeNanos = submitTimeNanos;
-  }
-
   public Type getType() {
     return type;
   }
@@ -264,15 +260,16 @@ public class MirroredSolrRequest<T extends SolrResponse> {
   public String toString() {
     final StringBuilder sb = new StringBuilder(getClass().getSimpleName() + 
"{type=");
     sb.append(type.toString());
-    sb.append(", method=" + solrRequest.getMethod());
-    sb.append(", params=" + solrRequest.getParams());
+    sb.append(", method=").append(solrRequest.getMethod());
+    sb.append(", params=").append(solrRequest.getParams());
     if (solrRequest instanceof UpdateRequest req) {
-      sb.append(", add=" + (req.getDocuments() != null ? 
req.getDocuments().size() : "0"));
-      sb.append(", del=" + (req.getDeleteByIdMap() != null ? 
req.getDeleteByIdMap().size() : "0"));
-      sb.append(", dbq=" + (req.getDeleteQuery() != null ? 
req.getDeleteQuery().size() : "0"));
+      sb.append(", add=").append(req.getDocuments() != null ? 
req.getDocuments().size() : "0");
+      sb.append(", del=")
+          .append(req.getDeleteByIdMap() != null ? 
req.getDeleteByIdMap().size() : "0");
+      sb.append(", dbq=").append(req.getDeleteQuery() != null ? 
req.getDeleteQuery().size() : "0");
     }
-    sb.append(", attempt=" + attempt);
-    sb.append(", submitTimeNanos=" + submitTimeNanos);
+    sb.append(", attempt=").append(attempt);
+    sb.append(", submitTimeNanos=").append(submitTimeNanos);
     sb.append('}');
     return sb.toString();
   }
diff --git 
a/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/update/processor/MirroringException.java
 
b/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/update/processor/MirroringException.java
deleted file mode 100644
index cf01fa83824..00000000000
--- 
a/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/update/processor/MirroringException.java
+++ /dev/null
@@ -1,36 +0,0 @@
-/*
- * 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
- *
- *     http://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.solr.crossdc.update.processor;
-
-/** Wrapper class for Mirroring exceptions. */
-public class MirroringException extends Exception {
-  public MirroringException() {
-    super();
-  }
-
-  public MirroringException(String message) {
-    super(message);
-  }
-
-  public MirroringException(String message, Throwable cause) {
-    super(message, cause);
-  }
-
-  public MirroringException(Throwable cause) {
-    super(cause);
-  }
-}
diff --git 
a/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/update/processor/MirroringUpdateProcessor.java
 
b/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/update/processor/MirroringUpdateProcessor.java
index 54ad07e3d24..284dd91ca7d 100644
--- 
a/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/update/processor/MirroringUpdateProcessor.java
+++ 
b/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/update/processor/MirroringUpdateProcessor.java
@@ -86,17 +86,11 @@ public class MirroringUpdateProcessor extends 
UpdateRequestProcessor {
   /** If true then commit commands are mirrored, otherwise they are processed 
only locally. */
   private final boolean mirrorCommits;
 
-  /** Controls the processing of Delete-By-Query requests.. */
+  /** Controls the processing of Delete-By-Query requests */
   private final CrossDcConf.ExpandDbq expandDbq;
 
   private final long maxMirroringDocSizeBytes;
 
-  /**
-   * The distributed processor downstream from us so we can establish if we're 
running on a leader
-   * shard
-   */
-  // private DistributedUpdateProcessor distProc;
-
   /** Distribution phase of the incoming requests */
   private DistributedUpdateProcessor.DistribPhase distribPhase;
 
@@ -255,7 +249,7 @@ public class MirroringUpdateProcessor extends 
UpdateRequestProcessor {
     super.processDelete(cmd); // let this throw to prevent mirroring invalid 
requests
     producerMetrics.getLocal().inc();
     if (doMirroring) {
-      boolean isLeader = false;
+      boolean isLeader;
       UpdateRequest mirrorRequest = createMirrorRequest();
       if (cmd.isDeleteById()) {
         // deleteById requests runs once per leader, so we just submit the 
request from the leader
diff --git 
a/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/update/processor/MirroringUpdateRequestProcessorFactory.java
 
b/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/update/processor/MirroringUpdateRequestProcessorFactory.java
index c4329b776f7..b159d060fc0 100644
--- 
a/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/update/processor/MirroringUpdateRequestProcessorFactory.java
+++ 
b/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/update/processor/MirroringUpdateRequestProcessorFactory.java
@@ -34,7 +34,6 @@ import java.io.IOException;
 import java.lang.invoke.MethodHandles;
 import java.util.HashMap;
 import java.util.Map;
-import java.util.Properties;
 import org.apache.solr.common.SolrException;
 import org.apache.solr.common.cloud.CollectionProperties;
 import org.apache.solr.common.cloud.SolrZkClient;
@@ -165,11 +164,7 @@ public class MirroringUpdateRequestProcessorFactory 
extends UpdateRequestProcess
       }
       String enabledVal = collectionProperties.get("crossdc.enabled");
       if (enabledVal != null) {
-        if (Boolean.parseBoolean(enabledVal.toString())) {
-          this.enabled = true;
-        } else {
-          this.enabled = false;
-        }
+        this.enabled = Boolean.parseBoolean(enabledVal);
       }
     } catch (Exception e) {
       log.error("Exception looking for CrossDC configuration in Zookeeper", e);
@@ -223,14 +218,6 @@ public class MirroringUpdateRequestProcessorFactory 
extends UpdateRequestProcess
     mirroringHandler = new KafkaRequestMirroringHandler(sink);
   }
 
-  private static Integer getIntegerPropValue(String name, Properties props) {
-    String value = props.getProperty(name);
-    if (value == null) {
-      return null;
-    }
-    return Integer.parseInt(value);
-  }
-
   @Override
   public UpdateRequestProcessor getInstance(
       final SolrQueryRequest req, final SolrQueryResponse rsp, final 
UpdateRequestProcessor next) {
diff --git 
a/solr/modules/cross-dc/src/test/org/apache/solr/crossdc/common/ConfUtilTest.java
 
b/solr/modules/cross-dc/src/test/org/apache/solr/crossdc/common/ConfUtilTest.java
index 28a98d6d7a1..53713cdfc12 100644
--- 
a/solr/modules/cross-dc/src/test/org/apache/solr/crossdc/common/ConfUtilTest.java
+++ 
b/solr/modules/cross-dc/src/test/org/apache/solr/crossdc/common/ConfUtilTest.java
@@ -291,7 +291,7 @@ public class ConfUtilTest extends SolrTestCaseJ4 {
   }
 
   @Test
-  public void testFillProperties_SecurityProperties() throws Exception {
+  public void testFillProperties_SecurityProperties() {
     Map<String, Object> properties = new HashMap<>();
 
     // Set security-related properties
diff --git 
a/solr/modules/cross-dc/src/test/org/apache/solr/crossdc/common/MirroredSolrRequestSerializerTest.java
 
b/solr/modules/cross-dc/src/test/org/apache/solr/crossdc/common/MirroredSolrRequestSerializerTest.java
index 2ce667fee64..e0afa71cd33 100644
--- 
a/solr/modules/cross-dc/src/test/org/apache/solr/crossdc/common/MirroredSolrRequestSerializerTest.java
+++ 
b/solr/modules/cross-dc/src/test/org/apache/solr/crossdc/common/MirroredSolrRequestSerializerTest.java
@@ -28,7 +28,7 @@ public class MirroredSolrRequestSerializerTest extends 
SolrTestCase {
   private static final byte[] EMPTY_ARR = new byte[3];
 
   @Test
-  public void testSerializationBufferOptimization() throws Exception {
+  public void testSerializationBufferOptimization() {
     MirroredSolrRequestSerializer serializer = new 
MirroredSolrRequestSerializer();
     UpdateRequest req = new UpdateRequest();
     SolrInputDocument doc = new SolrInputDocument();
@@ -50,7 +50,7 @@ public class MirroredSolrRequestSerializerTest extends 
SolrTestCase {
           (String)
               ((UpdateRequest) deserialized.getSolrRequest())
                   .getDocuments()
-                  .get(0)
+                  .getFirst()
                   .getFieldValue("test");
       assertEquals(fieldValue, deserValue);
     }
diff --git 
a/solr/modules/cross-dc/src/test/org/apache/solr/crossdc/handler/MirroringCollectionsHandlerTest.java
 
b/solr/modules/cross-dc/src/test/org/apache/solr/crossdc/handler/MirroringCollectionsHandlerTest.java
index d2a7dc33df2..a2d0537d57d 100644
--- 
a/solr/modules/cross-dc/src/test/org/apache/solr/crossdc/handler/MirroringCollectionsHandlerTest.java
+++ 
b/solr/modules/cross-dc/src/test/org/apache/solr/crossdc/handler/MirroringCollectionsHandlerTest.java
@@ -46,7 +46,6 @@ import org.mockito.ArgumentCaptor;
 import org.mockito.Mockito;
 
 @ThreadLeakFilters(
-    defaultFilters = true,
     filters = {
       SolrIgnoredThreadsFilter.class,
       QuickPatchThreadsFilter.class,
@@ -141,7 +140,7 @@ public class MirroringCollectionsHandlerTest extends 
SolrTestCaseJ4 {
       SolrParams mirroredParams = solrRequest.getParams();
       params.forEach(
           entry -> {
-            assertEquals(entry.getValue(), 
mirroredParams.getParams(entry.getKey()));
+            assertArrayEquals(entry.getValue(), 
mirroredParams.getParams(entry.getKey()));
           });
     } else {
       assertEquals(initialMirroredCount, captor.getAllValues().size());
diff --git 
a/solr/modules/cross-dc/src/test/org/apache/solr/crossdc/handler/MirroringConfigSetsHandlerTest.java
 
b/solr/modules/cross-dc/src/test/org/apache/solr/crossdc/handler/MirroringConfigSetsHandlerTest.java
index 275fa528385..b115f4cdcdb 100644
--- 
a/solr/modules/cross-dc/src/test/org/apache/solr/crossdc/handler/MirroringConfigSetsHandlerTest.java
+++ 
b/solr/modules/cross-dc/src/test/org/apache/solr/crossdc/handler/MirroringConfigSetsHandlerTest.java
@@ -20,7 +20,6 @@ import 
com.carrotsearch.randomizedtesting.annotations.ThreadLeakFilters;
 import com.carrotsearch.randomizedtesting.annotations.ThreadLeakLingering;
 import java.nio.charset.StandardCharsets;
 import java.nio.file.Path;
-import java.util.Arrays;
 import java.util.List;
 import org.apache.commons.io.IOUtils;
 import org.apache.lucene.tests.util.QuickPatchThreadsFilter;
@@ -52,7 +51,6 @@ import org.mockito.ArgumentCaptor;
 import org.mockito.Mockito;
 
 @ThreadLeakFilters(
-    defaultFilters = true,
     filters = {
       SolrIgnoredThreadsFilter.class,
       QuickPatchThreadsFilter.class,
@@ -149,7 +147,7 @@ public class MirroringConfigSetsHandlerTest extends 
SolrTestCaseJ4 {
       req.getParams()
           .forEach(
               entry -> {
-                assertEquals(entry.getValue(), 
mirroredParams.getParams(entry.getKey()));
+                assertArrayEquals(entry.getValue(), 
mirroredParams.getParams(entry.getKey()));
               });
       assertEquals("HTTP method", req.getHttpMethod(), 
solrRequest.getMethod().toString());
       if (expectStreams) {
@@ -167,7 +165,7 @@ public class MirroringConfigSetsHandlerTest extends 
SolrTestCaseJ4 {
               
MirroredSolrRequest.ExposedByteArrayContentStream.of(source).byteArray();
           byte[] mirroredContent =
               
MirroredSolrRequest.ExposedByteArrayContentStream.of(mirrored).byteArray();
-          assertTrue("different content", Arrays.equals(sourceContent, 
mirroredContent));
+          assertArrayEquals("different content", sourceContent, 
mirroredContent);
         }
       }
     } else {
diff --git 
a/solr/modules/cross-dc/src/test/org/apache/solr/crossdc/update/processor/MirroringUpdateProcessorTest.java
 
b/solr/modules/cross-dc/src/test/org/apache/solr/crossdc/update/processor/MirroringUpdateProcessorTest.java
index 77b3adf7a49..476b4529001 100644
--- 
a/solr/modules/cross-dc/src/test/org/apache/solr/crossdc/update/processor/MirroringUpdateProcessorTest.java
+++ 
b/solr/modules/cross-dc/src/test/org/apache/solr/crossdc/update/processor/MirroringUpdateProcessorTest.java
@@ -544,7 +544,7 @@ public class MirroringUpdateProcessorTest extends 
SolrTestCaseJ4 {
     processor.processDelete(deleteUpdateCommand);
     verify(requestMirroringHandler, times(1)).mirror(updateRequest);
     assertEquals("missing dbq", 1, updateRequest.getDeleteQuery().size());
-    assertEquals("dbq value", "id:test*", 
updateRequest.getDeleteQuery().get(0));
+    assertEquals("dbq value", "id:test*", 
updateRequest.getDeleteQuery().getFirst());
     // verify the metrics
     assertEquals(1, counters.get("local").get());
     assertEquals(1, counters.get("submittedDeleteByQuery").get());
@@ -575,13 +575,13 @@ public class MirroringUpdateProcessorTest extends 
SolrTestCaseJ4 {
 
   @Test
   public void testEstimateObjectSize() {
-    assertEquals(estimate(null), 0);
-    assertEquals(estimate("abc"), 6);
-    assertEquals(estimate("abcdefgh"), 16);
+    assertEquals(0, estimate(null));
+    assertEquals(6, estimate("abc"));
+    assertEquals(16, estimate("abcdefgh"));
     List<String> keys = List.of("int", "long", "double", "float", "str");
-    assertEquals(estimate(keys), 42);
+    assertEquals(42, estimate(keys));
     List<Object> values = List.of(12, 5L, 12.0, 5.0, "duck");
-    assertEquals(estimate(values), 8);
+    assertEquals(8, estimate(values));
 
     Map<String, Object> map = new HashMap<>();
     map.put("int", 12);
@@ -590,7 +590,7 @@ public class MirroringUpdateProcessorTest extends 
SolrTestCaseJ4 {
     map.put("float", 5.0f);
     map.put("str", "duck");
     map.put("short", null);
-    assertEquals(estimate(map), 60);
+    assertEquals(60, estimate(map));
 
     SolrInputDocument document = new SolrInputDocument();
     for (Map.Entry<String, Object> entry : map.entrySet()) {

Reply via email to