From 64df937a5878d710258a1413dce213683d801d91 Mon Sep 17 00:00:00 2001 From: markpollack Date: Sat, 25 Jul 2009 05:05:43 +0000 Subject: [PATCH] SPRNET-1233 - Support for sending to remote private queue. Changes to MessageQueueTemplate 1. Registers MessageQueueMetadataCache in ApplicationContext if not set explicitly. 2. Logs warnings if MetadataCache property is set to null, for example in stand alone API usage. 3. No longer has DefaultMessageQueueObjectName as a required . DefaultMessageQueueFactory updates MessageQueueMetadataCache in ApplicationContext when a new MessageQueue is registered. --- .../Core/DefaultMessageQueueFactory.cs | 5 ++ .../Messaging/Core/MessageQueueTemplate.cs | 89 +++++++++++++++---- ...AbstractPeekingMessageListenerContainer.cs | 5 +- .../Listener/MessageListenerAdapter.cs | 23 ++++- ...onTransactionalMessageListenerContainer.cs | 4 +- .../Core/MessageQueueMetadataCacheTests.cs | 2 +- .../Core/MessageQueueMetadataCacheTests.xml | 10 --- .../Core/MessageQueueTemplateTests.cs | 2 +- .../Core/MessageQueueTemplateTests.xml | 1 + ...ributedTxMessageListenerContainerTests.xml | 2 +- ...nsactionalMessageListenerContainerTests.cs | 11 +-- ...sactionalMessageListenerContainerTests.xml | 19 +++- .../Listener/SimpleExceptionHandler.cs | 1 + 13 files changed, 128 insertions(+), 46 deletions(-) diff --git a/src/Spring/Spring.Messaging/Messaging/Core/DefaultMessageQueueFactory.cs b/src/Spring/Spring.Messaging/Messaging/Core/DefaultMessageQueueFactory.cs index 6f272786..26f6e525 100644 --- a/src/Spring/Spring.Messaging/Messaging/Core/DefaultMessageQueueFactory.cs +++ b/src/Spring/Spring.Messaging/Messaging/Core/DefaultMessageQueueFactory.cs @@ -66,6 +66,11 @@ namespace Spring.Messaging.Core MessageQueueFactoryObject mqfo = new MessageQueueFactoryObject(); mqfo.MessageCreatorDelegate = messageQueueCreatorDelegate; applicationContext.ObjectFactory.RegisterSingleton(messageQueueObjectName, mqfo); + IDictionary caches = applicationContext.GetObjectsOfType(typeof(MessageQueueMetadataCache)); + foreach (DictionaryEntry entry in caches) + { + ((MessageQueueMetadataCache) entry.Value).Insert(mqfo.Path, new MessageQueueMetadata(mqfo.RemoteQueue, mqfo.RemoteQueueIsTransactional)); + } } /// diff --git a/src/Spring/Spring.Messaging/Messaging/Core/MessageQueueTemplate.cs b/src/Spring/Spring.Messaging/Messaging/Core/MessageQueueTemplate.cs index 22baa95e..8131184f 100644 --- a/src/Spring/Spring.Messaging/Messaging/Core/MessageQueueTemplate.cs +++ b/src/Spring/Spring.Messaging/Messaging/Core/MessageQueueTemplate.cs @@ -27,6 +27,7 @@ using Spring.Messaging.Support; using Spring.Messaging.Support.Converters; using Spring.Objects.Factory; using Spring.Objects.Factory.Config; +using Spring.Util; namespace Spring.Messaging.Core { @@ -82,12 +83,14 @@ namespace Spring.Messaging.Core private string messageConverterObjectName; private IMessageQueueFactory messageQueueFactory; - protected IApplicationContext applicationContext; + protected IConfigurableApplicationContext applicationContext; private TimeSpan timeout = MessageQueue.InfiniteTimeout; private MessageQueueMetadataCache metadataCache; + public const string METADATA_CACHE_NAME = "__MessageQueueMetadataCache__"; + #endregion #region Constructors @@ -241,7 +244,17 @@ namespace Spring.Messaging.Core public IApplicationContext ApplicationContext { get { return applicationContext; } - set { applicationContext = value; } + set { + AssertUtils.ArgumentNotNull(value, "An ApplicationContext instance is required"); + IConfigurableApplicationContext ctx = value as IConfigurableApplicationContext; + if (ctx == null) + { + throw new InvalidOperationException( + "Implementations of IApplicationContext must also implement IConfigurableApplicationContext"); + } + + applicationContext = ctx; + } } #endregion @@ -263,12 +276,10 @@ namespace Spring.Messaging.Core /// public void AfterPropertiesSet() { - if (DefaultMessageQueueObjectName == null) - { - throw new ArgumentException("DefaultMessageQueueObjectName is required."); - } if (MessageQueueFactory == null) { + AssertUtils.ArgumentNotNull(applicationContext, "MessageQueueTemplate requires the ApplicationContext property to be set if the MessageQueueFactory property is not set for automatic create of the DefaultMessageQueueFactory"); + DefaultMessageQueueFactory mqf = new DefaultMessageQueueFactory(); mqf.ApplicationContext = applicationContext; messageQueueFactory = mqf; @@ -278,12 +289,43 @@ namespace Spring.Messaging.Core messageConverterObjectName = QueueUtils.RegisterDefaultMessageConverter(applicationContext); } //If it has not been set by the user explicitly, then initialize. + CreateDefaultMetadataCache(); + } + + protected virtual void CreateDefaultMetadataCache() + { if (metadataCache == null) { - metadataCache = new MessageQueueMetadataCache(); - metadataCache.ApplicationContext = ApplicationContext; - metadataCache.AfterPropertiesSet(); - metadataCache.Initialize(); + if (applicationContext.ContainsObject(METADATA_CACHE_NAME)) + { + metadataCache = applicationContext.GetObject(METADATA_CACHE_NAME) as MessageQueueMetadataCache; + } + else + { + metadataCache = new MessageQueueMetadataCache(); + if (ApplicationContext != null) + { + metadataCache.ApplicationContext = ApplicationContext; + metadataCache.AfterPropertiesSet(); + metadataCache.Initialize(); + applicationContext.ObjectFactory.RegisterSingleton("__MessageQueueMetadataCache__", + metadataCache); + } + else + { + #region Logging + + if (LOG.IsWarnEnabled) + { + LOG.Warn( + "The ApplicationContext property has not been set, so the MessageQueueMetadataCache can not be automatically generated. " + + "Please explictly set the MessageQueueMetadataCache using the property MetadataCache or set the ApplicationContext property. " + + "This will only effect the use of MessageQueueTemplate when publishing to remote queues."); + } + + #endregion + } + } } } @@ -500,19 +542,30 @@ namespace Spring.Messaging.Core protected virtual void DoSendMessageQueue(MessageQueue mq, Message msg) { MessageQueueTransaction transactionToUse = QueueUtils.GetMessageQueueTransaction(null); - MessageQueueMetadata mqMetadata = metadataCache.Get(mq.Path); - if (mqMetadata != null) + if (metadataCache != null) { - if (mqMetadata.RemoteQueue) + MessageQueueMetadata mqMetadata = metadataCache.Get(mq.Path); + if (mqMetadata != null) { - if (mqMetadata.RemoteQueueIsTransactional) + if (mqMetadata.RemoteQueue) { - // DefaultMessageQueue transaction is externally managed. - DoSendMessageTransaction(mq, transactionToUse, msg); + if (mqMetadata.RemoteQueueIsTransactional) + { + // DefaultMessageQueue transaction is externally managed. + DoSendMessageTransaction(mq, transactionToUse, msg); + return; + } + DoSendMessageQueueNonTransactional(mq, transactionToUse, msg); return; } - DoSendMessageQueueNonTransactional(mq, transactionToUse, msg); - return; + } + } else + { + if (LOG.IsWarnEnabled) + { + LOG.Warn("MetadataCache has not been initialized. Set the MetadataCache explicitly in standalone usage and/or " + + "configure the MessageQueueTemplate in an ApplicationContext. If deployed in an ApplicationContext by default " + + "the MetadataCache will automaticaly populated."); } } // Handle assuming these are local queues. diff --git a/src/Spring/Spring.Messaging/Messaging/Listener/AbstractPeekingMessageListenerContainer.cs b/src/Spring/Spring.Messaging/Messaging/Listener/AbstractPeekingMessageListenerContainer.cs index 874c0c58..9b4ad481 100644 --- a/src/Spring/Spring.Messaging/Messaging/Listener/AbstractPeekingMessageListenerContainer.cs +++ b/src/Spring/Spring.Messaging/Messaging/Listener/AbstractPeekingMessageListenerContainer.cs @@ -445,11 +445,10 @@ namespace Spring.Messaging.Listener /// messageQueue.Receive(). /// /// It allows subclasses to modify the state of the MessageQueue - /// before receiving which maybe required when using remote queues, for example - /// to set a MessageFormatter. + /// before receiving which maybe required when using remote queues /// protected virtual void BeforeMessageReceived(MessageQueue messageQueue) - { + { } #endregion diff --git a/src/Spring/Spring.Messaging/Messaging/Listener/MessageListenerAdapter.cs b/src/Spring/Spring.Messaging/Messaging/Listener/MessageListenerAdapter.cs index 3deec1b2..1057c647 100644 --- a/src/Spring/Spring.Messaging/Messaging/Listener/MessageListenerAdapter.cs +++ b/src/Spring/Spring.Messaging/Messaging/Listener/MessageListenerAdapter.cs @@ -230,6 +230,12 @@ namespace Spring.Messaging.Listener { messageConverterObjectName = QueueUtils.RegisterDefaultMessageConverter(applicationContext); } + if (messageQueueTemplate == null) + { + messageQueueTemplate = new MessageQueueTemplate(); + messageQueueTemplate.ApplicationContext = ApplicationContext; + messageQueueTemplate.AfterPropertiesSet(); + } } #endregion @@ -351,6 +357,18 @@ namespace Spring.Messaging.Listener } } + /// + /// Sets the message queue template. + /// + /// If not set, will create one for it own internal use whne MessageListenerAdapter is constructed. + /// It maybe useful to share an existing instance if you have an extensively configured MessageQueueTemplate. + /// + /// The message queue template. + public MessageQueueTemplate MessageQueueTemplate + { + set { messageQueueTemplate = value; } + } + #region IMessageListener Members /// @@ -395,7 +413,6 @@ namespace Spring.Messaging.Listener protected virtual void InitDefaultStrategies() { processingExpression = Expression.Parse(defaultHandlerMethod + "(#convertedObject)"); - messageQueueTemplate = new MessageQueueTemplate(); } /// @@ -435,6 +452,10 @@ namespace Spring.Messaging.Listener protected virtual void SendResponse(MessageQueue destination, Message response) { //Will send with appropriate transaction semantics + if (logger.IsDebugEnabled) + { + logger.Debug("Sending response message to path = [" + destination.Path + "]"); + } messageQueueTemplate.Send(destination, response); } diff --git a/src/Spring/Spring.Messaging/Messaging/Listener/NonTransactionalMessageListenerContainer.cs b/src/Spring/Spring.Messaging/Messaging/Listener/NonTransactionalMessageListenerContainer.cs index e2501465..34b45139 100644 --- a/src/Spring/Spring.Messaging/Messaging/Listener/NonTransactionalMessageListenerContainer.cs +++ b/src/Spring/Spring.Messaging/Messaging/Listener/NonTransactionalMessageListenerContainer.cs @@ -73,7 +73,7 @@ namespace Spring.Messaging.Listener /// Perform a receive opertion on the message queue and execute the /// message listener /// - /// The DefaultMessageQueue. + /// The MessageQueue. /// /// true if received a message, false otherwise /// @@ -119,7 +119,7 @@ namespace Spring.Messaging.Listener if (LOG.IsErrorEnabled) { - LOG.Error("Error receiving message from DefaultMessageQueue [" + mq.Path + + LOG.Error("Error receiving message from MessageQueue [" + mq.Path + "], closing queue and clearing connection cache."); } diff --git a/test/Spring/Spring.Messaging.Tests/Messaging/Core/MessageQueueMetadataCacheTests.cs b/test/Spring/Spring.Messaging.Tests/Messaging/Core/MessageQueueMetadataCacheTests.cs index 9195d6a2..96dbdefb 100644 --- a/test/Spring/Spring.Messaging.Tests/Messaging/Core/MessageQueueMetadataCacheTests.cs +++ b/test/Spring/Spring.Messaging.Tests/Messaging/Core/MessageQueueMetadataCacheTests.cs @@ -36,7 +36,7 @@ namespace Spring.Messaging.Core { protected override string[] ConfigLocations { - get { return new[] {"assembly://Spring.Messaging.Tests/Spring.Messaging.Core/MessageQueueTemplateTests.xml"}; } + get { return new[] {"assembly://Spring.Messaging.Tests/Spring.Messaging.Core/MessageQueueMetadataCacheTests.xml"}; } } [Test] diff --git a/test/Spring/Spring.Messaging.Tests/Messaging/Core/MessageQueueMetadataCacheTests.xml b/test/Spring/Spring.Messaging.Tests/Messaging/Core/MessageQueueMetadataCacheTests.xml index da7def1c..5108625e 100644 --- a/test/Spring/Spring.Messaging.Tests/Messaging/Core/MessageQueueMetadataCacheTests.xml +++ b/test/Spring/Spring.Messaging.Tests/Messaging/Core/MessageQueueMetadataCacheTests.xml @@ -36,14 +36,4 @@ - - - - - - - - - - \ No newline at end of file diff --git a/test/Spring/Spring.Messaging.Tests/Messaging/Core/MessageQueueTemplateTests.cs b/test/Spring/Spring.Messaging.Tests/Messaging/Core/MessageQueueTemplateTests.cs index 7709d69c..f4bbf830 100644 --- a/test/Spring/Spring.Messaging.Tests/Messaging/Core/MessageQueueTemplateTests.cs +++ b/test/Spring/Spring.Messaging.Tests/Messaging/Core/MessageQueueTemplateTests.cs @@ -116,7 +116,7 @@ namespace Spring.Messaging.Core [Test] - [ExpectedException(typeof (ArgumentException), ExpectedMessage = "DefaultMessageQueueObjectName is required.")] + [ExpectedException(typeof (ArgumentNullException))] public void NoMessageQueueNameSpecified() { MessageQueueTemplate mqt = new MessageQueueTemplate(); diff --git a/test/Spring/Spring.Messaging.Tests/Messaging/Core/MessageQueueTemplateTests.xml b/test/Spring/Spring.Messaging.Tests/Messaging/Core/MessageQueueTemplateTests.xml index 305a027a..873708c3 100644 --- a/test/Spring/Spring.Messaging.Tests/Messaging/Core/MessageQueueTemplateTests.xml +++ b/test/Spring/Spring.Messaging.Tests/Messaging/Core/MessageQueueTemplateTests.xml @@ -34,6 +34,7 @@ + diff --git a/test/Spring/Spring.Messaging.Tests/Messaging/Listener/DistributedTxMessageListenerContainerTests.xml b/test/Spring/Spring.Messaging.Tests/Messaging/Listener/DistributedTxMessageListenerContainerTests.xml index 587ceb34..5de74533 100644 --- a/test/Spring/Spring.Messaging.Tests/Messaging/Listener/DistributedTxMessageListenerContainerTests.xml +++ b/test/Spring/Spring.Messaging.Tests/Messaging/Listener/DistributedTxMessageListenerContainerTests.xml @@ -77,7 +77,7 @@ - + diff --git a/test/Spring/Spring.Messaging.Tests/Messaging/Listener/NonTransactionalMessageListenerContainerTests.cs b/test/Spring/Spring.Messaging.Tests/Messaging/Listener/NonTransactionalMessageListenerContainerTests.cs index 11269512..1db7ce59 100644 --- a/test/Spring/Spring.Messaging.Tests/Messaging/Listener/NonTransactionalMessageListenerContainerTests.cs +++ b/test/Spring/Spring.Messaging.Tests/Messaging/Listener/NonTransactionalMessageListenerContainerTests.cs @@ -84,19 +84,20 @@ namespace Spring.Messaging.Listener [Test] public void SendAndAsyncReceive() { + //MessageQueueTemplate q = applicationContext["testQueueTemplate"] as MessageQueueTemplate; + MessageQueueTemplate q = applicationContext["testRemoteTemplate"] as MessageQueueTemplate; Assert.IsNotNull(q); - - /* + q.ConvertAndSend("Hello World 1"); q.ConvertAndSend("Hello World 2"); q.ConvertAndSend("Hello World 3"); q.ConvertAndSend("Hello World 4"); q.ConvertAndSend("Hello World 5"); - */ - + + Assert.AreEqual(0, listener.MessageCount); container.Start(); @@ -107,7 +108,7 @@ namespace Spring.Messaging.Listener container.Stop(); container.Shutdown(); - Thread.Sleep(2500); + Thread.Sleep(2500); } diff --git a/test/Spring/Spring.Messaging.Tests/Messaging/Listener/NonTransactionalMessageListenerContainerTests.xml b/test/Spring/Spring.Messaging.Tests/Messaging/Listener/NonTransactionalMessageListenerContainerTests.xml index cbe1fc27..b6b88940 100644 --- a/test/Spring/Spring.Messaging.Tests/Messaging/Listener/NonTransactionalMessageListenerContainerTests.xml +++ b/test/Spring/Spring.Messaging.Tests/Messaging/Listener/NonTransactionalMessageListenerContainerTests.xml @@ -40,34 +40,45 @@ - + - + + + + + + + + - + + + - + + + diff --git a/test/Spring/Spring.Messaging.Tests/Messaging/Listener/SimpleExceptionHandler.cs b/test/Spring/Spring.Messaging.Tests/Messaging/Listener/SimpleExceptionHandler.cs index 9dc3cf41..06a795b8 100644 --- a/test/Spring/Spring.Messaging.Tests/Messaging/Listener/SimpleExceptionHandler.cs +++ b/test/Spring/Spring.Messaging.Tests/Messaging/Listener/SimpleExceptionHandler.cs @@ -28,6 +28,7 @@ namespace Spring.Messaging.Listener public void OnException(Exception exception, Message message) { LOG.Error("Exception Handler processing message id = [" + message.Id + "]"); + LOG.Error("Exception = ", exception); messageCount++; }