summaryrefslogtreecommitdiff
path: root/qpid/java
diff options
context:
space:
mode:
authorRobert Godfrey <rgodfrey@apache.org>2014-04-25 10:12:06 +0000
committerRobert Godfrey <rgodfrey@apache.org>2014-04-25 10:12:06 +0000
commitacf84ebf5462342656a257ed978238c23fb1c900 (patch)
tree89895d4bb614e1e7f0fe4f692ba05fa7057ff961 /qpid/java
parent4eddea8954ba9342ab2bc35e495baa673a6015db (diff)
downloadqpid-python-acf84ebf5462342656a257ed978238c23fb1c900.tar.gz
QPID-5578 : Make TaskExecutor an interface and provide a test-only current thread implementation
git-svn-id: https://svn.apache.org/repos/asf/qpid/trunk@1589969 13f79535-47bb-0310-9956-ffa450edef68
Diffstat (limited to 'qpid/java')
-rw-r--r--qpid/java/bdbstore/src/test/java/org/apache/qpid/server/store/berkeleydb/BDBHAVirtualHostNodeTest.java10
-rw-r--r--qpid/java/broker-core/src/main/java/org/apache/qpid/server/Broker.java10
-rw-r--r--qpid/java/broker-core/src/main/java/org/apache/qpid/server/binding/BindingImpl.java4
-rw-r--r--qpid/java/broker-core/src/main/java/org/apache/qpid/server/configuration/updater/Task.java26
-rw-r--r--qpid/java/broker-core/src/main/java/org/apache/qpid/server/configuration/updater/TaskExecutor.java365
-rw-r--r--qpid/java/broker-core/src/main/java/org/apache/qpid/server/configuration/updater/TaskExecutorImpl.java380
-rw-r--r--qpid/java/broker-core/src/main/java/org/apache/qpid/server/configuration/updater/TaskWithException.java26
-rw-r--r--qpid/java/broker-core/src/main/java/org/apache/qpid/server/configuration/updater/VoidTask.java26
-rw-r--r--qpid/java/broker-core/src/main/java/org/apache/qpid/server/configuration/updater/VoidTaskWithException.java26
-rw-r--r--qpid/java/broker-core/src/main/java/org/apache/qpid/server/model/AbstractConfiguredObject.java42
-rw-r--r--qpid/java/broker-core/src/main/java/org/apache/qpid/server/model/BrokerModel.java9
-rw-r--r--qpid/java/broker-core/src/main/java/org/apache/qpid/server/model/Model.java2
-rw-r--r--qpid/java/broker-core/src/main/java/org/apache/qpid/server/model/SystemContextImpl.java14
-rw-r--r--qpid/java/broker-core/src/main/java/org/apache/qpid/server/model/adapter/BrokerAdapter.java19
-rw-r--r--qpid/java/broker-core/src/main/java/org/apache/qpid/server/security/auth/manager/ScramSHA1AuthenticationManager.java14
-rw-r--r--qpid/java/broker-core/src/main/java/org/apache/qpid/server/store/GenericRecoverer.java7
-rw-r--r--qpid/java/broker-core/src/test/java/org/apache/qpid/server/configuration/BrokerConfigurationStoreCreatorTest.java11
-rw-r--r--qpid/java/broker-core/src/test/java/org/apache/qpid/server/configuration/store/ConfigurationEntryStoreTestCase.java7
-rw-r--r--qpid/java/broker-core/src/test/java/org/apache/qpid/server/configuration/store/JsonConfigurationEntryStoreTest.java9
-rw-r--r--qpid/java/broker-core/src/test/java/org/apache/qpid/server/configuration/store/ManagementModeStoreHandlerTest.java9
-rw-r--r--qpid/java/broker-core/src/test/java/org/apache/qpid/server/configuration/updater/CurrentThreadTaskExecutor.java99
-rw-r--r--qpid/java/broker-core/src/test/java/org/apache/qpid/server/configuration/updater/TaskExecutorTest.java14
-rw-r--r--qpid/java/broker-core/src/test/java/org/apache/qpid/server/exchange/FanoutExchangeTest.java10
-rw-r--r--qpid/java/broker-core/src/test/java/org/apache/qpid/server/exchange/HeadersBindingTest.java8
-rw-r--r--qpid/java/broker-core/src/test/java/org/apache/qpid/server/exchange/HeadersExchangeTest.java3
-rw-r--r--qpid/java/broker-core/src/test/java/org/apache/qpid/server/model/VirtualHostTest.java3
-rw-r--r--qpid/java/broker-core/src/test/java/org/apache/qpid/server/model/adapter/FileSystemPreferencesProviderTest.java3
-rw-r--r--qpid/java/broker-core/src/test/java/org/apache/qpid/server/security/auth/manager/ScramSHA1AuthenticationManagerTest.java3
-rw-r--r--qpid/java/broker-core/src/test/java/org/apache/qpid/server/store/BrokerRecovererTest.java17
-rw-r--r--qpid/java/broker-core/src/test/java/org/apache/qpid/server/util/BrokerTestHelper.java12
-rw-r--r--qpid/java/broker-core/src/test/java/org/apache/qpid/server/virtualhost/VirtualHostQueueCreationTest.java5
-rw-r--r--qpid/java/broker-core/src/test/java/org/apache/qpid/server/virtualhostnode/AbstractStandardVirtualHostNodeTest.java3
32 files changed, 704 insertions, 492 deletions
diff --git a/qpid/java/bdbstore/src/test/java/org/apache/qpid/server/store/berkeleydb/BDBHAVirtualHostNodeTest.java b/qpid/java/bdbstore/src/test/java/org/apache/qpid/server/store/berkeleydb/BDBHAVirtualHostNodeTest.java
index d54ed354f0..6cb4e9ce11 100644
--- a/qpid/java/bdbstore/src/test/java/org/apache/qpid/server/store/berkeleydb/BDBHAVirtualHostNodeTest.java
+++ b/qpid/java/bdbstore/src/test/java/org/apache/qpid/server/store/berkeleydb/BDBHAVirtualHostNodeTest.java
@@ -20,7 +20,6 @@
*/
package org.apache.qpid.server.store.berkeleydb;
-import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.when;
import java.io.File;
@@ -31,7 +30,11 @@ import java.util.UUID;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
+import com.sleepycat.je.rep.ReplicatedEnvironment;
+import com.sleepycat.je.rep.ReplicationConfig;
+
import org.apache.qpid.server.configuration.updater.TaskExecutor;
+import org.apache.qpid.server.configuration.updater.TaskExecutorImpl;
import org.apache.qpid.server.model.Broker;
import org.apache.qpid.server.model.ConfigurationChangeListener;
import org.apache.qpid.server.model.ConfiguredObject;
@@ -46,9 +49,6 @@ import org.apache.qpid.server.virtualhostnode.berkeleydb.BDBHAVirtualHostNodeFac
import org.apache.qpid.test.utils.QpidTestCase;
import org.apache.qpid.util.FileUtils;
-import com.sleepycat.je.rep.ReplicatedEnvironment;
-import com.sleepycat.je.rep.ReplicationConfig;
-
public class BDBHAVirtualHostNodeTest extends QpidTestCase
{
@@ -64,7 +64,7 @@ public class BDBHAVirtualHostNodeTest extends QpidTestCase
_broker = BrokerTestHelper.createBrokerMock();
- _taskExecutor = new TaskExecutor();
+ _taskExecutor = new TaskExecutorImpl();
_taskExecutor.start();
when(_broker.getTaskExecutor()).thenReturn(_taskExecutor);
diff --git a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/Broker.java b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/Broker.java
index 3ccfc276fa..78adb46622 100644
--- a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/Broker.java
+++ b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/Broker.java
@@ -35,17 +35,16 @@ import javax.security.auth.Subject;
import org.apache.log4j.LogManager;
import org.apache.log4j.Logger;
import org.apache.log4j.PropertyConfigurator;
+
import org.apache.qpid.server.configuration.BrokerConfigurationStoreCreator;
import org.apache.qpid.server.configuration.store.ManagementModeStoreHandler;
import org.apache.qpid.server.configuration.updater.TaskExecutor;
+import org.apache.qpid.server.configuration.updater.TaskExecutorImpl;
import org.apache.qpid.server.logging.EventLogger;
import org.apache.qpid.server.logging.LogRecorder;
import org.apache.qpid.server.logging.SystemOutMessageLogger;
import org.apache.qpid.server.logging.log4j.LoggingManagementFacade;
import org.apache.qpid.server.logging.messages.BrokerMessages;
-import org.apache.qpid.server.model.BrokerModel;
-import org.apache.qpid.server.model.ConfiguredObjectFactory;
-import org.apache.qpid.server.model.ConfiguredObjectFactoryImpl;
import org.apache.qpid.server.model.SystemContext;
import org.apache.qpid.server.model.SystemContextImpl;
import org.apache.qpid.server.registry.ApplicationRegistry;
@@ -61,7 +60,7 @@ public class Broker
private volatile IApplicationRegistry _applicationRegistry;
private EventLogger _eventLogger;
private boolean _configuringOwnLogging = false;
- private final TaskExecutor _taskExecutor = new TaskExecutor();
+ private final TaskExecutor _taskExecutor = new TaskExecutorImpl();
protected static class InitException extends RuntimeException
{
@@ -140,8 +139,7 @@ public class Broker
LogRecorder logRecorder = new LogRecorder();
_taskExecutor.start();
- ConfiguredObjectFactory configuredObjectFactory = new ConfiguredObjectFactoryImpl(BrokerModel.getInstance());
- SystemContext systemContext = new SystemContextImpl(_taskExecutor, configuredObjectFactory, _eventLogger, logRecorder, options);
+ SystemContext systemContext = new SystemContextImpl(_taskExecutor, _eventLogger, logRecorder, options);
BrokerConfigurationStoreCreator storeCreator = new BrokerConfigurationStoreCreator();
DurableConfigurationStore store = storeCreator.createStore(systemContext, storeType, options.getInitialConfigurationLocation(),
diff --git a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/binding/BindingImpl.java b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/binding/BindingImpl.java
index 76826bfc56..3bf90248b9 100644
--- a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/binding/BindingImpl.java
+++ b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/binding/BindingImpl.java
@@ -29,7 +29,7 @@ import java.util.concurrent.CopyOnWriteArrayList;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicLong;
-import org.apache.qpid.server.configuration.updater.TaskExecutor;
+import org.apache.qpid.server.configuration.updater.VoidTask;
import org.apache.qpid.server.exchange.AbstractExchange;
import org.apache.qpid.server.exchange.ExchangeImpl;
import org.apache.qpid.server.logging.EventLogger;
@@ -256,7 +256,7 @@ public class BindingImpl
public void setArguments(final Map<String, Object> arguments)
{
- runTask(new TaskExecutor.VoidTask()
+ runTask(new VoidTask()
{
@Override
public void execute()
diff --git a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/configuration/updater/Task.java b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/configuration/updater/Task.java
new file mode 100644
index 0000000000..9fa6fdf8fa
--- /dev/null
+++ b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/configuration/updater/Task.java
@@ -0,0 +1,26 @@
+/*
+ *
+ * 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.configuration.updater;
+
+public interface Task<X>
+{
+ X execute();
+}
diff --git a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/configuration/updater/TaskExecutor.java b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/configuration/updater/TaskExecutor.java
index 29b03ff962..97d1c17f31 100644
--- a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/configuration/updater/TaskExecutor.java
+++ b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/configuration/updater/TaskExecutor.java
@@ -20,373 +20,26 @@
*/
package org.apache.qpid.server.configuration.updater;
-import java.security.AccessController;
-import java.security.PrivilegedAction;
-import java.util.List;
-import java.util.concurrent.Callable;
import java.util.concurrent.CancellationException;
-import java.util.concurrent.ExecutionException;
-import java.util.concurrent.ExecutorService;
-import java.util.concurrent.Executors;
-import java.util.concurrent.Future;
-import java.util.concurrent.RunnableFuture;
-import java.util.concurrent.ThreadFactory;
-import java.util.concurrent.TimeUnit;
-import java.util.concurrent.atomic.AtomicReference;
-import javax.security.auth.Subject;
-
-import org.apache.log4j.Logger;
import org.apache.qpid.server.model.State;
-import org.apache.qpid.server.util.ServerScopedRuntimeException;
-public class TaskExecutor
+public interface TaskExecutor
{
- private static final String TASK_EXECUTION_THREAD_NAME = "Broker-Configuration-Thread";
- private static final Logger LOGGER = Logger.getLogger(TaskExecutor.class);
-
- private volatile Thread _taskThread;
- private final AtomicReference<State> _state;
- private volatile ExecutorService _executor;
-
- public static interface Task<X>
- {
- X execute();
- }
-
- public static interface VoidTask
- {
- void execute();
- }
-
- public static interface TaskWithException<X,E extends Exception>
- {
- X execute() throws E;
- }
-
- public static interface VoidTaskWithException<E extends Exception>
- {
- void execute() throws E;
- }
-
-
- public TaskExecutor()
- {
- _state = new AtomicReference<State>(State.INITIALISING);
- }
-
- public State getState()
- {
- return _state.get();
- }
-
- public void start()
- {
- if (_state.compareAndSet(State.INITIALISING, State.ACTIVE))
- {
- LOGGER.debug("Starting task executor");
- _executor = Executors.newFixedThreadPool(1, new ThreadFactory()
- {
- @Override
- public Thread newThread(Runnable r)
- {
- _taskThread = new Thread(r, TASK_EXECUTION_THREAD_NAME);
- return _taskThread;
- }
- });
- LOGGER.debug("Task executor is started");
- }
- }
-
- public void stopImmediately()
- {
- if (_state.compareAndSet(State.ACTIVE, State.STOPPED))
- {
- ExecutorService executor = _executor;
- if (executor != null)
- {
- LOGGER.debug("Stopping task executor immediately");
- List<Runnable> cancelledTasks = executor.shutdownNow();
- for (Runnable runnable : cancelledTasks)
- {
- if (runnable instanceof RunnableFuture<?>)
- {
- ((RunnableFuture<?>) runnable).cancel(true);
- }
- }
-
- _executor = null;
- _taskThread = null;
- LOGGER.debug("Task executor was stopped immediately. Number of unfinished tasks: " + cancelledTasks.size());
- }
- }
- }
-
- public void stop()
- {
- if (_state.compareAndSet(State.ACTIVE, State.STOPPED))
- {
- ExecutorService executor = _executor;
- if (executor != null)
- {
- LOGGER.debug("Stopping task executor");
- executor.shutdown();
- _executor = null;
- _taskThread = null;
- LOGGER.debug("Task executor is stopped");
- }
- }
- }
-
- <T> Future<T> submit(Task<T> task)
- {
- checkState();
- if (LOGGER.isDebugEnabled())
- {
- LOGGER.debug("Submitting task: " + task);
- }
- Future<T> future = null;
- if (isTaskExecutorThread())
- {
- T result = executeTask(task);
- return new ImmediateFuture(result);
- }
- else
- {
- future = _executor.submit(new CallableWrapper(task));
- }
- return future;
- }
-
- public void run(final VoidTask task) throws CancellationException
- {
- run(new Task<Void>()
- {
- @Override
- public Void execute()
- {
- task.execute();
- return null;
- }
- });
- }
-
- private static class ExceptionTaskWrapper<T, E extends Exception> implements Task<T>
- {
- private final TaskWithException<T,E> _underlying;
- private E _exception;
-
- private ExceptionTaskWrapper(final TaskWithException<T, E> underlying)
- {
- _underlying = underlying;
- }
-
-
- @Override
- public T execute()
- {
- try
- {
- return _underlying.execute();
- }
- catch (Exception e)
- {
- _exception = (E) e;
- return null;
- }
- }
-
- E getException()
- {
- return _exception;
- }
- }
-
-
- private static class ExceptionVoidTaskWrapper<E extends Exception> implements Task<Void>
- {
- private final VoidTaskWithException<E> _underlying;
- private E _exception;
-
- private ExceptionVoidTaskWrapper(final VoidTaskWithException<E> underlying)
- {
- _underlying = underlying;
- }
-
-
- @Override
- public Void execute()
- {
- try
- {
- _underlying.execute();
-
- }
- catch (Exception e)
- {
- _exception = (E) e;
- }
- return null;
- }
-
- E getException()
- {
- return _exception;
- }
- }
-
- public <T, E extends Exception> T run(TaskWithException<T,E> task) throws CancellationException, E
- {
- ExceptionTaskWrapper<T,E> wrapper = new ExceptionTaskWrapper<T, E>(task);
- T result = run(wrapper);
- if(wrapper.getException() != null)
- {
- throw wrapper.getException();
- }
- else
- {
- return result;
- }
- }
-
-
- public <E extends Exception> void run(VoidTaskWithException<E> task) throws CancellationException, E
- {
- ExceptionVoidTaskWrapper<E> wrapper = new ExceptionVoidTaskWrapper<E>(task);
- run(wrapper);
- if(wrapper.getException() != null)
- {
- throw wrapper.getException();
- }
- }
-
- public <T> T run(Task<T> task) throws CancellationException
- {
- try
- {
- Future<T> future = submit(task);
- return future.get();
- }
- catch (InterruptedException e)
- {
- throw new ServerScopedRuntimeException("Task execution was interrupted: " + task, e);
- }
- catch (ExecutionException e)
- {
- Throwable cause = e.getCause();
- if (cause instanceof RuntimeException)
- {
- throw (RuntimeException) cause;
- }
- else if (cause instanceof Exception)
- {
- throw new ServerScopedRuntimeException("Failed to execute user task: " + task, cause);
- }
- else if (cause instanceof Error)
- {
- throw (Error) cause;
- }
- else
- {
- throw new ServerScopedRuntimeException("Failed to execute user task: " + task, cause);
- }
- }
- }
-
- public boolean isTaskExecutorThread()
- {
- return Thread.currentThread() == _taskThread;
- }
-
- private void checkState()
- {
- if (_state.get() != State.ACTIVE)
- {
- throw new IllegalStateException("Task executor is not in ACTIVE state");
- }
- }
-
- private <T> T executeTask(Task<T> userTask)
- {
- if (LOGGER.isDebugEnabled())
- {
- LOGGER.debug("Performing task " + userTask);
- }
- T result = userTask.execute();
- if (LOGGER.isDebugEnabled())
- {
- LOGGER.debug("Task " + userTask + " is performed successfully with result:" + result);
- }
- return result;
- }
-
- private class CallableWrapper<T> implements Callable<T>
- {
- private Task<T> _userTask;
- private Subject _contextSubject;
-
- public CallableWrapper(Task<T> userWork)
- {
- _userTask = userWork;
- _contextSubject = Subject.getSubject(AccessController.getContext());
- }
-
- @Override
- public T call()
- {
- T result = null;
- result = Subject.doAs(_contextSubject, new PrivilegedAction<T>()
- {
- @Override
- public T run()
- {
- return executeTask(_userTask);
- }
- });
-
+ State getState();
- return result;
- }
- }
+ void start();
- private static class ImmediateFuture<T> implements Future<T>
- {
- private T _result;
+ void stopImmediately();
- public ImmediateFuture(T result)
- {
- super();
- _result = result;
- }
+ void stop();
- @Override
- public boolean cancel(boolean mayInterruptIfRunning)
- {
- return false;
- }
+ void run(VoidTask task) throws CancellationException;
- @Override
- public boolean isCancelled()
- {
- return false;
- }
+ <T, E extends Exception> T run(TaskWithException<T, E> task) throws CancellationException, E;
- @Override
- public boolean isDone()
- {
- return true;
- }
+ <E extends Exception> void run(VoidTaskWithException<E> task) throws CancellationException, E;
- @Override
- public T get()
- {
- return _result;
- }
+ <T> T run(Task<T> task) throws CancellationException;
- @Override
- public T get(long timeout, TimeUnit unit)
- {
- return get();
- }
- }
}
diff --git a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/configuration/updater/TaskExecutorImpl.java b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/configuration/updater/TaskExecutorImpl.java
new file mode 100644
index 0000000000..5eb149b4ec
--- /dev/null
+++ b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/configuration/updater/TaskExecutorImpl.java
@@ -0,0 +1,380 @@
+/*
+ *
+ * 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.configuration.updater;
+
+import java.security.AccessController;
+import java.security.PrivilegedAction;
+import java.util.List;
+import java.util.concurrent.Callable;
+import java.util.concurrent.CancellationException;
+import java.util.concurrent.ExecutionException;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.Future;
+import java.util.concurrent.RunnableFuture;
+import java.util.concurrent.ThreadFactory;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicReference;
+
+import javax.security.auth.Subject;
+
+import org.apache.log4j.Logger;
+import org.apache.qpid.server.model.State;
+import org.apache.qpid.server.util.ServerScopedRuntimeException;
+
+public class TaskExecutorImpl implements TaskExecutor
+{
+ private static final String TASK_EXECUTION_THREAD_NAME = "Broker-Configuration-Thread";
+ private static final Logger LOGGER = Logger.getLogger(TaskExecutorImpl.class);
+
+ private volatile Thread _taskThread;
+ private final AtomicReference<State> _state;
+ private volatile ExecutorService _executor;
+
+
+ public TaskExecutorImpl()
+ {
+ _state = new AtomicReference<State>(State.INITIALISING);
+ }
+
+ @Override
+ public State getState()
+ {
+ return _state.get();
+ }
+
+ @Override
+ public void start()
+ {
+ if (_state.compareAndSet(State.INITIALISING, State.ACTIVE))
+ {
+ LOGGER.debug("Starting task executor");
+ _executor = Executors.newFixedThreadPool(1, new ThreadFactory()
+ {
+ @Override
+ public Thread newThread(Runnable r)
+ {
+ _taskThread = new Thread(r, TASK_EXECUTION_THREAD_NAME);
+ return _taskThread;
+ }
+ });
+ LOGGER.debug("Task executor is started");
+ }
+ }
+
+ @Override
+ public void stopImmediately()
+ {
+ if (_state.compareAndSet(State.ACTIVE, State.STOPPED))
+ {
+ ExecutorService executor = _executor;
+ if (executor != null)
+ {
+ LOGGER.debug("Stopping task executor immediately");
+ List<Runnable> cancelledTasks = executor.shutdownNow();
+ for (Runnable runnable : cancelledTasks)
+ {
+ if (runnable instanceof RunnableFuture<?>)
+ {
+ ((RunnableFuture<?>) runnable).cancel(true);
+ }
+ }
+
+ _executor = null;
+ _taskThread = null;
+ LOGGER.debug("Task executor was stopped immediately. Number of unfinished tasks: " + cancelledTasks.size());
+ }
+ }
+ }
+
+ @Override
+ public void stop()
+ {
+ if (_state.compareAndSet(State.ACTIVE, State.STOPPED))
+ {
+ ExecutorService executor = _executor;
+ if (executor != null)
+ {
+ LOGGER.debug("Stopping task executor");
+ executor.shutdown();
+ _executor = null;
+ _taskThread = null;
+ LOGGER.debug("Task executor is stopped");
+ }
+ }
+ }
+
+ <T> Future<T> submit(Task<T> task)
+ {
+ checkState();
+ if (LOGGER.isDebugEnabled())
+ {
+ LOGGER.debug("Submitting task: " + task);
+ }
+ Future<T> future = null;
+ if (isTaskExecutorThread())
+ {
+ T result = executeTask(task);
+ return new ImmediateFuture(result);
+ }
+ else
+ {
+ future = _executor.submit(new CallableWrapper(task));
+ }
+ return future;
+ }
+
+ @Override
+ public void run(final VoidTask task) throws CancellationException
+ {
+ run(new Task<Void>()
+ {
+ @Override
+ public Void execute()
+ {
+ task.execute();
+ return null;
+ }
+ });
+ }
+
+ private static class ExceptionTaskWrapper<T, E extends Exception> implements Task<T>
+ {
+ private final TaskWithException<T,E> _underlying;
+ private E _exception;
+
+ private ExceptionTaskWrapper(final TaskWithException<T, E> underlying)
+ {
+ _underlying = underlying;
+ }
+
+
+ @Override
+ public T execute()
+ {
+ try
+ {
+ return _underlying.execute();
+ }
+ catch (Exception e)
+ {
+ _exception = (E) e;
+ return null;
+ }
+ }
+
+ E getException()
+ {
+ return _exception;
+ }
+ }
+
+
+ private static class ExceptionVoidTaskWrapper<E extends Exception> implements Task<Void>
+ {
+ private final VoidTaskWithException<E> _underlying;
+ private E _exception;
+
+ private ExceptionVoidTaskWrapper(final VoidTaskWithException<E> underlying)
+ {
+ _underlying = underlying;
+ }
+
+
+ @Override
+ public Void execute()
+ {
+ try
+ {
+ _underlying.execute();
+
+ }
+ catch (Exception e)
+ {
+ _exception = (E) e;
+ }
+ return null;
+ }
+
+ E getException()
+ {
+ return _exception;
+ }
+ }
+
+ @Override
+ public <T, E extends Exception> T run(TaskWithException<T, E> task) throws CancellationException, E
+ {
+ ExceptionTaskWrapper<T,E> wrapper = new ExceptionTaskWrapper<T, E>(task);
+ T result = run(wrapper);
+ if(wrapper.getException() != null)
+ {
+ throw wrapper.getException();
+ }
+ else
+ {
+ return result;
+ }
+ }
+
+
+ @Override
+ public <E extends Exception> void run(VoidTaskWithException<E> task) throws CancellationException, E
+ {
+ ExceptionVoidTaskWrapper<E> wrapper = new ExceptionVoidTaskWrapper<E>(task);
+ run(wrapper);
+ if(wrapper.getException() != null)
+ {
+ throw wrapper.getException();
+ }
+ }
+
+ @Override
+ public <T> T run(Task<T> task) throws CancellationException
+ {
+ try
+ {
+ Future<T> future = submit(task);
+ return future.get();
+ }
+ catch (InterruptedException e)
+ {
+ throw new ServerScopedRuntimeException("Task execution was interrupted: " + task, e);
+ }
+ catch (ExecutionException e)
+ {
+ Throwable cause = e.getCause();
+ if (cause instanceof RuntimeException)
+ {
+ throw (RuntimeException) cause;
+ }
+ else if (cause instanceof Exception)
+ {
+ throw new ServerScopedRuntimeException("Failed to execute user task: " + task, cause);
+ }
+ else if (cause instanceof Error)
+ {
+ throw (Error) cause;
+ }
+ else
+ {
+ throw new ServerScopedRuntimeException("Failed to execute user task: " + task, cause);
+ }
+ }
+ }
+
+ private boolean isTaskExecutorThread()
+ {
+ return Thread.currentThread() == _taskThread;
+ }
+
+ private void checkState()
+ {
+ if (_state.get() != State.ACTIVE)
+ {
+ throw new IllegalStateException("Task executor is not in ACTIVE state");
+ }
+ }
+
+ private <T> T executeTask(Task<T> userTask)
+ {
+ if (LOGGER.isDebugEnabled())
+ {
+ LOGGER.debug("Performing task " + userTask);
+ }
+ T result = userTask.execute();
+ if (LOGGER.isDebugEnabled())
+ {
+ LOGGER.debug("Task " + userTask + " is performed successfully with result:" + result);
+ }
+ return result;
+ }
+
+ private class CallableWrapper<T> implements Callable<T>
+ {
+ private Task<T> _userTask;
+ private Subject _contextSubject;
+
+ public CallableWrapper(Task<T> userWork)
+ {
+ _userTask = userWork;
+ _contextSubject = Subject.getSubject(AccessController.getContext());
+ }
+
+ @Override
+ public T call()
+ {
+ T result = null;
+ result = Subject.doAs(_contextSubject, new PrivilegedAction<T>()
+ {
+ @Override
+ public T run()
+ {
+ return executeTask(_userTask);
+ }
+ });
+
+
+ return result;
+ }
+ }
+
+ private static class ImmediateFuture<T> implements Future<T>
+ {
+ private T _result;
+
+ public ImmediateFuture(T result)
+ {
+ super();
+ _result = result;
+ }
+
+ @Override
+ public boolean cancel(boolean mayInterruptIfRunning)
+ {
+ return false;
+ }
+
+ @Override
+ public boolean isCancelled()
+ {
+ return false;
+ }
+
+ @Override
+ public boolean isDone()
+ {
+ return true;
+ }
+
+ @Override
+ public T get()
+ {
+ return _result;
+ }
+
+ @Override
+ public T get(long timeout, TimeUnit unit)
+ {
+ return get();
+ }
+ }
+}
diff --git a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/configuration/updater/TaskWithException.java b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/configuration/updater/TaskWithException.java
new file mode 100644
index 0000000000..793480ca1f
--- /dev/null
+++ b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/configuration/updater/TaskWithException.java
@@ -0,0 +1,26 @@
+/*
+ *
+ * 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.configuration.updater;
+
+public interface TaskWithException<X,E extends Exception>
+{
+ X execute() throws E;
+}
diff --git a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/configuration/updater/VoidTask.java b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/configuration/updater/VoidTask.java
new file mode 100644
index 0000000000..34a4769a5a
--- /dev/null
+++ b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/configuration/updater/VoidTask.java
@@ -0,0 +1,26 @@
+/*
+ *
+ * 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.configuration.updater;
+
+public interface VoidTask
+{
+ void execute();
+}
diff --git a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/configuration/updater/VoidTaskWithException.java b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/configuration/updater/VoidTaskWithException.java
new file mode 100644
index 0000000000..45cb21dce4
--- /dev/null
+++ b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/configuration/updater/VoidTaskWithException.java
@@ -0,0 +1,26 @@
+/*
+ *
+ * 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.configuration.updater;
+
+public interface VoidTaskWithException<E extends Exception>
+{
+ void execute() throws E;
+}
diff --git a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/model/AbstractConfiguredObject.java b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/model/AbstractConfiguredObject.java
index 4860d72e26..c48f505259 100644
--- a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/model/AbstractConfiguredObject.java
+++ b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/model/AbstractConfiguredObject.java
@@ -48,7 +48,11 @@ import java.util.concurrent.atomic.AtomicBoolean;
import javax.security.auth.Subject;
import org.apache.qpid.server.configuration.IllegalConfigurationException;
+import org.apache.qpid.server.configuration.updater.Task;
import org.apache.qpid.server.configuration.updater.TaskExecutor;
+import org.apache.qpid.server.configuration.updater.TaskWithException;
+import org.apache.qpid.server.configuration.updater.VoidTask;
+import org.apache.qpid.server.configuration.updater.VoidTaskWithException;
import org.apache.qpid.server.security.SecurityManager;
import org.apache.qpid.server.security.auth.AuthenticatedPrincipal;
import org.apache.qpid.server.store.ConfiguredObjectRecord;
@@ -115,7 +119,7 @@ public abstract class AbstractConfiguredObject<X extends ConfiguredObject<X>> im
private final Class<? extends ConfiguredObject> _category;
private final Class<? extends ConfiguredObject> _bestFitInterface;
- private final ConfiguredObjectFactory _objectFactory;
+ private final Model _model;
@ManagedAttributeField
private long _createdTime;
@@ -173,16 +177,16 @@ public abstract class AbstractConfiguredObject<X extends ConfiguredObject<X>> im
Map<String, Object> attributes,
TaskExecutor taskExecutor)
{
- this(parents, attributes, taskExecutor, parents.values().iterator().next().getObjectFactory());
+ this(parents, attributes, taskExecutor, parents.values().iterator().next().getModel());
}
protected AbstractConfiguredObject(final Map<Class<? extends ConfiguredObject>, ConfiguredObject<?>> parents,
Map<String, Object> attributes,
TaskExecutor taskExecutor,
- ConfiguredObjectFactory objectFactory)
+ Model model)
{
_taskExecutor = taskExecutor;
- _objectFactory = objectFactory;
+ _model = model;
_category = Model.getCategory(getClass());
@@ -533,13 +537,13 @@ public abstract class AbstractConfiguredObject<X extends ConfiguredObject<X>> im
@Override
public final ConfiguredObjectFactory getObjectFactory()
{
- return _objectFactory;
+ return _model.getObjectFactory();
}
@Override
public final Model getModel()
{
- return _objectFactory.getModel();
+ return _model;
}
public Class<? extends ConfiguredObject> getCategoryClass()
@@ -563,7 +567,7 @@ public abstract class AbstractConfiguredObject<X extends ConfiguredObject<X>> im
{
- return runTask(new TaskExecutor.Task<State>()
+ return runTask(new Task<State>()
{
@Override
public State execute()
@@ -746,7 +750,7 @@ public abstract class AbstractConfiguredObject<X extends ConfiguredObject<X>> im
public Object setAttribute(final String name, final Object expected, final Object desired)
throws IllegalStateException, AccessControlException, IllegalArgumentException
{
- return _taskExecutor.run(new TaskExecutor.Task<Object>()
+ return _taskExecutor.run(new Task<Object>()
{
@Override
public Object execute()
@@ -883,7 +887,7 @@ public abstract class AbstractConfiguredObject<X extends ConfiguredObject<X>> im
public <C extends ConfiguredObject> C createChild(final Class<C> childClass, final Map<String, Object> attributes,
final ConfiguredObject... otherParents)
{
- return _taskExecutor.run(new TaskExecutor.Task<C>() {
+ return _taskExecutor.run(new Task<C>() {
@Override
public C execute()
@@ -980,22 +984,22 @@ public abstract class AbstractConfiguredObject<X extends ConfiguredObject<X>> im
return _taskExecutor;
}
- protected final <C> C runTask(TaskExecutor.Task<C> task)
+ protected final <C> C runTask(Task<C> task)
{
return _taskExecutor.run(task);
}
- protected void runTask(TaskExecutor.VoidTask task)
+ protected void runTask(VoidTask task)
{
_taskExecutor.run(task);
}
- protected final <T, E extends Exception> T runTask(TaskExecutor.TaskWithException<T,E> task) throws E
+ protected final <T, E extends Exception> T runTask(TaskWithException<T,E> task) throws E
{
return _taskExecutor.run(task);
}
- protected final <E extends Exception> void runTask(TaskExecutor.VoidTaskWithException<E> task) throws E
+ protected final <E extends Exception> void runTask(VoidTaskWithException<E> task) throws E
{
_taskExecutor.run(task);
}
@@ -1004,7 +1008,7 @@ public abstract class AbstractConfiguredObject<X extends ConfiguredObject<X>> im
@Override
public void setAttributes(final Map<String, Object> attributes) throws IllegalStateException, AccessControlException, IllegalArgumentException
{
- runTask(new TaskExecutor.VoidTask()
+ runTask(new VoidTask()
{
@Override
public void execute()
@@ -1189,7 +1193,7 @@ public abstract class AbstractConfiguredObject<X extends ConfiguredObject<X>> im
static String interpolate(ConfiguredObject<?> object, String value)
{
Map<String,String> inheritedContext = new HashMap<String, String>();
- generateInheritedContext(object, inheritedContext);
+ generateInheritedContext(object.getModel(), object, inheritedContext);
return Strings.expand(value, false,
new Strings.MapResolver(inheritedContext),
Strings.JAVA_SYS_PROPS_RESOLVER,
@@ -1197,17 +1201,17 @@ public abstract class AbstractConfiguredObject<X extends ConfiguredObject<X>> im
new Strings.MapResolver(_defaultContext));
}
- static void generateInheritedContext(final ConfiguredObject<?> object,
- final Map<String, String> inheritedContext)
+ static void generateInheritedContext(final Model model, final ConfiguredObject<?> object,
+ final Map<String, String> inheritedContext)
{
Collection<Class<? extends ConfiguredObject>> parents =
- object.getModel().getParentTypes(object.getCategoryClass());
+ model.getParentTypes(object.getCategoryClass());
if(parents != null && !parents.isEmpty())
{
ConfiguredObject parent = object.getParent(parents.iterator().next());
if(parent != null)
{
- generateInheritedContext(parent, inheritedContext);
+ generateInheritedContext(model, parent, inheritedContext);
}
}
if(object.getContext() != null)
diff --git a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/model/BrokerModel.java b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/model/BrokerModel.java
index bcfb413451..8d742b2bbd 100644
--- a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/model/BrokerModel.java
+++ b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/model/BrokerModel.java
@@ -54,6 +54,7 @@ public final class BrokerModel extends Model
new HashSet<Class<? extends ConfiguredObject>>();
private Class<? extends ConfiguredObject> _rootCategory;
+ private final ConfiguredObjectFactory _objectFactory;
private BrokerModel()
{
@@ -97,6 +98,8 @@ public final class BrokerModel extends Model
addRelationship(Session.class, Consumer.class);
addRelationship(Session.class, Publisher.class);
+
+ _objectFactory = new ConfiguredObjectFactoryImpl(this);
}
public static Model getInstance()
@@ -129,6 +132,12 @@ public final class BrokerModel extends Model
return MODEL_MINOR_VERSION;
}
+ @Override
+ public ConfiguredObjectFactory getObjectFactory()
+ {
+ return _objectFactory;
+ }
+
public Collection<Class<? extends ConfiguredObject>> getChildTypes(Class<? extends ConfiguredObject> parent)
{
Collection<Class<? extends ConfiguredObject>> childTypes = _children.get(parent);
diff --git a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/model/Model.java b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/model/Model.java
index f1201e3cb9..f04ce311b9 100644
--- a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/model/Model.java
+++ b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/model/Model.java
@@ -120,4 +120,6 @@ public abstract class Model
public abstract int getMajorVersion();
public abstract int getMinorVersion();
+ public abstract ConfiguredObjectFactory getObjectFactory();
+
}
diff --git a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/model/SystemContextImpl.java b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/model/SystemContextImpl.java
index d9e9950972..76a7de82b3 100644
--- a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/model/SystemContextImpl.java
+++ b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/model/SystemContextImpl.java
@@ -20,11 +20,8 @@
*/
package org.apache.qpid.server.model;
-import java.util.ArrayList;
-import java.util.Arrays;
import java.util.Collection;
import java.util.HashMap;
-import java.util.Iterator;
import java.util.Map;
import java.util.UUID;
@@ -34,17 +31,10 @@ import org.apache.qpid.server.configuration.updater.TaskExecutor;
import org.apache.qpid.server.logging.EventLogger;
import org.apache.qpid.server.logging.LogRecorder;
import org.apache.qpid.server.logging.messages.BrokerMessages;
-import org.apache.qpid.server.store.ConfiguredObjectDependency;
-import org.apache.qpid.server.store.ConfiguredObjectIdDependency;
-import org.apache.qpid.server.store.ConfiguredObjectNameDependency;
-import org.apache.qpid.server.store.ConfiguredObjectRecord;
-import org.apache.qpid.server.store.UnresolvedConfiguredObject;
-import org.apache.qpid.server.util.ServerScopedRuntimeException;
public class SystemContextImpl extends AbstractConfiguredObject<SystemContextImpl> implements SystemContext<SystemContextImpl>
{
private static final UUID SYSTEM_ID = new UUID(0l, 0l);
- private final ConfiguredObjectFactory _objectFactory;
private final EventLogger _eventLogger;
private final LogRecorder _logRecorder;
private final BrokerOptions _brokerOptions;
@@ -56,17 +46,15 @@ public class SystemContextImpl extends AbstractConfiguredObject<SystemContextImp
private String _storeType;
public SystemContextImpl(final TaskExecutor taskExecutor,
- final ConfiguredObjectFactory configuredObjectFactory,
final EventLogger eventLogger,
final LogRecorder logRecorder,
final BrokerOptions brokerOptions)
{
super(parentsMap(),
createAttributes(brokerOptions),
- taskExecutor, configuredObjectFactory);
+ taskExecutor, BrokerModel.getInstance());
_eventLogger = eventLogger;
getTaskExecutor().start();
- _objectFactory = configuredObjectFactory;
_logRecorder = logRecorder;
_brokerOptions = brokerOptions;
open();
diff --git a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/model/adapter/BrokerAdapter.java b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/model/adapter/BrokerAdapter.java
index 0a3ffe4595..2927bb1491 100644
--- a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/model/adapter/BrokerAdapter.java
+++ b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/model/adapter/BrokerAdapter.java
@@ -39,7 +39,8 @@ import org.apache.log4j.Logger;
import org.apache.qpid.common.QpidProperties;
import org.apache.qpid.server.BrokerOptions;
import org.apache.qpid.server.configuration.IllegalConfigurationException;
-import org.apache.qpid.server.configuration.updater.TaskExecutor;
+import org.apache.qpid.server.configuration.updater.Task;
+import org.apache.qpid.server.configuration.updater.VoidTask;
import org.apache.qpid.server.logging.EventLogger;
import org.apache.qpid.server.logging.LogRecorder;
import org.apache.qpid.server.logging.messages.BrokerMessages;
@@ -453,7 +454,7 @@ public class BrokerAdapter extends AbstractConfiguredObject<BrokerAdapter> imple
@Override
public <C extends ConfiguredObject> C addChild(final Class<C> childClass, final Map<String, Object> attributes, final ConfiguredObject... otherParents)
{
- return runTask( new TaskExecutor.Task<C>()
+ return runTask( new Task<C>()
{
@Override
public C execute()
@@ -516,8 +517,6 @@ public class BrokerAdapter extends AbstractConfiguredObject<BrokerAdapter> imple
private void addPort(final Port<?> port)
{
- assert getTaskExecutor().isTaskExecutorThread();
-
int portNumber = port.getPort();
String portName = port.getName();
@@ -543,8 +542,6 @@ public class BrokerAdapter extends AbstractConfiguredObject<BrokerAdapter> imple
private AccessControlProvider<?> createAccessControlProvider(final Map<String, Object> attributes)
{
- assert getTaskExecutor().isTaskExecutorThread();
-
AccessControlProvider<?> accessControlProvider = (AccessControlProvider<?>) createChild(AccessControlProvider.class, attributes);
addAccessControlProvider(accessControlProvider);
@@ -557,8 +554,6 @@ public class BrokerAdapter extends AbstractConfiguredObject<BrokerAdapter> imple
private void addAccessControlProvider(final AccessControlProvider<?> accessControlProvider)
{
- assert getTaskExecutor().isTaskExecutorThread();
-
accessControlProvider.addChangeListener(this);
accessControlProvider.addChangeListener(_securityManager);
@@ -574,7 +569,7 @@ public class BrokerAdapter extends AbstractConfiguredObject<BrokerAdapter> imple
private AuthenticationProvider createAuthenticationProvider(final Map<String, Object> attributes)
{
- return runTask(new TaskExecutor.Task<AuthenticationProvider>()
+ return runTask(new Task<AuthenticationProvider>()
{
@Override
public AuthenticationProvider execute()
@@ -604,14 +599,12 @@ public class BrokerAdapter extends AbstractConfiguredObject<BrokerAdapter> imple
*/
private void addAuthenticationProvider(AuthenticationProvider<?> authenticationProvider)
{
- assert getTaskExecutor().isTaskExecutorThread();
-
authenticationProvider.addChangeListener(this);
}
private GroupProvider<?> createGroupProvider(final Map<String, Object> attributes)
{
- return runTask(new TaskExecutor.Task<GroupProvider<?>>()
+ return runTask(new Task<GroupProvider<?>>()
{
@Override
public GroupProvider<?> execute()
@@ -755,7 +748,7 @@ public class BrokerAdapter extends AbstractConfiguredObject<BrokerAdapter> imple
final State desiredState,
final boolean swallowException)
{
- runTask(new TaskExecutor.VoidTask()
+ runTask(new VoidTask()
{
@Override
public void execute()
diff --git a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/security/auth/manager/ScramSHA1AuthenticationManager.java b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/security/auth/manager/ScramSHA1AuthenticationManager.java
index 85a0d632c2..5ff7f015b4 100644
--- a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/security/auth/manager/ScramSHA1AuthenticationManager.java
+++ b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/security/auth/manager/ScramSHA1AuthenticationManager.java
@@ -43,7 +43,9 @@ import javax.security.sasl.SaslException;
import javax.security.sasl.SaslServer;
import javax.xml.bind.DatatypeConverter;
-import org.apache.qpid.server.configuration.updater.TaskExecutor;
+import org.apache.qpid.server.configuration.updater.Task;
+import org.apache.qpid.server.configuration.updater.VoidTask;
+import org.apache.qpid.server.configuration.updater.VoidTaskWithException;
import org.apache.qpid.server.model.AbstractConfiguredObject;
import org.apache.qpid.server.model.Broker;
import org.apache.qpid.server.model.ConfiguredObject;
@@ -251,7 +253,7 @@ public class ScramSHA1AuthenticationManager
@Override
public boolean createUser(final String username, final String password, final Map<String, String> attributes)
{
- return runTask(new TaskExecutor.Task<Boolean>()
+ return runTask(new Task<Boolean>()
{
@Override
public Boolean execute()
@@ -289,7 +291,7 @@ public class ScramSHA1AuthenticationManager
@Override
public void deleteUser(final String user) throws AccountNotFoundException
{
- runTask(new TaskExecutor.VoidTaskWithException<AccountNotFoundException>()
+ runTask(new VoidTaskWithException<AccountNotFoundException>()
{
@Override
public void execute() throws AccountNotFoundException
@@ -310,7 +312,7 @@ public class ScramSHA1AuthenticationManager
@Override
public void setPassword(final String username, final String password) throws AccountNotFoundException
{
- runTask(new TaskExecutor.VoidTaskWithException<AccountNotFoundException>()
+ runTask(new VoidTaskWithException<AccountNotFoundException>()
{
@Override
public void execute() throws AccountNotFoundException
@@ -333,7 +335,7 @@ public class ScramSHA1AuthenticationManager
@Override
public Map<String, Map<String, String>> getUsers()
{
- return runTask(new TaskExecutor.Task<Map<String, Map<String, String>>>()
+ return runTask(new Task<Map<String, Map<String, String>>>()
{
@Override
public Map<String, Map<String, String>> execute()
@@ -422,7 +424,7 @@ public class ScramSHA1AuthenticationManager
public void setAttributes(final Map<String, Object> attributes)
throws IllegalStateException, AccessControlException, IllegalArgumentException
{
- runTask(new TaskExecutor.VoidTask()
+ runTask(new VoidTask()
{
@Override
diff --git a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/store/GenericRecoverer.java b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/store/GenericRecoverer.java
index 01be6cf556..38492310b5 100644
--- a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/store/GenericRecoverer.java
+++ b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/store/GenericRecoverer.java
@@ -30,7 +30,8 @@ import java.util.Map;
import java.util.UUID;
import org.apache.log4j.Logger;
-import org.apache.qpid.server.configuration.updater.TaskExecutor;
+
+import org.apache.qpid.server.configuration.updater.VoidTask;
import org.apache.qpid.server.model.ConfiguredObject;
import org.apache.qpid.server.model.ConfiguredObjectFactory;
import org.apache.qpid.server.util.ServerScopedRuntimeException;
@@ -50,7 +51,7 @@ public class GenericRecoverer
public void recover(final List<ConfiguredObjectRecord> records)
{
- _parentOfRoot.getTaskExecutor().run(new TaskExecutor.VoidTask()
+ _parentOfRoot.getTaskExecutor().run(new VoidTask()
{
@Override
public void execute()
@@ -215,4 +216,4 @@ public class GenericRecoverer
}
}
-} \ No newline at end of file
+}
diff --git a/qpid/java/broker-core/src/test/java/org/apache/qpid/server/configuration/BrokerConfigurationStoreCreatorTest.java b/qpid/java/broker-core/src/test/java/org/apache/qpid/server/configuration/BrokerConfigurationStoreCreatorTest.java
index 0068cab824..748a9f6433 100644
--- a/qpid/java/broker-core/src/test/java/org/apache/qpid/server/configuration/BrokerConfigurationStoreCreatorTest.java
+++ b/qpid/java/broker-core/src/test/java/org/apache/qpid/server/configuration/BrokerConfigurationStoreCreatorTest.java
@@ -36,12 +36,12 @@ import org.codehaus.jackson.map.SerializationConfig;
import org.apache.qpid.server.BrokerOptions;
import org.apache.qpid.server.configuration.store.JsonConfigurationEntryStore;
+import org.apache.qpid.server.configuration.updater.CurrentThreadTaskExecutor;
import org.apache.qpid.server.configuration.updater.TaskExecutor;
import org.apache.qpid.server.logging.EventLogger;
import org.apache.qpid.server.logging.LogRecorder;
import org.apache.qpid.server.model.Broker;
import org.apache.qpid.server.model.BrokerModel;
-import org.apache.qpid.server.model.ConfiguredObjectFactoryImpl;
import org.apache.qpid.server.model.SystemContext;
import org.apache.qpid.server.model.SystemContextImpl;
import org.apache.qpid.test.utils.QpidTestCase;
@@ -70,13 +70,12 @@ public class BrokerConfigurationStoreCreatorTest extends QpidTestCase
_userStoreLocation = new File(TMP_FOLDER, "_store_" + System.currentTimeMillis() + "_" + getTestName());
final BrokerOptions brokerOptions = mock(BrokerOptions.class);
when(brokerOptions.getConfigurationStoreLocation()).thenReturn(_userStoreLocation.getAbsolutePath());
- _taskExecutor = new TaskExecutor();
+ _taskExecutor = new CurrentThreadTaskExecutor();
_taskExecutor.start();
_systemContext = new SystemContextImpl(_taskExecutor,
- new ConfiguredObjectFactoryImpl(BrokerModel.getInstance()),
- mock(EventLogger.class),
- mock(LogRecorder.class),
- brokerOptions);
+ mock(EventLogger.class),
+ mock(LogRecorder.class),
+ brokerOptions);
}
public void tearDown() throws Exception
diff --git a/qpid/java/broker-core/src/test/java/org/apache/qpid/server/configuration/store/ConfigurationEntryStoreTestCase.java b/qpid/java/broker-core/src/test/java/org/apache/qpid/server/configuration/store/ConfigurationEntryStoreTestCase.java
index 691a4f9abf..85e2e08129 100644
--- a/qpid/java/broker-core/src/test/java/org/apache/qpid/server/configuration/store/ConfigurationEntryStoreTestCase.java
+++ b/qpid/java/broker-core/src/test/java/org/apache/qpid/server/configuration/store/ConfigurationEntryStoreTestCase.java
@@ -31,14 +31,13 @@ import java.util.UUID;
import org.apache.qpid.server.BrokerOptions;
import org.apache.qpid.server.configuration.ConfigurationEntry;
import org.apache.qpid.server.configuration.ConfigurationEntryImpl;
+import org.apache.qpid.server.configuration.updater.CurrentThreadTaskExecutor;
import org.apache.qpid.server.configuration.updater.TaskExecutor;
import org.apache.qpid.server.logging.EventLogger;
import org.apache.qpid.server.logging.LogRecorder;
import org.apache.qpid.server.model.AuthenticationProvider;
import org.apache.qpid.server.model.Broker;
-import org.apache.qpid.server.model.BrokerModel;
import org.apache.qpid.server.model.ConfiguredObject;
-import org.apache.qpid.server.model.ConfiguredObjectFactoryImpl;
import org.apache.qpid.server.model.GroupProvider;
import org.apache.qpid.server.model.KeyStore;
import org.apache.qpid.server.model.Port;
@@ -76,10 +75,10 @@ public abstract class ConfigurationEntryStoreTestCase extends QpidTestCase
super.setUp();
- _taskExecutor = new TaskExecutor();
+ _taskExecutor = new CurrentThreadTaskExecutor();
_taskExecutor.start();
- _systemContext = new SystemContextImpl(_taskExecutor, new ConfiguredObjectFactoryImpl(BrokerModel.getInstance()),
+ _systemContext = new SystemContextImpl(_taskExecutor,
mock(EventLogger.class), mock(LogRecorder.class),
new BrokerOptions());
diff --git a/qpid/java/broker-core/src/test/java/org/apache/qpid/server/configuration/store/JsonConfigurationEntryStoreTest.java b/qpid/java/broker-core/src/test/java/org/apache/qpid/server/configuration/store/JsonConfigurationEntryStoreTest.java
index 62697548d5..3563444f63 100644
--- a/qpid/java/broker-core/src/test/java/org/apache/qpid/server/configuration/store/JsonConfigurationEntryStoreTest.java
+++ b/qpid/java/broker-core/src/test/java/org/apache/qpid/server/configuration/store/JsonConfigurationEntryStoreTest.java
@@ -45,8 +45,6 @@ import org.apache.qpid.server.configuration.IllegalConfigurationException;
import org.apache.qpid.server.logging.EventLogger;
import org.apache.qpid.server.logging.LogRecorder;
import org.apache.qpid.server.model.Broker;
-import org.apache.qpid.server.model.BrokerModel;
-import org.apache.qpid.server.model.ConfiguredObjectFactoryImpl;
import org.apache.qpid.server.model.PreferencesProvider;
import org.apache.qpid.server.model.SystemContext;
import org.apache.qpid.server.model.SystemContextImpl;
@@ -100,10 +98,9 @@ public class JsonConfigurationEntryStoreTest extends ConfigurationEntryStoreTest
final BrokerOptions brokerOptions = mock(BrokerOptions.class);
when(brokerOptions.getConfigurationStoreLocation()).thenReturn(absolutePath);
SystemContext context = new SystemContextImpl(getTaskExecutor(),
- new ConfiguredObjectFactoryImpl(BrokerModel.getInstance()),
- mock(EventLogger.class),
- mock(LogRecorder.class),
- brokerOptions);
+ mock(EventLogger.class),
+ mock(LogRecorder.class),
+ brokerOptions);
JsonConfigurationEntryStore store = new JsonConfigurationEntryStore(context, initialStore, false,
Collections.<String,String>emptyMap());
diff --git a/qpid/java/broker-core/src/test/java/org/apache/qpid/server/configuration/store/ManagementModeStoreHandlerTest.java b/qpid/java/broker-core/src/test/java/org/apache/qpid/server/configuration/store/ManagementModeStoreHandlerTest.java
index b569825b46..45290d506d 100644
--- a/qpid/java/broker-core/src/test/java/org/apache/qpid/server/configuration/store/ManagementModeStoreHandlerTest.java
+++ b/qpid/java/broker-core/src/test/java/org/apache/qpid/server/configuration/store/ManagementModeStoreHandlerTest.java
@@ -40,12 +40,11 @@ import org.mockito.stubbing.Answer;
import org.apache.qpid.server.BrokerOptions;
import org.apache.qpid.server.configuration.IllegalConfigurationException;
+import org.apache.qpid.server.configuration.updater.CurrentThreadTaskExecutor;
import org.apache.qpid.server.configuration.updater.TaskExecutor;
import org.apache.qpid.server.logging.EventLogger;
import org.apache.qpid.server.logging.LogRecorder;
import org.apache.qpid.server.model.Broker;
-import org.apache.qpid.server.model.BrokerModel;
-import org.apache.qpid.server.model.ConfiguredObjectFactoryImpl;
import org.apache.qpid.server.model.Port;
import org.apache.qpid.server.model.Protocol;
import org.apache.qpid.server.model.State;
@@ -75,11 +74,11 @@ public class ManagementModeStoreHandlerTest extends QpidTestCase
_rootId = UUID.randomUUID();
_portEntryId = UUID.randomUUID();
_store = mock(DurableConfigurationStore.class);
- _taskExecutor = new TaskExecutor();
+ _taskExecutor = new CurrentThreadTaskExecutor();
_taskExecutor.start();
- _systemContext = new SystemContextImpl(_taskExecutor, new ConfiguredObjectFactoryImpl(BrokerModel.getInstance()), mock(
- EventLogger.class), mock(LogRecorder.class), new BrokerOptions());
+ _systemContext = new SystemContextImpl(_taskExecutor, mock(EventLogger.class),
+ mock(LogRecorder.class), new BrokerOptions());
ConfiguredObjectRecord systemContextRecord = _systemContext.asObjectRecord();
diff --git a/qpid/java/broker-core/src/test/java/org/apache/qpid/server/configuration/updater/CurrentThreadTaskExecutor.java b/qpid/java/broker-core/src/test/java/org/apache/qpid/server/configuration/updater/CurrentThreadTaskExecutor.java
new file mode 100644
index 0000000000..001a14e11e
--- /dev/null
+++ b/qpid/java/broker-core/src/test/java/org/apache/qpid/server/configuration/updater/CurrentThreadTaskExecutor.java
@@ -0,0 +1,99 @@
+/*
+ *
+ * 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.configuration.updater;
+
+import java.util.concurrent.CancellationException;
+import java.util.concurrent.atomic.AtomicReference;
+
+import org.apache.qpid.server.model.State;
+
+public class CurrentThreadTaskExecutor implements TaskExecutor
+{
+ private final AtomicReference<Thread> _thread = new AtomicReference<>();
+ private State _state;
+
+ @Override
+ public State getState()
+ {
+ return null;
+ }
+
+ @Override
+ public void start()
+ {
+ if(!_thread.compareAndSet(null, Thread.currentThread()))
+ {
+ checkThread();
+ }
+ _state = State.ACTIVE;
+ }
+
+ @Override
+ public void stopImmediately()
+ {
+ checkThread();
+ _state = State.STOPPED;
+
+ }
+
+ private void checkThread()
+ {
+ if(_thread.get() != Thread.currentThread())
+ {
+ throw new IllegalArgumentException("Can only access the thread executor from a single thread");
+ }
+ }
+
+ @Override
+ public void stop()
+ {
+ stopImmediately();
+ }
+
+ @Override
+ public void run(final VoidTask task) throws CancellationException
+ {
+ checkThread();
+ task.execute();
+ }
+
+ @Override
+ public <T, E extends Exception> T run(final TaskWithException<T, E> task) throws CancellationException, E
+ {
+ checkThread();
+ return task.execute();
+ }
+
+ @Override
+ public <E extends Exception> void run(final VoidTaskWithException<E> task) throws CancellationException, E
+ {
+ checkThread();
+ task.execute();
+ }
+
+ @Override
+ public <T> T run(final Task<T> task) throws CancellationException
+ {
+ checkThread();
+ return task.execute();
+ }
+
+}
diff --git a/qpid/java/broker-core/src/test/java/org/apache/qpid/server/configuration/updater/TaskExecutorTest.java b/qpid/java/broker-core/src/test/java/org/apache/qpid/server/configuration/updater/TaskExecutorTest.java
index fb331c73c1..04016d91bc 100644
--- a/qpid/java/broker-core/src/test/java/org/apache/qpid/server/configuration/updater/TaskExecutorTest.java
+++ b/qpid/java/broker-core/src/test/java/org/apache/qpid/server/configuration/updater/TaskExecutorTest.java
@@ -40,12 +40,12 @@ import org.apache.qpid.server.util.ServerScopedRuntimeException;
public class TaskExecutorTest extends TestCase
{
- private TaskExecutor _executor;
+ private TaskExecutorImpl _executor;
protected void setUp() throws Exception
{
super.setUp();
- _executor = new TaskExecutor();
+ _executor = new TaskExecutorImpl();
}
protected void tearDown() throws Exception
@@ -135,7 +135,7 @@ public class TaskExecutorTest extends TestCase
public void testSubmitAndWait() throws Exception
{
_executor.start();
- Object result = _executor.run(new TaskExecutor.Task<Object>()
+ Object result = _executor.run(new Task<Object>()
{
@Override
public String execute()
@@ -188,7 +188,7 @@ public class TaskExecutorTest extends TestCase
_executor.start();
try
{
- _executor.run(new TaskExecutor.Task<Object>()
+ _executor.run(new Task<Object>()
{
@Override
@@ -215,7 +215,7 @@ public class TaskExecutorTest extends TestCase
@Override
public Object run()
{
- _executor.run(new TaskExecutor.Task<Object>()
+ _executor.run(new Task<Object>()
{
@Override
public Void execute()
@@ -231,7 +231,7 @@ public class TaskExecutorTest extends TestCase
assertEquals("Unexpected security manager subject", subject, taskSubject.get());
}
- private class SubjectRetriever implements TaskExecutor.Task<Subject>
+ private class SubjectRetriever implements Task<Subject>
{
@Override
public Subject execute()
@@ -240,7 +240,7 @@ public class TaskExecutorTest extends TestCase
}
}
- private class NeverEndingCallable implements TaskExecutor.Task<Void>
+ private class NeverEndingCallable implements Task<Void>
{
private CountDownLatch _waitLatch;
diff --git a/qpid/java/broker-core/src/test/java/org/apache/qpid/server/exchange/FanoutExchangeTest.java b/qpid/java/broker-core/src/test/java/org/apache/qpid/server/exchange/FanoutExchangeTest.java
index 8ccd7af799..98c3e7ed8d 100644
--- a/qpid/java/broker-core/src/test/java/org/apache/qpid/server/exchange/FanoutExchangeTest.java
+++ b/qpid/java/broker-core/src/test/java/org/apache/qpid/server/exchange/FanoutExchangeTest.java
@@ -36,13 +36,13 @@ import org.mockito.invocation.InvocationOnMock;
import org.mockito.stubbing.Answer;
import org.apache.qpid.common.AMQPFilterTypes;
+import org.apache.qpid.server.configuration.updater.CurrentThreadTaskExecutor;
import org.apache.qpid.server.configuration.updater.TaskExecutor;
import org.apache.qpid.server.logging.EventLogger;
import org.apache.qpid.server.message.AMQMessageHeader;
import org.apache.qpid.server.message.InstanceProperties;
import org.apache.qpid.server.message.ServerMessage;
import org.apache.qpid.server.model.BrokerModel;
-import org.apache.qpid.server.model.ConfiguredObjectFactoryImpl;
import org.apache.qpid.server.model.Exchange;
import org.apache.qpid.server.model.Queue;
import org.apache.qpid.server.queue.AMQQueue;
@@ -55,7 +55,6 @@ public class FanoutExchangeTest extends TestCase
private FanoutExchange _exchange;
private VirtualHostImpl _virtualHost;
private TaskExecutor _taskExecutor;
- private ConfiguredObjectFactoryImpl _objectFactory;
public void setUp()
{
@@ -64,15 +63,14 @@ public class FanoutExchangeTest extends TestCase
attributes.put(Exchange.NAME, "test");
attributes.put(Exchange.DURABLE, false);
- _taskExecutor = new TaskExecutor();
+ _taskExecutor = new CurrentThreadTaskExecutor();
_taskExecutor.start();
_virtualHost = mock(VirtualHostImpl.class);
SecurityManager securityManager = mock(SecurityManager.class);
when(_virtualHost.getSecurityManager()).thenReturn(securityManager);
when(_virtualHost.getEventLogger()).thenReturn(new EventLogger());
when(_virtualHost.getTaskExecutor()).thenReturn(_taskExecutor);
- _objectFactory = new ConfiguredObjectFactoryImpl(BrokerModel.getInstance());
- when(_virtualHost.getObjectFactory()).thenReturn(_objectFactory);
+ when(_virtualHost.getModel()).thenReturn(BrokerModel.getInstance());
_exchange = new FanoutExchange(attributes, _virtualHost);
_exchange.open();
}
@@ -134,8 +132,6 @@ public class FanoutExchangeTest extends TestCase
AMQQueue queue = mock(AMQQueue.class);
when(queue.getVirtualHost()).thenReturn(_virtualHost);
when(queue.getCategoryClass()).thenReturn(Queue.class);
- when(queue.getObjectFactory()).thenReturn(_objectFactory);
- when(queue.getModel()).thenReturn(_objectFactory.getModel());
return queue;
}
diff --git a/qpid/java/broker-core/src/test/java/org/apache/qpid/server/exchange/HeadersBindingTest.java b/qpid/java/broker-core/src/test/java/org/apache/qpid/server/exchange/HeadersBindingTest.java
index ec901e5067..5723a004d9 100644
--- a/qpid/java/broker-core/src/test/java/org/apache/qpid/server/exchange/HeadersBindingTest.java
+++ b/qpid/java/broker-core/src/test/java/org/apache/qpid/server/exchange/HeadersBindingTest.java
@@ -36,7 +36,6 @@ import org.apache.qpid.server.logging.EventLogger;
import org.apache.qpid.server.message.AMQMessageHeader;
import org.apache.qpid.server.model.Binding;
import org.apache.qpid.server.model.BrokerModel;
-import org.apache.qpid.server.model.ConfiguredObjectFactoryImpl;
import org.apache.qpid.server.plugin.ExchangeType;
import org.apache.qpid.server.queue.AMQQueue;
import org.apache.qpid.server.virtualhost.VirtualHostImpl;
@@ -46,7 +45,6 @@ import org.apache.qpid.server.virtualhost.VirtualHostImpl;
public class HeadersBindingTest extends TestCase
{
- private ConfiguredObjectFactoryImpl _objectFactory;
private class MockHeader implements AMQMessageHeader
{
@@ -150,17 +148,17 @@ public class HeadersBindingTest extends TestCase
{
_count++;
_queue = mock(AMQQueue.class);
- _objectFactory = new ConfiguredObjectFactoryImpl(BrokerModel.getInstance());
+
VirtualHostImpl vhost = mock(VirtualHostImpl.class);
when(_queue.getVirtualHost()).thenReturn(vhost);
- when(_queue.getObjectFactory()).thenReturn(_objectFactory);
+ when(_queue.getModel()).thenReturn(BrokerModel.getInstance());
when(vhost.getSecurityManager()).thenReturn(mock(org.apache.qpid.server.security.SecurityManager.class));
final EventLogger eventLogger = new EventLogger();
when(vhost.getEventLogger()).thenReturn(eventLogger);
_exchange = mock(ExchangeImpl.class);
when(_exchange.getExchangeType()).thenReturn(mock(ExchangeType.class));
when(_exchange.getEventLogger()).thenReturn(eventLogger);
- when(_exchange.getObjectFactory()).thenReturn(_objectFactory);
+ when(_exchange.getModel()).thenReturn(BrokerModel.getInstance());
}
protected String getQueueName()
diff --git a/qpid/java/broker-core/src/test/java/org/apache/qpid/server/exchange/HeadersExchangeTest.java b/qpid/java/broker-core/src/test/java/org/apache/qpid/server/exchange/HeadersExchangeTest.java
index 86c5e23a0f..a5d2fbad57 100644
--- a/qpid/java/broker-core/src/test/java/org/apache/qpid/server/exchange/HeadersExchangeTest.java
+++ b/qpid/java/broker-core/src/test/java/org/apache/qpid/server/exchange/HeadersExchangeTest.java
@@ -39,6 +39,7 @@ import org.mockito.invocation.InvocationOnMock;
import org.mockito.stubbing.Answer;
import org.apache.qpid.common.AMQPFilterTypes;
+import org.apache.qpid.server.configuration.updater.CurrentThreadTaskExecutor;
import org.apache.qpid.server.configuration.updater.TaskExecutor;
import org.apache.qpid.server.logging.EventLogger;
import org.apache.qpid.server.message.AMQMessageHeader;
@@ -66,7 +67,7 @@ public class HeadersExchangeTest extends TestCase
{
super.setUp();
- _taskExecutor = new TaskExecutor();
+ _taskExecutor = new CurrentThreadTaskExecutor();
_taskExecutor.start();
_virtualHost = mock(VirtualHostImpl.class);
SecurityManager securityManager = mock(SecurityManager.class);
diff --git a/qpid/java/broker-core/src/test/java/org/apache/qpid/server/model/VirtualHostTest.java b/qpid/java/broker-core/src/test/java/org/apache/qpid/server/model/VirtualHostTest.java
index 3f4c3e9842..12cdb564ea 100644
--- a/qpid/java/broker-core/src/test/java/org/apache/qpid/server/model/VirtualHostTest.java
+++ b/qpid/java/broker-core/src/test/java/org/apache/qpid/server/model/VirtualHostTest.java
@@ -28,6 +28,7 @@ import java.util.HashMap;
import java.util.Map;
import java.util.UUID;
+import org.apache.qpid.server.configuration.updater.CurrentThreadTaskExecutor;
import org.apache.qpid.server.configuration.updater.TaskExecutor;
import org.apache.qpid.server.store.DurableConfigurationStore;
import org.apache.qpid.server.store.MessageStore;
@@ -48,7 +49,7 @@ public class VirtualHostTest extends QpidTestCase
super.setUp();
_broker = BrokerTestHelper.createBrokerMock();
- _taskExecutor = new TaskExecutor();
+ _taskExecutor = new CurrentThreadTaskExecutor();
_taskExecutor.start();
when(_broker.getTaskExecutor()).thenReturn(_taskExecutor);
diff --git a/qpid/java/broker-core/src/test/java/org/apache/qpid/server/model/adapter/FileSystemPreferencesProviderTest.java b/qpid/java/broker-core/src/test/java/org/apache/qpid/server/model/adapter/FileSystemPreferencesProviderTest.java
index fee31b3369..84c8498181 100644
--- a/qpid/java/broker-core/src/test/java/org/apache/qpid/server/model/adapter/FileSystemPreferencesProviderTest.java
+++ b/qpid/java/broker-core/src/test/java/org/apache/qpid/server/model/adapter/FileSystemPreferencesProviderTest.java
@@ -31,6 +31,7 @@ import java.util.Map;
import java.util.Set;
import java.util.UUID;
+import org.apache.qpid.server.configuration.updater.CurrentThreadTaskExecutor;
import org.apache.qpid.server.configuration.updater.TaskExecutor;
import org.apache.qpid.server.model.AuthenticationProvider;
import org.apache.qpid.server.model.Broker;
@@ -62,7 +63,7 @@ public class FileSystemPreferencesProviderTest extends QpidTestCase
_preferencesFile = TestFileUtils.createTempFile(this, ".prefs.json", TEST_PREFERENCES);
_broker = BrokerTestHelper.createBrokerMock();
- _taskExecutor = new TaskExecutor();
+ _taskExecutor = new CurrentThreadTaskExecutor();
_taskExecutor.start();
when(_broker.getTaskExecutor()).thenReturn(_taskExecutor);
diff --git a/qpid/java/broker-core/src/test/java/org/apache/qpid/server/security/auth/manager/ScramSHA1AuthenticationManagerTest.java b/qpid/java/broker-core/src/test/java/org/apache/qpid/server/security/auth/manager/ScramSHA1AuthenticationManagerTest.java
index 743e310414..2ef294b31d 100644
--- a/qpid/java/broker-core/src/test/java/org/apache/qpid/server/security/auth/manager/ScramSHA1AuthenticationManagerTest.java
+++ b/qpid/java/broker-core/src/test/java/org/apache/qpid/server/security/auth/manager/ScramSHA1AuthenticationManagerTest.java
@@ -30,6 +30,7 @@ import java.util.UUID;
import javax.security.auth.login.AccountNotFoundException;
+import org.apache.qpid.server.configuration.updater.CurrentThreadTaskExecutor;
import org.apache.qpid.server.configuration.updater.TaskExecutor;
import org.apache.qpid.server.model.AuthenticationProvider;
import org.apache.qpid.server.model.Broker;
@@ -51,7 +52,7 @@ public class ScramSHA1AuthenticationManagerTest extends QpidTestCase
public void setUp() throws Exception
{
super.setUp();
- _executor = new TaskExecutor();
+ _executor = new CurrentThreadTaskExecutor();
_executor.start();
_broker = BrokerTestHelper.createBrokerMock();
_securityManager = mock(SecurityManager.class);
diff --git a/qpid/java/broker-core/src/test/java/org/apache/qpid/server/store/BrokerRecovererTest.java b/qpid/java/broker-core/src/test/java/org/apache/qpid/server/store/BrokerRecovererTest.java
index 26aa99a481..36c5eec108 100644
--- a/qpid/java/broker-core/src/test/java/org/apache/qpid/server/store/BrokerRecovererTest.java
+++ b/qpid/java/broker-core/src/test/java/org/apache/qpid/server/store/BrokerRecovererTest.java
@@ -33,6 +33,7 @@ import junit.framework.TestCase;
import org.apache.qpid.server.BrokerOptions;
import org.apache.qpid.server.configuration.IllegalConfigurationException;
+import org.apache.qpid.server.configuration.updater.CurrentThreadTaskExecutor;
import org.apache.qpid.server.configuration.updater.TaskExecutor;
import org.apache.qpid.server.logging.EventLogger;
import org.apache.qpid.server.logging.LogRecorder;
@@ -40,16 +41,10 @@ import org.apache.qpid.server.model.AuthenticationProvider;
import org.apache.qpid.server.model.Broker;
import org.apache.qpid.server.model.BrokerModel;
import org.apache.qpid.server.model.ConfiguredObject;
-import org.apache.qpid.server.model.ConfiguredObjectFactory;
-import org.apache.qpid.server.model.ConfiguredObjectFactoryImpl;
import org.apache.qpid.server.model.GroupProvider;
import org.apache.qpid.server.model.Port;
import org.apache.qpid.server.model.SystemContext;
import org.apache.qpid.server.model.SystemContextImpl;
-import org.apache.qpid.server.store.ConfiguredObjectRecord;
-import org.apache.qpid.server.store.ConfiguredObjectRecordImpl;
-import org.apache.qpid.server.store.GenericRecoverer;
-import org.apache.qpid.server.store.UnresolvedConfiguredObject;
public class BrokerRecovererTest extends TestCase
{
@@ -59,7 +54,6 @@ public class BrokerRecovererTest extends TestCase
private AuthenticationProvider<?> _authenticationProvider1;
private UUID _authenticationProvider1Id = UUID.randomUUID();
private SystemContext<?> _systemContext;
- private ConfiguredObjectFactory _configuredObjectFactory;
private TaskExecutor _taskExecutor;
@Override
@@ -67,11 +61,10 @@ public class BrokerRecovererTest extends TestCase
{
super.setUp();
- _configuredObjectFactory = new ConfiguredObjectFactoryImpl(BrokerModel.getInstance());
- _taskExecutor = new TaskExecutor();
+ _taskExecutor = new CurrentThreadTaskExecutor();
_taskExecutor.start();
_systemContext = new SystemContextImpl(_taskExecutor,
- _configuredObjectFactory, mock(EventLogger.class), mock(LogRecorder.class), mock(BrokerOptions.class));
+ mock(EventLogger.class), mock(LogRecorder.class), mock(BrokerOptions.class));
when(_brokerEntry.getId()).thenReturn(_brokerId);
when(_brokerEntry.getType()).thenReturn(Broker.class.getSimpleName());
@@ -285,7 +278,7 @@ public class BrokerRecovererTest extends TestCase
try
{
UnresolvedConfiguredObject<? extends ConfiguredObject> recover =
- _configuredObjectFactory.recover(_brokerEntry, _systemContext);
+ _systemContext.getObjectFactory().recover(_brokerEntry, _systemContext);
Broker<?> broker = (Broker<?>) recover.resolve();
broker.open();
@@ -312,7 +305,7 @@ public class BrokerRecovererTest extends TestCase
try
{
UnresolvedConfiguredObject<? extends ConfiguredObject> recover =
- _configuredObjectFactory.recover(_brokerEntry, _systemContext);
+ _systemContext.getObjectFactory().recover(_brokerEntry, _systemContext);
Broker<?> broker = (Broker<?>) recover.resolve();
broker.open();
fail("The broker creation should fail due to unsupported model version");
diff --git a/qpid/java/broker-core/src/test/java/org/apache/qpid/server/util/BrokerTestHelper.java b/qpid/java/broker-core/src/test/java/org/apache/qpid/server/util/BrokerTestHelper.java
index da03509524..1584f0fb79 100644
--- a/qpid/java/broker-core/src/test/java/org/apache/qpid/server/util/BrokerTestHelper.java
+++ b/qpid/java/broker-core/src/test/java/org/apache/qpid/server/util/BrokerTestHelper.java
@@ -32,12 +32,11 @@ import java.util.UUID;
import javax.security.auth.Subject;
-import org.apache.qpid.server.BrokerOptions;
import org.apache.qpid.server.configuration.store.JsonConfigurationEntryStore;
+import org.apache.qpid.server.configuration.updater.CurrentThreadTaskExecutor;
import org.apache.qpid.server.configuration.updater.TaskExecutor;
import org.apache.qpid.server.exchange.ExchangeImpl;
import org.apache.qpid.server.logging.EventLogger;
-import org.apache.qpid.server.logging.LogRecorder;
import org.apache.qpid.server.model.Broker;
import org.apache.qpid.server.model.BrokerModel;
import org.apache.qpid.server.model.ConfiguredObjectFactory;
@@ -46,7 +45,6 @@ import org.apache.qpid.server.model.Exchange;
import org.apache.qpid.server.model.Queue;
import org.apache.qpid.server.model.State;
import org.apache.qpid.server.model.SystemContext;
-import org.apache.qpid.server.model.SystemContextImpl;
import org.apache.qpid.server.model.UUIDGenerator;
import org.apache.qpid.server.model.VirtualHost;
import org.apache.qpid.server.model.VirtualHostNode;
@@ -69,7 +67,7 @@ public class BrokerTestHelper
protected static final String BROKER_STORE_CLASS_NAME_KEY = "brokerstore.class.name";
protected static final String JSON_BROKER_STORE_CLASS_NAME = JsonConfigurationEntryStore.class.getName();
- private static final TaskExecutor TASK_EXECUTOR = new TaskExecutor();
+ private static final TaskExecutor TASK_EXECUTOR = new CurrentThreadTaskExecutor();
static
{
TASK_EXECUTOR.start();
@@ -116,14 +114,8 @@ public class BrokerTestHelper
throws Exception
{
- //VirtualHostFactory factory = new PluggableFactoryLoader<VirtualHostFactory>(VirtualHostFactory.class).get(hostType);
Broker<?> broker = createBrokerMock();
ConfiguredObjectFactory objectFactory = broker.getObjectFactory();
- SystemContext systemContext = new SystemContextImpl(TASK_EXECUTOR,
- objectFactory,
- mock(EventLogger.class),
- mock(LogRecorder.class),
- new BrokerOptions());
when(broker.getTaskExecutor()).thenReturn(TASK_EXECUTOR);
VirtualHostNode<?> virtualHostNode = mock(VirtualHostNode.class);
diff --git a/qpid/java/broker-core/src/test/java/org/apache/qpid/server/virtualhost/VirtualHostQueueCreationTest.java b/qpid/java/broker-core/src/test/java/org/apache/qpid/server/virtualhost/VirtualHostQueueCreationTest.java
index 37d8c2ca8c..05ba456b89 100644
--- a/qpid/java/broker-core/src/test/java/org/apache/qpid/server/virtualhost/VirtualHostQueueCreationTest.java
+++ b/qpid/java/broker-core/src/test/java/org/apache/qpid/server/virtualhost/VirtualHostQueueCreationTest.java
@@ -30,14 +30,15 @@ import java.util.UUID;
import org.apache.qpid.exchange.ExchangeDefaults;
import org.apache.qpid.server.configuration.BrokerProperties;
+import org.apache.qpid.server.configuration.updater.CurrentThreadTaskExecutor;
import org.apache.qpid.server.configuration.updater.TaskExecutor;
import org.apache.qpid.server.exchange.ExchangeImpl;
import org.apache.qpid.server.logging.EventLogger;
import org.apache.qpid.server.model.Broker;
import org.apache.qpid.server.model.BrokerModel;
-import org.apache.qpid.server.model.Exchange;
import org.apache.qpid.server.model.ConfiguredObjectFactory;
import org.apache.qpid.server.model.ConfiguredObjectFactoryImpl;
+import org.apache.qpid.server.model.Exchange;
import org.apache.qpid.server.model.LifetimePolicy;
import org.apache.qpid.server.model.Queue;
import org.apache.qpid.server.model.State;
@@ -70,7 +71,7 @@ public class VirtualHostQueueCreationTest extends QpidTestCase
SecurityManager securityManager = mock(SecurityManager.class);
ConfiguredObjectFactory objectFactory = new ConfiguredObjectFactoryImpl(BrokerModel.getInstance());
- _taskExecutor = new TaskExecutor();
+ _taskExecutor = new CurrentThreadTaskExecutor();
_taskExecutor.start();
SystemContext<?> context = mock(SystemContext.class);
diff --git a/qpid/java/broker-core/src/test/java/org/apache/qpid/server/virtualhostnode/AbstractStandardVirtualHostNodeTest.java b/qpid/java/broker-core/src/test/java/org/apache/qpid/server/virtualhostnode/AbstractStandardVirtualHostNodeTest.java
index 3154a1e524..dc13c24f5d 100644
--- a/qpid/java/broker-core/src/test/java/org/apache/qpid/server/virtualhostnode/AbstractStandardVirtualHostNodeTest.java
+++ b/qpid/java/broker-core/src/test/java/org/apache/qpid/server/virtualhostnode/AbstractStandardVirtualHostNodeTest.java
@@ -27,6 +27,7 @@ import java.util.HashMap;
import java.util.Map;
import java.util.UUID;
+import org.apache.qpid.server.configuration.updater.CurrentThreadTaskExecutor;
import org.apache.qpid.server.configuration.updater.TaskExecutor;
import org.apache.qpid.server.model.Broker;
import org.apache.qpid.server.model.BrokerModel;
@@ -69,7 +70,7 @@ public class AbstractStandardVirtualHostNodeTest extends QpidTestCase
SystemContext<?> systemContext = _broker.getParent(SystemContext.class);
when(systemContext.getObjectFactory()).thenReturn(new ConfiguredObjectFactoryImpl(mock(Model.class)));
- _taskExecutor = new TaskExecutor();
+ _taskExecutor = new CurrentThreadTaskExecutor();
_taskExecutor.start();
when(_broker.getTaskExecutor()).thenReturn(_taskExecutor);
}