Files
huachuang/docs/superpowers/plans/2026-07-14-task-core-implementation.md

1583 lines
53 KiB
Markdown
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

# 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` | 新增幂等查询方法、条件更新方法 |
---
### 任务 1API 模块 — 枚举与错误码
**文件:**
- 创建:`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, "业务回调发送失败");
```
- [ ] **步骤 4Commit**
```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: 新增回调状态枚举、事件类型枚举及错误码"
```
---
### 任务 2API 模块 — 请求/响应 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<String, Object> 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<String, Object> 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<String, Object> 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创建 TaskEventMessageMQ 消息体)**
```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<String, Object> payload;
}
```
- [ ] **步骤 7Commit**
```bash
git add nl-module-task-api/src/main/java/cn/code/nl/module/task/dto/
git commit -m "feat: 新增 task 服务核心 DTO创建/反馈/回调/消息)"
```
---
### 任务 3TransportTaskMapper — 新增幂等查询与条件更新
**文件:**
- 修改:`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<TransportTaskDO>()
.eq(TransportTaskDO::getBizType, bizType)
.eq(TransportTaskDO::getBizId, bizId)
.lt(TransportTaskDO::getTaskStatus, "067"));
}
/**
* 条件更新任务状态CAS 抢占)
*/
default boolean updateTaskStatus(Long taskId, String expectedStatus, String newStatus) {
LambdaUpdateWrapper<TransportTaskDO> 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;
```
- [ ] **步骤 2Commit**
```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 新增幂等查询与条件更新方法"
```
---
### 任务 4TransportTaskService — 扩展 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
- [ ] **步骤 6Commit**
```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 方法(含幂等)"
```
---
### 任务 5TransportTaskApiImpl — 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<Long> createTransportTask(@Valid @RequestBody TransportTaskCreateReqDTO reqDTO) {
return success(transportTaskService.createTransportTaskByRpc(reqDTO));
}
}
```
- [ ] **步骤 2Commit**
```bash
git add nl-module-task-server/src/main/java/cn/code/nl/module/task/api/TransportTaskApiImpl.java
git commit -m "feat: 新增 TransportTaskApiImpl RPC 创建任务接口"
```
---
### 任务 6TransportTaskStatusClient — 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 客户端
* <p>
* 通过 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<TaskStatusCallbackReqDTO> 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";
}
}
```
- [ ] **步骤 2Commit**
```bash
git add nl-module-task-server/src/main/java/cn/code/nl/module/task/client/TransportTaskStatusClient.java
git commit -m "feat: 新增 TransportTaskStatusClient 同步 HTTP 回调客户端"
```
---
### 任务 7TaskEventProducer — 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.mq.message.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 生产者
* <p>
* 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);
}
}
}
```
- [ ] **步骤 2Commit**
```bash
git add nl-module-task-server/src/main/java/cn/code/nl/module/task/mq/TaskEventProducer.java
git commit -m "feat: 新增 TaskEventProducer MQ 事件生产者"
```
---
### 任务 8TransportTaskFeedbackService — 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.producer.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.mq.message.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.producer.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<String, Object> 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
- [ ] **步骤 6Commit**
```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 反馈处理核心服务"
```
---
### 任务 9AcsFeedbackController — 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 状态反馈接口
* <p>
* 接收 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<Boolean> receiveAcsFeedback(@Valid @RequestBody AcsFeedbackReqDTO reqDTO) {
feedbackService.receiveAcsFeedback(reqDTO);
// 始终返回成功,避免 ACS 因业务异常而重试
return success(true);
}
}
```
- [ ] **步骤 2Commit**
```bash
git add nl-module-task-server/src/main/java/cn/code/nl/module/task/controller/AcsFeedbackController.java
git commit -m "feat: 新增 AcsFeedbackController ACS 状态反馈接口"
```
---
### 任务 10TransportTaskCallbackService — 业务回调结果处理
**文件:**
- 创建:`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<TransportTaskDO> 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
- [ ] **步骤 6Commit**
```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 业务回调结果处理"
```
---
### 任务 11TaskCallbackController — 业务回调结果入口
**文件:**
- 创建:`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<Boolean> receiveCallbackResult(@Valid @RequestBody TaskCallbackResultReqDTO reqDTO) {
callbackService.receiveCallbackResult(reqDTO);
return success(true);
}
}
```
- [ ] **步骤 2Commit**
```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
```
预期:所有测试通过
- [ ] **步骤 3Commit如有修正**
```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间一致