From 7c0aa128868ddcdeb2a774549cf95c1a18b19336 Mon Sep 17 00:00:00 2001 From: liyongde <1419499670@qq.com> Date: Tue, 14 Jul 2026 16:56:03 +0800 Subject: [PATCH] =?UTF-8?q?docs:=20task=E6=9C=8D=E5=8A=A1=E6=A0=B8?= =?UTF-8?q?=E5=BF=83=E4=B8=9A=E5=8A=A1=E5=AE=9E=E7=8E=B0=E8=AE=A1=E5=88=92?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../2026-07-14-task-core-implementation.md | 1582 +++++++++++++++++ 1 file changed, 1582 insertions(+) create mode 100644 docs/superpowers/plans/2026-07-14-task-core-implementation.md diff --git a/docs/superpowers/plans/2026-07-14-task-core-implementation.md b/docs/superpowers/plans/2026-07-14-task-core-implementation.md new file mode 100644 index 00000000..96eed23a --- /dev/null +++ b/docs/superpowers/plans/2026-07-14-task-core-implementation.md @@ -0,0 +1,1582 @@ +# Task 服务核心业务实现计划 + +> **面向 AI 代理的工作者:** 必需子技能:使用 superpowers:subagent-driven-development(推荐)或 superpowers:executing-plans 逐任务实现此计划。步骤使用复选框(`- [ ]`)语法来跟踪进度。 + +**目标:** 实现 Task 服务的任务创建、ACS 状态反馈处理、同步/异步回调 LMS/WMS、业务回调结果接收及幂等控制。 + +**架构:** 基于 Spring Boot + MyBatis Plus + RocketMQ + RestTemplate。ACS 反馈通过 HTTP 接收入口,根据状态路由:执行中直接改状态,取货完成同步 HTTP 调 LMS/WMS,完成/取消异步 MQ 通知。LMS/WMS 消费 MQ 后回调 Task 写入最终态。 + +**技术栈:** Java 17, Spring Boot 3, MyBatis Plus, RocketMQ, RestTemplate (LoadBalanced), Nacos 服务发现 + +--- + +## 文件结构 + +### 新建文件 + +| 文件 | 职责 | +|------|------| +| `nl-module-task-api/src/main/java/cn/code/nl/module/task/dto/TransportTaskCreateReqDTO.java` | LMS/WMS 创建任务请求 DTO | +| `nl-module-task-api/src/main/java/cn/code/nl/module/task/dto/AcsFeedbackReqDTO.java` | ACS 状态反馈请求 DTO | +| `nl-module-task-api/src/main/java/cn/code/nl/module/task/dto/TaskStatusCallbackReqDTO.java` | Task→LMS/WMS 同步回调请求 DTO | +| `nl-module-task-api/src/main/java/cn/code/nl/module/task/dto/TaskStatusCallbackRespDTO.java` | Task→LMS/WMS 同步回调响应 DTO | +| `nl-module-task-api/src/main/java/cn/code/nl/module/task/dto/TaskCallbackResultReqDTO.java` | LMS/WMS 业务结果回调请求 DTO | +| `nl-module-task-api/src/main/java/cn/code/nl/module/task/dto/TaskEventMessage.java` | MQ 任务事件消息 | +| `nl-module-task-api/src/main/java/cn/code/nl/module/task/enums/CallbackStatusEnum.java` | 回调状态枚举 PENDING/SUCCESS/FAILED | +| `nl-module-task-api/src/main/java/cn/code/nl/module/task/enums/TaskEventTypeEnum.java` | 事件类型枚举 TASK_FINISHED/TASK_CANCELLED | +| `nl-module-task-server/src/main/java/cn/code/nl/module/task/api/TransportTaskApiImpl.java` | RPC 接口实现(创建任务) | +| `nl-module-task-server/src/main/java/cn/code/nl/module/task/controller/AcsFeedbackController.java` | ACS 反馈 HTTP 接收入口 | +| `nl-module-task-server/src/main/java/cn/code/nl/module/task/controller/TaskCallbackController.java` | LMS/WMS 业务结果回调入口 | +| `nl-module-task-server/src/main/java/cn/code/nl/module/task/service/transporttask/TransportTaskFeedbackService.java` | ACS 反馈处理接口 | +| `nl-module-task-server/src/main/java/cn/code/nl/module/task/service/transporttask/TransportTaskFeedbackServiceImpl.java` | ACS 反馈处理实现 | +| `nl-module-task-server/src/main/java/cn/code/nl/module/task/service/transporttask/TransportTaskCallbackService.java` | 业务回调处理接口 | +| `nl-module-task-server/src/main/java/cn/code/nl/module/task/service/transporttask/TransportTaskCallbackServiceImpl.java` | 业务回调处理实现 | +| `nl-module-task-server/src/main/java/cn/code/nl/module/task/client/TransportTaskStatusClient.java` | HTTP 同步回调 LMS/WMS 客户端 | +| `nl-module-task-server/src/main/java/cn/code/nl/module/task/mq/TaskEventProducer.java` | MQ 完成/取消事件生产者 | + +### 修改文件 + +| 文件 | 变更 | +|------|------| +| `nl-module-task-api/.../enums/ErrorCodeConstants.java` | 新增错误码 | +| `nl-module-task-server/.../service/transporttask/TransportTaskService.java` | 新增 `createTransportTaskByRpc` 方法 | +| `nl-module-task-server/.../service/transporttask/TransportTaskServiceImpl.java` | 实现 `createTransportTaskByRpc`(含幂等) | +| `nl-module-task-server/.../dal/mysql/transporttask/TransportTaskMapper.java` | 新增幂等查询方法、条件更新方法 | + +--- + +### 任务 1:API 模块 — 枚举与错误码 + +**文件:** +- 创建:`nl-module-task-api/src/main/java/cn/code/nl/module/task/enums/CallbackStatusEnum.java` +- 创建:`nl-module-task-api/src/main/java/cn/code/nl/module/task/enums/TaskEventTypeEnum.java` +- 修改:`nl-module-task-api/src/main/java/cn/code/nl/module/task/enums/ErrorCodeConstants.java` + +- [ ] **步骤 1:创建 CallbackStatusEnum** + +```java +package cn.code.nl.module.task.enums; + +import lombok.Getter; + +/** + * 业务回调状态枚举 + * @Author: liyongde + * @Date: 2026/7/14 + */ +@Getter +public enum CallbackStatusEnum { + + PENDING("PENDING", "待回调"), + SUCCESS("SUCCESS", "回调成功"), + FAILED("FAILED", "回调失败"); + + private final String code; + private final String name; + + CallbackStatusEnum(String code, String name) { + this.code = code; + this.name = name; + } +} +``` + +- [ ] **步骤 2:创建 TaskEventTypeEnum** + +```java +package cn.code.nl.module.task.enums; + +import lombok.Getter; + +/** + * 任务事件类型枚举 + * @Author: liyongde + * @Date: 2026/7/14 + */ +@Getter +public enum TaskEventTypeEnum { + + TASK_FINISHED("TASK_FINISHED", "任务完成"), + TASK_CANCELLED("TASK_CANCELLED", "任务取消"); + + private final String code; + private final String name; + + TaskEventTypeEnum(String code, String name) { + this.code = code; + this.name = name; + } +} +``` + +- [ ] **步骤 3:补充 ErrorCodeConstants** + +在 `ErrorCodeConstants.java` 中新增错误码(在已有常量下方追加): + +```java +ErrorCode TRANSPORT_TASK_EXISTS_UNFINISHED = new ErrorCode(2, "同业务类型下存在未完结的搬运任务"); +ErrorCode TRANSPORT_TASK_STATUS_NOT_ALLOW = new ErrorCode(3, "当前任务状态不允许此操作"); +ErrorCode TRANSPORT_TASK_CALLBACK_SEND_FAILED = new ErrorCode(4, "业务回调发送失败"); +``` + +- [ ] **步骤 4:Commit** + +```bash +git add nl-module-task-api/src/main/java/cn/code/nl/module/task/enums/CallbackStatusEnum.java \ + nl-module-task-api/src/main/java/cn/code/nl/module/task/enums/TaskEventTypeEnum.java \ + nl-module-task-api/src/main/java/cn/code/nl/module/task/enums/ErrorCodeConstants.java +git commit -m "feat: 新增回调状态枚举、事件类型枚举及错误码" +``` + +--- + +### 任务 2:API 模块 — 请求/响应 DTO + +**文件:** +- 创建:`nl-module-task-api/src/main/java/cn/code/nl/module/task/dto/TransportTaskCreateReqDTO.java` +- 创建:`nl-module-task-api/src/main/java/cn/code/nl/module/task/dto/AcsFeedbackReqDTO.java` +- 创建:`nl-module-task-api/src/main/java/cn/code/nl/module/task/dto/TaskStatusCallbackReqDTO.java` +- 创建:`nl-module-task-api/src/main/java/cn/code/nl/module/task/dto/TaskStatusCallbackRespDTO.java` +- 创建:`nl-module-task-api/src/main/java/cn/code/nl/module/task/dto/TaskCallbackResultReqDTO.java` +- 创建:`nl-module-task-api/src/main/java/cn/code/nl/module/task/dto/TaskEventMessage.java` + +- [ ] **步骤 1:创建 TransportTaskCreateReqDTO** + +```java +package cn.code.nl.module.task.dto; + +import io.swagger.v3.oas.annotations.media.Schema; +import jakarta.validation.constraints.NotEmpty; +import lombok.Data; + +import java.util.Map; + +/** + * LMS/WMS 创建搬运任务请求 DTO + */ +@Schema(description = "RPC - 搬运任务创建 Request DTO") +@Data +public class TransportTaskCreateReqDTO { + + @Schema(description = "任务名称", requiredMode = Schema.RequiredMode.REQUIRED) + @NotEmpty(message = "任务名称不能为空") + private String taskName; + + @Schema(description = "业务归属服务:LMS/WMS", requiredMode = Schema.RequiredMode.REQUIRED) + @NotEmpty(message = "业务归属服务不能为空") + private String ownerService; + + @Schema(description = "业务类型", requiredMode = Schema.RequiredMode.REQUIRED) + @NotEmpty(message = "业务类型不能为空") + private String bizType; + + @Schema(description = "业务侧标识", requiredMode = Schema.RequiredMode.REQUIRED) + @NotEmpty(message = "业务侧标识不能为空") + private String bizId; + + @Schema(description = "业务回调处理器编码", requiredMode = Schema.RequiredMode.REQUIRED) + @NotEmpty(message = "业务回调处理器编码不能为空") + private String handleCode; + + @Schema(description = "ACS任务类型", requiredMode = Schema.RequiredMode.REQUIRED) + @NotEmpty(message = "ACS任务类型不能为空") + private String acsTaskType; + + @Schema(description = "AGV系统类型", requiredMode = Schema.RequiredMode.REQUIRED) + @NotEmpty(message = "AGV系统类型不能为空") + private String agvSystemType; + + @Schema(description = "取货点1") + private String pointCode1; + + @Schema(description = "放货点1") + private String pointCode2; + + @Schema(description = "取货点2") + private String pointCode3; + + @Schema(description = "放货点2") + private String pointCode4; + + @Schema(description = "载具编码") + private String vehicleCode; + + @Schema(description = "载具编码2") + private String vehicleCode2; + + @Schema(description = "优先级") + private String priority; + + @Schema(description = "生产区域") + private String productArea; + + @Schema(description = "是否自动下发") + private String isAutoIssue; + + @Schema(description = "任务类型") + private String taskType; + + @Schema(description = "载具类型") + private String vehicleType; + + @Schema(description = "载具数量") + private Long vehicleQty; + + @Schema(description = "车号") + private String carNo; + + @Schema(description = "任务组标识") + private Long taskGroupId; + + @Schema(description = "任务组顺序号") + private Long sortSeq; + + @Schema(description = "扩展参数(JSON)") + private Map requestParam; +} +``` + +- [ ] **步骤 2:创建 AcsFeedbackReqDTO** + +```java +package cn.code.nl.module.task.dto; + +import io.swagger.v3.oas.annotations.media.Schema; +import jakarta.validation.constraints.NotEmpty; +import jakarta.validation.constraints.NotNull; +import lombok.Data; + +import java.util.Map; + +/** + * ACS 状态反馈请求 DTO + */ +@Schema(description = "ACS 状态反馈 Request DTO") +@Data +public class AcsFeedbackReqDTO { + + @Schema(description = "任务ID", requiredMode = Schema.RequiredMode.REQUIRED) + @NotNull(message = "任务ID不能为空") + private Long taskId; + + @Schema(description = "ACS反馈状态:EXECUTING/PICKED/FINISHED/CANCELLED", requiredMode = Schema.RequiredMode.REQUIRED) + @NotEmpty(message = "反馈状态不能为空") + private String status; + + @Schema(description = "事件ID(幂等键)") + private String eventId; + + @Schema(description = "ACS反馈扩展数据") + private Map payload; +} +``` + +- [ ] **步骤 3:创建 TaskStatusCallbackReqDTO** + +```java +package cn.code.nl.module.task.dto; + +import lombok.Data; + +import java.util.Map; + +/** + * Task→LMS/WMS 同步状态回调请求 DTO + * LMS/WMS 需依此 DTO 实现回调端点 + */ +@Data +public class TaskStatusCallbackReqDTO { + + /** 任务ID */ + private Long taskId; + + /** 任务编码 */ + private String taskCode; + + /** 回调状态:当前为 PICKED(61) */ + private String status; + + /** 业务归属服务 */ + private String ownerService; + + /** 业务类型 */ + private String bizType; + + /** 业务侧标识 */ + private String bizId; + + /** 业务回调处理器编码(LMS/WMS 内部路由) */ + private String handleCode; + + /** 扩展数据 */ + private Map payload; +} +``` + +- [ ] **步骤 4:创建 TaskStatusCallbackRespDTO** + +```java +package cn.code.nl.module.task.dto; + +import lombok.AllArgsConstructor; +import lombok.Data; +import lombok.NoArgsConstructor; + +/** + * Task→LMS/WMS 同步状态回调响应 DTO + */ +@Data +@NoArgsConstructor +@AllArgsConstructor +public class TaskStatusCallbackRespDTO { + + /** 是否成功 */ + private Boolean success; + + /** 响应消息 */ + private String message; + + public static TaskStatusCallbackRespDTO success() { + return new TaskStatusCallbackRespDTO(true, "ok"); + } + + public static TaskStatusCallbackRespDTO fail(String message) { + return new TaskStatusCallbackRespDTO(false, message); + } +} +``` + +- [ ] **步骤 5:创建 TaskCallbackResultReqDTO** + +```java +package cn.code.nl.module.task.dto; + +import io.swagger.v3.oas.annotations.media.Schema; +import jakarta.validation.constraints.NotEmpty; +import jakarta.validation.constraints.NotNull; +import lombok.Data; + +/** + * LMS/WMS 业务处理结果回调请求 DTO + */ +@Schema(description = "RPC - 业务结果回调 Request DTO") +@Data +public class TaskCallbackResultReqDTO { + + @Schema(description = "任务ID", requiredMode = Schema.RequiredMode.REQUIRED) + @NotNull(message = "任务ID不能为空") + private Long taskId; + + @Schema(description = "事件ID", requiredMode = Schema.RequiredMode.REQUIRED) + @NotEmpty(message = "事件ID不能为空") + private String eventId; + + @Schema(description = "处理结果:SUCCESS/FAILED", requiredMode = Schema.RequiredMode.REQUIRED) + @NotEmpty(message = "处理结果不能为空") + private String result; + + @Schema(description = "结果描述") + private String message; +} +``` + +- [ ] **步骤 6:创建 TaskEventMessage(MQ 消息体)** + +```java +package cn.code.nl.module.task.dto; + +import lombok.Data; + +import java.util.Map; + +/** + * 任务事件 MQ 消息体 + */ +@Data +public class TaskEventMessage { + + /** 事件ID */ + private String eventId; + + /** 事件类型:TASK_FINISHED / TASK_CANCELLED */ + private String eventType; + + /** 任务ID */ + private Long taskId; + + /** 任务编码 */ + private String taskCode; + + /** 业务归属服务:LMS / WMS */ + private String ownerService; + + /** 业务类型 */ + private String bizType; + + /** 业务侧标识 */ + private String bizId; + + /** 业务回调处理器编码 */ + private String handleCode; + + /** 扩展数据 */ + private Map payload; +} +``` + +- [ ] **步骤 7:Commit** + +```bash +git add nl-module-task-api/src/main/java/cn/code/nl/module/task/dto/ +git commit -m "feat: 新增 task 服务核心 DTO(创建/反馈/回调/消息)" +``` + +--- + +### 任务 3:TransportTaskMapper — 新增幂等查询与条件更新 + +**文件:** +- 修改:`nl-module-task-server/src/main/java/cn/code/nl/module/task/dal/mysql/transporttask/TransportTaskMapper.java` + +- [ ] **步骤 1:新增幂等查询和条件更新方法** + +在 `TransportTaskMapper` 中追加以下方法: + +```java +/** + * 按业务归属查询未完结任务(幂等校验用) + */ +default TransportTaskDO selectUnfinishedByBiz(String bizType, String bizId) { + return selectOne(new LambdaQueryWrapperX() + .eq(TransportTaskDO::getBizType, bizType) + .eq(TransportTaskDO::getBizId, bizId) + .lt(TransportTaskDO::getTaskStatus, "067")); +} + +/** + * 条件更新任务状态(CAS 抢占) + */ +default boolean updateTaskStatus(Long taskId, String expectedStatus, String newStatus) { + LambdaUpdateWrapper wrapper = new LambdaUpdateWrapper<>(); + wrapper.eq(TransportTaskDO::getTaskId, taskId) + .eq(TransportTaskDO::getTaskStatus, expectedStatus) + .set(TransportTaskDO::getTaskStatus, newStatus); + return update(null, wrapper) > 0; +} +``` + +注:需在文件头部追加 import: + +```java +import com.baomidou.mybatisplus.core.conditions.update.LambdaUpdateWrapper; +``` + +- [ ] **步骤 2:Commit** + +```bash +git add nl-module-task-server/src/main/java/cn/code/nl/module/task/dal/mysql/transporttask/TransportTaskMapper.java +git commit -m "feat: TransportTaskMapper 新增幂等查询与条件更新方法" +``` + +--- + +### 任务 4:TransportTaskService — 扩展 RPC 创建方法 + +**文件:** +- 修改:`nl-module-task-server/src/main/java/cn/code/nl/module/task/service/transporttask/TransportTaskService.java` +- 修改:`nl-module-task-server/src/main/java/cn/code/nl/module/task/service/transporttask/TransportTaskServiceImpl.java` + +- [ ] **步骤 1:编写 TransportTaskServiceImpl 单元测试(先写测试)** + +创建 `nl-module-task-server/src/test/java/cn/code/nl/module/task/service/transporttask/TransportTaskServiceImplTest.java`: + +```java +package cn.code.nl.module.task.service.transporttask; + +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.TransportTaskCreateReqDTO; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.InjectMocks; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; + +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.*; + +@ExtendWith(MockitoExtension.class) +class TransportTaskServiceImplTest { + + @Mock + private TransportTaskMapper transportTaskMapper; + + @InjectMocks + private TransportTaskServiceImpl transportTaskService; + + @Test + void testCreateTransportTaskByRpc_shouldReturnTaskId() { + // given + TransportTaskCreateReqDTO req = new TransportTaskCreateReqDTO(); + req.setTaskName("测试任务"); + req.setOwnerService("LMS"); + req.setBizType("INBOUND"); + req.setBizId("INB-001"); + req.setHandleCode("inboundHandler"); + req.setAcsTaskType("MOVE"); + req.setAgvSystemType("AGV_TYPE_A"); + + when(transportTaskMapper.selectUnfinishedByBiz("INBOUND", "INB-001")).thenReturn(null); + when(transportTaskMapper.insert(any(TransportTaskDO.class))).thenAnswer(inv -> { + TransportTaskDO task = inv.getArgument(0); + task.setTaskId(1L); + return 1; + }); + + // when + Long taskId = transportTaskService.createTransportTaskByRpc(req); + + // then + assertNotNull(taskId); + verify(transportTaskMapper, times(1)).insert(any(TransportTaskDO.class)); + } + + @Test + void testCreateTransportTaskByRpc_shouldReturnExistingWhenDuplicated() { + // given + TransportTaskCreateReqDTO req = new TransportTaskCreateReqDTO(); + req.setTaskName("测试任务"); + req.setOwnerService("LMS"); + req.setBizType("INBOUND"); + req.setBizId("INB-001"); + req.setHandleCode("inboundHandler"); + req.setAcsTaskType("MOVE"); + req.setAgvSystemType("AGV_TYPE_A"); + + TransportTaskDO existing = TransportTaskDO.builder().taskId(99L).taskStatus("040").build(); + when(transportTaskMapper.selectUnfinishedByBiz("INBOUND", "INB-001")).thenReturn(existing); + + // when + Long taskId = transportTaskService.createTransportTaskByRpc(req); + + // then + assertNotNull(taskId); + org.assertj.core.api.Assertions.assertThat(taskId).isEqualTo(99L); + verify(transportTaskMapper, never()).insert(any()); + } +} +``` + +- [ ] **步骤 2:运行测试确认失败** + +```bash +cd nl-module-task && mvn test -pl nl-module-task-server -Dtest=TransportTaskServiceImplTest -DfailIfNoTests=false +``` + +- [ ] **步骤 3:在 TransportTaskService 接口中新增方法签名** + +```java +/** + * 通过 RPC 创建搬运任务(LMS/WMS 调用) + * + * @param reqDTO 创建请求 + * @return 任务ID + */ +Long createTransportTaskByRpc(@Valid TransportTaskCreateReqDTO reqDTO); +``` + +注:需在文件头部追加 import: + +```java +import cn.code.nl.module.task.dto.TransportTaskCreateReqDTO; +``` + +- [ ] **步骤 4:在 TransportTaskServiceImpl 中实现方法** + +```java +@Override +public Long createTransportTaskByRpc(TransportTaskCreateReqDTO reqDTO) { + // 1. 幂等检查:同 bizType + bizId 下是否有未完结任务 + TransportTaskDO existing = transportTaskMapper.selectUnfinishedByBiz( + reqDTO.getBizType(), reqDTO.getBizId()); + if (existing != null) { + return existing.getTaskId(); + } + + // 2. 构建 DO 并保存 + TransportTaskDO task = BeanUtils.toBean(reqDTO, TransportTaskDO.class); + // 参数完整 → 待下发(040),否则 → 生成(010) + boolean paramComplete = StrUtil.isNotBlank(reqDTO.getPointCode1()) + && StrUtil.isNotBlank(reqDTO.getPointCode2()); + task.setTaskStatus(paramComplete + ? String.format("%03d", TransportTaskStatusEnum.READY.getCode()) + : String.format("%03d", TransportTaskStatusEnum.CREATED.getCode())); + task.setCallbackStatus(CallbackStatusEnum.PENDING.getCode()); + task.setCallbackRetryCount(0); + task.setCreateMode("RPC"); + + transportTaskMapper.insert(task); + return task.getTaskId(); +} +``` + +注:需在文件头部追加 import: + +```java +import cn.code.nl.module.task.dto.TransportTaskCreateReqDTO; +import cn.code.nl.module.task.enums.CallbackStatusEnum; +import cn.code.nl.module.task.enums.TransportTaskStatusEnum; +import cn.hutool.core.util.StrUtil; +``` + +- [ ] **步骤 5:运行测试确认通过** + +```bash +cd nl-module-task && mvn test -pl nl-module-task-server -Dtest=TransportTaskServiceImplTest -DfailIfNoTests=false +``` + +预期:PASS + +- [ ] **步骤 6:Commit** + +```bash +git add nl-module-task-server/src/main/java/cn/code/nl/module/task/service/transporttask/TransportTaskService.java \ + nl-module-task-server/src/main/java/cn/code/nl/module/task/service/transporttask/TransportTaskServiceImpl.java \ + nl-module-task-server/src/test/java/cn/code/nl/module/task/service/transporttask/TransportTaskServiceImplTest.java +git commit -m "feat: TransportTaskService 新增 createTransportTaskByRpc 方法(含幂等)" +``` + +--- + +### 任务 5:TransportTaskApiImpl — RPC 创建任务接口 + +**文件:** +- 创建:`nl-module-task-server/src/main/java/cn/code/nl/module/task/api/TransportTaskApiImpl.java` + +- [ ] **步骤 1:创建 TransportTaskApiImpl** + +```java +package cn.code.nl.module.task.api; + +import cn.code.nl.framework.common.pojo.CommonResult; +import cn.code.nl.module.task.dto.TransportTaskCreateReqDTO; +import cn.code.nl.module.task.enums.ApiConstants; +import cn.code.nl.module.task.service.transporttask.TransportTaskService; +import jakarta.annotation.Resource; +import org.springframework.validation.annotation.Validated; +import org.springframework.web.bind.annotation.PostMapping; +import org.springframework.web.bind.annotation.RequestBody; +import org.springframework.web.bind.annotation.RequestMapping; +import org.springframework.web.bind.annotation.RestController; + +import jakarta.validation.Valid; + +import static cn.code.nl.framework.common.pojo.CommonResult.success; + +/** + * Task 服务 RPC 接口实现 + */ +@RestController +@RequestMapping(ApiConstants.PREFIX) +@Validated +public class TransportTaskApiImpl { + + @Resource + private TransportTaskService transportTaskService; + + @PostMapping("/transport/create") + public CommonResult createTransportTask(@Valid @RequestBody TransportTaskCreateReqDTO reqDTO) { + return success(transportTaskService.createTransportTaskByRpc(reqDTO)); + } +} +``` + +- [ ] **步骤 2:Commit** + +```bash +git add nl-module-task-server/src/main/java/cn/code/nl/module/task/api/TransportTaskApiImpl.java +git commit -m "feat: 新增 TransportTaskApiImpl RPC 创建任务接口" +``` + +--- + +### 任务 6:TransportTaskStatusClient — HTTP 同步回调客户端 + +**文件:** +- 创建:`nl-module-task-server/src/main/java/cn/code/nl/module/task/client/TransportTaskStatusClient.java` + +- [ ] **步骤 1:创建 TransportTaskStatusClient** + +```java +package cn.code.nl.module.task.client; + +import cn.code.nl.module.task.dto.TaskStatusCallbackReqDTO; +import cn.code.nl.module.task.dto.TaskStatusCallbackRespDTO; +import lombok.extern.slf4j.Slf4j; +import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.http.HttpEntity; +import org.springframework.http.HttpHeaders; +import org.springframework.http.MediaType; +import org.springframework.stereotype.Component; +import org.springframework.web.client.RestTemplate; + +import jakarta.annotation.Resource; + +/** + * LMS/WMS 同步状态回调 HTTP 客户端 + *

+ * 通过 Nacos 服务发现 + RestTemplate 调用目标服务 + */ +@Slf4j +@Component +public class TransportTaskStatusClient { + + @Resource + @Qualifier("loadBalancedRestTemplate") + private RestTemplate restTemplate; + + /** + * 同步回调 LMS/WMS,通知任务状态变更(取货完成) + * + * @param reqDTO 回调请求 + * @return 回调结果 + */ + public TaskStatusCallbackRespDTO notifyStatus(TaskStatusCallbackReqDTO reqDTO) { + String serviceName = resolveServiceName(reqDTO.getOwnerService()); + String url = "http://" + serviceName + "/rpc-api/transport-task/status-callback"; + + try { + HttpHeaders headers = new HttpHeaders(); + headers.setContentType(MediaType.APPLICATION_JSON); + HttpEntity request = new HttpEntity<>(reqDTO, headers); + + TaskStatusCallbackRespDTO resp = restTemplate.postForObject(url, request, TaskStatusCallbackRespDTO.class); + log.info("同步回调 {} 成功, taskId={}", serviceName, reqDTO.getTaskId()); + return resp != null ? resp : TaskStatusCallbackRespDTO.success(); + } catch (Exception e) { + log.error("同步回调 {} 失败, taskId={}, error={}", serviceName, reqDTO.getTaskId(), e.getMessage(), e); + return TaskStatusCallbackRespDTO.fail(e.getMessage()); + } + } + + /** + * 根据 ownerService 解析 Nacos 服务名 + * 约定:LMS → lms-server, WMS → wms-server + */ + private String resolveServiceName(String ownerService) { + if ("WMS".equalsIgnoreCase(ownerService)) { + return "wms-server"; + } + // 默认 LMS + return "lms-server"; + } +} +``` + +- [ ] **步骤 2:Commit** + +```bash +git add nl-module-task-server/src/main/java/cn/code/nl/module/task/client/TransportTaskStatusClient.java +git commit -m "feat: 新增 TransportTaskStatusClient 同步 HTTP 回调客户端" +``` + +--- + +### 任务 7:TaskEventProducer — MQ 事件生产者 + +**文件:** +- 创建:`nl-module-task-server/src/main/java/cn/code/nl/module/task/mq/TaskEventProducer.java` + +- [ ] **步骤 1:创建 TaskEventProducer** + +```java +package cn.code.nl.module.task.mq; + +import cn.code.nl.module.task.dto.TaskEventMessage; +import lombok.extern.slf4j.Slf4j; +import org.apache.rocketmq.spring.core.RocketMQTemplate; +import org.springframework.stereotype.Component; + +import jakarta.annotation.Resource; +import java.util.UUID; + +/** + * 任务事件 MQ 生产者 + *

+ * Topic: TASK_EVENT_TOPIC + * Tag: ownerService (LMS / WMS) + */ +@Slf4j +@Component +public class TaskEventProducer { + + private static final String TOPIC = "TASK_EVENT_TOPIC"; + + @Resource + private RocketMQTemplate rocketMQTemplate; + + /** + * 发布任务完成/取消事件 + * + * @param message 事件消息 + */ + public void publishEvent(TaskEventMessage message) { + // 补齐 eventId + if (message.getEventId() == null || message.getEventId().isEmpty()) { + message.setEventId(UUID.randomUUID().toString()); + } + + String destination = TOPIC + ":" + message.getOwnerService(); + try { + rocketMQTemplate.syncSend(destination, message); + log.info("MQ 发送成功, destination={}, taskId={}, eventType={}", + destination, message.getTaskId(), message.getEventType()); + } catch (Exception e) { + log.error("MQ 发送失败, destination={}, taskId={}, error={}", + destination, message.getTaskId(), e.getMessage(), e); + throw new RuntimeException("MQ 发送失败: " + e.getMessage(), e); + } + } +} +``` + +- [ ] **步骤 2:Commit** + +```bash +git add nl-module-task-server/src/main/java/cn/code/nl/module/task/mq/TaskEventProducer.java +git commit -m "feat: 新增 TaskEventProducer MQ 事件生产者" +``` + +--- + +### 任务 8:TransportTaskFeedbackService — ACS 反馈处理核心 + +**文件:** +- 创建:`nl-module-task-server/src/main/java/cn/code/nl/module/task/service/transporttask/TransportTaskFeedbackService.java` +- 创建:`nl-module-task-server/src/main/java/cn/code/nl/module/task/service/transporttask/TransportTaskFeedbackServiceImpl.java` + +- [ ] **步骤 1:创建 TransportTaskFeedbackService 接口** + +```java +package cn.code.nl.module.task.service.transporttask; + +import cn.code.nl.module.task.dto.AcsFeedbackReqDTO; +import jakarta.validation.Valid; + +/** + * ACS 反馈处理服务 + */ +public interface TransportTaskFeedbackService { + + /** + * 接收 ACS 状态反馈并处理 + * + * @param reqDTO 反馈请求 + */ + void receiveAcsFeedback(@Valid AcsFeedbackReqDTO reqDTO); +} +``` + +- [ ] **步骤 2:编写 TransportTaskFeedbackServiceImpl 单元测试** + +创建 `nl-module-task-server/src/test/java/cn/code/nl/module/task/service/transporttask/TransportTaskFeedbackServiceImplTest.java`: + +```java +package cn.code.nl.module.task.service.transporttask; + +import cn.code.nl.module.task.client.TransportTaskStatusClient; +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.mq.TaskEventProducer; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.InjectMocks; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; + +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.*; + +@ExtendWith(MockitoExtension.class) +class TransportTaskFeedbackServiceImplTest { + + @Mock + private TransportTaskMapper transportTaskMapper; + + @Mock + private TransportTaskStatusClient transportTaskStatusClient; + + @Mock + private TaskEventProducer taskEventProducer; + + @InjectMocks + private TransportTaskFeedbackServiceImpl feedbackService; + + @Test + void testExecutingFeedback_shouldOnlyUpdateStatus() { + // given + TransportTaskDO task = TransportTaskDO.builder() + .taskId(1L).taskCode("T001").taskStatus("040") + .ownerService("LMS").bizType("INBOUND").bizId("INB-001").handleCode("h1") + .build(); + when(transportTaskMapper.selectById(1L)).thenReturn(task); + + AcsFeedbackReqDTO req = new AcsFeedbackReqDTO(); + req.setTaskId(1L); + req.setStatus("EXECUTING"); + req.setEventId("evt-001"); + + // when + feedbackService.receiveAcsFeedback(req); + + // then: 只更新状态,不调 MQ,不调 HTTP 回调 + verify(transportTaskMapper).updateById(any(TransportTaskDO.class)); + verifyNoInteractions(transportTaskStatusClient); + verifyNoInteractions(taskEventProducer); + } + + @Test + void testFinishedFeedback_shouldPublishMq() { + // given + TransportTaskDO task = TransportTaskDO.builder() + .taskId(1L).taskCode("T001").taskStatus("060") + .ownerService("LMS").bizType("INBOUND").bizId("INB-001").handleCode("h1") + .build(); + when(transportTaskMapper.selectById(1L)).thenReturn(task); + + AcsFeedbackReqDTO req = new AcsFeedbackReqDTO(); + req.setTaskId(1L); + req.setStatus("FINISHED"); + req.setEventId("evt-002"); + + // when + feedbackService.receiveAcsFeedback(req); + + // then: 更新状态 + 发 MQ + verify(transportTaskMapper).updateById(any(TransportTaskDO.class)); + verify(taskEventProducer, times(1)).publishEvent(any()); + verifyNoInteractions(transportTaskStatusClient); + } +} +``` + +- [ ] **步骤 3:运行测试确认失败** + +```bash +cd nl-module-task && mvn test -pl nl-module-task-server -Dtest=TransportTaskFeedbackServiceImplTest -DfailIfNoTests=false +``` + +- [ ] **步骤 4:实现 TransportTaskFeedbackServiceImpl** + +```java +package cn.code.nl.module.task.service.transporttask; + +import cn.code.nl.module.task.client.TransportTaskStatusClient; +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.dto.TaskEventMessage; +import cn.code.nl.module.task.dto.TaskStatusCallbackReqDTO; +import cn.code.nl.module.task.dto.TaskStatusCallbackRespDTO; +import cn.code.nl.module.task.enums.CallbackStatusEnum; +import cn.code.nl.module.task.enums.TaskEventTypeEnum; +import cn.code.nl.module.task.enums.TransportTaskStatusEnum; +import cn.code.nl.module.task.mq.TaskEventProducer; +import cn.hutool.core.util.StrUtil; +import com.fasterxml.jackson.core.type.TypeReference; +import com.fasterxml.jackson.databind.ObjectMapper; +import lombok.extern.slf4j.Slf4j; +import org.springframework.stereotype.Service; +import org.springframework.validation.annotation.Validated; + +import jakarta.annotation.Resource; +import java.util.Map; + +import static cn.code.nl.framework.common.exception.util.ServiceExceptionUtil.exception; +import static cn.code.nl.module.task.enums.ErrorCodeConstants.*; + +/** + * ACS 反馈处理服务实现 + */ +@Slf4j +@Service +@Validated +public class TransportTaskFeedbackServiceImpl implements TransportTaskFeedbackService { + + @Resource + private TransportTaskMapper transportTaskMapper; + + @Resource + private TransportTaskStatusClient transportTaskStatusClient; + + @Resource + private TaskEventProducer taskEventProducer; + + private static final ObjectMapper OBJECT_MAPPER = new ObjectMapper(); + + @Override + public void receiveAcsFeedback(AcsFeedbackReqDTO reqDTO) { + // 1. 查询任务 + TransportTaskDO task = transportTaskMapper.selectById(reqDTO.getTaskId()); + if (task == null) { + log.warn("ACS 反馈 taskId 不存在, taskId={}, 直接返回成功", reqDTO.getTaskId()); + return; + } + + // 2. 幂等判断:已经是终态(完成/取消)则不重复处理 + String currentStatus = task.getTaskStatus(); + if (isFinalStatus(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. 按状态路由 + String feedbackStatus = reqDTO.getStatus(); + if ("EXECUTING".equalsIgnoreCase(feedbackStatus)) { + handleExecuting(task); + + } else if ("PICKED".equalsIgnoreCase(feedbackStatus)) { + handlePicked(task, reqDTO); + + } else if ("FINISHED".equalsIgnoreCase(feedbackStatus)) { + handleFinished(task, reqDTO); + + } else if ("CANCELLED".equalsIgnoreCase(feedbackStatus)) { + handleCancelled(task, reqDTO); + + } else { + log.warn("未知 ACS 反馈状态, taskId={}, status={}", reqDTO.getTaskId(), feedbackStatus); + } + } + + // ---------- 各状态处理方法 ---------- + + private void handleExecuting(TransportTaskDO task) { + // 前置状态校验 + if (!isAllowedFrom(task.getTaskStatus(), "040", "050", "060")) { + log.warn("执行中反馈状态校验不通过, taskId={}, currentStatus={}", task.getTaskId(), task.getTaskStatus()); + return; + } + task.setTaskStatus(statusCode(TransportTaskStatusEnum.EXECUTING)); + transportTaskMapper.updateById(task); + log.info("任务执行中, taskId={}", task.getTaskId()); + } + + private void handlePicked(TransportTaskDO task, AcsFeedbackReqDTO reqDTO) { + if (!isAllowedFrom(task.getTaskStatus(), "060", "061")) { + log.warn("取货完成反馈状态校验不通过, taskId={}, currentStatus={}", task.getTaskId(), task.getTaskStatus()); + return; + } + // 更新状态 + task.setTaskStatus(statusCode(TransportTaskStatusEnum.PICKED)); + transportTaskMapper.updateById(task); + log.info("任务已取货, taskId={}", task.getTaskId()); + + // 同步 HTTP 回调 LMS/WMS + TaskStatusCallbackReqDTO callbackReq = buildCallbackReq(task, reqDTO); + TaskStatusCallbackRespDTO resp = transportTaskStatusClient.notifyStatus(callbackReq); + + if (resp != null && Boolean.TRUE.equals(resp.getSuccess())) { + log.info("取货完成同步回调成功, taskId={}", task.getTaskId()); + } else { + // 回调失败:记录但不阻塞 ACS + task.setCallbackStatus(CallbackStatusEnum.FAILED.getCode()); + task.setCallbackErrorMsg(resp != null ? resp.getMessage() : "回调无响应"); + transportTaskMapper.updateById(task); + log.warn("取货完成同步回调失败, taskId={}, msg={}", task.getTaskId(), task.getCallbackErrorMsg()); + } + } + + private void handleFinished(TransportTaskDO task, AcsFeedbackReqDTO reqDTO) { + if (!isAllowedFrom(task.getTaskStatus(), "050", "060", "061", "067")) { + log.warn("完成反馈状态校验不通过, taskId={}, currentStatus={}", task.getTaskId(), task.getTaskStatus()); + return; + } + // 已是 067 → 只发 MQ,不改状态 + if (!statusCode(TransportTaskStatusEnum.FINISHED_CALLBACK_PENDING).equals(task.getTaskStatus())) { + task.setTaskStatus(statusCode(TransportTaskStatusEnum.FINISHED_CALLBACK_PENDING)); + } + task.setCallbackStatus(CallbackStatusEnum.PENDING.getCode()); + transportTaskMapper.updateById(task); + + // 发布 MQ + publishEvent(task, TaskEventTypeEnum.TASK_FINISHED, reqDTO.getPayload()); + } + + private void handleCancelled(TransportTaskDO task, AcsFeedbackReqDTO reqDTO) { + if (!isAllowedFrom(task.getTaskStatus(), "040", "050", "060", "061", "069")) { + log.warn("取消反馈状态校验不通过, taskId={}, currentStatus={}", task.getTaskId(), task.getTaskStatus()); + return; + } + if (!statusCode(TransportTaskStatusEnum.CANCEL_CALLBACK_PENDING).equals(task.getTaskStatus())) { + task.setTaskStatus(statusCode(TransportTaskStatusEnum.CANCEL_CALLBACK_PENDING)); + } + task.setCallbackStatus(CallbackStatusEnum.PENDING.getCode()); + transportTaskMapper.updateById(task); + + // 发布 MQ + publishEvent(task, TaskEventTypeEnum.TASK_CANCELLED, reqDTO.getPayload()); + } + + // ---------- 工具方法 ---------- + + 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.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); + } + } + + private TaskStatusCallbackReqDTO buildCallbackReq(TransportTaskDO task, AcsFeedbackReqDTO reqDTO) { + TaskStatusCallbackReqDTO req = new TaskStatusCallbackReqDTO(); + req.setTaskId(task.getTaskId()); + req.setTaskCode(task.getTaskCode()); + req.setStatus(statusCode(TransportTaskStatusEnum.PICKED)); + req.setOwnerService(task.getOwnerService()); + req.setBizType(task.getBizType()); + req.setBizId(task.getBizId()); + req.setHandleCode(task.getHandleCode()); + req.setPayload(reqDTO.getPayload()); + return req; + } + + private boolean isFinalStatus(String status) { + return statusCode(TransportTaskStatusEnum.FINISHED).equals(status) + || statusCode(TransportTaskStatusEnum.CANCELLED).equals(status); + } + + private boolean isAllowedFrom(String current, String... allowed) { + for (String s : allowed) { + if (s.equals(current)) return true; + } + return false; + } + + private String statusCode(TransportTaskStatusEnum statusEnum) { + return String.format("%03d", statusEnum.getCode()); + } +} +``` + +注:需在文件头部追加 import: + +```java +import java.util.Map; +``` + +- [ ] **步骤 5:运行测试确认通过** + +```bash +cd nl-module-task && mvn test -pl nl-module-task-server -Dtest=TransportTaskFeedbackServiceImplTest -DfailIfNoTests=false +``` + +预期:PASS + +- [ ] **步骤 6:Commit** + +```bash +git add nl-module-task-server/src/main/java/cn/code/nl/module/task/service/transporttask/TransportTaskFeedbackService.java \ + nl-module-task-server/src/main/java/cn/code/nl/module/task/service/transporttask/TransportTaskFeedbackServiceImpl.java \ + nl-module-task-server/src/test/java/cn/code/nl/module/task/service/transporttask/TransportTaskFeedbackServiceImplTest.java +git commit -m "feat: 新增 TransportTaskFeedbackService ACS 反馈处理核心服务" +``` + +--- + +### 任务 9:AcsFeedbackController — ACS 反馈 HTTP 入口 + +**文件:** +- 创建:`nl-module-task-server/src/main/java/cn/code/nl/module/task/controller/AcsFeedbackController.java` + +- [ ] **步骤 1:创建 AcsFeedbackController** + +```java +package cn.code.nl.module.task.controller; + +import cn.code.nl.framework.common.pojo.CommonResult; +import cn.code.nl.module.task.dto.AcsFeedbackReqDTO; +import cn.code.nl.module.task.service.transporttask.TransportTaskFeedbackService; +import io.swagger.v3.oas.annotations.Operation; +import io.swagger.v3.oas.annotations.tags.Tag; +import jakarta.annotation.Resource; +import jakarta.validation.Valid; +import org.springframework.validation.annotation.Validated; +import org.springframework.web.bind.annotation.PostMapping; +import org.springframework.web.bind.annotation.RequestBody; +import org.springframework.web.bind.annotation.RequestMapping; +import org.springframework.web.bind.annotation.RestController; + +import static cn.code.nl.framework.common.pojo.CommonResult.success; + +/** + * ACS 状态反馈接口 + *

+ * 接收 ACS 的任务状态回调,处理执行中/取货完成/完成/取消 + */ +@Tag(name = "ACS 反馈") +@RestController +@RequestMapping("/api/task/transport") +@Validated +public class AcsFeedbackController { + + @Resource + private TransportTaskFeedbackService feedbackService; + + @PostMapping("/acs-feedback") + @Operation(summary = "接收 ACS 状态反馈") + public CommonResult receiveAcsFeedback(@Valid @RequestBody AcsFeedbackReqDTO reqDTO) { + feedbackService.receiveAcsFeedback(reqDTO); + // 始终返回成功,避免 ACS 因业务异常而重试 + return success(true); + } +} +``` + +- [ ] **步骤 2:Commit** + +```bash +git add nl-module-task-server/src/main/java/cn/code/nl/module/task/controller/AcsFeedbackController.java +git commit -m "feat: 新增 AcsFeedbackController ACS 状态反馈接口" +``` + +--- + +### 任务 10:TransportTaskCallbackService — 业务回调结果处理 + +**文件:** +- 创建:`nl-module-task-server/src/main/java/cn/code/nl/module/task/service/transporttask/TransportTaskCallbackService.java` +- 创建:`nl-module-task-server/src/main/java/cn/code/nl/module/task/service/transporttask/TransportTaskCallbackServiceImpl.java` + +- [ ] **步骤 1:创建 TransportTaskCallbackService 接口** + +```java +package cn.code.nl.module.task.service.transporttask; + +import cn.code.nl.module.task.dto.TaskCallbackResultReqDTO; +import jakarta.validation.Valid; + +/** + * 业务回调结果处理服务 + */ +public interface TransportTaskCallbackService { + + /** + * 接收 LMS/WMS 的业务处理结果 + * + * @param reqDTO 回调结果 + */ + void receiveCallbackResult(@Valid TaskCallbackResultReqDTO reqDTO); +} +``` + +- [ ] **步骤 2:编写单元测试** + +创建 `nl-module-task-server/src/test/java/cn/code/nl/module/task/service/transporttask/TransportTaskCallbackServiceImplTest.java`: + +```java +package cn.code.nl.module.task.service.transporttask; + +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.TaskCallbackResultReqDTO; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.ArgumentCaptor; +import org.mockito.InjectMocks; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.mockito.Mockito.*; + +@ExtendWith(MockitoExtension.class) +class TransportTaskCallbackServiceImplTest { + + @Mock + private TransportTaskMapper transportTaskMapper; + + @InjectMocks + private TransportTaskCallbackServiceImpl callbackService; + + @Test + void testFinishSuccess_shouldUpdateToFinished() { + TransportTaskDO task = TransportTaskDO.builder() + .taskId(1L).taskStatus("067").callbackStatus("PENDING") + .build(); + when(transportTaskMapper.selectById(1L)).thenReturn(task); + + TaskCallbackResultReqDTO req = new TaskCallbackResultReqDTO(); + req.setTaskId(1L); + req.setEventId("evt-001"); + req.setResult("SUCCESS"); + + callbackService.receiveCallbackResult(req); + + ArgumentCaptor captor = ArgumentCaptor.forClass(TransportTaskDO.class); + verify(transportTaskMapper).updateById(captor.capture()); + TransportTaskDO updated = captor.getValue(); + assertEquals("070", updated.getTaskStatus()); + assertEquals("SUCCESS", updated.getCallbackStatus()); + } + + @Test + void testAlreadyFinal_shouldReturnDirectly() { + TransportTaskDO task = TransportTaskDO.builder() + .taskId(1L).taskStatus("070").callbackStatus("SUCCESS") + .build(); + when(transportTaskMapper.selectById(1L)).thenReturn(task); + + TaskCallbackResultReqDTO req = new TaskCallbackResultReqDTO(); + req.setTaskId(1L); + req.setEventId("evt-001"); + req.setResult("SUCCESS"); + + callbackService.receiveCallbackResult(req); + + // 已终态,不应再更新 + verify(transportTaskMapper, never()).updateById(any()); + } +} +``` + +- [ ] **步骤 3:运行测试确认失败** + +```bash +cd nl-module-task && mvn test -pl nl-module-task-server -Dtest=TransportTaskCallbackServiceImplTest -DfailIfNoTests=false +``` + +- [ ] **步骤 4:实现 TransportTaskCallbackServiceImpl** + +```java +package cn.code.nl.module.task.service.transporttask; + +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.TaskCallbackResultReqDTO; +import cn.code.nl.module.task.enums.CallbackStatusEnum; +import cn.code.nl.module.task.enums.TransportTaskStatusEnum; +import lombok.extern.slf4j.Slf4j; +import org.springframework.stereotype.Service; +import org.springframework.validation.annotation.Validated; + +import jakarta.annotation.Resource; + +/** + * 业务回调结果处理服务实现 + */ +@Slf4j +@Service +@Validated +public class TransportTaskCallbackServiceImpl implements TransportTaskCallbackService { + + @Resource + private TransportTaskMapper transportTaskMapper; + + @Override + public void receiveCallbackResult(TaskCallbackResultReqDTO reqDTO) { + TransportTaskDO task = transportTaskMapper.selectById(reqDTO.getTaskId()); + if (task == null) { + log.warn("回调 taskId 不存在, taskId={}", reqDTO.getTaskId()); + return; + } + + // 幂等:已是终态则直接返回 + String currentStatus = task.getTaskStatus(); + String finalFinished = statusCode(TransportTaskStatusEnum.FINISHED); + String finalCancelled = statusCode(TransportTaskStatusEnum.CANCELLED); + if (finalFinished.equals(currentStatus) || finalCancelled.equals(currentStatus)) { + log.info("任务已处终态, taskId={}, status={}, 忽略回调", reqDTO.getTaskId(), currentStatus); + return; + } + + boolean success = "SUCCESS".equalsIgnoreCase(reqDTO.getResult()); + + if (finalFinished.equals(currentStatus) != true) { + // 当前是 067(完成待业务处理)或 069(取消待业务处理) + if (success) { + if (statusCode(TransportTaskStatusEnum.FINISHED_CALLBACK_PENDING).equals(currentStatus)) { + task.setTaskStatus(finalFinished); + } else if (statusCode(TransportTaskStatusEnum.CANCEL_CALLBACK_PENDING).equals(currentStatus)) { + task.setTaskStatus(finalCancelled); + } + task.setCallbackStatus(CallbackStatusEnum.SUCCESS.getCode()); + } else { + task.setCallbackStatus(CallbackStatusEnum.FAILED.getCode()); + task.setCallbackErrorMsg(reqDTO.getMessage()); + // 增加重试次数 + task.setCallbackRetryCount(task.getCallbackRetryCount() != null + ? task.getCallbackRetryCount() + 1 : 1); + // 状态回退,等待下次重试 + if (statusCode(TransportTaskStatusEnum.FINISHED_CALLBACK_PENDING).equals(currentStatus)) { + // 保持 067 + } + } + + transportTaskMapper.updateById(task); + log.info("业务回调结果已处理, taskId={}, result={}", reqDTO.getTaskId(), reqDTO.getResult()); + } + } + + private String statusCode(TransportTaskStatusEnum statusEnum) { + return String.format("%03d", statusEnum.getCode()); + } +} +``` + +- [ ] **步骤 5:运行测试确认通过** + +```bash +cd nl-module-task && mvn test -pl nl-module-task-server -Dtest=TransportTaskCallbackServiceImplTest -DfailIfNoTests=false +``` + +预期:PASS + +- [ ] **步骤 6:Commit** + +```bash +git add nl-module-task-server/src/main/java/cn/code/nl/module/task/service/transporttask/TransportTaskCallbackService.java \ + nl-module-task-server/src/main/java/cn/code/nl/module/task/service/transporttask/TransportTaskCallbackServiceImpl.java \ + nl-module-task-server/src/test/java/cn/code/nl/module/task/service/transporttask/TransportTaskCallbackServiceImplTest.java +git commit -m "feat: 新增 TransportTaskCallbackService 业务回调结果处理" +``` + +--- + +### 任务 11:TaskCallbackController — 业务回调结果入口 + +**文件:** +- 创建:`nl-module-task-server/src/main/java/cn/code/nl/module/task/controller/TaskCallbackController.java` + +- [ ] **步骤 1:创建 TaskCallbackController** + +```java +package cn.code.nl.module.task.controller; + +import cn.code.nl.framework.common.pojo.CommonResult; +import cn.code.nl.module.task.dto.TaskCallbackResultReqDTO; +import cn.code.nl.module.task.enums.ApiConstants; +import cn.code.nl.module.task.service.transporttask.TransportTaskCallbackService; +import io.swagger.v3.oas.annotations.Operation; +import io.swagger.v3.oas.annotations.tags.Tag; +import jakarta.annotation.Resource; +import jakarta.validation.Valid; +import org.springframework.validation.annotation.Validated; +import org.springframework.web.bind.annotation.PostMapping; +import org.springframework.web.bind.annotation.RequestBody; +import org.springframework.web.bind.annotation.RequestMapping; +import org.springframework.web.bind.annotation.RestController; + +import static cn.code.nl.framework.common.pojo.CommonResult.success; + +/** + * LMS/WMS 业务回调结果接口 + */ +@Tag(name = "业务回调") +@RestController +@RequestMapping(ApiConstants.PREFIX) +@Validated +public class TaskCallbackController { + + @Resource + private TransportTaskCallbackService callbackService; + + @PostMapping("/transport/callback-result") + @Operation(summary = "接收 LMS/WMS 业务处理结果") + public CommonResult receiveCallbackResult(@Valid @RequestBody TaskCallbackResultReqDTO reqDTO) { + callbackService.receiveCallbackResult(reqDTO); + return success(true); + } +} +``` + +- [ ] **步骤 2:Commit** + +```bash +git add nl-module-task-server/src/main/java/cn/code/nl/module/task/controller/TaskCallbackController.java +git commit -m "feat: 新增 TaskCallbackController 业务回调结果接口" +``` + +--- + +### 任务 12:集成验证 — 启动服务并验证编译 + +**文件:** 无新建,整体编译验证。 + +- [ ] **步骤 1:编译全部模块** + +```bash +cd D:/Code/Work/nl/huachuang && mvn compile -pl nl-module-task -am -DskipTests +``` + +预期:BUILD SUCCESS + +- [ ] **步骤 2:运行全部 task 模块测试** + +```bash +cd D:/Code/Work/nl/huachuang && mvn test -pl nl-module-task -am +``` + +预期:所有测试通过 + +- [ ] **步骤 3:Commit(如有修正)** + +```bash +git add -A +git commit -m "chore: 编译验证与测试修复" +``` + +--- + +## 自检结果 + +### 1. 规格覆盖度 + +| 规格需求 | 对应任务 | +|----------|---------| +| LMS/WMS 创建任务 (RPC) | 任务 4、5 | +| 创建幂等(同 bizType+bizId) | 任务 3、4 | +| ACS 反馈:执行中 | 任务 8、9 | +| ACS 反馈:取货完成 + 同步 HTTP 回调 | 任务 6、8、9 | +| ACS 反馈:完成/取消 + MQ 通知 | 任务 7、8、9 | +| 业务回调结果接收(成功/失败) | 任务 10、11 | +| 回调结果幂等(已终态忽略) | 任务 10 | +| 状态前置校验 | 任务 8 | +| 回调失败记录重试信息 | 任务 8、10 | + +全部覆盖,无遗漏。 + +### 2. 占位符扫描 + +无 "TODO"、"待定"、"后续实现" 等占位符。 + +### 3. 类型一致性 + +- 所有 DTO 类名在接口、Service、测试中一致 +- 状态码统一使用 `String.format("%03d", enum.getCode())` 格式(如 "060"、"067") +- `TransportTaskDO` 字段与 DTO 映射通过 `BeanUtils.toBean()` 保持一致 +- 枚举值在创建(任务 1)、反馈(任务 8)、回调(任务 10)间一致