BATCH-1011: Refactor BatchStatus again so that FAILED has the same meaning as before (1.1.x etc.) and ABANDONED is introduced to signify a step that failed but does not need to be replayed on restart (formerly known as INCOMPLETE).

This commit is contained in:
dsyer
2009-02-15 10:59:44 +00:00
parent fa56955f8d
commit 7533f90480
70 changed files with 696 additions and 427 deletions

View File

@@ -37,7 +37,7 @@ public enum BatchStatus {
* steps that have finished processing, but were not successful, and where
* they should be skipped on a restart (so FAILED is the wrong status).
*/
COMPLETED, STARTING, STARTED, FAILED, INCOMPLETE, STOPPING, UNKNOWN;
COMPLETED, STARTING, STARTED, STOPPING, STOPPED, FAILED, ABANDONED, UNKNOWN;
public static BatchStatus max(BatchStatus status1, BatchStatus status2) {
if (status1.isLessThan(status2)) {
@@ -66,7 +66,7 @@ public enum BatchStatus {
* @return true if the status is FAILED or greater
*/
public boolean isUnsuccessful() {
return this == INCOMPLETE || this.isGreaterThan(INCOMPLETE);
return this == FAILED || this.isGreaterThan(FAILED);
}
/**

View File

@@ -66,7 +66,7 @@ public class ExitStatus implements Serializable, Comparable<ExitStatus> {
* Convenient constant value representing finished processing with
* interrupted status.
*/
public static final ExitStatus INTERRUPTED = new ExitStatus("INTERRUPTED");
public static final ExitStatus STOPPED = new ExitStatus("STOPPED");
private final String exitCode;
@@ -159,7 +159,7 @@ public class ExitStatus implements Serializable, Comparable<ExitStatus> {
if (status.exitCode.startsWith(NOOP.exitCode)) {
return 3;
}
if (status.exitCode.startsWith(INTERRUPTED.exitCode)) {
if (status.exitCode.startsWith(STOPPED.exitCode)) {
return 4;
}
if (status.exitCode.startsWith(FAILED.exitCode)) {

View File

@@ -41,20 +41,27 @@ import org.w3c.dom.NodeList;
public class FlowParser extends AbstractSingleBeanDefinitionParser {
private static final String NEXT = "next";
private static final String END = "end";
private static final String FAIL = "fail";
private static final String STOP = "stop";
// For generating unique state names for end transitions
private static int endCounter = 0;
private final String flowName;
private final String jobRepositoryRef;
/**
* Construct a {@link FlowParser} with the specified name and using the provided job repository ref.
* Construct a {@link FlowParser} with the specified name and using the
* provided job repository ref.
*
* @param flowName the name of the flow
* @param jobRepositoryRef the reference to the jobRepository from the enclosing tag
* @param jobRepositoryRef the reference to the jobRepository from the
* enclosing tag
*/
public FlowParser(String flowName, String jobRepositoryRef) {
this.flowName = flowName;
@@ -122,11 +129,26 @@ public class FlowParser extends AbstractSingleBeanDefinitionParser {
* @param stateDef The bean definition for the current state
* @param element the &lt;step/gt; element to parse
* @return a collection of
* {@link org.springframework.batch.core.job.flow.support.StateTransition}
* references
* {@link org.springframework.batch.core.job.flow.support.StateTransition}
* references
*/
protected static Collection<BeanDefinition> getNextElements(ParserContext parserContext, BeanDefinition stateDef,
Element element) {
return getNextElements(parserContext, null, stateDef, element);
}
/**
* @param parserContext the parser context for the bean factory
* @param stepId the id of the current state if it is a step state, null
* otherwise
* @param stateDef The bean definition for the current state
* @param element the &lt;step/gt; element to parse
* @return a collection of
* {@link org.springframework.batch.core.job.flow.support.StateTransition}
* references
*/
protected static Collection<BeanDefinition> getNextElements(ParserContext parserContext, String stepId,
BeanDefinition stateDef, Element element) {
Collection<BeanDefinition> list = new ArrayList<BeanDefinition>();
@@ -144,16 +166,16 @@ public class FlowParser extends AbstractSingleBeanDefinitionParser {
transitionName);
for (Element transitionElement : transitionElements) {
verifyUniquePattern(transitionElement, patterns, element, parserContext);
list.addAll(parseTransitionElement(transitionElement, stateDef, parserContext));
list.addAll(parseTransitionElement(transitionElement, stepId, stateDef, parserContext));
transitionElementExists = true;
}
}
if (!transitionElementExists) {
list.addAll(createTransition(BatchStatus.INCOMPLETE, ExitStatus.FAILED.getExitCode(), null, null,
stateDef, parserContext));
list.addAll(createTransition(BatchStatus.FAILED, ExitStatus.FAILED.getExitCode(), null, null, stateDef,
parserContext, false));
if (!hasNextAttribute) {
list.addAll(createTransition(BatchStatus.COMPLETED, null, null, null, stateDef, parserContext));
list.addAll(createTransition(BatchStatus.COMPLETED, null, null, null, stateDef, parserContext, false));
}
}
else if (hasNextAttribute) {
@@ -185,43 +207,50 @@ public class FlowParser extends AbstractSingleBeanDefinitionParser {
* @param stateDef The bean definition for the current state
* @param parserContext the parser context for the bean factory
* @param a collection of
* {@link org.springframework.batch.core.job.flow.support.StateTransition}
* references
* {@link org.springframework.batch.core.job.flow.support.StateTransition}
* references
*/
private static Collection<BeanDefinition> parseTransitionElement(Element transitionElement,
private static Collection<BeanDefinition> parseTransitionElement(Element transitionElement, String stateId,
BeanDefinition stateDef, ParserContext parserContext) {
BatchStatus batchStatus = getBatchStatusFromEndTransitionName(transitionElement.getNodeName());
String onAttribute = transitionElement.getAttribute("on");
String nextAttribute = transitionElement.getAttribute("to");
nextAttribute = StringUtils.hasText(nextAttribute) ? nextAttribute : transitionElement.getAttribute("restart");
String restartAttribute = transitionElement.getAttribute("restart");
nextAttribute = StringUtils.hasText(nextAttribute) ? nextAttribute : restartAttribute;
boolean abandon = false;
if (stateId != null && StringUtils.hasText(restartAttribute) && !restartAttribute.equals(stateId)) {
abandon = true;
}
String statusAttribute = transitionElement.getAttribute("status");
return createTransition(batchStatus, onAttribute, nextAttribute, statusAttribute, stateDef, parserContext);
return createTransition(batchStatus, onAttribute, nextAttribute, statusAttribute, stateDef, parserContext,
abandon);
}
/**
* @param batchStatus The batch status that this transition will set. Use
* BatchStatus.UNKNOWN if not applicable.
* BatchStatus.UNKNOWN if not applicable.
* @param on The pattern that this transition should match. Use null for
* "no restriction" (same as "*").
* "no restriction" (same as "*").
* @param next The state to which this transition should go. Use null if not
* applicable.
* applicable.
* @param exitCode The exit code that this transition will set. Use null to
* default to batchStatus.
* default to batchStatus.
* @param stateDef The bean definition for the current state
* @param parserContext the parser context for the bean factory
* @param a collection of
* {@link org.springframework.batch.core.job.flow.support.StateTransition}
* references
* {@link org.springframework.batch.core.job.flow.support.StateTransition}
* references
*/
private static Collection<BeanDefinition> createTransition(BatchStatus batchStatus, String on, String next,
String exitCode, BeanDefinition stateDef, ParserContext parserContext) {
String exitCode, BeanDefinition stateDef, ParserContext parserContext, boolean abandon) {
BeanDefinition endState = null;
if (batchStatus == BatchStatus.FAILED || batchStatus == BatchStatus.COMPLETED
|| batchStatus == BatchStatus.INCOMPLETE) {
// TODO: revise this for clarity
if (batchStatus == BatchStatus.STOPPED || batchStatus == BatchStatus.COMPLETED
|| batchStatus == BatchStatus.FAILED) {
BeanDefinitionBuilder endBuilder = BeanDefinitionBuilder
.genericBeanDefinition("org.springframework.batch.core.job.flow.support.state.EndState");
@@ -231,9 +260,12 @@ public class FlowParser extends AbstractSingleBeanDefinitionParser {
endBuilder.addConstructorArgValue(exitCodeExists ? new ExitStatus(exitCode)
: convertToExitStatus(batchStatus));
String endName = "end" + (endCounter++);
String endName = (batchStatus == BatchStatus.STOPPED ? STOP : batchStatus == BatchStatus.FAILED ? FAIL : END)
+ (endCounter++);
endBuilder.addConstructorArgValue(endName);
endBuilder.addConstructorArgValue(abandon);
String nextOnEnd = exitCodeExists ? null : next;
endState = getStateTransitionReference(parserContext, endBuilder.getBeanDefinition(), null, nextOnEnd);
next = endName;
@@ -258,7 +290,7 @@ public class FlowParser extends AbstractSingleBeanDefinitionParser {
*/
private static BatchStatus getBatchStatusFromEndTransitionName(String elementName) {
if (STOP.equals(elementName)) {
return BatchStatus.INCOMPLETE;
return BatchStatus.STOPPED;
}
else if (END.equals(elementName)) {
return BatchStatus.COMPLETED;
@@ -276,7 +308,7 @@ public class FlowParser extends AbstractSingleBeanDefinitionParser {
* @return the ExitStatus corresponding to the BatchStatus
*/
private static ExitStatus convertToExitStatus(BatchStatus batchStatus) {
if (batchStatus == BatchStatus.INCOMPLETE) {
if (batchStatus == BatchStatus.FAILED) {
return ExitStatus.FAILED;
}
else {
@@ -289,13 +321,14 @@ public class FlowParser extends AbstractSingleBeanDefinitionParser {
* @param stateDefinition a reference to the state implementation
* @param on the pattern value
* @param next the next step id
* @return a bean definition for a {@link org.springframework.batch.core.job.flow.support.StateTransition}
* @return a bean definition for a
* {@link org.springframework.batch.core.job.flow.support.StateTransition}
*/
public static BeanDefinition getStateTransitionReference(ParserContext parserContext,
BeanDefinition stateDefinition, String on, String next) {
BeanDefinitionBuilder nextBuilder =
BeanDefinitionBuilder.genericBeanDefinition("org.springframework.batch.core.job.flow.support.StateTransition");
BeanDefinitionBuilder nextBuilder = BeanDefinitionBuilder
.genericBeanDefinition("org.springframework.batch.core.job.flow.support.StateTransition");
nextBuilder.addConstructorArgValue(stateDefinition);
if (StringUtils.hasText(on)) {

View File

@@ -101,7 +101,7 @@ public class InlineStepParser extends AbstractStepParser {
else {
parserContext.getReaderContext().error("Incomplete configuration detected while creating step with name " + stepRef, element);
}
return FlowParser.getNextElements(parserContext, stateBuilder.getBeanDefinition(), element);
return FlowParser.getNextElements(parserContext, stepId, stateBuilder.getBeanDefinition(), element);
}

View File

@@ -84,7 +84,6 @@ public class SplitParser {
stateBuilder.addConstructorArgValue(managedList);
stateBuilder.addConstructorArgValue(idAttribute);
// TODO: allow TaskExecutor etc. to be set
return FlowParser.getNextElements(parserContext, stateBuilder.getBeanDefinition(), element);
}

View File

@@ -37,6 +37,7 @@ import org.springframework.batch.core.repository.JobRepository;
import org.springframework.batch.core.repository.JobRestartException;
import org.springframework.batch.core.step.StepLocator;
import org.springframework.batch.item.ExecutionContext;
import org.springframework.batch.repeat.RepeatException;
import org.springframework.beans.factory.BeanNameAware;
import org.springframework.beans.factory.InitializingBean;
import org.springframework.util.Assert;
@@ -235,28 +236,32 @@ public abstract class AbstractJob implements Job, StepLocator, BeanNameAware, In
listener.beforeJob(execution);
doExecute(execution);
try {
doExecute(execution);
} catch (RepeatException e) {
throw e.getCause();
}
}
else {
// The job was already stopped before we even got this far. Deal
// with it in the same way as any other interruption.
execution.setExitStatus(ExitStatus.FAILED);
execution.setStatus(BatchStatus.INCOMPLETE);
execution.setStatus(BatchStatus.FAILED);
execution.setExitStatus(ExitStatus.COMPLETED);
}
}
catch (JobInterruptedException e) {
logger.error("Encountered interruption executing job", e);
execution.setExitStatus(ExitStatus.FAILED);
execution.setStatus(BatchStatus.INCOMPLETE);
execution.setExitStatus(ExitStatus.STOPPED);
execution.setStatus(BatchStatus.STOPPED);
execution.addFailureException(e);
}
catch (Throwable t) {
logger.error("Encountered error executing job", t);
logger.error("Encountered fatal error executing job", t);
execution.setExitStatus(ExitStatus.FAILED);
execution.setStatus(BatchStatus.INCOMPLETE);
execution.setStatus(BatchStatus.FAILED);
execution.addFailureException(t);
}
finally {
@@ -274,8 +279,9 @@ public abstract class AbstractJob implements Job, StepLocator, BeanNameAware, In
catch (Exception e) {
logger.error("Exception encountered in afterStep callback", e);
}
jobRepository.update(execution);
}
@@ -293,7 +299,7 @@ public abstract class AbstractJob implements Job, StepLocator, BeanNameAware, In
* @return the {@link StepExecution} corresponding to this step
*
* @throws JobInterruptedException if the {@link JobExecution} has been
* interrupted, and in particular if {@link BatchStatus#INCOMPLETE} or
* interrupted, and in particular if {@link BatchStatus#ABANDONED} or
* {@link BatchStatus#STOPPING} is detected
* @throws StartLimitExceededException if the start limit has been exceeded
* for this step
@@ -302,7 +308,7 @@ public abstract class AbstractJob implements Job, StepLocator, BeanNameAware, In
*/
protected final StepExecution handleStep(Step step, JobExecution execution) throws JobInterruptedException,
JobRestartException, StartLimitExceededException {
if (execution.getStatus() == BatchStatus.STOPPING || execution.getStatus() == BatchStatus.INCOMPLETE) {
if (execution.isStopping()) {
throw new JobInterruptedException("JobExecution interrupted.");
}
@@ -327,6 +333,7 @@ public abstract class AbstractJob implements Job, StepLocator, BeanNameAware, In
jobRepository.add(currentStepExecution);
logger.info("Executing step: "+step);
step.execute(currentStepExecution);
jobRepository.updateExecutionContext(execution);
@@ -382,9 +389,10 @@ public abstract class AbstractJob implements Job, StepLocator, BeanNameAware, In
}
if ((stepStatus == BatchStatus.COMPLETED && step.isAllowStartIfComplete() == false)
|| stepStatus == BatchStatus.FAILED) {
|| stepStatus == BatchStatus.ABANDONED) {
// step is complete, false should be returned, indicating that the
// step should not be started
logger.info("Step already complete or not restartable, so no action to execute: "+lastStepExecution);
return false;
}

View File

@@ -101,6 +101,7 @@ public class SimpleJob extends AbstractJob {
// Update the job status to be the same as the last step
//
if(stepExecution != null) {
logger.debug("Upgrading JobExecution status: "+stepExecution);
execution.upgradeStatus(stepExecution.getStatus());
execution.setExitStatus(stepExecution.getExitStatus());
}

View File

@@ -15,59 +15,60 @@
*/
package org.springframework.batch.core.job.flow;
import org.springframework.batch.core.BatchStatus;
import org.springframework.batch.core.ExitStatus;
/**
* This class is used as a holder for a BatchStatus/ExitStatus pair.
*
* @author Dan Garrette
* @author Dave Syer
* @since 2.0
*/
public class FlowExecutionStatus implements Comparable<FlowExecutionStatus> {
private final BatchStatus batchStatus;
private final ExitStatus exitStatus;
/**
* Special well-known status value.
*/
public static final FlowExecutionStatus COMPLETED = new FlowExecutionStatus(Status.COMPLETED.toString());
/**
* Special well-known status value.
*/
public static final FlowExecutionStatus COMPLETED = new FlowExecutionStatus(BatchStatus.COMPLETED,
ExitStatus.COMPLETED);
public static final FlowExecutionStatus STOPPED = new FlowExecutionStatus(Status.STOPPED.toString());
/**
* Special well-known status value.
*/
public static final FlowExecutionStatus INCOMPLETE = new FlowExecutionStatus(BatchStatus.INCOMPLETE, ExitStatus.FAILED);
public static final FlowExecutionStatus FAILED = new FlowExecutionStatus(Status.FAILED.toString());
/**
* Special well-known status value.
*/
public static final FlowExecutionStatus FAILED = new FlowExecutionStatus(BatchStatus.FAILED, ExitStatus.FAILED);
public static final FlowExecutionStatus UNKNOWN = new FlowExecutionStatus(Status.UNKNOWN.toString());
private final String status;
private enum Status {
COMPLETED, STOPPED, FAILED, UNKNOWN;
static Status match(String value) {
for (int i = 0; i < values().length; i++) {
Status status = values()[i];
if (value.startsWith(status.toString())) {
return status;
}
}
// Default match should be the lowest priority
return COMPLETED;
}
};
/**
* Special well-known status value.
* @param status
*/
public static final FlowExecutionStatus UNKNOWN = new FlowExecutionStatus(BatchStatus.UNKNOWN, ExitStatus.UNKNOWN);
/**
* @param exitStatus
*/
public FlowExecutionStatus(String exitStatus) {
this.batchStatus = null;
this.exitStatus = new ExitStatus(exitStatus);
}
/**
* Convenience constructor that accepts a {@link BatchStatus} and
* {@link ExitStatus}.
*
* @param batchStatus
* @param exitStatus
*/
public FlowExecutionStatus(BatchStatus batchStatus, ExitStatus exitStatus) {
this.batchStatus = batchStatus;
this.exitStatus = exitStatus;
public FlowExecutionStatus(String status) {
this.status = status;
}
/**
@@ -80,13 +81,13 @@ public class FlowExecutionStatus implements Comparable<FlowExecutionStatus> {
* @return negative, zero or positive as per the contract
*/
public int compareTo(FlowExecutionStatus other) {
if (batchStatus != null && other.getBatchStatus() != null) {
int order = batchStatus.compareTo(other.getBatchStatus());
if (order != 0) {
return order;
}
Status one = Status.match(this.status);
Status two = Status.match(other.status);
int comparison = one.compareTo(two);
if (comparison == 0) {
return this.status.compareTo(other.status);
}
return exitStatus.compareTo(other.getExitStatus());
return comparison;
}
/**
@@ -94,26 +95,23 @@ public class FlowExecutionStatus implements Comparable<FlowExecutionStatus> {
*
* @see java.lang.Object#equals(java.lang.Object)
*/
public boolean equals(Object other) {
if (other == this) {
public boolean equals(Object object) {
if (object == this) {
return true;
}
if (!(other instanceof FlowExecutionStatus)) {
if (!(object instanceof FlowExecutionStatus)) {
return false;
}
FlowExecutionStatus flowExecutionStatus = (FlowExecutionStatus) other;
return batchStatus.equals(flowExecutionStatus.getBatchStatus());
FlowExecutionStatus other = (FlowExecutionStatus) object;
return status.equals(other.status);
}
public String toString() {
return "FlowExecutionStatus: status=[" + batchStatus + "] exitcode=[" + exitStatus.getExitCode() + "]";
return "FlowExecutionStatus: " + status;
}
public BatchStatus getBatchStatus() {
return batchStatus;
public String getStatus() {
return status;
}
public ExitStatus getExitStatus() {
return exitStatus;
}
}

View File

@@ -32,13 +32,15 @@ import org.springframework.batch.core.repository.JobRestartException;
public interface FlowExecutor {
/**
* @param step a {@link Step} to execute
* @param step
* a {@link Step} to execute
* @return the exit status that drives the surrounding {@link Flow}
* @throws StartLimitExceededException
* @throws JobRestartException
* @throws JobInterruptedException
*/
String executeStep(Step step) throws JobInterruptedException, JobRestartException, StartLimitExceededException;
String executeStep(Step step) throws JobInterruptedException,
JobRestartException, StartLimitExceededException;
/**
* @return the current {@link JobExecution}
@@ -49,11 +51,19 @@ public interface FlowExecutor {
* @return the latest {@link StepExecution} or null if there is none
*/
StepExecution getStepExecution();
/**
* Chance to clean up resources at the end of a flow (whether it completed successfully or not).
* @param result the final {@link FlowExecution}
* Chance to clean up resources at the end of a flow (whether it completed
* successfully or not).
*
* @param result
* the final {@link FlowExecution}
*/
void close(FlowExecution result);
/**
* Handle any status changes that might be needed at the start of a state.
*/
void updateStepExecutionStatus();
}

View File

@@ -51,7 +51,9 @@ public class FlowJob extends AbstractJob {
/**
* Public setter for the flow.
* @param flow the flow to set
*
* @param flow
* the flow to set
*/
public void setFlow(Flow flow) {
this.flow = flow;
@@ -75,21 +77,16 @@ public class FlowJob extends AbstractJob {
* @see AbstractJob#doExecute(JobExecution)
*/
@Override
protected void doExecute(final JobExecution execution) throws JobExecutionException {
protected void doExecute(final JobExecution execution)
throws JobExecutionException {
try {
FlowExecution flowExecution = flow.start(new JobFlowExecutor(execution));
synchronized (execution) {
FlowExecutionStatus status = flowExecution.getStatus();
execution.upgradeStatus(status.getBatchStatus());
execution.setExitStatus(status.getExitStatus());
}
}
catch (FlowExecutionException e) {
flow.start(new JobFlowExecutor(execution));
} catch (FlowExecutionException e) {
if (e.getCause() instanceof JobExecutionException) {
throw (JobExecutionException) e.getCause();
}
throw new JobExecutionException("Flow execution ended unexpectedly", e);
throw new JobExecutionException(
"Flow execution ended unexpectedly", e);
}
}
@@ -111,17 +108,22 @@ public class FlowJob extends AbstractJob {
stepExecutionHolder.set(null);
}
public String executeStep(Step step) throws JobInterruptedException, JobRestartException,
StartLimitExceededException {
StepExecution lastStepExecution = stepExecutionHolder.get();
if (lastStepExecution != null && lastStepExecution.getStatus() == BatchStatus.INCOMPLETE) {
lastStepExecution.setStatus(BatchStatus.FAILED);
updateStepExecution(lastStepExecution);
}
public String executeStep(Step step) throws JobInterruptedException,
JobRestartException, StartLimitExceededException {
StepExecution stepExecution = handleStep(step, execution);
stepExecutionHolder.set(stepExecution);
return stepExecution == null ? ExitStatus.COMPLETED.getExitCode() : stepExecution.getExitStatus()
.getExitCode();
return stepExecution == null ? ExitStatus.COMPLETED.getExitCode()
: stepExecution.getExitStatus().getExitCode();
}
public void updateStepExecutionStatus() {
StepExecution lastStepExecution = stepExecutionHolder.get();
if (lastStepExecution != null
&& lastStepExecution.getStatus().isGreaterThan(
BatchStatus.STOPPING)) {
lastStepExecution.upgradeStatus(BatchStatus.ABANDONED);
updateStepExecution(lastStepExecution);
}
}
public JobExecution getJobExecution() {

View File

@@ -121,7 +121,7 @@ public class SimpleFlow implements Flow, InitializingBean {
State state = stateMap.get(stateName);
// Terminate if there are no more states
while (state != null) {
while (state != null && status!=FlowExecutionStatus.STOPPED) {
stateName = state.getName();
@@ -158,7 +158,7 @@ public class SimpleFlow implements Flow, InitializingBean {
}
String next = null;
String exitCode = status.getExitStatus().getExitCode();
String exitCode = status.getStatus();
for (StateTransition stateTransition : set) {
if (stateTransition.matches(exitCode)) {
if (stateTransition.isEnd()) {

View File

@@ -33,17 +33,10 @@ import org.springframework.batch.core.job.flow.State;
public class EndState extends AbstractState {
private final BatchStatus status;
private final ExitStatus exitStatus;
/**
* ExitStatus will be defaulted to the given BatchStatus
*
* @param status The BatchStatus to end with
* @param name The name of the state
*/
public EndState(BatchStatus status, String name) {
this(status, new ExitStatus(status.toString()), name);
}
private final boolean abandon;
/**
* @param status The BatchStatus to end with
@@ -51,14 +44,27 @@ public class EndState extends AbstractState {
* @param name The name of the state
*/
public EndState(BatchStatus status, ExitStatus exitStatus, String name) {
this(status, exitStatus, name, false);
}
/**
* @param status The BatchStatus to end with
* @param exitStatus The ExitStatus to end with
* @param name The name of the state
* @param abandon flag to indicate that previous step execution can be
* marked as abandoned (if there is one)
*
*/
public EndState(BatchStatus status, ExitStatus exitStatus, String name, boolean abandon) {
super(name);
this.status = status;
this.exitStatus = exitStatus;
this.abandon = abandon;
}
/**
* Return the {@link BatchStatus} and {@link ExitStatus} stored. If the
* {@link BatchStatus} is {@link BatchStatus#INCOMPLETE}, then mark it on the
* {@link BatchStatus} is {@link BatchStatus#FAILED}, then mark it on the
* {@link JobExecution} so that the job will know to stop.
*
* @see State#handle(FlowExecutor)
@@ -66,25 +72,45 @@ public class EndState extends AbstractState {
@Override
public FlowExecutionStatus handle(FlowExecutor executor) throws Exception {
JobExecution jobExecution = executor.getJobExecution();
// If there are no step executions, then we are at the beginning of a
// restart
synchronized (jobExecution) {
if (!jobExecution.getStepExecutions().isEmpty()) {
if (status == BatchStatus.INCOMPLETE) {
jobExecution.upgradeStatus(status);
jobExecution.setExitStatus(exitStatus);
/*
* If there are step executions, then we are not at the
* beginning of a restart.
*
* N.B. EndState has to be able to set the status directly, but
* only because the internal flows inside SplitStates contain
* EndState (which maybe they should not, since the JobExecution
* is not ending).
*/
jobExecution.setStatus(status);
jobExecution.setExitStatus(exitStatus);
if (status == BatchStatus.STOPPED) {
/*
* If we are in flight (not a restart) and we are supposed
* to signal a stop, then make sure that happens
* irrespective of the exit status.
*/
if (abandon) {
// Only if instructed to do so upgrade the status of
// last step execution...
executor.updateStepExecutionStatus();
}
return FlowExecutionStatus.STOPPED;
}
}
return new FlowExecutionStatus(status, exitStatus);
return new FlowExecutionStatus(exitStatus.getExitCode());
}
}
/* (non-Javadoc)
* @see org.springframework.batch.core.job.flow.State#validate(java.lang.String)
/*
* (non-Javadoc)
*
* @see
* org.springframework.batch.core.job.flow.State#validate(java.lang.String)
*/
public void validate(String pattern, String nextState) {
if (status != BatchStatus.INCOMPLETE && nextState != null) {
if (status != BatchStatus.STOPPED && nextState != null) {
throw new IllegalStateException("The transition for " + getClass().getSimpleName() + " [" + getName()
+ "] may not have a 'next' state.");
}

View File

@@ -71,6 +71,8 @@ public class SplitState extends AbstractState {
@Override
public FlowExecutionStatus handle(final FlowExecutor executor) throws Exception {
// TODO: collect the last StepExecution from the flows as well, so they
// can be abandoned if necessary
Collection<Future<FlowExecution>> tasks = new ArrayList<Future<FlowExecution>>();
for (final Flow flow : flows) {
@@ -103,8 +105,11 @@ public class SplitState extends AbstractState {
}
/* (non-Javadoc)
* @see org.springframework.batch.core.job.flow.State#validate(java.lang.String)
/*
* (non-Javadoc)
*
* @see
* org.springframework.batch.core.job.flow.State#validate(java.lang.String)
*/
public void validate(String pattern, String nextState) {
if (nextState == null) {

View File

@@ -52,6 +52,7 @@ public class StepState extends AbstractState implements StepHolder {
@Override
public FlowExecutionStatus handle(FlowExecutor executor) throws Exception {
executor.updateStepExecutionStatus();
return new FlowExecutionStatus(executor.executeStep(step));
}

View File

@@ -112,7 +112,7 @@ public class SimpleJobLauncher implements JobLauncher, InitializingBean {
+ "] and the following status: [" + jobExecution.getStatus() + "]");
}
catch (Throwable t) {
logger.info("Job: [" + job + "] failed with the following parameters: [" + jobParameters + "]", t);
logger.info("Job: [" + job + "] failed unexpectedly and fatally with the following parameters: [" + jobParameters + "]", t);
rethrow(t);
}
}

View File

@@ -122,7 +122,7 @@ public class TaskExecutorPartitionHandler implements PartitionHandler, Initializ
* Set the status in case the caller is tracking it through the
* JobExecution.
*/
stepExecution.setStatus(BatchStatus.INCOMPLETE);
stepExecution.setStatus(BatchStatus.FAILED);
stepExecution.setExitStatus(exitStatus);
result.add(stepExecution);
}

View File

@@ -114,7 +114,7 @@ public class SimpleJobRepository implements JobRepository {
}
BatchStatus status = execution.getStatus();
if (status == BatchStatus.COMPLETED || status == BatchStatus.FAILED) {
if (status == BatchStatus.COMPLETED || status == BatchStatus.ABANDONED) {
throw new JobInstanceAlreadyCompleteException(
"A job instance already exists and is complete for parameters=" + jobParameters
+ ". If you want to run this job again, change the parameters.");
@@ -237,7 +237,7 @@ public class SimpleJobRepository implements JobRepository {
private void checkForInterruption(StepExecution stepExecution) {
JobExecution jobExecution = stepExecution.getJobExecution();
jobExecutionDao.synchronizeStatus(jobExecution);
if (jobExecution.getStatus() == BatchStatus.STOPPING) {
if (jobExecution.isStopping()) {
stepExecution.setTerminateOnly();
}
}

View File

@@ -34,6 +34,7 @@ import org.springframework.batch.core.listener.CompositeStepExecutionListener;
import org.springframework.batch.core.repository.JobRepository;
import org.springframework.batch.core.scope.context.StepSynchronizationManager;
import org.springframework.batch.item.ExecutionContext;
import org.springframework.batch.repeat.RepeatException;
import org.springframework.beans.factory.BeanNameAware;
import org.springframework.beans.factory.InitializingBean;
import org.springframework.util.Assert;
@@ -189,12 +190,15 @@ public abstract class AbstractStep implements Step, InitializingBean, BeanNameAw
getCompositeListener().beforeStep(stepExecution);
open(stepExecution.getExecutionContext());
doExecute(stepExecution);
try {
doExecute(stepExecution);
} catch (RepeatException e) {
throw e.getCause();
}
exitStatus = stepExecution.getExitStatus();
// Check if someone is trying to stop us
if (stepExecution.isTerminateOnly()) {
stepExecution.setStatus(BatchStatus.INCOMPLETE);
throw new JobInterruptedException("JobExecution interrupted.");
}
@@ -263,10 +267,10 @@ public abstract class AbstractStep implements Step, InitializingBean, BeanNameAw
return BatchStatus.UNKNOWN;
}
else if (e instanceof JobInterruptedException || e.getCause() instanceof JobInterruptedException) {
return BatchStatus.INCOMPLETE;
return BatchStatus.STOPPED;
}
else {
return BatchStatus.INCOMPLETE;
return BatchStatus.FAILED;
}
}
@@ -321,7 +325,7 @@ public abstract class AbstractStep implements Step, InitializingBean, BeanNameAw
private ExitStatus getDefaultExitStatusForFailure(Throwable ex) {
ExitStatus exitStatus;
if (ex instanceof JobInterruptedException || ex.getCause() instanceof JobInterruptedException) {
exitStatus = ExitStatus.INTERRUPTED.addExitDescription(JobInterruptedException.class.getName());
exitStatus = ExitStatus.STOPPED.addExitDescription(JobInterruptedException.class.getName());
}
else if (ex instanceof NoSuchJobException || ex.getCause() instanceof NoSuchJobException) {
exitStatus = new ExitStatus(ExitCodeMapper.NO_SUCH_JOB, ex.getClass().getName());

View File

@@ -257,7 +257,7 @@ public class TaskletStep extends AbstractStep {
locked = true;
}
catch (InterruptedException e) {
stepExecution.setStatus(BatchStatus.INCOMPLETE);
stepExecution.setStatus(BatchStatus.STOPPED);
Thread.currentThread().interrupt();
}

View File

@@ -40,22 +40,22 @@ public class BatchStatusTests {
*/
@Test
public void testToString() {
assertEquals("FAILED", BatchStatus.FAILED.toString());
assertEquals("ABANDONED", BatchStatus.ABANDONED.toString());
}
@Test
public void testMaxStatus() {
assertEquals(BatchStatus.INCOMPLETE, BatchStatus.max(BatchStatus.INCOMPLETE,BatchStatus.COMPLETED));
assertEquals(BatchStatus.INCOMPLETE, BatchStatus.max(BatchStatus.COMPLETED, BatchStatus.INCOMPLETE));
assertEquals(BatchStatus.INCOMPLETE, BatchStatus.max(BatchStatus.INCOMPLETE, BatchStatus.INCOMPLETE));
assertEquals(BatchStatus.FAILED, BatchStatus.max(BatchStatus.FAILED,BatchStatus.COMPLETED));
assertEquals(BatchStatus.FAILED, BatchStatus.max(BatchStatus.COMPLETED, BatchStatus.FAILED));
assertEquals(BatchStatus.FAILED, BatchStatus.max(BatchStatus.FAILED, BatchStatus.FAILED));
assertEquals(BatchStatus.STARTED, BatchStatus.max(BatchStatus.STARTED, BatchStatus.STARTING));
assertEquals(BatchStatus.STARTED, BatchStatus.max(BatchStatus.COMPLETED, BatchStatus.STARTED));
}
@Test
public void testUpgradeStatusFinished() {
assertEquals(BatchStatus.INCOMPLETE, BatchStatus.INCOMPLETE.upgradeTo(BatchStatus.COMPLETED));
assertEquals(BatchStatus.INCOMPLETE, BatchStatus.COMPLETED.upgradeTo(BatchStatus.INCOMPLETE));
assertEquals(BatchStatus.FAILED, BatchStatus.FAILED.upgradeTo(BatchStatus.COMPLETED));
assertEquals(BatchStatus.FAILED, BatchStatus.COMPLETED.upgradeTo(BatchStatus.FAILED));
}
@Test
@@ -68,7 +68,7 @@ public class BatchStatusTests {
@Test
public void testIsRunning() {
assertFalse(BatchStatus.INCOMPLETE.isRunning());
assertFalse(BatchStatus.FAILED.isRunning());
assertFalse(BatchStatus.COMPLETED.isRunning());
assertTrue(BatchStatus.STARTED.isRunning());
assertTrue(BatchStatus.STARTING.isRunning());
@@ -76,7 +76,7 @@ public class BatchStatusTests {
@Test
public void testIsUnsuccessful() {
assertTrue(BatchStatus.INCOMPLETE.isUnsuccessful());
assertTrue(BatchStatus.FAILED.isUnsuccessful());
assertFalse(BatchStatus.COMPLETED.isUnsuccessful());
assertFalse(BatchStatus.STARTED.isUnsuccessful());
assertFalse(BatchStatus.STARTING.isUnsuccessful());
@@ -84,7 +84,7 @@ public class BatchStatusTests {
@Test
public void testGetStatus() {
assertEquals(BatchStatus.INCOMPLETE, BatchStatus.valueOf(BatchStatus.INCOMPLETE.toString()));
assertEquals(BatchStatus.FAILED, BatchStatus.valueOf(BatchStatus.FAILED.toString()));
}
@Test

View File

@@ -113,9 +113,9 @@ public class JobExecutionTests {
*/
@Test
public void testDowngradeStatus() {
execution.setStatus(BatchStatus.INCOMPLETE);
execution.setStatus(BatchStatus.FAILED);
execution.upgradeStatus(BatchStatus.COMPLETED);
assertEquals(BatchStatus.INCOMPLETE, execution.getStatus());
assertEquals(BatchStatus.FAILED, execution.getStatus());
}
/**

View File

@@ -286,9 +286,9 @@ public class StepExecutionTests {
*/
@Test
public void testDowngradeStatus() {
execution.setStatus(BatchStatus.INCOMPLETE);
execution.setStatus(BatchStatus.FAILED);
execution.upgradeStatus(BatchStatus.COMPLETED);
assertEquals(BatchStatus.INCOMPLETE, execution.getStatus());
assertEquals(BatchStatus.FAILED, execution.getStatus());
}
private StepExecution newStepExecution(Step step, Long long2) {

View File

@@ -46,7 +46,7 @@ public class OsgiBundleXmlApplicationContextFactoryTests {
expect(bundleContext.getBundle()).andReturn(bundle).anyTimes();
replay(bundleContext, bundle);
factory.setBundleContext(bundleContext);
// TODO: finish this...
// TODO: mock out the OSGi bundle resource...
// factory.createApplicationContext();
verify(bundleContext, bundle);
}

View File

@@ -44,7 +44,7 @@ public class DefaultFailureJobParserTests extends AbstractJobParserTests {
assertTrue(stepNamesList.contains("s1"));
assertTrue(stepNamesList.contains("fail"));
assertEquals(BatchStatus.INCOMPLETE, jobExecution.getStatus());
assertEquals(BatchStatus.FAILED, jobExecution.getStatus());
assertEquals(ExitStatus.FAILED.getExitCode(), jobExecution.getExitStatus().getExitCode());
StepExecution stepExecution1 = getStepExecution(jobExecution, "s1");
@@ -52,7 +52,7 @@ public class DefaultFailureJobParserTests extends AbstractJobParserTests {
assertEquals(ExitStatus.COMPLETED, stepExecution1.getExitStatus());
StepExecution stepExecution2 = getStepExecution(jobExecution, "fail");
assertEquals(BatchStatus.INCOMPLETE, stepExecution2.getStatus());
assertEquals(BatchStatus.FAILED, stepExecution2.getStatus());
assertEquals(ExitStatus.FAILED.getExitCode(), stepExecution2.getExitStatus().getExitCode());
}

View File

@@ -47,7 +47,7 @@ public class EndTransitionDefaultStatusJobParserTests extends AbstractJobParserT
assertEquals(ExitStatus.COMPLETED, jobExecution.getExitStatus());
StepExecution stepExecution1 = getStepExecution(jobExecution, "fail");
assertEquals(BatchStatus.INCOMPLETE, stepExecution1.getStatus());
assertEquals(BatchStatus.FAILED, stepExecution1.getStatus());
assertEquals(ExitStatus.FAILED.getExitCode(), stepExecution1.getExitStatus().getExitCode());
}

View File

@@ -57,7 +57,7 @@ public class EndTransitionJobParserTests extends AbstractJobParserTests {
assertEquals(ExitStatus.COMPLETED, stepExecution1.getExitStatus());
StepExecution stepExecution2 = getStepExecution(jobExecution, "fail");
assertEquals(BatchStatus.INCOMPLETE, stepExecution2.getStatus());
assertEquals(BatchStatus.FAILED, stepExecution2.getStatus());
assertEquals(ExitStatus.FAILED.getExitCode(), stepExecution2.getExitStatus().getExitCode());
//

View File

@@ -17,7 +17,6 @@ package org.springframework.batch.core.configuration.xml;
import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertTrue;
import static org.junit.Assert.fail;
import org.junit.Test;
import org.junit.runner.RunWith;
@@ -25,7 +24,6 @@ import org.springframework.batch.core.BatchStatus;
import org.springframework.batch.core.ExitStatus;
import org.springframework.batch.core.JobExecution;
import org.springframework.batch.core.StepExecution;
import org.springframework.batch.core.repository.JobInstanceAlreadyCompleteException;
import org.springframework.test.context.ContextConfiguration;
import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;
@@ -50,28 +48,27 @@ public class FailTransitionJobParserTests extends AbstractJobParserTests {
assertTrue(stepNamesList.contains("fail"));
assertEquals(BatchStatus.FAILED, jobExecution.getStatus());
assertEquals("EARLY TERMINATION (FAIL)", jobExecution.getExitStatus().getExitCode());
assertEquals("EARLY TERMINATION (FAIL)", jobExecution.getExitStatus()
.getExitCode());
StepExecution stepExecution1 = getStepExecution(jobExecution, "s1");
assertEquals(BatchStatus.COMPLETED, stepExecution1.getStatus());
assertEquals(ExitStatus.COMPLETED, stepExecution1.getExitStatus());
StepExecution stepExecution2 = getStepExecution(jobExecution, "fail");
assertEquals(BatchStatus.INCOMPLETE, stepExecution2.getStatus());
assertEquals(ExitStatus.FAILED.getExitCode(), stepExecution2.getExitStatus().getExitCode());
assertEquals(BatchStatus.FAILED, stepExecution2.getStatus());
assertEquals(ExitStatus.FAILED.getExitCode(), stepExecution2
.getExitStatus().getExitCode());
//
// Second Launch
//
stepNamesList.clear();
try {
jobExecution = createJobExecution();
fail("JobInstanceAlreadyCompleteException expected");
} catch (JobInstanceAlreadyCompleteException e) {
//
// Expected
//
}
jobExecution = createJobExecution();
job.execute(jobExecution);
assertEquals(1, stepNamesList.size());
assertTrue(stepNamesList.contains("fail"));
assertEquals(BatchStatus.FAILED, jobExecution.getStatus());
}

View File

@@ -45,7 +45,7 @@ public class NextAttributeJobParserTests extends AbstractJobParserTests {
assertTrue(stepNamesList.contains("s1"));
assertTrue(stepNamesList.contains("fail"));
assertEquals(BatchStatus.INCOMPLETE, jobExecution.getStatus());
assertEquals(BatchStatus.FAILED, jobExecution.getStatus());
assertEquals("FAILED", jobExecution.getExitStatus().getExitCode());
StepExecution stepExecution1 = getStepExecution(jobExecution, "s1");
@@ -53,7 +53,7 @@ public class NextAttributeJobParserTests extends AbstractJobParserTests {
assertEquals(ExitStatus.COMPLETED, stepExecution1.getExitStatus());
StepExecution stepExecution2 = getStepExecution(jobExecution, "fail");
assertEquals(BatchStatus.INCOMPLETE, stepExecution2.getStatus());
assertEquals(BatchStatus.FAILED, stepExecution2.getStatus());
assertEquals(ExitStatus.FAILED.getExitCode(), stepExecution2.getExitStatus().getExitCode());
}

View File

@@ -44,7 +44,7 @@ public class SplitDifferentResultsFailFirstJobParserTests extends AbstractJobPar
assertTrue(stepNamesList.contains("s1"));
assertTrue(stepNamesList.contains("fail"));
assertEquals(BatchStatus.INCOMPLETE, jobExecution.getStatus());
assertEquals(BatchStatus.FAILED, jobExecution.getStatus());
assertEquals(ExitStatus.FAILED, jobExecution.getExitStatus());
StepExecution stepExecution1 = getStepExecution(jobExecution, "s1");
@@ -52,7 +52,7 @@ public class SplitDifferentResultsFailFirstJobParserTests extends AbstractJobPar
assertEquals(ExitStatus.COMPLETED, stepExecution1.getExitStatus());
StepExecution stepExecution2 = getStepExecution(jobExecution, "fail");
assertEquals(BatchStatus.INCOMPLETE, stepExecution2.getStatus());
assertEquals(BatchStatus.FAILED, stepExecution2.getStatus());
assertEquals(ExitStatus.FAILED.getExitCode(), stepExecution2.getExitStatus().getExitCode());
}

View File

@@ -40,7 +40,7 @@ public class SplitDifferentResultsFailSecondJobParserTests extends AbstractJobPa
JobExecution jobExecution = createJobExecution();
job.execute(jobExecution);
assertEquals(3, stepNamesList.size());
assertEquals("Wrong step anmes: "+stepNamesList, 3, stepNamesList.size());
assertTrue(stepNamesList.contains("s1"));
assertTrue(stepNamesList.contains("fail"));
assertTrue(stepNamesList.contains("s3"));
@@ -53,7 +53,7 @@ public class SplitDifferentResultsFailSecondJobParserTests extends AbstractJobPa
assertEquals(ExitStatus.COMPLETED, stepExecution1.getExitStatus());
StepExecution stepExecution2 = getStepExecution(jobExecution, "fail");
assertEquals(BatchStatus.INCOMPLETE, stepExecution2.getStatus());
assertEquals(BatchStatus.FAILED, stepExecution2.getStatus());
assertEquals(ExitStatus.FAILED.getExitCode(), stepExecution2.getExitStatus().getExitCode());
StepExecution stepExecution3 = getStepExecution(jobExecution, "s3");

View File

@@ -0,0 +1,74 @@
/*
* Copyright 2006-2007 the original author or authors.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.springframework.batch.core.configuration.xml;
import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertTrue;
import org.junit.Test;
import org.junit.runner.RunWith;
import org.springframework.batch.core.BatchStatus;
import org.springframework.batch.core.ExitStatus;
import org.springframework.batch.core.JobExecution;
import org.springframework.batch.core.StepExecution;
import org.springframework.test.context.ContextConfiguration;
import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;
/**
* @author Dave Syer
*
*/
@ContextConfiguration
@RunWith(SpringJUnit4ClassRunner.class)
public class StopAndRestartJobParserTests extends AbstractJobParserTests {
@Test
public void testStopIncomplete() throws Exception {
//
// First Launch
//
JobExecution jobExecution = createJobExecution();
job.execute(jobExecution);
assertEquals(1, stepNamesList.size());
assertTrue(stepNamesList.contains("s1"));
assertEquals(BatchStatus.STOPPED, jobExecution.getStatus());
assertEquals(ExitStatus.STOPPED.getExitCode(), jobExecution.getExitStatus().getExitCode());
StepExecution stepExecution1 = getStepExecution(jobExecution, "s1");
assertEquals(BatchStatus.COMPLETED, stepExecution1.getStatus());
assertEquals(ExitStatus.COMPLETED.getExitCode(), stepExecution1.getExitStatus().getExitCode());
//
// Second Launch
//
stepNamesList.clear();
jobExecution = createJobExecution();
job.execute(jobExecution);
assertEquals(1, stepNamesList.size()); // step1 is not executed
assertTrue(stepNamesList.contains("s2"));
assertEquals(BatchStatus.COMPLETED, jobExecution.getStatus());
assertEquals(ExitStatus.COMPLETED, jobExecution.getExitStatus());
StepExecution stepExecution2 = getStepExecution(jobExecution, "s2");
assertEquals(BatchStatus.COMPLETED, stepExecution2.getStatus());
assertEquals(ExitStatus.COMPLETED, stepExecution2.getExitStatus());
}
}

View File

@@ -24,7 +24,6 @@ import org.springframework.batch.core.BatchStatus;
import org.springframework.batch.core.ExitStatus;
import org.springframework.batch.core.JobExecution;
import org.springframework.batch.core.StepExecution;
import org.springframework.batch.core.job.flow.JobExecutionDecider;
import org.springframework.test.context.ContextConfiguration;
import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;
@@ -44,14 +43,14 @@ public class StopIncompleteJobParserTests extends AbstractJobParserTests {
//
JobExecution jobExecution = createJobExecution();
job.execute(jobExecution);
assertTrue("Wrong steps executed: "+stepNamesList, stepNamesList.contains("fail"));
assertEquals(1, stepNamesList.size());
assertTrue(stepNamesList.contains("fail"));
assertEquals(BatchStatus.INCOMPLETE, jobExecution.getStatus());
assertEquals(ExitStatus.FAILED.getExitCode(), jobExecution.getExitStatus().getExitCode());
assertEquals(BatchStatus.STOPPED, jobExecution.getStatus());
assertEquals(ExitStatus.STOPPED.getExitCode(), jobExecution.getExitStatus().getExitCode());
StepExecution stepExecution1 = getStepExecution(jobExecution, "fail");
assertEquals(BatchStatus.FAILED, stepExecution1.getStatus());
assertEquals(BatchStatus.ABANDONED, stepExecution1.getStatus());
assertEquals(ExitStatus.FAILED.getExitCode(), stepExecution1.getExitStatus().getExitCode());
//
@@ -60,8 +59,8 @@ public class StopIncompleteJobParserTests extends AbstractJobParserTests {
stepNamesList.clear();
jobExecution = createJobExecution();
job.execute(jobExecution);
assertTrue("Wrong steps executed: "+stepNamesList, stepNamesList.contains("s2"));
assertEquals(1, stepNamesList.size()); // step1 is not executed
assertTrue(stepNamesList.contains("s2"));
assertEquals(BatchStatus.COMPLETED, jobExecution.getStatus());
assertEquals(ExitStatus.COMPLETED, jobExecution.getExitStatus());
@@ -72,10 +71,4 @@ public class StopIncompleteJobParserTests extends AbstractJobParserTests {
}
public static class TestDecider implements JobExecutionDecider {
public String decide(JobExecution jobExecution, StepExecution stepExecution) {
return "FOO";
}
}
}

View File

@@ -47,8 +47,8 @@ public class StopJobParserTests extends AbstractJobParserTests {
assertEquals(1, stepNamesList.size());
assertTrue(stepNamesList.contains("s1"));
assertEquals(BatchStatus.INCOMPLETE, jobExecution.getStatus());
assertEquals(ExitStatus.FAILED.getExitCode(), jobExecution.getExitStatus().getExitCode());
assertEquals(BatchStatus.STOPPED, jobExecution.getStatus());
assertEquals(ExitStatus.STOPPED.getExitCode(), jobExecution.getExitStatus().getExitCode());
StepExecution stepExecution1 = getStepExecution(jobExecution, "s1");
assertEquals(BatchStatus.COMPLETED, stepExecution1.getStatus());

View File

@@ -61,8 +61,8 @@ public class StopRestartOnCompletedStepJobParserTests extends AbstractJobParserT
assertEquals(1, stepNamesList.size());
assertTrue(stepNamesList.contains("s1"));
assertEquals(BatchStatus.INCOMPLETE, jobExecution.getStatus());
assertEquals(ExitStatus.FAILED.getExitCode(), jobExecution.getExitStatus().getExitCode());
assertEquals(BatchStatus.STOPPED, jobExecution.getStatus());
assertEquals(ExitStatus.STOPPED.getExitCode(), jobExecution.getExitStatus().getExitCode());
StepExecution stepExecution1 = getStepExecution(jobExecution, "s1");
assertEquals(BatchStatus.COMPLETED, stepExecution1.getStatus());

View File

@@ -140,7 +140,7 @@ public class AbstractJobTests {
// simulate restart and check the job execution context's content survives
execution.setEndTime(new Date());
execution.setStatus(BatchStatus.INCOMPLETE);
execution.setStatus(BatchStatus.FAILED);
repository.update(execution);
JobExecution restarted = repository.createJobExecution("testHandleStepJob", new JobParameters());

View File

@@ -248,10 +248,10 @@ public class SimpleJobTests {
final JobInterruptedException exception = new JobInterruptedException("Interrupt!");
step1.setProcessException(exception);
job.execute(jobExecution);
assertEquals(2, jobExecution.getAllFailureExceptions().size());
assertEquals(1, jobExecution.getAllFailureExceptions().size());
assertEquals(exception, jobExecution.getStepExecutions().iterator().next().getFailureExceptions().get(0));
assertEquals(0, list.size());
checkRepository(BatchStatus.INCOMPLETE, ExitStatus.FAILED);
checkRepository(BatchStatus.STOPPED, ExitStatus.STOPPED);
}
@Test
@@ -265,8 +265,8 @@ public class SimpleJobTests {
assertEquals(1, jobExecution.getAllFailureExceptions().size());
assertEquals(exception, jobExecution.getAllFailureExceptions().get(0));
assertEquals(0, list.size());
assertEquals(BatchStatus.INCOMPLETE, jobExecution.getStatus());
checkRepository(BatchStatus.INCOMPLETE, ExitStatus.FAILED);
assertEquals(BatchStatus.FAILED, jobExecution.getStatus());
checkRepository(BatchStatus.FAILED, ExitStatus.FAILED);
}
@Test
@@ -283,7 +283,7 @@ public class SimpleJobTests {
assertEquals(1, jobExecution.getAllFailureExceptions().size());
assertEquals(exception, jobExecution.getAllFailureExceptions().get(0));
assertEquals(1, list.size());
checkRepository(BatchStatus.INCOMPLETE, ExitStatus.FAILED);
checkRepository(BatchStatus.FAILED, ExitStatus.FAILED);
}
@Test
@@ -297,7 +297,7 @@ public class SimpleJobTests {
assertEquals(1, jobExecution.getAllFailureExceptions().size());
assertEquals(exception, jobExecution.getAllFailureExceptions().get(0));
assertEquals(0, list.size());
checkRepository(BatchStatus.INCOMPLETE, ExitStatus.FAILED);
checkRepository(BatchStatus.FAILED, ExitStatus.FAILED);
}
@Test
@@ -352,7 +352,7 @@ public class SimpleJobTests {
job.execute(jobExecution);
assertEquals(0, list.size());
checkRepository(BatchStatus.INCOMPLETE, ExitStatus.NOOP);
checkRepository(BatchStatus.FAILED, ExitStatus.NOOP);
ExitStatus exitStatus = jobExecution.getExitStatus();
assertEquals(ExitStatus.NOOP.getExitCode(), exitStatus.getExitCode());
}
@@ -386,7 +386,7 @@ public class SimpleJobTests {
job.setJobExecutionListeners(new JobExecutionListener[] { listener });
job.execute(jobExecution);
assertEquals(BatchStatus.INCOMPLETE, jobExecution.getStatus());
assertEquals(BatchStatus.STOPPED, jobExecution.getStatus());
verify(listener);
}
@@ -438,7 +438,7 @@ public class SimpleJobTests {
job.setSteps(Arrays.asList(new Step[] { step1, step2 }));
job.execute(jobExecution);
assertEquals(BatchStatus.INCOMPLETE, jobExecution.getStatus());
assertEquals(BatchStatus.STOPPED, jobExecution.getStatus());
assertEquals(1, jobExecution.getAllFailureExceptions().size());
Throwable expected = jobExecution.getAllFailureExceptions().get(0);
assertTrue("Wrong exception " + expected, expected instanceof JobInterruptedException);
@@ -534,21 +534,27 @@ public class SimpleJobTests {
jobRepository.update(stepExecution);
jobRepository.updateExecutionContext(stepExecution);
if (exception instanceof JobInterruptedException) {
stepExecution.setExitStatus(ExitStatus.FAILED);
stepExecution.setStatus(BatchStatus.FAILED);
stepExecution.addFailureException(exception);
throw (JobInterruptedException)exception;
}
if (exception instanceof RuntimeException) {
stepExecution.setExitStatus(ExitStatus.FAILED);
stepExecution.setStatus(BatchStatus.INCOMPLETE);
stepExecution.setStatus(BatchStatus.FAILED);
stepExecution.addFailureException(exception);
return;
}
if (exception instanceof Error) {
stepExecution.setExitStatus(ExitStatus.FAILED);
stepExecution.setStatus(BatchStatus.INCOMPLETE);
stepExecution.setStatus(BatchStatus.FAILED);
stepExecution.addFailureException(exception);
return;
}
if (exception instanceof JobInterruptedException) {
stepExecution.setExitStatus(ExitStatus.FAILED);
stepExecution.setStatus(BatchStatus.STOPPING);
stepExecution.setStatus(BatchStatus.FAILED);
stepExecution.addFailureException(exception);
return;
}

View File

@@ -31,7 +31,7 @@ public class FlowExecutionTests {
public void testBasicProperties() throws Exception {
FlowExecution execution = new FlowExecution("foo", new FlowExecutionStatus("BAR"));
assertEquals("foo",execution.getName());
assertEquals("BAR",execution.getStatus().getExitStatus().getExitCode());
assertEquals("BAR",execution.getStatus().getStatus());
}
@Test
@@ -45,7 +45,7 @@ public class FlowExecutionTests {
@Test
public void testEnumOrdering() throws Exception {
FlowExecution first = new FlowExecution("foo", FlowExecutionStatus.COMPLETED);
FlowExecution second = new FlowExecution("foo", FlowExecutionStatus.INCOMPLETE);
FlowExecution second = new FlowExecution("foo", FlowExecutionStatus.FAILED);
assertTrue("Should be negative",first.compareTo(second)<0);
assertTrue("Should be positive",second.compareTo(first)>0);
}

View File

@@ -53,7 +53,7 @@ public class FlowJobTests {
private JobExecution jobExecution;
private JobRepository jobRepository;
private boolean fail = false;
@Before
@@ -64,18 +64,26 @@ public class FlowJobTests {
factory.afterPropertiesSet();
jobRepository = (JobRepository) factory.getObject();
job.setJobRepository(jobRepository);
jobExecution = jobRepository.createJobExecution("job", new JobParameters());
jobExecution = jobRepository.createJobExecution("job",
new JobParameters());
}
@Test
public void testTwoSteps() throws Exception {
SimpleFlow flow = new SimpleFlow("job");
List<StateTransition> transitions = new ArrayList<StateTransition>();
transitions.add(StateTransition.createStateTransition(new StepState(new StubStep("step1")), "step2"));
transitions.add(StateTransition.createStateTransition(new StepState(new StubStep("step2")), ExitStatus.FAILED.getExitCode(), "end0"));
transitions.add(StateTransition.createStateTransition(new StepState(new StubStep("step2")), ExitStatus.COMPLETED.getExitCode(), "end1"));
transitions.add(StateTransition.createEndStateTransition(new EndState(BatchStatus.INCOMPLETE, ExitStatus.FAILED, "end0")));
transitions.add(StateTransition.createEndStateTransition(new EndState(BatchStatus.COMPLETED, "end1")));
transitions.add(StateTransition.createStateTransition(new StepState(
new StubStep("step1")), "step2"));
transitions.add(StateTransition
.createStateTransition(new StepState(new StubStep("step2")),
ExitStatus.FAILED.getExitCode(), "end0"));
transitions.add(StateTransition.createStateTransition(new StepState(
new StubStep("step2")), ExitStatus.COMPLETED.getExitCode(),
"end1"));
transitions.add(StateTransition.createEndStateTransition(new EndState(
BatchStatus.FAILED, ExitStatus.FAILED, "end0")));
transitions.add(StateTransition.createEndStateTransition(new EndState(
BatchStatus.COMPLETED, ExitStatus.COMPLETED, "end1")));
flow.setStateTransitions(transitions);
job.setFlow(flow);
job.afterPropertiesSet();
@@ -89,11 +97,18 @@ public class FlowJobTests {
public void testFailedStep() throws Exception {
SimpleFlow flow = new SimpleFlow("job");
List<StateTransition> transitions = new ArrayList<StateTransition>();
transitions.add(StateTransition.createStateTransition(new StateSupport("step1", FlowExecutionStatus.INCOMPLETE), "step2"));
transitions.add(StateTransition.createStateTransition(new StepState(new StubStep("step2")), ExitStatus.FAILED.getExitCode(), "end0"));
transitions.add(StateTransition.createStateTransition(new StepState(new StubStep("step2")), ExitStatus.COMPLETED.getExitCode(), "end1"));
transitions.add(StateTransition.createEndStateTransition(new EndState(BatchStatus.INCOMPLETE, ExitStatus.FAILED, "end0")));
transitions.add(StateTransition.createEndStateTransition(new EndState(BatchStatus.COMPLETED, "end1")));
transitions.add(StateTransition.createStateTransition(new StateSupport(
"step1", FlowExecutionStatus.FAILED), "step2"));
transitions.add(StateTransition
.createStateTransition(new StepState(new StubStep("step2")),
ExitStatus.FAILED.getExitCode(), "end0"));
transitions.add(StateTransition.createStateTransition(new StepState(
new StubStep("step2")), ExitStatus.COMPLETED.getExitCode(),
"end1"));
transitions.add(StateTransition.createEndStateTransition(new EndState(
BatchStatus.FAILED, ExitStatus.FAILED, "end0")));
transitions.add(StateTransition.createEndStateTransition(new EndState(
BatchStatus.COMPLETED, ExitStatus.COMPLETED, "end1")));
flow.setStateTransitions(transitions);
job.setFlow(flow);
job.afterPropertiesSet();
@@ -108,24 +123,30 @@ public class FlowJobTests {
public void testFailedStepRestarted() throws Exception {
SimpleFlow flow = new SimpleFlow("job");
List<StateTransition> transitions = new ArrayList<StateTransition>();
transitions.add(StateTransition.createStateTransition(new StepState(new StubStep("step1")), "step2"));
transitions.add(StateTransition.createStateTransition(new StepState(
new StubStep("step1")), "step2"));
State step2State = new StateSupport("step2") {
@Override
public FlowExecutionStatus handle(FlowExecutor executor) throws Exception {
public FlowExecutionStatus handle(FlowExecutor executor)
throws Exception {
JobExecution jobExecution = executor.getJobExecution();
jobExecution.getStepExecutions().add(new StepExecution(getName(), jobExecution));
jobExecution.getStepExecutions().add(
new StepExecution(getName(), jobExecution));
if (fail) {
return FlowExecutionStatus.INCOMPLETE;
}
else {
return FlowExecutionStatus.FAILED;
} else {
return FlowExecutionStatus.COMPLETED;
}
}
};
transitions.add(StateTransition.createStateTransition(step2State, ExitStatus.COMPLETED.getExitCode(), "end0"));
transitions.add(StateTransition.createStateTransition(step2State, ExitStatus.FAILED.getExitCode(), "end1"));
transitions.add(StateTransition.createEndStateTransition(new EndState(BatchStatus.COMPLETED, "end0")));
transitions.add(StateTransition.createEndStateTransition(new EndState(BatchStatus.INCOMPLETE, ExitStatus.FAILED, "end1")));
transitions.add(StateTransition.createStateTransition(step2State,
ExitStatus.COMPLETED.getExitCode(), "end0"));
transitions.add(StateTransition.createStateTransition(step2State,
ExitStatus.FAILED.getExitCode(), "end1"));
transitions.add(StateTransition.createEndStateTransition(new EndState(
BatchStatus.COMPLETED, ExitStatus.COMPLETED, "end0")));
transitions.add(StateTransition.createEndStateTransition(new EndState(
BatchStatus.FAILED, ExitStatus.FAILED, "end1")));
flow.setStateTransitions(transitions);
job.setFlow(flow);
job.afterPropertiesSet();
@@ -134,7 +155,8 @@ public class FlowJobTests {
assertEquals(ExitStatus.FAILED, jobExecution.getExitStatus());
assertEquals(2, jobExecution.getStepExecutions().size());
jobRepository.update(jobExecution);
jobExecution = jobRepository.createJobExecution("job", new JobParameters());
jobExecution = jobRepository.createJobExecution("job",
new JobParameters());
fail = false;
job.execute(jobExecution);
assertEquals(ExitStatus.COMPLETED, jobExecution.getExitStatus());
@@ -145,64 +167,77 @@ public class FlowJobTests {
public void testStoppingStep() throws Exception {
SimpleFlow flow = new SimpleFlow("job");
List<StateTransition> transitions = new ArrayList<StateTransition>();
transitions.add(StateTransition.createStateTransition(new StepState(new StubStep("step1")), "step2"));
State state2 = new StateSupport("step2", FlowExecutionStatus.INCOMPLETE);
transitions.add(StateTransition.createStateTransition(state2, ExitStatus.FAILED.getExitCode(), "end0"));
transitions.add(StateTransition.createStateTransition(state2, ExitStatus.COMPLETED.getExitCode(), "end1"));
transitions.add(StateTransition.createStateTransition(new EndState(BatchStatus.INCOMPLETE, ExitStatus.FAILED, "end0"), "step3"));
transitions.add(StateTransition.createEndStateTransition(new EndState(BatchStatus.COMPLETED, "end1")));
transitions.add(StateTransition.createStateTransition(new StepState(new StubStep("step3")), "end2"));
transitions.add(StateTransition.createEndStateTransition(new EndState(BatchStatus.COMPLETED, "end2")));
transitions.add(StateTransition.createStateTransition(new StepState(
new StubStep("step1")), "step2"));
State state2 = new StateSupport("step2", FlowExecutionStatus.FAILED);
transitions.add(StateTransition.createStateTransition(state2,
ExitStatus.FAILED.getExitCode(), "end0"));
transitions.add(StateTransition.createStateTransition(state2,
ExitStatus.COMPLETED.getExitCode(), "end1"));
transitions.add(StateTransition.createStateTransition(new EndState(
BatchStatus.STOPPED, ExitStatus.STOPPED, "end0"), "step3"));
transitions.add(StateTransition.createEndStateTransition(new EndState(
BatchStatus.COMPLETED, ExitStatus.COMPLETED, "end1")));
transitions.add(StateTransition.createStateTransition(new StepState(
new StubStep("step3")), "end2"));
transitions.add(StateTransition.createEndStateTransition(new EndState(
BatchStatus.COMPLETED, ExitStatus.COMPLETED, "end2")));
flow.setStateTransitions(transitions);
job.setFlow(flow);
job.afterPropertiesSet();
try {
job.doExecute(jobExecution);
fail("Expected JobInterruptedException");
}
catch (JobInterruptedException e) {
// expected
}
job.doExecute(jobExecution);
assertEquals(2, jobExecution.getStepExecutions().size());
assertEquals(BatchStatus.STOPPED, jobExecution.getStatus());
}
@Test
public void testEndStateStopped() throws Exception {
SimpleFlow flow = new SimpleFlow("job");
List<StateTransition> transitions = new ArrayList<StateTransition>();
transitions.add(StateTransition.createStateTransition(new StepState(new StubStep("step1")), "end"));
transitions.add(StateTransition.createStateTransition(new EndState(BatchStatus.INCOMPLETE, ExitStatus.FAILED, "end"), "step2"));
transitions.add(StateTransition.createStateTransition(new StepState(new StubStep("step2")), ExitStatus.FAILED.getExitCode(), "end0"));
transitions.add(StateTransition.createStateTransition(new StepState(new StubStep("step2")), ExitStatus.COMPLETED.getExitCode(), "end1"));
transitions.add(StateTransition.createEndStateTransition(new EndState(BatchStatus.INCOMPLETE, ExitStatus.FAILED, "end0")));
transitions.add(StateTransition.createEndStateTransition(new EndState(BatchStatus.COMPLETED, "end1")));
transitions.add(StateTransition.createStateTransition(new StepState(
new StubStep("step1")), "end"));
transitions.add(StateTransition.createStateTransition(new EndState(
BatchStatus.STOPPED, ExitStatus.STOPPED, "end"), "step2"));
transitions.add(StateTransition
.createStateTransition(new StepState(new StubStep("step2")),
ExitStatus.FAILED.getExitCode(), "end0"));
transitions.add(StateTransition.createStateTransition(new StepState(
new StubStep("step2")), ExitStatus.COMPLETED.getExitCode(),
"end1"));
transitions.add(StateTransition.createEndStateTransition(new EndState(
BatchStatus.FAILED, ExitStatus.FAILED, "end0")));
transitions.add(StateTransition.createEndStateTransition(new EndState(
BatchStatus.COMPLETED, ExitStatus.COMPLETED, "end1")));
flow.setStateTransitions(transitions);
job.setFlow(flow);
job.afterPropertiesSet();
try {
job.doExecute(jobExecution);
fail("Expected JobInterruptedException");
}
catch (JobInterruptedException e) {
// expected
}
job.doExecute(jobExecution);
assertEquals(1, jobExecution.getStepExecutions().size());
assertEquals(BatchStatus.STOPPED, jobExecution.getStatus());
}
public void testEndStateFailed() throws Exception {
SimpleFlow flow = new SimpleFlow("job");
List<StateTransition> transitions = new ArrayList<StateTransition>();
transitions.add(StateTransition.createStateTransition(new StepState(new StubStep("step1")), "end"));
transitions.add(StateTransition.createStateTransition(new EndState(BatchStatus.INCOMPLETE, ExitStatus.FAILED, "end"), "step2"));
transitions.add(StateTransition.createStateTransition(new StepState(new StubStep("step2")), ExitStatus.FAILED.getExitCode(), "end0"));
transitions.add(StateTransition.createStateTransition(new StepState(new StubStep("step2")), ExitStatus.COMPLETED.getExitCode(), "end1"));
transitions.add(StateTransition.createEndStateTransition(new EndState(BatchStatus.INCOMPLETE, ExitStatus.FAILED, "end0")));
transitions.add(StateTransition.createEndStateTransition(new EndState(BatchStatus.COMPLETED, "end1")));
transitions.add(StateTransition.createStateTransition(new StepState(
new StubStep("step1")), "end"));
transitions.add(StateTransition.createStateTransition(new EndState(
BatchStatus.FAILED, ExitStatus.FAILED, "end"), "step2"));
transitions.add(StateTransition
.createStateTransition(new StepState(new StubStep("step2")),
ExitStatus.FAILED.getExitCode(), "end0"));
transitions.add(StateTransition.createStateTransition(new StepState(
new StubStep("step2")), ExitStatus.COMPLETED.getExitCode(),
"end1"));
transitions.add(StateTransition.createEndStateTransition(new EndState(
BatchStatus.FAILED, ExitStatus.FAILED, "end0")));
transitions.add(StateTransition.createEndStateTransition(new EndState(
BatchStatus.COMPLETED, ExitStatus.COMPLETED, "end1")));
flow.setStateTransitions(transitions);
job.setFlow(flow);
job.afterPropertiesSet();
job.doExecute(jobExecution);
assertEquals(BatchStatus.INCOMPLETE, jobExecution.getStatus());
assertEquals(BatchStatus.FAILED, jobExecution.getStatus());
assertEquals(1, jobExecution.getStepExecutions().size());
}
@@ -210,22 +245,31 @@ public class FlowJobTests {
public void testEndStateStoppedWithRestart() throws Exception {
SimpleFlow flow = new SimpleFlow("job");
List<StateTransition> transitions = new ArrayList<StateTransition>();
transitions.add(StateTransition.createStateTransition(new StepState(new StubStep("step1")), "end"));
transitions.add(StateTransition.createStateTransition(new EndState(BatchStatus.INCOMPLETE, ExitStatus.FAILED, "end"), "step2"));
transitions.add(StateTransition.createStateTransition(new StepState(new StubStep("step2")), ExitStatus.COMPLETED.getExitCode(), "end0"));
transitions.add(StateTransition.createStateTransition(new StepState(new StubStep("step2")), ExitStatus.FAILED.getExitCode(), "end1"));
transitions.add(StateTransition.createEndStateTransition(new EndState(BatchStatus.COMPLETED, "end0")));
transitions.add(StateTransition.createEndStateTransition(new EndState(BatchStatus.INCOMPLETE, ExitStatus.FAILED, "end1")));
transitions.add(StateTransition.createStateTransition(new StepState(
new StubStep("step1")), "end"));
transitions.add(StateTransition.createStateTransition(new EndState(
BatchStatus.STOPPED, ExitStatus.STOPPED, "end"), "step2"));
transitions.add(StateTransition.createStateTransition(new StepState(
new StubStep("step2")), ExitStatus.COMPLETED.getExitCode(),
"end0"));
transitions.add(StateTransition
.createStateTransition(new StepState(new StubStep("step2")),
ExitStatus.FAILED.getExitCode(), "end1"));
transitions.add(StateTransition.createEndStateTransition(new EndState(
BatchStatus.COMPLETED, ExitStatus.COMPLETED, "end0")));
transitions.add(StateTransition.createEndStateTransition(new EndState(
BatchStatus.FAILED, ExitStatus.FAILED, "end1")));
flow.setStateTransitions(transitions);
job.setFlow(flow);
job.afterPropertiesSet();
// To test a restart we have to use the AbstractJob.execute()...
job.execute(jobExecution);
assertEquals(BatchStatus.INCOMPLETE, jobExecution.getStatus());
assertEquals(BatchStatus.STOPPED, jobExecution.getStatus());
assertEquals(1, jobExecution.getStepExecutions().size());
jobExecution = jobRepository.createJobExecution("job", new JobParameters());
jobExecution = jobRepository.createJobExecution("job",
new JobParameters());
job.execute(jobExecution);
assertEquals(BatchStatus.COMPLETED, jobExecution.getStatus());
assertEquals(1, jobExecution.getStepExecutions().size());
@@ -236,16 +280,30 @@ public class FlowJobTests {
public void testBranching() throws Exception {
SimpleFlow flow = new SimpleFlow("job");
List<StateTransition> transitions = new ArrayList<StateTransition>();
transitions.add(StateTransition.createStateTransition(new StepState(new StubStep("step1")), "step2"));
transitions.add(StateTransition.createStateTransition(new StepState(new StubStep("step1")), "COMPLETED", "step3"));
transitions.add(StateTransition.createStateTransition(new StepState(new StubStep("step2")), ExitStatus.COMPLETED.getExitCode(), "end0"));
transitions.add(StateTransition.createStateTransition(new StepState(new StubStep("step2")), ExitStatus.FAILED.getExitCode(), "end1"));
transitions.add(StateTransition.createEndStateTransition(new EndState(BatchStatus.COMPLETED, "end0")));
transitions.add(StateTransition.createEndStateTransition(new EndState(BatchStatus.INCOMPLETE, ExitStatus.FAILED, "end1")));
transitions.add(StateTransition.createStateTransition(new StepState(new StubStep("step3")), ExitStatus.FAILED.getExitCode(), "end2"));
transitions.add(StateTransition.createStateTransition(new StepState(new StubStep("step3")), ExitStatus.COMPLETED.getExitCode(), "end3"));
transitions.add(StateTransition.createEndStateTransition(new EndState(BatchStatus.INCOMPLETE, ExitStatus.FAILED, "end2")));
transitions.add(StateTransition.createEndStateTransition(new EndState(BatchStatus.COMPLETED, "end3")));
transitions.add(StateTransition.createStateTransition(new StepState(
new StubStep("step1")), "step2"));
transitions.add(StateTransition.createStateTransition(new StepState(
new StubStep("step1")), "COMPLETED", "step3"));
transitions.add(StateTransition.createStateTransition(new StepState(
new StubStep("step2")), ExitStatus.COMPLETED.getExitCode(),
"end0"));
transitions.add(StateTransition
.createStateTransition(new StepState(new StubStep("step2")),
ExitStatus.FAILED.getExitCode(), "end1"));
transitions.add(StateTransition.createEndStateTransition(new EndState(
BatchStatus.COMPLETED, ExitStatus.COMPLETED, "end0")));
transitions.add(StateTransition.createEndStateTransition(new EndState(
BatchStatus.FAILED, ExitStatus.FAILED, "end1")));
transitions.add(StateTransition
.createStateTransition(new StepState(new StubStep("step3")),
ExitStatus.FAILED.getExitCode(), "end2"));
transitions.add(StateTransition.createStateTransition(new StepState(
new StubStep("step3")), ExitStatus.COMPLETED.getExitCode(),
"end3"));
transitions.add(StateTransition.createEndStateTransition(new EndState(
BatchStatus.FAILED, ExitStatus.FAILED, "end2")));
transitions.add(StateTransition.createEndStateTransition(new EndState(
BatchStatus.COMPLETED, ExitStatus.COMPLETED, "end3")));
flow.setStateTransitions(transitions);
job.setFlow(flow);
job.afterPropertiesSet();
@@ -259,8 +317,10 @@ public class FlowJobTests {
public void testBasicFlow() throws Throwable {
SimpleFlow flow = new SimpleFlow("job");
List<StateTransition> transitions = new ArrayList<StateTransition>();
transitions.add(StateTransition.createStateTransition(new StepState(new StubStep("step")), "end0"));
transitions.add(StateTransition.createEndStateTransition(new EndState(BatchStatus.COMPLETED, "end0")));
transitions.add(StateTransition.createStateTransition(new StepState(
new StubStep("step")), "end0"));
transitions.add(StateTransition.createEndStateTransition(new EndState(
BatchStatus.COMPLETED, ExitStatus.COMPLETED, "end0")));
flow.setStateTransitions(transitions);
job.setFlow(flow);
job.execute(jobExecution);
@@ -275,24 +335,40 @@ public class FlowJobTests {
SimpleFlow flow = new SimpleFlow("job");
JobExecutionDecider decider = new JobExecutionDecider() {
public String decide(JobExecution jobExecution, StepExecution stepExecution) {
public String decide(JobExecution jobExecution,
StepExecution stepExecution) {
assertNotNull(stepExecution);
return "SWITCH";
}
};
List<StateTransition> transitions = new ArrayList<StateTransition>();
transitions.add(StateTransition.createStateTransition(new StepState(new StubStep("step1")), "decision"));
transitions.add(StateTransition.createStateTransition(new DecisionState(decider, "decision"), "step2"));
transitions.add(StateTransition.createStateTransition(new DecisionState(decider, "decision"), "SWITCH", "step3"));
transitions.add(StateTransition.createStateTransition(new StepState(new StubStep("step2")), ExitStatus.COMPLETED.getExitCode(), "end0"));
transitions.add(StateTransition.createStateTransition(new StepState(new StubStep("step2")), ExitStatus.FAILED.getExitCode(), "end1"));
transitions.add(StateTransition.createEndStateTransition(new EndState(BatchStatus.COMPLETED, "end0")));
transitions.add(StateTransition.createEndStateTransition(new EndState(BatchStatus.INCOMPLETE, ExitStatus.FAILED, "end1")));
transitions.add(StateTransition.createStateTransition(new StepState(new StubStep("step3")), ExitStatus.FAILED.getExitCode(), "end2"));
transitions.add(StateTransition.createStateTransition(new StepState(new StubStep("step3")), ExitStatus.COMPLETED.getExitCode(), "end3"));
transitions.add(StateTransition.createEndStateTransition(new EndState(BatchStatus.INCOMPLETE, ExitStatus.FAILED, "end2")));
transitions.add(StateTransition.createEndStateTransition(new EndState(BatchStatus.COMPLETED, "end3")));
transitions.add(StateTransition.createStateTransition(new StepState(
new StubStep("step1")), "decision"));
transitions.add(StateTransition.createStateTransition(
new DecisionState(decider, "decision"), "step2"));
transitions.add(StateTransition.createStateTransition(
new DecisionState(decider, "decision"), "SWITCH", "step3"));
transitions.add(StateTransition.createStateTransition(new StepState(
new StubStep("step2")), ExitStatus.COMPLETED.getExitCode(),
"end0"));
transitions.add(StateTransition
.createStateTransition(new StepState(new StubStep("step2")),
ExitStatus.FAILED.getExitCode(), "end1"));
transitions.add(StateTransition.createEndStateTransition(new EndState(
BatchStatus.COMPLETED, ExitStatus.COMPLETED, "end0")));
transitions.add(StateTransition.createEndStateTransition(new EndState(
BatchStatus.FAILED, ExitStatus.FAILED, "end1")));
transitions.add(StateTransition
.createStateTransition(new StepState(new StubStep("step3")),
ExitStatus.FAILED.getExitCode(), "end2"));
transitions.add(StateTransition.createStateTransition(new StepState(
new StubStep("step3")), ExitStatus.COMPLETED.getExitCode(),
"end3"));
transitions.add(StateTransition.createEndStateTransition(new EndState(
BatchStatus.FAILED, ExitStatus.FAILED, "end2")));
transitions.add(StateTransition.createEndStateTransition(new EndState(
BatchStatus.COMPLETED, ExitStatus.COMPLETED, "end3")));
flow.setStateTransitions(transitions);
job.setFlow(flow);
@@ -311,14 +387,17 @@ public class FlowJobTests {
public void testGetStepExists() throws Exception {
SimpleFlow flow = new SimpleFlow("job");
List<StateTransition> transitions = new ArrayList<StateTransition>();
transitions.add(StateTransition.createStateTransition(new StepState(new StubStep("step1")), "step2"));
transitions.add(StateTransition.createStateTransition(new StepState(new StubStep("step2")), "end0"));
transitions.add(StateTransition.createEndStateTransition(new EndState(BatchStatus.COMPLETED, "end0")));
transitions.add(StateTransition.createStateTransition(new StepState(
new StubStep("step1")), "step2"));
transitions.add(StateTransition.createStateTransition(new StepState(
new StubStep("step2")), "end0"));
transitions.add(StateTransition.createEndStateTransition(new EndState(
BatchStatus.COMPLETED, ExitStatus.COMPLETED, "end0")));
flow.setStateTransitions(transitions);
flow.afterPropertiesSet();
job.setFlow(flow);
job.afterPropertiesSet();
Step step = job.getStep("step2");
assertNotNull(step);
assertEquals("step2", step.getName());
@@ -328,9 +407,12 @@ public class FlowJobTests {
public void testGetStepNotExists() throws Exception {
SimpleFlow flow = new SimpleFlow("job");
List<StateTransition> transitions = new ArrayList<StateTransition>();
transitions.add(StateTransition.createStateTransition(new StepState(new StubStep("step1")), "step2"));
transitions.add(StateTransition.createStateTransition(new StepState(new StubStep("step2")), "end0"));
transitions.add(StateTransition.createEndStateTransition(new EndState(BatchStatus.COMPLETED, "end0")));
transitions.add(StateTransition.createStateTransition(new StepState(
new StubStep("step1")), "step2"));
transitions.add(StateTransition.createStateTransition(new StepState(
new StubStep("step2")), "end0"));
transitions.add(StateTransition.createEndStateTransition(new EndState(
BatchStatus.COMPLETED, ExitStatus.COMPLETED, "end0")));
flow.setStateTransitions(transitions);
flow.afterPropertiesSet();
job.setFlow(flow);
@@ -339,14 +421,17 @@ public class FlowJobTests {
Step step = job.getStep("foo");
assertNull(step);
}
@Test
public void testGetStepNotStepState() throws Exception {
SimpleFlow flow = new SimpleFlow("job");
List<StateTransition> transitions = new ArrayList<StateTransition>();
transitions.add(StateTransition.createStateTransition(new StepState(new StubStep("step1")), "step2"));
transitions.add(StateTransition.createStateTransition(new StepState(new StubStep("step2")), "end0"));
transitions.add(StateTransition.createEndStateTransition(new EndState(BatchStatus.COMPLETED, "end0")));
transitions.add(StateTransition.createStateTransition(new StepState(
new StubStep("step1")), "step2"));
transitions.add(StateTransition.createStateTransition(new StepState(
new StubStep("step2")), "end0"));
transitions.add(StateTransition.createEndStateTransition(new EndState(
BatchStatus.COMPLETED, ExitStatus.COMPLETED, "end0")));
flow.setStateTransitions(transitions);
flow.afterPropertiesSet();
job.setFlow(flow);
@@ -355,7 +440,7 @@ public class FlowJobTests {
Step step = job.getStep("end0");
assertNull(step);
}
/**
* @author Dave Syer
*
@@ -366,7 +451,8 @@ public class FlowJobTests {
super(name);
}
public void execute(StepExecution stepExecution) throws JobInterruptedException {
public void execute(StepExecution stepExecution)
throws JobInterruptedException {
stepExecution.setStatus(BatchStatus.COMPLETED);
stepExecution.setExitStatus(ExitStatus.COMPLETED);
jobRepository.update(stepExecution);
@@ -379,15 +465,15 @@ public class FlowJobTests {
* @param stepName
* @return the StepExecution corresponding to the specified step
*/
private StepExecution getStepExecution(JobExecution jobExecution, String stepName)
{
for(StepExecution stepExecution : jobExecution.getStepExecutions()) {
if(stepExecution.getStepName().equals(stepName)) {
private StepExecution getStepExecution(JobExecution jobExecution,
String stepName) {
for (StepExecution stepExecution : jobExecution.getStepExecutions()) {
if (stepExecution.getStepName().equals(stepName)) {
return stepExecution;
}
}
fail("No stepExecution found with name: [" + stepName + "]");
return null;
}
}

View File

@@ -47,4 +47,7 @@ public class JobFlowExecutorSupport implements FlowExecutor {
public void close(FlowExecution result) {
}
public void updateStepExecutionStatus() {
}
}

View File

@@ -163,7 +163,7 @@ public class SimpleFlowTests {
flow.setStateTransitions(collect(StateTransition.createStateTransition(new StubState("step1") {
@Override
public FlowExecutionStatus handle(FlowExecutor executor) {
return FlowExecutionStatus.INCOMPLETE;
return FlowExecutionStatus.FAILED;
}
}, "step2"), StateTransition.createEndStateTransition(new StubState("step2"))));
flow.afterPropertiesSet();

View File

@@ -20,6 +20,7 @@ import static org.junit.Assert.assertEquals;
import org.junit.Before;
import org.junit.Test;
import org.springframework.batch.core.BatchStatus;
import org.springframework.batch.core.ExitStatus;
import org.springframework.batch.core.JobExecution;
import org.springframework.batch.core.job.flow.FlowExecutor;
import org.springframework.batch.core.job.flow.support.JobFlowExecutorSupport;
@@ -47,7 +48,7 @@ public class EndStateTests {
BatchStatus status = jobExecution.getStatus();
EndState state = new EndState(BatchStatus.UNKNOWN, "end");
EndState state = new EndState(BatchStatus.UNKNOWN, ExitStatus.UNKNOWN, "end");
state.handle(new JobFlowExecutorSupport() {
@Override
public JobExecution getJobExecution() {
@@ -68,7 +69,7 @@ public class EndStateTests {
jobExecution.createStepExecution("foo");
EndState state = new EndState(BatchStatus.UNKNOWN, "end");
EndState state = new EndState(BatchStatus.UNKNOWN, ExitStatus.UNKNOWN, "end");
state.handle(new JobFlowExecutorSupport() {
@Override
public JobExecution getJobExecution() {
@@ -76,7 +77,7 @@ public class EndStateTests {
}
});
assertEquals(BatchStatus.STARTING, jobExecution.getStatus());
assertEquals(BatchStatus.UNKNOWN, jobExecution.getStatus());
}
@@ -87,10 +88,10 @@ public class EndStateTests {
@Test
public void testHandleOngoingAttemptedDowngrade() throws Exception {
jobExecution.setStatus(BatchStatus.INCOMPLETE);
jobExecution.setStatus(BatchStatus.FAILED);
jobExecution.createStepExecution("foo");
EndState state = new EndState(BatchStatus.COMPLETED, "end");
EndState state = new EndState(BatchStatus.COMPLETED, ExitStatus.COMPLETED, "end");
state.handle(new JobFlowExecutorSupport() {
@Override
public JobExecution getJobExecution() {
@@ -98,8 +99,8 @@ public class EndStateTests {
}
});
// Can't downgrade a status - if it failed then it failed
assertEquals(BatchStatus.INCOMPLETE, jobExecution.getStatus());
// An EndState can downgrade a status - if it failed then it can be unfailed
assertEquals(BatchStatus.COMPLETED, jobExecution.getStatus());
}

View File

@@ -36,10 +36,10 @@ public class SimpleFlowExecutionAggregatorTests {
@Test
public void testFailed() throws Exception {
FlowExecution first = new FlowExecution("foo", FlowExecutionStatus.COMPLETED);
FlowExecution second = new FlowExecution("foo", FlowExecutionStatus.INCOMPLETE);
FlowExecution second = new FlowExecution("foo", FlowExecutionStatus.FAILED);
assertTrue("Should be negative", first.compareTo(second)<0);
assertTrue("Should be positive", second.compareTo(first)>0);
assertEquals(FlowExecutionStatus.INCOMPLETE, aggregator.aggregate(Arrays.asList(first, second)));
assertEquals(FlowExecutionStatus.FAILED, aggregator.aggregate(Arrays.asList(first, second)));
}
@Test

View File

@@ -78,7 +78,7 @@ public class RestartIntegrationTests {
int beforePartition = jdbcTemplate.queryForInt("SELECT COUNT(*) from BATCH_STEP_EXECUTION where STEP_NAME like 'step1:partition%'");
JobExecution execution = jobLauncher.run(job, jobParameters);
assertEquals(BatchStatus.INCOMPLETE,execution.getStatus());
assertEquals(BatchStatus.FAILED,execution.getStatus());
assertNotNull(jobLauncher.run(job, jobParameters));
int afterMaster = jdbcTemplate.queryForInt("SELECT COUNT(*) from BATCH_STEP_EXECUTION where STEP_NAME='step1:master'");

View File

@@ -88,7 +88,7 @@ public class PartitionStepTests {
throws Exception {
Set<StepExecution> executions = stepSplitter.split(stepExecution, 2);
for (StepExecution execution : executions) {
execution.setStatus(BatchStatus.INCOMPLETE);
execution.setStatus(BatchStatus.FAILED);
execution.setExitStatus(ExitStatus.FAILED);
}
return executions;
@@ -101,7 +101,7 @@ public class PartitionStepTests {
step.execute(stepExecution);
// one master and two workers
assertEquals(3, stepExecution.getJobExecution().getStepExecutions().size());
assertEquals(BatchStatus.INCOMPLETE, stepExecution.getStatus());
assertEquals(BatchStatus.FAILED, stepExecution.getStatus());
}
}

View File

@@ -45,21 +45,21 @@ public class StepExecutionAggregatorTests {
@Test
public void testAggregateStatusFromFailure() {
result.setStatus(BatchStatus.INCOMPLETE);
result.setStatus(BatchStatus.FAILED);
stepExecution1.setStatus(BatchStatus.COMPLETED);
stepExecution2.setStatus(BatchStatus.COMPLETED);
aggregator.aggregate(result, Arrays.<StepExecution> asList(stepExecution1, stepExecution2));
assertNotNull(result);
assertEquals(BatchStatus.INCOMPLETE, result.getStatus());
assertEquals(BatchStatus.FAILED, result.getStatus());
}
@Test
public void testAggregateStatusIncomplete() {
stepExecution1.setStatus(BatchStatus.COMPLETED);
stepExecution2.setStatus(BatchStatus.INCOMPLETE);
stepExecution2.setStatus(BatchStatus.FAILED);
aggregator.aggregate(result, Arrays.<StepExecution> asList(stepExecution1, stepExecution2));
assertNotNull(result);
assertEquals(BatchStatus.INCOMPLETE, result.getStatus());
assertEquals(BatchStatus.FAILED, result.getStatus());
}
@Test

View File

@@ -184,7 +184,7 @@ public abstract class AbstractStepExecutionDaoTests extends AbstractTransactiona
dao.saveStepExecution(stepExecution);
Integer versionAfterSave = stepExecution.getVersion();
stepExecution.setStatus(BatchStatus.FAILED);
stepExecution.setStatus(BatchStatus.ABANDONED);
stepExecution.setLastUpdated(new Date(System.currentTimeMillis()));
dao.updateStepExecution(stepExecution);
assertEquals(versionAfterSave + 1, stepExecution.getVersion().intValue());
@@ -192,7 +192,7 @@ public abstract class AbstractStepExecutionDaoTests extends AbstractTransactiona
StepExecution retrieved = dao.getStepExecution(jobExecution, stepExecution.getId());
assertEquals(stepExecution, retrieved);
assertEquals(stepExecution.getLastUpdated(), retrieved.getLastUpdated());
assertEquals(BatchStatus.FAILED, retrieved.getStatus());
assertEquals(BatchStatus.ABANDONED, retrieved.getStatus());
}
/**

View File

@@ -113,10 +113,10 @@ public class SimpleJobRepositoryIntegrationTests {
// first execution failed
firstJobExec.setStartTime(new Date(4));
firstStepExec.setStartTime(new Date(5));
firstStepExec.setStatus(BatchStatus.INCOMPLETE);
firstStepExec.setStatus(BatchStatus.FAILED);
firstStepExec.setEndTime(new Date(6));
jobRepository.update(firstStepExec);
firstJobExec.setStatus(BatchStatus.INCOMPLETE);
firstJobExec.setStatus(BatchStatus.FAILED);
firstJobExec.setEndTime(new Date(7));
jobRepository.update(firstJobExec);
@@ -182,7 +182,7 @@ public class SimpleJobRepositoryIntegrationTests {
@Test
public void testGetLastJobExecution() throws Exception {
JobExecution jobExecution = jobRepository.createJobExecution(job.getName(), jobParameters);
jobExecution.setStatus(BatchStatus.INCOMPLETE);
jobExecution.setStatus(BatchStatus.FAILED);
jobExecution.setEndTime(new Date());
jobRepository.update(jobExecution);
Thread.sleep(10);

View File

@@ -204,7 +204,7 @@ public class AbstractStepTests {
tested.setStepExecutionListeners(new StepExecutionListener[] { listener1, listener2 });
tested.execute(execution);
assertEquals(BatchStatus.INCOMPLETE, execution.getStatus());
assertEquals(BatchStatus.FAILED, execution.getStatus());
Throwable expected = execution.getFailureExceptions().get(0);
assertEquals("crash!", expected.getMessage());
@@ -242,7 +242,7 @@ public class AbstractStepTests {
tested.setStepExecutionListeners(new StepExecutionListener[] { listener1, listener2 });
tested.execute(execution);
assertEquals(BatchStatus.INCOMPLETE, execution.getStatus());
assertEquals(BatchStatus.STOPPED, execution.getStatus());
Throwable expected = execution.getFailureExceptions().get(0);
assertEquals("JobExecution interrupted.", expected.getMessage());
@@ -256,7 +256,7 @@ public class AbstractStepTests {
assertEquals("close", events.get(i++));
assertEquals(7, events.size());
assertEquals("INTERRUPTED", execution.getExitStatus().getExitCode());
assertEquals("STOPPED", execution.getExitStatus().getExitCode());
assertTrue("Execution context modifications made by listener should be persisted", repository.saved
.containsKey("afterStep"));

View File

@@ -343,7 +343,7 @@ public class FaultTolerantStepFactoryBeanRetryTests {
StepExecution stepExecution = new StepExecution(step.getName(), jobExecution);
repository.add(stepExecution);
step.execute(stepExecution);
assertEquals(BatchStatus.INCOMPLETE, stepExecution.getStatus());
assertEquals(BatchStatus.FAILED, stepExecution.getStatus());
List<String> expectedOutput = Arrays.asList(StringUtils.commaDelimitedListToStringArray(""));
assertEquals(expectedOutput, written);
@@ -439,7 +439,7 @@ public class FaultTolerantStepFactoryBeanRetryTests {
StepExecution stepExecution = new StepExecution(step.getName(), jobExecution);
repository.add(stepExecution);
step.execute(stepExecution);
assertEquals(BatchStatus.INCOMPLETE, stepExecution.getStatus());
assertEquals(BatchStatus.FAILED, stepExecution.getStatus());
List<String> expectedOutput = Arrays.asList(StringUtils.commaDelimitedListToStringArray(""));
assertEquals(expectedOutput, written);
@@ -488,7 +488,7 @@ public class FaultTolerantStepFactoryBeanRetryTests {
StepExecution stepExecution = new StepExecution(step.getName(), jobExecution);
repository.add(stepExecution);
step.execute(stepExecution);
assertEquals(BatchStatus.INCOMPLETE, stepExecution.getStatus());
assertEquals(BatchStatus.FAILED, stepExecution.getStatus());
// We added a bogus cache so no items are actually skipped
// because they aren't recognised as eligible

View File

@@ -112,7 +112,7 @@ public class FaultTolerantStepFactoryBeanTests {
Step step = (Step) factory.getObject();
step.execute(stepExecution);
assertEquals(BatchStatus.INCOMPLETE, stepExecution.getStatus());
assertEquals(BatchStatus.FAILED, stepExecution.getStatus());
assertEquals(ExitStatus.FAILED.getExitCode(), stepExecution.getExitStatus().getExitCode());
assertTrue(stepExecution.getExitStatus().getExitDescription().contains("Non-skippable exception during read"));
@@ -141,7 +141,7 @@ public class FaultTolerantStepFactoryBeanTests {
Step step = (Step) factory.getObject();
step.execute(stepExecution);
assertEquals(BatchStatus.INCOMPLETE, stepExecution.getStatus());
assertEquals(BatchStatus.FAILED, stepExecution.getStatus());
assertEquals(1, reader.processed.size());
assertEquals(ExitStatus.FAILED.getExitCode(), stepExecution.getExitStatus().getExitCode());
assertTrue(stepExecution.getExitStatus().getExitDescription().contains("non-skippable exception"));
@@ -335,7 +335,7 @@ public class FaultTolerantStepFactoryBeanTests {
Step step = (Step) factory.getObject();
step.execute(stepExecution);
assertEquals(BatchStatus.INCOMPLETE, stepExecution.getStatus());
assertEquals(BatchStatus.FAILED, stepExecution.getStatus());
assertEquals(3, stepExecution.getSkipCount());
assertEquals(2, stepExecution.getReadSkipCount());
@@ -374,7 +374,7 @@ public class FaultTolerantStepFactoryBeanTests {
Step step = (Step) factory.getObject();
step.execute(stepExecution);
assertEquals(BatchStatus.INCOMPLETE, stepExecution.getStatus());
assertEquals(BatchStatus.FAILED, stepExecution.getStatus());
assertEquals("oops", stepExecution.getFailureExceptions().get(0).getCause().getMessage());
// listeners are called only once chunk is about to commit, so
@@ -408,7 +408,7 @@ public class FaultTolerantStepFactoryBeanTests {
Step step = (Step) factory.getObject();
step.execute(stepExecution);
assertEquals(BatchStatus.INCOMPLETE, stepExecution.getStatus());
assertEquals(BatchStatus.FAILED, stepExecution.getStatus());
assertEquals("oops", stepExecution.getFailureExceptions().get(0).getCause().getMessage());
assertEquals(1, stepExecution.getSkipCount());
assertEquals(0, stepExecution.getReadSkipCount());
@@ -521,7 +521,7 @@ public class FaultTolerantStepFactoryBeanTests {
Step step = (Step) factory.getObject();
step.execute(stepExecution);
assertEquals(BatchStatus.INCOMPLETE, stepExecution.getStatus());
assertEquals(BatchStatus.FAILED, stepExecution.getStatus());
assertEquals("bad skip count", 3, stepExecution.getSkipCount());
assertEquals("bad read skip count", 2, stepExecution.getReadSkipCount());
assertEquals("bad write skip count", 1, stepExecution.getWriteSkipCount());

View File

@@ -159,7 +159,7 @@ public class SimpleStepFactoryBeanTests {
job.execute(jobExecution);
assertEquals("Error!", jobExecution.getAllFailureExceptions().get(0).getMessage());
assertEquals(BatchStatus.INCOMPLETE, jobExecution.getStatus());
assertEquals(BatchStatus.FAILED, jobExecution.getStatus());
assertEquals(0, written.size());
// provider should be at second item
assertEquals("bar", reader.read());
@@ -182,7 +182,7 @@ public class SimpleStepFactoryBeanTests {
job.execute(jobExecution);
assertEquals("Foo", jobExecution.getAllFailureExceptions().get(0).getMessage());
assertEquals(BatchStatus.INCOMPLETE, jobExecution.getStatus());
assertEquals(BatchStatus.FAILED, jobExecution.getStatus());
}
@Test

View File

@@ -7,8 +7,8 @@ import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertFalse;
import static org.junit.Assert.assertTrue;
import static org.springframework.batch.core.BatchStatus.COMPLETED;
import static org.springframework.batch.core.BatchStatus.INCOMPLETE;
import static org.springframework.batch.core.BatchStatus.FAILED;
import static org.springframework.batch.core.BatchStatus.STOPPED;
import static org.springframework.batch.core.BatchStatus.UNKNOWN;
import org.junit.Before;
@@ -73,14 +73,14 @@ public class TaskletStepExceptionTests {
public void testApplicationException() throws Exception {
taskletStep.execute(stepExecution);
assertEquals(INCOMPLETE, stepExecution.getStatus());
assertEquals(FAILED, stepExecution.getStatus());
}
@Test
public void testInterrupted() throws Exception {
taskletStep.setStepExecutionListeners(new StepExecutionListener[] { new InterruptionListener() });
taskletStep.execute(stepExecution);
assertEquals(INCOMPLETE, stepExecution.getStatus());
assertEquals(STOPPED, stepExecution.getStatus());
}
@Test
@@ -94,7 +94,7 @@ public class TaskletStepExceptionTests {
} });
taskletStep.execute(stepExecution);
assertEquals(INCOMPLETE, stepExecution.getStatus());
assertEquals(FAILED, stepExecution.getStatus());
assertTrue(stepExecution.getFailureExceptions().contains(exception));
assertEquals(2, jobRepository.getUpdateCount());
}
@@ -110,7 +110,7 @@ public class TaskletStepExceptionTests {
}
} });
taskletStep.execute(stepExecution);
assertEquals(INCOMPLETE, stepExecution.getStatus());
assertEquals(FAILED, stepExecution.getStatus());
assertTrue(stepExecution.getFailureExceptions().contains(exception));
assertEquals(2, jobRepository.getUpdateCount());
}
@@ -152,7 +152,7 @@ public class TaskletStepExceptionTests {
}
} });
taskletStep.execute(stepExecution);
assertEquals(INCOMPLETE, stepExecution.getStatus());
assertEquals(FAILED, stepExecution.getStatus());
assertTrue(stepExecution.getFailureExceptions().contains(taskletException));
assertFalse(stepExecution.getFailureExceptions().contains(exception));
assertEquals(2, jobRepository.getUpdateCount());
@@ -170,7 +170,7 @@ public class TaskletStepExceptionTests {
} });
taskletStep.execute(stepExecution);
assertEquals(INCOMPLETE, stepExecution.getStatus());
assertEquals(FAILED, stepExecution.getStatus());
assertTrue(stepExecution.getFailureExceptions().contains(taskletException));
assertTrue(stepExecution.getFailureExceptions().contains(exception));
assertEquals(2, jobRepository.getUpdateCount());

View File

@@ -114,7 +114,7 @@ public class StepExecutorInterruptionTests extends TestCase {
assertTrue("Timed out waiting for step to be interrupted.", count < 1000);
assertFalse(processingThread.isAlive());
assertEquals(BatchStatus.INCOMPLETE, stepExecution.getStatus());
assertEquals(BatchStatus.STOPPED, stepExecution.getStatus());
}
@@ -170,7 +170,7 @@ public class StepExecutorInterruptionTests extends TestCase {
assertTrue("Timed out waiting for step to be interrupted.", count < 1000);
assertFalse(processingThread.isAlive());
assertEquals(BatchStatus.INCOMPLETE, stepExecution.getStatus());
assertEquals(BatchStatus.STOPPED, stepExecution.getStatus());
}
@@ -201,7 +201,7 @@ public class StepExecutorInterruptionTests extends TestCase {
step.execute(stepExecution);
assertEquals("Planned!", stepExecution.getFailureExceptions().get(0).getMessage());
assertEquals(BatchStatus.INCOMPLETE, stepExecution.getStatus());
assertEquals(BatchStatus.FAILED, stepExecution.getStatus());
}
/**

View File

@@ -541,7 +541,7 @@ public class TaskletStepTests {
stepExecution.setExecutionContext(foobarEc);
step.execute(stepExecution);
assertEquals(BatchStatus.INCOMPLETE, stepExecution.getStatus());
assertEquals(BatchStatus.STOPPED, stepExecution.getStatus());
String msg = stepExecution.getExitStatus().getExitDescription();
assertTrue("Message does not contain 'JobInterruptedException': " + msg, contains(msg,
"JobInterruptedException"));
@@ -565,7 +565,7 @@ public class TaskletStepTests {
// step.setLastExecution(stepExecution);
step.execute(stepExecution);
assertEquals(BatchStatus.INCOMPLETE, stepExecution.getStatus());
assertEquals(BatchStatus.FAILED, stepExecution.getStatus());
// The original rollback was caused by this one:
assertEquals("Foo", stepExecution.getFailureExceptions().get(0).getMessage());
}
@@ -588,7 +588,7 @@ public class TaskletStepTests {
// step.setLastExecution(stepExecution);
step.execute(stepExecution);
assertEquals(BatchStatus.INCOMPLETE, stepExecution.getStatus());
assertEquals(BatchStatus.FAILED, stepExecution.getStatus());
// The original rollback was caused by this one:
assertEquals("Foo", stepExecution.getFailureExceptions().get(0).getMessage());
}
@@ -723,7 +723,7 @@ public class TaskletStepTests {
StepExecution stepExecution = new StepExecution(step.getName(), new JobExecution(jobInstance));
step.execute(stepExecution);
assertEquals(BatchStatus.INCOMPLETE, stepExecution.getStatus());
assertEquals(BatchStatus.FAILED, stepExecution.getStatus());
Throwable expected = stepExecution.getFailureExceptions().get(0);
assertEquals("CRASH!", expected.getMessage());
assertFalse(stepExecution.getExecutionContext().isEmpty());

View File

@@ -9,21 +9,24 @@
<beans:import resource="common-context.xml" />
<job id="job">
<split id="split1">
<next on="FAILED" to="s3"/>
<fail on="COMPLETED" />
<flow>
<step id="s1" ref="step1"/>
</flow>
<flow>
<step id="fail" ref="failingStep">
<fail on="FAILED"/>
<end on="*"/>
</step>
<step id="fail" ref="failingStep"/>
</flow>
</split>
<step id="s3" ref="step3"/>
</job>
</beans:beans>

View File

@@ -0,0 +1,16 @@
<?xml version="1.0" encoding="UTF-8"?>
<beans:beans xmlns="http://www.springframework.org/schema/batch" xmlns:beans="http://www.springframework.org/schema/beans"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://www.springframework.org/schema/batch http://www.springframework.org/schema/batch/spring-batch-2.0.xsd
http://www.springframework.org/schema/beans http://www.springframework.org/schema/beans/spring-beans-2.5.xsd">
<beans:import resource="common-context.xml" />
<job id="job">
<step id="s1" ref="step1">
<stop on="*" restart="s2"/>
</step>
<step id="s2" ref="step2"/>
</job>
</beans:beans>

View File

@@ -6,8 +6,6 @@
<beans:import resource="common-context.xml" />
<beans:bean id="decider" class="org.springframework.batch.core.configuration.xml.StopJobParserTests$TestDecider"/>
<job id="job">
<step id="fail" ref="failingStep">
<stop on="FAILED" restart="s2"/>

View File

@@ -8,12 +8,14 @@
<job id="job">
<step id="s1" ref="taskletStep">
<!-- On normal step status: stop the job, but restart with the same step. -->
<stop on="COMPLETED" restart="s1"/>
<end on="*"/>
</step>
</job>
<beans:bean id="taskletStep" parent="step1">
<beans:bean id="taskletStep" parent="step1">
<!-- On restart this step will be re-executed: effect is infinitely re-runnable job -->
<beans:property name="allowStartIfComplete" value="true"/>
</beans:bean>

View File

@@ -99,11 +99,11 @@ public class ChunkMessageChannelItemWriter<T> extends StepExecutionListenerSuppo
}
catch (RuntimeException e) {
logger.debug("Detected failure waiting for results in step listener.", e);
stepExecution.setStatus(BatchStatus.INCOMPLETE);
stepExecution.setStatus(BatchStatus.FAILED);
return ExitStatus.FAILED.addExitDescription(e.getClass().getName() + ": " + e.getMessage());
}
if (timedOut) {
stepExecution.setStatus(BatchStatus.INCOMPLETE);
stepExecution.setStatus(BatchStatus.FAILED);
throw new ItemStreamException("Timed out waiting for back log at end of step");
}
return ExitStatus.COMPLETED.addExitDescription("Waited for " + expecting + " results.");

View File

@@ -27,8 +27,8 @@ import org.springframework.batch.core.JobExecution;
* (generally a handler cannot determine if the whole job execution is complete,
* so this is just information about the step).<br/>
*
* If the incoming status is {@link BatchStatus#INCOMPLETE},
* {@link BatchStatus#FAILED} or {@link BatchStatus#STOPPING} the request
* If the incoming status is {@link BatchStatus#FAILED},
* {@link BatchStatus#ABANDONED} or {@link BatchStatus#STOPPING} the request
* should be ignored by handlers (passed on without modification).
*
* @author Dave Syer

View File

@@ -136,7 +136,7 @@ public class StepExecutionMessageHandler {
* @return
*/
private boolean isComplete(JobExecutionRequest request) {
return request.getStatus() == BatchStatus.INCOMPLETE || request.getStatus() == BatchStatus.FAILED
return request.getStatus() == BatchStatus.FAILED || request.getStatus() == BatchStatus.ABANDONED
|| request.getStatus() == BatchStatus.STOPPING;
}
@@ -146,7 +146,7 @@ public class StepExecutionMessageHandler {
*/
private void handleFailure(JobExecutionRequest request, Throwable e) {
request.registerThrowable(e);
request.setStatus(BatchStatus.INCOMPLETE);
request.setStatus(BatchStatus.FAILED);
}
/*

View File

@@ -170,7 +170,7 @@ public class ChunkMessageItemWriterIntegrationTests {
// And make the back log real
requests.send(getSimpleMessage("foo", 4321L));
step.execute(stepExecution);
assertEquals(BatchStatus.INCOMPLETE, stepExecution.getStatus());
assertEquals(BatchStatus.FAILED, stepExecution.getStatus());
assertEquals(ExitStatus.FAILED.getExitCode(), stepExecution.getExitStatus().getExitCode());
String message = stepExecution.getExitStatus().getExitDescription();
assertTrue("Message does not contain 'wrong job': " + message, message.contains("wrong job"));
@@ -207,7 +207,7 @@ public class ChunkMessageItemWriterIntegrationTests {
StepExecution stepExecution = getStepExecution(step);
step.execute(stepExecution);
assertEquals(BatchStatus.INCOMPLETE, stepExecution.getStatus());
assertEquals(BatchStatus.FAILED, stepExecution.getStatus());
assertEquals(ExitStatus.FAILED.getExitCode(), stepExecution.getExitStatus().getExitCode());
String message = stepExecution.getExitStatus().getExitDescription();
assertTrue("Message does not contain 'bad': " + message, message.contains("bad"));
@@ -219,7 +219,7 @@ public class ChunkMessageItemWriterIntegrationTests {
assertTrue(1 <= TestItemWriter.count);
assertTrue(6 >= TestItemWriter.count);
// But it should fail the step in any case
assertEquals(BatchStatus.INCOMPLETE, stepExecution.getStatus());
assertEquals(BatchStatus.FAILED, stepExecution.getStatus());
}
@@ -242,7 +242,7 @@ public class ChunkMessageItemWriterIntegrationTests {
* loop would be bad, so the best we can do is fail as fast as possible.
*/
step.execute(stepExecution);
assertEquals(BatchStatus.INCOMPLETE, stepExecution.getStatus());
assertEquals(BatchStatus.FAILED, stepExecution.getStatus());
assertEquals(ExitStatus.FAILED.getExitCode(), stepExecution.getExitStatus().getExitCode());
String message = stepExecution.getExitStatus().getExitDescription();
assertTrue("Message did not contain 'timed out': " + message, message.toLowerCase().contains("timed out"));
@@ -276,7 +276,7 @@ public class ChunkMessageItemWriterIntegrationTests {
assertTrue(1 <= TestItemWriter.count);
assertTrue(6 >= TestItemWriter.count);
assertEquals(BatchStatus.INCOMPLETE, stepExecution.getStatus());
assertEquals(BatchStatus.FAILED, stepExecution.getStatus());
assertEquals(ExitStatus.FAILED.getExitCode(), stepExecution.getExitStatus().getExitCode());
String exitDescription = stepExecution.getExitStatus().getExitDescription();

View File

@@ -105,7 +105,7 @@ public class MessageOrientedStepTests {
step.setPollingInterval(100);
StepExecution stepExecution = jobExecution.createStepExecution(step.getName());
step.execute(stepExecution);
assertEquals(BatchStatus.INCOMPLETE, stepExecution.getStatus());
assertEquals(BatchStatus.FAILED, stepExecution.getStatus());
assertEquals(ExitStatus.FAILED.getExitCode(), stepExecution.getExitStatus().getExitCode());
String message = stepExecution.getExitStatus().getExitDescription();
assertTrue("Wrong message: " + message, message.contains("StepExecutionTimeoutException"));
@@ -134,7 +134,7 @@ public class MessageOrientedStepTests {
});
StepExecution stepExecution = jobExecution.createStepExecution(step.getName());
step.execute(stepExecution);
assertEquals(BatchStatus.INCOMPLETE, stepExecution.getStatus());
assertEquals(BatchStatus.FAILED, stepExecution.getStatus());
assertEquals(ExitStatus.FAILED.getExitCode(), stepExecution.getExitStatus().getExitCode());
String message = stepExecution.getExitStatus().getExitDescription();
assertTrue("Wrong message: " + message, message.contains("Planned failure"));

View File

@@ -123,7 +123,7 @@ public class StepExecutionMessageHandlerTests {
JobRepositorySupport jobRepository = new JobRepositorySupport();
StepExecutionMessageHandler handler = createHandler(jobRepository);
JobExecution jobExecution = jobRepository.createJobExecution("job", new JobParameters());
jobExecution.setStatus(BatchStatus.INCOMPLETE);
jobExecution.setStatus(BatchStatus.FAILED);
JobExecutionRequest message = handler.handle(new JobExecutionRequest(jobExecution));
assertEquals(0, message.getJobExecution().getStepExecutions().size());
}
@@ -134,7 +134,7 @@ public class StepExecutionMessageHandlerTests {
@Override
public StepExecution getLastStepExecution(JobInstance jobInstance, String stepName) {
StepExecution stepExecution = new StepExecution(stepName, new JobExecution(jobInstance));
stepExecution.setStatus(BatchStatus.INCOMPLETE);
stepExecution.setStatus(BatchStatus.FAILED);
stepExecution.setExecutionContext(new ExecutionContext() {
{
put("foo", "bar");
@@ -205,7 +205,7 @@ public class StepExecutionMessageHandlerTests {
assertNotNull(message);
assertEquals(1, jobExecution.getStepExecutions().size());
JobExecutionRequest payload = message;
assertEquals(BatchStatus.INCOMPLETE, payload.getStatus());
assertEquals(BatchStatus.FAILED, payload.getStatus());
assertTrue(payload.hasErrors());
Throwable error = payload.getLastThrowable();
assertTrue(error instanceof StartLimitExceededException);

View File

@@ -27,6 +27,7 @@ public class DummyItemWriter implements ItemWriter<Object> {
public void write(List<? extends Object> item) throws Exception {
// NO-OP
Thread.sleep(500);
}
}

View File

@@ -69,7 +69,7 @@ public class DatabaseShutdownFunctionalTests extends AbstractBatchLauncherTests
}
assertFalse("Timed out waiting for job to end.", jobExecution.isRunning());
assertEquals(BatchStatus.INCOMPLETE, jobExecution.getStatus());
assertEquals(BatchStatus.STOPPED, jobExecution.getStatus());
}

View File

@@ -64,7 +64,7 @@ public class GracefulShutdownFunctionalTests extends AbstractBatchLauncherTests
}
assertFalse("Timed out waiting for job to end.", jobExecution.isRunning());
assertEquals(BatchStatus.INCOMPLETE, jobExecution.getStatus());
assertEquals(BatchStatus.STOPPED, jobExecution.getStatus());
}

View File

@@ -30,7 +30,7 @@ public class JobOperatorFunctionalTests {
private static final Log logger = LogFactory.getLog(JobOperatorFunctionalTests.class);
@Autowired
private JobOperator tested;
private JobOperator operator;
@Autowired
private Job job;
@@ -49,19 +49,19 @@ public class JobOperatorFunctionalTests {
String params = new JobParametersBuilder().addLong("jobOperatorTestParam", 7L).toJobParameters().toString();
long executionId = tested.start(job.getName(), params);
assertEquals(params, tested.getParameters(executionId));
long executionId = operator.start(job.getName(), params);
assertEquals(params, operator.getParameters(executionId));
stopAndCheckStatus(executionId);
long resumedExecutionId = tested.restart(executionId);
assertEquals(params, tested.getParameters(resumedExecutionId));
long resumedExecutionId = operator.restart(executionId);
assertEquals(params, operator.getParameters(resumedExecutionId));
stopAndCheckStatus(resumedExecutionId);
List<Long> instances = tested.getJobInstances(job.getName(), 0, 1);
List<Long> instances = operator.getJobInstances(job.getName(), 0, 1);
assertEquals(1, instances.size());
long instanceId = instances.get(0);
List<Long> executions = tested.getExecutions(instanceId);
List<Long> executions = operator.getExecutions(instanceId);
assertEquals(2, executions.size());
// latest execution is the first in the returned list
assertEquals(resumedExecutionId, executions.get(0).longValue());
@@ -77,47 +77,49 @@ public class JobOperatorFunctionalTests {
// wait to the job to get up and running
Thread.sleep(1000);
assertTrue(tested.getRunningExecutions(job.getName()).contains(executionId));
assertTrue(tested.getSummary(executionId).contains(BatchStatus.STARTED.toString()));
Set<Long> runningExecutions = operator.getRunningExecutions(job.getName());
assertTrue("Wrong executions: "+runningExecutions+" expected: "+executionId, runningExecutions.contains(executionId));
assertTrue("Wrong summary: "+operator.getSummary(executionId), operator.getSummary(executionId).contains(BatchStatus.STARTED.toString()));
tested.stop(executionId);
operator.stop(executionId);
int count = 0;
while (tested.getRunningExecutions(job.getName()).contains(executionId) && count <= 10) {
while (operator.getRunningExecutions(job.getName()).contains(executionId) && count <= 10) {
logger.info("Checking for running JobExecution: count=" + count);
Thread.sleep(100);
count++;
}
assertFalse(tested.getRunningExecutions(job.getName()).contains(executionId));
assertTrue(tested.getSummary(executionId).contains(BatchStatus.INCOMPLETE.toString()));
runningExecutions = operator.getRunningExecutions(job.getName());
assertFalse("Wrong executions: "+runningExecutions+" expected: "+executionId, runningExecutions.contains(executionId));
assertTrue("Wrong summary: "+operator.getSummary(executionId), operator.getSummary(executionId).contains(BatchStatus.STOPPED.toString()));
// there is just a single step in the test job
Map<Long, String> summaries = tested.getStepExecutionSummaries(executionId);
assertEquals(1, summaries.size());
assertTrue(summaries.values().toString().contains(BatchStatus.INCOMPLETE.toString()));
Map<Long, String> summaries = operator.getStepExecutionSummaries(executionId);
System.err.println(summaries);
assertTrue(summaries.values().toString().contains(BatchStatus.STOPPED.toString()));
}
@Test
public void testMultipleSimultaneousInstances() throws Exception {
String jobName = job.getName();
Set<String> names = tested.getJobNames();
Set<String> names = operator.getJobNames();
assertEquals(1, names.size());
assertTrue(names.contains(jobName));
long exec1 = tested.startNextInstance(jobName);
long exec2 = tested.startNextInstance(jobName);
long exec1 = operator.startNextInstance(jobName);
long exec2 = operator.startNextInstance(jobName);
assertTrue(exec1 != exec2);
assertTrue(tested.getParameters(exec1) != tested.getParameters(exec2));
assertTrue(operator.getParameters(exec1) != operator.getParameters(exec2));
Set<Long> executions = tested.getRunningExecutions(jobName);
Set<Long> executions = operator.getRunningExecutions(jobName);
assertTrue(executions.contains(exec1));
assertTrue(executions.contains(exec2));
tested.stop(exec1);
tested.stop(exec2);
operator.stop(exec1);
operator.stop(exec2);
}

View File

@@ -70,7 +70,7 @@ public class RestartFunctionalTests extends AbstractBatchLauncherTests {
int before = simpleJdbcTemplate.queryForInt("SELECT COUNT(*) FROM TRADE");
JobExecution jobExecution = runJobForRestartTest();
assertEquals(BatchStatus.INCOMPLETE, jobExecution.getStatus());
assertEquals(BatchStatus.FAILED, jobExecution.getStatus());
Throwable expected = jobExecution.getAllFailureExceptions().get(0);
assertTrue("Not planned exception: " + expected.getMessage(), expected.getMessage().toLowerCase().indexOf(

View File

@@ -166,7 +166,7 @@ public class JdbcJobRepositoryTests {
cacheJobIds(execution);
execution.setEndTime(new Timestamp(System.currentTimeMillis()));
repository.update(execution);
execution.setStatus(BatchStatus.INCOMPLETE);
execution.setStatus(BatchStatus.FAILED);
int before = simpleJdbcTemplate.queryForInt("SELECT COUNT(*) FROM BATCH_JOB_INSTANCE");
assertEquals(1, before);