Refactor document workflow and related UI state
This commit is contained in:
@@ -8,7 +8,7 @@ import lombok.Data;
|
||||
import java.util.Map;
|
||||
|
||||
/**
|
||||
* ACS 状态反馈请求 DTO
|
||||
* ACS 请求 DTO
|
||||
*/
|
||||
@Schema(description = "ACS 状态反馈 Request DTO")
|
||||
@Data
|
||||
@@ -21,7 +21,10 @@ public class AcsFeedbackReqDTO {
|
||||
@Schema(description = "任务编码")
|
||||
private String taskCode;
|
||||
|
||||
@Schema(description = "ACS反馈状态:EXECUTING/PICKED/FINISHED/CANCELLED/APPLYAGAIN", requiredMode = Schema.RequiredMode.REQUIRED)
|
||||
@Schema(description = "设备编码")
|
||||
private String deviceCode;
|
||||
|
||||
@Schema(description = "ACS反馈状态或业务操作:状态反馈为 EXECUTING/FINISHED/CANCELLED;业务操作为 PICKED/APPLY-AGAIN/REQUEST-RELEASE/REQUEST-PICK/REQUEST-LEAVE/REQUEST-ENTER", requiredMode = Schema.RequiredMode.REQUIRED)
|
||||
@NotEmpty(message = "反馈状态不能为空")
|
||||
private String status;
|
||||
|
||||
|
||||
@@ -0,0 +1,47 @@
|
||||
package cn.code.nl.module.task.enums;
|
||||
|
||||
import lombok.Getter;
|
||||
|
||||
/**
|
||||
* ACS 业务操作类型枚举
|
||||
*
|
||||
* @Author: liyongde
|
||||
* @Date: 2026/7/16
|
||||
*/
|
||||
@Getter
|
||||
public enum AcsBusinessOperationTypeEnum {
|
||||
|
||||
PICKED("PICKED", "取货完成"),
|
||||
|
||||
APPLY_AGAIN("APPLY-AGAIN", "二次请求"),
|
||||
|
||||
REQUEST_RELEASE("REQUEST-RELEASE", "请求放货"),
|
||||
|
||||
REQUEST_PICK("REQUEST-PICK", "请求取货"),
|
||||
|
||||
REQUEST_LEAVE("REQUEST-LEAVE", "请求离开"),
|
||||
|
||||
REQUEST_ENTER("REQUEST-ENTER", "请求进入");
|
||||
|
||||
/** 操作编码 */
|
||||
private final String code;
|
||||
/** 操作名称 */
|
||||
private final String name;
|
||||
|
||||
AcsBusinessOperationTypeEnum(String code, String name) {
|
||||
this.code = code;
|
||||
this.name = name;
|
||||
}
|
||||
|
||||
/**
|
||||
* 根据编码解析枚举,找不到返回 null
|
||||
*/
|
||||
public static AcsBusinessOperationTypeEnum getByCode(String code) {
|
||||
for (AcsBusinessOperationTypeEnum type : values()) {
|
||||
if (type.getCode().equalsIgnoreCase(code)) {
|
||||
return type;
|
||||
}
|
||||
}
|
||||
return null;
|
||||
}
|
||||
}
|
||||
@@ -8,6 +8,7 @@ import cn.code.nl.framework.common.exception.ErrorCode;
|
||||
* @Date: 2026/7/14 13:03
|
||||
*/
|
||||
public interface ErrorCodeConstants {
|
||||
ErrorCode TRANSPORT_TASK_NOT_EXISTS_ID = new ErrorCode(5000, "任务id={}不存在,请检查");
|
||||
ErrorCode TRANSPORT_TASK_NOT_EXISTS = new ErrorCode(5001, "搬运任务不存在");
|
||||
ErrorCode TRANSPORT_TASK_EXISTS_UNFINISHED = new ErrorCode(5002, "同业务类型下存在未完结的搬运任务");
|
||||
ErrorCode TRANSPORT_TASK_STATUS_NOT_ALLOW = new ErrorCode(5003, "当前任务状态不允许此操作");
|
||||
|
||||
@@ -15,14 +15,10 @@ public enum TaskOperationTypeEnum {
|
||||
|
||||
EXECUTING("EXECUTING", "执行中", true, false),
|
||||
|
||||
PICKED("PICKED", "取货完成", true, false),
|
||||
|
||||
FINISHED("FINISHED", "完成任务", true, true),
|
||||
|
||||
CANCELLED("CANCELLED", "取消任务", true, true),
|
||||
|
||||
APPLY_AGAIN("APPLY-AGAIN", "二次请求", true, false),
|
||||
|
||||
FORCE_FINISH("FORCE-FINISH", "强制完成", false, true);
|
||||
|
||||
/** 操作编码 */
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
package cn.code.nl.module.task.controller.admin.transporttask;
|
||||
|
||||
import cn.code.nl.framework.common.pojo.CommonResult;
|
||||
import cn.code.nl.framework.execute.biz.vo.AcsApplyActionRespVO;
|
||||
import cn.code.nl.module.task.dto.AcsFeedbackReqDTO;
|
||||
import cn.code.nl.module.task.service.transporttask.TransportTaskFeedbackService;
|
||||
import io.swagger.v3.oas.annotations.Operation;
|
||||
@@ -18,7 +19,7 @@ import static cn.code.nl.framework.common.pojo.CommonResult.success;
|
||||
/**
|
||||
* ACS 状态反馈接口
|
||||
* <p>
|
||||
* 接收 ACS 的任务状态回调,处理执行中/取货完成/完成/取消
|
||||
* 接收 ACS 的任务状态与业务操作回调
|
||||
*/
|
||||
@Tag(name = "ACS 反馈")
|
||||
@RestController
|
||||
@@ -35,4 +36,10 @@ public class AcsFeedbackController {
|
||||
feedbackService.receiveAcsFeedback(reqDTO);
|
||||
return success(true);
|
||||
}
|
||||
|
||||
@PostMapping("/acs-business-feedback")
|
||||
@Operation(summary = "接收 ACS 业务操作反馈")
|
||||
public CommonResult<AcsApplyActionRespVO> receiveAcsBusinessFeedback(@Valid @RequestBody AcsFeedbackReqDTO reqDTO) {
|
||||
return success(feedbackService.receiveAcsBusinessFeedback(reqDTO));
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,196 @@
|
||||
package cn.code.nl.module.task.manage;
|
||||
|
||||
import cn.code.nl.framework.common.exception.ServiceException;
|
||||
import cn.code.nl.framework.common.pojo.CommonResult;
|
||||
import cn.code.nl.framework.execute.biz.api.TaskCommonApi;
|
||||
import cn.code.nl.framework.execute.biz.dto.TaskStatusCallApiReqDTO;
|
||||
import cn.code.nl.framework.execute.biz.vo.AcsApplyActionRespVO;
|
||||
import cn.code.nl.framework.execute.core.TaskCommonApiFactory;
|
||||
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.enums.AcsBusinessOperationTypeEnum;
|
||||
import cn.code.nl.module.task.enums.TransportTaskStatusEnum;
|
||||
import com.alibaba.fastjson.JSON;
|
||||
import jakarta.annotation.PostConstruct;
|
||||
import jakarta.annotation.Resource;
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
import org.redisson.api.RLock;
|
||||
import org.redisson.api.RedissonClient;
|
||||
import org.springframework.stereotype.Component;
|
||||
|
||||
import java.util.EnumMap;
|
||||
import java.util.Map;
|
||||
import java.util.function.BiFunction;
|
||||
|
||||
import static cn.code.nl.module.task.framework.common.util.TaskUtil.isAllowedFrom;
|
||||
|
||||
/**
|
||||
* ACS 业务操作分发管理器
|
||||
*/
|
||||
@Slf4j
|
||||
@Component
|
||||
public class TransportTaskBusinessOperationManager {
|
||||
|
||||
@Resource
|
||||
private TransportTaskMapper transportTaskMapper;
|
||||
|
||||
@Resource
|
||||
private TaskCommonApiFactory taskCommonApiFactory;
|
||||
|
||||
@Resource
|
||||
private RedissonClient redissonClient;
|
||||
|
||||
/** ACS 业务操作路由表 */
|
||||
private final Map<AcsBusinessOperationTypeEnum, BiFunction<TransportTaskDO, AcsFeedbackReqDTO, Object>> operationHandlers =
|
||||
new EnumMap<>(AcsBusinessOperationTypeEnum.class);
|
||||
|
||||
/**
|
||||
* 初始化 ACS 业务操作路由表
|
||||
*/
|
||||
@PostConstruct
|
||||
public void initOperationHandlers() {
|
||||
operationHandlers.put(AcsBusinessOperationTypeEnum.PICKED, this::handlePicked);
|
||||
operationHandlers.put(AcsBusinessOperationTypeEnum.APPLY_AGAIN, this::handleApplyAgain);
|
||||
operationHandlers.put(AcsBusinessOperationTypeEnum.REQUEST_RELEASE, this::handleRequestRelease);
|
||||
operationHandlers.put(AcsBusinessOperationTypeEnum.REQUEST_PICK, this::handleRequestPick);
|
||||
operationHandlers.put(AcsBusinessOperationTypeEnum.REQUEST_LEAVE, this::handleRequestLeave);
|
||||
operationHandlers.put(AcsBusinessOperationTypeEnum.REQUEST_ENTER, this::handleRequestEnter);
|
||||
}
|
||||
|
||||
/**
|
||||
* 按 ACS 业务操作类型分发
|
||||
*/
|
||||
public AcsApplyActionRespVO dispatchOperation(TransportTaskDO task, AcsBusinessOperationTypeEnum type,
|
||||
AcsFeedbackReqDTO reqDTO) {
|
||||
RLock lock = redissonClient.getLock(String.valueOf(task.getTaskId()));
|
||||
if (lock.tryLock()) {
|
||||
try {
|
||||
BiFunction<TransportTaskDO, AcsFeedbackReqDTO, Object> handler = operationHandlers.get(type);
|
||||
if (handler == null) {
|
||||
log.warn("ACS 业务操作未注册处理器, taskId={}, type={}", task.getTaskId(), type.getCode());
|
||||
return buildResp(task, null);
|
||||
}
|
||||
Object data = handler.apply(task, reqDTO);
|
||||
return buildResp(task, data);
|
||||
} catch (Exception ex) {
|
||||
log.error("[acsBusinessOperation][执行异常][lockKey={}]", task.getTaskId(), ex);
|
||||
return buildResp(task, null);
|
||||
} finally {
|
||||
if (lock.isHeldByCurrentThread()) {
|
||||
lock.unlock();
|
||||
}
|
||||
}
|
||||
} else {
|
||||
throw new ServiceException(5007, "任务标识为:" + task.getTaskId() + "的任务正在操作中!");
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 处理取货完成
|
||||
*/
|
||||
private AcsApplyActionRespVO handlePicked(TransportTaskDO task, AcsFeedbackReqDTO reqDTO) {
|
||||
if (!isAllowedFrom(task.getTaskStatus(), TransportTaskStatusEnum.EXECUTING.getCode(),
|
||||
TransportTaskStatusEnum.PICKED.getCode())) {
|
||||
log.warn("取货完成反馈状态校验不通过, taskId={}, currentStatus={}",
|
||||
task.getTaskId(), task.getTaskStatus());
|
||||
return null;
|
||||
}
|
||||
task.setTaskStatus(TransportTaskStatusEnum.PICKED.getCode());
|
||||
transportTaskMapper.updateById(task);
|
||||
log.info("任务已取货, taskId={}", task.getTaskId());
|
||||
|
||||
// 调用具体的服务去执行取货完成操作。
|
||||
TaskCommonApi serverApi = taskCommonApiFactory.getByServerName(task.getOwnerService());
|
||||
CommonResult<AcsApplyActionRespVO> result =
|
||||
serverApi.doHandlePicked(buildReq(task, reqDTO));
|
||||
return result.getData();
|
||||
}
|
||||
|
||||
/**
|
||||
* 处理二次请求
|
||||
*/
|
||||
private AcsApplyActionRespVO handleApplyAgain(TransportTaskDO task, AcsFeedbackReqDTO reqDTO) {
|
||||
// TODO 二次请求业务待开发
|
||||
log.info("二次请求业务待开发, taskId={}", task.getTaskId());
|
||||
// 调用具体的服务去执行取货完成操作。
|
||||
TaskCommonApi serverApi = taskCommonApiFactory.getByServerName(task.getOwnerService());
|
||||
CommonResult<AcsApplyActionRespVO> result =
|
||||
serverApi.doHandleApplyAgain(buildReq(task, reqDTO));
|
||||
return result.getData();
|
||||
}
|
||||
|
||||
/**
|
||||
* 处理请求放货
|
||||
*/
|
||||
private Object handleRequestRelease(TransportTaskDO task, AcsFeedbackReqDTO reqDTO) {
|
||||
// TODO 请求放货业务待开发
|
||||
log.info("请求放货业务待开发, taskId={}", task.getTaskId());
|
||||
TaskCommonApi serverApi = taskCommonApiFactory.getByServerName(task.getOwnerService());
|
||||
CommonResult<AcsApplyActionRespVO> result =
|
||||
serverApi.doHandleRequestRelease(buildReq(task, reqDTO));
|
||||
return result.getData();
|
||||
}
|
||||
|
||||
/**
|
||||
* 处理请求取货
|
||||
*/
|
||||
private Object handleRequestPick(TransportTaskDO task, AcsFeedbackReqDTO reqDTO) {
|
||||
// TODO 请求取货业务待开发
|
||||
log.info("请求取货业务待开发, taskId={}", task.getTaskId());
|
||||
TaskCommonApi serverApi = taskCommonApiFactory.getByServerName(task.getOwnerService());
|
||||
CommonResult<AcsApplyActionRespVO> result =
|
||||
serverApi.doHandleRequestPick(buildReq(task, reqDTO));
|
||||
return result.getData();
|
||||
}
|
||||
|
||||
/**
|
||||
* 处理请求离开
|
||||
*/
|
||||
private Object handleRequestLeave(TransportTaskDO task, AcsFeedbackReqDTO reqDTO) {
|
||||
// TODO 请求离开业务待开发
|
||||
log.info("请求离开业务待开发, taskId={}", task.getTaskId());
|
||||
TaskCommonApi serverApi = taskCommonApiFactory.getByServerName(task.getOwnerService());
|
||||
CommonResult<AcsApplyActionRespVO> result =
|
||||
serverApi.doHandleRequestLeave(buildReq(task, reqDTO));
|
||||
return result.getData();
|
||||
}
|
||||
|
||||
/**
|
||||
* 处理请求进入
|
||||
*/
|
||||
private Object handleRequestEnter(TransportTaskDO task, AcsFeedbackReqDTO reqDTO) {
|
||||
// TODO 请求进入业务待开发
|
||||
log.info("请求进入业务待开发, taskId={}", task.getTaskId());
|
||||
TaskCommonApi serverApi = taskCommonApiFactory.getByServerName(task.getOwnerService());
|
||||
CommonResult<AcsApplyActionRespVO> result =
|
||||
serverApi.doHandleRequestEnter(buildReq(task, reqDTO));
|
||||
return result.getData();
|
||||
}
|
||||
|
||||
/**
|
||||
* 构建 ACS 业务操作返回对象
|
||||
*/
|
||||
private AcsApplyActionRespVO buildResp(TransportTaskDO task, Object data) {
|
||||
AcsApplyActionRespVO respVO = new AcsApplyActionRespVO();
|
||||
respVO.setTaskId(task.getTaskId());
|
||||
respVO.setTaskCode(task.getTaskCode());
|
||||
respVO.setData(JSON.parseObject(JSON.toJSONString(data)));
|
||||
return respVO;
|
||||
}
|
||||
/**
|
||||
* 构建 ACS 业务操作返回对象
|
||||
*/
|
||||
private TaskStatusCallApiReqDTO buildReq(TransportTaskDO task, AcsFeedbackReqDTO reqDTO) {
|
||||
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());
|
||||
return req;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,181 @@
|
||||
package cn.code.nl.module.task.manage;
|
||||
|
||||
import cn.code.nl.framework.common.exception.ServiceException;
|
||||
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.enums.CallbackStatusEnum;
|
||||
import cn.code.nl.module.task.enums.FinishedTypeEnum;
|
||||
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.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.redisson.api.RLock;
|
||||
import org.redisson.api.RedissonClient;
|
||||
import org.springframework.stereotype.Component;
|
||||
|
||||
import java.util.EnumMap;
|
||||
import java.util.Map;
|
||||
import java.util.function.BiConsumer;
|
||||
|
||||
import static cn.code.nl.module.task.framework.common.util.TaskUtil.isAllowedFrom;
|
||||
|
||||
/**
|
||||
* 搬运任务状态操作分发管理器
|
||||
*/
|
||||
@Slf4j
|
||||
@Component
|
||||
public class TransportTaskOperationManager {
|
||||
|
||||
@Resource
|
||||
private TransportTaskMapper transportTaskMapper;
|
||||
|
||||
@Resource
|
||||
private TaskEventProducer taskEventProducer;
|
||||
|
||||
@Resource
|
||||
private RedissonClient redissonClient;
|
||||
|
||||
/** 操作类型路由表 */
|
||||
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.FINISHED, this::handleFinished);
|
||||
operationHandlers.put(TaskOperationTypeEnum.CANCELLED, this::handleCancelled);
|
||||
operationHandlers.put(TaskOperationTypeEnum.FORCE_FINISH, (task, reqDTO) -> handleForceFinish(task));
|
||||
}
|
||||
|
||||
/**
|
||||
* 按操作类型查路由表分发
|
||||
*/
|
||||
public void dispatchOperation(TransportTaskDO task, TaskOperationTypeEnum type, AcsFeedbackReqDTO reqDTO) {
|
||||
RLock lock = redissonClient.getLock(String.valueOf(task.getTaskId()));
|
||||
if (lock.tryLock()) {
|
||||
try {
|
||||
BiConsumer<TransportTaskDO, AcsFeedbackReqDTO> handler = operationHandlers.get(type);
|
||||
if (handler == null) {
|
||||
log.warn("操作类型未注册处理器, taskId={}, type={}", task.getTaskId(), type.getCode());
|
||||
return;
|
||||
}
|
||||
handler.accept(task, reqDTO);
|
||||
} catch (Exception ex) {
|
||||
log.error("[messageResend][执行异常][lockKey={}]", task.getTaskId(), ex);
|
||||
} finally {
|
||||
if (lock.isHeldByCurrentThread()) {
|
||||
lock.unlock();
|
||||
}
|
||||
}
|
||||
} else {
|
||||
throw new ServiceException(5007, "任务标识为:" + task.getTaskId() + "的任务正在操作中!");
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 处理执行中
|
||||
*/
|
||||
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());
|
||||
}
|
||||
|
||||
/**
|
||||
* 处理完成
|
||||
*/
|
||||
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;
|
||||
}
|
||||
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());
|
||||
}
|
||||
|
||||
/**
|
||||
* 处理取消
|
||||
*/
|
||||
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());
|
||||
}
|
||||
|
||||
/**
|
||||
* 处理强制完成
|
||||
*/
|
||||
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 事件
|
||||
*/
|
||||
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);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1,5 +1,6 @@
|
||||
package cn.code.nl.module.task.service.transporttask;
|
||||
|
||||
import cn.code.nl.framework.execute.biz.vo.AcsApplyActionRespVO;
|
||||
import cn.code.nl.module.task.dto.AcsFeedbackReqDTO;
|
||||
import jakarta.validation.Valid;
|
||||
|
||||
@@ -16,4 +17,12 @@ public interface TransportTaskFeedbackService {
|
||||
* @param reqDTO 反馈请求
|
||||
*/
|
||||
void receiveAcsFeedback(@Valid AcsFeedbackReqDTO reqDTO);
|
||||
|
||||
/**
|
||||
* 接收 ACS 业务操作反馈并处理
|
||||
*
|
||||
* @param reqDTO 反馈请求
|
||||
*/
|
||||
AcsApplyActionRespVO receiveAcsBusinessFeedback(AcsFeedbackReqDTO reqDTO);
|
||||
|
||||
}
|
||||
|
||||
@@ -1,10 +1,15 @@
|
||||
package cn.code.nl.module.task.service.transporttask;
|
||||
|
||||
import cn.code.nl.framework.common.exception.ServiceException;
|
||||
import cn.code.nl.framework.execute.biz.vo.AcsApplyActionRespVO;
|
||||
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.enums.AcsBusinessOperationTypeEnum;
|
||||
import cn.code.nl.module.task.enums.TaskOperationTypeEnum;
|
||||
import cn.code.nl.module.task.enums.TransportTaskStatusEnum;
|
||||
import cn.code.nl.module.task.manage.TransportTaskBusinessOperationManager;
|
||||
import cn.code.nl.module.task.manage.TransportTaskOperationManager;
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
import org.springframework.stereotype.Service;
|
||||
@@ -13,6 +18,9 @@ import org.springframework.validation.annotation.Validated;
|
||||
|
||||
import jakarta.annotation.Resource;
|
||||
|
||||
import static cn.code.nl.framework.common.exception.util.ServiceExceptionUtil.exception;
|
||||
import static cn.code.nl.module.task.enums.ErrorCodeConstants.TRANSPORT_TASK_NOT_EXISTS_ID;
|
||||
|
||||
/**
|
||||
* ACS 反馈处理服务实现
|
||||
* <p>
|
||||
@@ -29,45 +37,96 @@ public class TransportTaskFeedbackServiceImpl implements TransportTaskFeedbackSe
|
||||
private TransportTaskMapper transportTaskMapper;
|
||||
|
||||
@Resource
|
||||
private TransportTaskService transportTaskService;
|
||||
private TransportTaskOperationManager transportTaskOperationManager;
|
||||
|
||||
@Resource
|
||||
private TransportTaskBusinessOperationManager transportTaskBusinessOperationManager;
|
||||
|
||||
private static final ObjectMapper OBJECT_MAPPER = new ObjectMapper();
|
||||
|
||||
@Override
|
||||
@Transactional(rollbackFor = Exception.class)
|
||||
public void receiveAcsFeedback(AcsFeedbackReqDTO reqDTO) {
|
||||
// 1. 查询任务
|
||||
TransportTaskDO task = transportTaskMapper.selectById(reqDTO.getTaskId());
|
||||
if (task == null) {
|
||||
log.warn("ACS 反馈 taskId 不存在, taskId={}, 直接返回成功", reqDTO.getTaskId());
|
||||
return;
|
||||
}
|
||||
TransportTaskDO task = getTaskForFeedback(reqDTO);
|
||||
saveResultParam(task, reqDTO);
|
||||
|
||||
// 2. 幂等判断:已经是终态(完成/取消)则不重复处理
|
||||
String currentStatus = task.getTaskStatus();
|
||||
if (TransportTaskStatusEnum.FINISHED.getCode().equals(currentStatus)
|
||||
|| TransportTaskStatusEnum.CANCELLED.getCode().equals(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. 解析操作类型并校验允许 ACS 触发,走统一路由表分发
|
||||
// 解析状态操作类型并校验允许 ACS 触发
|
||||
TaskOperationTypeEnum type = TaskOperationTypeEnum.getByCode(reqDTO.getStatus());
|
||||
if (type == null || !type.isAcsAllowed()) {
|
||||
log.warn("未知或不允许的 ACS 反馈状态, taskId={}, status={}",
|
||||
reqDTO.getTaskId(), reqDTO.getStatus());
|
||||
return;
|
||||
}
|
||||
transportTaskService.dispatchOperation(task, type, reqDTO);
|
||||
transportTaskOperationManager.dispatchOperation(task, type, reqDTO);
|
||||
}
|
||||
|
||||
@Override
|
||||
@Transactional(rollbackFor = Exception.class)
|
||||
public AcsApplyActionRespVO receiveAcsBusinessFeedback(AcsFeedbackReqDTO reqDTO) {
|
||||
TransportTaskDO task = getTaskForFeedback(reqDTO);
|
||||
if (isFinalTask(task, reqDTO)) {
|
||||
return buildResp(task);
|
||||
}
|
||||
|
||||
saveResultParam(task, reqDTO);
|
||||
|
||||
// 解析 ACS 业务操作类型,走业务操作路由表分发
|
||||
AcsBusinessOperationTypeEnum type = AcsBusinessOperationTypeEnum.getByCode(reqDTO.getStatus());
|
||||
if (type == null) {
|
||||
log.warn("未知的 ACS 业务操作类型, taskId={}, status={}",
|
||||
reqDTO.getTaskId(), reqDTO.getStatus());
|
||||
return buildResp(task);
|
||||
}
|
||||
return transportTaskBusinessOperationManager.dispatchOperation(task, type, reqDTO);
|
||||
}
|
||||
|
||||
/**
|
||||
* 查询反馈对应的任务
|
||||
*/
|
||||
private TransportTaskDO getTaskForFeedback(AcsFeedbackReqDTO reqDTO) {
|
||||
TransportTaskDO task = transportTaskMapper.selectById(reqDTO.getTaskId());
|
||||
if (task == null) {
|
||||
throw exception(TRANSPORT_TASK_NOT_EXISTS_ID, reqDTO.getTaskId());
|
||||
}
|
||||
return task;
|
||||
}
|
||||
|
||||
/**
|
||||
* 判断任务是否已经处于终态
|
||||
*/
|
||||
private boolean isFinalTask(TransportTaskDO task, AcsFeedbackReqDTO reqDTO) {
|
||||
String currentStatus = task.getTaskStatus();
|
||||
if (TransportTaskStatusEnum.FINISHED.getCode().equals(currentStatus)
|
||||
|| TransportTaskStatusEnum.CANCELLED.getCode().equals(currentStatus)) {
|
||||
log.info("任务已处终态, taskId={}, currentStatus={}, 忽略反馈",
|
||||
reqDTO.getTaskId(), currentStatus);
|
||||
return true;
|
||||
}
|
||||
return false;
|
||||
}
|
||||
|
||||
/**
|
||||
* 保存 ACS 反馈扩展数据
|
||||
*/
|
||||
private void saveResultParam(TransportTaskDO task, AcsFeedbackReqDTO reqDTO) {
|
||||
if (reqDTO.getPayload() == null) {
|
||||
return;
|
||||
}
|
||||
try {
|
||||
task.setResultParam(OBJECT_MAPPER.writeValueAsString(reqDTO.getPayload()));
|
||||
transportTaskMapper.updateById(task);
|
||||
} catch (Exception e) {
|
||||
log.warn("resultParam 序列化失败, taskId={}", reqDTO.getTaskId(), e);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 根据任务构建空数据返回
|
||||
*/
|
||||
private AcsApplyActionRespVO buildResp(TransportTaskDO task) {
|
||||
AcsApplyActionRespVO respVO = new AcsApplyActionRespVO();
|
||||
respVO.setTaskId(task.getTaskId());
|
||||
respVO.setTaskCode(task.getTaskCode());
|
||||
return respVO;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,106 +1,95 @@
|
||||
package cn.code.nl.module.task.service.transporttask;
|
||||
|
||||
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;
|
||||
|
||||
/**
|
||||
* 搬运任务 Service 接口
|
||||
*
|
||||
* @author 诺力管理员
|
||||
*/
|
||||
public interface TransportTaskService {
|
||||
|
||||
/**
|
||||
* 创建搬运任务
|
||||
*
|
||||
* @param createReqVO 创建信息
|
||||
* @return 编号
|
||||
*/
|
||||
Long createTransportTask(@Valid TransportTaskSaveReqVO createReqVO);
|
||||
|
||||
/**
|
||||
* 通过 RPC 创建搬运任务(LMS/WMS 调用)
|
||||
*
|
||||
* @param reqDTO 创建请求
|
||||
* @return 任务ID
|
||||
*/
|
||||
Long createTransportTaskByRpc(@Valid TransportTaskCreateReqDTO reqDTO);
|
||||
|
||||
/**
|
||||
* 更新搬运任务
|
||||
*
|
||||
* @param updateReqVO 更新信息
|
||||
*/
|
||||
void updateTransportTask(@Valid TransportTaskSaveReqVO updateReqVO);
|
||||
|
||||
/**
|
||||
* 删除搬运任务
|
||||
*
|
||||
* @param id 编号
|
||||
*/
|
||||
void deleteTransportTask(Long id);
|
||||
|
||||
/**
|
||||
* 批量删除搬运任务
|
||||
*
|
||||
* @param ids 编号
|
||||
*/
|
||||
void deleteTransportTaskListByIds(List<Long> ids);
|
||||
|
||||
/**
|
||||
* 获得搬运任务
|
||||
*
|
||||
* @param id 编号
|
||||
* @return 搬运任务
|
||||
*/
|
||||
TransportTaskDO getTransportTask(Long id);
|
||||
|
||||
/**
|
||||
* 获得搬运任务分页
|
||||
*
|
||||
* @param pageReqVO 分页查询
|
||||
* @return 搬运任务分页
|
||||
*/
|
||||
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);
|
||||
|
||||
}
|
||||
package cn.code.nl.module.task.service.transporttask;
|
||||
|
||||
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.TaskInfoDTO;
|
||||
import cn.code.nl.module.task.dto.TransportTaskCreateReqDTO;
|
||||
import cn.code.nl.framework.common.pojo.PageResult;
|
||||
import cn.code.nl.framework.common.pojo.PageParam;
|
||||
|
||||
/**
|
||||
* 搬运任务 Service 接口
|
||||
*
|
||||
* @author 诺力管理员
|
||||
*/
|
||||
public interface TransportTaskService {
|
||||
|
||||
/**
|
||||
* 创建搬运任务
|
||||
*
|
||||
* @param createReqVO 创建信息
|
||||
* @return 编号
|
||||
*/
|
||||
Long createTransportTask(@Valid TransportTaskSaveReqVO createReqVO);
|
||||
|
||||
/**
|
||||
* 通过 RPC 创建搬运任务(LMS/WMS 调用)
|
||||
*
|
||||
* @param reqDTO 创建请求
|
||||
* @return 任务ID
|
||||
*/
|
||||
Long createTransportTaskByRpc(@Valid TransportTaskCreateReqDTO reqDTO);
|
||||
|
||||
/**
|
||||
* 更新搬运任务
|
||||
*
|
||||
* @param updateReqVO 更新信息
|
||||
*/
|
||||
void updateTransportTask(@Valid TransportTaskSaveReqVO updateReqVO);
|
||||
|
||||
/**
|
||||
* 删除搬运任务
|
||||
*
|
||||
* @param id 编号
|
||||
*/
|
||||
void deleteTransportTask(Long id);
|
||||
|
||||
/**
|
||||
* 批量删除搬运任务
|
||||
*
|
||||
* @param ids 编号
|
||||
*/
|
||||
void deleteTransportTaskListByIds(List<Long> ids);
|
||||
|
||||
/**
|
||||
* 获得搬运任务
|
||||
*
|
||||
* @param id 编号
|
||||
* @return 搬运任务
|
||||
*/
|
||||
TransportTaskDO getTransportTask(Long id);
|
||||
|
||||
/**
|
||||
* 获得搬运任务分页
|
||||
*
|
||||
* @param pageReqVO 分页查询
|
||||
* @return 搬运任务分页
|
||||
*/
|
||||
PageResult<TransportTaskDO> getTransportTaskPage(TransportTaskPageReqVO pageReqVO);
|
||||
|
||||
/**
|
||||
* PC 端操作搬运任务(完成/取消/强制完成)
|
||||
*
|
||||
* @param reqVO 操作请求
|
||||
*/
|
||||
void operateTransportTask(@Valid TransportTaskOperateReqVO reqVO);
|
||||
|
||||
/**
|
||||
* 根据 taskId 查询任务全量信息(RPC 用,自动过滤逻辑删除,查不到返回 null)
|
||||
*
|
||||
* @param taskId 任务ID
|
||||
* @return 任务全量信息
|
||||
*/
|
||||
TaskInfoDTO getTaskInfoById(Long taskId);
|
||||
|
||||
/**
|
||||
* 根据 taskCode 查询任务全量信息(RPC 用,自动过滤逻辑删除,查不到返回 null)
|
||||
*
|
||||
* @param taskCode 任务编码
|
||||
* @return 任务全量信息
|
||||
*/
|
||||
TaskInfoDTO getTaskInfoByCode(String taskCode);
|
||||
|
||||
}
|
||||
|
||||
@@ -1,11 +1,7 @@
|
||||
package cn.code.nl.module.task.service.transporttask;
|
||||
|
||||
import cn.code.nl.framework.common.exception.ServiceException;
|
||||
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;
|
||||
@@ -16,28 +12,19 @@ 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.*;
|
||||
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 cn.code.nl.module.task.manage.TransportTaskOperationManager;
|
||||
import jakarta.annotation.Resource;
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
import org.redisson.api.RLock;
|
||||
import org.redisson.api.RedissonClient;
|
||||
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 实现类
|
||||
@@ -53,36 +40,7 @@ public class TransportTaskServiceImpl implements TransportTaskService {
|
||||
private TransportTaskMapper transportTaskMapper;
|
||||
|
||||
@Resource
|
||||
private TaskEventProducer taskEventProducer;
|
||||
|
||||
@Resource
|
||||
private TaskCommonApiFactory taskCommonApiFactory;
|
||||
|
||||
@Resource
|
||||
private RedissonClient redissonClient;
|
||||
|
||||
/**
|
||||
* 操作类型 -> 处理方法 路由表(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));
|
||||
}
|
||||
private TransportTaskOperationManager transportTaskOperationManager;
|
||||
|
||||
@Override
|
||||
public Long createTransportTask(TransportTaskSaveReqVO createReqVO) {
|
||||
@@ -181,33 +139,7 @@ public class TransportTaskServiceImpl implements TransportTaskService {
|
||||
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) {
|
||||
RLock lock = redissonClient.getLock(String.valueOf(task.getTaskId()));
|
||||
if (lock.tryLock()) {
|
||||
try {
|
||||
BiConsumer<TransportTaskDO, AcsFeedbackReqDTO> handler = operationHandlers.get(type);
|
||||
if (handler == null) {
|
||||
log.warn("操作类型未注册处理器, taskId={}, type={}", task.getTaskId(), type.getCode());
|
||||
return;
|
||||
}
|
||||
handler.accept(task, reqDTO);
|
||||
} catch (Exception ex) {
|
||||
log.error("[messageResend][执行异常][lockKey={}]", task.getTaskId(), ex);
|
||||
} finally {
|
||||
if (lock.isHeldByCurrentThread()) {
|
||||
lock.unlock();
|
||||
}
|
||||
}
|
||||
} else {
|
||||
throw new ServiceException(5007, "任务标识为:" + task.getTaskId() + "的任务正在操作中!");
|
||||
}
|
||||
transportTaskOperationManager.dispatchOperation(task, type, reqDTO);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -228,138 +160,4 @@ public class TransportTaskServiceImpl implements TransportTaskService {
|
||||
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);
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user