| 615 | } |
| 616 | |
| 617 | final class CollectPromise<T> extends Promise<List<T>> { |
| 618 | |
| 619 | private final Object[] results; |
| 620 | private final AtomicInteger count; |
| 621 | private final List<? extends Future<T>> list; |
| 622 | |
| 623 | public CollectPromise(final List<? extends Future<T>> list) { |
| 624 | final int size = list.size(); |
| 625 | this.results = new Object[size]; |
| 626 | this.count = new AtomicInteger(size); |
| 627 | this.list = list; |
| 628 | } |
| 629 | |
| 630 | @SuppressWarnings("unchecked") |
| 631 | public final void collect(final T value, final int i) { |
| 632 | results[i] = value; |
| 633 | if (count.decrementAndGet() == 0) |
| 634 | setValue((List<T>) Arrays.asList(results)); |
| 635 | } |
| 636 | |
| 637 | @Override |
| 638 | protected final InterruptHandler getInterruptHandler() { |
| 639 | return InterruptHandler.apply(list); |
| 640 | } |
| 641 | } |
| 642 | |
| 643 | final class JoinPromise<T> extends Promise<Void> implements Responder<T> { |
| 644 | private final AtomicInteger count; |
nothing calls this directly
no outgoing calls
no test coverage detected
searching dependent graphs…