summaryrefslogtreecommitdiff
path: root/java/broker/src/main
diff options
context:
space:
mode:
authorMartin Ritchie <ritchiem@apache.org>2009-08-12 18:05:34 +0000
committerMartin Ritchie <ritchiem@apache.org>2009-08-12 18:05:34 +0000
commit60276b6aec1492dcaeee934ad304060611208827 (patch)
tree142bf8af2d3e6d8454d4c6204f4b8d925039b6d7 /java/broker/src/main
parentae9928a6bf3e01c26d0ffae928c626bed771931d (diff)
downloadqpid-python-60276b6aec1492dcaeee934ad304060611208827.tar.gz
QPID-2002 : Addition of a QueueActor to be set during running of the processQueue thread
Made QueueLogSubject public so it can be reused by QueueActor Updated SAMQQ to create a QueueActor for use during the processQueue thread run git-svn-id: https://svn.apache.org/repos/asf/qpid/trunk/qpid@803638 13f79535-47bb-0310-9956-ffa450edef68
Diffstat (limited to 'java/broker/src/main')
-rw-r--r--java/broker/src/main/java/org/apache/qpid/server/logging/actors/QueueActor.java52
-rw-r--r--java/broker/src/main/java/org/apache/qpid/server/logging/subjects/QueueLogSubject.java4
-rw-r--r--java/broker/src/main/java/org/apache/qpid/server/queue/SimpleAMQQueue.java10
3 files changed, 64 insertions, 2 deletions
diff --git a/java/broker/src/main/java/org/apache/qpid/server/logging/actors/QueueActor.java b/java/broker/src/main/java/org/apache/qpid/server/logging/actors/QueueActor.java
new file mode 100644
index 0000000000..acac447ff6
--- /dev/null
+++ b/java/broker/src/main/java/org/apache/qpid/server/logging/actors/QueueActor.java
@@ -0,0 +1,52 @@
+/*
+ *
+ * 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.qpid.server.logging.actors;
+
+import org.apache.qpid.server.logging.RootMessageLogger;
+import org.apache.qpid.server.logging.subjects.QueueLogSubject;
+import org.apache.qpid.server.queue.AMQQueue;
+
+import java.text.MessageFormat;
+
+/**
+ * This Actor is used when while the queue is performing an asynchronous process
+ * of its queue.
+ */
+public class QueueActor extends AbstractActor
+{
+
+ /**
+ * Create an QueueLogSubject that Logs in the following format.
+ *
+ * @param queue The queue that this Actor is working for
+ * @param rootLogger the Root logger to use.
+ */
+ public QueueActor(AMQQueue queue, RootMessageLogger rootLogger)
+ {
+ super(rootLogger);
+
+ _logString = "[" + MessageFormat.format(QueueLogSubject.LOG_FORMAT,
+ queue.getVirtualHost().getName(),
+ queue.getName()) + "] ";
+
+ }
+}
+ \ No newline at end of file
diff --git a/java/broker/src/main/java/org/apache/qpid/server/logging/subjects/QueueLogSubject.java b/java/broker/src/main/java/org/apache/qpid/server/logging/subjects/QueueLogSubject.java
index 89f31ef477..b132d9e93f 100644
--- a/java/broker/src/main/java/org/apache/qpid/server/logging/subjects/QueueLogSubject.java
+++ b/java/broker/src/main/java/org/apache/qpid/server/logging/subjects/QueueLogSubject.java
@@ -33,12 +33,12 @@ public class QueueLogSubject extends AbstractLogSubject
* 0 - Virtualhost name
* 1 - queue name
*/
- protected static String BINDING_FORMAT = "vh(/{0})/qu({1})";
+ public static String LOG_FORMAT = "vh(/{0})/qu({1})";
/** Create an QueueLogSubject that Logs in the following format. */
public QueueLogSubject(AMQQueue queue)
{
- setLogStringWithFormat(BINDING_FORMAT,
+ setLogStringWithFormat(LOG_FORMAT,
queue.getVirtualHost().getName(),
queue.getName());
}
diff --git a/java/broker/src/main/java/org/apache/qpid/server/queue/SimpleAMQQueue.java b/java/broker/src/main/java/org/apache/qpid/server/queue/SimpleAMQQueue.java
index 82dd4b0195..b14b92b014 100644
--- a/java/broker/src/main/java/org/apache/qpid/server/queue/SimpleAMQQueue.java
+++ b/java/broker/src/main/java/org/apache/qpid/server/queue/SimpleAMQQueue.java
@@ -30,8 +30,10 @@ import org.apache.qpid.server.subscription.Subscription;
import org.apache.qpid.server.subscription.SubscriptionList;
import org.apache.qpid.server.virtualhost.VirtualHost;
import org.apache.qpid.server.logging.actors.CurrentActor;
+import org.apache.qpid.server.logging.actors.QueueActor;
import org.apache.qpid.server.logging.subjects.QueueLogSubject;
import org.apache.qpid.server.logging.LogSubject;
+import org.apache.qpid.server.logging.LogActor;
import org.apache.qpid.server.logging.messages.QueueMessages;
/*
@@ -118,6 +120,7 @@ public class SimpleAMQQueue implements AMQQueue, Subscription.StateListener
private AtomicInteger _deliveredMessages = new AtomicInteger();
private AtomicBoolean _stopped = new AtomicBoolean(false);
private LogSubject _logSubject;
+ private LogActor _logActor;
protected SimpleAMQQueue(AMQShortString name, boolean durable, AMQShortString owner, boolean autoDelete, VirtualHost virtualHost)
throws AMQException
@@ -154,6 +157,7 @@ public class SimpleAMQQueue implements AMQQueue, Subscription.StateListener
_asyncDelivery = ReferenceCountingExecutorService.getInstance().acquireExecutorService();
_logSubject = new QueueLogSubject(this);
+ _logActor = new QueueActor(this, CurrentActor.get().getRootMessageLogger());
// Log the correct creation message
@@ -1189,12 +1193,18 @@ public class SimpleAMQQueue implements AMQQueue, Subscription.StateListener
{
try
{
+ CurrentActor.set(_logActor);
processQueue(this);
}
catch (AMQException e)
{
_logger.error(e);
}
+ finally
+ {
+ CurrentActor.remove();
+ }
+
}