kotman12 commented on code in PR #4974: URL: https://github.com/apache/solr/pull/4974#discussion_r4160146726
########## solr/core/src/test/org/apache/solr/cloud/TlogLeaderElectionFrozenLeaderTest.java: ########## @@ -0,0 +1,343 @@ +/* + * 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.cloud; + +import com.carrotsearch.randomizedtesting.annotations.Name; +import com.carrotsearch.randomizedtesting.annotations.ParametersFactory; +import jakarta.servlet.Filter; +import jakarta.servlet.FilterChain; +import jakarta.servlet.ServletException; +import jakarta.servlet.ServletOutputStream; +import jakarta.servlet.ServletRequest; +import jakarta.servlet.ServletResponse; +import jakarta.servlet.http.HttpServletRequest; +import jakarta.servlet.http.HttpServletResponse; +import jakarta.servlet.http.HttpServletResponseWrapper; +import java.io.IOException; +import java.io.InterruptedIOException; +import java.lang.invoke.MethodHandles; +import java.util.List; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.ExecutionException; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; +import java.util.concurrent.atomic.AtomicReference; +import org.apache.solr.client.solrj.request.CollectionAdminRequest; +import org.apache.solr.common.SolrInputDocument; +import org.apache.solr.common.cloud.Replica; +import org.apache.solr.common.cloud.Slice; +import org.apache.solr.common.util.SuppressForbidden; +import org.apache.solr.embedded.JettySolrRunner; +import org.apache.solr.handler.ReplicationHandler; +import org.apache.solr.servlet.ServletOutputStreamWrapper; +import org.junit.After; +import org.junit.Before; +import org.junit.Test; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * Reproduces a TLOG leader election that stalls because the outgoing leader froze rather than died. + * + * <p>{@link ShardLeaderElectionContext#runLeaderProcess} calls {@link + * ZkController#stopReplicationFromLeader} inline on the election thread. That tears down the + * follower's replication process, which blocks in {@code ExecutorUtil.shutdownAndAwaitTermination} + * until the in-flight index fetch finishes. {@code IndexFetcher.abortFetch} only sets a flag that + * is polled while streaming file packets, so a fetch parked in the network phase is not cut short: + * the election parks for the 60s executor wait before {@code shutdownNow()} finally interrupts the + * poll thread. + */ +public class TlogLeaderElectionFrozenLeaderTest extends SolrCloudTestCase { + + private static final Logger log = LoggerFactory.getLogger(MethodHandles.lookup().lookupClass()); + + private static final String COLLECTION = "tlog_frozen_leader"; + private static final String SHARD = "shard1"; + + /** + * A replacement election sleeps a fixed 2.5s for in-flight updates to settle + * (ShardLeaderElectionContext), then syncs, replays its tlog and publishes in well under a + * second. The bug parks it for the 60s {@code ExecutorUtil.awaitTermination} wait instead. + */ + private static final long MAX_ACCEPTABLE_ELECTION_MS = 10_000; + + /** How long to wait for the follower's next poll to reach the stalled leader. */ + private static final long STALL_ARRIVAL_TIMEOUT_MS = 30_000; + + /** Where a follower's fetch can park against a frozen leader. */ + public enum TestStallPoint { Review Comment: There are several possible I/O stall points between new and former leader. To be thorough I test all in the nightly build. ########## solr/core/src/java/org/apache/solr/handler/IndexFetcher.java: ########## @@ -1666,11 +1710,6 @@ private int fetchPackets(FastInputStream fis) throws Exception { } return 0; } - if (stop) { - stop = false; Review Comment: This is a vestige from when individual fetch abortion was tracked at the `IndexFetcher` level (kind of strange). Now aborting a fetch only aborts the fetch-scoped client and that gets discarded in the `finally` block of `fetchLatestIndex` so the scope is more cleanly enforced. ########## solr/solrj-jetty/src/java/org/apache/solr/client/solrj/jetty/HttpJettySolrClient.java: ########## @@ -931,6 +961,68 @@ public void waitForComplete() { } } + /** + * Tracks a client's requests, synchronous and asynchronous, from queued to complete, so that + * {@link #abort(Throwable)} can reach them. + */ + private static class AbortTracker { + // requests queued and not yet complete + private final Set<Request> outstanding = ConcurrentHashMap.newKeySet(); + // non-null once abort() has been called; fails any request queued afterwards + private volatile Throwable abortCause; + + void track(Request request) { + // Pairs with abort(): add-then-check here and set-then-scan there, so a request queued + // concurrently with abort() is aborted by at least one of the two. + outstanding.add(request); + Throwable cause = abortCause; + if (cause != null) { + request.abort(cause); + } + } + + void untrack(Request request) { + outstanding.remove(request); + } + + void abort(Throwable cause) { + abortCause = Objects.requireNonNull(cause); + outstanding.forEach(request -> request.abort(cause)); + } + + boolean isAborted() { + return abortCause != null; + } + } + + /** + * Aborts every request in flight on this client, and fails any request made afterwards. Unlike + * {@link #close()}, this does not wait: each aborted request completes exceptionally straight + * away, so a subsequent {@link #close()} returns promptly. + * + * <p>Scoped to this client instance and to the clients derived from it by {@link + * #requestWithBaseUrl}, which share its tracking. Other clients wrapping the same Jetty {@link + * HttpClient} are unaffected, so a lightweight copy made with {@link + * Builder#withHttpClient(HttpJettySolrClient)} can be aborted without disturbing anyone else. + * + * @param cause reported as the failure of each aborted request, and rethrown as-is to anyone + * reading a streamed response body, so pass an {@link IOException} to keep {@link + * InputStream} callers on their usual error path; must not be null + * @throws IllegalStateException if the client was not built with {@link + * Builder#withAbortableRequests(boolean)} + */ + public void abort(Throwable cause) { + if (abortTracker == null) { + throw new IllegalStateException("Client was not built with withAbortableRequests(true)"); Review Comment: Throw instead of no-opping to dispel any illusions quickly. Since this is new functionality, it should be ok. ########## solr/core/src/test/org/apache/solr/cloud/TlogLeaderElectionFrozenLeaderTest.java: ########## @@ -0,0 +1,343 @@ +/* + * 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.cloud; + +import com.carrotsearch.randomizedtesting.annotations.Name; +import com.carrotsearch.randomizedtesting.annotations.ParametersFactory; +import jakarta.servlet.Filter; +import jakarta.servlet.FilterChain; +import jakarta.servlet.ServletException; +import jakarta.servlet.ServletOutputStream; +import jakarta.servlet.ServletRequest; +import jakarta.servlet.ServletResponse; +import jakarta.servlet.http.HttpServletRequest; +import jakarta.servlet.http.HttpServletResponse; +import jakarta.servlet.http.HttpServletResponseWrapper; +import java.io.IOException; +import java.io.InterruptedIOException; +import java.lang.invoke.MethodHandles; +import java.util.List; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.ExecutionException; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; +import java.util.concurrent.atomic.AtomicReference; +import org.apache.solr.client.solrj.request.CollectionAdminRequest; +import org.apache.solr.common.SolrInputDocument; +import org.apache.solr.common.cloud.Replica; +import org.apache.solr.common.cloud.Slice; +import org.apache.solr.common.util.SuppressForbidden; +import org.apache.solr.embedded.JettySolrRunner; +import org.apache.solr.handler.ReplicationHandler; +import org.apache.solr.servlet.ServletOutputStreamWrapper; +import org.junit.After; +import org.junit.Before; +import org.junit.Test; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * Reproduces a TLOG leader election that stalls because the outgoing leader froze rather than died. + * + * <p>{@link ShardLeaderElectionContext#runLeaderProcess} calls {@link + * ZkController#stopReplicationFromLeader} inline on the election thread. That tears down the + * follower's replication process, which blocks in {@code ExecutorUtil.shutdownAndAwaitTermination} + * until the in-flight index fetch finishes. {@code IndexFetcher.abortFetch} only sets a flag that + * is polled while streaming file packets, so a fetch parked in the network phase is not cut short: + * the election parks for the 60s executor wait before {@code shutdownNow()} finally interrupts the + * poll thread. + */ +public class TlogLeaderElectionFrozenLeaderTest extends SolrCloudTestCase { + + private static final Logger log = LoggerFactory.getLogger(MethodHandles.lookup().lookupClass()); + + private static final String COLLECTION = "tlog_frozen_leader"; + private static final String SHARD = "shard1"; + + /** + * A replacement election sleeps a fixed 2.5s for in-flight updates to settle + * (ShardLeaderElectionContext), then syncs, replays its tlog and publishes in well under a + * second. The bug parks it for the 60s {@code ExecutorUtil.awaitTermination} wait instead. + */ + private static final long MAX_ACCEPTABLE_ELECTION_MS = 10_000; + + /** How long to wait for the follower's next poll to reach the stalled leader. */ + private static final long STALL_ARRIVAL_TIMEOUT_MS = 30_000; + + /** Where a follower's fetch can park against a frozen leader. */ + public enum TestStallPoint { + INDEX_VERSION(ReplicationHandler.CMD_INDEX_VERSION, false), + FILE_LIST(ReplicationHandler.CMD_GET_FILE_LIST, false), + FILE_CONTENT(ReplicationHandler.CMD_GET_FILE, false), + /** + * Headers and part of the first packet arrive, then the body stops: the common production case. + */ + FILE_CONTENT_BODY(ReplicationHandler.CMD_GET_FILE, true); + + final String command; + final boolean midBody; + + TestStallPoint(String command, boolean midBody) { + this.command = command; + this.midBody = midBody; + } + } + + /** + * Nightly runs every stall point. Otherwise only the mid-body file download, the most likely + * failure in production, so every CI run covers it. + */ + @ParametersFactory + public static Iterable<Object[]> parameters() { + List<TestStallPoint> stallPoints = + TEST_NIGHTLY ? List.of(TestStallPoint.values()) : List.of(TestStallPoint.FILE_CONTENT_BODY); + return stallPoints.stream().map(stallPoint -> new Object[] {stallPoint}).toList(); + } + + private final TestStallPoint stallPoint; + + public TlogLeaderElectionFrozenLeaderTest(@Name("stallPoint") TestStallPoint stallPoint) { + this.stallPoint = stallPoint; + } + + @Before + public void setupCluster() throws Exception { + System.setProperty("solr.directoryFactory", "solr.StandardDirectoryFactory"); + + // extraFilters are installed ahead of SolrServlet and its filters (JettySolrRunner:336-338), so + // ours sees the request first. It is installed on every node but only acts on the armed core. + configureCluster(2) + .withJettyConfig(b -> b.withFilter(TestStallReplicationFilter.class, "/*")) + .addConfig("conf", configset("cloud-minimal")) + .configure(); + } + + @After + public void tearDownCluster() throws Exception { + TestStallChannel stallChannel = TestStallReplicationFilter.STALL_CHANNEL.getAndSet(null); + if (stallChannel != null) { + stallChannel.released().complete(null); + } + shutdownCluster(); + System.clearProperty("solr.directoryFactory"); + } + + @Test + public void testElectionIsNotBlockedByFrozenOldLeader() throws Exception { + CollectionAdminRequest.createCollection(COLLECTION, "conf", 1, 0, 2, 0) + .process(cluster.getSolrClient()); + cluster.waitForActiveCollection(COLLECTION, 1, 2); + + // Index something so the follower has a real index to poll against. + indexAndCommit(0, 10); + + Slice shard = getCollectionState(COLLECTION).getSlice(SHARD); + Replica oldLeader = shard.getLeader(); + Replica follower = + shard.getReplicas().stream() + .filter(r -> !r.getName().equals(oldLeader.getName())) + .findFirst() + .orElseThrow(); + JettySolrRunner leaderJetty = cluster.getReplicaJetty(oldLeader); + JettySolrRunner followerJetty = cluster.getReplicaJetty(follower); + + // Stall one of the leader's replication commands. The follower polls every second under + // jetty.testMode, so its next request for that command lands in the filter and never returns. + if (log.isInfoEnabled()) { + log.info("Stalling {} responses from leader core {}", stallPoint, oldLeader.getCoreName()); + } + TestStallChannel channel = new TestStallChannel(oldLeader.getCoreName(), stallPoint); + assertTrue( Review Comment: I don't think it is possible for two tests to use this at the same time but have a defensive check for if this ever changes. One could consider disambiguating by core name or something but figured it was not worth the extra code. ########## solr/core/src/java/org/apache/solr/handler/IndexFetcher.java: ########## @@ -1866,6 +1918,13 @@ private IOException closeStreamAndBuildIOE( } } + /** + * A fetch of one leader generation: the fetch's client plus the generation being fetched. A retry + * (when the leader discards the generation, or the full-copy retry) reads a fresh generation and + * starts a new {@code Fetch} with the same client. + */ + private record Fetch(HttpJettySolrClient client, long generation) {} Review Comment: This is sugar to prevent signature bloat. I guess it's pretty arbitrary though both of these data points are neatly scoped to a single `Fetch` and so it felt natural enough. ########## solr/core/src/java/org/apache/solr/handler/ReplicationHandler.java: ########## @@ -504,7 +504,7 @@ public IndexFetchResult doFetch(SolrParams solrParams, boolean forceReplication) if (!indexFetchLock.tryLock()) return IndexFetchResult.LOCK_OBTAIN_FAILED; if (core.getCoreContainer().isShutDown()) { log.warn("I was asked to replicate but CoreContainer is shutting down"); - return IndexFetchResult.CONTAINER_IS_SHUTTING_DOWN; + return IndexFetchResult.REPLICATION_SHUTTING_DOWN; Review Comment: more "technically correct" description ########## solr/solrj-jetty/src/java/org/apache/solr/client/solrj/jetty/HttpJettySolrClient.java: ########## @@ -119,6 +122,8 @@ public class HttpJettySolrClient extends HttpSolrClient { private final List<HttpListenerFactory> listenerFactory; protected AsyncTracker asyncTracker = new AsyncTracker(); + // null unless built withAbortableRequests(true) + protected AbortTracker abortTracker; Review Comment: `protected` for nested subclass assignment even though class is `private`. The field itself could be made `private` if we assign it via `super.abortTracker=...` in `NoCloseHttpJettySolrClient` though that seems equally strange. You can also pass it via the builder but then you need an `AbortTracker` reference in the builder. So all the options kind of suck. ########## AGENTS.md: ########## @@ -42,6 +42,10 @@ While README.md and CONTRIBUTING.md are mainly written for humans, this file is - For BATS shell integration tests in `solr/packaging/test/`: - Always use `run <command>` followed by `assert_output --partial "..."` or `refute_output --partial "..."` instead of capturing output into local variables and using `[[ ]]` comparisons - Avoid patterns like `local var=$(cmd | grep ...); [[ "$var" == *"..."* ]]` — use `run cmd` + `assert_output`/`refute_output` instead +- Avoid `Thread.sleep`; fixed sleeps waste build time on every run and can flake under load. Review Comment: I find from my personal experience writing _and_ reviewing that LLM-based agents _constantly_ have trouble with these concepts. ########## solr/solrj-jetty/src/test/org/apache/solr/client/solrj/jetty/HttpJettySolrClientTest.java: ########## @@ -698,6 +708,154 @@ public void testBuilder() { } } + /** + * How soon an aborted request must fail: well before the slow servlet would have answered on its + * own, so passing proves the abort, not the servlet, ended it. + */ + private static final long ABORT_DEADLINE_MS = SlowServlet.DELAY_MS / 2; + + @Test + public void testAbortFailsAsyncRequestInFlight() throws Exception { + try (var client = + new HttpJettySolrClient.Builder(solrTestRule.getBaseUrl() + SLOW_SERVLET_PATH) + .withAbortableRequests(true) + .build()) { + CompletableFuture<Void> arrived = new CompletableFuture<>(); + CompletableFuture<NamedList<Object>> future = client.requestAsync(stallRequest(arrived)); + arrived.get(SlowServlet.DELAY_MS, TimeUnit.MILLISECONDS); + client.abort(new IOException("test abort")); + expectThrows( + ExecutionException.class, () -> future.get(ABORT_DEADLINE_MS, TimeUnit.MILLISECONDS)); + assertAllRequestsComplete(client); + } + } + + @Test + public void testAbortFailsResponseMidBody() throws Exception { + try (var client = + new HttpJettySolrClient.Builder(solrTestRule.getBaseUrl() + SLOW_SERVLET_PATH) + .withAbortableRequests(true) + .build()) { + QueryRequest req = new QueryRequest(SolrParams.of(SlowServlet.FIRST_BYTE_PARAM, "true")); + req.setResponseParser(new InputStreamResponseParser(FILE_STREAM)); + NamedList<Object> response = client.requestAsync(req).get(10, TimeUnit.SECONDS); + try (InputStream is = (InputStream) response.get("stream")) { + assertEquals('0', is.read()); + // The servlet now sends nothing more; the abort must wake the read rather than wait it out. + client.abort(new IOException("test abort")); + long start = System.nanoTime(); + IOException e = expectThrows(IOException.class, is::read); + long elapsedMs = TimeUnit.NANOSECONDS.toMillis(System.nanoTime() - start); + assertEquals("test abort", e.getMessage()); + assertTrue("read took " + elapsedMs + "ms after abort", elapsedMs < ABORT_DEADLINE_MS); + } + assertAllRequestsComplete(client); + } + } + + @Test + public void testAbortFailsLaterAsyncRequests() throws Exception { + try (var client = + new HttpJettySolrClient.Builder(solrTestRule.getBaseUrl() + SLOW_STREAM_SERVLET_PATH) + .withAbortableRequests(true) + .build()) { + client.abort(new IOException("test abort")); + assertTrue(client.isAborted()); + QueryRequest req = new QueryRequest(SolrParams.of("count", "1")); + req.setResponseParser(new InputStreamResponseParser(FILE_STREAM)); + CompletableFuture<NamedList<Object>> future = client.requestAsync(req); + expectThrows( + ExecutionException.class, () -> future.get(ABORT_DEADLINE_MS, TimeUnit.MILLISECONDS)); + assertAllRequestsComplete(client); + } + } + + @Test + public void testAbortIsScopedToOneClient() throws Exception { + String url = solrTestRule.getBaseUrl() + SLOW_STREAM_SERVLET_PATH; + try (var base = new HttpJettySolrClient.Builder(url).build(); + var copy = + new HttpJettySolrClient.Builder(url) + .withHttpClient(base) + .withAbortableRequests(true) + .build()) { + copy.abort(new IOException("test abort")); + assertTrue(copy.isAborted()); + assertFalse(base.isAborted()); + // The copy shares base's Jetty client, but not its tracker, so base still works. + QueryRequest req = new QueryRequest(SolrParams.of("count", "1")); + req.setResponseParser(new InputStreamResponseParser(FILE_STREAM)); + NamedList<Object> response = base.requestAsync(req).get(10, TimeUnit.SECONDS); + try (InputStream is = (InputStream) response.get("stream")) { + assertEquals("0", new String(is.readAllBytes(), StandardCharsets.UTF_8)); + } + } + } + + @Test + public void testAbortFailsSyncRequestInFlight() throws Exception { + ExecutorService executor = + ExecutorUtil.newMDCAwareSingleThreadExecutor(new SolrNamedThreadFactory("syncRequest")); + try (var client = + new HttpJettySolrClient.Builder(solrTestRule.getBaseUrl() + SLOW_SERVLET_PATH) + .withAbortableRequests(true) + .build()) { + CompletableFuture<Void> arrived = new CompletableFuture<>(); + QueryRequest req = stallRequest(arrived); + Future<NamedList<Object>> call = executor.submit(() -> client.request(req)); + arrived.get(SlowServlet.DELAY_MS, TimeUnit.MILLISECONDS); + client.abort(new IOException("test abort")); + ExecutionException e = + expectThrows( + ExecutionException.class, () -> call.get(ABORT_DEADLINE_MS, TimeUnit.MILLISECONDS)); + assertTrue(e.getCause().toString(), e.getCause() instanceof SolrServerException); + } finally { + ExecutorUtil.shutdownAndAwaitTermination(executor); + } + } + + @Test + public void testAbortFailsLaterSyncRequests() { + try (var client = + new HttpJettySolrClient.Builder(solrTestRule.getBaseUrl() + SLOW_STREAM_SERVLET_PATH) + .withAbortableRequests(true) + .build()) { + client.abort(new IOException("test abort")); + QueryRequest req = new QueryRequest(SolrParams.of("count", "1")); + req.setResponseParser(new InputStreamResponseParser(FILE_STREAM)); + expectThrows(SolrServerException.class, () -> client.request(req)); + } + } + + @Test + public void testAbortRequiresOptIn() { + try (var client = new HttpJettySolrClient.Builder(solrTestRule.getBaseUrl()).build()) { + expectThrows(IllegalStateException.class, () -> client.abort(new IOException("test abort"))); Review Comment: throwing if we don't opt-in as previously explained. ########## AGENTS.md: ########## @@ -42,6 +42,10 @@ While README.md and CONTRIBUTING.md are mainly written for humans, this file is - For BATS shell integration tests in `solr/packaging/test/`: - Always use `run <command>` followed by `assert_output --partial "..."` or `refute_output --partial "..."` instead of capturing output into local variables and using `[[ ]]` comparisons - Avoid patterns like `local var=$(cmd | grep ...); [[ "$var" == *"..."* ]]` — use `run cmd` + `assert_output`/`refute_output` instead +- Avoid `Thread.sleep`; fixed sleeps waste build time on every run and can flake under load. Review Comment: maybe worth adding a bit about nightly vs regular builds so the price of tests is somewhat considered (actually how useful this will be is hard to say). -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
