Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -147,6 +147,8 @@ public class ActiveMQConnection implements Connection, TopicConnection, QueueCon
* This strictly rejects non-standard property types such as Character, Map, and List.
*/
private boolean strictCompliance = false;
private boolean deferPrefetchUntilStarted = false;
private final AtomicBoolean everStarted = new AtomicBoolean(false);

private boolean disableTimeStampsByDefault;
private boolean optimizedMessageDispatch = true;
Expand Down Expand Up @@ -579,6 +581,7 @@ public void setClientInternalExceptionListener(ClientInternalExceptionListener l
public void start() throws JMSException {
checkClosedOrFailed();
ensureConnectionInfoSent();
everStarted.set(true);
if (started.compareAndSet(false, true)) {
for (Iterator<ActiveMQSession> i = sessions.iterator(); i.hasNext();) {
ActiveMQSession session = i.next();
Expand Down Expand Up @@ -1057,6 +1060,25 @@ public void setStrictCompliance(boolean strictCompliance) {
this.strictCompliance = strictCompliance;
}

public boolean isDeferPrefetchUntilStarted() {
return deferPrefetchUntilStarted;
}

/**
* See {@link ActiveMQConnectionFactory#setDeferPrefetchUntilStarted(boolean)}.
*/
public void setDeferPrefetchUntilStarted(boolean deferPrefetchUntilStarted) {
this.deferPrefetchUntilStarted = deferPrefetchUntilStarted;
}

/**
* @return true once {@link #start()} has been called at least once,
* regardless of later {@link #stop()} calls.
*/
public boolean isEverStarted() {
return everStarted.get();
}

public boolean isExclusiveConsumer() {
return exclusiveConsumer;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -131,6 +131,7 @@ public class ActiveMQConnectionFactory extends JNDIBaseStorable implements Conne
* This strictly rejects non-standard property types such as Character, Map, and List.
*/
private boolean strictCompliance = false;
private boolean deferPrefetchUntilStarted = false;

private boolean disableTimeStampsByDefault;
private boolean optimizedMessageDispatch = true;
Expand Down Expand Up @@ -428,6 +429,7 @@ protected ActiveMQConnection createActiveMQConnection(Transport transport, JMSSt

protected void configureConnection(ActiveMQConnection connection) throws JMSException {
connection.setPrefetchPolicy(getPrefetchPolicy());
connection.setDeferPrefetchUntilStarted(isDeferPrefetchUntilStarted());
connection.setDisableTimeStampsByDefault(isDisableTimeStampsByDefault());
connection.setOptimizedMessageDispatch(isOptimizedMessageDispatch());
connection.setCopyMessageOnSend(isCopyMessageOnSend());
Expand Down Expand Up @@ -1059,6 +1061,22 @@ public void setStrictCompliance(boolean strictCompliance) {
this.strictCompliance = strictCompliance;
}

public boolean isDeferPrefetchUntilStarted() {
return deferPrefetchUntilStarted;
}

/**
* When enabled, queue consumers created before their connection is first
* started register with a prefetch of zero so the broker does not dispatch
* messages into a consumer that cannot deliver them. The configured
* prefetch is restored when the connection starts. Default is false to
* preserve the historical behavior of buffering pre-start dispatches in
* the client.
*/
public void setDeferPrefetchUntilStarted(boolean deferPrefetchUntilStarted) {
this.deferPrefetchUntilStarted = deferPrefetchUntilStarted;
}

public String getClientIDPrefix() {
return clientIDPrefix;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -44,6 +44,7 @@
import org.apache.activemq.command.ActiveMQMessage;
import org.apache.activemq.command.ActiveMQObjectMessage;
import org.apache.activemq.command.ActiveMQTempDestination;
import org.apache.activemq.command.ConsumerControl;
import org.apache.activemq.command.CommandTypes;
import org.apache.activemq.command.ConsumerId;
import org.apache.activemq.command.ConsumerInfo;
Expand Down Expand Up @@ -154,6 +155,9 @@ class PreviouslyDelivered {
private final String selector;
private boolean synchronizationRegistered;
private final AtomicBoolean started = new AtomicBoolean(false);
// Configured prefetch withheld until the connection first starts; zero when
// no deferral is pending. See ActiveMQConnectionFactory#setDeferPrefetchUntilStarted.
private volatile int deferredPrefetchSize;

private MessageAvailableListener availableListener;

Expand Down Expand Up @@ -293,6 +297,21 @@ public ActiveMQMessageConsumer(ActiveMQSession session, ConsumerId consumerId, A
|| this.nonBlockingRedelivery
|| session.connection.isMessagePrioritySupported();
this.consumerExpiryCheckEnabled = session.connection.isConsumerExpiryCheckEnabled();
// Register queue consumers created before the connection has ever been
// started with zero prefetch so the broker does not push messages into a
// consumer that cannot deliver them; such messages would park in the held
// dispatch channel and starve running consumers. The configured prefetch
// is restored when the connection starts.
if (this.session.connection.isDeferPrefetchUntilStarted()
&& !this.session.connection.isEverStarted()
&& dest.isQueue()
&& !browser
&& this.info.getPrefetchSize() > 0) {
this.deferredPrefetchSize = this.info.getPrefetchSize();
this.info.setPrefetchSize(0);
this.info.setCurrentPrefetchSize(0);
}

if (messageListener != null) {
setMessageListener(messageListener);
}
Expand Down Expand Up @@ -448,7 +467,7 @@ public MessageListener getMessageListener() throws JMSException {
@Override
public void setMessageListener(MessageListener listener) throws JMSException {
checkClosed();
if (info.getPrefetchSize() == 0) {
if (info.getPrefetchSize() == 0 && deferredPrefetchSize == 0) {
throw new JMSException("Illegal prefetch size of zero. This setting is not supported for asynchronous consumers please set a value of at least 1");
}
if (listener != null) {
Expand Down Expand Up @@ -665,7 +684,7 @@ public Message receive(long timeout) throws JMSException {
while (timeout > 0) {

MessageDispatch md;
if (info.getPrefetchSize() == 0) {
if (isPullConsumer()) {
md = dequeue(-1); // We let the broker let us know when we timeout.
} else {
md = dequeue(timeout);
Expand Down Expand Up @@ -697,7 +716,7 @@ public Message receiveNoWait() throws JMSException {
sendPullCommand(-1);

MessageDispatch md;
if (info.getPrefetchSize() == 0) {
if (isPullConsumer()) {
md = dequeue(-1); // We let the broker let us know when we
// timeout.
} else {
Expand Down Expand Up @@ -914,7 +933,9 @@ protected void checkClosed() throws IllegalStateException {
*/
protected void sendPullCommand(long timeout) throws JMSException {
clearDeliveredList();
if (info.getCurrentPrefetchSize() == 0 && unconsumedMessages.isEmpty()) {
// A consumer whose prefetch is deferred until connection start is not a
// pull consumer; a pull here could steal a message into the held channel.
if (info.getCurrentPrefetchSize() == 0 && deferredPrefetchSize == 0 && unconsumedMessages.isEmpty()) {
MessagePull messagePull = new MessagePull();
messagePull.configure(info);
messagePull.setTimeout(timeout);
Expand Down Expand Up @@ -1485,7 +1506,7 @@ public void dispatch(MessageDispatch md) {

// Pull consumer needs to check if pull timed out and send
// a new pull command if not.
if (info.getCurrentPrefetchSize() == 0) {
if (info.getCurrentPrefetchSize() == 0 && deferredPrefetchSize == 0) {
unconsumedMessages.enqueue(null);
}
}
Expand Down Expand Up @@ -1628,11 +1649,41 @@ public void start() throws JMSException {
if (unconsumedMessages.isClosed()) {
return;
}
restoreDeferredPrefetch();
started.set(true);
unconsumedMessages.start();
session.executor.wakeup();
}

/**
* A pull consumer relies on the broker to answer each MessagePull, including
* signalling the receive timeout. A consumer whose prefetch is merely deferred
* until the connection starts carries prefetch zero on the wire but sends no
* pulls, so it must use client-side timeouts like any push consumer.
*/
private boolean isPullConsumer() {
return info.getPrefetchSize() == 0 && deferredPrefetchSize == 0;
}

/**
* Restores the configured prefetch that was withheld while the connection
* had never been started and tells the broker so it re-credits this
* consumer and dispatches any pending messages.
*/
private void restoreDeferredPrefetch() throws JMSException {
var prefetch = deferredPrefetchSize;
if (prefetch > 0) {
deferredPrefetchSize = 0;
info.setPrefetchSize(prefetch);
info.setCurrentPrefetchSize(prefetch);
var control = new ConsumerControl();
control.setConsumerId(info.getConsumerId());
control.setDestination(info.getDestination());
control.setPrefetch(prefetch);
session.asyncSendPacket(control);
}
}

public void stop() {
started.set(false);
unconsumedMessages.stop();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -97,6 +97,11 @@ public Context getInitialContext(final Hashtable<?, ?> environment) throws Namin
private static ActiveMQConnectionFactory createConnectionFactory(final String clientId) {
final ActiveMQConnectionFactory factory = new ActiveMQConnectionFactory(BROKER_URL);
factory.setNestedMapAndListEnabled(false);
// JmsTool creates competing queue consumers on connections it never starts.
// Without deferral the broker round-robins messages into those consumers'
// prefetch, where they park in the held dispatch channel and the started
// consumer's receive() blocks forever (core20 jmsconsumertests queueReceiveTests).
factory.setDeferPrefetchUntilStarted(true);
if (clientId != null) {
factory.setClientID(clientId);
}
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,170 @@
/**
* 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.activemq;

import static org.junit.Assert.assertNotNull;
import static org.junit.Assert.assertNull;
import static org.junit.Assert.assertTrue;

import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;

import jakarta.jms.Connection;
import jakarta.jms.Session;

import org.apache.activemq.broker.BrokerService;
import org.junit.After;
import org.junit.Before;
import org.junit.Test;

/**
* A queue consumer whose connection has never been started must not attract
* push dispatch. Historically the broker round-robins messages into such a
* consumer's prefetch, where they sit in the client's held dispatch channel
* and starve the consumers that are actually running (seen as the TCK
* queueReceiveTests hang).
*/
public class UnstartedConnectionQueueDispatchTest {

private static final int MESSAGE_COUNT = 10;

private BrokerService broker;
private String connectionUri;
private Connection startedConnection;
private Connection unstartedConnection;

@Before
public void setUp() throws Exception {
broker = new BrokerService();
broker.setPersistent(false);
broker.setUseJmx(false);
broker.setAdvisorySupport(false);
broker.setSchedulerSupport(false);
broker.addConnector("vm://localhost");
broker.start();
broker.waitUntilStarted();
connectionUri = "vm://localhost";
}

@After
public void tearDown() throws Exception {
if (startedConnection != null) {
try { startedConnection.close(); } catch (Exception ignored) {}
}
if (unstartedConnection != null) {
try { unstartedConnection.close(); } catch (Exception ignored) {}
}
if (broker != null) {
broker.stop();
broker.waitUntilStopped();
}
}

private ActiveMQConnectionFactory createFactory() {
var factory = new ActiveMQConnectionFactory(connectionUri);
factory.setDeferPrefetchUntilStarted(true);
return factory;
}

@Test(timeout = 60000)
public void testStartedConsumerReceivesAllMessagesDespiteUnstartedCompetitor() throws Exception {
var factory = createFactory();

// competing consumer on a connection that is never started
unstartedConnection = factory.createConnection();
var unstartedSession = unstartedConnection.createSession(false, Session.AUTO_ACKNOWLEDGE);
var queue = unstartedSession.createQueue("test.unstarted.dispatch");
var neverStarted = unstartedSession.createConsumer(queue);
assertNotNull(neverStarted);

// active consumer on a started connection
startedConnection = factory.createConnection();
startedConnection.start();
var session = startedConnection.createSession(false, Session.AUTO_ACKNOWLEDGE);
var active = session.createConsumer(queue);

var producer = session.createProducer(queue);
for (int i = 0; i < MESSAGE_COUNT; i++) {
producer.send(session.createTextMessage("message-" + i));
}

// every message must reach the running consumer; none may park in the
// never-started consumer's prefetch
for (int i = 0; i < MESSAGE_COUNT; i++) {
var received = active.receive(5000);
assertNotNull("Message " + i + " was dispatched to the never-started consumer", received);
}
}

@Test(timeout = 60000)
public void testDeferredConsumerReceivesMessageReleasedByClosedConsumer() throws Exception {
// TCK core/queueConnection connNotStartedQueueTest shape: a started
// receiver prefetches two messages and consumes one; closing it returns
// the other to the queue; a receiver on a never-started connection must
// see nothing until start, then receive the released message.
var factory = createFactory();

startedConnection = factory.createConnection();
startedConnection.start();
var session = startedConnection.createSession(false, Session.AUTO_ACKNOWLEDGE);
var queue = session.createQueue("test.unstarted.released");
var first = session.createConsumer(queue);
var producer = session.createProducer(queue);
producer.send(session.createTextMessage("one"));
producer.send(session.createTextMessage("two"));
assertNotNull(first.receive(5000));
first.close();

unstartedConnection = factory.createConnection();
var unstartedSession = unstartedConnection.createSession(false, Session.AUTO_ACKNOWLEDGE);
var second = unstartedSession.createConsumer(queue);
assertNull("No delivery before the connection is started", second.receive(1000));

unstartedConnection.start();
assertNotNull("Released message must be delivered once the connection starts", second.receive(5000));
}

@Test(timeout = 60000)
public void testDeferredConsumerReceivesAfterConnectionStart() throws Exception {
var factory = createFactory();

unstartedConnection = factory.createConnection();
var session = unstartedConnection.createSession(false, Session.AUTO_ACKNOWLEDGE);
var queue = session.createQueue("test.unstarted.recovery");
var consumer = session.createConsumer(queue);

// the async pattern: listener registered before the connection starts
final var delivered = new CountDownLatch(MESSAGE_COUNT);
consumer.setMessageListener(message -> delivered.countDown());

startedConnection = factory.createConnection();
startedConnection.start();
var producerSession = startedConnection.createSession(false, Session.AUTO_ACKNOWLEDGE);
var producer = producerSession.createProducer(queue);
for (int i = 0; i < MESSAGE_COUNT; i++) {
producer.send(producerSession.createTextMessage("message-" + i));
}

// nothing may be delivered while the connection is not started
assertTrue("Messages must not be delivered before start", delivered.getCount() == MESSAGE_COUNT);

// starting the connection restores the prefetch credit and delivery flows
unstartedConnection.start();
assertTrue("Messages should be delivered after connection start",
delivered.await(10, TimeUnit.SECONDS));
}
}
Loading