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 b93e14fe2d9d4807933a5fa432b9337cd19ec91c Author: Quan Tran <[email protected]> AuthorDate: Thu Oct 1 13:21:52 2026 +0700 [IMPROVEMENT] One RabbitMQ dead letter queues health check fed by MonitoredDeadLetterQueue Each dead letter queue health check was its own class. Modules now only contribute a MonitoredDeadLetterQueue (RabbitMQ configuration, queue) with @ProvidesIntoSet, and a single RabbitMQDeadLetterQueuesHealthCheck checks them all, concurrently, each on its own RabbitMQ server. Downstream projects can monitor their dead letter queues without a dedicated class. The mailbox, JMAP and content deletion event bus and the mail queue dead letter queue health checks are replaced by contributions. Visible changes: - /healthcheck reports a single RabbitMQDeadLetterQueues component instead of RabbitMQMailboxEventBusDeadLetterQueueHealthCheck, RabbitMQJmapEventBusDeadLetterQueueHealthCheck, RabbitMQContentDeletionEventBusDeadLetterQueueHealthCheck and RabbitMQMailQueueDeadLetterQueueHealthCheck. Its cause names the affected queues. - an error on one queue makes the result unhealthy without hiding the other queues. Each error is logged with the queue it concerns, and all errors are kept, as suppressed exceptions, in the result. Co-Authored-By: Claude Opus 5.5 <[email protected]> --- .../james/backends/rabbitmq/MergedResults.java | 66 +++++++ .../rabbitmq/MonitoredDeadLetterQueue.java | 28 +++ .../RabbitMQDeadLetterQueuesHealthCheck.java | 95 ++++++++++ .../james/backends/rabbitmq/MergedResultsTest.java | 74 ++++++++ .../RabbitMQDeadLetterQueuesHealthCheckTest.java | 205 +++++++++++++++++++++ .../backends/rabbitmq/UnresponsiveServer.java | 76 ++++++++ ...DeletionEventBusDeadLetterQueueHealthCheck.java | 63 ------- ...itMQJmapEventBusDeadLetterQueueHealthCheck.java | 63 ------- ...QMailboxEventBusDeadLetterQueueHealthCheck.java | 65 ------- ...tionEventBusDeadLetterQueueHealthCheckTest.java | 131 ------------- ...JmapEventBusDeadLetterQueueHealthCheckTest.java | 131 ------------- ...lboxEventBusDeadLetterQueueHealthCheckTest.java | 132 ------------- .../event/ContentDeletionEventBusModule.java | 6 +- .../james/modules/event/JMAPEventBusModule.java | 6 +- .../james/modules/event/MailboxEventBusModule.java | 8 +- .../queue/rabbitmq/RabbitMQMailQueueModule.java | 12 +- .../modules/queue/rabbitmq/RabbitMQModule.java | 4 + ...itMQWebAdminServerIntegrationImmutableTest.java | 48 ++++- .../apache/james/queue/rabbitmq/MailQueueName.java | 2 +- ...abbitMQMailQueueDeadLetterQueueHealthCheck.java | 66 ------- ...tMQMailQueueDeadLetterQueueHealthCheckTest.java | 133 ------------- 21 files changed, 612 insertions(+), 802 deletions(-) diff --git a/backends-common/rabbitmq/src/main/java/org/apache/james/backends/rabbitmq/MergedResults.java b/backends-common/rabbitmq/src/main/java/org/apache/james/backends/rabbitmq/MergedResults.java new file mode 100644 index 0000000000..b8a9062740 --- /dev/null +++ b/backends-common/rabbitmq/src/main/java/org/apache/james/backends/rabbitmq/MergedResults.java @@ -0,0 +1,66 @@ +/**************************************************************** + * 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.Optional; +import java.util.stream.Collectors; + +import org.apache.james.core.healthcheck.ComponentName; +import org.apache.james.core.healthcheck.Result; +import org.apache.james.core.healthcheck.ResultStatus; + +final class MergedResults { + private static final String CAUSE_SEPARATOR = "; "; + + static Result merge(ComponentName componentName, List<Result> results) { + ResultStatus status = results.stream() + .map(Result::getStatus) + .reduce(ResultStatus.HEALTHY, ResultStatus::merge); + String cause = results.stream() + .map(Result::getCause) + .flatMap(Optional::stream) + .collect(Collectors.joining(CAUSE_SEPARATOR)); + + return switch (status) { + case HEALTHY -> Result.healthy(componentName); + case DEGRADED -> Result.degraded(componentName, cause); + case UNHEALTHY -> mergedError(results) + .map(error -> Result.unhealthy(componentName, cause, error)) + .orElseGet(() -> Result.unhealthy(componentName, cause)); + }; + } + + private static Optional<Throwable> mergedError(List<Result> results) { + List<Throwable> errors = results.stream() + .map(Result::getError) + .flatMap(Optional::stream) + .toList(); + if (errors.size() <= 1) { + return errors.stream().findFirst(); + } + RuntimeException mergedError = new RuntimeException(errors.size() + " errors, see the suppressed exceptions"); + errors.forEach(mergedError::addSuppressed); + return Optional.of(mergedError); + } + + private MergedResults() { + } +} diff --git a/backends-common/rabbitmq/src/main/java/org/apache/james/backends/rabbitmq/MonitoredDeadLetterQueue.java b/backends-common/rabbitmq/src/main/java/org/apache/james/backends/rabbitmq/MonitoredDeadLetterQueue.java new file mode 100644 index 0000000000..ea2acdf02f --- /dev/null +++ b/backends-common/rabbitmq/src/main/java/org/apache/james/backends/rabbitmq/MonitoredDeadLetterQueue.java @@ -0,0 +1,28 @@ +/**************************************************************** + * 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; + +/** + * A dead letter queue watched by {@link RabbitMQDeadLetterQueuesHealthCheck}. + * + * @param configuration of the RabbitMQ server hosting the queue + */ +public record MonitoredDeadLetterQueue(RabbitMQConfiguration configuration, String queue) { +} diff --git a/backends-common/rabbitmq/src/main/java/org/apache/james/backends/rabbitmq/RabbitMQDeadLetterQueuesHealthCheck.java b/backends-common/rabbitmq/src/main/java/org/apache/james/backends/rabbitmq/RabbitMQDeadLetterQueuesHealthCheck.java new file mode 100644 index 0000000000..5b9b3ccf1e --- /dev/null +++ b/backends-common/rabbitmq/src/main/java/org/apache/james/backends/rabbitmq/RabbitMQDeadLetterQueuesHealthCheck.java @@ -0,0 +1,95 @@ +/**************************************************************** + * 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.util.Map; +import java.util.Set; +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.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import com.google.common.collect.ImmutableList; +import com.google.common.collect.ImmutableMap; + +import reactor.core.publisher.Flux; +import reactor.core.publisher.Mono; +import reactor.core.scheduler.Schedulers; + +public class RabbitMQDeadLetterQueuesHealthCheck implements HealthCheck { + public static final ComponentName COMPONENT_NAME = new ComponentName("RabbitMQDeadLetterQueues"); + private static final Logger LOGGER = LoggerFactory.getLogger(RabbitMQDeadLetterQueuesHealthCheck.class); + private static final String DEFAULT_VHOST = "/"; + + private final ImmutableList<MonitoredDeadLetterQueue> deadLetterQueues; + private final Map<RabbitMQConfiguration, RabbitMQManagementAPI> managementAPIs; + + @Inject + public RabbitMQDeadLetterQueuesHealthCheck(Set<MonitoredDeadLetterQueue> deadLetterQueues) { + this.deadLetterQueues = ImmutableList.copyOf(deadLetterQueues); + this.managementAPIs = deadLetterQueues.stream() + .map(MonitoredDeadLetterQueue::configuration) + .distinct() + .collect(ImmutableMap.toImmutableMap(Function.identity(), RabbitMQManagementAPI::from)); + } + + @Override + public ComponentName componentName() { + return COMPONENT_NAME; + } + + @Override + public Mono<Result> check() { + return Flux.fromIterable(deadLetterQueues) + .flatMapSequential(this::check, DEFAULT_CONCURRENCY) + .collectList() + .map(results -> MergedResults.merge(COMPONENT_NAME, results)); + } + + private Mono<Result> check(MonitoredDeadLetterQueue deadLetterQueue) { + return Mono.fromCallable(() -> queueLength(deadLetterQueue)) + .map(queueLength -> { + if (queueLength != 0) { + return Result.degraded(COMPONENT_NAME, String.format("RabbitMQ dead letter queue %s contains %d messages. This might indicate transient failure on processing.", + deadLetterQueue.queue(), queueLength)); + } + return Result.healthy(COMPONENT_NAME); + }) + .onErrorResume(e -> { + LOGGER.warn("Error checking RabbitMQ dead letter queue {}", deadLetterQueue.queue(), e); + return Mono.just(Result.unhealthy(COMPONENT_NAME, "Error checking RabbitMQ dead letter queue " + deadLetterQueue.queue(), e)); + }) + .subscribeOn(Schedulers.boundedElastic()); // Reading the management API is blocking: each queue on its own thread + } + + private long queueLength(MonitoredDeadLetterQueue deadLetterQueue) { + RabbitMQConfiguration configuration = deadLetterQueue.configuration(); + return managementAPIs.get(configuration) + .queueDetails(configuration.getVhost().orElse(DEFAULT_VHOST), deadLetterQueue.queue()) + .getQueueLength(); + } +} diff --git a/backends-common/rabbitmq/src/test/java/org/apache/james/backends/rabbitmq/MergedResultsTest.java b/backends-common/rabbitmq/src/test/java/org/apache/james/backends/rabbitmq/MergedResultsTest.java new file mode 100644 index 0000000000..1e12f621c9 --- /dev/null +++ b/backends-common/rabbitmq/src/test/java/org/apache/james/backends/rabbitmq/MergedResultsTest.java @@ -0,0 +1,74 @@ +/**************************************************************** + * 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.assertj.core.api.Assertions.assertThat; + +import java.util.List; + +import org.apache.james.core.healthcheck.ComponentName; +import org.apache.james.core.healthcheck.Result; +import org.junit.jupiter.api.Test; + +class MergedResultsTest { + private static final ComponentName COMPONENT_NAME = new ComponentName("merged"); + private static final ComponentName ITEM = new ComponentName("item"); + + @Test + void mergeShouldBeHealthyWhenNoResult() { + assertThat(MergedResults.merge(COMPONENT_NAME, List.of()).isHealthy()).isTrue(); + } + + @Test + void mergeShouldKeepTheWorstStatusAndEveryCause() { + Result merged = MergedResults.merge(COMPONENT_NAME, List.of( + Result.healthy(ITEM), + Result.degraded(ITEM, "first cause"), + Result.unhealthy(ITEM, "second cause"))); + + assertThat(merged.isUnHealthy()).isTrue(); + assertThat(merged.getComponentName()).isEqualTo(COMPONENT_NAME); + assertThat(merged.getCause()).contains("first cause; second cause"); + } + + @Test + void mergeShouldKeepASingleErrorAsIs() { + RuntimeException error = new RuntimeException("error"); + + Result merged = MergedResults.merge(COMPONENT_NAME, List.of( + Result.unhealthy(ITEM, "cause", error), + Result.unhealthy(ITEM, "other cause"))); + + assertThat(merged.getError()).containsSame(error); + } + + @Test + void mergeShouldKeepEveryErrorWhenSeveral() { + RuntimeException firstError = new RuntimeException("first error"); + RuntimeException secondError = new RuntimeException("second error"); + + Result merged = MergedResults.merge(COMPONENT_NAME, List.of( + Result.unhealthy(ITEM, "first cause", firstError), + Result.unhealthy(ITEM, "second cause", secondError))); + + assertThat(merged.getError()).hasValueSatisfying(error -> assertThat(error.getSuppressed()) + .containsExactly(firstError, secondError)); + } +} diff --git a/backends-common/rabbitmq/src/test/java/org/apache/james/backends/rabbitmq/RabbitMQDeadLetterQueuesHealthCheckTest.java b/backends-common/rabbitmq/src/test/java/org/apache/james/backends/rabbitmq/RabbitMQDeadLetterQueuesHealthCheckTest.java new file mode 100644 index 0000000000..1f2b79c200 --- /dev/null +++ b/backends-common/rabbitmq/src/test/java/org/apache/james/backends/rabbitmq/RabbitMQDeadLetterQueuesHealthCheckTest.java @@ -0,0 +1,205 @@ +/**************************************************************** + * 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 java.nio.charset.StandardCharsets.UTF_8; +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.Optional; + +import org.apache.james.core.healthcheck.Result; +import org.awaitility.Awaitility; +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.Disposable; +import reactor.core.publisher.Mono; +import reactor.rabbitmq.OutboundMessage; +import reactor.rabbitmq.QueueSpecification; + +class RabbitMQDeadLetterQueuesHealthCheckTest { + private static final String DEFAULT_EXCHANGE = ""; + private static final String SECOND_VHOST = "second"; + private static final String FIRST_DEAD_LETTER_QUEUE = "first-dead-letter"; + private static final String SECOND_DEAD_LETTER_QUEUE = "second-dead-letter"; + + @RegisterExtension + static RabbitMQExtension rabbitMQExtension = RabbitMQExtension.singletonRabbitMQ() + .isolationPolicy(RabbitMQExtension.IsolationPolicy.STRONG); + + private RabbitMQConfiguration configuration; + private RabbitMQConfiguration secondVhostConfiguration; + private SimpleConnectionPool secondVhostConnectionPool; + + @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(), ".*", ".*", ".*"); + configuration = rabbitMQ.getConfiguration(); + secondVhostConfiguration = RabbitMQConfiguration.builder() + .amqpUri(URI.create(rabbitMQ.amqpUri() + "/" + SECOND_VHOST)) + .managementUri(rabbitMQ.managementUri()) + .managementCredentials(DEFAULT_MANAGEMENT_CREDENTIAL) + .vhost(Optional.of(SECOND_VHOST)) + .build(); + secondVhostConnectionPool = new SimpleConnectionPool(new RabbitMQConnectionFactory(secondVhostConfiguration), + SimpleConnectionPool.Configuration.DEFAULT); + } + + @AfterEach + void tearDown() { + secondVhostConnectionPool.close(); + } + + @Test + void checkShouldReturnHealthyWhenNoDeadLetterQueueIsMonitored() { + Result result = new RabbitMQDeadLetterQueuesHealthCheck(ImmutableSet.of()).check().block(); + + assertThat(result.isHealthy()).isTrue(); + } + + @Test + void checkShouldReturnHealthyWhenAllDeadLetterQueuesAreEmpty() throws Exception { + declareQueues(); + + Result result = twoQueuesHealthCheck().check().block(); + + assertThat(result.isHealthy()).isTrue(); + assertThat(result.getComponentName()).isEqualTo(RabbitMQDeadLetterQueuesHealthCheck.COMPONENT_NAME); + } + + @Test + void checkShouldReturnDegradedNamingOnlyTheFirstDeadLetterQueueWhenItIsTheOnlyNonEmptyOne() throws Exception { + declareQueues(); + publish(FIRST_DEAD_LETTER_QUEUE); + + assertDegradedNamingOnly(FIRST_DEAD_LETTER_QUEUE, SECOND_DEAD_LETTER_QUEUE); + } + + @Test + void checkShouldReturnDegradedNamingOnlyTheSecondDeadLetterQueueWhenItIsTheOnlyNonEmptyOne() throws Exception { + declareQueues(); + publishInSecondVhost(SECOND_DEAD_LETTER_QUEUE); + + assertDegradedNamingOnly(SECOND_DEAD_LETTER_QUEUE, FIRST_DEAD_LETTER_QUEUE); + } + + @Test + void checkShouldQueryTheDeadLetterQueuesConcurrently() throws Exception { + try (UnresponsiveServer firstServer = new UnresponsiveServer(); + UnresponsiveServer secondServer = new UnresponsiveServer()) { + Disposable check = new RabbitMQDeadLetterQueuesHealthCheck(ImmutableSet.of( + new MonitoredDeadLetterQueue(firstServer.configuration(), FIRST_DEAD_LETTER_QUEUE), + new MonitoredDeadLetterQueue(secondServer.configuration(), SECOND_DEAD_LETTER_QUEUE))) + .check() + .subscribe(); + + try { + // Well below the 60 seconds read timeout of the first, unanswered, request + Awaitility.await().atMost(Duration.ofSeconds(30)) + .untilAsserted(() -> { + assertThat(firstServer.acceptedConnections()).isPositive(); + assertThat(secondServer.acceptedConnections()).isPositive(); + }); + } finally { + check.dispose(); + } + } + } + + @Test + void checkShouldReturnUnhealthyAndStillReportOtherQueuesWhenADeadLetterQueueDoesNotExist() throws Exception { + declareQueueInSecondVhost(SECOND_DEAD_LETTER_QUEUE); + publishInSecondVhost(SECOND_DEAD_LETTER_QUEUE); + + awaitAtMostOneMinute.untilAsserted(() -> { + Result result = twoQueuesHealthCheck().check().block(); + + assertThat(result.isUnHealthy()).isTrue(); + assertThat(result.getCause()).hasValueSatisfying(cause -> assertThat(cause) + .contains("Error checking RabbitMQ dead letter queue " + FIRST_DEAD_LETTER_QUEUE) + .contains("RabbitMQ dead letter queue " + SECOND_DEAD_LETTER_QUEUE)); + }); + } + + private RabbitMQDeadLetterQueuesHealthCheck twoQueuesHealthCheck() { + return new RabbitMQDeadLetterQueuesHealthCheck(ImmutableSet.of( + new MonitoredDeadLetterQueue(configuration, FIRST_DEAD_LETTER_QUEUE), + new MonitoredDeadLetterQueue(secondVhostConfiguration, SECOND_DEAD_LETTER_QUEUE))); + } + + private void assertDegradedNamingOnly(String nonEmptyQueue, String emptyQueue) { + awaitAtMostOneMinute.untilAsserted(() -> { + Result result = twoQueuesHealthCheck().check().block(); + + assertThat(result.isDegraded()).isTrue(); + assertThat(result.getCause()).hasValueSatisfying(cause -> assertThat(cause) + .contains(nonEmptyQueue) + .doesNotContain(emptyQueue)); + }); + } + + /** + * The second dead letter queue only exists in the second vhost: checking it with the configuration of the first + * one fails. + */ + private void declareQueues() throws Exception { + declareQueue(FIRST_DEAD_LETTER_QUEUE); + declareQueueInSecondVhost(SECOND_DEAD_LETTER_QUEUE); + } + + private void declareQueue(String queue) { + rabbitMQExtension.getSender() + .declareQueue(QueueSpecification.queue(queue).durable(DURABLE)) + .block(); + } + + private void declareQueueInSecondVhost(String queue) throws Exception { + try (Channel channel = secondVhostConnectionPool.getResilientConnection().block().createChannel()) { + channel.queueDeclare(queue, DURABLE, !EXCLUSIVE, !AUTO_DELETE, ImmutableMap.of()); + } + } + + private void publishInSecondVhost(String queue) throws Exception { + try (Channel channel = secondVhostConnectionPool.getResilientConnection().block().createChannel()) { + channel.basicPublish(DEFAULT_EXCHANGE, queue, null, "message".getBytes(UTF_8)); + } + } + + private void publish(String queue) { + rabbitMQExtension.getSender() + .send(Mono.just(new OutboundMessage(DEFAULT_EXCHANGE, queue, "message".getBytes(UTF_8)))) + .block(); + } +} diff --git a/backends-common/rabbitmq/src/test/java/org/apache/james/backends/rabbitmq/UnresponsiveServer.java b/backends-common/rabbitmq/src/test/java/org/apache/james/backends/rabbitmq/UnresponsiveServer.java new file mode 100644 index 0000000000..fcbe5e031b --- /dev/null +++ b/backends-common/rabbitmq/src/test/java/org/apache/james/backends/rabbitmq/UnresponsiveServer.java @@ -0,0 +1,76 @@ +/**************************************************************** + * 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.io.IOException; +import java.net.ServerSocket; +import java.net.Socket; +import java.net.URI; +import java.util.List; +import java.util.concurrent.CopyOnWriteArrayList; + +/** + * Accepts TCP connections but never answers, as a stalled RabbitMQ server or management API. + */ +class UnresponsiveServer implements AutoCloseable { + private final ServerSocket serverSocket; + private final List<Socket> acceptedSockets = new CopyOnWriteArrayList<>(); + private final Thread acceptingThread; + + UnresponsiveServer() throws IOException { + serverSocket = new ServerSocket(0); + acceptingThread = new Thread(() -> { + try { + while (!serverSocket.isClosed()) { + acceptedSockets.add(serverSocket.accept()); + } + } catch (IOException e) { + // the server is closed + } + }); + acceptingThread.setDaemon(true); + acceptingThread.start(); + } + + URI uri(String scheme) { + return URI.create(scheme + "://localhost:" + serverSocket.getLocalPort()); + } + + RabbitMQConfiguration configuration() { + return RabbitMQConfiguration.builder() + .amqpUri(uri("amqp")) + .managementUri(uri("http")) + .managementCredentials(RabbitMQFixture.DEFAULT_MANAGEMENT_CREDENTIAL) + .build(); + } + + int acceptedConnections() { + return acceptedSockets.size(); + } + + @Override + public void close() throws IOException, InterruptedException { + serverSocket.close(); + acceptingThread.join(); + for (Socket socket : acceptedSockets) { + socket.close(); + } + } +} diff --git a/event-bus/distributed/src/main/java/org/apache/james/events/RabbitMQContentDeletionEventBusDeadLetterQueueHealthCheck.java b/event-bus/distributed/src/main/java/org/apache/james/events/RabbitMQContentDeletionEventBusDeadLetterQueueHealthCheck.java deleted file mode 100644 index 411279eac8..0000000000 --- a/event-bus/distributed/src/main/java/org/apache/james/events/RabbitMQContentDeletionEventBusDeadLetterQueueHealthCheck.java +++ /dev/null @@ -1,63 +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.events; - -import jakarta.inject.Inject; - -import org.apache.james.backends.rabbitmq.RabbitMQConfiguration; -import org.apache.james.backends.rabbitmq.RabbitMQManagementAPI; -import org.apache.james.core.healthcheck.ComponentName; -import org.apache.james.core.healthcheck.HealthCheck; -import org.apache.james.core.healthcheck.Result; - -import reactor.core.publisher.Mono; -import reactor.core.scheduler.Schedulers; - -public class RabbitMQContentDeletionEventBusDeadLetterQueueHealthCheck implements HealthCheck { - public static final ComponentName COMPONENT_NAME = new ComponentName("RabbitMQContentDeletionEventBusDeadLetterQueueHealthCheck"); - private static final String DEFAULT_VHOST = "/"; - - private final RabbitMQConfiguration configuration; - private final RabbitMQManagementAPI api; - - @Inject - public RabbitMQContentDeletionEventBusDeadLetterQueueHealthCheck(RabbitMQConfiguration configuration) { - this.configuration = configuration; - this.api = RabbitMQManagementAPI.from(configuration); - } - - @Override - public ComponentName componentName() { - return COMPONENT_NAME; - } - - @Override - public Mono<Result> check() { - return Mono.fromCallable(() -> api.queueDetails(configuration.getVhost().orElse(DEFAULT_VHOST), NamingStrategy.CONTENT_DELETION_NAMING_STRATEGY.deadLetterQueue().getName()).getQueueLength()) - .map(queueSize -> { - if (queueSize != 0) { - return Result.degraded(COMPONENT_NAME, "RabbitMQ dead letter queue of the content deletion event bus contain events. This might indicate transient failure on event processing."); - } - return Result.healthy(COMPONENT_NAME); - }) - .onErrorResume(e -> Mono.just(Result.unhealthy(COMPONENT_NAME, "Error checking RabbitMQContentDeletionEventBusDeadLetterQueueHealthCheck", e))) - .subscribeOn(Schedulers.boundedElastic()); - } -} diff --git a/event-bus/distributed/src/main/java/org/apache/james/events/RabbitMQJmapEventBusDeadLetterQueueHealthCheck.java b/event-bus/distributed/src/main/java/org/apache/james/events/RabbitMQJmapEventBusDeadLetterQueueHealthCheck.java deleted file mode 100644 index b8f9ab48a4..0000000000 --- a/event-bus/distributed/src/main/java/org/apache/james/events/RabbitMQJmapEventBusDeadLetterQueueHealthCheck.java +++ /dev/null @@ -1,63 +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.events; - -import jakarta.inject.Inject; - -import org.apache.james.backends.rabbitmq.RabbitMQConfiguration; -import org.apache.james.backends.rabbitmq.RabbitMQManagementAPI; -import org.apache.james.core.healthcheck.ComponentName; -import org.apache.james.core.healthcheck.HealthCheck; -import org.apache.james.core.healthcheck.Result; - -import reactor.core.publisher.Mono; -import reactor.core.scheduler.Schedulers; - -public class RabbitMQJmapEventBusDeadLetterQueueHealthCheck implements HealthCheck { - public static final ComponentName COMPONENT_NAME = new ComponentName("RabbitMQJmapEventBusDeadLetterQueueHealthCheck"); - private static final String DEFAULT_VHOST = "/"; - - private final RabbitMQConfiguration configuration; - private final RabbitMQManagementAPI api; - - @Inject - public RabbitMQJmapEventBusDeadLetterQueueHealthCheck(RabbitMQConfiguration configuration) { - this.configuration = configuration; - this.api = RabbitMQManagementAPI.from(configuration); - } - - @Override - public ComponentName componentName() { - return COMPONENT_NAME; - } - - @Override - public Mono<Result> check() { - return Mono.fromCallable(() -> api.queueDetails(configuration.getVhost().orElse(DEFAULT_VHOST), NamingStrategy.JMAP_NAMING_STRATEGY.deadLetterQueue().getName()).getQueueLength()) - .map(queueSize -> { - if (queueSize != 0) { - return Result.degraded(COMPONENT_NAME, "RabbitMQ dead letter queue of the JMAP event bus contain events. This might indicate transient failure on event processing."); - } - return Result.healthy(COMPONENT_NAME); - }) - .onErrorResume(e -> Mono.just(Result.unhealthy(COMPONENT_NAME, "Error checking RabbitMQJmapEventBusDeadLetterQueueHealthCheck", e))) - .subscribeOn(Schedulers.boundedElastic()); - } -} diff --git a/event-bus/distributed/src/main/java/org/apache/james/events/RabbitMQMailboxEventBusDeadLetterQueueHealthCheck.java b/event-bus/distributed/src/main/java/org/apache/james/events/RabbitMQMailboxEventBusDeadLetterQueueHealthCheck.java deleted file mode 100644 index 6af4176a4c..0000000000 --- a/event-bus/distributed/src/main/java/org/apache/james/events/RabbitMQMailboxEventBusDeadLetterQueueHealthCheck.java +++ /dev/null @@ -1,65 +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.events; - -import jakarta.inject.Inject; - -import org.apache.james.backends.rabbitmq.RabbitMQConfiguration; -import org.apache.james.backends.rabbitmq.RabbitMQManagementAPI; -import org.apache.james.core.healthcheck.ComponentName; -import org.apache.james.core.healthcheck.HealthCheck; -import org.apache.james.core.healthcheck.Result; - -import reactor.core.publisher.Mono; -import reactor.core.scheduler.Schedulers; - -public class RabbitMQMailboxEventBusDeadLetterQueueHealthCheck implements HealthCheck { - public static final ComponentName COMPONENT_NAME = new ComponentName("RabbitMQMailboxEventBusDeadLetterQueueHealthCheck"); - private static final String DEFAULT_VHOST = "/"; - - private final RabbitMQConfiguration configuration; - private final NamingStrategy mailboxEventNamingStrategy; - private final RabbitMQManagementAPI api; - - @Inject - public RabbitMQMailboxEventBusDeadLetterQueueHealthCheck(RabbitMQConfiguration configuration, NamingStrategy mailboxEventNamingStrategy) { - this.configuration = configuration; - this.mailboxEventNamingStrategy = mailboxEventNamingStrategy; - this.api = RabbitMQManagementAPI.from(configuration); - } - - @Override - public ComponentName componentName() { - return COMPONENT_NAME; - } - - @Override - public Mono<Result> check() { - return Mono.fromCallable(() -> api.queueDetails(configuration.getVhost().orElse(DEFAULT_VHOST), mailboxEventNamingStrategy.deadLetterQueue().getName()).getQueueLength()) - .map(queueSize -> { - if (queueSize != 0) { - return Result.degraded(COMPONENT_NAME, "RabbitMQ dead letter queue of the mailbox event bus contain events. This might indicate transient failure on event processing."); - } - return Result.healthy(COMPONENT_NAME); - }) - .onErrorResume(e -> Mono.just(Result.unhealthy(COMPONENT_NAME, "Error checking RabbitMQMailboxEventBusDeadLetterQueueHealthCheck", e))) - .subscribeOn(Schedulers.boundedElastic()); // Reading the management API is blocking - } -} diff --git a/event-bus/distributed/src/test/java/org/apache/james/events/RabbitMQContentDeletionEventBusDeadLetterQueueHealthCheckTest.java b/event-bus/distributed/src/test/java/org/apache/james/events/RabbitMQContentDeletionEventBusDeadLetterQueueHealthCheckTest.java deleted file mode 100644 index bf8e5282cb..0000000000 --- a/event-bus/distributed/src/test/java/org/apache/james/events/RabbitMQContentDeletionEventBusDeadLetterQueueHealthCheckTest.java +++ /dev/null @@ -1,131 +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.events; - -import static com.rabbitmq.client.MessageProperties.PERSISTENT_TEXT_PLAIN; -import static org.apache.james.backends.rabbitmq.Constants.AUTO_DELETE; -import static org.apache.james.backends.rabbitmq.Constants.DIRECT_EXCHANGE; -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.EXCHANGE_NAME; -import static org.apache.james.backends.rabbitmq.RabbitMQFixture.awaitAtMostOneMinute; -import static org.assertj.core.api.Assertions.assertThat; - -import java.io.IOException; -import java.net.URISyntaxException; -import java.nio.charset.StandardCharsets; -import java.util.Arrays; -import java.util.concurrent.TimeoutException; - -import org.apache.james.backends.rabbitmq.DockerRabbitMQ; -import org.apache.james.backends.rabbitmq.RabbitMQExtension; -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.rabbitmq.client.AMQP; -import com.rabbitmq.client.Channel; -import com.rabbitmq.client.Connection; -import com.rabbitmq.client.ConnectionFactory; - -class RabbitMQContentDeletionEventBusDeadLetterQueueHealthCheckTest { - @RegisterExtension - RabbitMQExtension rabbitMQExtension = RabbitMQExtension.singletonRabbitMQ() - .isolationPolicy(RabbitMQExtension.IsolationPolicy.STRONG); - - public static final ImmutableMap<String, Object> NO_QUEUE_DECLARE_ARGUMENTS = ImmutableMap.of(); - public static final String ROUTING_KEY_CONTENT_DELETION_EVENTS_EVENT_BUS = "contentDeletionRoutingKey"; - - private Connection connection; - private Channel channel; - private RabbitMQContentDeletionEventBusDeadLetterQueueHealthCheck testee; - - @BeforeEach - void setup(DockerRabbitMQ rabbitMQ) throws IOException, TimeoutException, URISyntaxException { - ConnectionFactory connectionFactory = rabbitMQ.connectionFactory(); - connectionFactory.setNetworkRecoveryInterval(1000); - connection = connectionFactory.newConnection(); - channel = connection.createChannel(); - testee = new RabbitMQContentDeletionEventBusDeadLetterQueueHealthCheck(rabbitMQ.getConfiguration()); - } - - @AfterEach - void tearDown(DockerRabbitMQ rabbitMQ) throws Exception { - closeQuietly(connection, channel); - rabbitMQ.reset(); - } - - @Test - void healthCheckShouldReturnUnhealthyWhenRabbitMQIsDown() throws Exception { - rabbitMQExtension.getRabbitMQ().stopApp(); - - assertThat(testee.check().block().isUnHealthy()).isTrue(); - } - - @Test - void healthCheckShouldReturnHealthyWhenContentDeletionEventBusDeadLetterQueueIsEmpty() throws Exception { - createDeadLetterQueue(channel, NamingStrategy.CONTENT_DELETION_NAMING_STRATEGY, ROUTING_KEY_CONTENT_DELETION_EVENTS_EVENT_BUS); - - assertThat(testee.check().block().isHealthy()).isTrue(); - } - - @Test - void healthCheckShouldReturnUnhealthyWhenThereIsNoDeadLetterQueue() { - assertThat(testee.check().block().isUnHealthy()).isTrue(); - } - - @Test - void healthCheckShouldReturnDegradedWhenContentDeletionEventBusDeadLetterQueueIsNotEmpty() throws Exception { - createDeadLetterQueue(channel, NamingStrategy.CONTENT_DELETION_NAMING_STRATEGY, ROUTING_KEY_CONTENT_DELETION_EVENTS_EVENT_BUS); - publishAMessage(channel, ROUTING_KEY_CONTENT_DELETION_EVENTS_EVENT_BUS); - - awaitAtMostOneMinute.until(() -> testee.check().block().isDegraded()); - } - - private void createDeadLetterQueue(Channel channel, NamingStrategy namingStrategy, String routingKey) throws IOException { - channel.exchangeDeclare(EXCHANGE_NAME, DIRECT_EXCHANGE, DURABLE); - channel.queueDeclare(namingStrategy.deadLetterQueue().getName(), DURABLE, !EXCLUSIVE, AUTO_DELETE, NO_QUEUE_DECLARE_ARGUMENTS).getQueue(); - channel.queueBind(namingStrategy.deadLetterQueue().getName(), EXCHANGE_NAME, routingKey); - } - - private void publishAMessage(Channel channel, String routingKey) throws IOException { - AMQP.BasicProperties basicProperties = new AMQP.BasicProperties.Builder() - .deliveryMode(PERSISTENT_TEXT_PLAIN.getDeliveryMode()) - .priority(PERSISTENT_TEXT_PLAIN.getPriority()) - .contentType(PERSISTENT_TEXT_PLAIN.getContentType()) - .build(); - - channel.basicPublish(EXCHANGE_NAME, routingKey, basicProperties, "Hello, world!".getBytes(StandardCharsets.UTF_8)); - } - - private void closeQuietly(AutoCloseable... closeables) { - Arrays.stream(closeables).forEach(this::closeQuietly); - } - - private void closeQuietly(AutoCloseable closeable) { - try { - closeable.close(); - } catch (Exception e) { - //ignore error - } - } -} diff --git a/event-bus/distributed/src/test/java/org/apache/james/events/RabbitMQJmapEventBusDeadLetterQueueHealthCheckTest.java b/event-bus/distributed/src/test/java/org/apache/james/events/RabbitMQJmapEventBusDeadLetterQueueHealthCheckTest.java deleted file mode 100644 index 472d713a53..0000000000 --- a/event-bus/distributed/src/test/java/org/apache/james/events/RabbitMQJmapEventBusDeadLetterQueueHealthCheckTest.java +++ /dev/null @@ -1,131 +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.events; - -import static com.rabbitmq.client.MessageProperties.PERSISTENT_TEXT_PLAIN; -import static org.apache.james.backends.rabbitmq.Constants.AUTO_DELETE; -import static org.apache.james.backends.rabbitmq.Constants.DIRECT_EXCHANGE; -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.EXCHANGE_NAME; -import static org.apache.james.backends.rabbitmq.RabbitMQFixture.awaitAtMostOneMinute; -import static org.assertj.core.api.Assertions.assertThat; - -import java.io.IOException; -import java.net.URISyntaxException; -import java.nio.charset.StandardCharsets; -import java.util.Arrays; -import java.util.concurrent.TimeoutException; - -import org.apache.james.backends.rabbitmq.DockerRabbitMQ; -import org.apache.james.backends.rabbitmq.RabbitMQExtension; -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.rabbitmq.client.AMQP; -import com.rabbitmq.client.Channel; -import com.rabbitmq.client.Connection; -import com.rabbitmq.client.ConnectionFactory; - -class RabbitMQJmapEventBusDeadLetterQueueHealthCheckTest { - @RegisterExtension - RabbitMQExtension rabbitMQExtension = RabbitMQExtension.singletonRabbitMQ() - .isolationPolicy(RabbitMQExtension.IsolationPolicy.STRONG); - - public static final ImmutableMap<String, Object> NO_QUEUE_DECLARE_ARGUMENTS = ImmutableMap.of(); - public static final String ROUTING_KEY_JMAP_EVENTS_EVENT_BUS = "jmapEventsRoutingKey"; - - private Connection connection; - private Channel channel; - private RabbitMQJmapEventBusDeadLetterQueueHealthCheck testee; - - @BeforeEach - void setup(DockerRabbitMQ rabbitMQ) throws IOException, TimeoutException, URISyntaxException { - ConnectionFactory connectionFactory = rabbitMQ.connectionFactory(); - connectionFactory.setNetworkRecoveryInterval(1000); - connection = connectionFactory.newConnection(); - channel = connection.createChannel(); - testee = new RabbitMQJmapEventBusDeadLetterQueueHealthCheck(rabbitMQ.getConfiguration()); - } - - @AfterEach - void tearDown(DockerRabbitMQ rabbitMQ) throws Exception { - closeQuietly(connection, channel); - rabbitMQ.reset(); - } - - @Test - void healthCheckShouldReturnUnhealthyWhenRabbitMQIsDown() throws Exception { - rabbitMQExtension.getRabbitMQ().stopApp(); - - assertThat(testee.check().block().isUnHealthy()).isTrue(); - } - - @Test - void healthCheckShouldReturnHealthyWhenJmapEventBusDeadLetterQueueIsEmpty() throws Exception { - createDeadLetterQueue(channel, NamingStrategy.JMAP_NAMING_STRATEGY, ROUTING_KEY_JMAP_EVENTS_EVENT_BUS); - - assertThat(testee.check().block().isHealthy()).isTrue(); - } - - @Test - void healthCheckShouldReturnUnhealthyWhenThereIsNoDeadLetterQueue() { - assertThat(testee.check().block().isUnHealthy()).isTrue(); - } - - @Test - void healthCheckShouldReturnDegradedWhenJmapEventBusDeadLetterQueueIsNotEmpty() throws Exception { - createDeadLetterQueue(channel, NamingStrategy.JMAP_NAMING_STRATEGY, ROUTING_KEY_JMAP_EVENTS_EVENT_BUS); - publishAMessage(channel, ROUTING_KEY_JMAP_EVENTS_EVENT_BUS); - - awaitAtMostOneMinute.until(() -> testee.check().block().isDegraded()); - } - - private void createDeadLetterQueue(Channel channel, NamingStrategy namingStrategy, String routingKey) throws IOException { - channel.exchangeDeclare(EXCHANGE_NAME, DIRECT_EXCHANGE, DURABLE); - channel.queueDeclare(namingStrategy.deadLetterQueue().getName(), DURABLE, !EXCLUSIVE, AUTO_DELETE, NO_QUEUE_DECLARE_ARGUMENTS).getQueue(); - channel.queueBind(namingStrategy.deadLetterQueue().getName(), EXCHANGE_NAME, routingKey); - } - - private void publishAMessage(Channel channel, String routingKey) throws IOException { - AMQP.BasicProperties basicProperties = new AMQP.BasicProperties.Builder() - .deliveryMode(PERSISTENT_TEXT_PLAIN.getDeliveryMode()) - .priority(PERSISTENT_TEXT_PLAIN.getPriority()) - .contentType(PERSISTENT_TEXT_PLAIN.getContentType()) - .build(); - - channel.basicPublish(EXCHANGE_NAME, routingKey, basicProperties, "Hello, world!".getBytes(StandardCharsets.UTF_8)); - } - - private void closeQuietly(AutoCloseable... closeables) { - Arrays.stream(closeables).forEach(this::closeQuietly); - } - - private void closeQuietly(AutoCloseable closeable) { - try { - closeable.close(); - } catch (Exception e) { - //ignore error - } - } -} diff --git a/event-bus/distributed/src/test/java/org/apache/james/events/RabbitMQMailboxEventBusDeadLetterQueueHealthCheckTest.java b/event-bus/distributed/src/test/java/org/apache/james/events/RabbitMQMailboxEventBusDeadLetterQueueHealthCheckTest.java deleted file mode 100644 index e33b8ba526..0000000000 --- a/event-bus/distributed/src/test/java/org/apache/james/events/RabbitMQMailboxEventBusDeadLetterQueueHealthCheckTest.java +++ /dev/null @@ -1,132 +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.events; - -import static com.rabbitmq.client.MessageProperties.PERSISTENT_TEXT_PLAIN; -import static org.apache.james.backends.rabbitmq.Constants.AUTO_DELETE; -import static org.apache.james.backends.rabbitmq.Constants.DIRECT_EXCHANGE; -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.EXCHANGE_NAME; -import static org.apache.james.backends.rabbitmq.RabbitMQFixture.awaitAtMostOneMinute; -import static org.assertj.core.api.Assertions.assertThat; - -import java.io.IOException; -import java.net.URISyntaxException; -import java.nio.charset.StandardCharsets; -import java.util.Arrays; -import java.util.concurrent.TimeoutException; - -import org.apache.james.backends.rabbitmq.DockerRabbitMQ; -import org.apache.james.backends.rabbitmq.RabbitMQExtension; -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.rabbitmq.client.AMQP; -import com.rabbitmq.client.Channel; -import com.rabbitmq.client.Connection; -import com.rabbitmq.client.ConnectionFactory; - -class RabbitMQMailboxEventBusDeadLetterQueueHealthCheckTest { - @RegisterExtension - RabbitMQExtension rabbitMQExtension = RabbitMQExtension.singletonRabbitMQ() - .isolationPolicy(RabbitMQExtension.IsolationPolicy.STRONG); - - public static final ImmutableMap<String, Object> NO_QUEUE_DECLARE_ARGUMENTS = ImmutableMap.of(); - public static final NamingStrategy MAILBOX_EVENTS_NAMING_STRATEGY = new DefaultNamingStrategy(new EventBusName("mailboxEvents")); - public static final String ROUTING_KEY_MAILBOX_EVENTS_EVENT_BUS = "mailboxEventsRoutingKey"; - - private Connection connection; - private Channel channel; - private RabbitMQMailboxEventBusDeadLetterQueueHealthCheck testee; - - @BeforeEach - void setup(DockerRabbitMQ rabbitMQ) throws IOException, TimeoutException, URISyntaxException { - ConnectionFactory connectionFactory = rabbitMQ.connectionFactory(); - connectionFactory.setNetworkRecoveryInterval(1000); - connection = connectionFactory.newConnection(); - channel = connection.createChannel(); - testee = new RabbitMQMailboxEventBusDeadLetterQueueHealthCheck(rabbitMQ.getConfiguration(), MAILBOX_EVENTS_NAMING_STRATEGY); - } - - @AfterEach - void tearDown(DockerRabbitMQ rabbitMQ) throws Exception { - closeQuietly(connection, channel); - rabbitMQ.reset(); - } - - @Test - void healthCheckShouldReturnUnhealthyWhenRabbitMQIsDown() throws Exception { - rabbitMQExtension.getRabbitMQ().stopApp(); - - assertThat(testee.check().block().isUnHealthy()).isTrue(); - } - - @Test - void healthCheckShouldReturnHealthyWhenMailboxEventBusDeadLetterQueueIsEmpty() throws Exception { - createDeadLetterQueue(channel, MAILBOX_EVENTS_NAMING_STRATEGY, ROUTING_KEY_MAILBOX_EVENTS_EVENT_BUS); - - assertThat(testee.check().block().isHealthy()).isTrue(); - } - - @Test - void healthCheckShouldReturnUnhealthyWhenThereIsNoDeadLetterQueue() { - assertThat(testee.check().block().isUnHealthy()).isTrue(); - } - - @Test - void healthCheckShouldReturnDegradedWhenMailboxEventBusDeadLetterQueueIsNotEmpty() throws Exception { - createDeadLetterQueue(channel, MAILBOX_EVENTS_NAMING_STRATEGY, ROUTING_KEY_MAILBOX_EVENTS_EVENT_BUS); - publishAMessage(channel, ROUTING_KEY_MAILBOX_EVENTS_EVENT_BUS); - - awaitAtMostOneMinute.until(() -> testee.check().block().isDegraded()); - } - - private void createDeadLetterQueue(Channel channel, NamingStrategy namingStrategy, String routingKey) throws IOException { - channel.exchangeDeclare(EXCHANGE_NAME, DIRECT_EXCHANGE, DURABLE); - channel.queueDeclare(namingStrategy.deadLetterQueue().getName(), DURABLE, !EXCLUSIVE, AUTO_DELETE, NO_QUEUE_DECLARE_ARGUMENTS).getQueue(); - channel.queueBind(namingStrategy.deadLetterQueue().getName(), EXCHANGE_NAME, routingKey); - } - - private void publishAMessage(Channel channel, String routingKey) throws IOException { - AMQP.BasicProperties basicProperties = new AMQP.BasicProperties.Builder() - .deliveryMode(PERSISTENT_TEXT_PLAIN.getDeliveryMode()) - .priority(PERSISTENT_TEXT_PLAIN.getPriority()) - .contentType(PERSISTENT_TEXT_PLAIN.getContentType()) - .build(); - - channel.basicPublish(EXCHANGE_NAME, routingKey, basicProperties, "Hello, world!".getBytes(StandardCharsets.UTF_8)); - } - - private void closeQuietly(AutoCloseable... closeables) { - Arrays.stream(closeables).forEach(this::closeQuietly); - } - - private void closeQuietly(AutoCloseable closeable) { - try { - closeable.close(); - } catch (Exception e) { - //ignore error - } - } -} 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 409dc08f40..b1940eac7b 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 @@ -25,6 +25,7 @@ import java.util.Set; import jakarta.inject.Named; +import org.apache.james.backends.rabbitmq.MonitoredDeadLetterQueue; import org.apache.james.backends.rabbitmq.RabbitMQConfiguration; import org.apache.james.backends.rabbitmq.SimpleConnectionPool; import org.apache.james.core.healthcheck.HealthCheck; @@ -36,7 +37,6 @@ 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.RabbitMQContentDeletionEventBusDeadLetterQueueHealthCheck; import org.apache.james.events.RabbitMQEventBus; import org.apache.james.events.RetryBackoffConfiguration; import org.apache.james.events.RoutingKeyConverter; @@ -92,8 +92,8 @@ public class ContentDeletionEventBusModule extends AbstractModule { } @ProvidesIntoSet - HealthCheck contentDeletionEventBusDeadLetterQueueHealthCheck(RabbitMQConfiguration rabbitMQConfiguration) { - return new RabbitMQContentDeletionEventBusDeadLetterQueueHealthCheck(rabbitMQConfiguration); + MonitoredDeadLetterQueue contentDeletionEventBusDeadLetterQueue(RabbitMQConfiguration rabbitMQConfiguration) { + return new MonitoredDeadLetterQueue(rabbitMQConfiguration, CONTENT_DELETION_NAMING_STRATEGY.deadLetterQueue().getName()); } @Provides 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 d06f83b356..0dfabc8319 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 @@ -23,6 +23,7 @@ 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.RabbitMQConfiguration; import org.apache.james.backends.rabbitmq.SimpleConnectionPool; import org.apache.james.core.healthcheck.HealthCheck; @@ -34,7 +35,6 @@ 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.RabbitMQJmapEventBusDeadLetterQueueHealthCheck; import org.apache.james.events.RetryBackoffConfiguration; import org.apache.james.events.RoutingKeyConverter; import org.apache.james.jmap.InjectionKeys; @@ -90,8 +90,8 @@ public class JMAPEventBusModule extends AbstractModule { } @ProvidesIntoSet - HealthCheck jmapEventBusDeadLetterQueueHealthCheck(RabbitMQConfiguration rabbitMQConfiguration) { - return new RabbitMQJmapEventBusDeadLetterQueueHealthCheck(rabbitMQConfiguration); + MonitoredDeadLetterQueue jmapEventBusDeadLetterQueue(RabbitMQConfiguration rabbitMQConfiguration) { + return new MonitoredDeadLetterQueue(rabbitMQConfiguration, JMAP_NAMING_STRATEGY.deadLetterQueue().getName()); } @Provides 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 bf7dbf63cc..e5fa0764f2 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 @@ -21,6 +21,7 @@ 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.RabbitMQConfiguration; import org.apache.james.backends.rabbitmq.SimpleConnectionPool; import org.apache.james.core.healthcheck.HealthCheck; @@ -33,7 +34,6 @@ 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.RabbitMQMailboxEventBusDeadLetterQueueHealthCheck; import org.apache.james.events.RegistrationKey; import org.apache.james.events.RetryBackoffConfiguration; import org.apache.james.events.RoutingKeyConverter; @@ -58,9 +58,11 @@ public class MailboxEventBusModule extends AbstractModule { bind(RetryBackoffConfiguration.class).toInstance(RetryBackoffConfiguration.DEFAULT); bind(EventBusId.class).toInstance(EventBusId.random()); + } - Multibinder.newSetBinder(binder(), HealthCheck.class) - .addBinding().to(RabbitMQMailboxEventBusDeadLetterQueueHealthCheck.class); + @ProvidesIntoSet + MonitoredDeadLetterQueue deadLetterQueue(NamingStrategy namingStrategy, RabbitMQConfiguration configuration) { + return new MonitoredDeadLetterQueue(configuration, namingStrategy.deadLetterQueue().getName()); } @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 a3688de353..0033ea812f 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 @@ -23,33 +23,39 @@ import static org.apache.james.modules.queue.rabbitmq.RabbitMQModule.RABBITMQ_CO import jakarta.inject.Named; import jakarta.inject.Singleton; +import org.apache.james.backends.rabbitmq.MonitoredDeadLetterQueue; +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.RabbitMQMailQueueDeadLetterQueueHealthCheck; import org.apache.james.queue.rabbitmq.RabbitMQMailQueueFactory; import org.apache.james.queue.rabbitmq.view.RabbitMQMailQueueConfiguration; import com.google.inject.AbstractModule; import com.google.inject.Provides; import com.google.inject.multibindings.Multibinder; +import com.google.inject.multibindings.ProvidesIntoSet; public class RabbitMQMailQueueModule extends AbstractModule { @Override protected void configure() { Multibinder<SimpleConnectionPool.ReconnectionHandler> reconnectionHandlerMultibinder = Multibinder.newSetBinder(binder(), SimpleConnectionPool.ReconnectionHandler.class); reconnectionHandlerMultibinder.addBinding().to(SpoolerReconnectionHandler.class); - Multibinder<HealthCheck> healthCheckMultiBinder = Multibinder.newSetBinder(binder(), HealthCheck.class); - healthCheckMultiBinder.addBinding().to(RabbitMQMailQueueDeadLetterQueueHealthCheck.class); Multibinder.newSetBinder(binder(), HealthCheck.class).addBinding() .to(RabbitMQMailQueueConsumerHealthCheck.class); } + @ProvidesIntoSet + MonitoredDeadLetterQueue spoolDeadLetterQueue(RabbitMQConfiguration configuration) { + return new MonitoredDeadLetterQueue(configuration, MailQueueName.fromString(MailQueueFactory.SPOOL.asString()).toDeadLetterQueueName()); + } + @Provides @Singleton public MailQueueFactory<RabbitMQMailQueue> provideRabbitMQMailQueueFactoryProxy(RabbitMQMailQueueFactory queueFactory) { 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 856761f726..1646751d0a 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 @@ -26,7 +26,9 @@ import jakarta.inject.Provider; import jakarta.inject.Singleton; import org.apache.commons.configuration2.ex.ConfigurationException; +import org.apache.james.backends.rabbitmq.MonitoredDeadLetterQueue; import org.apache.james.backends.rabbitmq.RabbitMQConfiguration; +import org.apache.james.backends.rabbitmq.RabbitMQDeadLetterQueuesHealthCheck; import org.apache.james.backends.rabbitmq.RabbitMQHealthCheck; import org.apache.james.backends.rabbitmq.ReactorRabbitMQChannelPool; import org.apache.james.backends.rabbitmq.ReceiverProvider; @@ -62,6 +64,8 @@ public class RabbitMQModule extends AbstractModule { Multibinder<HealthCheck> healthCheckMultiBinder = Multibinder.newSetBinder(binder(), HealthCheck.class); healthCheckMultiBinder.addBinding().to(RabbitMQHealthCheck.class); + healthCheckMultiBinder.addBinding().to(RabbitMQDeadLetterQueuesHealthCheck.class); + Multibinder.newSetBinder(binder(), MonitoredDeadLetterQueue.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 7ede034cf9..1d03765e61 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 @@ -29,19 +29,25 @@ import static org.hamcrest.Matchers.hasSize; import static org.hamcrest.Matchers.is; import java.util.List; +import java.util.Set; + +import jakarta.inject.Inject; import org.apache.james.CassandraExtension; import org.apache.james.CassandraRabbitMQJamesConfiguration; import org.apache.james.CassandraRabbitMQJamesServerMain; import org.apache.james.DockerOpenSearchExtension; +import org.apache.james.GuiceJamesServer; import org.apache.james.JamesServerBuilder; 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.junit.categories.BasicFeature; import org.apache.james.modules.AwsS3BlobStoreExtension; import org.apache.james.modules.RabbitMQExtension; import org.apache.james.modules.blobstore.BlobStoreConfiguration; +import org.apache.james.utils.GuiceProbe; import org.apache.james.webadmin.integration.WebAdminServerIntegrationImmutableTest; import org.apache.james.webadmin.routes.HealthCheckRoutes; import org.apache.james.webadmin.routes.TasksRoutes; @@ -50,8 +56,24 @@ import org.junit.jupiter.api.Tag; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.extension.RegisterExtension; +import com.google.inject.multibindings.Multibinder; + @Tag(BasicFeature.TAG) class RabbitMQWebAdminServerIntegrationImmutableTest extends WebAdminServerIntegrationImmutableTest { + public static class MonitoredRabbitMQProbe implements GuiceProbe { + private final Set<MonitoredDeadLetterQueue> deadLetterQueues; + + @Inject + public MonitoredRabbitMQProbe(Set<MonitoredDeadLetterQueue> deadLetterQueues) { + this.deadLetterQueues = deadLetterQueues; + } + + List<String> deadLetterQueues() { + return deadLetterQueues.stream() + .map(MonitoredDeadLetterQueue::queue) + .toList(); + } + } @RegisterExtension static JamesServerExtension testExtension = new JamesServerBuilder<CassandraRabbitMQJamesConfiguration>(tmpDir -> @@ -70,7 +92,9 @@ class RabbitMQWebAdminServerIntegrationImmutableTest extends WebAdminServerInteg .extension(new CassandraExtension()) .extension(new AwsS3BlobStoreExtension()) .extension(new RabbitMQExtension()) - .server(CassandraRabbitMQJamesServerMain::createServer) + .server(configuration -> CassandraRabbitMQJamesServerMain.createServer(configuration) + .overrideWith(binder -> Multibinder.newSetBinder(binder, GuiceProbe.class) + .addBinding().to(MonitoredRabbitMQProbe.class))) .lifeCycle(PER_CLASS) .build(); @@ -138,12 +162,26 @@ class RabbitMQWebAdminServerIntegrationImmutableTest extends WebAdminServerInteg .getList("checks.componentName", String.class); assertThat(listComponentNames).containsOnly("Guice application lifecycle", "EmptyErrorMailRepository", - "RabbitMQ backend", "RabbitMQMailQueueDeadLetterQueueHealthCheck", - "RabbitMQMailboxEventBusDeadLetterQueueHealthCheck", "MailReceptionCheck", + "RabbitMQ backend", "RabbitMQDeadLetterQueues", "MailReceptionCheck", "Cassandra backend", "EventDeadLettersHealthCheck", "MessageFastViewProjection", "RabbitMQMailQueue BrowseStart", "OpenSearch Backend", "ObjectStorage", "DistributedTaskManagerConsumers", "EventbusConsumers-jmapEvent", "MailQueueConsumers", "EventbusConsumers-mailboxEvent", - "RabbitMQJmapEventBusDeadLetterQueueHealthCheck", "IMAPHealthCheck", - "RabbitMQContentDeletionEventBusDeadLetterQueueHealthCheck", "EventbusConsumers-contentDeletionEvent"); + "IMAPHealthCheck", "EventbusConsumers-contentDeletionEvent"); + } + + @Test + void everyRabbitMQDeadLetterQueueShouldBeMonitored(GuiceJamesServer server) { + assertThat(server.getProbe(MonitoredRabbitMQProbe.class).deadLetterQueues()) + .containsExactlyInAnyOrder("mailboxEvent-dead-letter-queue", "jmapEvent-dead-letter-queue", "contentDeletionEvent-dead-letter-queue", + "JamesMailQueue-dead-letter-queue-spool"); + } + + @Test + void rabbitMQDeadLetterQueuesShouldBeHealthy() { + when() + .get(HealthCheckRoutes.HEALTHCHECK + "/checks/RabbitMQDeadLetterQueues") + .then() + .statusCode(HttpStatus.OK_200) + .body("status", is("healthy")); } } diff --git a/server/queue/queue-rabbitmq/src/main/java/org/apache/james/queue/rabbitmq/MailQueueName.java b/server/queue/queue-rabbitmq/src/main/java/org/apache/james/queue/rabbitmq/MailQueueName.java index f70151e21e..316b1b7252 100644 --- a/server/queue/queue-rabbitmq/src/main/java/org/apache/james/queue/rabbitmq/MailQueueName.java +++ b/server/queue/queue-rabbitmq/src/main/java/org/apache/james/queue/rabbitmq/MailQueueName.java @@ -140,7 +140,7 @@ public final class MailQueueName { return DEAD_LETTER_EXCHANGE_PREFIX + name; } - String toDeadLetterQueueName() { + public String toDeadLetterQueueName() { return DEAD_LETTER_QUEUE_PREFIX + name; } diff --git a/server/queue/queue-rabbitmq/src/main/java/org/apache/james/queue/rabbitmq/RabbitMQMailQueueDeadLetterQueueHealthCheck.java b/server/queue/queue-rabbitmq/src/main/java/org/apache/james/queue/rabbitmq/RabbitMQMailQueueDeadLetterQueueHealthCheck.java deleted file mode 100644 index 3ff86e60af..0000000000 --- a/server/queue/queue-rabbitmq/src/main/java/org/apache/james/queue/rabbitmq/RabbitMQMailQueueDeadLetterQueueHealthCheck.java +++ /dev/null @@ -1,66 +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.queue.api.MailQueueFactory.SPOOL; - -import jakarta.inject.Inject; - -import org.apache.james.backends.rabbitmq.RabbitMQConfiguration; -import org.apache.james.backends.rabbitmq.RabbitMQManagementAPI; -import org.apache.james.core.healthcheck.ComponentName; -import org.apache.james.core.healthcheck.HealthCheck; -import org.apache.james.core.healthcheck.Result; - -import reactor.core.publisher.Mono; -import reactor.core.scheduler.Schedulers; - -public class RabbitMQMailQueueDeadLetterQueueHealthCheck implements HealthCheck { - public static final MailQueueName JAMES_MAIL_QUEUE_NAME = MailQueueName.fromString(SPOOL.asString()); - public static final ComponentName COMPONENT_NAME = new ComponentName("RabbitMQMailQueueDeadLetterQueueHealthCheck"); - private static final String DEFAULT_VHOST = "/"; - - private final RabbitMQConfiguration configuration; - private final RabbitMQManagementAPI api; - - @Inject - public RabbitMQMailQueueDeadLetterQueueHealthCheck(RabbitMQConfiguration configuration) { - this.configuration = configuration; - this.api = RabbitMQManagementAPI.from(configuration); - } - - @Override - public ComponentName componentName() { - return COMPONENT_NAME; - } - - @Override - public Mono<Result> check() { - return Mono.fromCallable(() -> api.queueDetails(configuration.getVhost().orElse(DEFAULT_VHOST), JAMES_MAIL_QUEUE_NAME.toDeadLetterQueueName())) - .map(queueDetails -> { - if (queueDetails.getQueueLength() != 0) { - return Result.degraded(COMPONENT_NAME, "RabbitMQ dead letter queue of the mail queue contain messages. This might indicate transient failure on mail processing."); - } - return Result.healthy(COMPONENT_NAME); - }) - .onErrorResume(e -> Mono.just(Result.unhealthy(COMPONENT_NAME, "Error checking RabbitMQMailQueueDeadLetterQueueHealthCheck", e))) - .subscribeOn(Schedulers.boundedElastic()); - } -} diff --git a/server/queue/queue-rabbitmq/src/test/java/org/apache/james/queue/rabbitmq/RabbitMQMailQueueDeadLetterQueueHealthCheckTest.java b/server/queue/queue-rabbitmq/src/test/java/org/apache/james/queue/rabbitmq/RabbitMQMailQueueDeadLetterQueueHealthCheckTest.java deleted file mode 100644 index f32bd40a4e..0000000000 --- a/server/queue/queue-rabbitmq/src/test/java/org/apache/james/queue/rabbitmq/RabbitMQMailQueueDeadLetterQueueHealthCheckTest.java +++ /dev/null @@ -1,133 +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 com.rabbitmq.client.MessageProperties.PERSISTENT_TEXT_PLAIN; -import static org.apache.james.backends.rabbitmq.Constants.AUTO_DELETE; -import static org.apache.james.backends.rabbitmq.Constants.DIRECT_EXCHANGE; -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.EXCHANGE_NAME; -import static org.apache.james.backends.rabbitmq.RabbitMQFixture.ROUTING_KEY; -import static org.apache.james.backends.rabbitmq.RabbitMQFixture.awaitAtMostOneMinute; -import static org.apache.james.queue.rabbitmq.RabbitMQMailQueueDeadLetterQueueHealthCheck.JAMES_MAIL_QUEUE_NAME; -import static org.assertj.core.api.Assertions.assertThat; - -import java.io.IOException; -import java.net.URISyntaxException; -import java.nio.charset.StandardCharsets; -import java.util.Arrays; -import java.util.concurrent.TimeoutException; - -import org.apache.james.backends.rabbitmq.DockerRabbitMQ; -import org.apache.james.backends.rabbitmq.RabbitMQExtension; -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.rabbitmq.client.AMQP; -import com.rabbitmq.client.Channel; -import com.rabbitmq.client.Connection; -import com.rabbitmq.client.ConnectionFactory; - -class RabbitMQMailQueueDeadLetterQueueHealthCheckTest { - @RegisterExtension - RabbitMQExtension rabbitMQExtension = RabbitMQExtension.singletonRabbitMQ() - .isolationPolicy(RabbitMQExtension.IsolationPolicy.STRONG); - - public static final ImmutableMap<String, Object> NO_QUEUE_DECLARE_ARGUMENTS = ImmutableMap.of(); - - private Connection connection; - private Channel channel; - private RabbitMQMailQueueDeadLetterQueueHealthCheck testee; - - @BeforeEach - void setup(DockerRabbitMQ rabbitMQ) throws IOException, TimeoutException, URISyntaxException { - ConnectionFactory connectionFactory = rabbitMQ.connectionFactory(); - connectionFactory.setNetworkRecoveryInterval(1000); - connection = connectionFactory.newConnection(); - channel = connection.createChannel(); - testee = new RabbitMQMailQueueDeadLetterQueueHealthCheck(rabbitMQ.getConfiguration()); - } - - @AfterEach - void tearDown(DockerRabbitMQ rabbitMQ) throws Exception { - closeQuietly(connection, channel); - rabbitMQ.reset(); - } - - @Test - void healthCheckShouldReturnUnhealthyWhenRabbitMQIsDown() throws Exception { - rabbitMQExtension.getRabbitMQ().stopApp(); - - assertThat(testee.check().block().isUnHealthy()).isTrue(); - } - - @Test - void healthCheckShouldReturnHealthyWhenDeadLetterQueueIsEmpty() throws Exception { - createDeadLetterQueue(channel); - - assertThat(testee.check().block().isHealthy()).isTrue(); - } - - @Test - void healthCheckShouldReturnDegradedWhenDeadLetterQueueIsNotEmpty() throws Exception { - createDeadLetterQueue(channel); - publishAMessage(channel); - publishAMessage(channel); - - awaitAtMostOneMinute.until(() -> testee.check().block().isDegraded()); - } - - @Test - void healthCheckShouldReturnUnhealthyWhenThereIsNoDeadLetterQueue() { - assertThat(testee.check().block().isUnHealthy()).isTrue(); - } - - private void createDeadLetterQueue(Channel channel) throws IOException { - channel.exchangeDeclare(EXCHANGE_NAME, DIRECT_EXCHANGE, DURABLE); - channel.queueDeclare(JAMES_MAIL_QUEUE_NAME.toDeadLetterQueueName(), DURABLE, !EXCLUSIVE, AUTO_DELETE, NO_QUEUE_DECLARE_ARGUMENTS).getQueue(); - channel.queueBind(JAMES_MAIL_QUEUE_NAME.toDeadLetterQueueName(), EXCHANGE_NAME, ROUTING_KEY); - } - - private void publishAMessage(Channel channel) throws IOException { - AMQP.BasicProperties basicProperties = new AMQP.BasicProperties.Builder() - .deliveryMode(PERSISTENT_TEXT_PLAIN.getDeliveryMode()) - .priority(PERSISTENT_TEXT_PLAIN.getPriority()) - .contentType(PERSISTENT_TEXT_PLAIN.getContentType()) - .build(); - - channel.basicPublish(EXCHANGE_NAME, ROUTING_KEY, basicProperties, "Hello, world!".getBytes(StandardCharsets.UTF_8)); - } - - private void closeQuietly(AutoCloseable... closeables) { - Arrays.stream(closeables).forEach(this::closeQuietly); - } - - private void closeQuietly(AutoCloseable closeable) { - try { - closeable.close(); - } catch (Exception e) { - //ignore error - } - } -} --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
