git-svn-id: https://svn.apache.org/repos/asf/activemq/trunk@1195278 13f79535-47bb-0310-9956-ffa450edef68
This commit is contained in:
Timothy A. Bish 2011-10-30 23:27:52 +00:00
parent 15f641634b
commit 0a66b79332
4 changed files with 104 additions and 20 deletions

View File

@ -670,6 +670,9 @@ public class ActiveMQMessageConsumer implements MessageAvailableConsumer, StatsC
} }
void doClose() throws JMSException { void doClose() throws JMSException {
// Store interrupted state and clear so that Transport operations don't
// throw InterruptedException and we ensure that resources are clened up.
boolean interrupted = Thread.interrupted();
dispose(); dispose();
RemoveInfo removeCommand = info.createRemoveCommand(); RemoveInfo removeCommand = info.createRemoveCommand();
if (LOG.isDebugEnabled()) { if (LOG.isDebugEnabled()) {
@ -677,7 +680,9 @@ public class ActiveMQMessageConsumer implements MessageAvailableConsumer, StatsC
} }
removeCommand.setLastDeliveredSequenceId(lastDeliveredSequenceId); removeCommand.setLastDeliveredSequenceId(lastDeliveredSequenceId);
this.session.asyncSendPacket(removeCommand); this.session.asyncSendPacket(removeCommand);
} if (interrupted) {
Thread.currentThread().interrupt();
} }
void inProgressClearRequired() { void inProgressClearRequired() {
inProgressClearRequiredFlag = true; inProgressClearRequiredFlag = true;

View File

@ -637,10 +637,14 @@ public class ActiveMQSession implements Session, QueueSession, TopicSession, Sta
} }
private void doClose() throws JMSException { private void doClose() throws JMSException {
boolean interrupted = Thread.interrupted();
dispose(); dispose();
RemoveInfo removeCommand = info.createRemoveCommand(); RemoveInfo removeCommand = info.createRemoveCommand();
removeCommand.setLastDeliveredSequenceId(lastDeliveredSequenceId); removeCommand.setLastDeliveredSequenceId(lastDeliveredSequenceId);
connection.asyncSendPacket(removeCommand); connection.asyncSendPacket(removeCommand);
if (interrupted) {
Thread.currentThread().interrupt();
}
} }
void clearMessagesInProgress() { void clearMessagesInProgress() {
@ -1963,7 +1967,7 @@ public class ActiveMQSession implements Session, QueueSession, TopicSession, Sta
this.blobTransferPolicy = blobTransferPolicy; this.blobTransferPolicy = blobTransferPolicy;
} }
public List getUnconsumedMessages() { public List<MessageDispatch> getUnconsumedMessages() {
return executor.getUnconsumedMessages(); return executor.getUnconsumedMessages();
} }

View File

@ -31,7 +31,6 @@ import org.slf4j.LoggerFactory;
* A utility class used by the Session for dispatching messages asynchronously * A utility class used by the Session for dispatching messages asynchronously
* to consumers * to consumers
* *
*
* @see javax.jms.Session * @see javax.jms.Session
*/ */
public class ActiveMQSessionExecutor implements Task { public class ActiveMQSessionExecutor implements Task {
@ -125,9 +124,7 @@ public class ActiveMQSessionExecutor implements Task {
} }
void dispatch(MessageDispatch message) { void dispatch(MessageDispatch message) {
// TODO - we should use a Map for this indexed by consumerId // TODO - we should use a Map for this indexed by consumerId
for (ActiveMQMessageConsumer consumer : this.session.consumers) { for (ActiveMQMessageConsumer consumer : this.session.consumers) {
ConsumerId consumerId = message.getConsumerId(); ConsumerId consumerId = message.getConsumerId();
if (consumerId.equals(consumer.getConsumerId())) { if (consumerId.equals(consumer.getConsumerId())) {
@ -207,8 +204,7 @@ public class ActiveMQSessionExecutor implements Task {
} }
} }
List getUnconsumedMessages() { List<MessageDispatch> getUnconsumedMessages() {
return messageQueue.removeAll(); return messageQueue.removeAll();
} }
} }

View File

@ -13,29 +13,45 @@
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and * See the License for the specific language governing permissions and
* limitations under the License. * limitations under the License.
*
*/ */
package org.apache.activemq.bugs; package org.apache.activemq.bugs;
import static org.junit.Assert.*; import java.util.Properties;
import javax.jms.Connection; import javax.jms.Connection;
import javax.jms.ConnectionFactory; import javax.jms.ConnectionFactory;
import javax.jms.Destination;
import javax.jms.JMSException; import javax.jms.JMSException;
import javax.jms.MessageConsumer;
import javax.jms.Session; import javax.jms.Session;
import javax.naming.Context;
import javax.naming.InitialContext;
import javax.naming.NamingException;
import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertNotNull;
import static org.junit.Assert.fail;
import org.apache.activemq.ActiveMQConnectionFactory; import org.apache.activemq.ActiveMQConnectionFactory;
import org.apache.activemq.broker.BrokerService; import org.apache.activemq.broker.BrokerService;
import org.junit.After; import org.junit.After;
import org.junit.Before; import org.junit.Before;
import org.junit.Test; import org.junit.Test;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
public class AMQ3529Test { public class AMQ3529Test {
private static Logger LOG = LoggerFactory.getLogger(AMQ3529Test.class);
private ConnectionFactory connectionFactory; private ConnectionFactory connectionFactory;
private Connection connection; private Connection connection;
private Session session; private Session session;
private BrokerService broker; private BrokerService broker;
private String connectionUri; private String connectionUri;
private MessageConsumer consumer;
private Context ctx = null;
@Before @Before
public void startBroker() throws Exception { public void startBroker() throws Exception {
@ -72,20 +88,52 @@ public class AMQ3529Test {
connection = connectionFactory.createConnection(); connection = connectionFactory.createConnection();
session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE); session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE);
assertNotNull(session); assertNotNull(session);
} catch (JMSException e) {
fail(e.getMessage()); Properties props = new Properties();
props.setProperty(Context.INITIAL_CONTEXT_FACTORY, "org.apache.activemq.jndi.ActiveMQInitialContextFactory");
props.setProperty(Context.PROVIDER_URL, "tcp://0.0.0.0:0");
ctx = null;
try {
ctx = new InitialContext(props);
} catch (NoClassDefFoundError e) {
throw new NamingException(e.toString());
} catch (Exception e) {
throw new NamingException(e.toString());
}
Destination destination = (Destination) ctx.lookup("dynamicTopics/example.C");
consumer = session.createConsumer(destination);
consumer.receive(10000);
} catch (Exception e) {
// Expect an exception here from the interrupt.
} finally {
// next line is the nature of the test, if I remove this
// line, everything works OK
try {
consumer.close();
} catch (JMSException e) {
fail("Consumer Close failed with" + e.getMessage());
}
try {
session.close();
} catch (JMSException e) {
fail("Session Close failed with" + e.getMessage());
} }
// next line is the nature of the test, if I remove this line, everything works OK
Thread.currentThread().interrupt();
try { try {
connection.close(); connection.close();
} catch (JMSException e) { } catch (JMSException e) {
fail("Connection Close failed with" + e.getMessage());
}
try {
ctx.close();
} catch (Exception e) {
fail("Connection Close failed with" + e.getMessage());
}
} }
assertTrue(Thread.currentThread().isInterrupted());
} }
}; };
client.start(); client.start();
Thread.sleep(5000);
client.interrupt();
client.join(); client.join();
Thread.sleep(2000); Thread.sleep(2000);
Thread[] remainThreads = new Thread[tg.activeCount()]; Thread[] remainThreads = new Thread[tg.activeCount()];
@ -94,5 +142,36 @@ public class AMQ3529Test {
if (t.isAlive() && !t.isDaemon()) if (t.isAlive() && !t.isDaemon())
fail("Remaining thread: " + t.toString()); fail("Remaining thread: " + t.toString());
} }
ThreadGroup root = Thread.currentThread().getThreadGroup().getParent();
while (root.getParent() != null) {
root = root.getParent();
}
visit(root, 0);
}
// This method recursively visits all thread groups under `group'.
public static void visit(ThreadGroup group, int level) {
// Get threads in `group'
int numThreads = group.activeCount();
Thread[] threads = new Thread[numThreads * 2];
numThreads = group.enumerate(threads, false);
// Enumerate each thread in `group'
for (int i = 0; i < numThreads; i++) {
// Get thread
Thread thread = threads[i];
LOG.debug("Thread:" + thread.getName() + " is still running");
}
// Get thread subgroups of `group'
int numGroups = group.activeGroupCount();
ThreadGroup[] groups = new ThreadGroup[numGroups * 2];
numGroups = group.enumerate(groups, false);
// Recursively visit each subgroup
for (int i = 0; i < numGroups; i++) {
visit(groups[i], level + 1);
}
} }
} }