Hi guys!
Looks like i found solution for
https://issues.apache.org/jira/projects/JAMES/issues/JAMES-4192

We can use
Artemis _AMQ_SCHED_DELIVERY

If you think i found it right - tell me i can add this to
https://github.com/apache/james-project/pull/3194
or not?
================================================================================
JAMES-4192: Resolving the Delayed Mail Paging Deadlock in Apache ActiveMQ 
Artemis
================================================================================

1. PROBLEM ROOT CAUSE ANALYSIS
--------------------------------------------------------------------------------
In Apache James `JMSCacheableMailQueue`, message delay scheduling is currently
implemented using a custom message property `JAMES_NEXT_DELIVERY` combined with 
a
JMS Message Selector on the consumer side:

    String selector = "JAMES_NEXT_DELIVERY <= " + System.currentTimeMillis()
                    + " OR FORCE_DELIVERY = true";
    consumer = session.createConsumer(queue, selector);

Why this causes head-of-line blocking and deadlocks in message brokers:

- In ActiveMQ Classic:
  The broker scans up to `maxPageSize` (default: 200) messages into memory.
  If all 200 scanned messages are delayed retries, the selector rejects them 
all.
  The broker stops searching deeper into the queue, leaving ready messages stuck
  behind delayed ones (head-of-line blocking).

- In ActiveMQ Artemis:
  Artemis only evaluates JMS message selectors against messages held in memory,
  NOT against messages paged out to disk. Once message paging begins, memory
  fills up with delayed messages rejected by `JAMES_NEXT_DELIVERY <= now`.
  Because no messages are consumed, no space is freed to page in subsequent
  messages from disk. Ready messages residing in page files stay stuck 
indefinitely.


2. THE SOLUTION FOR ARTEMIS: NATIVE SCHEDULED DELIVERY (JMS 2.0)
--------------------------------------------------------------------------------
Artemis has a dedicated native scheduled message delivery engine built into its 
core
(`ScheduledDeliveryHandler`). Instead of keeping delayed messages in the visible
queue and relying on consumer selectors:

1. Messages are enqueued with a delivery delay via JMS 2.0 
`producer.setDeliveryDelay(ms)`
   or native Artemis property `_AMQ_SCHED_DELIVERY`.
2. Artemis intercepts these messages and places them in a dedicated priority 
scheduler
   tree. They are NOT placed in the visible queue or page files, and do NOT 
block consumers.
3. When the scheduled timestamp arrives, Artemis automatically moves the 
message into
   the ready queue.
4. The consumer reads from the queue with NO SELECTOR 
(`session.createConsumer(queue)`).
   It consumes messages at full wire speed without any paging or evaluation 
bottlenecks.


3. CODE IMPLEMENTATION EXAMPLE
--------------------------------------------------------------------------------

Below is an example showing how `ArtemisMailQueue` replaces the selector-based
pattern with native scheduled delivery.

```java
package org.apache.james.queue.artemis;

import java.time.Duration;
import java.util.Map;

import jakarta.jms.Connection;
import jakarta.jms.ConnectionFactory;
import jakarta.jms.JMSException;
import jakarta.jms.Message;
import jakarta.jms.MessageConsumer;
import jakarta.jms.MessageProducer;
import jakarta.jms.Queue;
import jakarta.jms.Session;

import org.apache.james.queue.api.MailQueue;
import org.apache.james.queue.api.MailQueueItem;
import org.apache.james.queue.api.MailQueueName;
import org.apache.james.queue.jms.JMSCacheableMailQueue;
import org.apache.mailet.Mail;

public class ArtemisMailQueue extends JMSCacheableMailQueue {

    /**
     * Artemis native scheduled delivery header key.
     * Value is absolute epoch timestamp in milliseconds when message should be 
delivered.
     * Using this is equivalent to JMS 2.0 setDeliveryDelay().
     */
    public static final String AMQ_SCHEDULED_DELIVERY = "_AMQ_SCHED_DELIVERY";

    private final Connection connection;
    private final MailQueueName queueName;

    public ArtemisMailQueue(ConnectionFactory connectionFactory,
                            MailQueueName queueName) throws JMSException {
        super(connectionFactory, queueName);
        this.connection = connectionFactory.createConnection();
        this.queueName = queueName;
    }

    /**
     * Enqueue a mail with an optional delay.
     *
     * Instead of setting JAMES_NEXT_DELIVERY property for consumer-side 
selector filtering,
     * we configure native Artemis delivery delay.
     */
    @Override
    public void enQueue(Mail mail, Duration delay) throws MailQueueException {
        try (Session session = connection.createSession(false, 
Session.AUTO_ACKNOWLEDGE)) {
            Queue queue = session.createQueue(queueName.asString());
            try (MessageProducer producer = session.createProducer(queue)) {

                long delayMillis = (delay != null && !delay.isNegative()) ? 
delay.toMillis() : 0L;

                if (delayMillis > 0) {
                    // Option A: Standard JMS 2.0 / Jakarta Messaging API
                    producer.setDeliveryDelay(delayMillis);

                    // Option B: Artemis-specific property (useful when 
producer pooling or compatibility wrapper is used)
                    // long deliverAt = System.currentTimeMillis() + 
delayMillis;
                    // message.setLongProperty(AMQ_SCHEDULED_DELIVERY, 
deliverAt);
                }

                Map<String, Object> props = getJMSProperties(mail, 0L);
                Message message = createJMSMessage(session, mail, props);

                producer.send(message);
            }
        } catch (Exception e) {
            throw new MailQueueException("Unable to enqueue mail " + 
mail.getName(), e);
        }
    }

    /**
     * Dequeue ready messages.
     *
     * CRITICAL FIX FOR JAMES-4192:
     * We DO NOT provide any message selector here!
     *
     * - The consumer receives messages directly without any selector filter.
     * - Delayed messages are held by Artemis's ScheduledDeliveryHandler and 
only
     *   released to the queue when their scheduled time arrives.
     * - Head-of-line blocking and memory paging deadlocks are completely 
eliminated.
     */
    @Override
    protected Mono<MailQueueItem> deQueueOneItem() {
        Session session = null;
        MessageConsumer consumer = null;
        try {
            session = connection.createSession(true, 
Session.SESSION_TRANSACTED);
            Queue queue = session.createQueue(queueName.asString());

            // NO SELECTOR USED HERE:
            // consumer = session.createConsumer(queue);
            // Rather than:
            // consumer = session.createConsumer(queue, getMessageSelector());
            consumer = session.createConsumer(queue);

            Message message = consumer.receive(10000);

            if (message != null) {
                return createMailQueueItem(session, consumer, message);
            } else {
                session.commit();
                closeConsumer(consumer);
                closeSession(session);
            }
        } catch (Exception e) {
            rollback(session);
            closeConsumer(consumer);
            closeSession(session);
            return Mono.error(new MailQueueException("Unable to dequeue next 
message", e));
        }
        return Mono.empty();
    }

    /**
     * Overridden to return null so that any base class callers do not inject
     * the obsolete JAMES_NEXT_DELIVERY selector.
     */
    @Override
    protected String getMessageSelector() {
        return null;
    }
}
```


4. WHY THIS DEFINITIVELY FIXES JAMES-4192
--------------------------------------------------------------------------------
1. No Selector Evaluation in Paging:
   Because the consumer selector is removed entirely, Artemis does not need
   to evaluate SQL selector expressions against memory or disk buffers.

2. Scheduled Messages Stored in Scheduler Index, Not in Queue Pages:
   In Artemis, messages with `_AMQ_SCHED_DELIVERY` or `deliveryDelay` are held
   in an in-memory red-black tree (or scheduler journal) and are not considered
   part of the active queue until triggered. Thus, they cannot saturate
   paging memory or block subsequent non-delayed emails.

3. Paging Works as Designed:
   When active ready messages exceed memory thresholds, Artemis pages them to 
disk.
   As consumers pull messages without filter rejections, new messages are read
   smoothly from page files in sequential order.

4. Performance & Simplicity:
   Compared to Apache Pulsar (which requires ZooKeeper, BookKeeper, separate
   requeue topics and filter actors), the Artemis native scheduled delivery
   leverages built-in JMS 2.0 semantics with zero extra infrastructure.
================================================================================
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to