From d429a09531aaf250db5aa49f16d353ad66af4e98 Mon Sep 17 00:00:00 2001 From: liyongde <1419499670@qq.com> Date: Thu, 16 Jul 2026 19:07:09 +0800 Subject: [PATCH] Refactor document workflow and related UI state --- .../execute/biz/api/TaskCommonApi.java | 27 ++- .../biz/dto/TaskStatusCallApiReqDTO.java | 4 +- .../execute/biz/vo/AcsApplyActionRespVO.java | 29 +++ .../framework/execute/core/AbstractTask.java | 2 +- .../framework/execute/core/TaskFactory.java | 6 +- nl-module-lms/nl-module-lms-api/pom.xml | 5 +- nl-module-lms/nl-module-lms-server/pom.xml | 4 - .../java/cn/code/nl/module/lms/DemoApi.java | 20 -- .../cn/code/nl/module/lms/demo/DemoApi.java | 80 +++++++ .../code/nl/module/lms/demo/DemoRespVO.java | 13 ++ .../nl/module/lms/{ => demo}/DemoTask.java | 7 +- .../nl/module/task/dto/AcsFeedbackReqDTO.java | 7 +- .../enums/AcsBusinessOperationTypeEnum.java | 47 ++++ .../module/task/enums/ErrorCodeConstants.java | 1 + .../task/enums/TaskOperationTypeEnum.java | 4 - .../transporttask/AcsFeedbackController.java | 9 +- ...TransportTaskBusinessOperationManager.java | 196 +++++++++++++++++ .../manage/TransportTaskOperationManager.java | 181 +++++++++++++++ .../TransportTaskFeedbackService.java | 9 + .../TransportTaskFeedbackServiceImpl.java | 113 +++++++--- .../transporttask/TransportTaskService.java | 201 ++++++++--------- .../TransportTaskServiceImpl.java | 208 +----------------- 22 files changed, 795 insertions(+), 378 deletions(-) create mode 100644 nl-framework/nl-spring-boot-starter-execute/src/main/java/cn/code/nl/framework/execute/biz/vo/AcsApplyActionRespVO.java delete mode 100644 nl-module-lms/nl-module-lms-server/src/main/java/cn/code/nl/module/lms/DemoApi.java create mode 100644 nl-module-lms/nl-module-lms-server/src/main/java/cn/code/nl/module/lms/demo/DemoApi.java create mode 100644 nl-module-lms/nl-module-lms-server/src/main/java/cn/code/nl/module/lms/demo/DemoRespVO.java rename nl-module-lms/nl-module-lms-server/src/main/java/cn/code/nl/module/lms/{ => demo}/DemoTask.java (77%) create mode 100644 nl-module-task/nl-module-task-api/src/main/java/cn/code/nl/module/task/enums/AcsBusinessOperationTypeEnum.java create mode 100644 nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/manage/TransportTaskBusinessOperationManager.java create mode 100644 nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/manage/TransportTaskOperationManager.java diff --git a/nl-framework/nl-spring-boot-starter-execute/src/main/java/cn/code/nl/framework/execute/biz/api/TaskCommonApi.java b/nl-framework/nl-spring-boot-starter-execute/src/main/java/cn/code/nl/framework/execute/biz/api/TaskCommonApi.java index f393fccb..66d9a473 100644 --- a/nl-framework/nl-spring-boot-starter-execute/src/main/java/cn/code/nl/framework/execute/biz/api/TaskCommonApi.java +++ b/nl-framework/nl-spring-boot-starter-execute/src/main/java/cn/code/nl/framework/execute/biz/api/TaskCommonApi.java @@ -1,6 +1,10 @@ package cn.code.nl.framework.execute.biz.api; +import cn.code.nl.framework.common.pojo.CommonResult; import cn.code.nl.framework.execute.biz.dto.TaskStatusCallApiReqDTO; +import cn.code.nl.framework.execute.biz.vo.AcsApplyActionRespVO; +import io.swagger.v3.oas.annotations.Operation; +import io.swagger.v3.oas.annotations.media.Schema; import org.springframework.web.bind.annotation.PostMapping; /** @@ -10,6 +14,27 @@ import org.springframework.web.bind.annotation.PostMapping; */ public interface TaskCommonApi { + @Operation(summary = "请求取货") @PostMapping("/do-handle-picked") - void doHandlePicked(TaskStatusCallApiReqDTO taskStatusCallApiReqDTO); + CommonResult doHandlePicked(TaskStatusCallApiReqDTO taskStatusCallApiReqDTO); + + @Operation(summary = "二次请求") + @PostMapping("/do-handle-apply-again") + CommonResult doHandleApplyAgain(TaskStatusCallApiReqDTO taskStatusCallApiReqDTO); + + @Operation(summary = "处理请求放货") + @PostMapping("/do-handle-request-release") + CommonResult doHandleRequestRelease(TaskStatusCallApiReqDTO taskStatusCallApiReqDTO); + + @Operation(summary = "处理请求取货") + @PostMapping("/do-handle-request-pick") + CommonResult doHandleRequestPick(TaskStatusCallApiReqDTO taskStatusCallApiReqDTO); + + @Operation(summary = "处理请求离开") + @PostMapping("/do-handle-request-leave") + CommonResult doHandleRequestLeave(TaskStatusCallApiReqDTO taskStatusCallApiReqDTO); + + @Operation(summary = "处理请求进入") + @PostMapping("/do-handle-request-enter") + CommonResult doHandleRequestEnter(TaskStatusCallApiReqDTO taskStatusCallApiReqDTO); } diff --git a/nl-framework/nl-spring-boot-starter-execute/src/main/java/cn/code/nl/framework/execute/biz/dto/TaskStatusCallApiReqDTO.java b/nl-framework/nl-spring-boot-starter-execute/src/main/java/cn/code/nl/framework/execute/biz/dto/TaskStatusCallApiReqDTO.java index fe2e4f7d..b3085140 100644 --- a/nl-framework/nl-spring-boot-starter-execute/src/main/java/cn/code/nl/framework/execute/biz/dto/TaskStatusCallApiReqDTO.java +++ b/nl-framework/nl-spring-boot-starter-execute/src/main/java/cn/code/nl/framework/execute/biz/dto/TaskStatusCallApiReqDTO.java @@ -18,7 +18,7 @@ public class TaskStatusCallApiReqDTO implements Serializable { /** 任务编码 */ private String taskCode; - /** 回调状态:当前为 PICKED(65) */ + /** 回调状态 */ private String status; /** 业务归属服务 */ @@ -30,7 +30,7 @@ public class TaskStatusCallApiReqDTO implements Serializable { /** 业务侧标识 */ private String bizId; - /** 业务回调处理器编码(LMS/WMS 内部路由) */ + /** 业务回调处理器编码(LMS/WMS 内部路由): 子类的beanName */ private String handleCode; /** 扩展数据 */ diff --git a/nl-framework/nl-spring-boot-starter-execute/src/main/java/cn/code/nl/framework/execute/biz/vo/AcsApplyActionRespVO.java b/nl-framework/nl-spring-boot-starter-execute/src/main/java/cn/code/nl/framework/execute/biz/vo/AcsApplyActionRespVO.java new file mode 100644 index 00000000..5096efe6 --- /dev/null +++ b/nl-framework/nl-spring-boot-starter-execute/src/main/java/cn/code/nl/framework/execute/biz/vo/AcsApplyActionRespVO.java @@ -0,0 +1,29 @@ +package cn.code.nl.framework.execute.biz.vo; + +import com.alibaba.fastjson.JSONObject; +import lombok.Data; + +/** + * ACS 请求行为的VO + * + * @Author: liyongde + * @Date: 2026/7/16 16:23 + */ +@Data +public class AcsApplyActionRespVO { + + /** + * 任务ID + */ + private Long taskId; + + /** + * 任务编码 + */ + private String taskCode; + + /** + * 数据 + */ + private JSONObject data; +} diff --git a/nl-framework/nl-spring-boot-starter-execute/src/main/java/cn/code/nl/framework/execute/core/AbstractTask.java b/nl-framework/nl-spring-boot-starter-execute/src/main/java/cn/code/nl/framework/execute/core/AbstractTask.java index 2e3bceb7..4480a102 100644 --- a/nl-framework/nl-spring-boot-starter-execute/src/main/java/cn/code/nl/framework/execute/core/AbstractTask.java +++ b/nl-framework/nl-spring-boot-starter-execute/src/main/java/cn/code/nl/framework/execute/core/AbstractTask.java @@ -31,7 +31,7 @@ public abstract class AbstractTask { return null; } - /** 请求放货*/ + /** 请求放货 */ public Boolean requestPutAway(TaskExecuteDTO taskExecuteDTO) { return true; } diff --git a/nl-framework/nl-spring-boot-starter-execute/src/main/java/cn/code/nl/framework/execute/core/TaskFactory.java b/nl-framework/nl-spring-boot-starter-execute/src/main/java/cn/code/nl/framework/execute/core/TaskFactory.java index 57c987c4..4b16f8c9 100644 --- a/nl-framework/nl-spring-boot-starter-execute/src/main/java/cn/code/nl/framework/execute/core/TaskFactory.java +++ b/nl-framework/nl-spring-boot-starter-execute/src/main/java/cn/code/nl/framework/execute/core/TaskFactory.java @@ -30,10 +30,10 @@ public class TaskFactory implements BeanPostProcessor { return bean; } - public AbstractTask getTask(String taskType) { - if (taskType == null) { + public AbstractTask getTask(String handleCode) { + if (handleCode == null) { return null; } - return taskMap.get(taskType); + return taskMap.get(handleCode); } } diff --git a/nl-module-lms/nl-module-lms-api/pom.xml b/nl-module-lms/nl-module-lms-api/pom.xml index 7e232dc1..4cfe7411 100644 --- a/nl-module-lms/nl-module-lms-api/pom.xml +++ b/nl-module-lms/nl-module-lms-api/pom.xml @@ -22,7 +22,10 @@ cn.nl.cloud nl-common - + + cn.nl.cloud + nl-spring-boot-starter-execute + org.springdoc diff --git a/nl-module-lms/nl-module-lms-server/pom.xml b/nl-module-lms/nl-module-lms-server/pom.xml index 8f91845d..562d3f5f 100644 --- a/nl-module-lms/nl-module-lms-server/pom.xml +++ b/nl-module-lms/nl-module-lms-server/pom.xml @@ -32,10 +32,6 @@ nl-module-lms-api ${revision} - - cn.nl.cloud - nl-spring-boot-starter-execute - diff --git a/nl-module-lms/nl-module-lms-server/src/main/java/cn/code/nl/module/lms/DemoApi.java b/nl-module-lms/nl-module-lms-server/src/main/java/cn/code/nl/module/lms/DemoApi.java deleted file mode 100644 index fcc5b893..00000000 --- a/nl-module-lms/nl-module-lms-server/src/main/java/cn/code/nl/module/lms/DemoApi.java +++ /dev/null @@ -1,20 +0,0 @@ -package cn.code.nl.module.lms; - -import cn.code.nl.framework.execute.biz.api.lms.LmsTaskCommonApi; -import cn.code.nl.framework.execute.biz.dto.TaskStatusCallApiReqDTO; -import org.springframework.validation.annotation.Validated; -import org.springframework.web.bind.annotation.RestController; - -/** - * - * @Author: liyongde - * @Date: 2026/7/15 15:25 - */ -@RestController // 提供 RESTful API 接口,给 Feign 调用 -@Validated -public class DemoApi implements LmsTaskCommonApi { - @Override - public void doHandlePicked(TaskStatusCallApiReqDTO taskStatusCallApiReqDTO) { - - } -} diff --git a/nl-module-lms/nl-module-lms-server/src/main/java/cn/code/nl/module/lms/demo/DemoApi.java b/nl-module-lms/nl-module-lms-server/src/main/java/cn/code/nl/module/lms/demo/DemoApi.java new file mode 100644 index 00000000..a1b44f70 --- /dev/null +++ b/nl-module-lms/nl-module-lms-server/src/main/java/cn/code/nl/module/lms/demo/DemoApi.java @@ -0,0 +1,80 @@ +package cn.code.nl.module.lms.demo; + +import cn.code.nl.framework.common.pojo.CommonResult; +import cn.code.nl.framework.execute.biz.api.lms.LmsTaskCommonApi; +import cn.code.nl.framework.execute.biz.dto.TaskStatusCallApiReqDTO; +import cn.code.nl.framework.execute.biz.vo.AcsApplyActionRespVO; +import com.alibaba.fastjson.JSON; +import org.springframework.validation.annotation.Validated; +import org.springframework.web.bind.annotation.RestController; + +import static cn.code.nl.framework.common.pojo.CommonResult.success; + +/** + * + * @Author: liyongde + * @Date: 2026/7/15 15:25 + */ +@RestController // 提供 RESTful API 接口,给 Feign 调用 +@Validated +public class DemoApi implements LmsTaskCommonApi { + @Override + public CommonResult doHandlePicked(TaskStatusCallApiReqDTO taskStatusCallApiReqDTO) { + // 构建外层VO + AcsApplyActionRespVO respVO = new AcsApplyActionRespVO(); + respVO.setTaskId(10001L); + respVO.setTaskCode("ACS20260716001"); + // 内层data赋值 "success" + return success(respVO); + } + + @Override + public CommonResult doHandleApplyAgain(TaskStatusCallApiReqDTO taskStatusCallApiReqDTO) { + DemoRespVO demoRespVO = new DemoRespVO(); + demoRespVO.setTargetPoint("A_10001"); + + // 构建外层VO + AcsApplyActionRespVO respVO = new AcsApplyActionRespVO(); + respVO.setTaskId(10001L); + respVO.setTaskCode("ACS20260716001"); + // 内层data赋值 "success" + respVO.setData(JSON.parseObject(JSON.toJSONString(demoRespVO))); + return success(respVO); + } + + @Override + public CommonResult doHandleRequestRelease(TaskStatusCallApiReqDTO taskStatusCallApiReqDTO) { + // 构建外层VO + AcsApplyActionRespVO respVO = new AcsApplyActionRespVO(); + respVO.setTaskId(10001L); + respVO.setTaskCode("ACS20260716001"); + return success(respVO); + } + + @Override + public CommonResult doHandleRequestPick(TaskStatusCallApiReqDTO taskStatusCallApiReqDTO) { + // 构建外层VO + AcsApplyActionRespVO respVO = new AcsApplyActionRespVO(); + respVO.setTaskId(10001L); + respVO.setTaskCode("ACS20260716001"); + return success(respVO); + } + + @Override + public CommonResult doHandleRequestLeave(TaskStatusCallApiReqDTO taskStatusCallApiReqDTO) { + // 构建外层VO + AcsApplyActionRespVO respVO = new AcsApplyActionRespVO(); + respVO.setTaskId(10001L); + respVO.setTaskCode("ACS20260716001"); + return success(respVO); + } + + @Override + public CommonResult doHandleRequestEnter(TaskStatusCallApiReqDTO taskStatusCallApiReqDTO) { + // 构建外层VO + AcsApplyActionRespVO respVO = new AcsApplyActionRespVO(); + respVO.setTaskId(10001L); + respVO.setTaskCode("ACS20260716001"); + return success(respVO); + } +} diff --git a/nl-module-lms/nl-module-lms-server/src/main/java/cn/code/nl/module/lms/demo/DemoRespVO.java b/nl-module-lms/nl-module-lms-server/src/main/java/cn/code/nl/module/lms/demo/DemoRespVO.java new file mode 100644 index 00000000..29ea7eb0 --- /dev/null +++ b/nl-module-lms/nl-module-lms-server/src/main/java/cn/code/nl/module/lms/demo/DemoRespVO.java @@ -0,0 +1,13 @@ +package cn.code.nl.module.lms.demo; + +import lombok.Data; + +/** + * + * @Author: liyongde + * @Date: 2026/7/16 17:51 + */ +@Data +public class DemoRespVO { + private String targetPoint; +} diff --git a/nl-module-lms/nl-module-lms-server/src/main/java/cn/code/nl/module/lms/DemoTask.java b/nl-module-lms/nl-module-lms-server/src/main/java/cn/code/nl/module/lms/demo/DemoTask.java similarity index 77% rename from nl-module-lms/nl-module-lms-server/src/main/java/cn/code/nl/module/lms/DemoTask.java rename to nl-module-lms/nl-module-lms-server/src/main/java/cn/code/nl/module/lms/demo/DemoTask.java index c4fdf2ac..78c2c42b 100644 --- a/nl-module-lms/nl-module-lms-server/src/main/java/cn/code/nl/module/lms/DemoTask.java +++ b/nl-module-lms/nl-module-lms-server/src/main/java/cn/code/nl/module/lms/demo/DemoTask.java @@ -1,4 +1,4 @@ -package cn.code.nl.module.lms; +package cn.code.nl.module.lms.demo; import cn.code.nl.framework.execute.core.AbstractTask; import cn.code.nl.framework.execute.core.dto.TaskExecuteDTO; @@ -23,4 +23,9 @@ public class DemoTask extends AbstractTask { public void doHandleCancel(TaskExecuteDTO taskExecuteDTO) { } + + @Override + public String againApply(TaskExecuteDTO taskExecuteDTO) { + return "A_0001"; + } } diff --git a/nl-module-task/nl-module-task-api/src/main/java/cn/code/nl/module/task/dto/AcsFeedbackReqDTO.java b/nl-module-task/nl-module-task-api/src/main/java/cn/code/nl/module/task/dto/AcsFeedbackReqDTO.java index c8050209..4d630acc 100644 --- a/nl-module-task/nl-module-task-api/src/main/java/cn/code/nl/module/task/dto/AcsFeedbackReqDTO.java +++ b/nl-module-task/nl-module-task-api/src/main/java/cn/code/nl/module/task/dto/AcsFeedbackReqDTO.java @@ -8,7 +8,7 @@ import lombok.Data; import java.util.Map; /** - * ACS 状态反馈请求 DTO + * ACS 请求 DTO */ @Schema(description = "ACS 状态反馈 Request DTO") @Data @@ -21,7 +21,10 @@ public class AcsFeedbackReqDTO { @Schema(description = "任务编码") private String taskCode; - @Schema(description = "ACS反馈状态:EXECUTING/PICKED/FINISHED/CANCELLED/APPLYAGAIN", requiredMode = Schema.RequiredMode.REQUIRED) + @Schema(description = "设备编码") + private String deviceCode; + + @Schema(description = "ACS反馈状态或业务操作:状态反馈为 EXECUTING/FINISHED/CANCELLED;业务操作为 PICKED/APPLY-AGAIN/REQUEST-RELEASE/REQUEST-PICK/REQUEST-LEAVE/REQUEST-ENTER", requiredMode = Schema.RequiredMode.REQUIRED) @NotEmpty(message = "反馈状态不能为空") private String status; diff --git a/nl-module-task/nl-module-task-api/src/main/java/cn/code/nl/module/task/enums/AcsBusinessOperationTypeEnum.java b/nl-module-task/nl-module-task-api/src/main/java/cn/code/nl/module/task/enums/AcsBusinessOperationTypeEnum.java new file mode 100644 index 00000000..9a896208 --- /dev/null +++ b/nl-module-task/nl-module-task-api/src/main/java/cn/code/nl/module/task/enums/AcsBusinessOperationTypeEnum.java @@ -0,0 +1,47 @@ +package cn.code.nl.module.task.enums; + +import lombok.Getter; + +/** + * ACS 业务操作类型枚举 + * + * @Author: liyongde + * @Date: 2026/7/16 + */ +@Getter +public enum AcsBusinessOperationTypeEnum { + + PICKED("PICKED", "取货完成"), + + APPLY_AGAIN("APPLY-AGAIN", "二次请求"), + + REQUEST_RELEASE("REQUEST-RELEASE", "请求放货"), + + REQUEST_PICK("REQUEST-PICK", "请求取货"), + + REQUEST_LEAVE("REQUEST-LEAVE", "请求离开"), + + REQUEST_ENTER("REQUEST-ENTER", "请求进入"); + + /** 操作编码 */ + private final String code; + /** 操作名称 */ + private final String name; + + AcsBusinessOperationTypeEnum(String code, String name) { + this.code = code; + this.name = name; + } + + /** + * 根据编码解析枚举,找不到返回 null + */ + public static AcsBusinessOperationTypeEnum getByCode(String code) { + for (AcsBusinessOperationTypeEnum type : values()) { + if (type.getCode().equalsIgnoreCase(code)) { + return type; + } + } + return null; + } +} diff --git a/nl-module-task/nl-module-task-api/src/main/java/cn/code/nl/module/task/enums/ErrorCodeConstants.java b/nl-module-task/nl-module-task-api/src/main/java/cn/code/nl/module/task/enums/ErrorCodeConstants.java index 96d36590..41e2a96a 100644 --- a/nl-module-task/nl-module-task-api/src/main/java/cn/code/nl/module/task/enums/ErrorCodeConstants.java +++ b/nl-module-task/nl-module-task-api/src/main/java/cn/code/nl/module/task/enums/ErrorCodeConstants.java @@ -8,6 +8,7 @@ import cn.code.nl.framework.common.exception.ErrorCode; * @Date: 2026/7/14 13:03 */ public interface ErrorCodeConstants { + ErrorCode TRANSPORT_TASK_NOT_EXISTS_ID = new ErrorCode(5000, "任务id={}不存在,请检查"); ErrorCode TRANSPORT_TASK_NOT_EXISTS = new ErrorCode(5001, "搬运任务不存在"); ErrorCode TRANSPORT_TASK_EXISTS_UNFINISHED = new ErrorCode(5002, "同业务类型下存在未完结的搬运任务"); ErrorCode TRANSPORT_TASK_STATUS_NOT_ALLOW = new ErrorCode(5003, "当前任务状态不允许此操作"); diff --git a/nl-module-task/nl-module-task-api/src/main/java/cn/code/nl/module/task/enums/TaskOperationTypeEnum.java b/nl-module-task/nl-module-task-api/src/main/java/cn/code/nl/module/task/enums/TaskOperationTypeEnum.java index 27cba873..508e2e22 100644 --- a/nl-module-task/nl-module-task-api/src/main/java/cn/code/nl/module/task/enums/TaskOperationTypeEnum.java +++ b/nl-module-task/nl-module-task-api/src/main/java/cn/code/nl/module/task/enums/TaskOperationTypeEnum.java @@ -15,14 +15,10 @@ public enum TaskOperationTypeEnum { EXECUTING("EXECUTING", "执行中", true, false), - PICKED("PICKED", "取货完成", true, false), - FINISHED("FINISHED", "完成任务", true, true), CANCELLED("CANCELLED", "取消任务", true, true), - APPLY_AGAIN("APPLY-AGAIN", "二次请求", true, false), - FORCE_FINISH("FORCE-FINISH", "强制完成", false, true); /** 操作编码 */ diff --git a/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/controller/admin/transporttask/AcsFeedbackController.java b/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/controller/admin/transporttask/AcsFeedbackController.java index 3e9a9590..77f962f5 100644 --- a/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/controller/admin/transporttask/AcsFeedbackController.java +++ b/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/controller/admin/transporttask/AcsFeedbackController.java @@ -1,6 +1,7 @@ package cn.code.nl.module.task.controller.admin.transporttask; import cn.code.nl.framework.common.pojo.CommonResult; +import cn.code.nl.framework.execute.biz.vo.AcsApplyActionRespVO; import cn.code.nl.module.task.dto.AcsFeedbackReqDTO; import cn.code.nl.module.task.service.transporttask.TransportTaskFeedbackService; import io.swagger.v3.oas.annotations.Operation; @@ -18,7 +19,7 @@ import static cn.code.nl.framework.common.pojo.CommonResult.success; /** * ACS 状态反馈接口 *

- * 接收 ACS 的任务状态回调,处理执行中/取货完成/完成/取消 + * 接收 ACS 的任务状态与业务操作回调 */ @Tag(name = "ACS 反馈") @RestController @@ -35,4 +36,10 @@ public class AcsFeedbackController { feedbackService.receiveAcsFeedback(reqDTO); return success(true); } + + @PostMapping("/acs-business-feedback") + @Operation(summary = "接收 ACS 业务操作反馈") + public CommonResult receiveAcsBusinessFeedback(@Valid @RequestBody AcsFeedbackReqDTO reqDTO) { + return success(feedbackService.receiveAcsBusinessFeedback(reqDTO)); + } } diff --git a/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/manage/TransportTaskBusinessOperationManager.java b/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/manage/TransportTaskBusinessOperationManager.java new file mode 100644 index 00000000..1d24db11 --- /dev/null +++ b/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/manage/TransportTaskBusinessOperationManager.java @@ -0,0 +1,196 @@ +package cn.code.nl.module.task.manage; + +import cn.code.nl.framework.common.exception.ServiceException; +import cn.code.nl.framework.common.pojo.CommonResult; +import cn.code.nl.framework.execute.biz.api.TaskCommonApi; +import cn.code.nl.framework.execute.biz.dto.TaskStatusCallApiReqDTO; +import cn.code.nl.framework.execute.biz.vo.AcsApplyActionRespVO; +import cn.code.nl.framework.execute.core.TaskCommonApiFactory; +import cn.code.nl.module.task.dal.dataobject.transporttask.TransportTaskDO; +import cn.code.nl.module.task.dal.mysql.transporttask.TransportTaskMapper; +import cn.code.nl.module.task.dto.AcsFeedbackReqDTO; +import cn.code.nl.module.task.enums.AcsBusinessOperationTypeEnum; +import cn.code.nl.module.task.enums.TransportTaskStatusEnum; +import com.alibaba.fastjson.JSON; +import jakarta.annotation.PostConstruct; +import jakarta.annotation.Resource; +import lombok.extern.slf4j.Slf4j; +import org.redisson.api.RLock; +import org.redisson.api.RedissonClient; +import org.springframework.stereotype.Component; + +import java.util.EnumMap; +import java.util.Map; +import java.util.function.BiFunction; + +import static cn.code.nl.module.task.framework.common.util.TaskUtil.isAllowedFrom; + +/** + * ACS 业务操作分发管理器 + */ +@Slf4j +@Component +public class TransportTaskBusinessOperationManager { + + @Resource + private TransportTaskMapper transportTaskMapper; + + @Resource + private TaskCommonApiFactory taskCommonApiFactory; + + @Resource + private RedissonClient redissonClient; + + /** ACS 业务操作路由表 */ + private final Map> operationHandlers = + new EnumMap<>(AcsBusinessOperationTypeEnum.class); + + /** + * 初始化 ACS 业务操作路由表 + */ + @PostConstruct + public void initOperationHandlers() { + operationHandlers.put(AcsBusinessOperationTypeEnum.PICKED, this::handlePicked); + operationHandlers.put(AcsBusinessOperationTypeEnum.APPLY_AGAIN, this::handleApplyAgain); + operationHandlers.put(AcsBusinessOperationTypeEnum.REQUEST_RELEASE, this::handleRequestRelease); + operationHandlers.put(AcsBusinessOperationTypeEnum.REQUEST_PICK, this::handleRequestPick); + operationHandlers.put(AcsBusinessOperationTypeEnum.REQUEST_LEAVE, this::handleRequestLeave); + operationHandlers.put(AcsBusinessOperationTypeEnum.REQUEST_ENTER, this::handleRequestEnter); + } + + /** + * 按 ACS 业务操作类型分发 + */ + public AcsApplyActionRespVO dispatchOperation(TransportTaskDO task, AcsBusinessOperationTypeEnum type, + AcsFeedbackReqDTO reqDTO) { + RLock lock = redissonClient.getLock(String.valueOf(task.getTaskId())); + if (lock.tryLock()) { + try { + BiFunction handler = operationHandlers.get(type); + if (handler == null) { + log.warn("ACS 业务操作未注册处理器, taskId={}, type={}", task.getTaskId(), type.getCode()); + return buildResp(task, null); + } + Object data = handler.apply(task, reqDTO); + return buildResp(task, data); + } catch (Exception ex) { + log.error("[acsBusinessOperation][执行异常][lockKey={}]", task.getTaskId(), ex); + return buildResp(task, null); + } finally { + if (lock.isHeldByCurrentThread()) { + lock.unlock(); + } + } + } else { + throw new ServiceException(5007, "任务标识为:" + task.getTaskId() + "的任务正在操作中!"); + } + } + + /** + * 处理取货完成 + */ + private AcsApplyActionRespVO handlePicked(TransportTaskDO task, AcsFeedbackReqDTO reqDTO) { + if (!isAllowedFrom(task.getTaskStatus(), TransportTaskStatusEnum.EXECUTING.getCode(), + TransportTaskStatusEnum.PICKED.getCode())) { + log.warn("取货完成反馈状态校验不通过, taskId={}, currentStatus={}", + task.getTaskId(), task.getTaskStatus()); + return null; + } + task.setTaskStatus(TransportTaskStatusEnum.PICKED.getCode()); + transportTaskMapper.updateById(task); + log.info("任务已取货, taskId={}", task.getTaskId()); + + // 调用具体的服务去执行取货完成操作。 + TaskCommonApi serverApi = taskCommonApiFactory.getByServerName(task.getOwnerService()); + CommonResult result = + serverApi.doHandlePicked(buildReq(task, reqDTO)); + return result.getData(); + } + + /** + * 处理二次请求 + */ + private AcsApplyActionRespVO handleApplyAgain(TransportTaskDO task, AcsFeedbackReqDTO reqDTO) { + // TODO 二次请求业务待开发 + log.info("二次请求业务待开发, taskId={}", task.getTaskId()); + // 调用具体的服务去执行取货完成操作。 + TaskCommonApi serverApi = taskCommonApiFactory.getByServerName(task.getOwnerService()); + CommonResult result = + serverApi.doHandleApplyAgain(buildReq(task, reqDTO)); + return result.getData(); + } + + /** + * 处理请求放货 + */ + private Object handleRequestRelease(TransportTaskDO task, AcsFeedbackReqDTO reqDTO) { + // TODO 请求放货业务待开发 + log.info("请求放货业务待开发, taskId={}", task.getTaskId()); + TaskCommonApi serverApi = taskCommonApiFactory.getByServerName(task.getOwnerService()); + CommonResult result = + serverApi.doHandleRequestRelease(buildReq(task, reqDTO)); + return result.getData(); + } + + /** + * 处理请求取货 + */ + private Object handleRequestPick(TransportTaskDO task, AcsFeedbackReqDTO reqDTO) { + // TODO 请求取货业务待开发 + log.info("请求取货业务待开发, taskId={}", task.getTaskId()); + TaskCommonApi serverApi = taskCommonApiFactory.getByServerName(task.getOwnerService()); + CommonResult result = + serverApi.doHandleRequestPick(buildReq(task, reqDTO)); + return result.getData(); + } + + /** + * 处理请求离开 + */ + private Object handleRequestLeave(TransportTaskDO task, AcsFeedbackReqDTO reqDTO) { + // TODO 请求离开业务待开发 + log.info("请求离开业务待开发, taskId={}", task.getTaskId()); + TaskCommonApi serverApi = taskCommonApiFactory.getByServerName(task.getOwnerService()); + CommonResult result = + serverApi.doHandleRequestLeave(buildReq(task, reqDTO)); + return result.getData(); + } + + /** + * 处理请求进入 + */ + private Object handleRequestEnter(TransportTaskDO task, AcsFeedbackReqDTO reqDTO) { + // TODO 请求进入业务待开发 + log.info("请求进入业务待开发, taskId={}", task.getTaskId()); + TaskCommonApi serverApi = taskCommonApiFactory.getByServerName(task.getOwnerService()); + CommonResult result = + serverApi.doHandleRequestEnter(buildReq(task, reqDTO)); + return result.getData(); + } + + /** + * 构建 ACS 业务操作返回对象 + */ + private AcsApplyActionRespVO buildResp(TransportTaskDO task, Object data) { + AcsApplyActionRespVO respVO = new AcsApplyActionRespVO(); + respVO.setTaskId(task.getTaskId()); + respVO.setTaskCode(task.getTaskCode()); + respVO.setData(JSON.parseObject(JSON.toJSONString(data))); + return respVO; + } + /** + * 构建 ACS 业务操作返回对象 + */ + private TaskStatusCallApiReqDTO buildReq(TransportTaskDO task, AcsFeedbackReqDTO reqDTO) { + TaskStatusCallApiReqDTO req = new TaskStatusCallApiReqDTO(); + req.setTaskId(task.getTaskId()); + req.setTaskCode(task.getTaskCode()); + req.setStatus(TransportTaskStatusEnum.PICKED.getCode()); + req.setOwnerService(task.getOwnerService()); + req.setBizType(task.getBizType()); + req.setBizId(task.getBizId()); + req.setHandleCode(task.getHandleCode()); + req.setPayload(reqDTO.getPayload()); + return req; + } +} diff --git a/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/manage/TransportTaskOperationManager.java b/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/manage/TransportTaskOperationManager.java new file mode 100644 index 00000000..e0c3785f --- /dev/null +++ b/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/manage/TransportTaskOperationManager.java @@ -0,0 +1,181 @@ +package cn.code.nl.module.task.manage; + +import cn.code.nl.framework.common.exception.ServiceException; +import cn.code.nl.module.task.dal.dataobject.transporttask.TransportTaskDO; +import cn.code.nl.module.task.dal.mysql.transporttask.TransportTaskMapper; +import cn.code.nl.module.task.dto.AcsFeedbackReqDTO; +import cn.code.nl.module.task.enums.CallbackStatusEnum; +import cn.code.nl.module.task.enums.FinishedTypeEnum; +import cn.code.nl.module.task.enums.TaskEventTypeEnum; +import cn.code.nl.module.task.enums.TaskOperationTypeEnum; +import cn.code.nl.module.task.enums.TransportTaskStatusEnum; +import cn.code.nl.module.task.mq.message.TaskEventMessage; +import cn.code.nl.module.task.mq.producer.TaskEventProducer; +import cn.hutool.core.util.StrUtil; +import jakarta.annotation.PostConstruct; +import jakarta.annotation.Resource; +import lombok.extern.slf4j.Slf4j; +import org.redisson.api.RLock; +import org.redisson.api.RedissonClient; +import org.springframework.stereotype.Component; + +import java.util.EnumMap; +import java.util.Map; +import java.util.function.BiConsumer; + +import static cn.code.nl.module.task.framework.common.util.TaskUtil.isAllowedFrom; + +/** + * 搬运任务状态操作分发管理器 + */ +@Slf4j +@Component +public class TransportTaskOperationManager { + + @Resource + private TransportTaskMapper transportTaskMapper; + + @Resource + private TaskEventProducer taskEventProducer; + + @Resource + private RedissonClient redissonClient; + + /** 操作类型路由表 */ + private final Map> operationHandlers = + new EnumMap<>(TaskOperationTypeEnum.class); + + /** + * 初始化操作类型路由表 + */ + @PostConstruct + public void initOperationHandlers() { + operationHandlers.put(TaskOperationTypeEnum.EXECUTING, this::handleExecuting); + operationHandlers.put(TaskOperationTypeEnum.FINISHED, this::handleFinished); + operationHandlers.put(TaskOperationTypeEnum.CANCELLED, this::handleCancelled); + operationHandlers.put(TaskOperationTypeEnum.FORCE_FINISH, (task, reqDTO) -> handleForceFinish(task)); + } + + /** + * 按操作类型查路由表分发 + */ + public void dispatchOperation(TransportTaskDO task, TaskOperationTypeEnum type, AcsFeedbackReqDTO reqDTO) { + RLock lock = redissonClient.getLock(String.valueOf(task.getTaskId())); + if (lock.tryLock()) { + try { + BiConsumer handler = operationHandlers.get(type); + if (handler == null) { + log.warn("操作类型未注册处理器, taskId={}, type={}", task.getTaskId(), type.getCode()); + return; + } + handler.accept(task, reqDTO); + } catch (Exception ex) { + log.error("[messageResend][执行异常][lockKey={}]", task.getTaskId(), ex); + } finally { + if (lock.isHeldByCurrentThread()) { + lock.unlock(); + } + } + } else { + throw new ServiceException(5007, "任务标识为:" + task.getTaskId() + "的任务正在操作中!"); + } + } + + /** + * 处理执行中 + */ + private void handleExecuting(TransportTaskDO task, AcsFeedbackReqDTO reqDTO) { + if (!isAllowedFrom(task.getTaskStatus(), TransportTaskStatusEnum.READY.getCode(), + TransportTaskStatusEnum.ISSUED.getCode())) { + log.warn("执行中反馈状态校验不通过, taskId={}, currentStatus={}", + task.getTaskId(), task.getTaskStatus()); + return; + } + task.setTaskStatus(TransportTaskStatusEnum.EXECUTING.getCode()); + task.setCarNo(reqDTO.getCarNo()); + transportTaskMapper.updateById(task); + log.info("任务执行中, taskId={}", task.getTaskId()); + } + + /** + * 处理完成 + */ + private void handleFinished(TransportTaskDO task, AcsFeedbackReqDTO reqDTO) { + if (!isAllowedFrom(task.getTaskStatus(), TransportTaskStatusEnum.CREATED.getCode(), + TransportTaskStatusEnum.READY.getCode(), TransportTaskStatusEnum.ISSUING.getCode(), + TransportTaskStatusEnum.ISSUED.getCode(), TransportTaskStatusEnum.EXECUTING.getCode(), + TransportTaskStatusEnum.PICKED.getCode(), + TransportTaskStatusEnum.FINISHED_CALLBACK_PENDING.getCode())) { + log.warn("完成反馈状态校验不通过, taskId={}, currentStatus={}", + task.getTaskId(), task.getTaskStatus()); + return; + } + if (!TransportTaskStatusEnum.FINISHED_CALLBACK_PENDING.getCode().equals(task.getTaskStatus())) { + task.setTaskStatus(TransportTaskStatusEnum.FINISHED_CALLBACK_PENDING.getCode()); + } + task.setCallbackStatus(CallbackStatusEnum.PENDING.getCode()); + transportTaskMapper.updateById(task); + + publishEvent(task, TaskEventTypeEnum.TASK_FINISHED, reqDTO.getPayload()); + } + + /** + * 处理取消 + */ + private void handleCancelled(TransportTaskDO task, AcsFeedbackReqDTO reqDTO) { + if (!isAllowedFrom(task.getTaskStatus(), TransportTaskStatusEnum.CREATED.getCode(), + TransportTaskStatusEnum.READY.getCode(), TransportTaskStatusEnum.ISSUING.getCode(), + TransportTaskStatusEnum.ISSUED.getCode(), TransportTaskStatusEnum.EXECUTING.getCode(), + TransportTaskStatusEnum.PICKED.getCode(), + TransportTaskStatusEnum.FINISHED_CALLBACK_PENDING.getCode(), + TransportTaskStatusEnum.CANCEL_CALLBACK_PENDING.getCode())) { + log.warn("取消反馈状态校验不通过, taskId={}, currentStatus={}", + task.getTaskId(), task.getTaskStatus()); + return; + } + if (!TransportTaskStatusEnum.CANCEL_CALLBACK_PENDING.getCode().equals(task.getTaskStatus())) { + task.setTaskStatus(TransportTaskStatusEnum.CANCEL_CALLBACK_PENDING.getCode()); + } + task.setCallbackStatus(CallbackStatusEnum.PENDING.getCode()); + transportTaskMapper.updateById(task); + + publishEvent(task, TaskEventTypeEnum.TASK_CANCELLED, reqDTO.getPayload()); + } + + /** + * 处理强制完成 + */ + private void handleForceFinish(TransportTaskDO task) { + task.setTaskStatus(TransportTaskStatusEnum.FINISHED.getCode()); + task.setFinishedType(FinishedTypeEnum.MANUAL_FORCE.getCode()); + transportTaskMapper.updateById(task); + log.info("任务强制完成, taskId={}", task.getTaskId()); + } + + /** + * 发布 MQ 事件 + */ + private void publishEvent(TransportTaskDO task, TaskEventTypeEnum eventType, + Map payload) { + TaskEventMessage msg = new TaskEventMessage(); + msg.setEventType(eventType.getCode()); + msg.setTaskId(task.getTaskId()); + msg.setTaskCode(task.getTaskCode()); + msg.setOwnerService(task.getOwnerService()); + msg.setBizType(task.getBizType()); + msg.setBizId(task.getBizId()); + msg.setTaskType(task.getTaskType()); + msg.setHandleCode(task.getHandleCode()); + msg.setPayload(payload); + + try { + taskEventProducer.publishEvent(msg); + } catch (Exception e) { + log.error("MQ 发布失败, taskId={}, eventType={}", + task.getTaskId(), eventType.getCode(), e); + task.setCallbackStatus(CallbackStatusEnum.FAILED.getCode()); + task.setCallbackErrorMsg(StrUtil.maxLength(e.getMessage(), 500)); + transportTaskMapper.updateById(task); + } + } +} diff --git a/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/service/transporttask/TransportTaskFeedbackService.java b/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/service/transporttask/TransportTaskFeedbackService.java index fc3b21db..0a010be1 100644 --- a/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/service/transporttask/TransportTaskFeedbackService.java +++ b/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/service/transporttask/TransportTaskFeedbackService.java @@ -1,5 +1,6 @@ package cn.code.nl.module.task.service.transporttask; +import cn.code.nl.framework.execute.biz.vo.AcsApplyActionRespVO; import cn.code.nl.module.task.dto.AcsFeedbackReqDTO; import jakarta.validation.Valid; @@ -16,4 +17,12 @@ public interface TransportTaskFeedbackService { * @param reqDTO 反馈请求 */ void receiveAcsFeedback(@Valid AcsFeedbackReqDTO reqDTO); + + /** + * 接收 ACS 业务操作反馈并处理 + * + * @param reqDTO 反馈请求 + */ + AcsApplyActionRespVO receiveAcsBusinessFeedback(AcsFeedbackReqDTO reqDTO); + } diff --git a/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/service/transporttask/TransportTaskFeedbackServiceImpl.java b/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/service/transporttask/TransportTaskFeedbackServiceImpl.java index 19032375..d1f519de 100644 --- a/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/service/transporttask/TransportTaskFeedbackServiceImpl.java +++ b/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/service/transporttask/TransportTaskFeedbackServiceImpl.java @@ -1,10 +1,15 @@ package cn.code.nl.module.task.service.transporttask; +import cn.code.nl.framework.common.exception.ServiceException; +import cn.code.nl.framework.execute.biz.vo.AcsApplyActionRespVO; import cn.code.nl.module.task.dal.dataobject.transporttask.TransportTaskDO; import cn.code.nl.module.task.dal.mysql.transporttask.TransportTaskMapper; import cn.code.nl.module.task.dto.AcsFeedbackReqDTO; +import cn.code.nl.module.task.enums.AcsBusinessOperationTypeEnum; import cn.code.nl.module.task.enums.TaskOperationTypeEnum; import cn.code.nl.module.task.enums.TransportTaskStatusEnum; +import cn.code.nl.module.task.manage.TransportTaskBusinessOperationManager; +import cn.code.nl.module.task.manage.TransportTaskOperationManager; import com.fasterxml.jackson.databind.ObjectMapper; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Service; @@ -13,6 +18,9 @@ import org.springframework.validation.annotation.Validated; import jakarta.annotation.Resource; +import static cn.code.nl.framework.common.exception.util.ServiceExceptionUtil.exception; +import static cn.code.nl.module.task.enums.ErrorCodeConstants.TRANSPORT_TASK_NOT_EXISTS_ID; + /** * ACS 反馈处理服务实现 *

@@ -29,45 +37,96 @@ public class TransportTaskFeedbackServiceImpl implements TransportTaskFeedbackSe private TransportTaskMapper transportTaskMapper; @Resource - private TransportTaskService transportTaskService; + private TransportTaskOperationManager transportTaskOperationManager; + + @Resource + private TransportTaskBusinessOperationManager transportTaskBusinessOperationManager; private static final ObjectMapper OBJECT_MAPPER = new ObjectMapper(); @Override @Transactional(rollbackFor = Exception.class) public void receiveAcsFeedback(AcsFeedbackReqDTO reqDTO) { - // 1. 查询任务 - TransportTaskDO task = transportTaskMapper.selectById(reqDTO.getTaskId()); - if (task == null) { - log.warn("ACS 反馈 taskId 不存在, taskId={}, 直接返回成功", reqDTO.getTaskId()); - return; - } + TransportTaskDO task = getTaskForFeedback(reqDTO); + saveResultParam(task, reqDTO); - // 2. 幂等判断:已经是终态(完成/取消)则不重复处理 - String currentStatus = task.getTaskStatus(); - if (TransportTaskStatusEnum.FINISHED.getCode().equals(currentStatus) - || TransportTaskStatusEnum.CANCELLED.getCode().equals(currentStatus)) { - log.info("任务已处终态, taskId={}, currentStatus={}, 忽略反馈", - reqDTO.getTaskId(), currentStatus); - return; - } - - // 3. 保存 resultParam - if (reqDTO.getPayload() != null) { - try { - task.setResultParam(OBJECT_MAPPER.writeValueAsString(reqDTO.getPayload())); - } catch (Exception e) { - log.warn("resultParam 序列化失败, taskId={}", reqDTO.getTaskId(), e); - } - } - - // 4. 解析操作类型并校验允许 ACS 触发,走统一路由表分发 + // 解析状态操作类型并校验允许 ACS 触发 TaskOperationTypeEnum type = TaskOperationTypeEnum.getByCode(reqDTO.getStatus()); if (type == null || !type.isAcsAllowed()) { log.warn("未知或不允许的 ACS 反馈状态, taskId={}, status={}", reqDTO.getTaskId(), reqDTO.getStatus()); return; } - transportTaskService.dispatchOperation(task, type, reqDTO); + transportTaskOperationManager.dispatchOperation(task, type, reqDTO); + } + + @Override + @Transactional(rollbackFor = Exception.class) + public AcsApplyActionRespVO receiveAcsBusinessFeedback(AcsFeedbackReqDTO reqDTO) { + TransportTaskDO task = getTaskForFeedback(reqDTO); + if (isFinalTask(task, reqDTO)) { + return buildResp(task); + } + + saveResultParam(task, reqDTO); + + // 解析 ACS 业务操作类型,走业务操作路由表分发 + AcsBusinessOperationTypeEnum type = AcsBusinessOperationTypeEnum.getByCode(reqDTO.getStatus()); + if (type == null) { + log.warn("未知的 ACS 业务操作类型, taskId={}, status={}", + reqDTO.getTaskId(), reqDTO.getStatus()); + return buildResp(task); + } + return transportTaskBusinessOperationManager.dispatchOperation(task, type, reqDTO); + } + + /** + * 查询反馈对应的任务 + */ + private TransportTaskDO getTaskForFeedback(AcsFeedbackReqDTO reqDTO) { + TransportTaskDO task = transportTaskMapper.selectById(reqDTO.getTaskId()); + if (task == null) { + throw exception(TRANSPORT_TASK_NOT_EXISTS_ID, reqDTO.getTaskId()); + } + return task; + } + + /** + * 判断任务是否已经处于终态 + */ + private boolean isFinalTask(TransportTaskDO task, AcsFeedbackReqDTO reqDTO) { + String currentStatus = task.getTaskStatus(); + if (TransportTaskStatusEnum.FINISHED.getCode().equals(currentStatus) + || TransportTaskStatusEnum.CANCELLED.getCode().equals(currentStatus)) { + log.info("任务已处终态, taskId={}, currentStatus={}, 忽略反馈", + reqDTO.getTaskId(), currentStatus); + return true; + } + return false; + } + + /** + * 保存 ACS 反馈扩展数据 + */ + private void saveResultParam(TransportTaskDO task, AcsFeedbackReqDTO reqDTO) { + if (reqDTO.getPayload() == null) { + return; + } + try { + task.setResultParam(OBJECT_MAPPER.writeValueAsString(reqDTO.getPayload())); + transportTaskMapper.updateById(task); + } catch (Exception e) { + log.warn("resultParam 序列化失败, taskId={}", reqDTO.getTaskId(), e); + } + } + + /** + * 根据任务构建空数据返回 + */ + private AcsApplyActionRespVO buildResp(TransportTaskDO task) { + AcsApplyActionRespVO respVO = new AcsApplyActionRespVO(); + respVO.setTaskId(task.getTaskId()); + respVO.setTaskCode(task.getTaskCode()); + return respVO; } } diff --git a/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/service/transporttask/TransportTaskService.java b/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/service/transporttask/TransportTaskService.java index 68d7d413..b7fdb0ef 100644 --- a/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/service/transporttask/TransportTaskService.java +++ b/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/service/transporttask/TransportTaskService.java @@ -1,106 +1,95 @@ -package cn.code.nl.module.task.service.transporttask; - -import java.util.*; -import jakarta.validation.*; -import cn.code.nl.module.task.controller.admin.transporttask.vo.*; -import cn.code.nl.module.task.dal.dataobject.transporttask.TransportTaskDO; -import cn.code.nl.module.task.dto.AcsFeedbackReqDTO; -import cn.code.nl.module.task.dto.TaskInfoDTO; -import cn.code.nl.module.task.dto.TransportTaskCreateReqDTO; -import cn.code.nl.module.task.enums.TaskOperationTypeEnum; -import cn.code.nl.framework.common.pojo.PageResult; -import cn.code.nl.framework.common.pojo.PageParam; - -/** - * 搬运任务 Service 接口 - * - * @author 诺力管理员 - */ -public interface TransportTaskService { - - /** - * 创建搬运任务 - * - * @param createReqVO 创建信息 - * @return 编号 - */ - Long createTransportTask(@Valid TransportTaskSaveReqVO createReqVO); - - /** - * 通过 RPC 创建搬运任务(LMS/WMS 调用) - * - * @param reqDTO 创建请求 - * @return 任务ID - */ - Long createTransportTaskByRpc(@Valid TransportTaskCreateReqDTO reqDTO); - - /** - * 更新搬运任务 - * - * @param updateReqVO 更新信息 - */ - void updateTransportTask(@Valid TransportTaskSaveReqVO updateReqVO); - - /** - * 删除搬运任务 - * - * @param id 编号 - */ - void deleteTransportTask(Long id); - - /** - * 批量删除搬运任务 - * - * @param ids 编号 - */ - void deleteTransportTaskListByIds(List ids); - - /** - * 获得搬运任务 - * - * @param id 编号 - * @return 搬运任务 - */ - TransportTaskDO getTransportTask(Long id); - - /** - * 获得搬运任务分页 - * - * @param pageReqVO 分页查询 - * @return 搬运任务分页 - */ - PageResult getTransportTaskPage(TransportTaskPageReqVO pageReqVO); - - /** - * PC 端操作搬运任务(完成/取消/强制完成) - * - * @param reqVO 操作请求 - */ - void operateTransportTask(@Valid TransportTaskOperateReqVO reqVO); - - /** - * 按操作类型分发处理任务(ACS 反馈与 PC 端操作共用路由表) - * - * @param task 任务 - * @param type 操作类型 - * @param reqDTO 反馈请求(PC 端为构造的精简对象) - */ - void dispatchOperation(TransportTaskDO task, TaskOperationTypeEnum type, AcsFeedbackReqDTO reqDTO); - - /** - * 根据 taskId 查询任务全量信息(RPC 用,自动过滤逻辑删除,查不到返回 null) - * - * @param taskId 任务ID - * @return 任务全量信息 - */ - TaskInfoDTO getTaskInfoById(Long taskId); - - /** - * 根据 taskCode 查询任务全量信息(RPC 用,自动过滤逻辑删除,查不到返回 null) - * - * @param taskCode 任务编码 - * @return 任务全量信息 - */ - TaskInfoDTO getTaskInfoByCode(String taskCode); - -} \ No newline at end of file +package cn.code.nl.module.task.service.transporttask; + +import java.util.*; +import jakarta.validation.*; +import cn.code.nl.module.task.controller.admin.transporttask.vo.*; +import cn.code.nl.module.task.dal.dataobject.transporttask.TransportTaskDO; +import cn.code.nl.module.task.dto.TaskInfoDTO; +import cn.code.nl.module.task.dto.TransportTaskCreateReqDTO; +import cn.code.nl.framework.common.pojo.PageResult; +import cn.code.nl.framework.common.pojo.PageParam; + +/** + * 搬运任务 Service 接口 + * + * @author 诺力管理员 + */ +public interface TransportTaskService { + + /** + * 创建搬运任务 + * + * @param createReqVO 创建信息 + * @return 编号 + */ + Long createTransportTask(@Valid TransportTaskSaveReqVO createReqVO); + + /** + * 通过 RPC 创建搬运任务(LMS/WMS 调用) + * + * @param reqDTO 创建请求 + * @return 任务ID + */ + Long createTransportTaskByRpc(@Valid TransportTaskCreateReqDTO reqDTO); + + /** + * 更新搬运任务 + * + * @param updateReqVO 更新信息 + */ + void updateTransportTask(@Valid TransportTaskSaveReqVO updateReqVO); + + /** + * 删除搬运任务 + * + * @param id 编号 + */ + void deleteTransportTask(Long id); + + /** + * 批量删除搬运任务 + * + * @param ids 编号 + */ + void deleteTransportTaskListByIds(List ids); + + /** + * 获得搬运任务 + * + * @param id 编号 + * @return 搬运任务 + */ + TransportTaskDO getTransportTask(Long id); + + /** + * 获得搬运任务分页 + * + * @param pageReqVO 分页查询 + * @return 搬运任务分页 + */ + PageResult getTransportTaskPage(TransportTaskPageReqVO pageReqVO); + + /** + * PC 端操作搬运任务(完成/取消/强制完成) + * + * @param reqVO 操作请求 + */ + void operateTransportTask(@Valid TransportTaskOperateReqVO reqVO); + + /** + * 根据 taskId 查询任务全量信息(RPC 用,自动过滤逻辑删除,查不到返回 null) + * + * @param taskId 任务ID + * @return 任务全量信息 + */ + TaskInfoDTO getTaskInfoById(Long taskId); + + /** + * 根据 taskCode 查询任务全量信息(RPC 用,自动过滤逻辑删除,查不到返回 null) + * + * @param taskCode 任务编码 + * @return 任务全量信息 + */ + TaskInfoDTO getTaskInfoByCode(String taskCode); + +} diff --git a/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/service/transporttask/TransportTaskServiceImpl.java b/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/service/transporttask/TransportTaskServiceImpl.java index d013fa2e..f8e26b69 100644 --- a/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/service/transporttask/TransportTaskServiceImpl.java +++ b/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/service/transporttask/TransportTaskServiceImpl.java @@ -1,11 +1,7 @@ package cn.code.nl.module.task.service.transporttask; -import cn.code.nl.framework.common.exception.ServiceException; import cn.code.nl.framework.common.pojo.PageResult; import cn.code.nl.framework.common.util.object.BeanUtils; -import cn.code.nl.framework.execute.biz.api.TaskCommonApi; -import cn.code.nl.framework.execute.biz.dto.TaskStatusCallApiReqDTO; -import cn.code.nl.framework.execute.core.TaskCommonApiFactory; import cn.code.nl.module.task.controller.admin.transporttask.vo.TransportTaskOperateReqVO; import cn.code.nl.module.task.controller.admin.transporttask.vo.TransportTaskPageReqVO; import cn.code.nl.module.task.controller.admin.transporttask.vo.TransportTaskSaveReqVO; @@ -16,28 +12,19 @@ import cn.code.nl.module.task.dto.AcsFeedbackReqDTO; import cn.code.nl.module.task.dto.TaskInfoDTO; import cn.code.nl.module.task.dto.TransportTaskCreateReqDTO; import cn.code.nl.module.task.enums.*; -import cn.code.nl.module.task.mq.message.TaskEventMessage; -import cn.code.nl.module.task.mq.producer.TaskEventProducer; -import cn.hutool.core.util.StrUtil; -import jakarta.annotation.PostConstruct; +import cn.code.nl.module.task.manage.TransportTaskOperationManager; import jakarta.annotation.Resource; import lombok.extern.slf4j.Slf4j; -import org.redisson.api.RLock; -import org.redisson.api.RedissonClient; import org.springframework.stereotype.Service; import org.springframework.transaction.annotation.Transactional; import org.springframework.validation.annotation.Validated; -import java.util.EnumMap; import java.util.List; -import java.util.Map; -import java.util.function.BiConsumer; import static cn.code.nl.framework.common.exception.util.ServiceExceptionUtil.exception; import static cn.code.nl.module.task.enums.ErrorCodeConstants.TRANSPORT_TASK_ALREADY_FINAL; import static cn.code.nl.module.task.enums.ErrorCodeConstants.TRANSPORT_TASK_NOT_EXISTS; import static cn.code.nl.module.task.enums.ErrorCodeConstants.TRANSPORT_TASK_OPERATION_NOT_SUPPORTED; -import static cn.code.nl.module.task.framework.common.util.TaskUtil.isAllowedFrom; /** * 搬运任务 Service 实现类 @@ -53,36 +40,7 @@ public class TransportTaskServiceImpl implements TransportTaskService { private TransportTaskMapper transportTaskMapper; @Resource - private TaskEventProducer taskEventProducer; - - @Resource - private TaskCommonApiFactory taskCommonApiFactory; - - @Resource - private RedissonClient redissonClient; - - /** - * 操作类型 -> 处理方法 路由表(ACS 反馈与 PC 端操作共用) - */ - private final Map> operationHandlers = - new EnumMap<>(TaskOperationTypeEnum.class); - - /** - * 初始化操作类型路由表 - */ - @PostConstruct - public void initOperationHandlers() { - operationHandlers.put(TaskOperationTypeEnum.EXECUTING, this::handleExecuting); - operationHandlers.put(TaskOperationTypeEnum.PICKED, this::handlePicked); - operationHandlers.put(TaskOperationTypeEnum.FINISHED, this::handleFinished); - operationHandlers.put(TaskOperationTypeEnum.CANCELLED, this::handleCancelled); - // todo: 二次请求业务未开发 - operationHandlers.put(TaskOperationTypeEnum.APPLY_AGAIN, - (task, reqDTO) -> log.info("二次请求业务未开发, taskId={}", task.getTaskId())); - // todo: 强制完成业务未开发 - operationHandlers.put(TaskOperationTypeEnum.FORCE_FINISH, - (task, reqDTO) -> handleForceFinish(task)); - } + private TransportTaskOperationManager transportTaskOperationManager; @Override public Long createTransportTask(TransportTaskSaveReqVO createReqVO) { @@ -181,33 +139,7 @@ public class TransportTaskServiceImpl implements TransportTaskService { AcsFeedbackReqDTO reqDTO = new AcsFeedbackReqDTO(); reqDTO.setTaskId(task.getTaskId()); reqDTO.setStatus(type.getCode()); - dispatchOperation(task, type, reqDTO); - } - - /** - * 按操作类型查路由表分发 - */ - @Override - public void dispatchOperation(TransportTaskDO task, TaskOperationTypeEnum type, AcsFeedbackReqDTO reqDTO) { - RLock lock = redissonClient.getLock(String.valueOf(task.getTaskId())); - if (lock.tryLock()) { - try { - BiConsumer handler = operationHandlers.get(type); - if (handler == null) { - log.warn("操作类型未注册处理器, taskId={}, type={}", task.getTaskId(), type.getCode()); - return; - } - handler.accept(task, reqDTO); - } catch (Exception ex) { - log.error("[messageResend][执行异常][lockKey={}]", task.getTaskId(), ex); - } finally { - if (lock.isHeldByCurrentThread()) { - lock.unlock(); - } - } - } else { - throw new ServiceException(5007, "任务标识为:" + task.getTaskId() + "的任务正在操作中!"); - } + transportTaskOperationManager.dispatchOperation(task, type, reqDTO); } /** @@ -228,138 +160,4 @@ public class TransportTaskServiceImpl implements TransportTaskService { return TransportTaskConvert.INSTANCE.convert(task); } - // ==================== 各状态处理方法(自 TransportTaskFeedbackServiceImpl 迁入,逻辑未改动) ==================== - - /** - * 处理执行中(60):直接更新状态,不做业务回调 - */ - private void handleExecuting(TransportTaskDO task, AcsFeedbackReqDTO reqDTO) { - if (!isAllowedFrom(task.getTaskStatus(), TransportTaskStatusEnum.READY.getCode(), - TransportTaskStatusEnum.ISSUED.getCode())) { - log.warn("执行中反馈状态校验不通过, taskId={}, currentStatus={}", - task.getTaskId(), task.getTaskStatus()); - return; - } - task.setTaskStatus(TransportTaskStatusEnum.EXECUTING.getCode()); - task.setCarNo(reqDTO.getCarNo()); - transportTaskMapper.updateById(task); - log.info("任务执行中, taskId={}", task.getTaskId()); - } - - /** - * 处理取货完成(65):更新状态 + 同步回调 LMS/WMS - * ownerService: 服务提供商 - */ - private void handlePicked(TransportTaskDO task, AcsFeedbackReqDTO reqDTO) { - if (!isAllowedFrom(task.getTaskStatus(), TransportTaskStatusEnum.EXECUTING.getCode(), - TransportTaskStatusEnum.PICKED.getCode())) { - log.warn("取货完成反馈状态校验不通过, taskId={}, currentStatus={}", - task.getTaskId(), task.getTaskStatus()); - return; - } - // 更新状态 - task.setTaskStatus(TransportTaskStatusEnum.PICKED.getCode()); - transportTaskMapper.updateById(task); - log.info("任务已取货, taskId={}", task.getTaskId()); - - TaskStatusCallApiReqDTO req = new TaskStatusCallApiReqDTO(); - req.setTaskId(task.getTaskId()); - req.setTaskCode(task.getTaskCode()); - req.setStatus(TransportTaskStatusEnum.PICKED.getCode()); - req.setOwnerService(task.getOwnerService()); - req.setBizType(task.getBizType()); - req.setBizId(task.getBizId()); - req.setHandleCode(task.getHandleCode()); - req.setPayload(reqDTO.getPayload()); - - // 同步请求业务 - TaskCommonApi serverApi = taskCommonApiFactory.getByServerName(task.getOwnerService()); - serverApi.doHandlePicked(req); - } - - /** - * 处理完成(75):更新状态 + 异步 MQ 通知 LMS/WMS - */ - private void handleFinished(TransportTaskDO task, AcsFeedbackReqDTO reqDTO) { - if (!isAllowedFrom(task.getTaskStatus(), TransportTaskStatusEnum.CREATED.getCode(), - TransportTaskStatusEnum.READY.getCode(), TransportTaskStatusEnum.ISSUING.getCode(), - TransportTaskStatusEnum.ISSUED.getCode(), TransportTaskStatusEnum.EXECUTING.getCode(), - TransportTaskStatusEnum.PICKED.getCode(), - TransportTaskStatusEnum.FINISHED_CALLBACK_PENDING.getCode())) { - log.warn("完成反馈状态校验不通过, taskId={}, currentStatus={}", - task.getTaskId(), task.getTaskStatus()); - return; - } - // 已是 75 则只发 MQ,不改状态 - if (!TransportTaskStatusEnum.FINISHED_CALLBACK_PENDING.getCode().equals(task.getTaskStatus())) { - task.setTaskStatus(TransportTaskStatusEnum.FINISHED_CALLBACK_PENDING.getCode()); - } - task.setCallbackStatus(CallbackStatusEnum.PENDING.getCode()); - transportTaskMapper.updateById(task); - - publishEvent(task, TaskEventTypeEnum.TASK_FINISHED, reqDTO.getPayload()); - } - - /** - * 处理取消(85):更新状态 + 异步 MQ 通知 LMS/WMS - */ - private void handleCancelled(TransportTaskDO task, AcsFeedbackReqDTO reqDTO) { - if (!isAllowedFrom(task.getTaskStatus(), TransportTaskStatusEnum.CREATED.getCode(), - TransportTaskStatusEnum.READY.getCode(), TransportTaskStatusEnum.ISSUING.getCode(), - TransportTaskStatusEnum.ISSUED.getCode(), TransportTaskStatusEnum.EXECUTING.getCode(), - TransportTaskStatusEnum.PICKED.getCode(), - TransportTaskStatusEnum.FINISHED_CALLBACK_PENDING.getCode(), - TransportTaskStatusEnum.CANCEL_CALLBACK_PENDING.getCode())) { - log.warn("取消反馈状态校验不通过, taskId={}, currentStatus={}", - task.getTaskId(), task.getTaskStatus()); - return; - } - if (!TransportTaskStatusEnum.CANCEL_CALLBACK_PENDING.getCode().equals(task.getTaskStatus())) { - task.setTaskStatus(TransportTaskStatusEnum.CANCEL_CALLBACK_PENDING.getCode()); - } - task.setCallbackStatus(CallbackStatusEnum.PENDING.getCode()); - transportTaskMapper.updateById(task); - - publishEvent(task, TaskEventTypeEnum.TASK_CANCELLED, reqDTO.getPayload()); - } - - /** - * 处理强制完成(79):直接更新状态,不做业务回调 - */ - private void handleForceFinish(TransportTaskDO task) { - task.setTaskStatus(TransportTaskStatusEnum.FINISHED.getCode()); - task.setFinishedType(FinishedTypeEnum.MANUAL_FORCE.getCode()); - transportTaskMapper.updateById(task); - log.info("任务强制完成, taskId={}", task.getTaskId()); - } - - // ==================== 工具方法 ==================== - - /** - * 发布 MQ 事件。发送失败时记录 callbackStatus=FAILED - */ - private void publishEvent(TransportTaskDO task, TaskEventTypeEnum eventType, - Map payload) { - TaskEventMessage msg = new TaskEventMessage(); - msg.setEventType(eventType.getCode()); - msg.setTaskId(task.getTaskId()); - msg.setTaskCode(task.getTaskCode()); - msg.setOwnerService(task.getOwnerService()); - msg.setBizType(task.getBizType()); - msg.setBizId(task.getBizId()); - msg.setTaskType(task.getTaskType()); - msg.setHandleCode(task.getHandleCode()); - msg.setPayload(payload); - - try { - taskEventProducer.publishEvent(msg); - } catch (Exception e) { - log.error("MQ 发布失败, taskId={}, eventType={}", - task.getTaskId(), eventType.getCode(), e); - task.setCallbackStatus(CallbackStatusEnum.FAILED.getCode()); - task.setCallbackErrorMsg(StrUtil.maxLength(e.getMessage(), 500)); - transportTaskMapper.updateById(task); - } - } - }