[jira] [Commented] (JAMES-4154) Generalize RabbitMQ work queue code used for async deletion callback(s)

"Matthieu Baechler (Jira)" <[email protected]>
Newsgroups gmane.comp.jakarta.james.devel
Message-ID <[email protected]>
    [ https://issues.apache.org/jira/browse/JAMES-4154?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=18043469#comment-18043469 ] 

Matthieu Baechler commented on JAMES-4154:
------------------------------------------

Reading the Event Bus ADR, it sounds like it's the right component for the goal of this issue.

 

> Generalize RabbitMQ work queue code used for async deletion callback(s)
> -----------------------------------------------------------------------
>
>                 Key: JAMES-4154
>                 URL: https://issues.apache.org/jira/browse/JAMES-4154
>             Project: James Server
>          Issue Type: Improvement
>          Components: deletedMessageVault, Queue
>    Affects Versions: master
>            Reporter: Tran Hong Quan
>            Priority: Minor
>
> h2. Why
> Today, we have the deleted message vault that handles deleted message callback asynchronously, relying on the RabbitMQ work queue to avoid timeout consuming e.g. DTM for deletion of a mailbox having 1 million emails.
> We may want to have other async deletion callback(s) that could rely on RabbiMQ work queue too.
> We should find a way to mutualize the shared RabbitMQ work queue setup, and plug different callback(s) easily.
> h2. How
> - Introduce `AsyncDeletionCallback` interface
> ```java
> public interface AsyncDeletionCallback extends DeleteMessageListener.DeletionCallback {
> }
> ```
> - Implement `AggregatedAsyncDeletionCallback` that uses the RabbitMQ work queue code similar to current`DistributedDeletedMessageVaultDeletionCallback`, but takes an `AsyncDeletionCallback` set as an argument.
>    The idea is to share the same RabbitMQ work queue e.g. `async-deletion-work-queue`, to handle all the async deletion callback(s).
> ```java
>     private Mono<Void> handleMessage(AcknowledgableDelivery delivery) {
>         try {
>             CopyCommandDTO copyCommandDTO = objectMapper.readValue(delivery.getBody(), CopyCommandDTO.class);
>             return Flux.fromIterable(asyncDeletionCallbacks)
>                 .flatMap(callback -> callback.forMessage(copyCommandDTO.asPojo(mailboxIdFactory, messageIdFactory, blobIdFactory)))
>                 .then()
>                 .timeout(Duration.ofMinutes(5))
>                 .doOnSuccess(any -> delivery.ack())
>                 .doOnCancel(() -> delivery.nack(REQUEUE))
>                 .onErrorResume(e -> {
>                     LOGGER.error("Failed executing async deletion callbacks for {}", copyCommandDTO.messageId, e);
>                     delivery.nack(REQUEUE);
>                     return Mono.empty();
>                 });
> ```
> - Refactor `DistributedDeletedMessageVaultDeletionCallback`: strip all the rabbitmq code, just implement `AsyncDeletionCallback`
> - Refactor `DeletedMessageVaultWorkQueueReconnectionHandler` so it handles reconnection for `AggregatedAsyncDeletionCallback` instead.
> - Guice binding to plug DTM as an `AsyncDeletionCallback`



--
This message was sent by Atlassian Jira
(v8.20.10#820010)
lmpx.com only provides a reader for public news (NNTP) servers. It is not affiliated with the servers or forums shown here and is not responsible for the content of articles, which is written by their respective authors.