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.CopyOnWriteArrayList;
012    import java.util.concurrent.Future;
013    import java.util.concurrent.Semaphore;
014    import java.util.concurrent.TimeUnit;
015    import java.util.concurrent.TimeoutException;
016    
017    import org.slf4j.Logger;
018    import org.slf4j.LoggerFactory;
019    
020    /**
021     * Groups a number of {@link FutureCallbackTask} results together, and executes
022     * a single callback only after all FutureCallbacks have completed (succeeded or
023     * failed.)
024     * 
025     * This class does not support direct access to succeed() or fail(), and will
026     * throw and UnsupportedOperationException for them.
027     * 
028     * @author jeremy
029     * @see FutureCallbackTask
030     * @param <E>
031     *            The deferred result type
032     */
033    public class FutureCallbackList<E> implements
034                    FutureCallback<List<E>, FutureCallbackList<E>> {
035    
036            private final static Logger log = LoggerFactory
037                            .getLogger(FutureCallbackList.class);
038    
039            private final FutureListener<E> listener = new FutureListener<E>() {
040    
041                    @Override
042                    public void resultAvailable(final Future<E> result_) {
043                            onResponse(result_);
044                    }
045    
046            };
047    
048            private final List<FutureListener<List<E>>> listeners =
049                            new CopyOnWriteArrayList<FutureListener<List<E>>>();
050            private final List<E> responses = new CopyOnWriteArrayList<E>();
051            private final List<? extends FutureCallback<E, ?>> results;
052            private final Semaphore resultMonitor = new Semaphore(1);
053    
054            private volatile int completed = 0;
055            private volatile boolean fired = false;
056            private volatile boolean cancelled = false;
057    
058            /**
059             * Create a FutureCallbackList from the specified FutureCallback objects.
060             * 
061             * @param deferreds_
062             *            The list of deferred results to monitor
063             */
064            public FutureCallbackList(
065                            final List<? extends FutureCallback<E, ?>> results_) {
066                    results = results_;
067                    if (results_ == null || results_.size() == 0) {
068                            fired = true;
069                    } else {
070                            // Lock the semaphore until a result is available
071                            resultMonitor.acquireUninterruptibly();
072                            for (final FutureCallback<E, ?> r : results_) {
073                                    r.addResultListener(listener);
074                            }
075                    }
076            }
077    
078            /**
079             * Add a callback listener that will be called when all FutureCallbacks have
080             * completed. The callback will be passed a list of results collected from
081             * the FutureCallbacks. A null member may indicate that a FutureCallback
082             * failed, but does not guarantee it (since a FutureCallback callback could
083             * return a null value on purpose.)
084             * 
085             * @param callback
086             *            The object to notify of the deferred results
087             */
088            @Override
089            public FutureCallbackList<E> addResultListener(
090                            final FutureListener<List<E>> callback) {
091                    listeners.add(callback);
092                    if (fired) {
093                            try {
094                                    callback.resultAvailable(this);
095                            } catch (final Exception e) {
096                                    log.warn("Unhandled exception in callback", e);
097                            }
098                    }
099                    return this;
100            }
101    
102            /**
103             * Register a response from a FutureCallback.
104             */
105            private void onResponse(final Future<E> result) {
106                    if (fired) {
107                            return;
108                    }
109                    completed++;
110                    if (completed == results.size()) {
111                            for (final FutureCallback<E, ?> r : results) {
112                                    try {
113                                            responses.add(r.get());
114                                    } catch (final Exception e) {
115                                            responses.add(null);
116                                    }
117                            }
118                            fired = true;
119                            resultMonitor.release();
120                            for (final FutureListener<List<E>> l : listeners) {
121                                    try {
122                                            l.resultAvailable(this);
123                                    } catch (final Exception e) {
124                                            log.warn("Unhandled exception in callback", e);
125                                    }
126                            }
127                    }
128            }
129    
130            @Override
131            public boolean cancel(final boolean interrupt) {
132                    if (fired) {
133                            return false;
134                    }
135                    cancelled = true;
136                    for (final FutureCallback<E, ?> r : results) {
137                            r.cancel(interrupt);
138                    }
139                    return true;
140            }
141    
142            @Override
143            public List<E> get() {
144                    if (!fired) {
145                            resultMonitor.acquireUninterruptibly();
146                    }
147                    return responses;
148            }
149    
150            @Override
151            public List<E> getUnchecked() {
152                    return get();
153            }
154    
155            @Override
156            public List<E> get(final long timeout, final TimeUnit unit)
157                            throws InterruptedException, TimeoutException {
158                    if (!fired && !resultMonitor.tryAcquire(timeout, unit)) {
159                            throw new TimeoutException("Timed out waiting for result");
160                    }
161                    return responses;
162            }
163    
164            @Override
165            public boolean isCancelled() {
166                    return cancelled;
167            }
168    
169            @Override
170            public boolean isDone() {
171                    return fired;
172            }
173    
174            @Override
175            public FutureCallbackList<E> succeed(final List<E> result) {
176                    throw new UnsupportedOperationException();
177            }
178    
179            @Override
180            public FutureCallbackList<E> fail(final Throwable error) {
181                    throw new UnsupportedOperationException();
182            }
183    
184    }