diff --git a/cmvr-iot-admin/src/main/java/com/cmvr/web/controller/test/TeTaskOrchestrationController.java b/cmvr-iot-admin/src/main/java/com/cmvr/web/controller/test/TeTaskOrchestrationController.java index 087674d..be5dc9e 100644 --- a/cmvr-iot-admin/src/main/java/com/cmvr/web/controller/test/TeTaskOrchestrationController.java +++ b/cmvr-iot-admin/src/main/java/com/cmvr/web/controller/test/TeTaskOrchestrationController.java @@ -33,7 +33,6 @@ class TeTaskOrchestrationController extends BaseController { @NotEmpty(message = "任务ID不能为空") @RequestParam("taskId") String taskId ) { - startPage(); List list = teTaskOrchestrationService.queryByTaskId(taskId); return getDataTable(list); } diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowTaskEngine.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowTaskEngine.java index a10077f..8d72639 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowTaskEngine.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowTaskEngine.java @@ -41,6 +41,9 @@ public class FlowTaskEngine { Map itemOccurrences = new HashMap<>(); for (int i = 0; i < totalItemCount; i++) { + // Scheduler deduplication is scoped to one orchestration item. Different + // items may legitimately contain the same node IDs (for example, copied flows). + flowTaskScheduler.clear(instId); TeQueryTaskDetailDTO item = items.get(i); // 等待当前 item 完成后再执行下一个 CountDownLatch latch = new CountDownLatch(1); diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/support/FlowSchedulerCache.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/support/FlowSchedulerCache.java index e2977bc..e86e999 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/support/FlowSchedulerCache.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/support/FlowSchedulerCache.java @@ -24,7 +24,8 @@ public class FlowSchedulerCache { } public void clearContextCache(String instId) { - scheduled.keySet().removeIf(k -> k.startsWith(instId + "_")); - completedPreMap.keySet().removeIf(k -> k.startsWith(instId + "_")); + String prefix = instId + ":"; + scheduled.keySet().removeIf(k -> k.startsWith(prefix)); + completedPreMap.keySet().removeIf(k -> k.startsWith(prefix)); } } diff --git a/cmvr-iot-test/src/test/java/com/cmvr/test/flow/runtime/engine/support/FlowSchedulerCacheTest.java b/cmvr-iot-test/src/test/java/com/cmvr/test/flow/runtime/engine/support/FlowSchedulerCacheTest.java new file mode 100644 index 0000000..dce68a3 --- /dev/null +++ b/cmvr-iot-test/src/test/java/com/cmvr/test/flow/runtime/engine/support/FlowSchedulerCacheTest.java @@ -0,0 +1,33 @@ +package com.cmvr.test.flow.runtime.engine.support; + +import org.junit.Test; + +import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertTrue; + +public class FlowSchedulerCacheTest { + + @Test + public void clearAllowsCopiedNodeIdToBeScheduledByNextItem() { + FlowSchedulerCache cache = new FlowSchedulerCache(); + String key = "inst-1:copied-node"; + + assertTrue(cache.trySchedule(key)); + assertFalse(cache.trySchedule(key)); + + cache.clearContextCache("inst-1"); + + assertTrue(cache.trySchedule(key)); + } + + @Test + public void clearDoesNotAffectAnotherInstanceWithSimilarId() { + FlowSchedulerCache cache = new FlowSchedulerCache(); + String otherInstanceKey = "inst-10:node-1"; + assertTrue(cache.trySchedule(otherInstanceKey)); + + cache.clearContextCache("inst-1"); + + assertFalse(cache.trySchedule(otherInstanceKey)); + } +}