[java-metadata-aggregator COMMIT] in /trunk/aggregator-pipeline/src: main/java/net/shibboleth/metadata/pipeline/Futur...
noreply at shibboleth.net
noreply at shibboleth.net
Wed Sep 18 10:12:56 EDT 2013
Author: iay
Date: Wed Sep 18 10:12:55 2013
New Revision: 260
URL: http://svn.shibboleth.net/view/java-metadata-aggregator?rev=260&view=rev
Log:
MDA-44: exceptions should be propagated up from subordinate pipelines
Added:
trunk/aggregator-pipeline/src/main/java/net/shibboleth/metadata/pipeline/FutureSupport.java (with props)
trunk/aggregator-pipeline/src/test/java/net/shibboleth/metadata/pipeline/TerminatingStage.java (with props)
Modified:
trunk/aggregator-pipeline/src/main/java/net/shibboleth/metadata/pipeline/PipelineCallable.java
trunk/aggregator-pipeline/src/main/java/net/shibboleth/metadata/pipeline/PipelineDemultiplexerStage.java
trunk/aggregator-pipeline/src/main/java/net/shibboleth/metadata/pipeline/PipelineMergeStage.java
trunk/aggregator-pipeline/src/main/java/net/shibboleth/metadata/pipeline/SplitMergeStage.java
trunk/aggregator-pipeline/src/test/java/net/shibboleth/metadata/pipeline/PipelineDemultiplexerStageTest.java
trunk/aggregator-pipeline/src/test/java/net/shibboleth/metadata/pipeline/PipelineMergeStageTest.java
trunk/aggregator-pipeline/src/test/java/net/shibboleth/metadata/pipeline/SplitMergeStageTest.java
Modified: trunk/aggregator-pipeline/src/main/java/net/shibboleth/metadata/pipeline/PipelineCallable.java
URL: http://svn.shibboleth.net/view/java-metadata-aggregator/trunk/aggregator-pipeline/src/main/java/net/shibboleth/metadata/pipeline/PipelineCallable.java?rev=260&r1=259&r2=260&view=diff
==============================================================================
--- trunk/aggregator-pipeline/src/main/java/net/shibboleth/metadata/pipeline/PipelineCallable.java (original)
+++ trunk/aggregator-pipeline/src/main/java/net/shibboleth/metadata/pipeline/PipelineCallable.java Wed Sep 18 10:12:55 2013
@@ -55,16 +55,15 @@
itemCollection = Constraint.isNotNull(items, "Item collection can not be null");
}
- /** {@inheritDoc} */
- @Nonnull @NonnullElements public Collection<? extends Item> call() {
- try {
- log.debug("Executing pipeline {} on an item collection containing {} items", pipeline.getId(),
- itemCollection.size());
- pipeline.execute(itemCollection);
- return itemCollection;
- } catch (PipelineProcessingException e) {
- log.error("Execution of pipeline {} failed.", pipeline.getId(), e);
- return null;
- }
+ /**
+ * {@inheritDoc}
+ *
+ * @throws PipelineProcessingException
+ */
+ @Nonnull @NonnullElements public Collection<? extends Item> call() throws PipelineProcessingException {
+ log.debug("Executing pipeline {} on an item collection containing {} items", pipeline.getId(),
+ itemCollection.size());
+ pipeline.execute(itemCollection);
+ return itemCollection;
}
}
Modified: trunk/aggregator-pipeline/src/main/java/net/shibboleth/metadata/pipeline/PipelineDemultiplexerStage.java
URL: http://svn.shibboleth.net/view/java-metadata-aggregator/trunk/aggregator-pipeline/src/main/java/net/shibboleth/metadata/pipeline/PipelineDemultiplexerStage.java?rev=260&r1=259&r2=260&view=diff
==============================================================================
--- trunk/aggregator-pipeline/src/main/java/net/shibboleth/metadata/pipeline/PipelineDemultiplexerStage.java (original)
+++ trunk/aggregator-pipeline/src/main/java/net/shibboleth/metadata/pipeline/PipelineDemultiplexerStage.java Wed Sep 18 10:12:55 2013
@@ -21,7 +21,6 @@
import java.util.Collection;
import java.util.Collections;
import java.util.List;
-import java.util.concurrent.ExecutionException;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.Future;
@@ -38,9 +37,6 @@
import net.shibboleth.utilities.java.support.component.ComponentSupport;
import net.shibboleth.utilities.java.support.logic.Constraint;
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
-
import com.google.common.base.Predicate;
import com.google.common.base.Supplier;
import com.google.common.collect.ImmutableList.Builder;
@@ -67,9 +63,6 @@
*/
@ThreadSafe
public class PipelineDemultiplexerStage<ItemType extends Item<?>> extends BaseStage<ItemType> {
-
- /** Class logger. */
- private final Logger log = LoggerFactory.getLogger(PipelineDemultiplexerStage.class);
/** Service used to execute the selected and/or non-selected item pipelines. */
private ExecutorService executorService = Executors.newSingleThreadExecutor();
@@ -207,13 +200,7 @@
if (isWaitingForPipelines()) {
for (Future pipelineFuture : pipelineFutures) {
- try {
- pipelineFuture.get();
- } catch (ExecutionException e) {
- log.error("Pipeline threw an unexpected exception", e);
- } catch (InterruptedException e) {
- log.error("Execution service was interrupted", e);
- }
+ FutureSupport.futureItems(pipelineFuture);
}
}
}
[... 263 lines stripped ...]
More information about the commits
mailing list