feat: 完成task服务核心业务——ACS反馈处理、MQ/HTTP回调、业务结果接收
This commit is contained in:
@@ -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 状态反馈接口
|
||||
* <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);
|
||||
}
|
||||
}
|
||||
@@ -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<Boolean> receiveCallbackResult(@Valid @RequestBody TaskCallbackResultReqDTO reqDTO) {
|
||||
callbackService.receiveCallbackResult(reqDTO);
|
||||
return success(true);
|
||||
}
|
||||
}
|
||||
@@ -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 生产者
|
||||
* <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);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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);
|
||||
}
|
||||
@@ -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());
|
||||
}
|
||||
}
|
||||
@@ -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);
|
||||
}
|
||||
@@ -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<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);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 构建同步回调请求 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());
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user