From 2b96dd750aeaab8f6171f59962dec1e215d0ef3d Mon Sep 17 00:00:00 2001 From: Roman Kichasov Date: Thu, 23 Jul 2026 20:53:52 +0300 Subject: [PATCH 1/4] fix(process-engine): survive optimistic-lock races on task, process, and context rows MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Under load the aggregator's backup flow kept dying with VersionMismatchException ('Version in repository: 2, Task Version 1' / 'Current version are less than saved') from BackupDatabaseTask: several actors hold independent in-memory copies of the same row — the POProcess tick (every 2s), the executing task, ProcessTaskFailureHandler, and the TaskExecutorService timeout watchdog — and the loser of a milliseconds-wide race aborted its completion, leaving the db-scheduler execution stuck picked until dead-execution revival. Externally that surfaced as hung backups and HTTP 500s. Three fixes: - All three repositories now guard the UPDATE itself (WHERE version=?) instead of a check-then-act compare, closing the silent lost-update window, and bump the in-memory version only on success. - Must-win writes recover instead of aborting: terminal COMPLETED/FAILED saves go through TaskInstanceImpl.saveResolvingConflict (reload, re-apply, save), and DataContext.apply re-applies its idempotent mutation on a fresh copy after a conflict. - TaskExecutionWrapper ran its callback twice on failure (catch plus finally); it now runs exactly once and logs the error. --- .../helpers/TaskExecutorService.java | 2 +- .../core/scheduler/po/DataContext.java | 17 +- .../po/model/pojo/TaskInstanceImpl.java | 25 +++ .../impl/ContextRepositoryImpl.java | 25 ++- .../impl/ProcessInstanceRepositoryImpl.java | 24 ++- .../impl/TaskInstanceRepositoryImpl.java | 52 ++--- .../po/runnable/TaskExecutionWrapper.java | 10 +- .../po/task/ProcessTaskFailureHandler.java | 2 +- .../task/templates/AbstractProcessTask.java | 2 +- .../task/templates/AsyncTaskWithPolling.java | 4 +- .../scheduler/po/VersionConflictTest.java | 181 ++++++++++++++++++ 11 files changed, 295 insertions(+), 49 deletions(-) create mode 100644 core-process-orchestrator/src/test/java/com/netcracker/core/scheduler/po/VersionConflictTest.java 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..23afb4464 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 @@ -123,7 +123,7 @@ public void execute(@NotNull Runnable command) { processInstance.save(); TaskInstanceImpl task = ProcessOrchestrator.getInstance().getTaskInstanceRepository().getTaskInstance(taskId); task.setState(TaskState.FAILED); - task.save(); + task.saveResolvingConflict(t -> t.setState(TaskState.FAILED)); } return Boolean.TRUE; 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..50424aa8e 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; @@ -71,8 +72,22 @@ 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(); + } } } 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..d53bd62bb 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,29 @@ 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. Returns the instance that was + * actually persisted. + */ + public TaskInstanceImpl saveResolvingConflict(Consumer reapply) { + try { + save(); + return this; + } catch (VersionMismatchException e) { + TaskInstanceImpl fresh = reload(); + reapply.accept(fresh); + fresh.save(); + return fresh; + } + } + /**/ 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..dbe365a25 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(); + 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..2cb50fd26 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 @@ -42,7 +42,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..e7eaf026d 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) { @@ -79,7 +79,7 @@ public CompletionHandler execute(TaskInstance t.setState(TaskState.FAILED)); when.accept(null); } 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..b5286329f --- /dev/null +++ b/core-process-orchestrator/src/test/java/com/netcracker/core/scheduler/po/VersionConflictTest.java @@ -0,0 +1,181 @@ +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()); + Assertions.assertEquals(TaskState.COMPLETED, repository.getTaskInstance("task-1").getState()); + } 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"); + } + + @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()); + } +} From 5a3352e6fadbc2efdb2a7303120fa3b0e4ec5bf5 Mon Sep 17 00:00:00 2001 From: Roman Kichasov Date: Thu, 23 Jul 2026 21:06:09 +0300 Subject: [PATCH 2/4] fix(process-engine): recover process terminal saves and sync callers after recovery MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Review findings on the race fix: - Process-level terminal writes (POProcess completion FAILED/COMPLETED, AsyncTaskWithPolling failure, the timeout watchdog) still called raw save() and could abort on the same version race the task writes were protected from. ProcessInstanceImpl gains the same saveResolvingConflict, used at all four sites; the pre-existing reload() is reused. - After conflict recovery both saveResolvingConflict variants and DataContext.apply now sync the caller's instance with the persisted state (contents, fields, version, clean dirty flag). Without that the caller kept a stale copy and the very next save() — for example the trailing process save in the POProcess completion handler — failed again despite the recovery having succeeded. --- .../helpers/TaskExecutorService.java | 2 +- .../core/scheduler/po/DataContext.java | 9 +++++ .../netcracker/core/scheduler/po/Process.java | 10 +++-- .../po/model/pojo/ProcessInstanceImpl.java | 24 ++++++++++++ .../po/model/pojo/TaskInstanceImpl.java | 12 ++++-- .../task/templates/AsyncTaskWithPolling.java | 2 +- .../scheduler/po/VersionConflictTest.java | 39 ++++++++++++++++++- 7 files changed, 88 insertions(+), 10 deletions(-) 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 23afb4464..a65359329 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 @@ -120,7 +120,7 @@ public void execute(@NotNull Runnable command) { future.cancel(true); ProcessInstanceImpl processInstance = ProcessOrchestrator.getInstance().getProcessInstance(taskInstance.getProcessID()); processInstance.setState(TaskState.FAILED); - processInstance.save(); + 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)); 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 50424aa8e..157718383 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 @@ -88,6 +88,15 @@ public void apply(Consumer function) { 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/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 d53bd62bb..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 @@ -118,19 +118,23 @@ public TaskInstanceImpl reload() { * 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. Returns the instance that was - * actually persisted. + * 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(); - return this; } catch (VersionMismatchException e) { TaskInstanceImpl fresh = reload(); reapply.accept(fresh); fresh.save(); - return fresh; + setState(fresh.getState()); + setName(fresh.getName()); + setType(fresh.getType()); + setDependsOn(fresh.getDependsOn()); + setVersion(fresh.getVersion()); } + return this; } /**/ 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 e7eaf026d..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 @@ -78,7 +78,7 @@ 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/po/VersionConflictTest.java b/core-process-orchestrator/src/test/java/com/netcracker/core/scheduler/po/VersionConflictTest.java index b5286329f..aa6aa9129 100644 --- 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 @@ -130,7 +130,37 @@ void saveResolvingConflictReappliesTheTerminalStateOnAFreshCopy() throws Excepti TaskInstanceImpl persisted = worker.saveResolvingConflict(t -> t.setState(TaskState.COMPLETED)); Assertions.assertEquals(TaskState.COMPLETED, persisted.getState()); - Assertions.assertEquals(TaskState.COMPLETED, repository.getTaskInstance("task-1").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(); } @@ -157,6 +187,13 @@ void contextApplyRecoversFromAConcurrentWriter() { 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 From 36b4a12cdd5c124dcd4626862448298fd525418b Mon Sep 17 00:00:00 2001 From: Roman Kichasov Date: Thu, 23 Jul 2026 22:01:22 +0300 Subject: [PATCH 3/4] fix(process-engine): address Sonar findings in touched code - FutureKey.equals no longer accepts String/UUID (equals-contract violation); terminate() now filters by getTaskId() explicitly (S2159) - re-interrupt the thread on InterruptedException in the watchdog and in terminate()'s wait loop (S2142) - terminate() returns Future instead of Future (S1452) - TaskExecutionWrapper catches Exception, not Throwable (S1181) - ProcessTaskFailureHandler reuses a single Optional instance (S3655) - suppress S3011 (reflective access into db-scheduler internals) and S2160 (HashMap content equality is intentional) with justifications --- .../helpers/TaskExecutorService.java | 21 ++++++++++++++----- .../core/scheduler/po/DataContext.java | 3 +++ .../core/scheduler/po/FutureKey.java | 11 +++------- .../po/runnable/TaskExecutionWrapper.java | 2 +- .../po/task/ProcessTaskFailureHandler.java | 5 +++-- 5 files changed, 26 insertions(+), 16 deletions(-) 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 a65359329..446f4ca1b 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 @@ -116,6 +116,9 @@ public void execute(@NotNull Runnable command) { try { future.get(timeout, TimeUnit.SECONDS); } catch (InterruptedException | ExecutionException | TimeoutException e) { + if (e instanceof InterruptedException) { + Thread.currentThread().interrupt(); + } log.info("Task {} execution was interrupted by timeout", taskId); future.cancel(true); ProcessInstanceImpl processInstance = ProcessOrchestrator.getInstance().getProcessInstance(taskInstance.getProcessID()); @@ -135,6 +138,10 @@ public void execute(@NotNull Runnable command) { } + // 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 +166,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 +178,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 157718383..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 @@ -18,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 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/runnable/TaskExecutionWrapper.java b/core-process-orchestrator/src/main/java/com/netcracker/core/scheduler/po/runnable/TaskExecutionWrapper.java index dbe365a25..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 @@ -24,7 +24,7 @@ public Boolean call() throws Exception { // task's own failure handling is done by the scheduler's failure hooks. try { task.run(); - } catch (Throwable e) { + } 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 2cb50fd26..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(); From 30d7341a1e8bb7f62293e00c48dbae0a0dd86e61 Mon Sep 17 00:00:00 2001 From: Roman Kichasov Date: Thu, 23 Jul 2026 22:16:27 +0300 Subject: [PATCH 4/4] fix(process-engine): do not fail a task when its watchdog is interrupted MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit An interrupt is the watchdog's normal shutdown signal: the completion callback in execute() cancels it with cancel(true) once the task finishes in time. The shared catch treated that interrupt as a timeout and marked the completed task and its process FAILED — and saveResolvingConflict made the bogus overwrite reliable instead of losing the version race. Split the catch: InterruptedException restores the interrupt flag and stops watching without touching state; only ExecutionException and TimeoutException cancel the worker and mark the task and process FAILED. The watchdog body moved to a package-private watchTask() so tests can drive both paths directly. --- .../helpers/TaskExecutorService.java | 43 ++++---- .../helpers/TaskExecutorServiceTest.java | 98 +++++++++++++++++++ 2 files changed, 122 insertions(+), 19 deletions(-) create mode 100644 core-process-orchestrator/src/test/java/com/netcracker/core/scheduler/helpers/TaskExecutorServiceTest.java 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 446f4ca1b..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,31 +112,36 @@ 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) { - if (e instanceof InterruptedException) { - Thread.currentThread().interrupt(); - } - log.info("Task {} execution was interrupted by timeout", taskId); - future.cancel(true); - ProcessInstanceImpl processInstance = ProcessOrchestrator.getInstance().getProcessInstance(taskInstance.getProcessID()); - 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; - })); + 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 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"); + } +}