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 }