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

53 KiB
Raw Blame History

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

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
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 中新增错误码(在已有常量下方追加):

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
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

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
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
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
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
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 消息体)
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
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 中追加以下方法:

/**
 * 按业务归属查询未完结任务(幂等校验用)
 */
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

import com.baomidou.mybatisplus.core.conditions.update.LambdaUpdateWrapper;
  • 步骤 2Commit
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

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运行测试确认失败
cd nl-module-task && mvn test -pl nl-module-task-server -Dtest=TransportTaskServiceImplTest -DfailIfNoTests=false
  • 步骤 3在 TransportTaskService 接口中新增方法签名
/**
 * 通过 RPC 创建搬运任务LMS/WMS 调用)
 *
 * @param reqDTO 创建请求
 * @return 任务ID
 */
Long createTransportTaskByRpc(@Valid TransportTaskCreateReqDTO reqDTO);

注:需在文件头部追加 import

import cn.code.nl.module.task.dto.TransportTaskCreateReqDTO;
  • 步骤 4在 TransportTaskServiceImpl 中实现方法
@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

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运行测试确认通过
cd nl-module-task && mvn test -pl nl-module-task-server -Dtest=TransportTaskServiceImplTest -DfailIfNoTests=false

预期PASS

  • 步骤 6Commit
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

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
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

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
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

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 生产者
 * <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
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 接口

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

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运行测试确认失败
cd nl-module-task && mvn test -pl nl-module-task-server -Dtest=TransportTaskFeedbackServiceImplTest -DfailIfNoTests=false
  • 步骤 4实现 TransportTaskFeedbackServiceImpl
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<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

import java.util.Map;
  • 步骤 5运行测试确认通过
cd nl-module-task && mvn test -pl nl-module-task-server -Dtest=TransportTaskFeedbackServiceImplTest -DfailIfNoTests=false

预期PASS

  • 步骤 6Commit
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

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
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 接口

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

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运行测试确认失败
cd nl-module-task && mvn test -pl nl-module-task-server -Dtest=TransportTaskCallbackServiceImplTest -DfailIfNoTests=false
  • 步骤 4实现 TransportTaskCallbackServiceImpl
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运行测试确认通过
cd nl-module-task && mvn test -pl nl-module-task-server -Dtest=TransportTaskCallbackServiceImplTest -DfailIfNoTests=false

预期PASS

  • 步骤 6Commit
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

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
git add nl-module-task-server/src/main/java/cn/code/nl/module/task/controller/TaskCallbackController.java
git commit -m "feat: 新增 TaskCallbackController 业务回调结果接口"

任务 12集成验证 — 启动服务并验证编译

文件: 无新建,整体编译验证。

  • 步骤 1编译全部模块
cd D:/Code/Work/nl/huachuang && mvn compile -pl nl-module-task -am -DskipTests

预期BUILD SUCCESS

  • 步骤 2运行全部 task 模块测试
cd D:/Code/Work/nl/huachuang && mvn test -pl nl-module-task -am

预期:所有测试通过

  • 步骤 3Commit如有修正
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间一致