From 8d8d77ffdaa4d8aa64c77708d68b3e4f6bd49f6d Mon Sep 17 00:00:00 2001
From: =?UTF-8?q?=E9=99=88=E5=AD=90=E9=BB=98?= <925456043@qq.com>
Date: Sun, 9 Aug 2026 21:24:40 +0800
Subject: [PATCH] =?UTF-8?q?feat:=20=E5=AE=8C=E5=96=84=E5=B7=A5=E4=BD=9C?=
=?UTF-8?q?=E6=B5=81=E5=AE=9E=E4=BE=8B=E7=8A=B6=E6=80=81=E6=9F=A5=E8=AF=A2?=
=?UTF-8?q?=E4=B8=8E=E6=81=A2=E5=A4=8D?=
MIME-Version: 1.0
Content-Type: text/plain; charset=UTF-8
Content-Transfer-Encoding: 8bit
- 增加暂停态原子恢复守卫,避免重复恢复覆盖实例状态
- 支持从实例定义快照读取节点名称
---
.../com/easyagents/flow/core/chain/Chain.java | 40 ++++++++
.../core/chain/runtime/ChainExecutor.java | 76 +++++++++++++++
.../ChainExecutorInstanceMetadataTest.java | 68 +++++++++++++
.../flow/core/test/ChainResumeGuardTest.java | 96 +++++++++++++++++++
4 files changed, 280 insertions(+)
create mode 100644 easy-agents-flow/src/test/java/com/easyagents/flow/core/test/ChainExecutorInstanceMetadataTest.java
create mode 100644 easy-agents-flow/src/test/java/com/easyagents/flow/core/test/ChainResumeGuardTest.java
diff --git a/easy-agents-flow/src/main/java/com/easyagents/flow/core/chain/Chain.java b/easy-agents-flow/src/main/java/com/easyagents/flow/core/chain/Chain.java
index 9fc1fb1..b04db36 100644
--- a/easy-agents-flow/src/main/java/com/easyagents/flow/core/chain/Chain.java
+++ b/easy-agents-flow/src/main/java/com/easyagents/flow/core/chain/Chain.java
@@ -1732,7 +1732,47 @@ public class Chain {
}
+ /**
+ * 仅在工作流处于暂停状态时恢复执行。
+ *
+ *
状态判断与恢复动作在同一个实例锁内完成,避免并发恢复请求重复注入变量
+ * 或把终态实例重新改为运行中。
+ *
+ * @param variables 恢复时注入的变量
+ * @return 本次是否完成了暂停态到运行态的转换
+ */
+ public boolean resumeIfSuspended(Map variables) {
+ return executeWithLock(
+ stateInstanceId,
+ 10L,
+ TimeUnit.SECONDS,
+ () -> {
+ ChainState current =
+ chainStateRepository.load(stateInstanceId);
+ if (current == null
+ || current.getStatus() != ChainStatus.SUSPEND) {
+ return false;
+ }
+ resumeSuspended(variables);
+ return true;
+ });
+ }
+
+ /**
+ * 恢复暂停中的工作流。
+ *
+ * @param variables 恢复时注入的变量
+ */
public void resume(Map variables) {
+ resumeIfSuspended(variables);
+ }
+
+ /**
+ * 在调用方持有实例锁且已确认暂停状态后执行恢复动作。
+ *
+ * @param variables 恢复时注入的变量
+ */
+ private void resumeSuspended(Map variables) {
ChainState newState = updateStateSafely(state -> {
if (variables != null) {
state.getMemory().putAll(variables);
diff --git a/easy-agents-flow/src/main/java/com/easyagents/flow/core/chain/runtime/ChainExecutor.java b/easy-agents-flow/src/main/java/com/easyagents/flow/core/chain/runtime/ChainExecutor.java
index c2491e2..a752ada 100644
--- a/easy-agents-flow/src/main/java/com/easyagents/flow/core/chain/runtime/ChainExecutor.java
+++ b/easy-agents-flow/src/main/java/com/easyagents/flow/core/chain/runtime/ChainExecutor.java
@@ -896,6 +896,82 @@ public class ChainExecutor {
chain.resume(variables);
}
+ /**
+ * 仅在工作流实例处于暂停状态时恢复执行。
+ *
+ * 状态判断和恢复由 {@link Chain} 在同一个实例锁内完成,可安全处理并发恢复请求。
+ *
+ * @param stateInstanceId 工作流实例 ID
+ * @param variables 恢复时注入的变量
+ * @return 本次是否完成了暂停态到运行态的转换
+ */
+ public boolean resumeAsyncIfSuspended(
+ String stateInstanceId,
+ Map variables) {
+ ChainState state = chainStateRepository.load(stateInstanceId);
+ if (state == null) {
+ return false;
+ }
+ ChainDefinition definition = getDefinitionForInstance(state);
+ if (definition == null) {
+ return false;
+ }
+ Chain chain = configureChain(
+ definition,
+ state.getInstanceId(),
+ state);
+ return chain.resumeIfSuspended(variables);
+ }
+
+ /**
+ * 获取工作流实例启动时定义快照中的节点名称。
+ *
+ * @param stateInstanceId 工作流实例 ID
+ * @return 按定义顺序排列的节点 ID 与名称;实例或定义不存在时返回空映射
+ */
+ public Map getInstanceNodeNames(
+ String stateInstanceId) {
+ ChainState state = chainStateRepository.load(stateInstanceId);
+ if (state == null) {
+ return Collections.emptyMap();
+ }
+ return getInstanceNodeNames(state);
+ }
+
+ /**
+ * 使用调用方已经加载的状态获取实例定义快照中的节点名称。
+ *
+ * @param state 已加载的工作流状态
+ * @return 按定义顺序排列的节点 ID 与名称;状态或定义不存在时返回空映射
+ */
+ public Map getInstanceNodeNames(
+ ChainState state) {
+ if (state == null) {
+ return Collections.emptyMap();
+ }
+ ChainDefinition definition = getDefinitionForInstance(state);
+ if (definition == null
+ || definition.getNodes() == null
+ || definition.getNodes().isEmpty()) {
+ return Collections.emptyMap();
+ }
+ Map nodeNames = new LinkedHashMap<>();
+ for (Node node : definition.getNodes()) {
+ if (node == null
+ || node.getId() == null
+ || node.getId().isBlank()) {
+ continue;
+ }
+ String nodeName = node.getName();
+ nodeNames.put(
+ node.getId(),
+ nodeName == null || nodeName.isBlank()
+ ? node.getId()
+ : nodeName);
+ }
+ return Collections.unmodifiableMap(nodeNames);
+ }
+
private Chain createChain(String definitionId) {
ChainDefinition definition = definitionRepository.getChainDefinitionById(definitionId);
diff --git a/easy-agents-flow/src/test/java/com/easyagents/flow/core/test/ChainExecutorInstanceMetadataTest.java b/easy-agents-flow/src/test/java/com/easyagents/flow/core/test/ChainExecutorInstanceMetadataTest.java
new file mode 100644
index 0000000..9be3d60
--- /dev/null
+++ b/easy-agents-flow/src/test/java/com/easyagents/flow/core/test/ChainExecutorInstanceMetadataTest.java
@@ -0,0 +1,68 @@
+package com.easyagents.flow.core.test;
+
+import com.easyagents.flow.core.chain.ChainDefinition;
+import com.easyagents.flow.core.chain.repository.InMemoryChainStateRepository;
+import com.easyagents.flow.core.chain.repository.InMemoryNodeStateRepository;
+import com.easyagents.flow.core.chain.runtime.ChainExecutor;
+import com.easyagents.flow.core.chain.runtime.InMemoryTriggerStore;
+import com.easyagents.flow.core.chain.runtime.TriggerScheduler;
+import com.easyagents.flow.core.node.StartNode;
+import org.junit.Assert;
+import org.junit.Test;
+
+import java.util.Collections;
+import java.util.Map;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.ScheduledExecutorService;
+
+/**
+ * {@link ChainExecutor} 实例定义元数据查询测试。
+ */
+public class ChainExecutorInstanceMetadataTest {
+
+ /**
+ * 验证节点名称来自实例启动时可恢复的定义快照。
+ */
+ @Test
+ public void shouldResolveNodeNamesFromInstanceDefinition() {
+ ChainDefinition definition = new ChainDefinition();
+ definition.setId("metadata-definition");
+ StartNode start = new StartNode();
+ start.setId("start");
+ start.setName("开始节点");
+ definition.setNodes(Collections.singletonList(start));
+ definition.setEdges(Collections.emptyList());
+
+ InMemoryChainStateRepository stateRepository =
+ new InMemoryChainStateRepository();
+ stateRepository.create("metadata-instance")
+ .setChainDefinitionId(definition.getId());
+ ScheduledExecutorService schedulerPool =
+ Executors.newSingleThreadScheduledExecutor();
+ ExecutorService workerPool =
+ Executors.newSingleThreadExecutor();
+ TriggerScheduler scheduler = new TriggerScheduler(
+ new InMemoryTriggerStore(),
+ schedulerPool,
+ workerPool,
+ 1_000L);
+ ChainExecutor executor = new ChainExecutor(
+ ignored -> definition,
+ stateRepository,
+ new InMemoryNodeStateRepository(),
+ scheduler);
+
+ try {
+ Map nodeNames =
+ executor.getInstanceNodeNames(
+ "metadata-instance");
+
+ Assert.assertEquals(
+ Map.of("start", "开始节点"),
+ nodeNames);
+ } finally {
+ scheduler.shutdown();
+ }
+ }
+}
diff --git a/easy-agents-flow/src/test/java/com/easyagents/flow/core/test/ChainResumeGuardTest.java b/easy-agents-flow/src/test/java/com/easyagents/flow/core/test/ChainResumeGuardTest.java
new file mode 100644
index 0000000..6013c34
--- /dev/null
+++ b/easy-agents-flow/src/test/java/com/easyagents/flow/core/test/ChainResumeGuardTest.java
@@ -0,0 +1,96 @@
+package com.easyagents.flow.core.test;
+
+import com.easyagents.flow.core.chain.Chain;
+import com.easyagents.flow.core.chain.ChainDefinition;
+import com.easyagents.flow.core.chain.ChainStatus;
+import com.easyagents.flow.core.chain.EventManager;
+import com.easyagents.flow.core.chain.repository.InMemoryChainStateRepository;
+import com.easyagents.flow.core.chain.repository.InMemoryNodeStateRepository;
+import org.junit.Assert;
+import org.junit.Test;
+
+import java.util.Collections;
+import java.util.Map;
+
+/**
+ * {@link Chain} 暂停恢复状态守卫测试。
+ */
+public class ChainResumeGuardTest {
+
+ /**
+ * 验证只有暂停中的实例可以恢复,重复恢复不会再次注入变量。
+ */
+ @Test
+ public void shouldResumeOnlyOnceFromSuspendedState() {
+ InMemoryChainStateRepository stateRepository =
+ new InMemoryChainStateRepository();
+ Chain chain = createChain(stateRepository, "resume-once");
+ chain.suspend();
+
+ boolean resumed = chain.resumeIfSuspended(
+ Map.of("approved", true));
+ boolean resumedAgain = chain.resumeIfSuspended(
+ Map.of("unexpected", true));
+
+ Assert.assertTrue(resumed);
+ Assert.assertFalse(resumedAgain);
+ Assert.assertEquals(
+ ChainStatus.RUNNING,
+ stateRepository.load("resume-once").getStatus());
+ Assert.assertEquals(
+ Boolean.TRUE,
+ stateRepository.load("resume-once")
+ .getMemory()
+ .get("approved"));
+ Assert.assertFalse(
+ stateRepository.load("resume-once")
+ .getMemory()
+ .containsKey("unexpected"));
+ }
+
+ /**
+ * 验证成功终态不会被恢复操作改回运行中。
+ */
+ @Test
+ public void shouldKeepTerminalStateUnchanged() {
+ InMemoryChainStateRepository stateRepository =
+ new InMemoryChainStateRepository();
+ Chain chain = createChain(stateRepository, "resume-terminal");
+ stateRepository.load("resume-terminal")
+ .setStatus(ChainStatus.SUCCEEDED);
+
+ boolean resumed = chain.resumeIfSuspended(
+ Map.of("unexpected", true));
+
+ Assert.assertFalse(resumed);
+ Assert.assertEquals(
+ ChainStatus.SUCCEEDED,
+ stateRepository.load("resume-terminal").getStatus());
+ Assert.assertFalse(
+ stateRepository.load("resume-terminal")
+ .getMemory()
+ .containsKey("unexpected"));
+ }
+
+ /**
+ * 创建使用进程内状态仓储的最小工作流。
+ *
+ * @param stateRepository 状态仓储
+ * @param instanceId 工作流实例 ID
+ * @return 已配置工作流
+ */
+ private Chain createChain(
+ InMemoryChainStateRepository stateRepository,
+ String instanceId) {
+ ChainDefinition definition = new ChainDefinition();
+ definition.setId("resume-definition");
+ definition.setNodes(Collections.emptyList());
+ definition.setEdges(Collections.emptyList());
+ Chain chain = new Chain(definition, instanceId);
+ chain.setChainStateRepository(stateRepository);
+ chain.setNodeStateRepository(
+ new InMemoryNodeStateRepository());
+ chain.setEventManager(new EventManager());
+ return chain;
+ }
+}