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]
