diff --git a/core-process-orchestrator/src/main/java/com/netcracker/core/scheduler/helpers/TaskExecutorService.java b/core-process-orchestrator/src/main/java/com/netcracker/core/scheduler/helpers/TaskExecutorService.java index 4854d1d2a..b58771e2f 100644 --- a/core-process-orchestrator/src/main/java/com/netcracker/core/scheduler/helpers/TaskExecutorService.java +++ b/core-process-orchestrator/src/main/java/com/netcracker/core/scheduler/helpers/TaskExecutorService.java @@ -112,29 +112,41 @@ public void execute(@NotNull Runnable command) { final Future future = delegate.submit(wrapper); tasks.put(fk, future); if (timeout != 0L) { - tasks.put(fk2, delegate.submit(() -> { - try { - future.get(timeout, TimeUnit.SECONDS); - } catch (InterruptedException | ExecutionException | TimeoutException e) { - log.info("Task {} execution was interrupted by timeout", taskId); - future.cancel(true); - ProcessInstanceImpl processInstance = ProcessOrchestrator.getInstance().getProcessInstance(taskInstance.getProcessID()); - processInstance.setState(TaskState.FAILED); - processInstance.save(); - TaskInstanceImpl task = ProcessOrchestrator.getInstance().getTaskInstanceRepository().getTaskInstance(taskId); - task.setState(TaskState.FAILED); - task.save(); - } - - return Boolean.TRUE; - })); + tasks.put(fk2, delegate.submit(() -> watchTask(future, timeout, taskId, taskInstance.getProcessID()))); } } } + // Watchdog for tasks with a sync timeout. The completion callback in execute() + // cancels the watchdog with cancel(true) once the task finishes, so an + // interrupt means "stop watching", not "the task timed out" — only a real + // timeout or a failed worker future may mark the task and process FAILED. + // Package-private so tests can drive it directly. + Boolean watchTask(Future future, long timeout, String taskId, String processId) { + try { + future.get(timeout, TimeUnit.SECONDS); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } catch (ExecutionException | TimeoutException e) { + log.info("Task {} execution was interrupted by timeout", taskId); + future.cancel(true); + ProcessInstanceImpl processInstance = ProcessOrchestrator.getInstance().getProcessInstance(processId); + processInstance.setState(TaskState.FAILED); + processInstance.saveResolvingConflict(pi -> pi.setState(TaskState.FAILED)); + TaskInstanceImpl task = ProcessOrchestrator.getInstance().getTaskInstanceRepository().getTaskInstance(taskId); + task.setState(TaskState.FAILED); + task.saveResolvingConflict(t -> t.setState(TaskState.FAILED)); + } + return Boolean.TRUE; + } + + // Reflection into db-scheduler's lambda internals is the only way to reach the + // Execution from the submitted Runnable; on any failure getId degrades to null + // and the task runs without the wrapper bookkeeping. + @SuppressWarnings("java:S3011") @Nullable private Execution getId(Runnable task) { try { @@ -159,11 +171,11 @@ private Execution getId(Runnable task) { } } - public Future terminate(String key) { + public Future terminate(String key) { List> fs = tasks .entrySet() .stream() - .filter(e -> e.getKey().equals(key)) + .filter(e -> e.getKey().getTaskId().equals(key)) .map(f -> woreService.submit(new TerminateRunnable(f.getValue(), key))) .toList(); return woreService.submit(() -> { @@ -171,11 +183,15 @@ public Future terminate(String key) { fs.forEach(f -> { try { f.get(); - } catch (InterruptedException | ExecutionException e) { - throw new RuntimeException(e); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new IllegalStateException(e); + } catch (ExecutionException e) { + throw new IllegalStateException(e); } }); - } + }, + null ); } } diff --git a/core-process-orchestrator/src/main/java/com/netcracker/core/scheduler/po/DataContext.java b/core-process-orchestrator/src/main/java/com/netcracker/core/scheduler/po/DataContext.java index 2bcdb025e..f8b76a1be 100644 --- a/core-process-orchestrator/src/main/java/com/netcracker/core/scheduler/po/DataContext.java +++ b/core-process-orchestrator/src/main/java/com/netcracker/core/scheduler/po/DataContext.java @@ -6,6 +6,7 @@ import com.fasterxml.jackson.databind.annotation.JsonDeserialize; import com.fasterxml.jackson.databind.annotation.JsonSerialize; import com.netcracker.core.scheduler.po.repository.ContextRepository; +import com.netcracker.core.scheduler.po.repository.VersionMismatchException; import com.netcracker.core.scheduler.po.serializers.DataContextDeserializer; import com.netcracker.core.scheduler.po.serializers.DataContextSerializer; import lombok.Getter; @@ -17,6 +18,9 @@ @Getter @JsonSerialize(using = DataContextSerializer.class) @JsonDeserialize(using = DataContextDeserializer.class) +// Content equality inherited from HashMap is intentional: id, version, and the +// dirty flag are persistence bookkeeping, not part of the context's identity. +@SuppressWarnings("java:S2160") public class DataContext extends HashMap { @Setter @@ -71,8 +75,31 @@ public void save() { repository.putContext(this); } + /** + * Applies the mutation and saves. On a version conflict the fresh context is + * reloaded, the same mutation is re-applied to it, and the fresh copy is + * saved — so a fixed-value update (task state description, start/end time) + * survives a concurrent writer instead of aborting the whole execution. + * The mutation must be idempotent. + */ public void apply(Consumer function) { function.accept(this); - save(); + try { + save(); + } catch (VersionMismatchException e) { + DataContext fresh = repository.getContext(getId()); + fresh.setRepository(repository); + function.accept(fresh); + fresh.save(); + // Sync this instance with the persisted state: contents, version, and + // a clean dirty flag — otherwise the caller keeps a stale copy and the + // very next save() fails again despite apply() having succeeded. + // super-level access bypasses the dirty-marking overrides. + super.clear(); + for (Entry entry : fresh.entrySet()) { + super.put(entry.getKey(), entry.getValue()); + } + setVersion(fresh.getVersion()); + } } } diff --git a/core-process-orchestrator/src/main/java/com/netcracker/core/scheduler/po/FutureKey.java b/core-process-orchestrator/src/main/java/com/netcracker/core/scheduler/po/FutureKey.java index bfa8988d0..2ea46a1b9 100644 --- a/core-process-orchestrator/src/main/java/com/netcracker/core/scheduler/po/FutureKey.java +++ b/core-process-orchestrator/src/main/java/com/netcracker/core/scheduler/po/FutureKey.java @@ -19,14 +19,9 @@ public FutureKey(String taskId) { @Override public boolean equals(Object obj) { - if (obj instanceof String) { - return taskId.equals(obj); - } else if (obj instanceof UUID) { - return uuid.equals(obj); - } else if (obj instanceof FutureKey futureKey) { - return taskId.equals(futureKey.taskId) && uuid.equals(futureKey.uuid); - } - return false; + return obj instanceof FutureKey futureKey + && taskId.equals(futureKey.taskId) + && uuid.equals(futureKey.uuid); } @Override diff --git a/core-process-orchestrator/src/main/java/com/netcracker/core/scheduler/po/Process.java b/core-process-orchestrator/src/main/java/com/netcracker/core/scheduler/po/Process.java index 2955e9f3f..e90dc330e 100644 --- a/core-process-orchestrator/src/main/java/com/netcracker/core/scheduler/po/Process.java +++ b/core-process-orchestrator/src/main/java/com/netcracker/core/scheduler/po/Process.java @@ -56,7 +56,7 @@ public CompletionHandler execute(com.github.kagkarlsson.schedule logger.info("Process {} are failed", taskInstance.getId()); taskInstance.getData().getProcess().setState(TaskState.FAILED); executionOperations.remove(); - taskInstance.getData().getProcess().save(); + taskInstance.getData().getProcess().saveResolvingConflict(pi -> pi.setState(TaskState.FAILED)); return; } completed = completed && task.getState() == TaskState.COMPLETED; @@ -65,9 +65,13 @@ public CompletionHandler execute(com.github.kagkarlsson.schedule executionOperations.reschedule(executionComplete, Instant.now().plusSeconds(2), taskInstance.getData()); else { logger.info("Process {} are completed", taskInstance.getId()); - taskInstance.getData().getProcess().setEndTime(Calendar.getInstance().getTime()); + java.util.Date endTime = Calendar.getInstance().getTime(); + taskInstance.getData().getProcess().setEndTime(endTime); taskInstance.getData().getProcess().setState(TaskState.COMPLETED); - taskInstance.getData().getProcess().save(); + taskInstance.getData().getProcess().saveResolvingConflict(pi -> { + pi.setEndTime(endTime); + pi.setState(TaskState.COMPLETED); + }); executionOperations.remove(); } taskInstance.getData().getProcess().save(); diff --git a/core-process-orchestrator/src/main/java/com/netcracker/core/scheduler/po/model/pojo/ProcessInstanceImpl.java b/core-process-orchestrator/src/main/java/com/netcracker/core/scheduler/po/model/pojo/ProcessInstanceImpl.java index 8c5dd6431..3aced3d96 100644 --- a/core-process-orchestrator/src/main/java/com/netcracker/core/scheduler/po/model/pojo/ProcessInstanceImpl.java +++ b/core-process-orchestrator/src/main/java/com/netcracker/core/scheduler/po/model/pojo/ProcessInstanceImpl.java @@ -3,11 +3,13 @@ import com.fasterxml.jackson.annotation.JsonIgnoreType; import com.netcracker.core.scheduler.po.DataContext; import com.netcracker.core.scheduler.po.ProcessOrchestrator; +import com.netcracker.core.scheduler.po.repository.VersionMismatchException; import com.netcracker.core.scheduler.po.repository.TaskInstanceRepository; import com.netcracker.core.scheduler.po.task.TaskState; import lombok.Getter; import java.util.Date; +import java.util.function.Consumer; import java.util.List; import java.util.Objects; @@ -87,6 +89,28 @@ public int hashCode() { return Objects.hash(getName(), getId(), getProcessDefinitionID(), getStartTime(), getEndTime(), getState(), getVersion()); } + /** + * Saves, and on a version conflict reloads the row, re-applies the mutation + * to the fresh copy, saves it, and syncs this instance with the persisted + * state — so callers (and any later raw save()) keep working with a clean, + * current object. Meant for must-win writes such as the terminal + * COMPLETED/FAILED transitions; the mutation must be idempotent. + */ + public ProcessInstanceImpl saveResolvingConflict(Consumer reapply) { + try { + save(); + } catch (VersionMismatchException e) { + ProcessInstanceImpl fresh = reload(); + reapply.accept(fresh); + fresh.save(); + setState(fresh.getState()); + setStartTime(fresh.getStartTime()); + setEndTime(fresh.getEndTime()); + setVersion(fresh.getVersion()); + } + return this; + } + public void save() { if (isDirty) ProcessOrchestrator.getInstance().getProcessInstanceRepository().putProcessInstance(this); diff --git a/core-process-orchestrator/src/main/java/com/netcracker/core/scheduler/po/model/pojo/TaskInstanceImpl.java b/core-process-orchestrator/src/main/java/com/netcracker/core/scheduler/po/model/pojo/TaskInstanceImpl.java index 821f57e00..4446e6a34 100644 --- a/core-process-orchestrator/src/main/java/com/netcracker/core/scheduler/po/model/pojo/TaskInstanceImpl.java +++ b/core-process-orchestrator/src/main/java/com/netcracker/core/scheduler/po/model/pojo/TaskInstanceImpl.java @@ -2,6 +2,7 @@ import com.netcracker.core.scheduler.po.DataContext; import com.netcracker.core.scheduler.po.ProcessOrchestrator; +import com.netcracker.core.scheduler.po.repository.VersionMismatchException; import com.netcracker.core.scheduler.po.task.NamedTask; import com.netcracker.core.scheduler.po.task.TaskState; import lombok.Getter; @@ -9,6 +10,7 @@ import java.util.Date; import java.util.List; import java.util.Objects; +import java.util.function.Consumer; public class TaskInstanceImpl { @@ -108,6 +110,33 @@ public TaskInstanceImpl reload() { return ProcessOrchestrator.getInstance().getTaskInstanceRepository().getTaskInstance(id); } + /** + * Saves, and on a version conflict reloads the row, re-applies the mutation + * to the fresh copy, and saves that copy. Meant for must-win writes such as + * the terminal COMPLETED/FAILED transitions: several actors (the process + * tick, the executing task, the failure handler, the timeout watchdog) hold + * independent in-memory copies of the same row, and a completion that loses + * the version race by milliseconds must not abort — an aborted completion + * leaves the db-scheduler execution stuck picked until dead-execution + * revival. The mutation must be idempotent. This instance is synced with + * the persisted state afterwards, so callers keep a clean, current object. + */ + public TaskInstanceImpl saveResolvingConflict(Consumer reapply) { + try { + save(); + } catch (VersionMismatchException e) { + TaskInstanceImpl fresh = reload(); + reapply.accept(fresh); + fresh.save(); + setState(fresh.getState()); + setName(fresh.getName()); + setType(fresh.getType()); + setDependsOn(fresh.getDependsOn()); + setVersion(fresh.getVersion()); + } + return this; + } + /**/ public DataContext getContext() { diff --git a/core-process-orchestrator/src/main/java/com/netcracker/core/scheduler/po/repository/impl/ContextRepositoryImpl.java b/core-process-orchestrator/src/main/java/com/netcracker/core/scheduler/po/repository/impl/ContextRepositoryImpl.java index 882df0190..daf9efa9f 100644 --- a/core-process-orchestrator/src/main/java/com/netcracker/core/scheduler/po/repository/impl/ContextRepositoryImpl.java +++ b/core-process-orchestrator/src/main/java/com/netcracker/core/scheduler/po/repository/impl/ContextRepositoryImpl.java @@ -53,24 +53,31 @@ public void putContext(DataContext context) { if (version == null) { jdbcRunner.execute(insertQuery, (PreparedStatement p) -> assignParameters(context, p)); context.setDirty(false); - } else if (Objects.equals(context.getVersion(), version)) { - context.setVersion(version + 1); - jdbcRunner.execute( + } else { + // The version check must live in the UPDATE itself: a check-then-act + // compare leaves a window where a concurrent writer slips between the + // SELECT above and the UPDATE, silently losing one of the writes. The + // in-memory version is bumped only after the guarded UPDATE succeeded. + int expectedVersion = context.getVersion(); + int updated = jdbcRunner.execute( "update " + tableName + " set " + "version = ?, " + - "context_data = ? where id=?" + "context_data = ? where id=? and version=?" , (PreparedStatement p) -> { - p.setString(3, context.getId()); + p.setInt(1, expectedVersion + 1); p.setObject(2, serializer.serialize(context)); - p.setInt(1, context.getVersion()); + p.setString(3, context.getId()); + p.setInt(4, expectedVersion); } ); + if (updated == 0) { + throw new VersionMismatchException("Current version are less than saved"); + } + context.setVersion(expectedVersion + 1); context.setDirty(false); - } else throw new - - VersionMismatchException("Current version are less than saved"); + } } @Override diff --git a/core-process-orchestrator/src/main/java/com/netcracker/core/scheduler/po/repository/impl/ProcessInstanceRepositoryImpl.java b/core-process-orchestrator/src/main/java/com/netcracker/core/scheduler/po/repository/impl/ProcessInstanceRepositoryImpl.java index 5b9a8141c..fd5dff887 100644 --- a/core-process-orchestrator/src/main/java/com/netcracker/core/scheduler/po/repository/impl/ProcessInstanceRepositoryImpl.java +++ b/core-process-orchestrator/src/main/java/com/netcracker/core/scheduler/po/repository/impl/ProcessInstanceRepositoryImpl.java @@ -37,17 +37,20 @@ public void putProcessInstance(ProcessInstanceImpl processInstance) { if (version == null) { insertInstance(processInstance); } else { - if (version.equals(processInstance.getVersion())) { - processInstance.setVersion(version + 1); - - updateInstance(processInstance); - } else throw new VersionMismatchException("Current version are less than saved"); - + // The version check must live in the UPDATE itself: a check-then-act + // compare leaves a window where a concurrent writer slips between the + // SELECT above and the UPDATE, silently losing one of the writes. The + // in-memory version is bumped only after the guarded UPDATE succeeded. + int expectedVersion = processInstance.getVersion(); + if (updateInstance(processInstance, expectedVersion) == 0) { + throw new VersionMismatchException("Current version are less than saved"); + } + processInstance.setVersion(expectedVersion + 1); } } - private void updateInstance(ProcessInstanceImpl processInstance) { - jdbcRunner.execute( + private int updateInstance(ProcessInstanceImpl processInstance, int expectedVersion) { + return jdbcRunner.execute( "update " + tableName + " set " + @@ -55,7 +58,7 @@ private void updateInstance(ProcessInstanceImpl processInstance) { " start_time=?," + " end_time=?," + " version=?" + - " where pi_id=?" + " where pi_id=? and version=?" , (PreparedStatement p) -> { p.setString(1, processInstance.getState().toString()); @@ -67,9 +70,10 @@ private void updateInstance(ProcessInstanceImpl processInstance) { if (processInstance.getEndTime() != null) { p.setLong(3, processInstance.getEndTime().getTime()); } else p.setLong(3, 0L); - p.setInt(4, processInstance.getVersion()); + p.setInt(4, expectedVersion + 1); p.setString(5, processInstance.getId()); + p.setInt(6, expectedVersion); } ); } diff --git a/core-process-orchestrator/src/main/java/com/netcracker/core/scheduler/po/repository/impl/TaskInstanceRepositoryImpl.java b/core-process-orchestrator/src/main/java/com/netcracker/core/scheduler/po/repository/impl/TaskInstanceRepositoryImpl.java index 4a011e154..3b2d8b165 100644 --- a/core-process-orchestrator/src/main/java/com/netcracker/core/scheduler/po/repository/impl/TaskInstanceRepositoryImpl.java +++ b/core-process-orchestrator/src/main/java/com/netcracker/core/scheduler/po/repository/impl/TaskInstanceRepositoryImpl.java @@ -41,29 +41,35 @@ public void putTaskInstance(TaskInstanceImpl taskInstance) { " (task_id,name,def_id,version, state,type,pi_id,depends_on) values(?,?,?,?,?,?,?,?)" , (PreparedStatement p) -> assignParameters(taskInstance, p)); } else { - if (version.equals(taskInstance.getVersion())) { - taskInstance.setVersion(version + 1); - jdbcRunner.execute( - "update " + - tableName + - " set " + - " state=?," + - " type=?," + - " version=?," + - " name=?," + - " depends_on=?" + - " where task_id=?" - , - (PreparedStatement p) -> { - p.setString(1, taskInstance.getState().toString()); - p.setString(2, taskInstance.getType()); - p.setInt(3, taskInstance.getVersion()); - p.setString(4, taskInstance.getName()); - p.setObject(5, serializer.serialize(taskInstance.getDependsOn())); - p.setString(6, taskInstance.getId()); - - }); - } else throw new VersionMismatchException(String.format("Version in repository: %s, Task Version %s",version,taskInstance.getVersion())); + // The version check must live in the UPDATE itself: a check-then-act + // compare leaves a window where a concurrent writer slips between the + // SELECT above and the UPDATE, silently losing one of the writes. The + // in-memory version is bumped only after the guarded UPDATE succeeded. + int expectedVersion = taskInstance.getVersion(); + int updated = jdbcRunner.execute( + "update " + + tableName + + " set " + + " state=?," + + " type=?," + + " version=?," + + " name=?," + + " depends_on=?" + + " where task_id=? and version=?" + , + (PreparedStatement p) -> { + p.setString(1, taskInstance.getState().toString()); + p.setString(2, taskInstance.getType()); + p.setInt(3, expectedVersion + 1); + p.setString(4, taskInstance.getName()); + p.setObject(5, serializer.serialize(taskInstance.getDependsOn())); + p.setString(6, taskInstance.getId()); + p.setInt(7, expectedVersion); + }); + if (updated == 0) { + throw new VersionMismatchException(String.format("Version in repository: %s, Task Version %s", version, expectedVersion)); + } + taskInstance.setVersion(expectedVersion + 1); } } diff --git a/core-process-orchestrator/src/main/java/com/netcracker/core/scheduler/po/runnable/TaskExecutionWrapper.java b/core-process-orchestrator/src/main/java/com/netcracker/core/scheduler/po/runnable/TaskExecutionWrapper.java index 327bd3fa4..a98358ebc 100644 --- a/core-process-orchestrator/src/main/java/com/netcracker/core/scheduler/po/runnable/TaskExecutionWrapper.java +++ b/core-process-orchestrator/src/main/java/com/netcracker/core/scheduler/po/runnable/TaskExecutionWrapper.java @@ -1,9 +1,14 @@ package com.netcracker.core.scheduler.po.runnable; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + import java.util.concurrent.Callable; public class TaskExecutionWrapper implements Callable { + private static final Logger logger = LoggerFactory.getLogger(TaskExecutionWrapper.class); + private final Runnable callback; private final Runnable task; @@ -14,10 +19,13 @@ public TaskExecutionWrapper(Runnable task, Runnable callback) { @Override public Boolean call() throws Exception { + // The callback must run exactly once: the previous catch-plus-finally + // shape ran it twice on failure. Errors are logged, not rethrown — the + // task's own failure handling is done by the scheduler's failure hooks. try { task.run(); - } catch (Throwable e) { - callback.run(); + } catch (Exception e) { + logger.error("Task execution failed", e); } finally { callback.run(); } diff --git a/core-process-orchestrator/src/main/java/com/netcracker/core/scheduler/po/task/ProcessTaskFailureHandler.java b/core-process-orchestrator/src/main/java/com/netcracker/core/scheduler/po/task/ProcessTaskFailureHandler.java index 0f8068e7b..0ede34e98 100644 --- a/core-process-orchestrator/src/main/java/com/netcracker/core/scheduler/po/task/ProcessTaskFailureHandler.java +++ b/core-process-orchestrator/src/main/java/com/netcracker/core/scheduler/po/task/ProcessTaskFailureHandler.java @@ -30,8 +30,9 @@ public ProcessTaskFailureHandler(FailureHandler delegate) public void onFailure(ExecutionComplete executionComplete, ExecutionOperations executionOperations) { List list = new ArrayList<>(); - if (executionComplete.getCause().isPresent()) { - Throwable ex = executionComplete.getCause().get(); + java.util.Optional cause = executionComplete.getCause(); + if (cause.isPresent()) { + Throwable ex = cause.get(); while (ex != null) { list.add(ex); ex = ex.getCause(); @@ -42,7 +43,7 @@ public void onFailure(ExecutionComplete executionComplete, ExecutionOperations t.setState(TaskState.FAILED)); executionOperations.remove(); logger.error("Task {} with ID:{} failed", executionComplete.getExecution().taskInstance.getTaskName(), executionComplete.getExecution().taskInstance.getId()); } diff --git a/core-process-orchestrator/src/main/java/com/netcracker/core/scheduler/po/task/templates/AbstractProcessTask.java b/core-process-orchestrator/src/main/java/com/netcracker/core/scheduler/po/task/templates/AbstractProcessTask.java index a47511eb7..14f52c33a 100644 --- a/core-process-orchestrator/src/main/java/com/netcracker/core/scheduler/po/task/templates/AbstractProcessTask.java +++ b/core-process-orchestrator/src/main/java/com/netcracker/core/scheduler/po/task/templates/AbstractProcessTask.java @@ -53,7 +53,7 @@ public CompletionHandler execute(TaskInstance t.setState(TaskState.COMPLETED)); logger.info("Task {} with id {} completed.", getTaskName(), taskInstance.getId()); } return (executionComplete, executionOperations) -> executionOperations.remove(); diff --git a/core-process-orchestrator/src/main/java/com/netcracker/core/scheduler/po/task/templates/AsyncTaskWithPolling.java b/core-process-orchestrator/src/main/java/com/netcracker/core/scheduler/po/task/templates/AsyncTaskWithPolling.java index b5156754e..a4a69f487 100644 --- a/core-process-orchestrator/src/main/java/com/netcracker/core/scheduler/po/task/templates/AsyncTaskWithPolling.java +++ b/core-process-orchestrator/src/main/java/com/netcracker/core/scheduler/po/task/templates/AsyncTaskWithPolling.java @@ -70,7 +70,7 @@ public CompletionHandler execute(TaskInstance t.setState(TaskState.COMPLETED)); when.accept(null); } } catch (Exception e) { @@ -78,8 +78,8 @@ public CompletionHandler execute(TaskInstance pi.setState(TaskState.FAILED)); + task.saveResolvingConflict(t -> t.setState(TaskState.FAILED)); when.accept(null); } diff --git a/core-process-orchestrator/src/test/java/com/netcracker/core/scheduler/helpers/TaskExecutorServiceTest.java b/core-process-orchestrator/src/test/java/com/netcracker/core/scheduler/helpers/TaskExecutorServiceTest.java new file mode 100644 index 000000000..aee85cfaf --- /dev/null +++ b/core-process-orchestrator/src/test/java/com/netcracker/core/scheduler/helpers/TaskExecutorServiceTest.java @@ -0,0 +1,98 @@ +package com.netcracker.core.scheduler.helpers; + +import com.netcracker.core.scheduler.po.ProcessOrchestrator; +import com.netcracker.core.scheduler.po.model.pojo.ProcessInstanceImpl; +import com.netcracker.core.scheduler.po.model.pojo.TaskInstanceImpl; +import com.netcracker.core.scheduler.po.samples.tasks.DummyTask2; +import com.netcracker.core.scheduler.po.task.TaskState; +import com.zaxxer.hikari.HikariDataSource; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +import javax.sql.DataSource; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.Future; + +/** + * Covers the sync-timeout watchdog. The completion callback cancels the + * watchdog with an interrupt when the task finishes in time, so an interrupt + * must leave the task and process state untouched; only a real timeout may + * mark them FAILED. + */ +class TaskExecutorServiceTest { + + DataSource dataSource; + ProcessOrchestrator orchestrator; + TaskExecutorService service; + + @BeforeEach + void setup() throws Exception { + dataSource = SchedulerUtils.initDatabase(); + orchestrator = new ProcessOrchestrator(dataSource); + service = new TaskExecutorService(2); + } + + @AfterEach + void teardown() throws Exception { + service.shutdown(); + orchestrator.stop(); + ((HikariDataSource) dataSource).close(); + } + + private void seedTaskAndProcess(TaskState state) { + TaskInstanceImpl task = new TaskInstanceImpl("task-1", "TestTask", DummyTask2.class.getName(), "proc-1"); + task.setState(state); + orchestrator.getTaskInstanceRepository().putTaskInstance(task); + + ProcessInstanceImpl process = new ProcessInstanceImpl("Test Instance", "proc-1", "def-1"); + process.setState(state); + orchestrator.getProcessInstanceRepository().putProcessInstance(process); + } + + @Test + void interruptedWatchdogLeavesACompletedTaskUntouched() throws Exception { + seedTaskAndProcess(TaskState.COMPLETED); + + // The worker future is still pending when the interrupt arrives, exactly + // like a completion callback that cancels the watchdog from the worker's + // finally block before the worker future itself is marked done. + Future worker = new CompletableFuture<>(); + Thread watchdog = new Thread(() -> service.watchTask(worker, 60L, "task-1", "proc-1")); + watchdog.start(); + + long deadline = System.currentTimeMillis() + 5000; + while (watchdog.getState() != Thread.State.TIMED_WAITING && System.currentTimeMillis() < deadline) { + Thread.sleep(10); + } + Assertions.assertEquals(Thread.State.TIMED_WAITING, watchdog.getState(), + "the watchdog must be parked in future.get before the interrupt"); + + watchdog.interrupt(); + watchdog.join(5000); + Assertions.assertFalse(watchdog.isAlive(), "the watchdog must stop after the interrupt"); + + Assertions.assertEquals(TaskState.COMPLETED, + orchestrator.getTaskInstanceRepository().getTaskInstance("task-1").getState(), + "an interrupted watchdog must not mark a completed task FAILED"); + Assertions.assertEquals(TaskState.COMPLETED, + orchestrator.getProcessInstanceRepository().getProcess("proc-1").getState(), + "an interrupted watchdog must not mark the process FAILED"); + Assertions.assertFalse(worker.isCancelled(), "the worker future must not be cancelled"); + } + + @Test + void timedOutWatchdogMarksTaskAndProcessFailed() { + seedTaskAndProcess(TaskState.IN_PROGRESS); + + Future worker = new CompletableFuture<>(); + service.watchTask(worker, 1L, "task-1", "proc-1"); + + Assertions.assertEquals(TaskState.FAILED, + orchestrator.getTaskInstanceRepository().getTaskInstance("task-1").getState()); + Assertions.assertEquals(TaskState.FAILED, + orchestrator.getProcessInstanceRepository().getProcess("proc-1").getState()); + Assertions.assertTrue(worker.isCancelled(), "the hung worker must be cancelled on timeout"); + } +} diff --git a/core-process-orchestrator/src/test/java/com/netcracker/core/scheduler/po/VersionConflictTest.java b/core-process-orchestrator/src/test/java/com/netcracker/core/scheduler/po/VersionConflictTest.java new file mode 100644 index 000000000..aa6aa9129 --- /dev/null +++ b/core-process-orchestrator/src/test/java/com/netcracker/core/scheduler/po/VersionConflictTest.java @@ -0,0 +1,218 @@ +package com.netcracker.core.scheduler.po; + +import com.netcracker.core.scheduler.helpers.SchedulerUtils; +import com.netcracker.core.scheduler.po.model.pojo.ProcessInstanceImpl; +import com.netcracker.core.scheduler.po.model.pojo.TaskInstanceImpl; +import com.netcracker.core.scheduler.po.repository.ContextRepository; +import com.netcracker.core.scheduler.po.repository.ProcessInstanceRepository; +import com.netcracker.core.scheduler.po.repository.TaskInstanceRepository; +import com.netcracker.core.scheduler.po.repository.VersionMismatchException; +import com.netcracker.core.scheduler.po.repository.impl.ContextRepositoryImpl; +import com.netcracker.core.scheduler.po.repository.impl.ProcessInstanceRepositoryImpl; +import com.netcracker.core.scheduler.po.repository.impl.TaskInstanceRepositoryImpl; +import com.netcracker.core.scheduler.po.runnable.TaskExecutionWrapper; +import com.netcracker.core.scheduler.po.samples.tasks.DummyTask2; +import com.netcracker.core.scheduler.po.task.TaskState; +import com.zaxxer.hikari.HikariDataSource; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +import javax.sql.DataSource; +import java.util.concurrent.atomic.AtomicInteger; + +/** + * Covers the optimistic-lock behavior under interleaved writers: two in-memory + * copies of the same row exist (the process tick, the executing task, the + * failure handler, and the timeout watchdog all load their own), one saves + * first, and the loser must either fail atomically or recover by reloading. + */ +class VersionConflictTest { + + DataSource dataSource; + + @BeforeEach + void setup() { + dataSource = SchedulerUtils.initDatabase(); + } + + @AfterEach + void teardown() { + ((HikariDataSource) dataSource).close(); + } + + private TaskInstanceImpl storedTask(TaskInstanceRepository repository) { + TaskInstanceImpl task = new TaskInstanceImpl("task-1", "TestTask", DummyTask2.class.getName(), "proc-1"); + task.setState(TaskState.NOT_STARTED); + repository.putTaskInstance(task); + return task; + } + + @Test + void staleTaskCopyFailsWithoutCorruptingTheRow() { + TaskInstanceRepository repository = new TaskInstanceRepositoryImpl(dataSource); + storedTask(repository); + + TaskInstanceImpl copyA = repository.getTaskInstance("task-1"); + TaskInstanceImpl copyB = repository.getTaskInstance("task-1"); + + copyA.setState(TaskState.IN_PROGRESS); + repository.putTaskInstance(copyA); + + copyB.setState(TaskState.COMPLETED); + Assertions.assertThrows(VersionMismatchException.class, () -> repository.putTaskInstance(copyB)); + + TaskInstanceImpl stored = repository.getTaskInstance("task-1"); + Assertions.assertEquals(TaskState.IN_PROGRESS, stored.getState(), "the winning write must survive"); + Assertions.assertEquals(copyA.getVersion(), stored.getVersion()); + } + + @Test + void staleProcessCopyFailsWithoutCorruptingTheRow() { + ProcessInstanceRepository repository = new ProcessInstanceRepositoryImpl(dataSource); + ProcessInstanceImpl process = new ProcessInstanceImpl("Test Instance", "proc-1", "def-1"); + process.setState(TaskState.NOT_STARTED); + repository.putProcessInstance(process); + + ProcessInstanceImpl copyA = repository.getProcess("proc-1"); + ProcessInstanceImpl copyB = repository.getProcess("proc-1"); + + copyA.setState(TaskState.IN_PROGRESS); + repository.putProcessInstance(copyA); + + copyB.setState(TaskState.FAILED); + Assertions.assertThrows(VersionMismatchException.class, () -> repository.putProcessInstance(copyB)); + + Assertions.assertEquals(TaskState.IN_PROGRESS, repository.getProcess("proc-1").getState()); + } + + @Test + void staleContextCopyFailsWithoutCorruptingTheRow() { + ContextRepository repository = new ContextRepositoryImpl(dataSource); + DataContext context = new DataContext("ctx-1"); + context.setRepository(repository); + context.put("k", "v0"); + repository.putContext(context); + + DataContext copyA = repository.getContext("ctx-1"); + copyA.setRepository(repository); + DataContext copyB = repository.getContext("ctx-1"); + copyB.setRepository(repository); + + copyA.put("k", "vA"); + repository.putContext(copyA); + + copyB.put("k", "vB"); + Assertions.assertThrows(VersionMismatchException.class, () -> repository.putContext(copyB)); + + Assertions.assertEquals("vA", repository.getContext("ctx-1").get("k")); + } + + @Test + void saveResolvingConflictReappliesTheTerminalStateOnAFreshCopy() throws Exception { + // save()/reload() resolve the repository through the orchestrator singleton, + // so bootstrap a real one the same way SchedulerTest does. + ProcessOrchestrator orchestrator = new ProcessOrchestrator(dataSource); + try { + TaskInstanceRepository repository = orchestrator.getTaskInstanceRepository(); + storedTask(repository); + + TaskInstanceImpl worker = repository.getTaskInstance("task-1"); + TaskInstanceImpl tick = repository.getTaskInstance("task-1"); + + // The process tick wins the race with an unrelated bump. + tick.setState(TaskState.IN_PROGRESS); + repository.putTaskInstance(tick); + + // The worker's terminal save must recover and win. + worker.setState(TaskState.COMPLETED); + TaskInstanceImpl persisted = worker.saveResolvingConflict(t -> t.setState(TaskState.COMPLETED)); + + Assertions.assertEquals(TaskState.COMPLETED, persisted.getState()); + TaskInstanceImpl stored = repository.getTaskInstance("task-1"); + Assertions.assertEquals(TaskState.COMPLETED, stored.getState()); + // The caller's copy is synced with the persisted row, so a later + // save on the same object does not fail again. + Assertions.assertEquals(stored.getVersion(), worker.getVersion()); + } finally { + orchestrator.stop(); + } + } + + @Test + void processSaveResolvingConflictReappliesTheTerminalStateOnAFreshCopy() throws Exception { + ProcessOrchestrator orchestrator = new ProcessOrchestrator(dataSource); + try { + ProcessInstanceRepository repository = orchestrator.getProcessInstanceRepository(); + ProcessInstanceImpl seed = new ProcessInstanceImpl("Test Instance", "proc-1", "def-1"); + seed.setState(TaskState.IN_PROGRESS); + repository.putProcessInstance(seed); + + ProcessInstanceImpl worker = repository.getProcess("proc-1"); + ProcessInstanceImpl tick = repository.getProcess("proc-1"); + + tick.setStartTime(java.util.Calendar.getInstance().getTime()); + repository.putProcessInstance(tick); + + worker.setState(TaskState.FAILED); + worker.saveResolvingConflict(pi -> pi.setState(TaskState.FAILED)); + + ProcessInstanceImpl stored = repository.getProcess("proc-1"); + Assertions.assertEquals(TaskState.FAILED, stored.getState()); + Assertions.assertEquals(stored.getVersion(), worker.getVersion(), "the caller's copy must be synced"); + } finally { + orchestrator.stop(); + } + } + + @Test + void contextApplyRecoversFromAConcurrentWriter() { + ContextRepository repository = new ContextRepositoryImpl(dataSource); + DataContext seed = new DataContext("ctx-1"); + seed.setRepository(repository); + seed.put("seed", "1"); + repository.putContext(seed); + + DataContext worker = repository.getContext("ctx-1"); + worker.setRepository(repository); + DataContext tick = repository.getContext("ctx-1"); + tick.setRepository(repository); + + tick.put("tick", "yes"); + repository.putContext(tick); + + worker.apply(c -> c.put("stateDescription", "Done")); + + DataContext stored = repository.getContext("ctx-1"); + Assertions.assertEquals("Done", stored.get("stateDescription"), "the re-applied mutation must land"); + Assertions.assertEquals("yes", stored.get("tick"), "the concurrent writer's data must survive"); + + // The caller's copy is synced after recovery: it sees the merged + // contents and its next save() succeeds instead of conflicting again. + Assertions.assertEquals("yes", worker.get("tick")); + worker.put("after", "recovery"); + Assertions.assertDoesNotThrow(worker::save); + Assertions.assertEquals("recovery", repository.getContext("ctx-1").get("after")); + } + + @Test + void wrapperRunsTheCallbackExactlyOnceOnFailure() throws Exception { + AtomicInteger callbacks = new AtomicInteger(); + TaskExecutionWrapper wrapper = new TaskExecutionWrapper( + () -> { throw new IllegalStateException("boom"); }, + callbacks::incrementAndGet); + + Assertions.assertEquals(Boolean.TRUE, wrapper.call()); + Assertions.assertEquals(1, callbacks.get()); + } + + @Test + void wrapperRunsTheCallbackExactlyOnceOnSuccess() throws Exception { + AtomicInteger callbacks = new AtomicInteger(); + TaskExecutionWrapper wrapper = new TaskExecutionWrapper(() -> { }, callbacks::incrementAndGet); + + Assertions.assertEquals(Boolean.TRUE, wrapper.call()); + Assertions.assertEquals(1, callbacks.get()); + } +}