20260928171336

This commit is contained in:
oneao committed 2026-09-28 17:13:37 +08:00
1 parent f22e67fc85
commit cbd49ded6f
34 files changed
+738 -7842

No files matched your search

@@ -25,6 +25,7 @@ import java.util.LinkedHashMap;
import java.util.List;
import java.util.Locale;
import java.util.Map;
import java.util.Set;
import java.util.StringJoiner;
@Service
@@ -46,6 +47,19 @@ public class DataSaveService {
private static final String AUDIT_UPDATED_BY = "b_updated_by";
private static final String AUDIT_UPDATED_AT = "b_updated_at";
/**
* Workflow 运行表只能由 WorkflowRuntimeService 写入。若允许通用 saveobjt 直写,
* 客户端可以绕过任务归属、动作权限、实例锁和 outbox,伪造审批状态。
*/
private static final Set<String> WORKFLOW_RUNTIME_TABLES = Set.of(
"wf_instance",
"wf_node_run",
"wf_task",
"wf_action",
"wf_cc",
"wf_outbox"
);
@Resource
private DataSource dataSource;
@@ -482,6 +496,11 @@ public class DataSaveService {
String tableName = ParamUtils.getRequiredString(request, "table");
String keyField = ParamUtils.getRequiredString(request, "key_field");
if (isWorkflowRuntimeTable(tableName)) {
throw new BusinessException(
"Workflow 运行表必须通过流程服务写入,不能使用通用 saveobjt: " + tableName
);
}
DbUtils.TableMetadata table = loadMetadata(
connection,
tableName,
@@ -594,6 +613,14 @@ public class DataSaveService {
);
}
private boolean isWorkflowRuntimeTable(String tableName) {
String normalized = tableName == null ? "" : tableName.trim().toLowerCase(Locale.ROOT);
if (normalized.startsWith("dbo.")) {
normalized = normalized.substring("dbo.".length());
}
return WORKFLOW_RUNTIME_TABLES.contains(normalized);
}
/**
* 执行请求声明的依赖行删除(设计《删除策略》§6.2)。
*
@@ -59,6 +59,7 @@ public class WorkflowConfigService {
new ActionPower("approve", "通过"),
new ActionPower("reject", "驳回"),
new ActionPower("return", "退回"),
new ActionPower("claim", "抢占"),
new ActionPower("withdraw", "撤回"),
new ActionPower("transfer", "转交"),
new ActionPower("delegate", "委托"),
@@ -335,7 +336,7 @@ public class WorkflowConfigService {
connection.setAutoCommit(false);
try {
Map<String, Object> version = queryOne(connection,
"select b_id, b_definition_id, b_status from dbo.wf_version where b_id = ?", versionId);
"select b_id, b_definition_id, b_status from dbo.wf_version with (UPDLOCK, ROWLOCK) where b_id = ?", versionId);
if (version == null) {
throw new BusinessException("版本「" + versionId + "」不存在");
}
@@ -387,7 +388,7 @@ public class WorkflowConfigService {
connection.setAutoCommit(false);
try {
Map<String, Object> version = queryOne(connection,
"select b_status from dbo.wf_version where b_id = ?", versionId);
"select b_status from dbo.wf_version with (UPDLOCK, ROWLOCK) where b_id = ?", versionId);
if (version == null) {
throw new BusinessException("版本「" + versionId + "」不存在");
}
@@ -740,10 +741,13 @@ public class WorkflowConfigService {
) throws SQLException {
if (versionId != null && !versionId.isBlank()) {
Map<String, Object> version = queryOne(connection,
"select b_id, b_version_no, b_status from dbo.wf_version where b_id = ?", versionId);
"select b_id, b_definition_id, b_version_no, b_status from dbo.wf_version where b_id = ?", versionId);
if (version == null) {
throw new BusinessException("版本「" + versionId + "」不存在");
}
if (!definitionId.equals(text(version.get("b_definition_id")))) {
throw new BusinessException("版本「" + versionId + "」不属于流程定义「" + definitionId + "」");
}
if (!STATUS_DRAFT.equals(String.valueOf(version.get("b_status")))) {
throw new BusinessException("已发布或已停用的版本不能修改,请先复制出新版本");
}
@@ -751,8 +755,9 @@ public class WorkflowConfigService {
}
Map<String, Object> draft = queryOne(connection, """
select top 1 b_id, b_version_no, b_status from dbo.wf_version
where b_definition_id = ? and b_status = ? order by b_version_no desc
select top 1 b_id, b_version_no, b_status from dbo.wf_version with (UPDLOCK, HOLDLOCK)
where b_definition_id = ? and b_status = ?
order by b_version_no desc
""", definitionId, STATUS_DRAFT);
if (draft != null) {
return draft;
@@ -162,7 +162,19 @@ public class WorkflowRuntimeService {
connection.commit();
return submittedResult(connection, instance.instanceId(), false);
} catch (SQLException | RuntimeException exception) {
} catch (SQLException exception) {
rollback(connection, exception);
if (isUniqueViolation(exception) && !idempotencyKey.isEmpty()) {
Map<String, Object> duplicate = WorkflowSql.queryOne(connection, """
select top 1 b_id from dbo.wf_instance
where b_module_id = ? and b_business_id = ? and b_idempotency_key = ?
""", moduleId, businessId, idempotencyKey);
if (duplicate != null) {
return submittedResult(connection, String.valueOf(duplicate.get("b_id")), true);
}
}
throw exception;
} catch (RuntimeException exception) {
rollback(connection, exception);
throw exception;
} finally {
@@ -273,12 +285,12 @@ public class WorkflowRuntimeService {
cancelSiblingTasks(connection, instance, text(nodeRun.get("b_id")), text(taskId),
"同节点其余待办已自动取消");
writeAction(connection, instance, text(nodeRun.get("b_id")), taskId, "approve",
userId, "pending", "approved", comment, idempotencyKey, null);
userId, status, "approved", comment, idempotencyKey, null);
advance(connection, instance, text(nodeRun.get("b_node_key")), variables, 0);
projectBusinessStatus(connection, instance, currentStatus(connection, instance.instanceId()));
} else {
writeAction(connection, instance, text(nodeRun.get("b_id")), taskId, "approve",
userId, "pending", "running", comment, idempotencyKey, null);
userId, status, "running", comment, idempotencyKey, null);
updateInstanceSummary(connection, instance, "审批中:" + nodeName);
}
} else {
@@ -295,22 +307,22 @@ public class WorkflowRuntimeService {
if ("return".equals(action)) {
writeAction(connection, instance, text(nodeRun.get("b_id")), taskId, "return",
userId, "pending", "returned", comment, idempotencyKey, null);
userId, status, "returned", comment, idempotencyKey, null);
finishInstance(connection, instance, "returned", nodeName);
} else if ("return_previous".equals(rejectMode)) {
writeAction(connection, instance, text(nodeRun.get("b_id")), taskId, "reject",
userId, "pending", "rejected", comment, idempotencyKey,
userId, status, "rejected", comment, idempotencyKey,
toJson(Map.of("rejectMode", rejectMode)));
reenterPreviousNode(connection, instance, text(nodeRun.get("b_node_key")), variables);
} else if ("return_initiator".equals(rejectMode)) {
writeAction(connection, instance, text(nodeRun.get("b_id")), taskId, "reject",
userId, "pending", "returned", comment, idempotencyKey,
userId, status, "returned", comment, idempotencyKey,
toJson(Map.of("rejectMode", rejectMode)));
finishInstance(connection, instance, "returned", nodeName);
} else {
completeNodeRun(connection, text(nodeRun.get("b_id")), "completed");
writeAction(connection, instance, text(nodeRun.get("b_id")), taskId, "reject",
userId, "pending", "rejected", comment, idempotencyKey,
userId, status, "rejected", comment, idempotencyKey,
toJson(Map.of("rejectMode", "terminate")));
finishInstance(connection, instance, "rejected", nodeName);
}
@@ -318,7 +330,18 @@ public class WorkflowRuntimeService {
connection.commit();
return submittedResult(connection, instance.instanceId(), false);
} catch (SQLException | RuntimeException exception) {
} catch (SQLException exception) {
rollback(connection, exception);
if (isUniqueViolation(exception) && !idempotencyKey.isEmpty()) {
Map<String, Object> duplicate = WorkflowSql.queryOne(connection,
"select top 1 b_id from dbo.wf_action where b_request_id = ?",
idempotencyKey);
if (duplicate != null) {
return Map.of("duplicated", true, "message", "重复请求已忽略");
}
}
throw exception;
} catch (RuntimeException exception) {
rollback(connection, exception);
throw exception;
} finally {
@@ -415,7 +438,8 @@ public class WorkflowRuntimeService {
connection.setAutoCommit(false);
try {
Map<String, Object> task = WorkflowSql.queryOne(connection, """
select b_id, b_instance_id, b_status, b_claimed_by from dbo.wf_task with (UPDLOCK, ROWLOCK)
select b_id, b_instance_id, b_node_key, b_node_run_id, b_status, b_claimed_by
from dbo.wf_task with (UPDLOCK, ROWLOCK)
where b_id = ?
""", uuidParam(taskId));
if (task == null) {
@@ -425,12 +449,23 @@ public class WorkflowRuntimeService {
throw new BusinessException("任务已被处理或已被占用");
}
InstanceRow instance = lockInstance(connection, text(task.get("b_instance_id")));
if (!"running".equals(instance.status())) {
throw new BusinessException("流程已结束,不能抢占待办");
}
requireAction(connection, instance.moduleId(), userId, "claim");
Map<String, Object> node = WorkflowSql.queryOne(connection, """
select b_config_json from dbo.wf_node
where b_version_id = ? and b_node_key = ?
""", instance.versionId(), text(task.get("b_node_key")));
if (!configFlag(node == null ? null : node.get("b_config_json"), "allowClaim")) {
throw new BusinessException("当前节点未开启抢占待办");
}
WorkflowSql.update(connection, """
update dbo.wf_task
set b_status = 'claimed', b_claimed_by = ?, b_claimed_at = sysutcdatetime(), b_updated_at = sysutcdatetime()
where b_id = ?
""", userId, uuidParam(taskId));
writeAction(connection, instance, null, taskId, "claim", userId, "pending", "pending",
writeAction(connection, instance, text(task.get("b_node_run_id")), taskId, "claim", userId, "pending", "claimed",
"抢占待办", null, null);
connection.commit();
return Map.of("taskId", taskId, "status", "claimed");
@@ -472,6 +507,10 @@ public class WorkflowRuntimeService {
}
requireEnabledUser(connection, toUserId);
InstanceRow instance = lockInstance(connection, text(task.get("b_instance_id")));
if (!"running".equals(instance.status())) {
throw new BusinessException("流程已结束,不能转交待办");
}
requireAction(connection, instance.moduleId(), userId, "transfer");
WorkflowSql.update(connection, """
update dbo.wf_task set b_status = 'cancelled', b_updated_at = sysutcdatetime() where b_id = ?
@@ -542,10 +581,21 @@ public class WorkflowRuntimeService {
throw new BusinessException("该节点的待办已处理完,不能加签");
}
InstanceRow instance = lockInstance(connection, text(task.get("b_instance_id")));
if (!"running".equals(instance.status())) {
throw new BusinessException("流程已结束,不能加签");
}
requireAction(connection, instance.moduleId(), userId, "add_approver");
int seq = task.get("b_seq") instanceof Number number ? number.intValue() : 0;
for (String target : userIds) {
requireEnabledUser(connection, target);
Map<String, Object> activeTask = WorkflowSql.queryOne(connection, """
select top 1 b_id from dbo.wf_task
where b_node_run_id = ? and b_assignee_id = ?
and b_status in ('pending', 'claimed', 'waiting')
""", uuidParam(text(task.get("b_node_run_id"))), target);
if (activeTask != null) {
throw new BusinessException("用户「" + target + "」已经是该节点的活动审批人");
}
seq += 1;
WorkflowSql.update(connection, """
insert into dbo.wf_task
@@ -582,15 +632,22 @@ public class WorkflowRuntimeService {
try {
Map<String, Object> task = WorkflowSql.queryOne(connection, """
select b_id, b_instance_id, b_node_run_id, b_assignee_id, b_status
from dbo.wf_task where b_id = ?
from dbo.wf_task with (UPDLOCK, ROWLOCK) where b_id = ?
""", uuidParam(taskId));
if (task == null) {
throw new BusinessException("任务不存在");
}
String status = text(task.get("b_status"));
if (!"pending".equals(status) && !"claimed".equals(status)) {
throw new BusinessException("只有待处理任务可以催办");
}
InstanceRow instance = lockInstance(connection, text(task.get("b_instance_id")));
if (!"running".equals(instance.status())) {
throw new BusinessException("流程已结束,不能催办");
}
requireAction(connection, instance.moduleId(), userId, "urge");
writeAction(connection, instance, text(task.get("b_node_run_id")), taskId, "urge", userId,
text(task.get("b_status")), text(task.get("b_status")), comment, null, null);
status, status, comment, null, null);
writeOutbox(connection, "task.urged", "task", taskId,
Map.of("assigneeId", text(task.get("b_assignee_id"))));
connection.commit();
@@ -731,8 +788,12 @@ public class WorkflowRuntimeService {
String nodeRunId = recordNodeRun(connection, instance, nodeKey, "running", variables, snapshot, required);
Integer dueMinutes = dueMinutesOf(node);
List<String> taskUsers = resolved.userIds();
if ("single".equals(mode) && taskUsers.size() > 1) {
taskUsers = List.of(taskUsers.get(0));
}
int seq = 0;
for (String userId : resolved.userIds()) {
for (String userId : taskUsers) {
seq += 10;
// 串行会签:同一时刻只有第一名是 pending,其余 waiting,通过一位激活下一位(§9.6)
String status = "serial".equals(mode) && seq > 10 ? "waiting" : "pending";
@@ -1256,6 +1317,17 @@ public class WorkflowRuntimeService {
return fallback;
}
private boolean configFlag(Object json, String key) {
Object value = parseJsonMap(json).get(key);
if (value instanceof Boolean flag) {
return flag;
}
if (value instanceof Number number) {
return number.intValue() != 0;
}
return value instanceof String text && Boolean.parseBoolean(text.trim());
}
private Object uuidParam(String value) {
String text = text(value);
if (text.isEmpty()) {
@@ -1322,6 +1394,18 @@ public class WorkflowRuntimeService {
}
}
private boolean isUniqueViolation(SQLException exception) {
SQLException current = exception;
while (current != null) {
if (current.getErrorCode() == 2601 || current.getErrorCode() == 2627
|| "23000".equals(current.getSQLState())) {
return true;
}
current = current.getNextException();
}
return false;
}
/** 实例运行上下文(一次动作内复用,避免反复回表) */
private record InstanceRow(
String instanceId,
@@ -107,6 +107,22 @@ class DataSaveServiceTests {
verify(fixture.insertStatement).setObject(2, "");
}
@Test
void rejectsWorkflowRuntimeTablesBeforeLoadingMetadata() throws Exception {
JdbcFixture fixture = new JdbcFixture("b_id", "b_status");
DataSaveService service = fixture.createService();
assertThatThrownBy(() -> service.save(List.of(Map.of(
"table", "dbo.wf_task",
"key_field", "b_id",
"updates", List.of(Map.of("b_id", "T001", "b_status", "approved"))
)))).isInstanceOf(BusinessException.class)
.hasMessageContaining("Workflow 运行表必须通过流程服务写入");
verify(fixture.connection).rollback();
verify(fixture.connection, org.mockito.Mockito.never()).createStatement();
}
@Test
void sqlFailureRollsBackAndReportsOperationLocation() throws Exception {
JdbcFixture fixture = new JdbcFixture("b_id", "b_name");