diff --git a/fineract-provider/src/main/java/org/apache/fineract/infrastructure/jobs/domain/JobExecutionRepository.java b/fineract-provider/src/main/java/org/apache/fineract/infrastructure/jobs/domain/JobExecutionRepository.java index bbfff4ca7ae..0f811167b87 100644 --- a/fineract-provider/src/main/java/org/apache/fineract/infrastructure/jobs/domain/JobExecutionRepository.java +++ b/fineract-provider/src/main/java/org/apache/fineract/infrastructure/jobs/domain/JobExecutionRepository.java @@ -121,7 +121,7 @@ HAVING COUNT(BJI.JOB_INSTANCE_ID) <= :threshold List.of(COMPLETED.name(), FAILED.name(), UNKNOWN.name()), "threshold", threshold), Long.class); } - public Long getNotCompletedPartitionsCount(Long jobExecutionId, String partitionerStepName) { + public Long getNotCompletedPartitionsCount(Long jobExecutionId, String partitionerStepName, String exitCode) { return namedParameterJdbcTemplate.queryForObject(""" SELECT COUNT(bse.STEP_EXECUTION_ID) FROM BATCH_STEP_EXECUTION BSE @@ -131,7 +131,11 @@ SELECT COUNT(bse.STEP_EXECUTION_ID) BSE.STEP_NAME <> :stepName AND BSE.status <> :status - """, Map.of("jobExecutionId", jobExecutionId, "stepName", partitionerStepName, "status", COMPLETED.name()), Long.class); + AND + BSE.exit_code <> :exitCode + """, + Map.of("jobExecutionId", jobExecutionId, "stepName", partitionerStepName, "status", COMPLETED.name(), "exitCode", exitCode), + Long.class); } public void updateJobStatusToFailed(Long stuckJobId, String partitionerStepName) { diff --git a/fineract-provider/src/main/java/org/apache/fineract/infrastructure/jobs/service/StuckJobExecutorServiceImpl.java b/fineract-provider/src/main/java/org/apache/fineract/infrastructure/jobs/service/StuckJobExecutorServiceImpl.java index 5f46c88f763..68da4f082da 100644 --- a/fineract-provider/src/main/java/org/apache/fineract/infrastructure/jobs/service/StuckJobExecutorServiceImpl.java +++ b/fineract-provider/src/main/java/org/apache/fineract/infrastructure/jobs/service/StuckJobExecutorServiceImpl.java @@ -20,12 +20,20 @@ import static org.springframework.transaction.TransactionDefinition.PROPAGATION_REQUIRES_NEW; +import java.time.LocalDateTime; import java.util.List; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; +import org.apache.fineract.infrastructure.core.service.DateUtils; import org.apache.fineract.infrastructure.jobs.data.partitionedjobs.PartitionedJob; import org.apache.fineract.infrastructure.jobs.domain.JobExecutionRepository; +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.explore.JobExplorer; import org.springframework.batch.core.launch.JobOperator; +import org.springframework.batch.core.repository.JobRepository; import org.springframework.stereotype.Service; import org.springframework.transaction.TransactionStatus; import org.springframework.transaction.support.TransactionCallbackWithoutResult; @@ -39,14 +47,30 @@ public class StuckJobExecutorServiceImpl implements StuckJobExecutorService { private final JobExecutionRepository jobExecutionRepository; private final TransactionTemplate transactionTemplate; private final JobOperator jobOperator; + private final JobExplorer jobExplorer; + private final JobRepository jobRepository; + + private static final String EXIT_CODE_VALUE = "FAILED_AFTER_MANAGER_RESTARTED"; @Override public void resumeStuckJob(String jobName) { - List stuckJobIds = getStuckJobIds(jobName); - if (isPartitionedJob(jobName) && areThereStuckJobs(jobName)) { - restartPartitionedJobs(jobName, stuckJobIds); - } else { - restartTaskletJobs(stuckJobIds); + if (areThereStuckJobs(jobName)) { + // Exit Status (Tag) to know If the Job and Steps were marked as Failed + final ExitStatus cleanupTaskExitStatus = new ExitStatus(EXIT_CODE_VALUE); + final List stuckJobIds = getStuckJobIds(jobName); + + // Mark as Failed the Steps Executions (If exists) after Restart + stuckJobIds.forEach(stuckJobId -> verifyBatchStepsExecutions(stuckJobId, cleanupTaskExitStatus)); + + // Wait the other Job Partitions (If exists) to be completed + if (isPartitionedJob(jobName)) { + restartPartitionedJobs(jobName, stuckJobIds); + } else { + restartTaskletJobs(stuckJobIds); + } + + // Mark as Failed the Job Executions (If exists) after Restart + stuckJobIds.forEach(stuckJobId -> verifyBatchJob(stuckJobId, cleanupTaskExitStatus)); } } @@ -108,7 +132,67 @@ private void waitUntilAllPartitionsFinished(Long stuckJobId, String partitionerS } private boolean areAllPartitionsCompleted(Long stuckJobId, String partitionerStepName) { - Long notCompletedPartitions = jobExecutionRepository.getNotCompletedPartitionsCount(stuckJobId, partitionerStepName); + Long notCompletedPartitions = jobExecutionRepository.getNotCompletedPartitionsCount(stuckJobId, partitionerStepName, + EXIT_CODE_VALUE); return notCompletedPartitions == 0L; } + + private void verifyBatchJob(long executionId, ExitStatus exitStatus) { + JobExecution jobExecution = jobExplorer.getJobExecution(executionId); + if (isRunning(jobExecution)) { + log.info("Job Execution {} will be mark as Failed", executionId); + final LocalDateTime now = DateUtils.getLocalDateTimeOfSystem(); + + // Mark the Job Execution as Failed + jobExecution.setStatus(BatchStatus.FAILED); + jobExecution.setExitStatus(exitStatus); + jobExecution.setEndTime(now); + jobRepository.update(jobExecution); + } + } + + private void verifyBatchStepsExecutions(long executionId, ExitStatus exitStatus) { + final JobExecution jobExecution = jobExplorer.getJobExecution(executionId); + if (isRunning(jobExecution)) { + log.info("Step Execution {} will be mark as Failed", executionId); + final LocalDateTime now = DateUtils.getLocalDateTimeOfSystem(); + + // Mark the Step Executions as Failed + for (StepExecution stepExecution : jobExecution.getStepExecutions()) { + if (isRunning(stepExecution)) { + stepExecution.setStatus(BatchStatus.FAILED); + stepExecution.setExitStatus(exitStatus); + stepExecution.setEndTime(now); + jobRepository.update(stepExecution); + } + } + } + } + + private boolean isRunning(JobExecution jobExecution) { + switch (jobExecution.getStatus()) { + case STARTED: + case STARTING: + case STOPPING: + case UNKNOWN: + return true; + + default: + return jobExecution.getEndTime() == null; + } + } + + private boolean isRunning(StepExecution stepExecution) { + switch (stepExecution.getStatus()) { + case STARTED: + case STARTING: + case STOPPING: + case UNKNOWN: + return true; + + default: + return stepExecution.getEndTime() == null; + } + } + }