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    }