diff --git a/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/controller/AcsFeedbackController.java b/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/controller/AcsFeedbackController.java new file mode 100644 index 00000000..6ca2d92e --- /dev/null +++ b/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/controller/AcsFeedbackController.java @@ -0,0 +1,39 @@ +package cn.code.nl.module.task.controller; + +import cn.code.nl.framework.common.pojo.CommonResult; +import cn.code.nl.module.task.dto.AcsFeedbackReqDTO; +import cn.code.nl.module.task.service.transporttask.TransportTaskFeedbackService; +import io.swagger.v3.oas.annotations.Operation; +import io.swagger.v3.oas.annotations.tags.Tag; +import jakarta.annotation.Resource; +import jakarta.validation.Valid; +import org.springframework.validation.annotation.Validated; +import org.springframework.web.bind.annotation.PostMapping; +import org.springframework.web.bind.annotation.RequestBody; +import org.springframework.web.bind.annotation.RequestMapping; +import org.springframework.web.bind.annotation.RestController; + +import static cn.code.nl.framework.common.pojo.CommonResult.success; + +/** + * ACS 状态反馈接口 + *

+ * 接收 ACS 的任务状态回调,处理执行中/取货完成/完成/取消 + */ +@Tag(name = "ACS 反馈") +@RestController +@RequestMapping("/api/task/transport") +@Validated +public class AcsFeedbackController { + + @Resource + private TransportTaskFeedbackService feedbackService; + + @PostMapping("/acs-feedback") + @Operation(summary = "接收 ACS 状态反馈") + public CommonResult receiveAcsFeedback(@Valid @RequestBody AcsFeedbackReqDTO reqDTO) { + feedbackService.receiveAcsFeedback(reqDTO); + // 始终返回成功,避免 ACS 因业务异常而重试 + return success(true); + } +} diff --git a/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/controller/TaskCallbackController.java b/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/controller/TaskCallbackController.java new file mode 100644 index 00000000..6c8fb27f --- /dev/null +++ b/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/controller/TaskCallbackController.java @@ -0,0 +1,37 @@ +package cn.code.nl.module.task.controller; + +import cn.code.nl.framework.common.pojo.CommonResult; +import cn.code.nl.module.task.dto.TaskCallbackResultReqDTO; +import cn.code.nl.module.task.enums.ApiConstants; +import cn.code.nl.module.task.service.transporttask.TransportTaskCallbackService; +import io.swagger.v3.oas.annotations.Operation; +import io.swagger.v3.oas.annotations.tags.Tag; +import jakarta.annotation.Resource; +import jakarta.validation.Valid; +import org.springframework.validation.annotation.Validated; +import org.springframework.web.bind.annotation.PostMapping; +import org.springframework.web.bind.annotation.RequestBody; +import org.springframework.web.bind.annotation.RequestMapping; +import org.springframework.web.bind.annotation.RestController; + +import static cn.code.nl.framework.common.pojo.CommonResult.success; + +/** + * LMS/WMS 业务回调结果接口 + */ +@Tag(name = "业务回调") +@RestController +@RequestMapping(ApiConstants.PREFIX) +@Validated +public class TaskCallbackController { + + @Resource + private TransportTaskCallbackService callbackService; + + @PostMapping("/transport/callback-result") + @Operation(summary = "接收 LMS/WMS 业务处理结果") + public CommonResult receiveCallbackResult(@Valid @RequestBody TaskCallbackResultReqDTO reqDTO) { + callbackService.receiveCallbackResult(reqDTO); + return success(true); + } +} diff --git a/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/mq/TaskEventProducer.java b/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/mq/TaskEventProducer.java new file mode 100644 index 00000000..ab119a52 --- /dev/null +++ b/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/mq/TaskEventProducer.java @@ -0,0 +1,48 @@ +package cn.code.nl.module.task.mq; + +import cn.code.nl.module.task.dto.TaskEventMessage; +import lombok.extern.slf4j.Slf4j; +import org.apache.rocketmq.spring.core.RocketMQTemplate; +import org.springframework.stereotype.Component; + +import jakarta.annotation.Resource; +import java.util.UUID; + +/** + * 任务事件 MQ 生产者 + *

+ * Topic: TASK_EVENT_TOPIC + * Tag: ownerService (LMS / WMS) + */ +@Slf4j +@Component +public class TaskEventProducer { + + private static final String TOPIC = "TASK_EVENT_TOPIC"; + + @Resource + private RocketMQTemplate rocketMQTemplate; + + /** + * 发布任务完成/取消事件 + * + * @param message 事件消息 + */ + public void publishEvent(TaskEventMessage message) { + // 补齐 eventId + if (message.getEventId() == null || message.getEventId().isEmpty()) { + message.setEventId(UUID.randomUUID().toString()); + } + + String destination = TOPIC + ":" + message.getOwnerService(); + try { + rocketMQTemplate.syncSend(destination, message); + log.info("MQ 发送成功, destination={}, taskId={}, eventType={}", + destination, message.getTaskId(), message.getEventType()); + } catch (Exception e) { + log.error("MQ 发送失败, destination={}, taskId={}, error={}", + destination, message.getTaskId(), e.getMessage(), e); + throw new RuntimeException("MQ 发送失败: " + e.getMessage(), e); + } + } +} diff --git a/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/service/transporttask/TransportTaskCallbackService.java b/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/service/transporttask/TransportTaskCallbackService.java new file mode 100644 index 00000000..6cfa1716 --- /dev/null +++ b/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/service/transporttask/TransportTaskCallbackService.java @@ -0,0 +1,19 @@ +package cn.code.nl.module.task.service.transporttask; + +import cn.code.nl.module.task.dto.TaskCallbackResultReqDTO; +import jakarta.validation.Valid; + +/** + * 业务回调结果处理服务接口 + * + * @author 诺力管理员 + */ +public interface TransportTaskCallbackService { + + /** + * 接收 LMS/WMS 的业务处理结果 + * + * @param reqDTO 回调结果 + */ + void receiveCallbackResult(@Valid TaskCallbackResultReqDTO reqDTO); +} diff --git a/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/service/transporttask/TransportTaskCallbackServiceImpl.java b/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/service/transporttask/TransportTaskCallbackServiceImpl.java new file mode 100644 index 00000000..1c5e53c4 --- /dev/null +++ b/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/service/transporttask/TransportTaskCallbackServiceImpl.java @@ -0,0 +1,72 @@ +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; + +/** + * 业务回调结果处理服务实现 + * + * @author 诺力管理员 + */ +@Slf4j +@Service +@Validated +public class TransportTaskCallbackServiceImpl implements TransportTaskCallbackService { + + @Resource + private TransportTaskMapper transportTaskMapper; + + @Override + public void receiveCallbackResult(TaskCallbackResultReqDTO reqDTO) { + // 1. 查询任务 + TransportTaskDO task = transportTaskMapper.selectById(reqDTO.getTaskId()); + if (task == null) { + log.warn("回调 taskId 不存在, taskId={}", reqDTO.getTaskId()); + return; + } + + // 2. 幂等:已是终态则直接返回 + 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; + } + + // 3. 根据结果处理 + boolean success = "SUCCESS".equalsIgnoreCase(reqDTO.getResult()); + 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); + } + + transportTaskMapper.updateById(task); + log.info("业务回调结果已处理, taskId={}, result={}, newStatus={}", + reqDTO.getTaskId(), reqDTO.getResult(), task.getTaskStatus()); + } + + private String statusCode(TransportTaskStatusEnum statusEnum) { + return String.format("%03d", statusEnum.getCode()); + } +} diff --git a/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/service/transporttask/TransportTaskFeedbackService.java b/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/service/transporttask/TransportTaskFeedbackService.java new file mode 100644 index 00000000..fc3b21db --- /dev/null +++ b/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/service/transporttask/TransportTaskFeedbackService.java @@ -0,0 +1,19 @@ +package cn.code.nl.module.task.service.transporttask; + +import cn.code.nl.module.task.dto.AcsFeedbackReqDTO; +import jakarta.validation.Valid; + +/** + * ACS 反馈处理服务接口 + * + * @author 诺力管理员 + */ +public interface TransportTaskFeedbackService { + + /** + * 接收 ACS 状态反馈并处理 + * + * @param reqDTO 反馈请求 + */ + void receiveAcsFeedback(@Valid AcsFeedbackReqDTO reqDTO); +} diff --git a/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/service/transporttask/TransportTaskFeedbackServiceImpl.java b/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/service/transporttask/TransportTaskFeedbackServiceImpl.java new file mode 100644 index 00000000..051d9be9 --- /dev/null +++ b/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/service/transporttask/TransportTaskFeedbackServiceImpl.java @@ -0,0 +1,239 @@ +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.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; + +/** + * ACS 反馈处理服务实现 + * + * @author 诺力管理员 + */ +@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); + } + } + + // ==================== 各状态处理方法 ==================== + + /** + * 处理执行中(60):直接更新状态,不做业务回调 + */ + 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()); + } + + /** + * 处理取货完成(61):更新状态 + 同步 HTTP 回调 LMS/WMS + */ + 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()); + } + } + + /** + * 处理完成(67):更新状态 + 异步 MQ 通知 LMS/WMS + */ + 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); + + publishEvent(task, TaskEventTypeEnum.TASK_FINISHED, reqDTO.getPayload()); + } + + /** + * 处理取消(69):更新状态 + 异步 MQ 通知 LMS/WMS + */ + 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); + + publishEvent(task, TaskEventTypeEnum.TASK_CANCELLED, reqDTO.getPayload()); + } + + // ==================== 工具方法 ==================== + + /** + * 发布 MQ 事件。发送失败时记录 callbackStatus=FAILED + */ + private void publishEvent(TransportTaskDO task, TaskEventTypeEnum eventType, + Map payload) { + TaskEventMessage msg = new TaskEventMessage(); + msg.setEventType(eventType.getCode()); + msg.setTaskId(task.getTaskId()); + msg.setTaskCode(task.getTaskCode()); + msg.setOwnerService(task.getOwnerService()); + msg.setBizType(task.getBizType()); + msg.setBizId(task.getBizId()); + msg.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); + } + } + + /** + * 构建同步回调请求 DTO + */ + 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()); + } +}