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]
