[jboss-cvs] JBoss Messaging SVN: r5433 - in trunk: src/main/org/jboss/messaging/core/config and 10 other directories.

jboss-cvs-commits at lists.jboss.org jboss-cvs-commits at lists.jboss.org
Tue Nov 25 07:51:55 EST 2008


Author: ataylor
Date: 2008-11-25 07:51:55 -0500 (Tue, 25 Nov 2008)
New Revision: 5433

Added:
   trunk/tests/src/org/jboss/messaging/tests/integration/queue/ExpiryRunnerTest.java
Modified:
   trunk/src/config/jbm-configuration.xml
   trunk/src/main/org/jboss/messaging/core/config/Configuration.java
   trunk/src/main/org/jboss/messaging/core/config/impl/ConfigurationImpl.java
   trunk/src/main/org/jboss/messaging/core/config/impl/FileConfiguration.java
   trunk/src/main/org/jboss/messaging/core/postoffice/PostOffice.java
   trunk/src/main/org/jboss/messaging/core/postoffice/impl/PostOfficeImpl.java
   trunk/src/main/org/jboss/messaging/core/server/Queue.java
   trunk/src/main/org/jboss/messaging/core/server/impl/MessagingServerImpl.java
   trunk/src/main/org/jboss/messaging/core/server/impl/QueueImpl.java
   trunk/src/main/org/jboss/messaging/util/JBMThreadFactory.java
   trunk/tests/config/ConfigurationTest-config.xml
   trunk/tests/src/org/jboss/messaging/tests/performance/persistence/FakePostOffice.java
   trunk/tests/src/org/jboss/messaging/tests/unit/core/config/impl/ConfigurationImplTest.java
   trunk/tests/src/org/jboss/messaging/tests/unit/core/config/impl/FileConfigurationTest.java
Log:
added an expired message reaper

Modified: trunk/src/config/jbm-configuration.xml
===================================================================
--- trunk/src/config/jbm-configuration.xml	2008-11-25 12:18:19 UTC (rev 5432)
+++ trunk/src/config/jbm-configuration.xml	2008-11-25 12:51:55 UTC (rev 5433)
@@ -65,6 +65,11 @@
       <transaction-timeout>60000</transaction-timeout>
       <!--how often to scan for timedout transactions-->
       <transaction-timeout-scan-period>1000</transaction-timeout-scan-period>
+
+      <!-- how often do we scan the queues for expired messages-->
+      <message-expiry-scan-period>30000</message-expiry-scan-period>
+      <!-- the priority of the thread that expires th emessages (between 1 - 10 inclusive)-->
+      <message-expiry-thread-priority>3</message-expiry-thread-priority>
       
       <!-- Example interceptors 
       <remoting-interceptors>

Modified: trunk/src/main/org/jboss/messaging/core/config/Configuration.java
===================================================================
--- trunk/src/main/org/jboss/messaging/core/config/Configuration.java	2008-11-25 12:18:19 UTC (rev 5432)
+++ trunk/src/main/org/jboss/messaging/core/config/Configuration.java	2008-11-25 12:51:55 UTC (rev 5433)
@@ -197,4 +197,12 @@
    long getTransactionTimeoutScanPeriod();
 
    void setTransactionTimeoutScanPeriod(long period);
+
+   long getMessageExpiryScanPeriod();
+
+   void setMessageExpiryScanPeriod(long messageExpiryScanPeriod);
+
+   int getMessageExpiryThreadPriority();
+
+   void setMessageExpiryThreadPriority(int messageExpiryThreadPriority);
 }

Modified: trunk/src/main/org/jboss/messaging/core/config/impl/ConfigurationImpl.java
===================================================================
--- trunk/src/main/org/jboss/messaging/core/config/impl/ConfigurationImpl.java	2008-11-25 12:18:19 UTC (rev 5432)
+++ trunk/src/main/org/jboss/messaging/core/config/impl/ConfigurationImpl.java	2008-11-25 12:51:55 UTC (rev 5433)
@@ -103,6 +103,10 @@
 
    public static final long DEFAULT_MAX_FORWARD_BATCH_TIME = -1;
 
+   public static final long DEFAULT_MESSAGE_EXPIRY_SCAN_PERIOD = 30000;
+
+   public static final int DEFAULT_MESSAGE_EXPIRY_THREAD_PRIORITY = 3;
+
    // Attributes -----------------------------------------------------------------------------
 
    protected boolean clustered = DEFAULT_CLUSTERED;
@@ -123,6 +127,10 @@
 
    protected long connectionScanPeriod = DEFAULT_CONNECTION_SCAN_PERIOD;
 
+   protected long messageExpiryScanPeriod = DEFAULT_MESSAGE_EXPIRY_SCAN_PERIOD;
+
+   protected int messageExpiryThreadPriority = DEFAULT_MESSAGE_EXPIRY_THREAD_PRIORITY;
+
    protected List<String> interceptorClassNames = new ArrayList<String>();
    
    protected Map<String, TransportConfiguration> connectorConfigs = new HashMap<String, TransportConfiguration>();
@@ -465,6 +473,26 @@
       transactionTimeoutScanPeriod = period;
    }
 
+   public long getMessageExpiryScanPeriod()
+   {
+      return messageExpiryScanPeriod;
+   }
+
+   public void setMessageExpiryScanPeriod(long messageExpiryScanPeriod)
+   {
+      this.messageExpiryScanPeriod = messageExpiryScanPeriod;
+   }
+
+   public int getMessageExpiryThreadPriority()
+   {
+      return messageExpiryThreadPriority;
+   }
+
+   public void setMessageExpiryThreadPriority(int messageExpiryThreadPriority)
+   {
+      this.messageExpiryThreadPriority = messageExpiryThreadPriority;
+   }
+
    public boolean isSecurityEnabled()
    {
       return securityEnabled;

Modified: trunk/src/main/org/jboss/messaging/core/config/impl/FileConfiguration.java
===================================================================
--- trunk/src/main/org/jboss/messaging/core/config/impl/FileConfiguration.java	2008-11-25 12:18:19 UTC (rev 5432)
+++ trunk/src/main/org/jboss/messaging/core/config/impl/FileConfiguration.java	2008-11-25 12:51:55 UTC (rev 5433)
@@ -98,6 +98,10 @@
 
       transactionTimeoutScanPeriod = getLong(e, "transaction-timeout-scan-period", transactionTimeoutScanPeriod);
 
+      messageExpiryScanPeriod = getLong(e, "message-expiry-scan-period", messageExpiryScanPeriod);
+
+      messageExpiryThreadPriority = getInteger(e, "message-expiry-thread-priority", messageExpiryThreadPriority);
+
       managementAddress = new SimpleString(getString(e, "management-address", managementAddress.toString()));
 
       managementNotificationAddress = new SimpleString(getString(e, "management-notification-address", managementNotificationAddress.toString()));

Modified: trunk/src/main/org/jboss/messaging/core/postoffice/PostOffice.java
===================================================================
--- trunk/src/main/org/jboss/messaging/core/postoffice/PostOffice.java	2008-11-25 12:18:19 UTC (rev 5432)
+++ trunk/src/main/org/jboss/messaging/core/postoffice/PostOffice.java	2008-11-25 12:51:55 UTC (rev 5433)
@@ -78,6 +78,8 @@
    Set<SimpleString> listAllDestinations();
    
    List<Queue> activate();
+
+   List<Queue> getQueues();
    
    PagingManager getPagingManager();
    

Modified: trunk/src/main/org/jboss/messaging/core/postoffice/impl/PostOfficeImpl.java
===================================================================
--- trunk/src/main/org/jboss/messaging/core/postoffice/impl/PostOfficeImpl.java	2008-11-25 12:18:19 UTC (rev 5432)
+++ trunk/src/main/org/jboss/messaging/core/postoffice/impl/PostOfficeImpl.java	2008-11-25 12:51:55 UTC (rev 5433)
@@ -375,6 +375,21 @@
       return queues;
    }
 
+   public List<Queue> getQueues()
+   {
+      Map<SimpleString, Binding> nameMap = addressManager.getBindings();
+
+      List<Queue> queues = new ArrayList<Queue>();
+
+      for (Binding binding : nameMap.values())
+      {
+         Queue queue = binding.getQueue();
+         queues.add(queue);
+      }
+
+      return queues;
+   }
+
    public synchronized SendLock getAddressLock(final SimpleString address)
    {
       SendLock lock = addressLocks.get(address);

Modified: trunk/src/main/org/jboss/messaging/core/server/Queue.java
===================================================================
--- trunk/src/main/org/jboss/messaging/core/server/Queue.java	2008-11-25 12:18:19 UTC (rev 5432)
+++ trunk/src/main/org/jboss/messaging/core/server/Queue.java	2008-11-25 12:51:55 UTC (rev 5433)
@@ -118,6 +118,10 @@
          HierarchicalRepository<QueueSettings> queueSettingsRepository)
          throws Exception;
 
+   void expireMessages(final StorageManager storageManager,
+                                final PostOffice postOffice,
+                                final HierarchicalRepository<QueueSettings> queueSettingsRepository) throws Exception;
+
    boolean sendMessageToDeadLetterAddress(long messageID, StorageManager storageManager,
          PostOffice postOffice,
          HierarchicalRepository<QueueSettings> queueSettingsRepository)

Modified: trunk/src/main/org/jboss/messaging/core/server/impl/MessagingServerImpl.java
===================================================================
--- trunk/src/main/org/jboss/messaging/core/server/impl/MessagingServerImpl.java	2008-11-25 12:18:19 UTC (rev 5432)
+++ trunk/src/main/org/jboss/messaging/core/server/impl/MessagingServerImpl.java	2008-11-25 12:51:55 UTC (rev 5433)
@@ -16,6 +16,7 @@
 import java.util.List;
 import java.util.Map;
 import java.util.Set;
+import java.util.Collection;
 import java.util.concurrent.ConcurrentHashMap;
 import java.util.concurrent.ExecutorService;
 import java.util.concurrent.Executors;
@@ -39,6 +40,7 @@
 import org.jboss.messaging.core.paging.impl.PagingManagerImpl;
 import org.jboss.messaging.core.persistence.StorageManager;
 import org.jboss.messaging.core.postoffice.PostOffice;
+import org.jboss.messaging.core.postoffice.Binding;
 import org.jboss.messaging.core.postoffice.impl.PostOfficeImpl;
 import org.jboss.messaging.core.remoting.Channel;
 import org.jboss.messaging.core.remoting.ChannelHandler;
@@ -57,6 +59,7 @@
 import org.jboss.messaging.core.server.Queue;
 import org.jboss.messaging.core.server.QueueFactory;
 import org.jboss.messaging.core.server.ServerSession;
+import org.jboss.messaging.core.server.MessageReference;
 import org.jboss.messaging.core.server.cluster.ClusterManager;
 import org.jboss.messaging.core.server.cluster.impl.ClusterManagerImpl;
 import org.jboss.messaging.core.settings.HierarchicalRepository;
@@ -136,6 +139,8 @@
 
    private ConnectionManager replicatingConnectionManager;
 
+   private ScheduledThreadPoolExecutor messageExpiryExecutor;
+
    // Constructors
    // ---------------------------------------------------------------------------------
 
@@ -243,6 +248,10 @@
                                                           this);
 
       postOffice.start();
+      MessageExpiryRunner messageExpiryRunner = new MessageExpiryRunner(postOffice);
+      messageExpiryRunner.setPriority(3);
+      messageExpiryExecutor = new ScheduledThreadPoolExecutor(1, new JBMThreadFactory("JBM-scheduled-threads", configuration.getMessageExpiryThreadPriority()) );
+      messageExpiryExecutor.scheduleAtFixedRate(messageExpiryRunner, configuration.getMessageExpiryScanPeriod(), configuration.getMessageExpiryScanPeriod(), TimeUnit.MILLISECONDS);
       resourceManager.start();
 
       // FIXME the destination corresponding to the notification address is always created
@@ -336,6 +345,7 @@
       securityStore = null;
       resourceManager.stop();
       resourceManager = null;
+      messageExpiryExecutor.shutdown();
       postOffice.stop();
       postOffice = null;
       securityRepository = null;
@@ -724,4 +734,30 @@
          queue.activateNow(asyncDeliveryPool);
       }
    }
+
+   private class MessageExpiryRunner extends Thread
+   {
+      private final PostOffice postOffice;
+
+      public MessageExpiryRunner(PostOffice postOffice)
+      {
+         this.postOffice = postOffice;
+      }
+
+      public void run()
+      {
+         List<Queue> queues = postOffice.getQueues();
+         for (Queue queue : queues)
+         {
+            try
+            {
+               queue.expireMessages(storageManager, postOffice, queueSettingsRepository);
+            }
+            catch (Exception e)
+            {
+               log.error("failed to expire messages for queue " + queue.getName(), e);
+            }
+         }
+      }
+   }
 }

Modified: trunk/src/main/org/jboss/messaging/core/server/impl/QueueImpl.java
===================================================================
--- trunk/src/main/org/jboss/messaging/core/server/impl/QueueImpl.java	2008-11-25 12:18:19 UTC (rev 5432)
+++ trunk/src/main/org/jboss/messaging/core/server/impl/QueueImpl.java	2008-11-25 12:51:55 UTC (rev 5433)
@@ -450,7 +450,7 @@
       return deleted;
    }
 
-   public boolean expireMessage(final long messageID,
+   public synchronized boolean expireMessage(final long messageID,
                                 final StorageManager storageManager,
                                 final PostOffice postOffice,
                                 final HierarchicalRepository<QueueSettings> queueSettingsRepository) throws Exception
@@ -471,6 +471,27 @@
       return false;
    }
 
+   /**
+    * todo- at present we need the whole method syncronized until the iterator is fail safe.
+    */
+   public synchronized void expireMessages(final StorageManager storageManager,
+                                final PostOffice postOffice,
+                                final HierarchicalRepository<QueueSettings> queueSettingsRepository) throws Exception
+   {
+      Iterator<MessageReference> iter = messageReferences.iterator();
+
+      while (iter.hasNext())
+      {
+         MessageReference ref = iter.next();
+         if (ref.getMessage().isExpired())
+         {
+            deliveringCount.incrementAndGet();
+            ref.expire(storageManager, postOffice, queueSettingsRepository);
+            iter.remove();
+         }
+      }
+   }
+
    public boolean sendMessageToDeadLetterAddress(final long messageID,
                                    final StorageManager storageManager,
                                    final PostOffice postOffice,

Modified: trunk/src/main/org/jboss/messaging/util/JBMThreadFactory.java
===================================================================
--- trunk/src/main/org/jboss/messaging/util/JBMThreadFactory.java	2008-11-25 12:18:19 UTC (rev 5432)
+++ trunk/src/main/org/jboss/messaging/util/JBMThreadFactory.java	2008-11-25 12:51:55 UTC (rev 5433)
@@ -37,9 +37,16 @@
 
    private final AtomicInteger threadCount = new AtomicInteger(0);
 
+   private final int threadPriority;
    public JBMThreadFactory(final String groupName)
    {
+      this(groupName, Thread.NORM_PRIORITY);
+   }
+
+   public JBMThreadFactory(String groupName, int priority)
+   {
       group = new ThreadGroup(groupName + "-" + System.identityHashCode(this));
+      threadPriority = priority;
    }
 
    public Thread newThread(final Runnable command)
@@ -51,7 +58,7 @@
 
       // Don't want to prevent VM from exiting
       t.setDaemon(true);
-
+      t.setPriority(threadPriority);
       return t;
    }
 }

Modified: trunk/tests/config/ConfigurationTest-config.xml
===================================================================
--- trunk/tests/config/ConfigurationTest-config.xml	2008-11-25 12:18:19 UTC (rev 5432)
+++ trunk/tests/config/ConfigurationTest-config.xml	2008-11-25 12:51:55 UTC (rev 5433)
@@ -13,6 +13,8 @@
       <connection-scan-period>6543</connection-scan-period>
       <transaction-timeout>98765</transaction-timeout>
       <transaction-timeout-scan-period>56789</transaction-timeout-scan-period>
+      <message-expiry-scan-period>10111213</message-expiry-scan-period>
+      <message-expiry-thread-priority>8</message-expiry-thread-priority>
       <management-address>Giraffe</management-address>
       <remoting-interceptors>
          <class-name>org.jboss.messaging.tests.unit.core.config.impl.TestInterceptor1</class-name>

Added: trunk/tests/src/org/jboss/messaging/tests/integration/queue/ExpiryRunnerTest.java
===================================================================
--- trunk/tests/src/org/jboss/messaging/tests/integration/queue/ExpiryRunnerTest.java	                        (rev 0)
+++ trunk/tests/src/org/jboss/messaging/tests/integration/queue/ExpiryRunnerTest.java	2008-11-25 12:51:55 UTC (rev 5433)
@@ -0,0 +1,348 @@
+/*
+ * JBoss, Home of Professional Open Source
+ * Copyright 2005-2008, Red Hat Middleware LLC, and individual contributors
+ * by the @authors tag. See the copyright.txt in the distribution for a
+ * full listing of individual contributors.
+ *
+ * This is free software; you can redistribute it and/or modify it
+ * under the terms of the GNU Lesser General Public License as
+ * published by the Free Software Foundation; either version 2.1 of
+ * the License, or (at your option) any later version.
+ *
+ * This software is distributed in the hope that it will be useful,
+ * but WITHOUT ANY WARRANTY; without even the implied warranty of
+ * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the GNU
+ * Lesser General Public License for more details.
+ *
+ * You should have received a copy of the GNU Lesser General Public
+ * License along with this software; if not, write to the Free
+ * Software Foundation, Inc., 51 Franklin St, Fifth Floor, Boston, MA
+ * 02110-1301 USA, or see the FSF site: http://www.fsf.org.
+ */
+package org.jboss.messaging.tests.integration.queue;
+
+import org.jboss.messaging.tests.util.UnitTestCase;
+import org.jboss.messaging.core.server.MessagingService;
+import org.jboss.messaging.core.server.impl.MessagingServiceImpl;
+import org.jboss.messaging.core.client.ClientSession;
+import org.jboss.messaging.core.client.ClientSessionFactory;
+import org.jboss.messaging.core.client.ClientProducer;
+import org.jboss.messaging.core.client.ClientMessage;
+import org.jboss.messaging.core.client.ClientConsumer;
+import org.jboss.messaging.core.client.MessageHandler;
+import org.jboss.messaging.core.client.impl.ClientSessionFactoryImpl;
+import org.jboss.messaging.core.config.impl.ConfigurationImpl;
+import org.jboss.messaging.core.config.TransportConfiguration;
+import org.jboss.messaging.core.exception.MessagingException;
+import org.jboss.messaging.core.settings.impl.QueueSettings;
+import org.jboss.messaging.util.SimpleString;
+
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.TimeUnit;
+import java.util.List;
+import java.util.ArrayList;
+
+import junit.framework.TestCase;
+import junit.framework.TestResult;
+import junit.framework.TestSuite;
+import junit.textui.TestRunner;
+
+/**
+ * @author <a href="mailto:andy.taylor at jboss.org">Andy Taylor</a>
+ */
+public class ExpiryRunnerTest extends UnitTestCase
+{
+   private MessagingService messagingService;
+
+   private ClientSession clientSession;
+
+   private SimpleString qName = new SimpleString("ExpiryRunnerTestQ");
+
+   private SimpleString qName2 = new SimpleString("ExpiryRunnerTestQ2");
+
+   private SimpleString expiryQueue;
+
+   private SimpleString expiryAddress;
+
+   public void testBasicExpire() throws Exception
+   {
+      ClientProducer producer = clientSession.createProducer(qName);
+      int numMessages = 100;
+      long expiration = System.currentTimeMillis();
+      for (int i = 0; i < numMessages; i++)
+      {
+         ClientMessage m = createTextMessage("m" + i, clientSession);
+         m.setExpiration(expiration);
+         producer.send(m);
+      }
+      Thread.sleep(1600);
+      assertEquals(0, messagingService.getServer().getPostOffice().getBinding(qName).getQueue().getMessageCount());
+      assertEquals(0, messagingService.getServer().getPostOffice().getBinding(qName).getQueue().getDeliveringCount());
+
+      ClientConsumer consumer = clientSession.createConsumer(expiryQueue);
+      clientSession.start();
+      for (int i = 0; i < numMessages; i++)
+      {
+         ClientMessage cm = consumer.receive(500);
+         assertNotNull(cm);
+         assertEquals("m" + i, cm.getBody().getString());
+      }
+      consumer.close();
+   }
+
+   public void testExpireHalf() throws Exception
+   {
+      ClientProducer producer = clientSession.createProducer(qName);
+      int numMessages = 100;
+      long expiration = System.currentTimeMillis();
+      for (int i = 0; i < numMessages; i++)
+      {
+         ClientMessage m = createTextMessage("m" + i, clientSession);
+         if (i % 2 == 0)
+         {
+            m.setExpiration(expiration);
+         }
+         producer.send(m);
+      }
+      Thread.sleep(1600);
+      assertEquals(50, messagingService.getServer().getPostOffice().getBinding(qName).getQueue().getMessageCount());
+      assertEquals(0, messagingService.getServer().getPostOffice().getBinding(qName).getQueue().getDeliveringCount());
+
+      ClientConsumer consumer = clientSession.createConsumer(expiryQueue);
+      clientSession.start();
+      for (int i = 0; i < numMessages; i += 2)
+      {
+         ClientMessage cm = consumer.receive(500);
+         assertNotNull(cm);
+         assertEquals("m" + i, cm.getBody().getString());
+      }
+      consumer.close();
+   }
+
+   public void testExpireConsumeHalf() throws Exception
+   {
+      ClientProducer producer = clientSession.createProducer(qName);
+      int numMessages = 100;
+      long expiration = System.currentTimeMillis() + 1000;
+      for (int i = 0; i < numMessages; i++)
+      {
+         ClientMessage m = createTextMessage("m" + i, clientSession);
+         m.setExpiration(expiration);
+         producer.send(m);
+      }
+      ClientConsumer consumer = clientSession.createConsumer(qName);
+      clientSession.start();
+      for (int i = 0; i < numMessages / 2; i++)
+      {
+         ClientMessage cm = consumer.receive(500);
+         assertNotNull("message not received " + i, cm);
+         cm.acknowledge();
+         assertEquals("m" + i, cm.getBody().getString());
+      }
+      consumer.close();
+      Thread.sleep(2100);
+      assertEquals(0, messagingService.getServer().getPostOffice().getBinding(qName).getQueue().getMessageCount());
+      assertEquals(0, messagingService.getServer().getPostOffice().getBinding(qName).getQueue().getDeliveringCount());
+
+      consumer = clientSession.createConsumer(expiryQueue);
+      clientSession.start();
+      for (int i = 50; i < numMessages; i++)
+      {
+         ClientMessage cm = consumer.receive(500);
+         assertNotNull(cm);
+         assertEquals("m" + i, cm.getBody().getString());
+      }
+      consumer.close();
+   }
+
+   public void testExpireFromMultipleQueues() throws Exception
+   {
+      clientSession.createQueue(qName, qName2, null, false, false, true);
+      QueueSettings queueSettings = new QueueSettings();
+      queueSettings.setExpiryAddress(expiryAddress);
+      messagingService.getServer().getQueueSettingsRepository().addMatch(qName2.toString(), queueSettings);
+      ClientProducer producer = clientSession.createProducer(qName);
+      int numMessages = 100;
+      long expiration = System.currentTimeMillis();
+      for (int i = 0; i < numMessages; i++)
+      {
+         ClientMessage m = createTextMessage("m" + i, clientSession);
+         m.setExpiration(expiration);
+         producer.send(m);
+      }
+      Thread.sleep(1600);
+      assertEquals(0, messagingService.getServer().getPostOffice().getBinding(qName).getQueue().getMessageCount());
+      assertEquals(0, messagingService.getServer().getPostOffice().getBinding(qName).getQueue().getDeliveringCount());
+
+      ClientConsumer consumer = clientSession.createConsumer(expiryQueue);
+      clientSession.start();
+      for (int i = 0; i < numMessages; i++)
+      {
+         ClientMessage cm = consumer.receive(500);
+         assertNotNull(cm);
+         assertEquals("m" + i, cm.getBody().getString());
+      }
+      for (int i = 0; i < numMessages; i++)
+      {
+         ClientMessage cm = consumer.receive(500);
+         assertNotNull(cm);
+         assertEquals("m" + i, cm.getBody().getString());
+      }
+      consumer.close();
+   }
+
+   public void testExpireWhilstConsuming() throws Exception
+   {
+      ClientProducer producer = clientSession.createProducer(qName);
+      ClientConsumer consumer = clientSession.createConsumer(qName);
+      CountDownLatch latch = new CountDownLatch(1);
+      DummyMessageHandler dummyMessageHandler = new DummyMessageHandler(consumer, latch);
+      clientSession.start();
+      new Thread(dummyMessageHandler).start();
+      long expiration = System.currentTimeMillis() + 1000;
+      int numMessages = 0;
+      long sendMessagesUntil = System.currentTimeMillis() + 2000;
+      do
+      {
+         ClientMessage m = createTextMessage("m" + (numMessages++), clientSession);
+         m.setExpiration(expiration);
+         producer.send(m);
+         Thread.sleep(100);
+      }
+      while (System.currentTimeMillis() < sendMessagesUntil);
+      assertTrue(latch.await(10000, TimeUnit.MILLISECONDS));
+      consumer.close();
+
+      consumer = clientSession.createConsumer(expiryQueue);
+      do
+      {
+         ClientMessage cm = consumer.receive(2000);
+         if(cm == null)
+         {
+            break;
+         }
+         String text = cm.getBody().getString();
+         cm.acknowledge();
+         assertFalse(dummyMessageHandler.payloads.contains(text));
+         dummyMessageHandler.payloads.add(text);
+      } while(true);
+
+      for(int i = 0; i < numMessages; i++)
+      {
+         assertTrue(dummyMessageHandler.payloads.remove("m" + i));
+      }
+      assertTrue(dummyMessageHandler.payloads.isEmpty());
+      consumer.close();
+   }
+
+   public static void main(String[] args) throws Exception
+   {
+      for (int i = 0; i < 1000; i++)
+      {
+         TestSuite suite = new TestSuite();
+         ExpiryRunnerTest expiryRunnerTest = new ExpiryRunnerTest();
+         expiryRunnerTest.setName("testExpireWhilstConsuming");
+         suite.addTest(expiryRunnerTest);
+
+         TestResult result = TestRunner.run(suite);
+         if(result.errorCount() > 0 || result.failureCount() > 0)
+         {
+            System.exit(1);
+         }
+      }
+   }
+
+   @Override
+   protected void setUp() throws Exception
+   {
+      ConfigurationImpl configuration = new ConfigurationImpl();
+      configuration.setSecurityEnabled(false);
+      configuration.setMessageExpiryScanPeriod(1000);
+      TransportConfiguration transportConfig = new TransportConfiguration(INVM_ACCEPTOR_FACTORY);
+      configuration.getAcceptorConfigurations().add(transportConfig);
+      messagingService = MessagingServiceImpl.newNullStorageMessagingServer(configuration);
+      // start the server
+      messagingService.start();
+      // then we create a client as normal
+      ClientSessionFactory sessionFactory = new ClientSessionFactoryImpl(new TransportConfiguration(INVM_CONNECTOR_FACTORY));
+      sessionFactory.setBlockOnAcknowledge(true);
+      clientSession = sessionFactory.createSession(false, true, true);
+      clientSession.createQueue(qName, qName, null, false, false, true);
+      expiryAddress = new SimpleString("EA");
+      expiryQueue = new SimpleString("expiryQ");
+      QueueSettings queueSettings = new QueueSettings();
+      queueSettings.setExpiryAddress(expiryAddress);
+      messagingService.getServer().getQueueSettingsRepository().addMatch(qName.toString(), queueSettings);
+      messagingService.getServer().getQueueSettingsRepository().addMatch(qName2.toString(), queueSettings);
+      clientSession.createQueue(expiryAddress, expiryQueue, null, false, false, true);
+   }
+
+   @Override
+   protected void tearDown() throws Exception
+   {
+      if (clientSession != null)
+      {
+         try
+         {
+            clientSession.close();
+         }
+         catch (MessagingException e1)
+         {
+            //
+         }
+      }
+      if (messagingService != null && messagingService.isStarted())
+      {
+         try
+         {
+            messagingService.stop();
+         }
+         catch (Exception e1)
+         {
+            //
+         }
+      }
+      messagingService = null;
+      clientSession = null;
+   }
+
+   private static class DummyMessageHandler implements Runnable
+   {
+      List<String> payloads = new ArrayList<String>();
+
+      private final ClientConsumer consumer;
+
+      private final CountDownLatch latch;
+
+      public DummyMessageHandler(ClientConsumer consumer, CountDownLatch latch)
+      {
+         this.consumer = consumer;
+         this.latch = latch;
+      }
+
+      public void run()
+      {
+         while (true)
+         {
+            try
+            {
+               ClientMessage message = consumer.receive(5000);
+               if (message == null)
+               {
+                  break;
+               }
+               message.acknowledge();
+               payloads.add(message.getBody().getString());
+
+               Thread.sleep(110);
+            }
+            catch (Exception e)
+            {
+               e.printStackTrace();
+            }
+         }
+         latch.countDown();
+
+      }
+   }
+}

Modified: trunk/tests/src/org/jboss/messaging/tests/performance/persistence/FakePostOffice.java
===================================================================
--- trunk/tests/src/org/jboss/messaging/tests/performance/persistence/FakePostOffice.java	2008-11-25 12:18:19 UTC (rev 5432)
+++ trunk/tests/src/org/jboss/messaging/tests/performance/persistence/FakePostOffice.java	2008-11-25 12:51:55 UTC (rev 5433)
@@ -89,6 +89,12 @@
       return null;
    }
 
+
+   public List<Queue> getQueues()
+   {
+      return null;
+   }
+
    public Map<SimpleString, List<Binding>> getMappings()
    {
       return null;

Modified: trunk/tests/src/org/jboss/messaging/tests/unit/core/config/impl/ConfigurationImplTest.java
===================================================================
--- trunk/tests/src/org/jboss/messaging/tests/unit/core/config/impl/ConfigurationImplTest.java	2008-11-25 12:18:19 UTC (rev 5432)
+++ trunk/tests/src/org/jboss/messaging/tests/unit/core/config/impl/ConfigurationImplTest.java	2008-11-25 12:51:55 UTC (rev 5433)
@@ -72,6 +72,8 @@
       assertEquals(ConfigurationImpl.DEFAULT_JOURNAL_MAX_AIO, conf.getJournalMaxAIO());
       assertEquals(ConfigurationImpl.DEFAULT_WILDCARD_ROUTING_ENABLED, conf.isWildcardRoutingEnabled());
       assertEquals(ConfigurationImpl.DEFAULT_TRANSACTION_TIMEOUT, conf.getTransactionTimeout());
+      assertEquals(ConfigurationImpl.DEFAULT_MESSAGE_EXPIRY_SCAN_PERIOD, conf.getMessageExpiryScanPeriod());
+      assertEquals(ConfigurationImpl.DEFAULT_MESSAGE_EXPIRY_THREAD_PRIORITY, conf.getMessageExpiryThreadPriority());
       assertEquals(ConfigurationImpl.DEFAULT_TRANSACTION_TIMEOUT_SCAN_PERIOD, conf.getTransactionTimeoutScanPeriod());
       assertEquals(ConfigurationImpl.DEFAULT_MANAGEMENT_ADDRESS, conf.getManagementAddress());
    }
@@ -156,6 +158,15 @@
          s = randomString();
          conf.setManagementAddress(new SimpleString(s));
          assertEquals(s, conf.getManagementAddress().toString());
+
+         i = randomInt();
+
+         conf.setMessageExpiryThreadPriority(i);
+         assertEquals(i, conf.getMessageExpiryThreadPriority());
+
+         l = randomLong();
+         conf.setMessageExpiryScanPeriod(l);
+         assertEquals(l, conf.getMessageExpiryScanPeriod());
       }
    }
    

Modified: trunk/tests/src/org/jboss/messaging/tests/unit/core/config/impl/FileConfigurationTest.java
===================================================================
--- trunk/tests/src/org/jboss/messaging/tests/unit/core/config/impl/FileConfigurationTest.java	2008-11-25 12:18:19 UTC (rev 5432)
+++ trunk/tests/src/org/jboss/messaging/tests/unit/core/config/impl/FileConfigurationTest.java	2008-11-25 12:51:55 UTC (rev 5433)
@@ -59,6 +59,8 @@
       assertEquals(true, conf.isWildcardRoutingEnabled());
       assertEquals(98765, conf.getTransactionTimeout());
       assertEquals(56789, conf.getTransactionTimeoutScanPeriod());
+      assertEquals(10111213, conf.getMessageExpiryScanPeriod());
+      assertEquals(8, conf.getMessageExpiryThreadPriority());
       assertEquals(new SimpleString("Giraffe"), conf.getManagementAddress());
       assertEquals(2, conf.getInterceptorClassNames().size());
       assertTrue(conf.getInterceptorClassNames().contains("org.jboss.messaging.tests.unit.core.config.impl.TestInterceptor1"));




More information about the jboss-cvs-commits mailing list