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

github-merge-queue[bot] pushed a commit to branch 
gh-readonly-queue/dev/pr-12133-c57dfaa8d08dba8a95567fdb77fbfb75c1621293
in repository https://gitbox.apache.org/repos/asf/seatunnel.git

commit c20fefe3bf05dc636bc70488a11a7ab439a1d2e5
Author: Daniel <[email protected]>
AuthorDate: Sun Sep 6 03:42:30 2026 +0000

    [Fix][Zeta] Add E2E coverage for telemetry startup deadlock regression 
(apache/seatunnel#10840) (#12133)
    
    Co-authored-by: Claude Sonnet 5 <[email protected]>
---
 .../e2e/telemetry/TelemetryStartupDeadlockIT.java  | 473 +++++++++++++++++++++
 1 file changed, 473 insertions(+)

diff --git 
a/seatunnel-e2e/seatunnel-engine-e2e/connector-seatunnel-e2e-base/src/test/java/org/apache/seatunnel/engine/e2e/telemetry/TelemetryStartupDeadlockIT.java
 
b/seatunnel-e2e/seatunnel-engine-e2e/connector-seatunnel-e2e-base/src/test/java/org/apache/seatunnel/engine/e2e/telemetry/TelemetryStartupDeadlockIT.java
new file mode 100644
index 0000000000..573c27dd7a
--- /dev/null
+++ 
b/seatunnel-e2e/seatunnel-engine-e2e/connector-seatunnel-e2e-base/src/test/java/org/apache/seatunnel/engine/e2e/telemetry/TelemetryStartupDeadlockIT.java
@@ -0,0 +1,473 @@
+/*
+ * 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.seatunnel.engine.e2e.telemetry;
+
+import org.apache.seatunnel.engine.client.SeaTunnelClient;
+import org.apache.seatunnel.engine.client.job.ClientJobExecutionEnvironment;
+import org.apache.seatunnel.engine.client.job.ClientJobProxy;
+import org.apache.seatunnel.engine.common.config.ConfigProvider;
+import org.apache.seatunnel.engine.common.config.JobConfig;
+import org.apache.seatunnel.engine.common.config.SeaTunnelConfig;
+import org.apache.seatunnel.engine.common.job.JobStatus;
+import org.apache.seatunnel.engine.e2e.TestUtils;
+import org.apache.seatunnel.engine.server.SeaTunnelServerStarter;
+import 
org.apache.seatunnel.engine.server.telemetry.metrics.exports.JobMetricExports;
+import 
org.apache.seatunnel.engine.server.telemetry.metrics.exports.JobThreadPoolStatusExports;
+
+import org.awaitility.Awaitility;
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+
+import com.hazelcast.client.config.ClientConfig;
+import com.hazelcast.instance.impl.HazelcastInstanceImpl;
+import io.prometheus.client.Collector;
+import lombok.extern.slf4j.Slf4j;
+
+import java.util.ArrayList;
+import java.util.List;
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.CopyOnWriteArrayList;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.ScheduledExecutorService;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
+
+/**
+ * Regression test for the Zeta engine master-node cold-start deadlock 
described in <a
+ * href="https://github.com/apache/seatunnel/issues/10840";>issue #10840</a> 
and fixed by <a
+ * href="https://github.com/apache/seatunnel/pull/10841";>PR #10841</a>.
+ *
+ * <p><b>The regression.</b> With {@code telemetry.metric.enabled=true},
+ * SeaTunnelServerStarter#initTelemetryInstance(Node) registers the Prometheus 
collectors
+ * synchronously right after a node joins the Hazelcast cluster - before that 
node's {@code
+ * CoordinatorService} has any chance to finish its own asynchronous 
activation (driven by a
+ * separate {@code masterActiveListener} thread that polls every 100ms and 
itself needs Hazelcast
+ * operation threads to read its coordination {@code IMap}s). Before the fix, 
{@code
+ * JobMetricExports#collect()} and {@code 
JobThreadPoolStatusExports#collect()} called {@code
+ * SeaTunnelServer#getCoordinatorService()} unconditionally whenever {@code 
isMaster()==true}, which
+ * blocks the calling thread for up to 1.5s (3 retries x 500ms) while the 
coordinator is still
+ * initializing. Under a full cluster force-restart (every node restarting 
simultaneously, e.g. a
+ * rolling upgrade or a datacenter power cycle), enough concurrently-blocked 
scrape calls pile onto
+ * the small Hazelcast operation-thread pool that {@code 
initCoordinatorService()} can never obtain
+ * a thread to finish its own {@code IMap} work, permanently deadlocking the 
cluster. The fix added
+ * a non-blocking {@code isCoordinatorReady()} guard so both collectors return 
empty immediately
+ * instead of blocking while the coordinator is not yet active.
+ *
+ * <p><b>The gap this closes.</b> The merged unit test for the guard itself,
+ * TelemetryCollectorCoordinatorGuardTest, proves the guard logic in isolation 
against mocked {@code
+ * SeaTunnelServer}/{@code CoordinatorService} objects, one collector 
invocation at a time. Nothing
+ * previously started a real multi-node cluster with telemetry enabled and 
proved that a genuinely
+ * concurrent full-cluster restart does not deadlock end-to-end through the 
actual startup sequence
+ * - Hazelcast join, telemetry registration, and {@code CoordinatorService} 
activation all racing
+ * for real, under a live metrics scraper hammering the exact code path the 
bug lived in. This test
+ * fills that gap.
+ *
+ * <p><b>Why this lives here, and why it looks different from {@link
+ * MasterWorkerClusterSeaTunnelWithTelemetryIT}.</b> That sibling test starts 
two
+ * Docker-containerized nodes SEQUENTIALLY (master, then worker) to prove 
per-node metric reporting
+ * and lite-worker promotion timing over the real REST API - a materially 
different scenario (timing
+ * of one worker's metrics becoming visible) from this test's target (whether 
a simultaneous cold
+ * start of the whole cluster deadlocks). Reproducing *this* regression 
requires precise control
+ * over concurrent startup timing that a cross-process Testcontainers/Docker 
harness cannot offer
+ * cheaply, so this test instead reuses the lighter in-process {@code 
HazelcastInstanceImpl} harness
+ * already established in this module by {@code 
org.apache.seatunnel.engine.e2e.ClusterIT} and
+ * {@code org.apache.seatunnel.engine.e2e.ClusterFailureNoRestoreIT}: real 
{@code
+ * HazelcastInstanceImpl} nodes created via {@code SeaTunnelServerStarter}, 
driven from the same JVM
+ * so multiple {@code createHazelcastInstance()} calls can be launched 
genuinely concurrently.
+ * Unlike those two classes (which only ever call {@code 
createHazelcastInstance()} sequentially and
+ * can therefore safely share one mutable {@code SeaTunnelConfig}), this test 
gives every node its
+ * own independently-loaded config, since concurrent construction from a 
single shared config object
+ * is not a pattern already proven safe anywhere else in this codebase.
+ *
+ * <p><b>A second, narrower race this test surfaced along the way.</b> With no 
prior cluster to
+ * join, {@value #NODE_COUNT} fully-concurrent nodes can each briefly form 
their own singleton
+ * Hazelcast cluster - and, being that singleton's sole member, briefly BE its 
master and activate
+ * their own {@code CoordinatorService} - before discovering each other and 
merging into the real
+ * {@value #NODE_COUNT}-member cluster, at which point some of them step down. 
A scrape can land in
+ * the narrow instant between {@code isCoordinatorReady()} reading true and 
{@code
+ * getReadyCoordinatorService()} re-checking it, right as that node steps 
down, paying one
+ * legitimate ~500-700ms retry even with the fix fully intact. This is a 
separate, narrow,
+ * at-most-once-per-node TOCTOU gap that PR #10841 does not claim to close, 
not the
+ * unconditional-blocking regression it fixed - it just happens to be visible 
through the same
+ * scrape-latency window this test measures, which shapes how {@link 
#MAX_ACCEPTABLE_SCRAPE_MILLIS}
+ * is calibrated below.
+ *
+ * <p><b>Regression-verification methodology and its limits.</b> Per this 
initiative's standing
+ * requirement to confirm a new regression test would actually have caught the 
bug it targets,
+ * {@code JobMetricExports#collect()} was temporarily reverted to the pre-fix 
pattern (unconditional
+ * {@code isMaster()} check calling the blocking {@code 
getCoordinatorService()}), {@code
+ * seatunnel-engine-server} was rebuilt, and this test was re-run repeatedly 
against both the
+ * reverted and the real fix, comparing scrape-latency and slow-scrape-count 
statistics. The honest
+ * result: with a tight per-scrape ceiling (200ms) the reverted build reliably 
failed (observed
+ * ~530-555ms); but a wider sweep then showed the REAL fix can occasionally 
hit a similar ~500-700ms
+ * latency too, via the legitimate step-down TOCTOU race above, and the two 
builds' slow-scrape
+ * counts and totals were statistically indistinguishable across repeated runs 
(4-7 slow scrapes out
+ * of roughly 1200-1550 total, in both the fixed and the reverted build). The 
reason is structural,
+ * not a flaw specific to one threshold choice: this scrape hammer invokes 
{@code collect()}
+ * directly from test-owned threads, not through Hazelcast's own 
operation-thread pool the way a
+ * real Prometheus scrape or the "/hazelcast/rest/instance/metrics" handler 
would, so it cannot
+ * reproduce the genuine operation-thread-pool exhaustion that turned 
individually-slow scrapes into
+ * the permanent cluster-wide deadlock in issue #10840 - it can only observe 
that {@code collect()}
+ * itself remains fast and exception-free. Given that, {@link 
#MAX_ACCEPTABLE_SCRAPE_MILLIS} is
+ * deliberately set as a generous sanity backstop (comfortably above the 
observed legitimate range,
+ * still well below sustained-starvation territory) rather than a precise 
discriminator, and this
+ * test's actual, reliable regression signal is the Awaitility timeouts above: 
the bug report
+ * describes the real deadlock as permanent ("never recovers"), so if 
operation threads were
+ * genuinely exhausted, cluster convergence and job scheduling would not 
merely be slow, they would
+ * never complete, and those bounded waits would fail outright.
+ */
+@Slf4j
+public class TelemetryStartupDeadlockIT {
+
+    /**
+     * Number of nodes started concurrently to simulate a "full cluster 
force-restart". More than
+     * the bare minimum of 2 so the race involves genuine multi-node 
contention (one node eventually
+     * wins Hazelcast master election while the others do not), closer to the 
production incident
+     * than a minimal two-node reproduction would be.
+     */
+    private static final int NODE_COUNT = 3;
+
+    /**
+     * Bound for each individual node's {@code createHazelcastInstance()} 
call. Generous relative to
+     * a normal single-node startup (well under a second on an idle machine) 
to absorb this shared
+     * machine's disk/CPU contention and the concurrent metrics-scrape hammer 
started alongside it;
+     * still finite, so a hung node-start fails the test instead of hanging CI.
+     */
+    private static final long NODE_START_TIMEOUT_SECONDS = 60L;
+
+    /**
+     * Bound for all {@value #NODE_COUNT} nodes to observe each other as 
cluster members.
+     * ClusterIT's equivalent two-node, sequential-start check uses 10s; this 
budget is 6x more
+     * generous to cover a concurrent {@value #NODE_COUNT}-node start plus the 
concurrent scrape
+     * hammer under contention, while remaining clearly finite. The pre-fix 
deadlock is described in
+     * issue #10840 as permanent (the cluster "never recovers"), so any finite 
timeout is sufficient
+     * to catch a regression here - the generosity is purely to avoid false 
failures from legitimate
+     * slowness on this shared machine, not because the fix is expected to be 
slow.
+     */
+    private static final long CLUSTER_READY_TIMEOUT_SECONDS = 60L;
+
+    /**
+     * Bound for the submitted job to reach RUNNING. Matches the 2-minute 
budget already used by the
+     * Docker-based MasterWorkerClusterSeaTunnelWithTelemetryIT and 
JobLogUrlPortIT for the same
+     * kind of check, so this in-process test does not invent a materially 
different tolerance for
+     * what is functionally the same "job scheduling actually works" assertion.
+     */
+    private static final long JOB_RUNNING_TIMEOUT_MINUTES = 2L;
+
+    /**
+     * Interval between metrics scrapes fired at each node while the cluster 
is starting. Tight
+     * enough (50 scrapes/second/node) to have a realistic chance of landing 
calls inside the short
+     * window between a node winning Hazelcast master election and its {@code 
CoordinatorService}
+     * finishing activation.
+     */
+    private static final long SCRAPE_INTERVAL_MILLIS = 20L;
+
+    /**
+     * Ceiling for a single scrape's wall-clock time, deliberately generous - 
see the class Javadoc
+     * section "Regression-verification methodology and its limits" for the 
full reasoning and the
+     * experimental data behind this number. Summary: the pre-fix blocking 
retry loop in {@code
+     * SeaTunnelServer#getCoordinatorService()} sleeps in increments of 500ms 
(up to 3 times, 1500ms
+     * total); but this test independently discovered that starting {@value 
#NODE_COUNT} nodes fully
+     * concurrently with no prior cluster to join makes each one briefly form 
its own singleton
+     * cluster and activate its own {@code CoordinatorService} before merging 
and (for some)
+     * stepping down, and a scrape landing on that legitimate step-down TOCTOU 
race pays one
+     * ~500-700ms retry even with the fix fully intact - observed consistently 
across both the fixed
+     * code and an artificially-reverted pre-fix build, so a tight ceiling 
cannot reliably separate
+     * the two in this harness. This ceiling is set comfortably above that 
observed legitimate range
+     * while still well below what sustained operation-thread-pool starvation 
would look like, so it
+     * exists as a sanity backstop; the Awaitility timeouts above (cluster 
convergence, job reaching
+     * RUNNING) are this test's primary, reliable regression signal.
+     */
+    private static final long MAX_ACCEPTABLE_SCRAPE_MILLIS = 1200L;
+
+    /**
+     * Starts {@value #NODE_COUNT} Hazelcast nodes concurrently with telemetry 
metrics enabled -
+     * reproducing the "full cluster force-restart" scenario from issue #10840 
- while a concurrent
+     * metrics-scrape hammer exercises the exact collectors the fix guards, 
then proves the cluster
+     * still becomes fully healthy and functional within a bounded timeout.
+     *
+     * @throws Exception if job submission or cluster startup fails 
unexpectedly
+     */
+    @Test
+    public void testConcurrentClusterStartupWithTelemetryDoesNotDeadlock() 
throws Exception {
+        String testClusterName = 
TestUtils.getClusterName("TelemetryStartupDeadlockIT");
+
+        List<HazelcastInstanceImpl> nodes = new CopyOnWriteArrayList<>();
+        SeaTunnelClient engineClient = null;
+        ExecutorService nodeStartExecutor = 
Executors.newFixedThreadPool(NODE_COUNT);
+        ScheduledExecutorService scrapeExecutor = 
Executors.newScheduledThreadPool(NODE_COUNT);
+        AtomicBoolean keepScraping = new AtomicBoolean(true);
+        List<Long> scrapeDurationsMillis = new CopyOnWriteArrayList<>();
+        List<Throwable> scrapeFailures = new CopyOnWriteArrayList<>();
+
+        try {
+            // Build every node's config up front, sequentially, before any 
concurrency starts -
+            // see the class Javadoc on why each node gets its own independent 
config object
+            // rather than sharing one across the concurrent 
createHazelcastInstance() calls
+            // below.
+            List<SeaTunnelConfig> nodeConfigs = new ArrayList<>();
+            for (int i = 0; i < NODE_COUNT; i++) {
+                nodeConfigs.add(buildNodeConfig(testClusterName));
+            }
+
+            // Genuinely concurrent startup: submit all node-creation calls at 
once so Hazelcast
+            // join, telemetry registration, and CoordinatorService activation 
race across nodes
+            // exactly like the issue's "all nodes restart simultaneously" 
scenario. A sequential
+            // start (as ClusterIT/ClusterFailureNoRestoreIT do for their own 
unrelated scenarios)
+            // would let node1 finish activating long before node2 even exists 
and would never
+            // reproduce the race.
+            List<CompletableFuture<HazelcastInstanceImpl>> nodeFutures = new 
ArrayList<>();
+            for (SeaTunnelConfig nodeConfig : nodeConfigs) {
+                CompletableFuture<HazelcastInstanceImpl> future =
+                        CompletableFuture.supplyAsync(
+                                () -> 
SeaTunnelServerStarter.createHazelcastInstance(nodeConfig),
+                                nodeStartExecutor);
+                // The instant this node finishes joining Hazelcast and 
registering its telemetry
+                // collectors, start hammering ITS metrics at a tight fixed 
interval. This is what
+                // actually exercises the vulnerable path: registering a 
collector never invokes
+                // collect() by itself, so the original deadlock required a 
live scraper hitting
+                // the metrics endpoint throughout the startup race window, 
the same way a
+                // Prometheus server or the "/hazelcast/rest/instance/metrics" 
health check would
+                // in production.
+                future.thenAcceptAsync(
+                        node ->
+                                startMetricsScrapeHammer(
+                                        node,
+                                        scrapeExecutor,
+                                        keepScraping,
+                                        scrapeDurationsMillis,
+                                        scrapeFailures),
+                        scrapeExecutor);
+                nodeFutures.add(future);
+            }
+            for (CompletableFuture<HazelcastInstanceImpl> future : 
nodeFutures) {
+                nodes.add(future.get(NODE_START_TIMEOUT_SECONDS, 
TimeUnit.SECONDS));
+            }
+
+            Awaitility.await()
+                    .atMost(CLUSTER_READY_TIMEOUT_SECONDS, TimeUnit.SECONDS)
+                    .untilAsserted(
+                            () ->
+                                    Assertions.assertEquals(
+                                            NODE_COUNT,
+                                            
nodes.get(0).getCluster().getMembers().size(),
+                                            "all nodes must converge into a 
single cluster"));
+
+            List<HazelcastInstanceImpl> masterNodes = new ArrayList<>();
+            for (HazelcastInstanceImpl node : nodes) {
+                if (node.node.isMaster()) {
+                    masterNodes.add(node);
+                }
+            }
+            Assertions.assertEquals(
+                    1, masterNodes.size(), "exactly one node must hold 
Hazelcast master role");
+            HazelcastInstanceImpl masterNode = masterNodes.get(0);
+
+            // Prove the cluster is genuinely functional, not just "didn't 
crash": submit and run
+            // a real job. Scheduling a job requires 
CoordinatorService#isCoordinatorActive() to
+            // be true on the master node, so reaching RUNNING is direct proof 
it finished
+            // initializing despite racing telemetry init and the other nodes' 
concurrent joins.
+            ClientConfig clientConfig = 
ConfigProvider.locateAndGetClientConfig();
+            clientConfig.setClusterName(testClusterName);
+            engineClient = new SeaTunnelClient(clientConfig);
+
+            String confPath = 
TestUtils.getResource("stream_fakesource_to_console.conf");
+            JobConfig jobConfig = new JobConfig();
+            jobConfig.setName("TelemetryStartupDeadlockIT_job");
+            ClientJobExecutionEnvironment jobExecutionEnv =
+                    engineClient.createExecutionContext(confPath, jobConfig, 
nodeConfigs.get(0));
+            ClientJobProxy clientJobProxy = jobExecutionEnv.execute();
+
+            Awaitility.await()
+                    .atMost(JOB_RUNNING_TIMEOUT_MINUTES, TimeUnit.MINUTES)
+                    .untilAsserted(
+                            () ->
+                                    Assertions.assertEquals(
+                                            JobStatus.RUNNING, 
clientJobProxy.getJobStatus()));
+
+            // Stop hammering only after the whole race window (cluster join, 
coordinator
+            // activation, and first job scheduling) is long over, so the 
assertions below read a
+            // stable, complete set of samples.
+            keepScraping.set(false);
+            scrapeExecutor.shutdown();
+            Assertions.assertTrue(
+                    scrapeExecutor.awaitTermination(10, TimeUnit.SECONDS),
+                    "scrape hammer executor did not terminate cleanly");
+
+            Assertions.assertTrue(
+                    scrapeFailures.isEmpty(),
+                    "collect() must never throw while the coordinator is 
initializing, but got: "
+                            + scrapeFailures);
+            Assertions.assertFalse(
+                    scrapeDurationsMillis.isEmpty(), "the metrics scrape 
hammer never ran");
+            long maxObservedScrapeMillis =
+                    
scrapeDurationsMillis.stream().mapToLong(Long::longValue).max().orElse(0L);
+            // See the class Javadoc "Regression-verification methodology and 
its limits" section
+            // for why this ceiling is generous rather than tight, and why the 
Awaitility timeouts
+            // above - not this check - are this test's primary regression 
signal. In short: this
+            // scrape hammer calls collect() directly from test-owned threads, 
not through
+            // Hazelcast's own operation-thread pool the way a real 
Prometheus/HTTP scrape would,
+            // so it cannot reproduce genuine operation-thread-pool 
exhaustion; and experimentally,
+            // a single legitimate mastership-step-down TOCTOU race (present 
even with the fix
+            // intact, see below) already costs one ~500-700ms retry 
indistinguishable in isolation
+            // from the guard being absent. This ceiling is set well above 
that observed legitimate
+            // range while still well below what sustained thread-pool 
starvation would look like
+            // (many scrapes each pinned near the full 1500ms retry ceiling 
for as long as the
+            // pool stays exhausted, not one or two isolated ~500-700ms blips).
+            Assertions.assertTrue(
+                    maxObservedScrapeMillis < MAX_ACCEPTABLE_SCRAPE_MILLIS,
+                    "a metrics scrape took "
+                            + maxObservedScrapeMillis
+                            + "ms during the startup race (out of "
+                            + scrapeDurationsMillis.size()
+                            + " total scrapes); expected under "
+                            + MAX_ACCEPTABLE_SCRAPE_MILLIS
+                            + "ms");
+
+            // Finally, assert the telemetry/metrics path itself is reachable 
and reports the
+            // running job, following the same job_count{type="running"} 
assertion style already
+            // used by 
MasterWorkerClusterSeaTunnelWithTelemetryIT#testGetMetrics.
+            List<Collector.MetricFamilySamples> jobMetrics =
+                    new JobMetricExports(masterNode.node).collect();
+            assertRunningJobCount(jobMetrics, 1.0D);
+
+            List<Collector.MetricFamilySamples> poolMetrics =
+                    new JobThreadPoolStatusExports(masterNode.node).collect();
+            Assertions.assertFalse(
+                    poolMetrics.isEmpty(),
+                    "job_thread_pool_* metrics must be reachable on the active 
master node");
+
+            // The guard must still correctly suppress these same metrics on 
non-master nodes,
+            // matching TelemetryCollectorCoordinatorGuardTest's mock-based 
coverage but now
+            // proven against a real multi-node cluster.
+            for (HazelcastInstanceImpl node : nodes) {
+                if (node == masterNode) {
+                    continue;
+                }
+                Assertions.assertTrue(
+                        new JobMetricExports(node.node).collect().isEmpty(),
+                        "non-master nodes must never report job metrics");
+            }
+        } finally {
+            keepScraping.set(false);
+            scrapeExecutor.shutdownNow();
+            nodeStartExecutor.shutdownNow();
+            if (engineClient != null) {
+                engineClient.close();
+            }
+            for (HazelcastInstanceImpl node : nodes) {
+                if (node != null) {
+                    node.shutdown();
+                }
+            }
+        }
+    }
+
+    /**
+     * Builds an independent {@link SeaTunnelConfig} for one node, with the 
trigger condition from
+     * issue #10840 ({@code telemetry.metric.enabled=true}) turned on and 
SeaTunnel's own Jetty REST
+     * server turned off (it is irrelevant to this scenario and would 
otherwise make all {@value
+     * #NODE_COUNT} in-process nodes race to bind the same fixed HTTP port).
+     *
+     * @param clusterName the shared Hazelcast cluster name every node must 
join
+     * @return a freshly-loaded config, independent from any other node's
+     */
+    private SeaTunnelConfig buildNodeConfig(String clusterName) {
+        SeaTunnelConfig seaTunnelConfig = 
ConfigProvider.locateAndGetSeaTunnelConfig();
+        seaTunnelConfig.getHazelcastConfig().setClusterName(clusterName);
+        seaTunnelConfig.getEngineConfig().getHttpConfig().setEnabled(false);
+        
seaTunnelConfig.getEngineConfig().getTelemetryConfig().getMetric().setEnabled(true);
+        return seaTunnelConfig;
+    }
+
+    /**
+     * Repeatedly invokes the two collectors issue #10840 fixed - {@link 
JobMetricExports} and
+     * {@link JobThreadPoolStatusExports} - at a tight fixed interval, 
starting the instant the
+     * given node finishes joining the cluster, recording each call's 
wall-clock duration and any
+     * thrown exception for later assertion.
+     *
+     * @param node the node to scrape
+     * @param scrapeExecutor executor the periodic scrape task runs on
+     * @param keepScraping flag the caller flips to stop future scrapes 
without cancelling in-flight
+     *     ones
+     * @param scrapeDurationsMillis thread-safe sink for each successful 
scrape's duration
+     * @param scrapeFailures thread-safe sink for any exception a scrape throws
+     */
+    private void startMetricsScrapeHammer(
+            HazelcastInstanceImpl node,
+            ScheduledExecutorService scrapeExecutor,
+            AtomicBoolean keepScraping,
+            List<Long> scrapeDurationsMillis,
+            List<Throwable> scrapeFailures) {
+        scrapeExecutor.scheduleAtFixedRate(
+                () -> {
+                    if (!keepScraping.get()) {
+                        return;
+                    }
+                    long startNanos = System.nanoTime();
+                    try {
+                        new JobMetricExports(node.node).collect();
+                        new JobThreadPoolStatusExports(node.node).collect();
+                        scrapeDurationsMillis.add(
+                                
TimeUnit.NANOSECONDS.toMillis(System.nanoTime() - startNanos));
+                    } catch (Throwable t) {
+                        // Caught deliberately so a scrape failure never 
cancels this periodic
+                        // task's future executions (ScheduledExecutorService 
drops a periodic
+                        // task after an uncaught exception) - we want every 
failure recorded, not
+                        // just the first.
+                        scrapeFailures.add(t);
+                    }
+                },
+                0,
+                SCRAPE_INTERVAL_MILLIS,
+                TimeUnit.MILLISECONDS);
+    }
+
+    /**
+     * Asserts the {@code job_count} metric family is present and reports the 
expected value for
+     * {@code type="running"}, mirroring the assertion style already used by 
{@link
+     * MasterWorkerClusterSeaTunnelWithTelemetryIT#testGetMetrics}.
+     *
+     * @param jobMetrics samples returned by {@link JobMetricExports#collect()}
+     * @param expectedRunningCount expected value of the {@code 
job_count{type="running"}} sample
+     */
+    private void assertRunningJobCount(
+            List<Collector.MetricFamilySamples> jobMetrics, double 
expectedRunningCount) {
+        Collector.MetricFamilySamples jobCountFamily =
+                jobMetrics.stream()
+                        .filter(family -> "job_count".equals(family.name))
+                        .findFirst()
+                        .orElse(null);
+        Assertions.assertNotNull(jobCountFamily, "job_count metric family must 
be present");
+
+        Collector.MetricFamilySamples.Sample runningSample =
+                jobCountFamily.samples.stream()
+                        .filter(
+                                sample -> {
+                                    int typeIndex = 
sample.labelNames.indexOf("type");
+                                    return typeIndex >= 0
+                                            && 
"running".equals(sample.labelValues.get(typeIndex));
+                                })
+                        .findFirst()
+                        .orElse(null);
+        Assertions.assertNotNull(
+                runningSample, "job_count{type=\"running\"} sample must be 
present");
+        Assertions.assertEquals(expectedRunningCount, runningSample.value, 
0.0001D);
+    }
+}

Reply via email to