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

zqr10159 pushed a commit to branch 2.0.0
in repository https://gitbox.apache.org/repos/asf/hertzbeat.git

commit 7e67aca4966014e9ae09b564b5cdd7bdd9106e53
Author: Logic <[email protected]>
AuthorDate: Sun Oct 11 17:37:58 2026 +0800

    fix(manager): preserve JDBC reuse across bounded worker handoff
---
 .../workflow/TargetJdbcConnectionFactory.java      | 119 ++++++----
 .../workflow/TargetJdbcConnectionHandoffTest.java  | 239 +++++++++++++++++++++
 2 files changed, 321 insertions(+), 37 deletions(-)

diff --git 
a/hertzbeat-manager/src/main/java/org/apache/hertzbeat/manager/setup/workflow/TargetJdbcConnectionFactory.java
 
b/hertzbeat-manager/src/main/java/org/apache/hertzbeat/manager/setup/workflow/TargetJdbcConnectionFactory.java
index db8c15cea3..ecd2c7b6e1 100644
--- 
a/hertzbeat-manager/src/main/java/org/apache/hertzbeat/manager/setup/workflow/TargetJdbcConnectionFactory.java
+++ 
b/hertzbeat-manager/src/main/java/org/apache/hertzbeat/manager/setup/workflow/TargetJdbcConnectionFactory.java
@@ -25,21 +25,19 @@ import java.util.concurrent.RejectedExecutionException;
 import java.util.concurrent.SynchronousQueue;
 import java.util.concurrent.ThreadPoolExecutor;
 import java.util.concurrent.TimeUnit;
-import java.util.concurrent.locks.LockSupport;
 import org.apache.hertzbeat.manager.setup.config.MetadataDatabaseSettings;
 import org.apache.hertzbeat.manager.setup.config.SecretValue;
 
 /** Owns one bounded target JDBC acquisition worker and any failed cleanup 
handle. */
 final class TargetJdbcConnectionFactory implements AutoCloseable, 
TargetJdbcConnectionAttemptOwner {
 
-    private static final long AVAILABILITY_POLL_NANOS = 
TimeUnit.MILLISECONDS.toNanos(1);
-
     private final ThreadPoolExecutor worker;
     private final TargetJdbcCleanupLane cleanupLane;
     private final TargetJdbcConnector connector;
     private final TargetJdbcConnectionVerifier verifier;
     private final TargetJdbcResultWaiter resultWaiter;
     private TargetJdbcConnectionAttempt active;
+    private boolean previousAttemptFinished;
     private Error lateFailure;
     private boolean closed;
 
@@ -87,59 +85,43 @@ final class TargetJdbcConnectionFactory implements 
AutoCloseable, TargetJdbcConn
         TargetJdbcConnectionAttempt attempt = new TargetJdbcConnectionAttempt(
                 this, connector, verifier, target, settings.username(),
                 attemptPassword, deadline, resultWaiter);
+        boolean handoffExpected;
         try {
-            claim(attempt);
+            handoffExpected = claim(attempt);
         } catch (RuntimeException | Error claimFailure) {
             Arrays.fill(attemptPassword, '\0');
             throw claimFailure;
         }
+        boolean submitted;
         try {
-            submitWhenAvailable(attempt, deadline);
+            submitted = submit(attempt, deadline, handoffExpected);
+        } catch (InterruptedException interrupted) {
+            Arrays.fill(attemptPassword, '\0');
+            releaseUnsubmitted(attempt, handoffExpected);
+            Thread.currentThread().interrupt();
+            throw failure(TargetJdbcConnectionErrorCode.TIMEOUT);
         } catch (RejectedExecutionException rejected) {
             Arrays.fill(attemptPassword, '\0');
             poison(null);
-            finished(attempt);
+            releaseUnsubmitted(attempt, false);
             throw failure(TargetJdbcConnectionErrorCode.FACTORY_CLOSED);
         } catch (RuntimeException submitFailure) {
             Arrays.fill(attemptPassword, '\0');
             poison(null);
-            finished(attempt);
+            releaseUnsubmitted(attempt, false);
             throw failure(TargetJdbcConnectionErrorCode.FACTORY_CLOSED);
         } catch (Error fatalSubmission) {
             Arrays.fill(attemptPassword, '\0');
             lateFatal(fatalSubmission, null);
-            finished(attempt);
+            releaseUnsubmitted(attempt, false);
             throw fatalSubmission;
         }
-        return attempt.await();
-    }
-
-    private void submitWhenAvailable(
-            TargetJdbcConnectionAttempt attempt,
-            JdbcMetadataMigrationDeadline deadline) {
-        boolean interrupted = false;
-        try {
-            while (true) {
-                try {
-                    worker.execute(attempt);
-                    return;
-                } catch (RejectedExecutionException transientRejection) {
-                    // the zero-queue worker may still be unwinding the 
previous attempt
-                    if (worker.isShutdown() || deadline.remainingNanos() <= 0) 
{
-                        throw transientRejection;
-                    }
-                    LockSupport.parkNanos(AVAILABILITY_POLL_NANOS);
-                    if (Thread.interrupted()) {
-                        interrupted = true;
-                        throw transientRejection;
-                    }
-                }
-            }
-        } finally {
-            if (interrupted) {
-                Thread.currentThread().interrupt();
-            }
+        if (!submitted) {
+            Arrays.fill(attemptPassword, '\0');
+            releaseUnsubmitted(attempt, handoffExpected);
+            throw failure(TargetJdbcConnectionErrorCode.TIMEOUT);
         }
+        return attempt.await();
     }
 
     void retryCleanup(JdbcMetadataMigrationDeadline deadline) {
@@ -171,6 +153,7 @@ final class TargetJdbcConnectionFactory implements 
AutoCloseable, TargetJdbcConn
     @Override
     public synchronized void close() {
         closed = true;
+        notifyAll();
         if (active != null) {
             active.abandon(TargetJdbcConnectionErrorCode.FACTORY_CLOSED);
         } else {
@@ -200,6 +183,7 @@ final class TargetJdbcConnectionFactory implements 
AutoCloseable, TargetJdbcConn
     @Override
     public synchronized void finished(TargetJdbcConnectionAttempt attempt) {
         if (active == attempt) {
+            previousAttemptFinished = true;
             active = null;
             notifyAll();
             if (closed) {
@@ -230,7 +214,7 @@ final class TargetJdbcConnectionFactory implements 
AutoCloseable, TargetJdbcConn
         }
     }
 
-    private synchronized void claim(TargetJdbcConnectionAttempt attempt) {
+    private synchronized boolean claim(TargetJdbcConnectionAttempt attempt) {
         if (lateFailure != null) {
             throw lateFailure;
         }
@@ -245,6 +229,67 @@ final class TargetJdbcConnectionFactory implements 
AutoCloseable, TargetJdbcConn
             throw failure(TargetJdbcConnectionErrorCode.OPERATION_CONFLICT);
         }
         active = attempt;
+        boolean handoffExpected = previousAttemptFinished;
+        previousAttemptFinished = false;
+        return handoffExpected;
+    }
+
+    /**
+     * A finished attempt releases the logical slot before its executor worker 
becomes idle.
+     * The admitted successor owns that slot while submitting with its 
original acquisition budget.
+     * Always submit through execute: shutdownNow can drain even a 
SynchronousQueue timed offer.
+     */
+    private boolean submit(TargetJdbcConnectionAttempt attempt,
+                           JdbcMetadataMigrationDeadline deadline, boolean 
handoffExpected) throws InterruptedException {
+        try {
+            worker.execute(attempt);
+            return true;
+        } catch (RejectedExecutionException rejection) {
+            if (!handoffExpected || worker.isShutdown()
+                    || !(worker.getQueue() instanceof SynchronousQueue<?>)) {
+                throw rejection;
+            }
+        }
+        while (true) {
+            synchronized (this) {
+                if (closed || worker.isShutdown()) {
+                    throw new RejectedExecutionException();
+                }
+                long remaining = deadline.remainingNanos();
+                if (remaining <= 0) {
+                    return false;
+                }
+                // close notifies this condition; executor readiness is 
retried within the same budget.
+                TimeUnit.NANOSECONDS.timedWait(this,
+                        Math.min(remaining, 
TimeUnit.MILLISECONDS.toNanos(10)));
+                if (closed || worker.isShutdown()) {
+                    throw new RejectedExecutionException();
+                }
+                if (deadline.remainingNanos() <= 0) {
+                    return false;
+                }
+            }
+            try {
+                worker.execute(attempt);
+                return true;
+            } catch (RejectedExecutionException rejection) {
+                if (worker.isShutdown()) {
+                    throw rejection;
+                }
+            }
+        }
+    }
+
+    private synchronized void releaseUnsubmitted(
+            TargetJdbcConnectionAttempt attempt, boolean handoffExpected) {
+        if (active == attempt) {
+            previousAttemptFinished = handoffExpected;
+            active = null;
+            notifyAll();
+            if (closed) {
+                cleanupLane.close();
+            }
+        }
     }
 
     private static TargetJdbcConnectionException 
failure(TargetJdbcConnectionErrorCode code) {
diff --git 
a/hertzbeat-manager/src/test/java/org/apache/hertzbeat/manager/setup/workflow/TargetJdbcConnectionHandoffTest.java
 
b/hertzbeat-manager/src/test/java/org/apache/hertzbeat/manager/setup/workflow/TargetJdbcConnectionHandoffTest.java
new file mode 100644
index 0000000000..99e5cdf55c
--- /dev/null
+++ 
b/hertzbeat-manager/src/test/java/org/apache/hertzbeat/manager/setup/workflow/TargetJdbcConnectionHandoffTest.java
@@ -0,0 +1,239 @@
+/*
+ * 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.hertzbeat.manager.setup.workflow;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+import java.sql.SQLException;
+import java.time.Duration;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.Future;
+import java.util.concurrent.SynchronousQueue;
+import java.util.concurrent.ThreadPoolExecutor;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicInteger;
+import java.util.concurrent.atomic.AtomicLong;
+import java.util.concurrent.atomic.AtomicReference;
+import 
org.apache.hertzbeat.manager.setup.api.SetupApiContract.MetadataDatabaseKind;
+import org.apache.hertzbeat.manager.setup.config.MetadataDatabaseSettings;
+import org.apache.hertzbeat.manager.setup.config.SecretValue;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.Timeout;
+
+@Timeout(15)
+class TargetJdbcConnectionHandoffTest {
+
+    @Test
+    void reusableSuccessorWaitsForPhysicalHandoffAndRejectsConcurrentAcquire() 
throws Exception {
+        try (Fixture fixture = new Fixture(); ExecutorService caller = 
Executors.newSingleThreadExecutor()) {
+            fixture.firstFailure();
+            Future<TargetJdbcConnectionErrorCode> next = caller.submit(() -> 
acquire(fixture.factory, deadline()));
+            fixture.awaitHandoff();
+            assertThat(next.isDone()).isFalse();
+            assertThat(acquire(fixture.factory, 
deadline())).isEqualTo(TargetJdbcConnectionErrorCode.OPERATION_CONFLICT);
+            fixture.release.countDown();
+            assertThat(next.get(5, 
TimeUnit.SECONDS)).isEqualTo(TargetJdbcConnectionErrorCode.UNAVAILABLE);
+            assertThat(fixture.calls.get()).isEqualTo(2);
+            
assertThat(fixture.factory.settleFailedAcquire(deadline())).isEqualTo(TargetJdbcFailedAcquireSettlement.REUSABLE);
+            assertThat(acquire(fixture.factory, 
deadline())).isEqualTo(TargetJdbcConnectionErrorCode.UNAVAILABLE);
+            assertThat(fixture.calls.get()).isEqualTo(3);
+            fixture.assertWiped();
+        }
+    }
+
+    @Test
+    void handoffExpiryUsesOriginalDeadlineAndLeavesFactoryReusable() throws 
Exception {
+        try (Fixture fixture = new Fixture(); ExecutorService caller = 
Executors.newSingleThreadExecutor()) {
+            fixture.firstFailure();
+            AtomicLong clock = new AtomicLong();
+            JdbcMetadataMigrationDeadline budget = 
JdbcMetadataMigrationDeadline.start(Duration.ofSeconds(5), clock::get);
+            Future<TargetJdbcConnectionErrorCode> next = caller.submit(() -> 
acquire(fixture.factory, budget));
+            fixture.awaitHandoff();
+            clock.set(TimeUnit.SECONDS.toNanos(5));
+            assertThat(next.get(5, 
TimeUnit.SECONDS)).isEqualTo(TargetJdbcConnectionErrorCode.TIMEOUT);
+            assertThat(fixture.calls.get()).isEqualTo(1);
+            
assertThat(fixture.factory.settleFailedAcquire(deadline())).isEqualTo(TargetJdbcFailedAcquireSettlement.REUSABLE);
+            fixture.release.countDown();
+            assertThat(acquire(fixture.factory, 
deadline())).isEqualTo(TargetJdbcConnectionErrorCode.UNAVAILABLE);
+            fixture.assertWiped();
+        }
+    }
+
+    @Test
+    void interruptDuringHandoffRestoresInterruptWithoutPoisoningFactory() 
throws Exception {
+        try (Fixture fixture = new Fixture(); ExecutorService caller = 
Executors.newSingleThreadExecutor()) {
+            fixture.firstFailure();
+            AtomicReference<Thread> thread = new AtomicReference<>();
+            Future<Boolean> next = caller.submit(() -> {
+                thread.set(Thread.currentThread());
+                return acquire(fixture.factory, deadline()) == 
TargetJdbcConnectionErrorCode.TIMEOUT
+                        && Thread.currentThread().isInterrupted();
+            });
+            fixture.awaitHandoff();
+            thread.get().interrupt();
+            assertThat(next.get(5, TimeUnit.SECONDS)).isTrue();
+            
assertThat(fixture.factory.settleFailedAcquire(deadline())).isEqualTo(TargetJdbcFailedAcquireSettlement.REUSABLE);
+            fixture.release.countDown();
+            assertThat(acquire(fixture.factory, 
deadline())).isEqualTo(TargetJdbcConnectionErrorCode.UNAVAILABLE);
+        }
+    }
+
+    @Test
+    void closeDuringHandoffClearsSlotWithoutExecutingTheSuccessor() throws 
Exception {
+        try (Fixture fixture = new Fixture(); ExecutorService caller = 
Executors.newSingleThreadExecutor()) {
+            fixture.firstFailure();
+            Future<TargetJdbcConnectionErrorCode> next = caller.submit(() -> 
acquire(fixture.factory, deadline()));
+            fixture.awaitHandoff();
+            fixture.factory.close();
+            assertThat(next.get(5, 
TimeUnit.SECONDS)).isEqualTo(TargetJdbcConnectionErrorCode.FACTORY_CLOSED);
+            assertThat(fixture.calls.get()).isEqualTo(1);
+            
assertThat(fixture.factory.settleFailedAcquire(deadline())).isEqualTo(TargetJdbcFailedAcquireSettlement.TERMINAL_CLOSED);
+            assertThat(acquire(fixture.factory, 
deadline())).isEqualTo(TargetJdbcConnectionErrorCode.FACTORY_CLOSED);
+        }
+    }
+
+    @Test
+    void executorShutdownDuringHandoffRemainsTerminal() throws Exception {
+        try (Fixture fixture = new Fixture(); ExecutorService caller = 
Executors.newSingleThreadExecutor()) {
+            fixture.firstFailure();
+            Future<TargetJdbcConnectionErrorCode> next = caller.submit(() -> 
acquire(fixture.factory, deadline()));
+            fixture.awaitHandoff();
+            fixture.worker.shutdown();
+            assertThat(next.get(5, 
TimeUnit.SECONDS)).isEqualTo(TargetJdbcConnectionErrorCode.FACTORY_CLOSED);
+            
assertThat(fixture.factory.settleFailedAcquire(deadline())).isEqualTo(TargetJdbcFailedAcquireSettlement.TERMINAL_CLOSED);
+            assertThat(fixture.calls.get()).isEqualTo(1);
+        }
+    }
+
+    @Test
+    void rejectionWithoutCompletedPredecessorRemainsTerminal() {
+        ThreadPoolExecutor rejecting = new ThreadPoolExecutor(0, 1, 30, 
TimeUnit.SECONDS, new SynchronousQueue<>()) {
+            @Override
+            public void execute(Runnable task) {
+                throw new 
java.util.concurrent.RejectedExecutionException("private failure");
+            }
+        };
+        try (TargetJdbcConnectionFactory factory = new 
TargetJdbcConnectionFactory(rejecting, worker(), Runnable::run,
+                (target, username, password, budget) -> { throw new 
SQLException(); },
+                new TargetJdbcConnectionVerifier(Runnable::run))) {
+            assertThat(acquire(factory, 
deadline())).isEqualTo(TargetJdbcConnectionErrorCode.FACTORY_CLOSED);
+            assertThat(acquire(factory, 
deadline())).isEqualTo(TargetJdbcConnectionErrorCode.FACTORY_CLOSED);
+        }
+    }
+
+    private static TargetJdbcConnectionErrorCode acquire(
+            TargetJdbcConnectionFactory factory, JdbcMetadataMigrationDeadline 
budget) {
+        try (SecretValue password = SecretValue.of("unused-fixture")) {
+            factory.acquire(new 
MetadataDatabaseSettings(MetadataDatabaseKind.MYSQL,
+                    "jdbc:mysql://example.invalid/hertzbeat?sslMode=REQUIRED", 
"unused-fixture"), password, budget);
+            throw new AssertionError("Expected isolated connector failure");
+        } catch (TargetJdbcConnectionException failure) {
+            assertThat(failure).hasNoCause();
+            return failure.code();
+        }
+    }
+
+    private static JdbcMetadataMigrationDeadline deadline() {
+        return JdbcMetadataMigrationDeadline.start(Duration.ofSeconds(5), 
System::nanoTime);
+    }
+
+    private static ThreadPoolExecutor worker() {
+        return new ThreadPoolExecutor(0, 1, 30, TimeUnit.SECONDS, new 
SynchronousQueue<>());
+    }
+
+    private static final class Fixture implements AutoCloseable {
+        private final CountDownLatch unwinding = new CountDownLatch(1);
+        private final CountDownLatch release = new CountDownLatch(1);
+        private final CountDownLatch offering = new CountDownLatch(1);
+        private final AtomicReference<Throwable> hookFailure = new 
AtomicReference<>();
+        private final AtomicReference<char[]> seenPassword = new 
AtomicReference<>();
+        private final AtomicInteger calls = new AtomicInteger();
+        private final ThreadPoolExecutor worker;
+        private final TargetJdbcConnectionFactory factory;
+
+        private Fixture() {
+            worker = new ThreadPoolExecutor(0, 1, 30, TimeUnit.SECONDS, new 
SynchronousQueue<>()) {
+                @Override
+                public void execute(Runnable task) {
+                    try {
+                        super.execute(task);
+                    } catch (java.util.concurrent.RejectedExecutionException 
rejection) {
+                        offering.countDown();
+                        throw rejection;
+                    }
+                }
+
+                @Override
+                protected void afterExecute(Runnable task, Throwable failure) {
+                    if (failure != null) {
+                        hookFailure.compareAndSet(null, failure);
+                    }
+                    if (unwinding.getCount() != 0) {
+                        unwinding.countDown();
+                        boolean interrupted = false;
+                        while (release.getCount() != 0) {
+                            try {
+                                if (!release.await(5, TimeUnit.SECONDS)) {
+                                    hookFailure.compareAndSet(null, new 
AssertionError("release timed out"));
+                                    break;
+                                }
+                            } catch (InterruptedException ignored) {
+                                interrupted = true;
+                            }
+                        }
+                        if (interrupted) {
+                            Thread.currentThread().interrupt();
+                        }
+                    }
+                }
+            };
+            factory = new TargetJdbcConnectionFactory(worker, worker(), 
Runnable::run,
+                    (target, username, password, budget) -> {
+                        seenPassword.set(password);
+                        calls.incrementAndGet();
+                        throw new SQLException("isolated unavailable 
connector");
+                    }, new TargetJdbcConnectionVerifier(Runnable::run));
+        }
+
+        private void firstFailure() throws Exception {
+            assertThat(acquire(factory, 
deadline())).isEqualTo(TargetJdbcConnectionErrorCode.UNAVAILABLE);
+            assertThat(unwinding.await(5, TimeUnit.SECONDS)).isTrue();
+            
assertThat(factory.settleFailedAcquire(deadline())).isEqualTo(TargetJdbcFailedAcquireSettlement.REUSABLE);
+            assertWiped();
+        }
+
+        private void awaitHandoff() throws Exception {
+            assertThat(offering.await(5, TimeUnit.SECONDS)).isTrue();
+            assertThat(worker.getQueue()).isEmpty();
+        }
+
+        private void assertWiped() {
+            assertThat(seenPassword.get()).containsOnly('\0');
+        }
+
+        @Override
+        public void close() throws Exception {
+            release.countDown();
+            factory.close();
+            assertThat(worker.awaitTermination(5, TimeUnit.SECONDS)).isTrue();
+            assertThat(hookFailure.get()).isNull();
+        }
+    }
+}


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to