+ * Use of this software is governed by the Commercial License Agreement + * obtained after purchasing a license from BladeX. + *
+ * 1. This software is for development use only under a valid license + * from BladeX. + *
+ * 2. Redistribution of this software's source code to any third party + * without a commercial license is strictly prohibited. + *
+ * 3. Licensees may copyright their own code but cannot use segments + * from this software for such purposes. Copyright of this software + * remains with BladeX. + *
+ * Using this software signifies agreement to this License, and the software + * must not be used for illegal purposes. + *
+ * THIS SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY. The author is + * not liable for any claims arising from secondary or illegal development. + *
+ * Author: Chill Zhuang (bladejava@qq.com) + */ +package org.springblade.process.pojo.entity; + +import com.baomidou.mybatisplus.annotation.*; +import com.fasterxml.jackson.annotation.JsonFormat; +import com.fasterxml.jackson.databind.annotation.JsonSerialize; +import com.fasterxml.jackson.databind.ser.std.ToStringSerializer; +import io.swagger.v3.oas.annotations.media.Schema; +import lombok.Data; +import org.springblade.core.tool.utils.DateUtil; +import org.springframework.format.annotation.DateTimeFormat; + +import java.io.Serial; +import java.io.Serializable; +import java.util.Date; + +/** + * 业务流程关联表 实体类 + * + * @author BladeX + * @since 2024-09-19 + */ +@Data +@TableName("blade_business_process") +@Schema(description = "BusinessProcess对象") +public class BusinessProcess implements Serializable { + + @Serial + private static final long serialVersionUID = 1L; + + /** + * 主键 + */ + @JsonSerialize(using = ToStringSerializer.class) + @Schema(description = "主键") + @TableId(value = "id", type = IdType.ASSIGN_ID) + private Long id; + + /** + * 业务id + */ + @Schema(description = "业务id") + private Long bizId; + /** + * 流程实例id + */ + @Schema(description = "流程实例id") + private String processInstanceId; + /** + * 流程类型 + */ + @Schema(description = "流程类型") + private String processType; + /** + * 文档编号 + */ + @Schema(description = "文档编号") + private String docCode; + /** + * 标题 + */ + @Schema(description = "标题") + private String subject; + /** + * 发起人id + */ + @Schema(description = "发起人id") + private Long promoterId; + /** + * 发起人名称 + */ + @Schema(description = "发起人名称") + private String promoterName; + /** + * 发起人登录名 + */ + @Schema(description = "发起人登录名") + private String promoterLoginName; + /** + * 提交时间 + */ + @Schema(description = "提交时间") + private Date submitTime; + /** + * 完成时间 + */ + @Schema(description = "完成时间") + private Date completeTime; + /** + * 当前节点id,多个用逗号拼接 + */ + @Schema(description = "当前节点id,多个用逗号拼接") + private String currentNodeIds; + /** + * 当前节点名称,多个用逗号拼接 + */ + @Schema(description = "当前节点名称,多个用逗号拼接") + private String currentNodeNames; + /** + * 当前处理人,多个用逗号拼接 + */ + @Schema(description = "当前处理人,多个用逗号拼接") + private String currentHandlers; + /** + * 接收时间 + */ + @Schema(description = "接收时间") + private Date receiveTime; + /** + * 是否已完成(0:未完成, 1:已完成) + */ + @Schema(description = "是否已完成") + private Integer isCompleted; + /** + * 审批状态 + */ + @Schema(description = "审批状态") + private String approveStatus; + /** + * 租户ID + */ + @Schema(description = "租户ID") + private String tenantId; + /** + * 创建时间 + */ + @DateTimeFormat(pattern = DateUtil.PATTERN_DATETIME) + @JsonFormat(pattern = DateUtil.PATTERN_DATE) + @Schema(description = "创建时间", hidden = true) + @TableField(fill = FieldFill.INSERT) + private Date createTime; + /** + * 更新时间 + */ + @DateTimeFormat(pattern = DateUtil.PATTERN_DATETIME) + @JsonFormat(pattern = DateUtil.PATTERN_DATE) + @Schema(description = "更新时间", hidden = true) + @TableField(fill = FieldFill.INSERT_UPDATE) + private Date updateTime; +} diff --git a/blade-service-api/blade-process-api/src/main/java/org/springblade/process/pojo/enums/ApproveStatusEnum.java b/blade-service-api/blade-process-api/src/main/java/org/springblade/process/pojo/enums/ApproveStatusEnum.java new file mode 100644 index 0000000..9b9aafe --- /dev/null +++ b/blade-service-api/blade-process-api/src/main/java/org/springblade/process/pojo/enums/ApproveStatusEnum.java @@ -0,0 +1,165 @@ +/** + * BladeX Commercial License Agreement + * Copyright (c) 2018-2099, https://bladex.cn. All rights reserved. + *
+ * Use of this software is governed by the Commercial License Agreement + * obtained after purchasing a license from BladeX. + *
+ * 1. This software is for development use only under a valid license + * from BladeX. + *
+ * 2. Redistribution of this software's source code to any third party + * without a commercial license is strictly prohibited. + *
+ * 3. Licensees may copyright their own code but cannot use segments + * from this software for such purposes. Copyright of this software + * remains with BladeX. + *
+ * Using this software signifies agreement to this License, and the software + * must not be used for illegal purposes. + *
+ * THIS SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY. The author is + * not liable for any claims arising from secondary or illegal development. + *
+ * Author: Chill Zhuang (bladejava@qq.com)
+ */
+package org.springblade.process.pojo.enums;
+
+import lombok.AllArgsConstructor;
+import lombok.Getter;
+
+import java.util.List;
+import java.util.Objects;
+
+/**
+ * 审批状态
+ *
+ * @author LiuXinjie
+ * @apiNote 合同审批状态
+ */
+@Getter
+@AllArgsConstructor
+public enum ApproveStatusEnum {
+
+ /**
+ * 默认编号
+ */
+ DRAFT("draft", "草稿"), //可提交
+ APPROVING("approval", "审批中"),
+ APPROVED("pass", "审批通过"),
+ REJECTED("reject", "审批驳回"), //通用的流程 驳回可编辑
+ REVOCATION("revocation", "已撤回"), //可重新提交
+ ABANDON("abandon", "废弃"),
+ ;
+
+ final String value;
+ final String text;
+
+ public boolean match(String value){
+ return this.value.equals(value);
+ }
+
+ public static String getValueByText(String text) {
+ if (text == null) {
+ return null;
+ }
+ for (ApproveStatusEnum item : values()) {
+ if (Objects.equals(item.getText(), text)) {
+ return item.getValue();
+ }
+ }
+ return null;
+ }
+
+ public static String getTextByValue(String value) {
+ if (value == null) {
+ return null;
+ }
+ for (ApproveStatusEnum item : values()) {
+ if (Objects.equals(item.getValue(), value)) {
+ return item.getText();
+ }
+ }
+ return null;
+ }
+
+ /**
+ * 是否可以撤回
+ *
+ * @param value
+ * @return
+ */
+ public static boolean canRevoke(String value) {
+ return APPROVING.getValue().equals(value);
+ }
+
+ /**
+ * 能不能删除审批流
+ * @param value
+ * @return
+ */
+ public static boolean canDelAuditFlow(String value){
+ return REJECTED.getValue().equals(value) || REVOCATION.getValue().equals(value);
+ }
+
+ /**
+ * 驳回或撤回
+ * @return
+ */
+ public static boolean rejectedOrRevocation(String approveStatus) {
+ return canDelAuditFlow(approveStatus);
+ }
+
+ /**
+ * 能不能删除数据
+ * @param value
+ * @return
+ */
+ public static boolean canDeleteData(String value){
+ return ABANDON.getValue().equals(value) || DRAFT.getValue().equals(value) || canDelAuditFlow(value);
+ }
+
+ public static String getNameStr(String value) {
+ for (ApproveStatusEnum state : values()) {
+ if (state.value.equals(value)) {
+ return state.text;
+ }
+ }
+ return null;
+ }
+
+
+ /**
+ * 是否可以编辑表单
+ *
+ * @param value
+ * @return
+ */
+ public static boolean canEdit(String value) {
+ return DRAFT.getValue().equals(value) || REJECTED.getValue().equals(value) || REVOCATION.getValue().equals(value);
+ }
+
+ /**
+ * 获取可驳回状态
+ *
+ * @return
+ */
+ public static List
+ * callbackParam 仅保留 MK 原始回调参数, + * 其余字段为 openapi 在处理过程中补充的上下文参数。 + *
+ *+ * 设计目的: + * 1. 避免把内部推导字段继续堆到 MK 原始回调 DTO 上; + * 2. 对外保留原始回调对象,便于排查问题、记录日志和后续扩展; + * 3. 通过代理 getter 尽量兼容原来直接读取 DTO 字段的使用习惯,降低老流程和后续分支合并成本。 + *
+ * + * @author bfhuange + * @date 2026/4/9 + */ +@Data +@Builder +public class ProcessOperationContext implements Serializable { + @Serial + private static final long serialVersionUID = 1L; + + /** + * MK 原始回调参数 + */ + private Api4MKProcessApprovalDTO callbackParam; + + /** + * 流程类型 + */ + private String processType; + + /** + * 业务审批状态。 + * 这是系统内部按事件语义统一补充的状态,不属于 MK 原始回调参数。 + */ + private String approveStatus; + + /** + * 是否流程已完成。 + * 这是系统内部按事件语义统一补充的状态,不属于 MK 原始回调参数。 + */ + private boolean complete; + + /** + * 是否异步刷新当前处理人。 + * 用于控制当前处理人更新是走同步刷新还是异步调度任务。 + */ + private boolean async; + + @JsonIgnore + @JSONField(serialize = false) + public String getProcessInstanceId() { + return callbackParam == null ? null : callbackParam.getProcessInstanceId(); + } + + @JsonIgnore + @JSONField(serialize = false) + public String getFormInstanceId() { + return callbackParam == null ? null : callbackParam.getFormInstanceId(); + } + + @JsonIgnore + @JSONField(serialize = false) + public String getTemplateId() { + return callbackParam == null ? null : callbackParam.getTemplateId(); + } + + @JsonIgnore + @JSONField(serialize = false) + public String getTemplateCode() { + return callbackParam == null ? null : callbackParam.getTemplateCode(); + } + + @JsonIgnore + @JSONField(serialize = false) + public String getProcessStatus() { + return callbackParam == null ? null : callbackParam.getProcessStatus(); + } + + @JsonIgnore + @JSONField(serialize = false) + public String getApplicantLoginName() { + return callbackParam == null ? null : callbackParam.getApplicantLoginName(); + } + + @JsonIgnore + @JSONField(serialize = false) + public String getRejectNodeId() { + return callbackParam == null ? null : callbackParam.getRejectNodeId(); + } + + @JsonIgnore + @JSONField(serialize = false) + public String getCurrentNodeId() { + return callbackParam == null ? null : callbackParam.getCurrentNodeId(); + } + + @JsonIgnore + @JSONField(serialize = false) + public String getCurrentNodeNumber() { + return callbackParam == null ? null : callbackParam.getCurrentNodeNumber(); + } + + @JsonIgnore + @JSONField(serialize = false) + public String getOperation() { + return callbackParam == null ? null : callbackParam.getOperation(); + } + + @JsonIgnore + @JSONField(serialize = false) + public String getOperationName() { + return callbackParam == null ? null : callbackParam.getOperationName(); + } + + @JsonIgnore + @JSONField(serialize = false) + public String getApprovalOpinion() { + return callbackParam == null ? null : callbackParam.getApprovalOpinion(); + } + + @JsonIgnore + @JSONField(serialize = false) + public String getOperatorLoginName() { + return callbackParam == null ? null : callbackParam.getOperatorLoginName(); + } +} diff --git a/blade-service/blade-openapi/src/main/java/org/springblade/openapi/mk/support/base/ProcessOperationHandler.java b/blade-service/blade-openapi/src/main/java/org/springblade/openapi/mk/support/base/ProcessOperationHandler.java new file mode 100644 index 0000000..c914acd --- /dev/null +++ b/blade-service/blade-openapi/src/main/java/org/springblade/openapi/mk/support/base/ProcessOperationHandler.java @@ -0,0 +1,45 @@ +package org.springblade.openapi.mk.support.base; + +/** + * 流程操作处理器 + * @author bfhuange + * @since 2024/11/25 + */ +public interface ProcessOperationHandler { + + /** + * 提交 + * @param param + */ + void submit(ProcessOperationContext param); + + /** + * 审批结束 + * @param param + */ + void approveFinish(ProcessOperationContext param); + + /** + * 审批同意 + * @param param + */ + void approvePass(ProcessOperationContext param); + + /** + * 审批拒绝 + * @param param + */ + void approveReject(ProcessOperationContext param); + + /** + * 审批撤销 + * @param param + */ + void approveRevoke(ProcessOperationContext param); + + /** + * 审批废弃 + * @param param + */ + void approveAbandon(ProcessOperationContext param); +} diff --git a/blade-service/blade-openapi/src/main/java/org/springblade/openapi/mk/support/handler/AsyncService.java b/blade-service/blade-openapi/src/main/java/org/springblade/openapi/mk/support/handler/AsyncService.java new file mode 100644 index 0000000..87250c5 --- /dev/null +++ b/blade-service/blade-openapi/src/main/java/org/springblade/openapi/mk/support/handler/AsyncService.java @@ -0,0 +1,113 @@ +package org.springblade.openapi.mk.support.handler; + +import lombok.extern.slf4j.Slf4j; +import org.dromara.dynamictp.core.DtpRegistry; +import org.dromara.dynamictp.core.aware.TaskEnhanceAware; +import org.springblade.core.log.exception.ServiceException; +import org.springblade.openapi.mk.config.AsyncExecutorProperties; +import org.springframework.boot.context.event.ApplicationReadyEvent; +import org.springframework.context.event.EventListener; +import org.springframework.core.Ordered; +import org.springframework.core.annotation.Order; +import org.springframework.stereotype.Service; + +import java.util.concurrent.Executor; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.TimeUnit; + +/** + * @author bfhuange + * @date 2024/9/20 + */ +@Slf4j +@Service +public class AsyncService { + private final AsyncExecutorProperties asyncExecutorProperties; + private Executor workerExecutor; + private ScheduledExecutorService schedulerExecutor; + + public AsyncService(AsyncExecutorProperties asyncExecutorProperties) { + this.asyncExecutorProperties = asyncExecutorProperties; + } + + /** + * 启动完成后预加载并校验线程池配置,避免等到第一次真正执行任务时才发现线程池缺失或类型配置错误。 + */ + @Order(Ordered.HIGHEST_PRECEDENCE) + @EventListener(ApplicationReadyEvent.class) + public void initExecutors() { + this.workerExecutor = resolveWorkerExecutor(); + this.schedulerExecutor = resolveSchedulerExecutor(); + log.info("当前处理人刷新异步线程池初始化完成,workerExecutorName:{},schedulerExecutorName:{}", + asyncExecutorProperties.getWorkerExecutorName(), asyncExecutorProperties.getSchedulerExecutorName()); + } + + /** + * 立即异步执行 + * @param runnable 任务 + */ + public void execute(Runnable runnable) { + getWorkerExecutor().execute(runnable); + } + + /** + * 延迟执行指定毫秒数。 + *+ * 这里改为使用 ScheduledDtpExecutor 做真正的定时调度, + * 避免再通过线程池线程 sleep 的方式占用工作线程,导致真正的业务任务迟迟无法启动。 + *
+ * + * @param delayMillis 延迟毫秒数 + * @param runnable 任务 + */ + public void delayExecute(long delayMillis, Runnable runnable) { + if (delayMillis <= 0) { + execute(runnable); + return; + } + Runnable dispatchRunnable = wrapWithConfiguredTaskWrappers( + asyncExecutorProperties.getSchedulerExecutorName(), + () -> execute(runnable) + ); + getSchedulerExecutor().schedule(dispatchRunnable, delayMillis, TimeUnit.MILLISECONDS); + } + + private Executor getWorkerExecutor() { + return workerExecutor != null ? workerExecutor : resolveWorkerExecutor(); + } + + private ScheduledExecutorService getSchedulerExecutor() { + return schedulerExecutor != null ? schedulerExecutor : resolveSchedulerExecutor(); + } + + private Executor resolveWorkerExecutor() { + return DtpRegistry.getExecutor(asyncExecutorProperties.getWorkerExecutorName()); + } + + private ScheduledExecutorService resolveSchedulerExecutor() { + String schedulerExecutorName = asyncExecutorProperties.getSchedulerExecutorName(); + Executor executor = DtpRegistry.getExecutor(schedulerExecutorName); + if (executor instanceof ScheduledExecutorService scheduledExecutorService) { + return scheduledExecutorService; + } + String message = "线程池未按 ScheduledExecutorService 注册,name: " + schedulerExecutorName; + log.error(message); + throw new ServiceException(message); + } + + /** + * 按线程池已配置的 task wrappers 手动包装任务。 + *+ * 当前使用的 dynamic-tp 版本下,ScheduledDtpExecutor 对 taskWrapper 的透传存在缺口, + * 这里直接读取线程池上已生效的 wrappers,按框架默认增强链顺序主动包装一次, + * 这样既能复用现有配置,又避免手写 mdc 透传逻辑与框架实现产生偏差。 + *
+ */ + private Runnable wrapWithConfiguredTaskWrappers(String executorName, Runnable runnable) { + Executor executor = DtpRegistry.getExecutor(executorName); + if (executor instanceof TaskEnhanceAware taskEnhanceAware) { + return taskEnhanceAware.getEnhancedTask(runnable, taskEnhanceAware.getTaskWrappers()); + } + return runnable; + } +} diff --git a/blade-service/blade-openapi/src/main/java/org/springblade/openapi/mk/support/handler/ProcessCurrentHandlerRefreshService.java b/blade-service/blade-openapi/src/main/java/org/springblade/openapi/mk/support/handler/ProcessCurrentHandlerRefreshService.java new file mode 100644 index 0000000..057f541 --- /dev/null +++ b/blade-service/blade-openapi/src/main/java/org/springblade/openapi/mk/support/handler/ProcessCurrentHandlerRefreshService.java @@ -0,0 +1,903 @@ +package org.springblade.openapi.mk.support.handler; + +import cn.hutool.core.util.IdUtil; +import com.alibaba.fastjson2.JSON; +import lombok.extern.slf4j.Slf4j; +import org.apache.commons.lang3.StringUtils; +import org.redisson.api.RLock; +import org.redisson.api.RMapCache; +import org.springblade.core.redis.cache.BladeRedis; +import org.springblade.core.redis.lock.RedisLockClient; +import org.springblade.core.tool.api.FR; +import org.springblade.openapi.mk.config.CurrentHandlerRefreshProperties; +import org.springblade.openapi.mk.constant.ProcessLockKeyConstant; +import org.springblade.openapi.mk.pojo.enums.ProcessCallbackType; +import org.springblade.openapi.mk.support.base.AbstractProcessOperationHandler; +import org.springblade.openapi.mk.support.base.ProcessOperationContext; +import org.springblade.process.feign.IBusinessProcessClient; +import org.springblade.process.pojo.dto.BusinessProcessCurrentHandlerRefreshDTO; +import org.springblade.process.pojo.vo.BusinessProcessVO; +import org.springframework.beans.factory.ObjectProvider; +import org.springframework.boot.context.event.ApplicationReadyEvent; +import org.springframework.context.event.EventListener; +import org.springframework.core.annotation.Order; +import org.springframework.stereotype.Service; + +import java.time.Duration; +import java.util.Collections; +import java.util.List; +import java.util.Map; +import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.TimeUnit; +import java.util.stream.Collectors; +import java.util.stream.Stream; + +/** + * 当前处理人刷新调度服务。 + *+ * 背景: + * 流程引擎回调业务系统时,流程往往还没有真正流转到下一个激活节点, + * 此时立即查询当前节点/当前处理人,拿到的仍可能是上一节点的旧结果。 + * 因此这里不再依赖一次性的固定延迟,而是改成“按流程实例维度入队 + 固定间隔轮询刷新”的调度模型。 + *
+ *+ * 整体流程: + * 1. openapi 收到流程事件后,先同步更新业务流程状态; + * 2. 如果当前事件要求异步刷新当前处理人,则调用 {@link #enqueue(ProcessOperationContext)} 写入刷新任务; + * 3. 任务以流程实例 id 为唯一主记录保存在 Redis,记录期望版本、执行版本、基线快照、最近回调参数、重试次数等信息; + * 4. 同一个流程实例只保留一条主任务记录,新的回调不会重复创建任务,只会提升 {@code desiredVersion} 并覆盖最近一次回调参数; + * 5. 等待中的流程实例 id 会放入 Redis ZSet,score 为下次重试时间,用于按时间顺序派工; + * 6. 调度器 {@link #tryDispatch()} 会在集群范围内抢占派工锁,按配置的最大 worker 数拉起异步 worker; + * 7. worker 执行时调用 system 侧“只刷新当前节点/当前处理人”接口,并比较“当前节点 + 当前处理人”快照是否相对基线发生变化; + * 8. 如果快照未变化,说明流程大概率还没流转完成,则按固定间隔重新入队重试; + * 9. 如果快照发生变化,则回调对应业务处理器 {@link AbstractProcessOperationHandler#handleCurrentHandlerRefresh(ProcessOperationContext, BusinessProcessVO)}; + * 10. 若执行期间又收到同一流程的新回调,则旧版本执行完后会把最新快照提升为新基线,并重新排到队尾,避免同一流程长期占用 worker; + * 11. 当达到最大重试次数后,任务进入失败态并保留一段时间,便于排查; + * 12. 成功完成的任务进入完成态并短期保留,随后自动过期。 + *
+ *+ * 集群与并发约束: + * 1. 流程实例级别使用分布式锁,保证同一流程实例的任务状态变更串行化; + * 2. 派工使用全局分布式锁,保证多个实例不会同时超发 worker; + * 3. 活跃 worker 数通过 Redis 租约控制,服务异常中断后,租约超时即可视为 worker 失活; + * 4. 运行中的任务会持续更新心跳,如果服务升级、中断或线程异常退出,超时恢复逻辑会把任务重新转回等待态; + * 5. 启动时不会全量恢复运行中任务,避免在集群环境中误伤其他实例上仍在执行的任务。 + *
+ *+ * 成功判定规则: + * 不再区分终态/非终态,也不依赖回调里传入的 complete true/false 单独判定是否成功, + * 统一以“当前节点变化 + 当前处理人变化后的最新快照”是否相对基线发生变化作为刷新成功依据。 + *
+ * + * @author bfhuange + * @date 2026/4/9 + */ +@Slf4j +@Service +public class ProcessCurrentHandlerRefreshService { + + + private final AsyncService asyncService; + private final BladeRedis bladeRedis; + private final RedisLockClient redisLockClient; + private final IBusinessProcessClient processClient; + private final CurrentHandlerRefreshProperties refreshProperties; + private final ObjectProvider+ * 这里只针对提交事件生效,不能推广到审批通过/会签等场景: + * 会签节点在部分人审批完成后,当前节点可能仍然不变,但当前处理人已经发生变化, + * 此时应当允许按“快照变化”判定成功,而不是继续等待节点变化。 + *
+ * + * @param callbackParam 回调参数 + * @param businessProcessVO 最新流程快照 + * @return 是否继续等待下一节点 + */ + private boolean shouldWaitForNextNode(ProcessOperationContext callbackParam, BusinessProcessVO businessProcessVO) { + if (callbackParam == null || businessProcessVO == null) { + return false; + } + if (ProcessCallbackType.SUBMIT != ProcessCallbackType.getCallbackType(callbackParam.getOperation())) { + return false; + } + String callbackNodeId = callbackParam.getCurrentNodeId(); + if (StringUtils.isBlank(callbackNodeId)) { + return false; + } + return containsCsvValue(businessProcessVO.getCurrentNodeIds(), callbackNodeId); + } + + /** + * 统一规范逗号拼接字段,避免比较时顺序影响结果 + * + * @param value 原始值 + * @return 规范化后的字符串 + */ + private String normalizeCsv(String value) { + if (StringUtils.isBlank(value)) { + return ""; + } + return Stream.of(value.split(",")) + .map(String::trim) + .filter(StringUtils::isNotBlank) + .distinct() + .sorted() + .collect(Collectors.joining(",")); + } + + /** + * 判断逗号分隔字段中是否包含指定值 + * + * @param csv 逗号分隔字段 + * @param target 目标值 + * @return 是否包含 + */ + private boolean containsCsvValue(String csv, String target) { + if (StringUtils.isAnyBlank(csv, target)) { + return false; + } + return Stream.of(csv.split(",")) + .map(String::trim) + .anyMatch(target::equals); + } + + /** + * 判断任务token是否仍然有效 + * + * @param processInstanceId 流程实例id + * @param runToken 运行token + * @return 是否匹配 + */ + private boolean isTaskTokenMatched(String processInstanceId, String runToken) { + ProcessCurrentHandlerRefreshTask latestTask = getTask(processInstanceId); + return latestTask != null + && TaskState.STATE_RUNNING.equals(latestTask.getState()) + && StringUtils.equals(runToken, latestTask.getRunToken()); + } + + /** + * 更新任务心跳,表示当前worker仍然存活 + * + * @param processInstanceId 流程实例id + * @param runToken 运行token + */ + private void updateHeartbeat(String processInstanceId, String runToken) { + RLock lock = getProcessLock(processInstanceId); + boolean locked = false; + try { + locked = lock.tryLock(refreshProperties.getLockWaitSeconds(), TimeUnit.SECONDS); + if (!locked) { + return; + } + ProcessCurrentHandlerRefreshTask task = getTask(processInstanceId); + if (task == null || !StringUtils.equals(runToken, task.getRunToken())) { + return; + } + task.setHeartbeatAt(System.currentTimeMillis()); + saveTask(task); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + log.error("更新刷新任务心跳被中断,流程实例id:{}", processInstanceId, e); + } finally { + unlock(lock); + } + } + + /** + * 按默认方式保存任务 + * + * @param task 任务 + */ + private void saveTask(ProcessCurrentHandlerRefreshTask task) { + bladeRedis.getStringRedisTemplate().opsForValue().set(getTaskKey(task.getProcessInstanceId()), JSON.toJSONString(task)); + } + + /** + * 按TTL保存任务 + * + * @param task 任务 + * @param ttl TTL + */ + private void saveTask(ProcessCurrentHandlerRefreshTask task, Duration ttl) { + bladeRedis.getStringRedisTemplate().opsForValue() + .set(getTaskKey(task.getProcessInstanceId()), JSON.toJSONString(task), ttl); + } + + /** + * 获取任务 + * + * @param processInstanceId 流程实例id + * @return 任务 + */ + private ProcessCurrentHandlerRefreshTask getTask(String processInstanceId) { + return getTaskByKey(getTaskKey(processInstanceId)); + } + + /** + * 通过key读取任务 + * + * @param taskKey 任务key + * @return 任务 + */ + private ProcessCurrentHandlerRefreshTask getTaskByKey(String taskKey) { + String content = bladeRedis.getStringRedisTemplate().opsForValue().get(taskKey); + if (StringUtils.isBlank(content)) { + return null; + } + return JSON.parseObject(content, ProcessCurrentHandlerRefreshTask.class); + } + + /** + * 放入等待队列 + * + * @param processInstanceId 流程实例id + * @param nextRetryAt 下次执行时间 + */ + private void putWaitingTask(String processInstanceId, long nextRetryAt) { + bladeRedis.getStringRedisTemplate().opsForZSet().add(ProcessLockKeyConstant.WAITING_KEY, processInstanceId, nextRetryAt); + } + + /** + * 移除等待队列中的任务 + * + * @param processInstanceId 流程实例id + */ + private void removeWaitingTask(String processInstanceId) { + bladeRedis.getStringRedisTemplate().opsForZSet().remove(ProcessLockKeyConstant.WAITING_KEY, processInstanceId); + } + + /** + * 是否存在等待任务 + * + * @return 是否存在 + */ + private boolean hasWaitingTask() { + Long size = bladeRedis.getStringRedisTemplate().opsForZSet().zCard(ProcessLockKeyConstant.WAITING_KEY); + return size != null && size > 0; + } + + /** + * 获取worker租约map + * + * @return worker租约map + */ + private RMapCache