fix: ACS、PC对任务状态的变化

This commit is contained in:
2026-07-16 14:11:30 +08:00
parent f24e01c7dd
commit be46cf49c5
37 changed files with 2349 additions and 268 deletions

View File

@@ -1,6 +1,7 @@
package cn.code.nl.module.task.api;
import cn.code.nl.framework.common.pojo.CommonResult;
import cn.code.nl.module.task.dto.TaskInfoDTO;
import cn.code.nl.module.task.dto.TransportTaskCreateReqDTO;
import cn.code.nl.module.task.dto.TaskCallbackResultReqDTO;
import cn.code.nl.module.task.service.transporttask.TransportTaskCallbackService;
@@ -34,4 +35,14 @@ public class TransportTaskApiImpl implements TransportTaskApi {
callbackService.receiveCallbackResult(reqDTO);
return success(true);
}
@Override
public CommonResult<TaskInfoDTO> getTaskById(Long taskId) {
return success(transportTaskService.getTaskInfoById(taskId));
}
@Override
public CommonResult<TaskInfoDTO> getTaskByCode(String taskCode) {
return success(transportTaskService.getTaskInfoByCode(taskCode));
}
}

View File

@@ -16,6 +16,7 @@ import jakarta.annotation.Resource;
* LMS/WMS 同步状态回调 HTTP 客户端
* 通过 Nacos 服务发现 + RestTemplate 调用目标服务
*/
@Deprecated
@Slf4j
@Component
public class TransportTaskStatusClient {

View File

@@ -33,7 +33,6 @@ public class AcsFeedbackController {
@Operation(summary = "接收 ACS 状态反馈")
public CommonResult<Boolean> receiveAcsFeedback(@Valid @RequestBody AcsFeedbackReqDTO reqDTO) {
feedbackService.receiveAcsFeedback(reqDTO);
// 始终返回成功,避免 ACS 因业务异常而重试
return success(true);
}
}

View File

@@ -53,6 +53,14 @@ public class TransportTaskController {
return success(true);
}
@PostMapping("/operate")
@Operation(summary = "PC 端操作搬运任务(完成/取消/强制完成)")
@PreAuthorize("@ss.hasPermission('task:transport-task:operate')")
public CommonResult<Boolean> operateTransportTask(@Valid @RequestBody TransportTaskOperateReqVO reqVO) {
transportTaskService.operateTransportTask(reqVO);
return success(true);
}
@DeleteMapping("/delete")
@Operation(summary = "删除搬运任务")
@Parameter(name = "id", description = "编号", required = true)

View File

@@ -0,0 +1,25 @@
package cn.code.nl.module.task.controller.admin.transporttask.vo;
import io.swagger.v3.oas.annotations.media.Schema;
import jakarta.validation.constraints.NotEmpty;
import jakarta.validation.constraints.NotNull;
import lombok.Data;
/**
* PC 端搬运任务操作 Request VO
*
* @Author: liyongde
* @Date: 2026/7/15
*/
@Schema(description = "管理后台 - PC 端搬运任务操作 Request VO")
@Data
public class TransportTaskOperateReqVO {
@Schema(description = "任务ID", requiredMode = Schema.RequiredMode.REQUIRED)
@NotNull(message = "任务ID不能为空")
private Long taskId;
@Schema(description = "操作类型FINISHED/CANCELLED/FORCE-FINISH", requiredMode = Schema.RequiredMode.REQUIRED)
@NotEmpty(message = "操作类型不能为空")
private String operationType;
}

View File

@@ -34,8 +34,8 @@ public class TransportTaskPageReqVO extends PageParam {
@Schema(description = "任务类型", example = "2")
private String taskType;
@Schema(description = "任务状态", example = "2")
private String taskStatus;
@Schema(description = "任务状态(支持多选)", example = "[\"10\", \"50\"]")
private List<String> taskStatus;
@Schema(description = "ACS任务类型", example = "1")
private String acsTaskType;
@@ -112,7 +112,7 @@ public class TransportTaskPageReqVO extends PageParam {
@Schema(description = "ACS反馈参数")
private String resultParam;
@Schema(description = "备注", example = "你说的对")
@Schema(description = "备注")
private String remark;
@Schema(description = "创建时间")

View File

@@ -0,0 +1,23 @@
package cn.code.nl.module.task.convert.transporttask;
import cn.code.nl.module.task.dal.dataobject.transporttask.TransportTaskDO;
import cn.code.nl.module.task.dto.TaskInfoDTO;
import org.mapstruct.Mapper;
import org.mapstruct.factory.Mappers;
/**
* 搬运任务 Convert
*
* @Author: liyongde
* @Date: 2026/7/16
*/
@Mapper
public interface TransportTaskConvert {
TransportTaskConvert INSTANCE = Mappers.getMapper(TransportTaskConvert.class);
/**
* DO 转 RPC 全量信息 DTO字段同名零配置映射
*/
TaskInfoDTO convert(TransportTaskDO bean);
}

View File

@@ -28,7 +28,7 @@ public interface TransportTaskMapper extends BaseMapperX<TransportTaskDO> {
.eqIfPresent(TransportTaskDO::getBizId, reqVO.getBizId())
.eqIfPresent(TransportTaskDO::getHandleCode, reqVO.getHandleCode())
.eqIfPresent(TransportTaskDO::getTaskType, reqVO.getTaskType())
.eqIfPresent(TransportTaskDO::getTaskStatus, reqVO.getTaskStatus())
.inIfPresent(TransportTaskDO::getTaskStatus, reqVO.getTaskStatus())
.eqIfPresent(TransportTaskDO::getAcsTaskType, reqVO.getAcsTaskType())
.eqIfPresent(TransportTaskDO::getAgvSystemType, reqVO.getAgvSystemType())
.eqIfPresent(TransportTaskDO::getExternalTaskNo, reqVO.getExternalTaskNo())
@@ -59,6 +59,13 @@ public interface TransportTaskMapper extends BaseMapperX<TransportTaskDO> {
.orderByDesc(TransportTaskDO::getTaskId));
}
/**
* 根据任务编码查询任务taskCode 业务唯一)
*/
default TransportTaskDO selectByTaskCode(String taskCode) {
return selectOne(TransportTaskDO::getTaskCode, taskCode);
}
/**
* 按业务归属查询未完结任务(幂等校验用)
*/

View File

@@ -25,6 +25,9 @@ public class TaskEventMessage {
/** 业务归属服务LMS / WMS */
private String ownerService;
/** 任务类型 */
private String taskType;
/** 业务类型 */
private String bizType;

View File

@@ -1,18 +1,10 @@
package cn.code.nl.module.task.service.transporttask;
import cn.code.nl.framework.execute.biz.TaskCommonApi;
import cn.code.nl.framework.execute.core.TaskCommonApiFactory;
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.enums.CallbackStatusEnum;
import cn.code.nl.module.task.enums.TaskEventTypeEnum;
import cn.code.nl.module.task.enums.TaskOperationTypeEnum;
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.databind.ObjectMapper;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
@@ -20,12 +12,11 @@ import org.springframework.transaction.annotation.Transactional;
import org.springframework.validation.annotation.Validated;
import jakarta.annotation.Resource;
import java.util.Map;
import static cn.code.nl.module.task.framework.common.util.TaskUtil.isAllowedFrom;
/**
* ACS 反馈处理服务实现
* <p>
* 只做 ACS 前置处理(查任务/幂等/保存反馈参数),具体操作业务统一由 {@link TransportTaskService} 路由表分发
*
* @author 诺力管理员
*/
@@ -38,13 +29,7 @@ public class TransportTaskFeedbackServiceImpl implements TransportTaskFeedbackSe
private TransportTaskMapper transportTaskMapper;
@Resource
private TransportTaskStatusClient transportTaskStatusClient;
@Resource
private TaskEventProducer taskEventProducer;
@Resource
private TaskCommonApiFactory taskCommonApiFactory;
private TransportTaskService transportTaskService;
private static final ObjectMapper OBJECT_MAPPER = new ObjectMapper();
@@ -76,153 +61,13 @@ public class TransportTaskFeedbackServiceImpl implements TransportTaskFeedbackSe
}
}
// 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(), "40", "50")) {
log.warn("执行中反馈状态校验不通过, taskId={}, currentStatus={}",
task.getTaskId(), task.getTaskStatus());
// 4. 解析操作类型并校验允许 ACS 触发,走统一路由表分发
TaskOperationTypeEnum type = TaskOperationTypeEnum.getByCode(reqDTO.getStatus());
if (type == null || !type.isAcsAllowed()) {
log.warn("未知或不允许的 ACS 反馈状态, taskId={}, status={}",
reqDTO.getTaskId(), reqDTO.getStatus());
return;
}
task.setTaskStatus(TransportTaskStatusEnum.EXECUTING.getCode());
transportTaskMapper.updateById(task);
log.info("任务执行中, taskId={}", task.getTaskId());
}
/**
* 处理取货完成(61):更新状态 + 同步 HTTP 回调 LMS/WMS
* ownerService 服务提供商
*/
private void handlePicked(TransportTaskDO task, AcsFeedbackReqDTO reqDTO) {
if (!isAllowedFrom(task.getTaskStatus(), "60", "61")) {
log.warn("取货完成反馈状态校验不通过, taskId={}, currentStatus={}",
task.getTaskId(), task.getTaskStatus());
return;
}
// 更新状态
task.setTaskStatus(TransportTaskStatusEnum.PICKED.getCode());
transportTaskMapper.updateById(task);
log.info("任务已取货, taskId={}", task.getTaskId());
// 考虑可行性
TaskCommonApi byServerName = taskCommonApiFactory.getByServerName("lms-server");
byServerName.doHandlePicked();
// todo: 改成 feign 的 rpc 调用
// 同步 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(), "10", "40", "45", "50", "60", "61", "67")) {
log.warn("完成反馈状态校验不通过, taskId={}, currentStatus={}",
task.getTaskId(), task.getTaskStatus());
return;
}
// 已是 67 则只发 MQ不改状态
if (!TransportTaskStatusEnum.FINISHED_CALLBACK_PENDING.getCode().equals(task.getTaskStatus())) {
task.setTaskStatus(TransportTaskStatusEnum.FINISHED_CALLBACK_PENDING.getCode());
}
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(), "10", "40", "45", "50", "60", "61", "67", "69")) {
log.warn("取消反馈状态校验不通过, taskId={}, currentStatus={}",
task.getTaskId(), task.getTaskStatus());
return;
}
if (!TransportTaskStatusEnum.CANCEL_CALLBACK_PENDING.getCode().equals(task.getTaskStatus())) {
task.setTaskStatus(TransportTaskStatusEnum.CANCEL_CALLBACK_PENDING.getCode());
}
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(TransportTaskStatusEnum.PICKED.getCode());
req.setOwnerService(task.getOwnerService());
req.setBizType(task.getBizType());
req.setBizId(task.getBizId());
req.setHandleCode(task.getHandleCode());
req.setPayload(reqDTO.getPayload());
return req;
transportTaskService.dispatchOperation(task, type, reqDTO);
}
}

View File

@@ -4,7 +4,10 @@ import java.util.*;
import jakarta.validation.*;
import cn.code.nl.module.task.controller.admin.transporttask.vo.*;
import cn.code.nl.module.task.dal.dataobject.transporttask.TransportTaskDO;
import cn.code.nl.module.task.dto.AcsFeedbackReqDTO;
import cn.code.nl.module.task.dto.TaskInfoDTO;
import cn.code.nl.module.task.dto.TransportTaskCreateReqDTO;
import cn.code.nl.module.task.enums.TaskOperationTypeEnum;
import cn.code.nl.framework.common.pojo.PageResult;
import cn.code.nl.framework.common.pojo.PageParam;
@@ -68,4 +71,36 @@ public interface TransportTaskService {
*/
PageResult<TransportTaskDO> getTransportTaskPage(TransportTaskPageReqVO pageReqVO);
/**
* PC 端操作搬运任务(完成/取消/强制完成)
*
* @param reqVO 操作请求
*/
void operateTransportTask(@Valid TransportTaskOperateReqVO reqVO);
/**
* 按操作类型分发处理任务ACS 反馈与 PC 端操作共用路由表)
*
* @param task 任务
* @param type 操作类型
* @param reqDTO 反馈请求PC 端为构造的精简对象)
*/
void dispatchOperation(TransportTaskDO task, TaskOperationTypeEnum type, AcsFeedbackReqDTO reqDTO);
/**
* 根据 taskId 查询任务全量信息RPC 用,自动过滤逻辑删除,查不到返回 null
*
* @param taskId 任务ID
* @return 任务全量信息
*/
TaskInfoDTO getTaskInfoById(Long taskId);
/**
* 根据 taskCode 查询任务全量信息RPC 用,自动过滤逻辑删除,查不到返回 null
*
* @param taskCode 任务编码
* @return 任务全量信息
*/
TaskInfoDTO getTaskInfoByCode(String taskCode);
}

View File

@@ -2,27 +2,46 @@ package cn.code.nl.module.task.service.transporttask;
import cn.code.nl.framework.common.pojo.PageResult;
import cn.code.nl.framework.common.util.object.BeanUtils;
import cn.code.nl.framework.execute.biz.api.TaskCommonApi;
import cn.code.nl.framework.execute.biz.dto.TaskStatusCallApiReqDTO;
import cn.code.nl.framework.execute.core.TaskCommonApiFactory;
import cn.code.nl.module.task.controller.admin.transporttask.vo.TransportTaskOperateReqVO;
import cn.code.nl.module.task.controller.admin.transporttask.vo.TransportTaskPageReqVO;
import cn.code.nl.module.task.controller.admin.transporttask.vo.TransportTaskSaveReqVO;
import cn.code.nl.module.task.convert.transporttask.TransportTaskConvert;
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.TaskInfoDTO;
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.code.nl.module.task.enums.*;
import cn.code.nl.module.task.mq.message.TaskEventMessage;
import cn.code.nl.module.task.mq.producer.TaskEventProducer;
import cn.hutool.core.util.StrUtil;
import jakarta.annotation.PostConstruct;
import jakarta.annotation.Resource;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
import org.springframework.validation.annotation.Validated;
import java.util.EnumMap;
import java.util.List;
import java.util.Map;
import java.util.function.BiConsumer;
import static cn.code.nl.framework.common.exception.util.ServiceExceptionUtil.exception;
import static cn.code.nl.module.task.enums.ErrorCodeConstants.TRANSPORT_TASK_ALREADY_FINAL;
import static cn.code.nl.module.task.enums.ErrorCodeConstants.TRANSPORT_TASK_NOT_EXISTS;
import static cn.code.nl.module.task.enums.ErrorCodeConstants.TRANSPORT_TASK_OPERATION_NOT_SUPPORTED;
import static cn.code.nl.module.task.framework.common.util.TaskUtil.isAllowedFrom;
/**
* 搬运任务 Service 实现类
*
* @author 诺力管理员
*/
@Slf4j
@Service
@Validated
public class TransportTaskServiceImpl implements TransportTaskService {
@@ -30,6 +49,35 @@ public class TransportTaskServiceImpl implements TransportTaskService {
@Resource
private TransportTaskMapper transportTaskMapper;
@Resource
private TaskEventProducer taskEventProducer;
@Resource
private TaskCommonApiFactory taskCommonApiFactory;
/**
* 操作类型 -> 处理方法 路由表ACS 反馈与 PC 端操作共用)
*/
private final Map<TaskOperationTypeEnum, BiConsumer<TransportTaskDO, AcsFeedbackReqDTO>> operationHandlers =
new EnumMap<>(TaskOperationTypeEnum.class);
/**
* 初始化操作类型路由表
*/
@PostConstruct
public void initOperationHandlers() {
operationHandlers.put(TaskOperationTypeEnum.EXECUTING, this::handleExecuting);
operationHandlers.put(TaskOperationTypeEnum.PICKED, this::handlePicked);
operationHandlers.put(TaskOperationTypeEnum.FINISHED, this::handleFinished);
operationHandlers.put(TaskOperationTypeEnum.CANCELLED, this::handleCancelled);
// todo: 二次请求业务未开发
operationHandlers.put(TaskOperationTypeEnum.APPLY_AGAIN,
(task, reqDTO) -> log.info("二次请求业务未开发, taskId={}", task.getTaskId()));
// todo: 强制完成业务未开发
operationHandlers.put(TaskOperationTypeEnum.FORCE_FINISH,
(task, reqDTO) -> handleForceFinish(task));
}
@Override
public Long createTransportTask(TransportTaskSaveReqVO createReqVO) {
// 插入
@@ -96,4 +144,203 @@ public class TransportTaskServiceImpl implements TransportTaskService {
return transportTaskMapper.selectPage(pageReqVO);
}
/**
* PC 端操作搬运任务(完成/取消/强制完成):先校验再走统一路由表分发
*/
@Override
@Transactional(rollbackFor = Exception.class)
public void operateTransportTask(TransportTaskOperateReqVO reqVO) {
// 1. 校验任务存在
TransportTaskDO task = transportTaskMapper.selectById(reqVO.getTaskId());
if (task == null) {
throw exception(TRANSPORT_TASK_NOT_EXISTS);
}
// 2. 校验操作类型允许 PC 端触发
TaskOperationTypeEnum type = TaskOperationTypeEnum.getByCode(reqVO.getOperationType());
if (type == null || !type.isPcAllowed()) {
throw exception(TRANSPORT_TASK_OPERATION_NOT_SUPPORTED);
}
// 3. 终态校验:已完成/已取消不允许再操作
if (TransportTaskStatusEnum.FINISHED.getCode().equals(task.getTaskStatus())
|| TransportTaskStatusEnum.CANCELLED.getCode().equals(task.getTaskStatus())) {
throw exception(TRANSPORT_TASK_ALREADY_FINAL);
}
// 4. PC 端与 ACS 的差异点:记录完成类型
if (type == TaskOperationTypeEnum.FINISHED) {
task.setFinishedType(FinishedTypeEnum.MANUAL.getCode());
} else if (type == TaskOperationTypeEnum.FORCE_FINISH) {
task.setFinishedType(FinishedTypeEnum.MANUAL_FORCE.getCode());
}
// 5. 构造精简反馈对象,走统一路由表分发
AcsFeedbackReqDTO reqDTO = new AcsFeedbackReqDTO();
reqDTO.setTaskId(task.getTaskId());
reqDTO.setStatus(type.getCode());
dispatchOperation(task, type, reqDTO);
}
/**
* 按操作类型查路由表分发
*/
@Override
public void dispatchOperation(TransportTaskDO task, TaskOperationTypeEnum type, AcsFeedbackReqDTO reqDTO) {
BiConsumer<TransportTaskDO, AcsFeedbackReqDTO> handler = operationHandlers.get(type);
if (handler == null) {
log.warn("操作类型未注册处理器, taskId={}, type={}", task.getTaskId(), type.getCode());
return;
}
handler.accept(task, reqDTO);
}
/**
* 根据 taskId 查询任务全量信息:查不到返回 null由调用方判断
*/
@Override
public TaskInfoDTO getTaskInfoById(Long taskId) {
TransportTaskDO task = transportTaskMapper.selectById(taskId);
return TransportTaskConvert.INSTANCE.convert(task);
}
/**
* 根据 taskCode 查询任务全量信息:查不到返回 null由调用方判断
*/
@Override
public TaskInfoDTO getTaskInfoByCode(String taskCode) {
TransportTaskDO task = transportTaskMapper.selectByTaskCode(taskCode);
return TransportTaskConvert.INSTANCE.convert(task);
}
// ==================== 各状态处理方法(自 TransportTaskFeedbackServiceImpl 迁入,逻辑未改动) ====================
/**
* 处理执行中(60):直接更新状态,不做业务回调
*/
private void handleExecuting(TransportTaskDO task, AcsFeedbackReqDTO reqDTO) {
if (!isAllowedFrom(task.getTaskStatus(), TransportTaskStatusEnum.READY.getCode(),
TransportTaskStatusEnum.ISSUED.getCode())) {
log.warn("执行中反馈状态校验不通过, taskId={}, currentStatus={}",
task.getTaskId(), task.getTaskStatus());
return;
}
task.setTaskStatus(TransportTaskStatusEnum.EXECUTING.getCode());
task.setCarNo(reqDTO.getCarNo());
transportTaskMapper.updateById(task);
log.info("任务执行中, taskId={}", task.getTaskId());
}
/**
* 处理取货完成(65):更新状态 + 同步回调 LMS/WMS
* ownerService 服务提供商
*/
private void handlePicked(TransportTaskDO task, AcsFeedbackReqDTO reqDTO) {
if (!isAllowedFrom(task.getTaskStatus(), TransportTaskStatusEnum.EXECUTING.getCode(),
TransportTaskStatusEnum.PICKED.getCode())) {
log.warn("取货完成反馈状态校验不通过, taskId={}, currentStatus={}",
task.getTaskId(), task.getTaskStatus());
return;
}
// 更新状态
task.setTaskStatus(TransportTaskStatusEnum.PICKED.getCode());
transportTaskMapper.updateById(task);
log.info("任务已取货, taskId={}", task.getTaskId());
TaskStatusCallApiReqDTO req = new TaskStatusCallApiReqDTO();
req.setTaskId(task.getTaskId());
req.setTaskCode(task.getTaskCode());
req.setStatus(TransportTaskStatusEnum.PICKED.getCode());
req.setOwnerService(task.getOwnerService());
req.setBizType(task.getBizType());
req.setBizId(task.getBizId());
req.setHandleCode(task.getHandleCode());
req.setPayload(reqDTO.getPayload());
// 同步请求业务
TaskCommonApi serverApi = taskCommonApiFactory.getByServerName(task.getOwnerService());
serverApi.doHandlePicked(req);
}
/**
* 处理完成(75):更新状态 + 异步 MQ 通知 LMS/WMS
*/
private void handleFinished(TransportTaskDO task, AcsFeedbackReqDTO reqDTO) {
if (!isAllowedFrom(task.getTaskStatus(), TransportTaskStatusEnum.CREATED.getCode(),
TransportTaskStatusEnum.READY.getCode(), TransportTaskStatusEnum.ISSUING.getCode(),
TransportTaskStatusEnum.ISSUED.getCode(), TransportTaskStatusEnum.EXECUTING.getCode(),
TransportTaskStatusEnum.PICKED.getCode(),
TransportTaskStatusEnum.FINISHED_CALLBACK_PENDING.getCode())) {
log.warn("完成反馈状态校验不通过, taskId={}, currentStatus={}",
task.getTaskId(), task.getTaskStatus());
return;
}
// 已是 75 则只发 MQ不改状态
if (!TransportTaskStatusEnum.FINISHED_CALLBACK_PENDING.getCode().equals(task.getTaskStatus())) {
task.setTaskStatus(TransportTaskStatusEnum.FINISHED_CALLBACK_PENDING.getCode());
}
task.setCallbackStatus(CallbackStatusEnum.PENDING.getCode());
transportTaskMapper.updateById(task);
publishEvent(task, TaskEventTypeEnum.TASK_FINISHED, reqDTO.getPayload());
}
/**
* 处理取消(85):更新状态 + 异步 MQ 通知 LMS/WMS
*/
private void handleCancelled(TransportTaskDO task, AcsFeedbackReqDTO reqDTO) {
if (!isAllowedFrom(task.getTaskStatus(), TransportTaskStatusEnum.CREATED.getCode(),
TransportTaskStatusEnum.READY.getCode(), TransportTaskStatusEnum.ISSUING.getCode(),
TransportTaskStatusEnum.ISSUED.getCode(), TransportTaskStatusEnum.EXECUTING.getCode(),
TransportTaskStatusEnum.PICKED.getCode(),
TransportTaskStatusEnum.FINISHED_CALLBACK_PENDING.getCode(),
TransportTaskStatusEnum.CANCEL_CALLBACK_PENDING.getCode())) {
log.warn("取消反馈状态校验不通过, taskId={}, currentStatus={}",
task.getTaskId(), task.getTaskStatus());
return;
}
if (!TransportTaskStatusEnum.CANCEL_CALLBACK_PENDING.getCode().equals(task.getTaskStatus())) {
task.setTaskStatus(TransportTaskStatusEnum.CANCEL_CALLBACK_PENDING.getCode());
}
task.setCallbackStatus(CallbackStatusEnum.PENDING.getCode());
transportTaskMapper.updateById(task);
publishEvent(task, TaskEventTypeEnum.TASK_CANCELLED, reqDTO.getPayload());
}
/**
* 处理强制完成(79):直接更新状态,不做业务回调
*/
private void handleForceFinish(TransportTaskDO task) {
task.setTaskStatus(TransportTaskStatusEnum.FINISHED.getCode());
task.setFinishedType(FinishedTypeEnum.MANUAL_FORCE.getCode());
transportTaskMapper.updateById(task);
log.info("任务强制完成, taskId={}", task.getTaskId());
}
// ==================== 工具方法 ====================
/**
* 发布 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.setTaskType(task.getTaskType());
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);
}
}
}