diff options
| author | Robert Godfrey <rgodfrey@apache.org> | 2014-04-25 10:12:06 +0000 |
|---|---|---|
| committer | Robert Godfrey <rgodfrey@apache.org> | 2014-04-25 10:12:06 +0000 |
| commit | acf84ebf5462342656a257ed978238c23fb1c900 (patch) | |
| tree | 89895d4bb614e1e7f0fe4f692ba05fa7057ff961 /qpid/java | |
| parent | 4eddea8954ba9342ab2bc35e495baa673a6015db (diff) | |
| download | qpid-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')
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); } |
