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

quantranhong1999 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/james-project.git

commit 41c8d9e48824e1a18e3f60acff15a85fd69612eb
Author: Quan Tran <[email protected]>
AuthorDate: Thu Oct 1 13:29:58 2026 +0700

    [IMPROVEMENT] One RabbitMQ consumers health check fed by 
MonitoredRabbitMQConsumers
    
    Each consumer health check was its own class. Modules now only contribute
    MonitoredRabbitMQConsumers (name, connection pool of the RabbitMQ server,
    queues, restart) with @ProvidesIntoSet, and a single
    RabbitMQConsumersHealthCheck checks them all: only the contributions having
    a queue without consumer are restarted. Downstream projects can monitor
    their consumers, possibly on another RabbitMQ server, without a dedicated
    class.
    
    Contributions are checked concurrently, each on its own thread, so that an
    unresponsive RabbitMQ server does not delay the others, but restarted one
    at a time, as restarting the same component concurrently may not be safe. A
    restart is given the connection used by the check, composed reactively so
    that the check never blocks on another connection lookup, and run on the
    blocking call scheduler as restarts may block. Each restart and each error
    is logged with the contribution it concerns.
    
    The event bus, distributed task manager and mail queue consumer health
    checks are replaced by contributions (RabbitMQEventBusConsumers,
    RabbitMQMailQueueConsumers, and MonitoredRabbitMQConsumers.of for the task
    manager). As before, restarting the mail queue consumers runs every
    ReconnectionHandler, hence also restarts the other James consumers.
    
    Visible changes:
     - /healthcheck reports a single RabbitMQConsumers component instead of
       EventbusConsumers-mailboxEvent, EventbusConsumers-jmapEvent,
       EventbusConsumers-contentDeletionEvent, DistributedTaskManagerConsumers
       and MailQueueConsumers. Its cause names the queue without consumer.
     - errors are reported as unhealthy, without hiding the other consumers,
       instead of failing the whole health check call.
    
    Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
 .../rabbitmq/MonitoredRabbitMQConsumers.java       |  79 ++++++
 .../rabbitmq/RabbitMQConsumersHealthCheck.java     | 110 ++++++++
 .../rabbitmq/RabbitMQConsumersHealthCheckTest.java | 278 +++++++++++++++++++++
 ...thCheck.java => RabbitMQEventBusConsumers.java} |  62 ++---
 .../events/RabbitMQEventBusConsumersTest.java      |  80 ++++++
 .../modules/task/DistributedTaskManagerModule.java |  14 +-
 .../modules/DistributedTaskManagerModule.java      |  14 +-
 .../event/ContentDeletionEventBusModule.java       |  11 +-
 .../james/modules/event/JMAPEventBusModule.java    |  11 +-
 .../james/modules/event/MailboxEventBusModule.java |  11 +-
 .../queue/rabbitmq/RabbitMQMailQueueModule.java    |  10 +-
 .../modules/queue/rabbitmq/RabbitMQModule.java     |   4 +
 ...itMQWebAdminServerIntegrationImmutableTest.java |  29 ++-
 ...hCheck.java => RabbitMQMailQueueConsumers.java} |  56 ++---
 .../RabbitMQMailQueueConsumerHealthCheckTest.java  |  92 -------
 .../DistributedTaskManagerHealthCheck.java         |  77 ------
 .../distributed/RabbitMQWorkQueue.java             |   2 +-
 17 files changed, 665 insertions(+), 275 deletions(-)

diff --git 
a/backends-common/rabbitmq/src/main/java/org/apache/james/backends/rabbitmq/MonitoredRabbitMQConsumers.java
 
b/backends-common/rabbitmq/src/main/java/org/apache/james/backends/rabbitmq/MonitoredRabbitMQConsumers.java
new file mode 100644
index 0000000000..655afb8114
--- /dev/null
+++ 
b/backends-common/rabbitmq/src/main/java/org/apache/james/backends/rabbitmq/MonitoredRabbitMQConsumers.java
@@ -0,0 +1,79 @@
+/****************************************************************
+ * 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.james.backends.rabbitmq;
+
+import java.util.List;
+import java.util.function.Function;
+import java.util.function.Supplier;
+
+import org.reactivestreams.Publisher;
+
+import com.rabbitmq.client.Connection;
+
+/**
+ * RabbitMQ consumers watched by {@link RabbitMQConsumersHealthCheck}: they 
are restarted when one of their queues has
+ * no consumer.
+ */
+public interface MonitoredRabbitMQConsumers {
+    static MonitoredRabbitMQConsumers of(String name, SimpleConnectionPool 
connectionPool,
+                                         Supplier<List<String>> queues, 
Function<Connection, Publisher<Void>> restart) {
+        return new MonitoredRabbitMQConsumers() {
+            @Override
+            public String name() {
+                return name;
+            }
+
+            @Override
+            public SimpleConnectionPool connectionPool() {
+                return connectionPool;
+            }
+
+            @Override
+            public List<String> queues() {
+                return queues.get();
+            }
+
+            @Override
+            public Publisher<Void> restart(Connection connection) {
+                return restart.apply(connection);
+            }
+        };
+    }
+
+    /**
+     * Shown in the health check cause.
+     */
+    String name();
+
+    /**
+     * Connection pool of the RabbitMQ server hosting the queues.
+     */
+    SimpleConnectionPool connectionPool();
+
+    /**
+     * Called on each check, as the monitored queues can change over time.
+     */
+    List<String> queues();
+
+    /**
+     * Restarts the consumers, given the connection used by the check.
+     */
+    Publisher<Void> restart(Connection connection);
+}
diff --git 
a/backends-common/rabbitmq/src/main/java/org/apache/james/backends/rabbitmq/RabbitMQConsumersHealthCheck.java
 
b/backends-common/rabbitmq/src/main/java/org/apache/james/backends/rabbitmq/RabbitMQConsumersHealthCheck.java
new file mode 100644
index 0000000000..afabf11dc3
--- /dev/null
+++ 
b/backends-common/rabbitmq/src/main/java/org/apache/james/backends/rabbitmq/RabbitMQConsumersHealthCheck.java
@@ -0,0 +1,110 @@
+/****************************************************************
+ * 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.james.backends.rabbitmq;
+
+import static org.apache.james.util.ReactorUtils.DEFAULT_CONCURRENCY;
+
+import java.io.IOException;
+import java.util.Optional;
+import java.util.Set;
+import java.util.concurrent.TimeoutException;
+import java.util.function.Function;
+
+import jakarta.inject.Inject;
+
+import org.apache.james.core.healthcheck.ComponentName;
+import org.apache.james.core.healthcheck.HealthCheck;
+import org.apache.james.core.healthcheck.Result;
+import org.apache.james.util.ReactorUtils;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import com.github.fge.lambdas.Throwing;
+import com.google.common.collect.ImmutableList;
+import com.rabbitmq.client.Channel;
+import com.rabbitmq.client.Connection;
+
+import reactor.core.publisher.Flux;
+import reactor.core.publisher.Mono;
+
+public class RabbitMQConsumersHealthCheck implements HealthCheck {
+    public static final ComponentName COMPONENT_NAME = new 
ComponentName("RabbitMQConsumers");
+    private static final Logger LOGGER = 
LoggerFactory.getLogger(RabbitMQConsumersHealthCheck.class);
+
+    private final ImmutableList<MonitoredRabbitMQConsumers> monitoredConsumers;
+
+    @Inject
+    public RabbitMQConsumersHealthCheck(Set<MonitoredRabbitMQConsumers> 
monitoredConsumers) {
+        this.monitoredConsumers = ImmutableList.copyOf(monitoredConsumers);
+    }
+
+    @Override
+    public ComponentName componentName() {
+        return COMPONENT_NAME;
+    }
+
+    @Override
+    public Mono<Result> check() {
+        return Flux.fromIterable(monitoredConsumers)
+            .flatMap(this::detect, DEFAULT_CONCURRENCY)
+            .concatMap(Function.identity())
+            .collectList()
+            .map(results -> MergedResults.merge(COMPONENT_NAME, results));
+    }
+
+    /**
+     * Detections run concurrently, each on its own thread, so that an 
unresponsive RabbitMQ server does not delay the
+     * other consumers. The returned result restarts the consumers once 
subscribed, and {@link #check()} subscribes to
+     * those one at a time: restarting the same component concurrently is not 
safe.
+     */
+    private Mono<Mono<Result>> detect(MonitoredRabbitMQConsumers consumers) {
+        return consumers.connectionPool().getResilientConnection()
+            .flatMap(connection -> Mono.fromCallable(() -> 
queueWithoutConsumers(consumers, connection))
+                .map(queueWithoutConsumers -> queueWithoutConsumers
+                    .map(queue -> restart(consumers, connection, queue))
+                    .orElseGet(() -> 
Mono.just(Result.healthy(COMPONENT_NAME)))))
+            .onErrorResume(e -> Mono.just(Mono.just(unhealthy(consumers, e))))
+            .subscribeOn(ReactorUtils.BLOCKING_CALL_WRAPPER);
+    }
+
+    private Mono<Result> restart(MonitoredRabbitMQConsumers consumers, 
Connection connection, String queue) {
+        return Mono.defer(() -> {
+                LOGGER.warn("No consumers on {} of {}, restarting them", 
queue, consumers.name());
+                return Mono.from(consumers.restart(connection));
+            })
+            // Restarts may block, and the previous restart may have completed 
on a non-blocking thread
+            .subscribeOn(ReactorUtils.BLOCKING_CALL_WRAPPER)
+            .thenReturn(Result.degraded(COMPONENT_NAME, String.format("No 
consumers on %s of %s", queue, consumers.name())))
+            .onErrorResume(e -> Mono.just(unhealthy(consumers, e)));
+    }
+
+    private Result unhealthy(MonitoredRabbitMQConsumers consumers, Throwable 
error) {
+        LOGGER.warn("Error checking the consumers of {}", consumers.name(), 
error);
+        return Result.unhealthy(COMPONENT_NAME, "Error checking the consumers 
of " + consumers.name(), error);
+    }
+
+    private Optional<String> queueWithoutConsumers(MonitoredRabbitMQConsumers 
consumers, Connection connection) throws IOException, TimeoutException {
+        try (Channel channel = connection.createChannel()) {
+            return consumers.queues().stream()
+                .filter(Throwing.predicate(queue -> 
channel.consumerCount(queue) == 0))
+                .findFirst();
+        }
+    }
+}
diff --git 
a/backends-common/rabbitmq/src/test/java/org/apache/james/backends/rabbitmq/RabbitMQConsumersHealthCheckTest.java
 
b/backends-common/rabbitmq/src/test/java/org/apache/james/backends/rabbitmq/RabbitMQConsumersHealthCheckTest.java
new file mode 100644
index 0000000000..014fb5720c
--- /dev/null
+++ 
b/backends-common/rabbitmq/src/test/java/org/apache/james/backends/rabbitmq/RabbitMQConsumersHealthCheckTest.java
@@ -0,0 +1,278 @@
+/****************************************************************
+ * 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.james.backends.rabbitmq;
+
+import static org.apache.james.backends.rabbitmq.Constants.AUTO_ACK;
+import static org.apache.james.backends.rabbitmq.Constants.AUTO_DELETE;
+import static org.apache.james.backends.rabbitmq.Constants.DURABLE;
+import static org.apache.james.backends.rabbitmq.Constants.EXCLUSIVE;
+import static 
org.apache.james.backends.rabbitmq.RabbitMQFixture.DEFAULT_MANAGEMENT_CREDENTIAL;
+import static 
org.apache.james.backends.rabbitmq.RabbitMQFixture.awaitAtMostOneMinute;
+import static org.assertj.core.api.Assertions.assertThat;
+
+import java.net.URI;
+import java.time.Duration;
+import java.util.List;
+import java.util.Optional;
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.CopyOnWriteArrayList;
+import java.util.concurrent.atomic.AtomicInteger;
+import java.util.concurrent.atomic.AtomicReference;
+import java.util.function.Function;
+import java.util.function.Supplier;
+
+import org.apache.james.core.healthcheck.Result;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.extension.RegisterExtension;
+import org.reactivestreams.Publisher;
+
+import com.google.common.collect.ImmutableMap;
+import com.google.common.collect.ImmutableSet;
+import com.rabbitmq.client.Channel;
+import com.rabbitmq.client.Connection;
+
+import reactor.core.Disposable;
+import reactor.core.publisher.Mono;
+import reactor.core.scheduler.Schedulers;
+
+class RabbitMQConsumersHealthCheckTest {
+    private static final String SECOND_VHOST = "second";
+    private static final String FIRST_QUEUE = "first-queue";
+    private static final String SECOND_QUEUE = "second-queue";
+
+    @RegisterExtension
+    static RabbitMQExtension rabbitMQExtension = 
RabbitMQExtension.singletonRabbitMQ()
+        .isolationPolicy(RabbitMQExtension.IsolationPolicy.STRONG);
+
+    private final List<String> firstQueues = new 
CopyOnWriteArrayList<>(List.of(FIRST_QUEUE));
+    private final AtomicInteger firstRestarts = new AtomicInteger();
+    private final AtomicInteger secondRestarts = new AtomicInteger();
+    private final AtomicReference<Connection> firstRestartConnection = new 
AtomicReference<>();
+    private final AtomicReference<Connection> secondRestartConnection = new 
AtomicReference<>();
+    private SimpleConnectionPool secondVhostConnectionPool;
+    private Channel firstChannel;
+    private Channel secondChannel;
+    private MonitoredRabbitMQConsumers first;
+    private MonitoredRabbitMQConsumers second;
+
+    /**
+     * The second consumers' queue only exists in the second vhost: checking 
it with the connection pool of the first
+     * consumers fails.
+     */
+    @BeforeEach
+    void setUp() throws Exception {
+        DockerRabbitMQ rabbitMQ = rabbitMQExtension.getRabbitMQ();
+        rabbitMQ.container().execInContainer("rabbitmqctl", "add_vhost", 
SECOND_VHOST);
+        rabbitMQ.container().execInContainer("rabbitmqctl", "set_permissions", 
"-p", SECOND_VHOST, rabbitMQ.getUsername(), ".*", ".*", ".*");
+        secondVhostConnectionPool = new SimpleConnectionPool(new 
RabbitMQConnectionFactory(RabbitMQConfiguration.builder()
+                .amqpUri(URI.create(rabbitMQ.amqpUri() + "/" + SECOND_VHOST))
+                .managementUri(rabbitMQ.managementUri())
+                .managementCredentials(DEFAULT_MANAGEMENT_CREDENTIAL)
+                .vhost(Optional.of(SECOND_VHOST))
+                .build()),
+            SimpleConnectionPool.Configuration.DEFAULT);
+
+        firstChannel = 
rabbitMQExtension.getConnectionPool().getResilientConnection().block().createChannel();
+        secondChannel = 
secondVhostConnectionPool.getResilientConnection().block().createChannel();
+        declareQueue(firstChannel, FIRST_QUEUE);
+        declareQueue(secondChannel, SECOND_QUEUE);
+
+        first = MonitoredRabbitMQConsumers.of("first consumers", 
rabbitMQExtension.getConnectionPool(), () -> List.copyOf(firstQueues),
+            connection -> Mono.fromRunnable(() -> {
+                firstRestarts.incrementAndGet();
+                firstRestartConnection.set(connection);
+            }));
+        second = MonitoredRabbitMQConsumers.of("second consumers", 
secondVhostConnectionPool, () -> List.of(SECOND_QUEUE),
+            connection -> Mono.fromRunnable(() -> {
+                secondRestarts.incrementAndGet();
+                secondRestartConnection.set(connection);
+            }));
+    }
+
+    @AfterEach
+    void tearDown() throws Exception {
+        firstChannel.close();
+        secondChannel.close();
+        secondVhostConnectionPool.close();
+    }
+
+    @Test
+    void checkShouldReturnHealthyWhenNoConsumerIsMonitored() {
+        Result result = new 
RabbitMQConsumersHealthCheck(ImmutableSet.of()).check().block();
+
+        assertThat(result.isHealthy()).isTrue();
+    }
+
+    @Test
+    void checkShouldReturnHealthyWithoutRestartWhenAllQueuesHaveConsumers() 
throws Exception {
+        consume(firstChannel, FIRST_QUEUE);
+        consume(secondChannel, SECOND_QUEUE);
+
+        Result result = new 
RabbitMQConsumersHealthCheck(ImmutableSet.of(first, second)).check().block();
+
+        assertThat(result.isHealthy()).isTrue();
+        
assertThat(result.getComponentName()).isEqualTo(RabbitMQConsumersHealthCheck.COMPONENT_NAME);
+        assertThat(firstRestarts).hasValue(0);
+        assertThat(secondRestarts).hasValue(0);
+    }
+
+    @Test
+    void checkShouldRestartOnlyTheConsumersWhoseQueueHasNoConsumer() throws 
Exception {
+        consume(secondChannel, SECOND_QUEUE);
+
+        Result result = new 
RabbitMQConsumersHealthCheck(ImmutableSet.of(first, second)).check().block();
+
+        assertThat(result.isDegraded()).isTrue();
+        assertThat(result.getCause()).contains("No consumers on " + 
FIRST_QUEUE + " of first consumers");
+        assertThat(firstRestarts).hasValue(1);
+        assertThat(secondRestarts).hasValue(0);
+    }
+
+    @Test
+    void checkShouldRestartConsumersWithAConnectionOfTheirOwnPool() {
+        new RabbitMQConsumersHealthCheck(ImmutableSet.of(first, 
second)).check().block();
+
+        
assertThat(firstRestartConnection.get()).isSameAs(rabbitMQExtension.getConnectionPool().getResilientConnection().block());
+        
assertThat(secondRestartConnection.get()).isSameAs(secondVhostConnectionPool.getResilientConnection().block());
+    }
+
+    @Test
+    void checkShouldStillCheckOtherConsumersWhenARestartFails() {
+        MonitoredRabbitMQConsumers failingRestart = 
MonitoredRabbitMQConsumers.of("failing consumers", 
rabbitMQExtension.getConnectionPool(),
+            () -> List.of(FIRST_QUEUE), connection -> Mono.error(new 
RuntimeException("Restart failure")));
+
+        Result result = new 
RabbitMQConsumersHealthCheck(ImmutableSet.of(failingRestart, 
second)).check().block();
+
+        assertThat(result.isUnHealthy()).isTrue();
+        assertThat(result.getCause()).hasValueSatisfying(cause -> 
assertThat(cause)
+            .contains("Error checking the consumers of failing consumers")
+            .contains("No consumers on " + SECOND_QUEUE + " of second 
consumers"));
+        assertThat(secondRestarts).hasValue(1);
+    }
+
+    @Test
+    void checkShouldRestartOtherConsumersWhenARabbitMQServerIsUnresponsive() 
throws Exception {
+        try (UnresponsiveServer unresponsiveServer = new UnresponsiveServer();
+             SimpleConnectionPool unresponsivePool = new 
SimpleConnectionPool(new 
RabbitMQConnectionFactory(unresponsiveServer.configuration()),
+                 SimpleConnectionPool.Configuration.DEFAULT)) {
+            MonitoredRabbitMQConsumers unresponsive = 
MonitoredRabbitMQConsumers.of("unresponsive consumers", unresponsivePool,
+                () -> List.of(FIRST_QUEUE), connection -> Mono.empty());
+
+            Disposable check = new 
RabbitMQConsumersHealthCheck(ImmutableSet.of(unresponsive, 
second)).check().subscribe();
+
+            try {
+                awaitAtMostOneMinute.untilAsserted(() -> {
+                    
assertThat(unresponsiveServer.acceptedConnections()).isPositive();
+                    assertThat(secondRestarts).hasValue(1);
+                });
+            } finally {
+                check.dispose();
+            }
+        }
+    }
+
+    @Test
+    void checkShouldNotStartARestartOnANonBlockingThread() {
+        List<Boolean> restartsStartedOnNonBlockingThread = new 
CopyOnWriteArrayList<>();
+        AtomicInteger listedQueues = new AtomicInteger();
+        CompletableFuture<Void> bothDetectionsDone = new CompletableFuture<>();
+        Supplier<List<String>> firstQueue = () -> listQueue(FIRST_QUEUE, 
listedQueues, bothDetectionsDone);
+        Supplier<List<String>> secondQueue = () -> listQueue(SECOND_QUEUE, 
listedQueues, bothDetectionsDone);
+        // Completes on a non-blocking thread once both detections are done: 
the second restart then starts from there
+        Function<Connection, Publisher<Void>> asynchronousRestart = connection 
-> Mono.fromRunnable(
+                () -> 
restartsStartedOnNonBlockingThread.add(Schedulers.isInNonBlockingThread()))
+            .then(Mono.fromFuture(bothDetectionsDone))
+            .then(Mono.delay(Duration.ofMillis(200)))
+            .then();
+        MonitoredRabbitMQConsumers asynchronousFirst = 
MonitoredRabbitMQConsumers.of("asynchronous first consumers",
+            rabbitMQExtension.getConnectionPool(), firstQueue, 
asynchronousRestart);
+        MonitoredRabbitMQConsumers asynchronousSecond = 
MonitoredRabbitMQConsumers.of("asynchronous second consumers",
+            secondVhostConnectionPool, secondQueue, asynchronousRestart);
+
+        new RabbitMQConsumersHealthCheck(ImmutableSet.of(asynchronousFirst, 
asynchronousSecond)).check().block();
+
+        assertThat(restartsStartedOnNonBlockingThread).containsExactly(false, 
false);
+    }
+
+    @Test
+    void checkShouldNotRestartConsumersConcurrently() {
+        AtomicInteger runningRestarts = new AtomicInteger();
+        AtomicInteger maxRunningRestarts = new AtomicInteger();
+        MonitoredRabbitMQConsumers slowFirst = 
MonitoredRabbitMQConsumers.of("slow first consumers", 
rabbitMQExtension.getConnectionPool(),
+            () -> List.of(FIRST_QUEUE), connection -> 
slowRestart(runningRestarts, maxRunningRestarts));
+        MonitoredRabbitMQConsumers slowSecond = 
MonitoredRabbitMQConsumers.of("slow second consumers", 
secondVhostConnectionPool,
+            () -> List.of(SECOND_QUEUE), connection -> 
slowRestart(runningRestarts, maxRunningRestarts));
+
+        Result result = new 
RabbitMQConsumersHealthCheck(ImmutableSet.of(slowFirst, 
slowSecond)).check().block();
+
+        assertThat(result.getCause()).hasValueSatisfying(cause -> 
assertThat(cause)
+            .contains("slow first consumers")
+            .contains("slow second consumers"));
+        assertThat(maxRunningRestarts).hasValue(1);
+    }
+
+    @Test
+    void checkShouldReturnUnhealthyWhenAQueueDoesNotExist() throws Exception {
+        consume(firstChannel, FIRST_QUEUE);
+        firstQueues.add("missing-queue");
+
+        Result result = new 
RabbitMQConsumersHealthCheck(ImmutableSet.of(first)).check().block();
+
+        assertThat(result.isUnHealthy()).isTrue();
+    }
+
+    @Test
+    void checkShouldMonitorQueuesAddedAfterCreation() throws Exception {
+        consume(firstChannel, FIRST_QUEUE);
+        consume(secondChannel, SECOND_QUEUE);
+        RabbitMQConsumersHealthCheck testee = new 
RabbitMQConsumersHealthCheck(ImmutableSet.of(first, second));
+        String thirdQueue = "third-queue";
+        declareQueue(firstChannel, thirdQueue);
+        firstQueues.add(thirdQueue);
+
+        Result result = testee.check().block();
+
+        assertThat(result.isDegraded()).isTrue();
+        assertThat(result.getCause()).contains("No consumers on " + thirdQueue 
+ " of first consumers");
+    }
+
+    private List<String> listQueue(String queue, AtomicInteger listedQueues, 
CompletableFuture<Void> bothDetectionsDone) {
+        if (listedQueues.incrementAndGet() == 2) {
+            bothDetectionsDone.complete(null);
+        }
+        return List.of(queue);
+    }
+
+    private Mono<Void> slowRestart(AtomicInteger runningRestarts, 
AtomicInteger maxRunningRestarts) {
+        return Mono.fromRunnable(() -> 
maxRunningRestarts.accumulateAndGet(runningRestarts.incrementAndGet(), 
Math::max))
+            .then(Mono.delay(Duration.ofSeconds(1)))
+            .then(Mono.fromRunnable(runningRestarts::decrementAndGet));
+    }
+
+    private void declareQueue(Channel channel, String queue) throws Exception {
+        channel.queueDeclare(queue, DURABLE, !EXCLUSIVE, !AUTO_DELETE, 
ImmutableMap.of());
+    }
+
+    private void consume(Channel channel, String queue) throws Exception {
+        channel.basicConsume(queue, AUTO_ACK, (consumerTag, delivery) -> { }, 
consumerTag -> { });
+    }
+}
diff --git 
a/event-bus/distributed/src/main/java/org/apache/james/events/RabbitEventBusConsumerHealthCheck.java
 
b/event-bus/distributed/src/main/java/org/apache/james/events/RabbitMQEventBusConsumers.java
similarity index 51%
rename from 
event-bus/distributed/src/main/java/org/apache/james/events/RabbitEventBusConsumerHealthCheck.java
rename to 
event-bus/distributed/src/main/java/org/apache/james/events/RabbitMQEventBusConsumers.java
index bed6973870..78e715c20f 100644
--- 
a/event-bus/distributed/src/main/java/org/apache/james/events/RabbitEventBusConsumerHealthCheck.java
+++ 
b/event-bus/distributed/src/main/java/org/apache/james/events/RabbitMQEventBusConsumers.java
@@ -19,31 +19,30 @@
 
 package org.apache.james.events;
 
-import java.util.Optional;
+import java.util.List;
 import java.util.stream.Stream;
 
+import org.apache.james.backends.rabbitmq.MonitoredRabbitMQConsumers;
 import org.apache.james.backends.rabbitmq.SimpleConnectionPool;
-import org.apache.james.core.healthcheck.ComponentName;
-import org.apache.james.core.healthcheck.HealthCheck;
-import org.apache.james.core.healthcheck.Result;
-import org.apache.james.util.ReactorUtils;
+import org.reactivestreams.Publisher;
 
-import com.github.fge.lambdas.Throwing;
-import com.rabbitmq.client.Channel;
+import com.google.common.collect.ImmutableList;
+import com.rabbitmq.client.Connection;
 
 import reactor.core.publisher.Mono;
 
-public class RabbitEventBusConsumerHealthCheck implements HealthCheck {
-    public static final String COMPONENT = "EventbusConsumers";
-
+/**
+ * The group consumers of a RabbitMQ event bus, restarted with the event bus.
+ */
+public class RabbitMQEventBusConsumers implements MonitoredRabbitMQConsumers {
     private final EventBus eventBus;
     private final NamingStrategy namingStrategy;
     private final SimpleConnectionPool connectionPool;
     private final Group groupRegistrationHandlerGroup;
 
-    public RabbitEventBusConsumerHealthCheck(EventBus eventBus, NamingStrategy 
namingStrategy,
-                                             SimpleConnectionPool 
connectionPool,
-                                             Group 
groupRegistrationHandlerGroup) {
+    public RabbitMQEventBusConsumers(EventBus eventBus, NamingStrategy 
namingStrategy,
+                                     SimpleConnectionPool connectionPool,
+                                     Group groupRegistrationHandlerGroup) {
         this.eventBus = eventBus;
         this.namingStrategy = namingStrategy;
         this.connectionPool = connectionPool;
@@ -51,36 +50,27 @@ public class RabbitEventBusConsumerHealthCheck implements 
HealthCheck {
     }
 
     @Override
-    public ComponentName componentName() {
-        return new ComponentName(COMPONENT + "-" + 
namingStrategy.getEventBusName().value());
+    public String name() {
+        return namingStrategy.getEventBusName().value() + " event bus";
     }
 
     @Override
-    public Mono<Result> check() {
-        return connectionPool.getResilientConnection()
-            .map(Throwing.function(connection -> {
-                try (Channel channel = connection.createChannel()) {
-                    return check(channel);
-                }
-            })).subscribeOn(ReactorUtils.BLOCKING_CALL_WRAPPER);
+    public SimpleConnectionPool connectionPool() {
+        return connectionPool;
     }
 
-    private Result check(Channel channel) {
-        Stream<Group> groups = Stream.concat(
-            eventBus.listRegisteredGroups().stream(),
-            Stream.of(groupRegistrationHandlerGroup));
-
-        Optional<String> queueWithoutConsumers = groups
+    @Override
+    public List<String> queues() {
+        return Stream.concat(
+                eventBus.listRegisteredGroups().stream(),
+                Stream.of(groupRegistrationHandlerGroup))
             .map(namingStrategy::workQueue)
             .map(GroupRegistration.WorkQueueName::asString)
-            .filter(Throwing.predicate(queue -> channel.consumerCount(queue) 
== 0))
-            .findAny();
+            .collect(ImmutableList.toImmutableList());
+    }
 
-        if (queueWithoutConsumers.isPresent()) {
-            eventBus.restart();
-            return Result.degraded(componentName(), "No consumers on " + 
queueWithoutConsumers.get());
-        } else {
-            return Result.healthy(componentName());
-        }
+    @Override
+    public Publisher<Void> restart(Connection connection) {
+        return Mono.fromRunnable(eventBus::restart);
     }
 }
diff --git 
a/event-bus/distributed/src/test/java/org/apache/james/events/RabbitMQEventBusConsumersTest.java
 
b/event-bus/distributed/src/test/java/org/apache/james/events/RabbitMQEventBusConsumersTest.java
new file mode 100644
index 0000000000..5f525fffc5
--- /dev/null
+++ 
b/event-bus/distributed/src/test/java/org/apache/james/events/RabbitMQEventBusConsumersTest.java
@@ -0,0 +1,80 @@
+/****************************************************************
+ * 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.james.events;
+
+import static 
org.apache.james.backends.rabbitmq.RabbitMQFixture.awaitAtMostOneMinute;
+import static org.apache.james.events.EventBusTestFixture.GROUP_A;
+import static 
org.apache.james.events.EventBusTestFixture.RETRY_BACKOFF_CONFIGURATION;
+import static org.apache.james.events.EventBusTestFixture.newListener;
+import static org.assertj.core.api.Assertions.assertThat;
+
+import org.apache.james.backends.rabbitmq.RabbitMQConsumersHealthCheck;
+import org.apache.james.backends.rabbitmq.RabbitMQExtension;
+import org.apache.james.events.EventBusTestFixture.TestEventSerializer;
+import org.apache.james.events.EventBusTestFixture.TestRegistrationKeyFactory;
+import org.apache.james.metrics.tests.RecordingMetricFactory;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.extension.RegisterExtension;
+
+import com.google.common.collect.ImmutableSet;
+
+class RabbitMQEventBusConsumersTest {
+    private static final NamingStrategy NAMING_STRATEGY = new 
DefaultNamingStrategy(new EventBusName("consumersTest"));
+
+    @RegisterExtension
+    static RabbitMQExtension rabbitMQExtension = 
RabbitMQExtension.singletonRabbitMQ()
+        .isolationPolicy(RabbitMQExtension.IsolationPolicy.STRONG);
+
+    private RabbitMQEventBus eventBus;
+    private RabbitMQEventBusConsumers testee;
+
+    @BeforeEach
+    void setUp() throws Exception {
+        eventBus = new RabbitMQEventBus(NAMING_STRATEGY, 
rabbitMQExtension.getSender(), rabbitMQExtension.getReceiverProvider(),
+            new TestEventSerializer(), RoutingKeyConverter.forFactories(new 
TestRegistrationKeyFactory()), new MemoryEventDeadLetters(),
+            new RecordingMetricFactory(), 
rabbitMQExtension.getRabbitChannelPool(), EventBusId.random(),
+            new 
RabbitMQEventBus.Configurations(rabbitMQExtension.getRabbitMQ().getConfiguration(),
 RETRY_BACKOFF_CONFIGURATION));
+        eventBus.start();
+        eventBus.register(newListener(), GROUP_A);
+
+        testee = new RabbitMQEventBusConsumers(eventBus, NAMING_STRATEGY, 
rabbitMQExtension.getConnectionPool(), GroupRegistrationHandler.GROUP);
+    }
+
+    @AfterEach
+    void tearDown() {
+        eventBus.stop();
+    }
+
+    @Test
+    void 
queuesShouldBeTheWorkQueuesOfTheRegisteredGroupsAndOfTheGroupRegistrationHandler()
 {
+        assertThat(testee.queues()).containsExactlyInAnyOrder(
+            NAMING_STRATEGY.workQueue(GROUP_A).asString(),
+            
NAMING_STRATEGY.workQueue(GroupRegistrationHandler.GROUP).asString());
+    }
+
+    @Test
+    void consumersOfAStartedEventBusShouldBeHealthy() {
+        RabbitMQConsumersHealthCheck healthCheck = new 
RabbitMQConsumersHealthCheck(ImmutableSet.of(testee));
+
+        awaitAtMostOneMinute.untilAsserted(() -> 
assertThat(healthCheck.check().block().isHealthy()).isTrue());
+    }
+}
diff --git 
a/server/apps/postgres-app/src/main/java/org/apache/james/modules/task/DistributedTaskManagerModule.java
 
b/server/apps/postgres-app/src/main/java/org/apache/james/modules/task/DistributedTaskManagerModule.java
index 694158b409..67f6b25a5e 100644
--- 
a/server/apps/postgres-app/src/main/java/org/apache/james/modules/task/DistributedTaskManagerModule.java
+++ 
b/server/apps/postgres-app/src/main/java/org/apache/james/modules/task/DistributedTaskManagerModule.java
@@ -27,8 +27,8 @@ import jakarta.inject.Singleton;
 
 import org.apache.commons.configuration2.Configuration;
 import org.apache.commons.configuration2.ex.ConfigurationException;
+import org.apache.james.backends.rabbitmq.MonitoredRabbitMQConsumers;
 import org.apache.james.backends.rabbitmq.SimpleConnectionPool;
-import org.apache.james.core.healthcheck.HealthCheck;
 import org.apache.james.modules.server.HostnameModule;
 import org.apache.james.modules.server.TaskSerializationModule;
 import org.apache.james.task.TaskManager;
@@ -36,7 +36,6 @@ import 
org.apache.james.task.eventsourcing.EventSourcingTaskManager;
 import org.apache.james.task.eventsourcing.TerminationSubscriber;
 import org.apache.james.task.eventsourcing.WorkQueueSupplier;
 import org.apache.james.task.eventsourcing.distributed.CancelRequestQueueName;
-import 
org.apache.james.task.eventsourcing.distributed.DistributedTaskManagerHealthCheck;
 import 
org.apache.james.task.eventsourcing.distributed.RabbitMQTerminationSubscriber;
 import org.apache.james.task.eventsourcing.distributed.RabbitMQWorkQueue;
 import 
org.apache.james.task.eventsourcing.distributed.RabbitMQWorkQueueConfiguration;
@@ -49,12 +48,15 @@ import org.apache.james.utils.InitializationOperation;
 import org.apache.james.utils.InitilizationOperationBuilder;
 import org.apache.james.utils.PropertiesProvider;
 
+import com.google.common.collect.ImmutableList;
 import com.google.inject.AbstractModule;
 import com.google.inject.Provides;
 import com.google.inject.Scopes;
 import com.google.inject.multibindings.Multibinder;
 import com.google.inject.multibindings.ProvidesIntoSet;
 
+import reactor.core.publisher.Mono;
+
 public class DistributedTaskManagerModule extends AbstractModule {
 
     @Override
@@ -75,10 +77,12 @@ public class DistributedTaskManagerModule extends 
AbstractModule {
         Multibinder<SimpleConnectionPool.ReconnectionHandler> 
reconnectionHandlerMultibinder = Multibinder.newSetBinder(binder(), 
SimpleConnectionPool.ReconnectionHandler.class);
         
reconnectionHandlerMultibinder.addBinding().to(RabbitMQWorkQueueReconnectionHandler.class);
         
reconnectionHandlerMultibinder.addBinding().to(TerminationReconnectionHandler.class);
+    }
 
-        Multibinder.newSetBinder(binder(), HealthCheck.class)
-            .addBinding()
-            .to(DistributedTaskManagerHealthCheck.class);
+    @ProvidesIntoSet
+    MonitoredRabbitMQConsumers taskManagerConsumers(EventSourcingTaskManager 
taskManager, SimpleConnectionPool connectionPool) {
+        return MonitoredRabbitMQConsumers.of("task manager", connectionPool, 
() -> ImmutableList.of(RabbitMQWorkQueue.QUEUE_NAME),
+            connection -> Mono.fromRunnable(taskManager::restart));
     }
 
     @Provides
diff --git 
a/server/container/guice/distributed/src/main/java/org/apache/james/modules/DistributedTaskManagerModule.java
 
b/server/container/guice/distributed/src/main/java/org/apache/james/modules/DistributedTaskManagerModule.java
index 20ea4613fc..c6b1552979 100644
--- 
a/server/container/guice/distributed/src/main/java/org/apache/james/modules/DistributedTaskManagerModule.java
+++ 
b/server/container/guice/distributed/src/main/java/org/apache/james/modules/DistributedTaskManagerModule.java
@@ -28,8 +28,8 @@ import jakarta.inject.Singleton;
 import org.apache.commons.configuration2.Configuration;
 import org.apache.commons.configuration2.ex.ConfigurationException;
 import org.apache.james.backends.cassandra.components.CassandraDataDefinition;
+import org.apache.james.backends.rabbitmq.MonitoredRabbitMQConsumers;
 import org.apache.james.backends.rabbitmq.SimpleConnectionPool;
-import org.apache.james.core.healthcheck.HealthCheck;
 import org.apache.james.modules.server.HostnameModule;
 import org.apache.james.modules.server.TaskSerializationModule;
 import org.apache.james.task.TaskManager;
@@ -40,7 +40,6 @@ import org.apache.james.task.eventsourcing.WorkQueueSupplier;
 import 
org.apache.james.task.eventsourcing.cassandra.CassandraTaskExecutionDetailsProjection;
 import 
org.apache.james.task.eventsourcing.cassandra.CassandraTaskExecutionDetailsProjectionModule;
 import org.apache.james.task.eventsourcing.distributed.CancelRequestQueueName;
-import 
org.apache.james.task.eventsourcing.distributed.DistributedTaskManagerHealthCheck;
 import 
org.apache.james.task.eventsourcing.distributed.RabbitMQTerminationSubscriber;
 import org.apache.james.task.eventsourcing.distributed.RabbitMQWorkQueue;
 import 
org.apache.james.task.eventsourcing.distributed.RabbitMQWorkQueueConfiguration;
@@ -53,12 +52,15 @@ import org.apache.james.utils.InitializationOperation;
 import org.apache.james.utils.InitilizationOperationBuilder;
 import org.apache.james.utils.PropertiesProvider;
 
+import com.google.common.collect.ImmutableList;
 import com.google.inject.AbstractModule;
 import com.google.inject.Provides;
 import com.google.inject.Scopes;
 import com.google.inject.multibindings.Multibinder;
 import com.google.inject.multibindings.ProvidesIntoSet;
 
+import reactor.core.publisher.Mono;
+
 public class DistributedTaskManagerModule extends AbstractModule {
 
     @Override
@@ -83,10 +85,12 @@ public class DistributedTaskManagerModule extends 
AbstractModule {
         Multibinder<SimpleConnectionPool.ReconnectionHandler> 
reconnectionHandlerMultibinder = Multibinder.newSetBinder(binder(), 
SimpleConnectionPool.ReconnectionHandler.class);
         
reconnectionHandlerMultibinder.addBinding().to(RabbitMQWorkQueueReconnectionHandler.class);
         
reconnectionHandlerMultibinder.addBinding().to(TerminationReconnectionHandler.class);
+    }
 
-        Multibinder.newSetBinder(binder(), HealthCheck.class)
-            .addBinding()
-            .to(DistributedTaskManagerHealthCheck.class);
+    @ProvidesIntoSet
+    MonitoredRabbitMQConsumers taskManagerConsumers(EventSourcingTaskManager 
taskManager, SimpleConnectionPool connectionPool) {
+        return MonitoredRabbitMQConsumers.of("task manager", connectionPool, 
() -> ImmutableList.of(RabbitMQWorkQueue.QUEUE_NAME),
+            connection -> Mono.fromRunnable(taskManager::restart));
     }
 
     @Provides
diff --git 
a/server/container/guice/distributed/src/main/java/org/apache/james/modules/event/ContentDeletionEventBusModule.java
 
b/server/container/guice/distributed/src/main/java/org/apache/james/modules/event/ContentDeletionEventBusModule.java
index b1940eac7b..d6c30e644d 100644
--- 
a/server/container/guice/distributed/src/main/java/org/apache/james/modules/event/ContentDeletionEventBusModule.java
+++ 
b/server/container/guice/distributed/src/main/java/org/apache/james/modules/event/ContentDeletionEventBusModule.java
@@ -26,9 +26,9 @@ import java.util.Set;
 import jakarta.inject.Named;
 
 import org.apache.james.backends.rabbitmq.MonitoredDeadLetterQueue;
+import org.apache.james.backends.rabbitmq.MonitoredRabbitMQConsumers;
 import org.apache.james.backends.rabbitmq.RabbitMQConfiguration;
 import org.apache.james.backends.rabbitmq.SimpleConnectionPool;
-import org.apache.james.core.healthcheck.HealthCheck;
 import org.apache.james.event.json.MailboxEventSerializer;
 import org.apache.james.events.EventBus;
 import org.apache.james.events.EventBusId;
@@ -36,8 +36,8 @@ import org.apache.james.events.EventBusReconnectionHandler;
 import org.apache.james.events.EventListener;
 import org.apache.james.events.GroupRegistrationHandler;
 import org.apache.james.events.KeyReconnectionHandler;
-import org.apache.james.events.RabbitEventBusConsumerHealthCheck;
 import org.apache.james.events.RabbitMQEventBus;
+import org.apache.james.events.RabbitMQEventBusConsumers;
 import org.apache.james.events.RetryBackoffConfiguration;
 import org.apache.james.events.RoutingKeyConverter;
 import org.apache.james.jmap.change.Factory;
@@ -85,10 +85,9 @@ public class ContentDeletionEventBusModule extends 
AbstractModule {
     }
 
     @ProvidesIntoSet
-    HealthCheck healthCheck(@Named(CONTENT_DELETION) RabbitMQEventBus eventBus,
-                            SimpleConnectionPool connectionPool) {
-        return new RabbitEventBusConsumerHealthCheck(eventBus, 
CONTENT_DELETION_NAMING_STRATEGY, connectionPool,
-            GroupRegistrationHandler.GROUP);
+    MonitoredRabbitMQConsumers eventBusConsumers(@Named(CONTENT_DELETION) 
RabbitMQEventBus eventBus,
+                                                 SimpleConnectionPool 
connectionPool) {
+        return new RabbitMQEventBusConsumers(eventBus, 
CONTENT_DELETION_NAMING_STRATEGY, connectionPool, 
GroupRegistrationHandler.GROUP);
     }
 
     @ProvidesIntoSet
diff --git 
a/server/container/guice/distributed/src/main/java/org/apache/james/modules/event/JMAPEventBusModule.java
 
b/server/container/guice/distributed/src/main/java/org/apache/james/modules/event/JMAPEventBusModule.java
index 0dfabc8319..120690b289 100644
--- 
a/server/container/guice/distributed/src/main/java/org/apache/james/modules/event/JMAPEventBusModule.java
+++ 
b/server/container/guice/distributed/src/main/java/org/apache/james/modules/event/JMAPEventBusModule.java
@@ -24,17 +24,17 @@ import static 
org.apache.james.events.NamingStrategy.JMAP_NAMING_STRATEGY;
 import jakarta.inject.Named;
 
 import org.apache.james.backends.rabbitmq.MonitoredDeadLetterQueue;
+import org.apache.james.backends.rabbitmq.MonitoredRabbitMQConsumers;
 import org.apache.james.backends.rabbitmq.RabbitMQConfiguration;
 import org.apache.james.backends.rabbitmq.SimpleConnectionPool;
-import org.apache.james.core.healthcheck.HealthCheck;
 import org.apache.james.events.EventBus;
 import org.apache.james.events.EventBusId;
 import org.apache.james.events.EventBusReconnectionHandler;
 import org.apache.james.events.EventSerializer;
 import org.apache.james.events.GroupRegistrationHandler;
 import org.apache.james.events.KeyReconnectionHandler;
-import org.apache.james.events.RabbitEventBusConsumerHealthCheck;
 import org.apache.james.events.RabbitMQEventBus;
+import org.apache.james.events.RabbitMQEventBusConsumers;
 import org.apache.james.events.RetryBackoffConfiguration;
 import org.apache.james.events.RoutingKeyConverter;
 import org.apache.james.jmap.InjectionKeys;
@@ -83,10 +83,9 @@ public class JMAPEventBusModule extends AbstractModule {
     }
 
     @ProvidesIntoSet
-    HealthCheck healthCheck(@Named(InjectionKeys.JMAP) RabbitMQEventBus 
eventBus,
-                            SimpleConnectionPool connectionPool) {
-        return new RabbitEventBusConsumerHealthCheck(eventBus, 
JMAP_NAMING_STRATEGY, connectionPool,
-            GroupRegistrationHandler.GROUP);
+    MonitoredRabbitMQConsumers eventBusConsumers(@Named(InjectionKeys.JMAP) 
RabbitMQEventBus eventBus,
+                                                 SimpleConnectionPool 
connectionPool) {
+        return new RabbitMQEventBusConsumers(eventBus, JMAP_NAMING_STRATEGY, 
connectionPool, GroupRegistrationHandler.GROUP);
     }
 
     @ProvidesIntoSet
diff --git 
a/server/container/guice/distributed/src/main/java/org/apache/james/modules/event/MailboxEventBusModule.java
 
b/server/container/guice/distributed/src/main/java/org/apache/james/modules/event/MailboxEventBusModule.java
index e5fa0764f2..591d9f5b67 100644
--- 
a/server/container/guice/distributed/src/main/java/org/apache/james/modules/event/MailboxEventBusModule.java
+++ 
b/server/container/guice/distributed/src/main/java/org/apache/james/modules/event/MailboxEventBusModule.java
@@ -22,9 +22,9 @@ package org.apache.james.modules.event;
 import static 
org.apache.james.events.NamingStrategy.MAILBOX_EVENT_NAMING_STRATEGY;
 
 import org.apache.james.backends.rabbitmq.MonitoredDeadLetterQueue;
+import org.apache.james.backends.rabbitmq.MonitoredRabbitMQConsumers;
 import org.apache.james.backends.rabbitmq.RabbitMQConfiguration;
 import org.apache.james.backends.rabbitmq.SimpleConnectionPool;
-import org.apache.james.core.healthcheck.HealthCheck;
 import org.apache.james.event.json.MailboxEventSerializer;
 import org.apache.james.events.EventBus;
 import org.apache.james.events.EventBusId;
@@ -32,8 +32,8 @@ import org.apache.james.events.EventBusReconnectionHandler;
 import org.apache.james.events.GroupRegistrationHandler;
 import org.apache.james.events.KeyReconnectionHandler;
 import org.apache.james.events.NamingStrategy;
-import org.apache.james.events.RabbitEventBusConsumerHealthCheck;
 import org.apache.james.events.RabbitMQEventBus;
+import org.apache.james.events.RabbitMQEventBusConsumers;
 import org.apache.james.events.RegistrationKey;
 import org.apache.james.events.RetryBackoffConfiguration;
 import org.apache.james.events.RoutingKeyConverter;
@@ -66,10 +66,9 @@ public class MailboxEventBusModule extends AbstractModule {
     }
 
     @ProvidesIntoSet
-    HealthCheck healthCheck(RabbitMQEventBus eventBus, NamingStrategy 
namingStrategy,
-                            SimpleConnectionPool connectionPool) {
-        return new RabbitEventBusConsumerHealthCheck(eventBus, namingStrategy, 
connectionPool,
-            GroupRegistrationHandler.GROUP);
+    MonitoredRabbitMQConsumers eventBusConsumers(RabbitMQEventBus eventBus, 
NamingStrategy namingStrategy,
+                                                 SimpleConnectionPool 
connectionPool) {
+        return new RabbitMQEventBusConsumers(eventBus, namingStrategy, 
connectionPool, GroupRegistrationHandler.GROUP);
     }
 
     @ProvidesIntoSet
diff --git 
a/server/container/guice/queue/rabbitmq/src/main/java/org/apache/james/modules/queue/rabbitmq/RabbitMQMailQueueModule.java
 
b/server/container/guice/queue/rabbitmq/src/main/java/org/apache/james/modules/queue/rabbitmq/RabbitMQMailQueueModule.java
index 0033ea812f..cdd50603c5 100644
--- 
a/server/container/guice/queue/rabbitmq/src/main/java/org/apache/james/modules/queue/rabbitmq/RabbitMQMailQueueModule.java
+++ 
b/server/container/guice/queue/rabbitmq/src/main/java/org/apache/james/modules/queue/rabbitmq/RabbitMQMailQueueModule.java
@@ -24,15 +24,15 @@ import jakarta.inject.Named;
 import jakarta.inject.Singleton;
 
 import org.apache.james.backends.rabbitmq.MonitoredDeadLetterQueue;
+import org.apache.james.backends.rabbitmq.MonitoredRabbitMQConsumers;
 import org.apache.james.backends.rabbitmq.RabbitMQConfiguration;
 import org.apache.james.backends.rabbitmq.SimpleConnectionPool;
-import org.apache.james.core.healthcheck.HealthCheck;
 import org.apache.james.queue.api.MailQueue;
 import org.apache.james.queue.api.MailQueueFactory;
 import org.apache.james.queue.api.ManageableMailQueue;
 import org.apache.james.queue.rabbitmq.MailQueueName;
 import org.apache.james.queue.rabbitmq.RabbitMQMailQueue;
-import org.apache.james.queue.rabbitmq.RabbitMQMailQueueConsumerHealthCheck;
+import org.apache.james.queue.rabbitmq.RabbitMQMailQueueConsumers;
 import org.apache.james.queue.rabbitmq.RabbitMQMailQueueFactory;
 import org.apache.james.queue.rabbitmq.view.RabbitMQMailQueueConfiguration;
 
@@ -46,9 +46,11 @@ public class RabbitMQMailQueueModule extends AbstractModule {
     protected void configure() {
         Multibinder<SimpleConnectionPool.ReconnectionHandler> 
reconnectionHandlerMultibinder = Multibinder.newSetBinder(binder(), 
SimpleConnectionPool.ReconnectionHandler.class);
         
reconnectionHandlerMultibinder.addBinding().to(SpoolerReconnectionHandler.class);
+    }
 
-        Multibinder.newSetBinder(binder(), HealthCheck.class).addBinding()
-            .to(RabbitMQMailQueueConsumerHealthCheck.class);
+    @ProvidesIntoSet
+    MonitoredRabbitMQConsumers mailQueueConsumers(RabbitMQMailQueueConsumers 
mailQueueConsumers) {
+        return mailQueueConsumers;
     }
 
     @ProvidesIntoSet
diff --git 
a/server/container/guice/queue/rabbitmq/src/main/java/org/apache/james/modules/queue/rabbitmq/RabbitMQModule.java
 
b/server/container/guice/queue/rabbitmq/src/main/java/org/apache/james/modules/queue/rabbitmq/RabbitMQModule.java
index 1646751d0a..69f8c5632d 100644
--- 
a/server/container/guice/queue/rabbitmq/src/main/java/org/apache/james/modules/queue/rabbitmq/RabbitMQModule.java
+++ 
b/server/container/guice/queue/rabbitmq/src/main/java/org/apache/james/modules/queue/rabbitmq/RabbitMQModule.java
@@ -27,7 +27,9 @@ import jakarta.inject.Singleton;
 
 import org.apache.commons.configuration2.ex.ConfigurationException;
 import org.apache.james.backends.rabbitmq.MonitoredDeadLetterQueue;
+import org.apache.james.backends.rabbitmq.MonitoredRabbitMQConsumers;
 import org.apache.james.backends.rabbitmq.RabbitMQConfiguration;
+import org.apache.james.backends.rabbitmq.RabbitMQConsumersHealthCheck;
 import org.apache.james.backends.rabbitmq.RabbitMQDeadLetterQueuesHealthCheck;
 import org.apache.james.backends.rabbitmq.RabbitMQHealthCheck;
 import org.apache.james.backends.rabbitmq.ReactorRabbitMQChannelPool;
@@ -65,7 +67,9 @@ public class RabbitMQModule extends AbstractModule {
         Multibinder<HealthCheck> healthCheckMultiBinder = 
Multibinder.newSetBinder(binder(), HealthCheck.class);
         healthCheckMultiBinder.addBinding().to(RabbitMQHealthCheck.class);
         
healthCheckMultiBinder.addBinding().to(RabbitMQDeadLetterQueuesHealthCheck.class);
+        
healthCheckMultiBinder.addBinding().to(RabbitMQConsumersHealthCheck.class);
         Multibinder.newSetBinder(binder(), MonitoredDeadLetterQueue.class);
+        Multibinder.newSetBinder(binder(), MonitoredRabbitMQConsumers.class);
 
         Multibinder<SimpleConnectionPool.ReconnectionHandler> 
reconnectionHandlerMultibinder = Multibinder.newSetBinder(binder(), 
SimpleConnectionPool.ReconnectionHandler.class);
     }
diff --git 
a/server/protocols/webadmin-integration-test/distributed-webadmin-integration-test/src/test/java/org/apache/james/webadmin/integration/rabbitmq/RabbitMQWebAdminServerIntegrationImmutableTest.java
 
b/server/protocols/webadmin-integration-test/distributed-webadmin-integration-test/src/test/java/org/apache/james/webadmin/integration/rabbitmq/RabbitMQWebAdminServerIntegrationImmutableTest.java
index 1d03765e61..0cffa69c0b 100644
--- 
a/server/protocols/webadmin-integration-test/distributed-webadmin-integration-test/src/test/java/org/apache/james/webadmin/integration/rabbitmq/RabbitMQWebAdminServerIntegrationImmutableTest.java
+++ 
b/server/protocols/webadmin-integration-test/distributed-webadmin-integration-test/src/test/java/org/apache/james/webadmin/integration/rabbitmq/RabbitMQWebAdminServerIntegrationImmutableTest.java
@@ -43,6 +43,7 @@ import org.apache.james.JamesServerExtension;
 import org.apache.james.SearchConfiguration;
 import 
org.apache.james.backends.cassandra.versions.CassandraSchemaVersionManager;
 import org.apache.james.backends.rabbitmq.MonitoredDeadLetterQueue;
+import org.apache.james.backends.rabbitmq.MonitoredRabbitMQConsumers;
 import org.apache.james.junit.categories.BasicFeature;
 import org.apache.james.modules.AwsS3BlobStoreExtension;
 import org.apache.james.modules.RabbitMQExtension;
@@ -55,6 +56,8 @@ import org.eclipse.jetty.http.HttpStatus;
 import org.junit.jupiter.api.Tag;
 import org.junit.jupiter.api.Test;
 import org.junit.jupiter.api.extension.RegisterExtension;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.ValueSource;
 
 import com.google.inject.multibindings.Multibinder;
 
@@ -62,10 +65,12 @@ import com.google.inject.multibindings.Multibinder;
 class RabbitMQWebAdminServerIntegrationImmutableTest extends 
WebAdminServerIntegrationImmutableTest {
     public static class MonitoredRabbitMQProbe implements GuiceProbe {
         private final Set<MonitoredDeadLetterQueue> deadLetterQueues;
+        private final Set<MonitoredRabbitMQConsumers> consumers;
 
         @Inject
-        public MonitoredRabbitMQProbe(Set<MonitoredDeadLetterQueue> 
deadLetterQueues) {
+        public MonitoredRabbitMQProbe(Set<MonitoredDeadLetterQueue> 
deadLetterQueues, Set<MonitoredRabbitMQConsumers> consumers) {
             this.deadLetterQueues = deadLetterQueues;
+            this.consumers = consumers;
         }
 
         List<String> deadLetterQueues() {
@@ -73,6 +78,12 @@ class RabbitMQWebAdminServerIntegrationImmutableTest extends 
WebAdminServerInteg
                 .map(MonitoredDeadLetterQueue::queue)
                 .toList();
         }
+
+        List<String> consumerNames() {
+            return consumers.stream()
+                .map(MonitoredRabbitMQConsumers::name)
+                .toList();
+        }
     }
 
     @RegisterExtension
@@ -164,9 +175,8 @@ class RabbitMQWebAdminServerIntegrationImmutableTest 
extends WebAdminServerInteg
         assertThat(listComponentNames).containsOnly("Guice application 
lifecycle", "EmptyErrorMailRepository",
             "RabbitMQ backend", "RabbitMQDeadLetterQueues", 
"MailReceptionCheck",
             "Cassandra backend", "EventDeadLettersHealthCheck", 
"MessageFastViewProjection",
-            "RabbitMQMailQueue BrowseStart", "OpenSearch Backend", 
"ObjectStorage", "DistributedTaskManagerConsumers",
-            "EventbusConsumers-jmapEvent", "MailQueueConsumers", 
"EventbusConsumers-mailboxEvent",
-            "IMAPHealthCheck", "EventbusConsumers-contentDeletionEvent");
+            "RabbitMQMailQueue BrowseStart", "OpenSearch Backend", 
"ObjectStorage", "RabbitMQConsumers",
+            "IMAPHealthCheck");
     }
 
     @Test
@@ -177,9 +187,16 @@ class RabbitMQWebAdminServerIntegrationImmutableTest 
extends WebAdminServerInteg
     }
 
     @Test
-    void rabbitMQDeadLetterQueuesShouldBeHealthy() {
+    void everyRabbitMQConsumerShouldBeMonitored(GuiceJamesServer server) {
+        
assertThat(server.getProbe(MonitoredRabbitMQProbe.class).consumerNames())
+            .containsExactlyInAnyOrder("mailboxEvent event bus", "jmapEvent 
event bus", "contentDeletionEvent event bus", "task manager", "mail queues");
+    }
+
+    @ParameterizedTest
+    @ValueSource(strings = {"RabbitMQDeadLetterQueues", "RabbitMQConsumers"})
+    void rabbitMQHealthChecksShouldBeHealthy(String componentName) {
         when()
-            .get(HealthCheckRoutes.HEALTHCHECK + 
"/checks/RabbitMQDeadLetterQueues")
+            .get(HealthCheckRoutes.HEALTHCHECK + "/checks/" + componentName)
         .then()
             .statusCode(HttpStatus.OK_200)
             .body("status", is("healthy"));
diff --git 
a/server/queue/queue-rabbitmq/src/main/java/org/apache/james/queue/rabbitmq/RabbitMQMailQueueConsumerHealthCheck.java
 
b/server/queue/queue-rabbitmq/src/main/java/org/apache/james/queue/rabbitmq/RabbitMQMailQueueConsumers.java
similarity index 51%
rename from 
server/queue/queue-rabbitmq/src/main/java/org/apache/james/queue/rabbitmq/RabbitMQMailQueueConsumerHealthCheck.java
rename to 
server/queue/queue-rabbitmq/src/main/java/org/apache/james/queue/rabbitmq/RabbitMQMailQueueConsumers.java
index 6adf86f5e9..0509e5763d 100644
--- 
a/server/queue/queue-rabbitmq/src/main/java/org/apache/james/queue/rabbitmq/RabbitMQMailQueueConsumerHealthCheck.java
+++ 
b/server/queue/queue-rabbitmq/src/main/java/org/apache/james/queue/rabbitmq/RabbitMQMailQueueConsumers.java
@@ -19,66 +19,60 @@
 
 package org.apache.james.queue.rabbitmq;
 
+import java.util.List;
 import java.util.Set;
 
 import jakarta.inject.Inject;
 
+import org.apache.james.backends.rabbitmq.MonitoredRabbitMQConsumers;
 import org.apache.james.backends.rabbitmq.SimpleConnectionPool;
-import org.apache.james.core.healthcheck.ComponentName;
-import org.apache.james.core.healthcheck.HealthCheck;
-import org.apache.james.core.healthcheck.Result;
+import org.reactivestreams.Publisher;
 
-import com.github.fge.lambdas.Throwing;
-import com.rabbitmq.client.Channel;
+import com.google.common.collect.ImmutableList;
 import com.rabbitmq.client.Connection;
 
 import reactor.core.publisher.Flux;
-import reactor.core.publisher.Mono;
-import reactor.core.scheduler.Schedulers;
-
-public class RabbitMQMailQueueConsumerHealthCheck implements HealthCheck {
-    public static final ComponentName COMPONENT_NAME = new 
ComponentName("RabbitMQMailQueueConsumersHealthCheck");
-    public static final ComponentName COMPONENT = new 
ComponentName("MailQueueConsumers");
 
+/**
+ * The consumers of the mail queues. Restarting them runs every {@link 
SimpleConnectionPool.ReconnectionHandler}.
+ */
+public class RabbitMQMailQueueConsumers implements MonitoredRabbitMQConsumers {
     private final RabbitMQMailQueueFactory queueFactory;
     private final Set<SimpleConnectionPool.ReconnectionHandler> 
reconnectionHandlers;
     private final SimpleConnectionPool connectionPool;
 
     @Inject
-    public RabbitMQMailQueueConsumerHealthCheck(RabbitMQMailQueueFactory 
queueFactory, Set<SimpleConnectionPool.ReconnectionHandler> 
reconnectionHandlers, SimpleConnectionPool connectionPool) {
+    public RabbitMQMailQueueConsumers(RabbitMQMailQueueFactory queueFactory, 
Set<SimpleConnectionPool.ReconnectionHandler> reconnectionHandlers,
+                                      SimpleConnectionPool connectionPool) {
         this.queueFactory = queueFactory;
         this.reconnectionHandlers = reconnectionHandlers;
         this.connectionPool = connectionPool;
     }
 
     @Override
-    public ComponentName componentName() {
-        return COMPONENT_NAME;
+    public String name() {
+        return "mail queues";
     }
 
     @Override
-    public Mono<Result> check() {
-        return connectionPool.getResilientConnection()
-            .flatMap(connection -> Mono.using(connection::createChannel,
-                channel -> check(connection, channel),
-                Throwing.consumer(Channel::close)))
-            .subscribeOn(Schedulers.boundedElastic());
+    public SimpleConnectionPool connectionPool() {
+        return connectionPool;
     }
 
-    private Mono<Result> check(Connection connection, Channel channel) {
-        boolean queueWithoutConsumers = queueFactory.listCreatedMailQueues()
+    @Override
+    public List<String> queues() {
+        return queueFactory.listCreatedMailQueues()
             .stream()
             .map(org.apache.james.queue.api.MailQueueName::asString)
             .map(MailQueueName::fromString)
-            .map(m -> m.toWorkQueueName().asString())
-            .anyMatch(Throwing.predicate(queue -> channel.consumerCount(queue) 
== 0));
+            .map(mailQueueName -> mailQueueName.toWorkQueueName().asString())
+            .collect(ImmutableList.toImmutableList());
+    }
 
-        if (queueWithoutConsumers) {
-            return Flux.fromIterable(reconnectionHandlers)
-                .concatMap(reconnectionHandler -> 
reconnectionHandler.handleReconnection(connection))
-                .then(Mono.just(Result.degraded(COMPONENT, "No consumers")));
-        } else {
-            return Mono.just(Result.healthy(COMPONENT));
-        }
+    @Override
+    public Publisher<Void> restart(Connection connection) {
+        return Flux.fromIterable(reconnectionHandlers)
+            .concatMap(reconnectionHandler -> 
reconnectionHandler.handleReconnection(connection))
+            .then();
     }
 }
diff --git 
a/server/queue/queue-rabbitmq/src/test/java/org/apache/james/queue/rabbitmq/RabbitMQMailQueueConsumerHealthCheckTest.java
 
b/server/queue/queue-rabbitmq/src/test/java/org/apache/james/queue/rabbitmq/RabbitMQMailQueueConsumerHealthCheckTest.java
deleted file mode 100644
index 9106ab08d8..0000000000
--- 
a/server/queue/queue-rabbitmq/src/test/java/org/apache/james/queue/rabbitmq/RabbitMQMailQueueConsumerHealthCheckTest.java
+++ /dev/null
@@ -1,92 +0,0 @@
-/****************************************************************
- * Licensed to the Apache Software Foundation (ASF) under one   *
- * or more contributor license agreements.  See the NOTICE file *
- * distributed with this work for additional information        *
- * regarding copyright ownership.  The ASF licenses this file   *
- * to you under the Apache License, Version 2.0 (the            *
- * "License"); you may not use this file except in compliance   *
- * with the License.  You may obtain a copy of the License at   *
- *                                                              *
- *   http://www.apache.org/licenses/LICENSE-2.0                 *
- *                                                              *
- * Unless required by applicable law or agreed to in writing,   *
- * software distributed under the License is distributed on an  *
- * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY       *
- * KIND, either express or implied.  See the License for the    *
- * specific language governing permissions and limitations      *
- * under the License.                                           *
- ****************************************************************/
-
-package org.apache.james.queue.rabbitmq;
-
-import static org.apache.james.backends.rabbitmq.Constants.AUTO_ACK;
-import static org.apache.james.backends.rabbitmq.Constants.AUTO_DELETE;
-import static org.apache.james.backends.rabbitmq.Constants.DURABLE;
-import static org.apache.james.backends.rabbitmq.Constants.EXCLUSIVE;
-import static org.apache.james.queue.api.MailQueueFactory.SPOOL;
-import static org.assertj.core.api.Assertions.assertThat;
-import static org.mockito.Mockito.mock;
-import static org.mockito.Mockito.when;
-
-import java.util.concurrent.atomic.AtomicInteger;
-
-import org.apache.james.backends.rabbitmq.RabbitMQExtension;
-import org.apache.james.backends.rabbitmq.SimpleConnectionPool;
-import org.apache.james.core.healthcheck.Result;
-import org.junit.jupiter.api.AfterEach;
-import org.junit.jupiter.api.BeforeEach;
-import org.junit.jupiter.api.Test;
-import org.junit.jupiter.api.extension.RegisterExtension;
-
-import com.google.common.collect.ImmutableMap;
-import com.google.common.collect.ImmutableSet;
-import com.rabbitmq.client.Channel;
-
-import reactor.core.publisher.Mono;
-
-class RabbitMQMailQueueConsumerHealthCheckTest {
-    private static final String SPOOL_WORK_QUEUE = 
MailQueueName.fromString(SPOOL.asString()).toWorkQueueName().asString();
-
-    @RegisterExtension
-    static RabbitMQExtension rabbitMQExtension = 
RabbitMQExtension.singletonRabbitMQ()
-        .isolationPolicy(RabbitMQExtension.IsolationPolicy.STRONG);
-
-    private final AtomicInteger reconnections = new AtomicInteger();
-    private Channel channel;
-    private RabbitMQMailQueueConsumerHealthCheck testee;
-
-    @BeforeEach
-    void setUp() throws Exception {
-        channel = 
rabbitMQExtension.getConnectionPool().getResilientConnection().block().createChannel();
-        channel.queueDeclare(SPOOL_WORK_QUEUE, DURABLE, !EXCLUSIVE, 
!AUTO_DELETE, ImmutableMap.of());
-
-        RabbitMQMailQueueFactory queueFactory = 
mock(RabbitMQMailQueueFactory.class);
-        
when(queueFactory.listCreatedMailQueues()).thenReturn(ImmutableSet.of(SPOOL));
-        SimpleConnectionPool.ReconnectionHandler countingHandler = connection 
-> Mono.fromRunnable(reconnections::incrementAndGet);
-
-        testee = new RabbitMQMailQueueConsumerHealthCheck(queueFactory, 
ImmutableSet.of(countingHandler), rabbitMQExtension.getConnectionPool());
-    }
-
-    @AfterEach
-    void tearDown() throws Exception {
-        channel.close();
-    }
-
-    @Test
-    void checkShouldRunReconnectionHandlersWhenAWorkQueueHasNoConsumer() {
-        Result result = testee.check().block();
-
-        assertThat(result.isDegraded()).isTrue();
-        assertThat(reconnections).hasValue(1);
-    }
-
-    @Test
-    void checkShouldNotRunReconnectionHandlersWhenWorkQueuesHaveConsumers() 
throws Exception {
-        channel.basicConsume(SPOOL_WORK_QUEUE, AUTO_ACK, (consumerTag, 
delivery) -> { }, consumerTag -> { });
-
-        Result result = testee.check().block();
-
-        assertThat(result.isHealthy()).isTrue();
-        assertThat(reconnections).hasValue(0);
-    }
-}
diff --git 
a/server/task/task-distributed/src/main/java/org/apache/james/task/eventsourcing/distributed/DistributedTaskManagerHealthCheck.java
 
b/server/task/task-distributed/src/main/java/org/apache/james/task/eventsourcing/distributed/DistributedTaskManagerHealthCheck.java
deleted file mode 100644
index 323f24a124..0000000000
--- 
a/server/task/task-distributed/src/main/java/org/apache/james/task/eventsourcing/distributed/DistributedTaskManagerHealthCheck.java
+++ /dev/null
@@ -1,77 +0,0 @@
-/****************************************************************
- * Licensed to the Apache Software Foundation (ASF) under one   *
- * or more contributor license agreements.  See the NOTICE file *
- * distributed with this work for additional information        *
- * regarding copyright ownership.  The ASF licenses this file   *
- * to you under the Apache License, Version 2.0 (the            *
- * "License"); you may not use this file except in compliance   *
- * with the License.  You may obtain a copy of the License at   *
- *                                                              *
- *   http://www.apache.org/licenses/LICENSE-2.0                 *
- *                                                              *
- * Unless required by applicable law or agreed to in writing,   *
- * software distributed under the License is distributed on an  *
- * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY       *
- * KIND, either express or implied.  See the License for the    *
- * specific language governing permissions and limitations      *
- * under the License.                                           *
- ****************************************************************/
-
-package org.apache.james.task.eventsourcing.distributed;
-
-import static 
org.apache.james.task.eventsourcing.distributed.RabbitMQWorkQueue.QUEUE_NAME;
-
-import java.io.IOException;
-
-import jakarta.inject.Inject;
-
-import org.apache.james.backends.rabbitmq.SimpleConnectionPool;
-import org.apache.james.core.healthcheck.ComponentName;
-import org.apache.james.core.healthcheck.HealthCheck;
-import org.apache.james.core.healthcheck.Result;
-import org.apache.james.task.eventsourcing.EventSourcingTaskManager;
-import org.apache.james.util.ReactorUtils;
-
-import com.github.fge.lambdas.Throwing;
-import com.rabbitmq.client.Channel;
-
-import reactor.core.publisher.Mono;
-
-public class DistributedTaskManagerHealthCheck implements HealthCheck {
-    public static final ComponentName COMPONENT_NAME = new 
ComponentName("DistributedTaskManagerConsumersHealthCheck");
-    public static final ComponentName COMPONENT = new 
ComponentName("DistributedTaskManagerConsumers");
-
-    private final EventSourcingTaskManager taskManager;
-    private final SimpleConnectionPool connectionPool;
-
-    @Inject
-    public DistributedTaskManagerHealthCheck(EventSourcingTaskManager 
taskManager, SimpleConnectionPool connectionPool) {
-        this.taskManager = taskManager;
-        this.connectionPool = connectionPool;
-    }
-
-    @Override
-    public ComponentName componentName() {
-        return COMPONENT_NAME;
-    }
-
-    @Override
-    public Mono<Result> check() {
-        return connectionPool.getResilientConnection()
-            .map(Throwing.function(connection -> {
-                try (Channel channel = connection.createChannel()) {
-                    return check(channel);
-                }
-            })).subscribeOn(ReactorUtils.BLOCKING_CALL_WRAPPER);
-    }
-
-    private Result check(Channel channel) throws IOException {
-        if (channel.consumerCount(QUEUE_NAME) == 0) {
-            taskManager.restart();
-
-            return Result.degraded(COMPONENT, "No consumers");
-        } else {
-            return Result.healthy(COMPONENT);
-        }
-    }
-}
diff --git 
a/server/task/task-distributed/src/main/java/org/apache/james/task/eventsourcing/distributed/RabbitMQWorkQueue.java
 
b/server/task/task-distributed/src/main/java/org/apache/james/task/eventsourcing/distributed/RabbitMQWorkQueue.java
index 193cd5e0b4..e31200b1d3 100644
--- 
a/server/task/task-distributed/src/main/java/org/apache/james/task/eventsourcing/distributed/RabbitMQWorkQueue.java
+++ 
b/server/task/task-distributed/src/main/java/org/apache/james/task/eventsourcing/distributed/RabbitMQWorkQueue.java
@@ -74,7 +74,7 @@ public class RabbitMQWorkQueue implements WorkQueue {
         .orElse(1);
 
     static final String EXCHANGE_NAME = "taskManagerWorkQueueExchange";
-    static final String QUEUE_NAME = "taskManagerWorkQueue";
+    public static final String QUEUE_NAME = "taskManagerWorkQueue";
     static final String ROUTING_KEY = "taskManagerWorkQueueRoutingKey";
 
     static final String CANCEL_REQUESTS_EXCHANGE_NAME = 
"taskManagerCancelRequestsExchange";


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

Reply via email to