001 /**
002 * Copyright (C) 2011-2012 Barchart, Inc. <http://www.barchart.com/>
003 *
004 * All rights reserved. Licensed under the OSI BSD License.
005 *
006 * http://www.opensource.org/licenses/bsd-license.php
007 */
008 package com.barchart.util.concurrent;
009
010 import java.util.List;
011 import java.util.concurrent.CancellationException;
012 import java.util.concurrent.CopyOnWriteArrayList;
013 import java.util.concurrent.ExecutionException;
014 import java.util.concurrent.Future;
015 import java.util.concurrent.TimeUnit;
016 import java.util.concurrent.TimeoutException;
017 import java.util.concurrent.locks.AbstractQueuedSynchronizer;
018 import java.util.concurrent.locks.Lock;
019 import java.util.concurrent.locks.ReentrantLock;
020
021 import org.slf4j.Logger;
022 import org.slf4j.LoggerFactory;
023
024 /**
025 * An implementation of FutureCallback that does not do any actual computation,
026 * but just serves as a communication proxy between the executor and listeners.
027 *
028 * This is a base class suitable for subclassing to return the correct object
029 * type in succeed() and fail(). Ad-hoc notifiers can use the FutureNotifier
030 * subclass to avoid additional parameterization.
031 *
032 * @param <V>
033 * The result type
034 * @param <T>
035 * The subclass being defined (for return value parameterization)
036 */
037 public class FutureNotifierBase<V, T extends FutureCallback<V, T>> implements
038 FutureCallback<V, T>, FutureListener<V> {
039
040 private final static Logger log = LoggerFactory
041 .getLogger(FutureNotifierBase.class);
042
043 /** Callback listeners */
044 private final List<FutureListener<V>> listeners =
045 new CopyOnWriteArrayList<FutureListener<V>>();
046 private final Lock callbackLock = new ReentrantLock();
047
048 /** Synchronization control for FutureNotifier */
049 private final Sync sync;
050
051 /**
052 * Creates a <tt>FutureNotifier</tt> that notifies any callbacks upon
053 * completion.
054 */
055 public FutureNotifierBase() {
056 sync = new Sync();
057 }
058
059 /**
060 * Creates a <tt>FutureNotifier</tt> that notifies any callbacks upon
061 * completion. If a runner thread is specified here, cancel() will attempt
062 * to interrupt it if requested.
063 */
064 public FutureNotifierBase(final Thread runner) {
065 sync = new Sync(runner);
066 }
067
068 @Override
069 public boolean isCancelled() {
070 return sync.innerIsCancelled();
071 }
072
073 @Override
074 public boolean isDone() {
075 return sync.innerIsDone();
076 }
077
078 @Override
079 public boolean cancel(final boolean mayInterruptIfRunning) {
080 return sync.innerCancel(mayInterruptIfRunning);
081 }
082
083 /**
084 * @throws CancellationException
085 * {@inheritDoc}
086 */
087 @Override
088 public V get() throws InterruptedException, ExecutionException {
089 return sync.innerGet();
090 }
091
092 @Override
093 public V getUnchecked() {
094 try {
095 return sync.innerGet();
096 } catch (final Exception e) {
097 return null;
098 }
099 }
100
101 /**
102 * @throws CancellationException
103 * {@inheritDoc}
104 */
105 @Override
106 public V get(final long timeout, final TimeUnit unit)
107 throws InterruptedException, ExecutionException, TimeoutException {
108 return sync.innerGet(unit.toNanos(timeout));
109 }
110
111 /**
112 * Result listener for chaining callbacks together.
113 */
114 @Override
115 public void resultAvailable(final Future<V> result) {
116 try {
117 succeed(result.get());
118 } catch (final ExecutionException e) {
119 fail(e.getCause());
120 } catch (final Exception e) {
121 fail(e);
122 }
123 }
124
125 /**
126 * Protected method invoked when this task transitions to state
127 * <tt>isDone</tt> (whether normally or via cancellation). The default
128 * implementation does nothing. Subclasses may override this method to
129 * invoke completion callbacks or perform bookkeeping. Note that you can
130 * query status inside the implementation of this method to determine
131 * whether this task has been cancelled.
132 */
133 protected void done() {
134 callbackLock.lock();
135 try {
136 for (final FutureListener<V> l : listeners) {
137 try {
138 l.resultAvailable(this);
139 } catch (final Exception ex) {
140 log.warn("Unhandled exception in callback", ex);
141 }
142 }
143 } finally {
144 callbackLock.unlock();
145 }
146 }
147
148 /**
149 * Sets the result of this Future to the given value unless this future has
150 * already been set or has been cancelled. This method is invoked internally
151 * by the <tt>run</tt> method upon successful completion of the computation.
152 *
153 * @param v
154 * the value
155 */
156 protected void set(final V v) {
157 if (isDone()) {
158 throw new IllegalStateException("Future already completed");
159 }
160 sync.innerSet(v);
161 }
162
163 /**
164 * Causes this future to report an <tt>ExecutionException</tt> with the
165 * given throwable as its cause, unless this Future has already been set or
166 * has been cancelled. This method is invoked internally by the <tt>run</tt>
167 * method upon failure of the computation.
168 *
169 * @param t
170 * the cause of failure
171 */
172 protected void setException(final Throwable t) {
173 if (isDone()) {
174 throw new IllegalStateException("Future already completed");
175 }
176 sync.innerSetException(t);
177 }
178
179 @SuppressWarnings("unchecked")
180 @Override
181 public T addResultListener(final FutureListener<V> listener) {
182 callbackLock.lock();
183 try {
184 listeners.add(listener);
185 if (isDone()) {
186 try {
187 listener.resultAvailable(this);
188 } catch (final Exception ex) {
189 log.warn("Unhandled exception in callback", ex);
190 }
191 } else if (listener instanceof CancellableFutureNotifier) {
192 // MAGIC - register as parent for cancel calls
193 ((CancellableFutureNotifier<?, ?>) listener)
194 .setCancelCallback(this);
195 }
196 } finally {
197 callbackLock.unlock();
198 }
199 return (T) this;
200 }
201
202 @SuppressWarnings("unchecked")
203 @Override
204 public T fail(final Throwable error) {
205 setException(error);
206 return (T) this;
207 }
208
209 @SuppressWarnings("unchecked")
210 @Override
211 public T succeed(final V result) {
212 set(result);
213 return (T) this;
214 }
215
216 /**
217 * Synchronization control for FutureNotifier. Note that this must be a
218 * non-static inner class in order to invoke the protected <tt>done</tt>
219 * method. For clarity, all inner class support methods are same as outer,
220 * prefixed with "inner".
221 *
222 * Uses AQS sync state to represent run status
223 */
224 private final class Sync extends AbstractQueuedSynchronizer {
225 private static final long serialVersionUID = -7828117401763700385L;
226
227 /** State value representing that task ran */
228 private static final int RAN = 1;
229 /** State value representing that task was cancelled */
230 private static final int CANCELLED = 2;
231
232 /** The result to return from get() */
233 private V result = null;
234 /** The exception to throw from get() */
235 private Throwable exception = null;
236
237 /**
238 * The thread running task. When nulled after set/cancel, this indicates
239 * that the results are accessible. Must be volatile, to ensure
240 * visibility upon completion.
241 */
242 private volatile Thread runner = null;
243
244 Sync() {
245 }
246
247 Sync(final Thread runner_) {
248 runner = runner_;
249 }
250
251 private boolean ranOrCancelled(final int state) {
252 return (state & (RAN | CANCELLED)) != 0;
253 }
254
255 /**
256 * Implements AQS base acquire to succeed if ran or cancelled
257 */
258 @Override
259 protected int tryAcquireShared(final int ignore) {
260 return innerIsDone() ? 1 : -1;
261 }
262
263 /**
264 * Implements AQS base release to always signal after setting final done
265 * status by nulling runner thread.
266 */
267 @Override
268 protected boolean tryReleaseShared(final int ignore) {
269 runner = null;
270 return true;
271 }
272
273 boolean innerIsCancelled() {
274 return getState() == CANCELLED;
275 }
276
277 boolean innerIsDone() {
278 return ranOrCancelled(getState()) && runner == null;
279 }
280
281 V innerGet() throws InterruptedException, ExecutionException {
282 acquireSharedInterruptibly(0);
283 if (getState() == CANCELLED) {
284 throw new CancellationException();
285 }
286 if (exception != null) {
287 throw new ExecutionException(exception);
288 }
289 return result;
290 }
291
292 V innerGet(final long nanosTimeout) throws InterruptedException,
293 ExecutionException, TimeoutException {
294 if (!tryAcquireSharedNanos(0, nanosTimeout)) {
295 throw new TimeoutException();
296 }
297 if (getState() == CANCELLED) {
298 throw new CancellationException();
299 }
300 if (exception != null) {
301 throw new ExecutionException(exception);
302 }
303 return result;
304 }
305
306 void innerSet(final V v) {
307 for (;;) {
308 final int s = getState();
309 if (s == RAN) {
310 return;
311 }
312 if (s == CANCELLED) {
313 // aggressively release to set runner to null,
314 // in case we are racing with a cancel request
315 // that will try to interrupt runner
316 releaseShared(0);
317 return;
318 }
319 if (compareAndSetState(s, RAN)) {
320 result = v;
321 releaseShared(0);
322 done();
323 return;
324 }
325 }
326 }
327
328 void innerSetException(final Throwable t) {
329 for (;;) {
330 final int s = getState();
331 if (s == RAN) {
332 return;
333 }
334 if (s == CANCELLED) {
335 // aggressively release to set runner to null,
336 // in case we are racing with a cancel request
337 // that will try to interrupt runner
338 releaseShared(0);
339 return;
340 }
341 if (compareAndSetState(s, RAN)) {
342 exception = t;
343 result = null;
344 releaseShared(0);
345 done();
346 return;
347 }
348 }
349 }
350
351 boolean innerCancel(final boolean mayInterruptIfRunning) {
352 for (;;) {
353 final int s = getState();
354 if (ranOrCancelled(s)) {
355 return false;
356 }
357 if (compareAndSetState(s, CANCELLED)) {
358 break;
359 }
360 }
361 if (mayInterruptIfRunning) {
362 final Thread r = runner;
363 if (r != null) {
364 r.interrupt();
365 }
366 }
367 releaseShared(0);
368 done();
369 return true;
370 }
371
372 }
373
374 }