Compare commits
2 Commits
c31b925c93
...
feature/20
| Author | SHA1 | Date | |
|---|---|---|---|
| d9f21147aa | |||
| 1e3f95eace |
67
docs/superpowers/plans/2026-08-03-task-callback-mode.md
Normal file
67
docs/superpowers/plans/2026-08-03-task-callback-mode.md
Normal file
@@ -0,0 +1,67 @@
|
||||
# 任务完成与取消双通道回调实现计划
|
||||
|
||||
> **面向 AI 代理的工作者:** 在当前会话逐任务实现;步骤使用复选框跟踪。当前工作区存在用户改动,不自动提交、不覆盖无关文件。
|
||||
|
||||
**目标:** 实现 YAML 驱动的 MQ 异步与 Feign 同步完成/取消链路,并在业务成功后统一推进任务最终状态。
|
||||
|
||||
**架构:** Task 操作管理器负责配置分流,execute starter 提供统一 Feign 业务入口,MQ Consumer 在业务处理后回调 Task。最终状态仅由 `TransportTaskCallbackService` 收口更新。
|
||||
|
||||
**技术栈:** Java 17、Spring Boot、OpenFeign、RocketMQ、MyBatis、JUnit 5
|
||||
|
||||
---
|
||||
|
||||
### 任务 1:补齐 execute starter 的完成与取消入口
|
||||
|
||||
**文件:**
|
||||
- 修改:`nl-framework/nl-spring-boot-starter-execute/pom.xml`
|
||||
- 修改:`nl-framework/nl-spring-boot-starter-execute/src/main/java/cn/code/nl/framework/execute/biz/api/TaskCommonApi.java`
|
||||
- 修改:`nl-framework/nl-spring-boot-starter-execute/src/main/java/cn/code/nl/framework/execute/biz/api/AbstractTaskCommonApiImpl.java`
|
||||
- 创建:`nl-framework/nl-spring-boot-starter-execute/src/test/java/cn/code/nl/framework/execute/biz/api/AbstractTaskCommonApiImplLocalSpringTest.java`
|
||||
|
||||
- [x] 编写 Spring 集成测试,调用 `doHandleFinished`、`doHandleCancelled` 并断言测试任务处理器收到正确 taskId 与 payload。
|
||||
- [x] 运行定向测试,确认因接口缺失而失败。
|
||||
- [x] 在 `TaskCommonApi` 和公共实现中增加两个接口,复用现有 `execute` 路由。
|
||||
- [x] 重跑定向测试并确认通过。
|
||||
|
||||
### 任务 2:增加 YAML 模式分流并统一同步状态收口
|
||||
|
||||
**文件:**
|
||||
- 创建:`nl-module-task/nl-module-task-api/src/main/java/cn/code/nl/module/task/enums/TaskCallbackTypeEnum.java`
|
||||
- 修改:`nl-module-task/nl-module-task-server/src/main/resources/application.yaml`
|
||||
- 修改:`nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/manage/TransportTaskOperationManager.java`
|
||||
- 修改:`nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/service/transporttask/TransportTaskCallbackServiceImpl.java`
|
||||
|
||||
- [x] 定义 `mq`、`feign` 配置枚举及严格解析方法。
|
||||
- [x] 为 Task 服务增加 `nl.task.callback-type: mq`。
|
||||
- [x] 将完成、取消共同的待处理状态保存后,按配置选择事务后发 MQ 或同步调用 `TaskCommonApi`。
|
||||
- [x] Feign 成功后构造成功回调请求并调用 `TransportTaskCallbackService`。
|
||||
- [x] 删除字符串状态码的错误格式化,直接使用枚举 code。
|
||||
|
||||
### 任务 3:补齐 MQ 成功与失败回调
|
||||
|
||||
**文件:**
|
||||
- 修改:`nl-module-lms/nl-module-lms-server/src/main/java/cn/code/nl/module/lms/mq/consumer/LmsTaskStatusChangeConsumer.java`
|
||||
- 修改:`nl-module-wms/nl-module-wms-server/src/main/java/cn/code/nl/module/wms/mq/consumer/WmsTaskStatusChangeConsumer.java`
|
||||
|
||||
- [x] 完成或取消业务正常结束后,通过 `TransportTaskApi.receiveCallbackResult` 回调 `SUCCESS` 并检查返回结果。
|
||||
- [x] 业务异常时回调 `FAILED` 记录异常信息,再抛出原异常触发 MQ 重试。
|
||||
- [x] 未找到处理器时抛出异常,避免消息被确认后任务永久停留在待处理状态。
|
||||
- [x] 统一事件标识并使用 Redis 成功标识隔离业务执行与成功回调重试。
|
||||
|
||||
### 任务 4:验证与审查
|
||||
|
||||
**文件:**
|
||||
- 检查:以上全部生产与测试文件
|
||||
|
||||
- [ ] 运行 execute starter 定向 Spring 集成测试。
|
||||
- [ ] 编译 execute starter、task-api、task-server、lms-server、wms-server 及依赖模块。
|
||||
- [ ] 检查 git diff,确认没有覆盖工作区既有改动。
|
||||
- [ ] 按 Java 规范检查中文注释、`@Resource` 注入、Feign 返回值检查及异常日志。
|
||||
|
||||
### 最终验证结果
|
||||
|
||||
- [x] execute starter 4 个 Spring 定向测试通过。
|
||||
- [x] task-api 3 个业务归属与 Topic 映射 Spring 定向测试通过。
|
||||
- [x] Execute、Task、LMS、WMS 及依赖共 33 个模块编译通过。
|
||||
- [x] 聚合启动工程 `nl-server` 及依赖共 38 个模块编译通过。
|
||||
- [x] `git diff --check` 通过,仅有工作区行尾转换提示。
|
||||
@@ -0,0 +1,72 @@
|
||||
# 任务完成与取消双通道回调设计
|
||||
|
||||
## 目标
|
||||
|
||||
任务模块根据 YAML 配置选择 MQ 异步或 Feign 同步方式通知 LMS/WMS 执行完成、取消业务;业务处理成功后,两条链路都由 Task 服务统一把任务从回调待处理状态推进到最终完成或取消状态。
|
||||
|
||||
## 当前问题
|
||||
|
||||
- `TransportTaskOperationManager#handleFinished`、`handleCancelled` 固定发布 RocketMQ 消息,无法切换为同步 Feign。
|
||||
- LMS/WMS 的 MQ Consumer 调用 `AbstractTask#doHandleFinish`、`doHandleCancel` 后没有回调 Task,任务停留在 `FINISHED_CALLBACK_PENDING` 或 `CANCEL_CALLBACK_PENDING`。
|
||||
- `TransportTaskCallbackServiceImpl` 使用 `%03d` 格式化字符串状态码,回调处理存在运行时类型异常。
|
||||
- `TaskCommonApi` 只暴露取货、进出等动作,没有完成和取消接口。
|
||||
|
||||
## 配置
|
||||
|
||||
在 Task 服务 YAML 中增加必填配置:
|
||||
|
||||
```yaml
|
||||
nl:
|
||||
task:
|
||||
callback-type: mq
|
||||
```
|
||||
|
||||
允许值为 `mq`、`feign`。使用枚举解析配置;配置缺失或值非法时启动失败,不做静默兜底。默认项目配置使用 `mq`,保持现有部署行为。
|
||||
|
||||
## 执行链路
|
||||
|
||||
### MQ 异步模式
|
||||
|
||||
1. Task 将任务置为对应的回调待处理状态,并把 `callbackStatus` 置为 `PENDING`。
|
||||
2. 当前事务提交后发布完成或取消消息。
|
||||
3. LMS/WMS Consumer 校验任务仍处于当前事件对应的待处理状态,然后路由到 `AbstractTask` 子类执行业务。
|
||||
4. 业务成功后,Consumer 写入带有效期的 Redis 成功标识,再通过 `TransportTaskApi.receiveCallbackResult` 回调 `SUCCESS`;若成功回调失败,MQ 重试时只重试回调,不重复执行业务。
|
||||
5. Task 的统一回调服务把任务更新为 `FINISHED` 或 `CANCELLED`,并把 `callbackStatus` 置为 `SUCCESS`。
|
||||
6. 业务异常时,Consumer 尝试回调 `FAILED` 记录错误与重试次数,再重新抛出异常,让 MQ 执行重试。
|
||||
|
||||
### Feign 同步模式
|
||||
|
||||
1. Task 将任务置为对应的回调待处理状态,并把 `callbackStatus` 置为 `PENDING`。
|
||||
2. Task 根据 `ownerService` 从 `TaskCommonApiFactory` 获取 LMS/WMS Feign Client。
|
||||
3. 调用新增的 `doHandleFinished` 或 `doHandleCancelled`,业务服务通过 `TaskFactory` 路由到相同的 `AbstractTask` 子类。
|
||||
4. Feign 返回成功后,Task 直接调用统一回调服务,以相同规则更新最终状态。
|
||||
5. Feign 返回失败或抛出异常时不进入最终状态,异常沿现有任务操作链路上抛。
|
||||
|
||||
## 接口与数据
|
||||
|
||||
- 在 `TaskCommonApi` 增加 `doHandleFinished`、`doHandleCancelled`。
|
||||
- 两个接口继续复用 `TaskStatusCallApiReqDTO`,携带任务标识、业务归属、业务标识、处理器编码和 payload。
|
||||
- 同步与异步链路都复用 `TaskCallbackResultReqDTO` 和 `TransportTaskCallbackService` 完成最终状态收口。
|
||||
- `eventId` 使用任务 ID 与事件类型组成的稳定值;Task 按当前 pending 状态校验 eventId,拒绝延迟或错序事件推进错误终态。
|
||||
|
||||
## 幂等与错误处理
|
||||
|
||||
- Task 已处于 `FINISHED` 或 `CANCELLED` 时忽略重复成功回调。
|
||||
- 只有 `FINISHED_CALLBACK_PENDING` 能推进到 `FINISHED`,只有 `CANCEL_CALLBACK_PENDING` 能推进到 `CANCELLED`。
|
||||
- MQ 重复消息在消费前查询 Task 状态;状态已不匹配时直接跳过。
|
||||
- MQ 业务成功后使用七天 Redis 标识隔离业务执行与 Task 成功回调重试;成功回调完成后删除标识。
|
||||
- Feign `CommonResult` 必须调用 `getCheckedData` 或检查 `isSuccess`,业务失败不允许误更新最终状态。
|
||||
- MQ 回调失败不吞掉业务异常,保留 RocketMQ 重试能力。
|
||||
- 任务状态分布式锁延迟到本地事务完成后释放,避免提交前窗口发生并发重复处理。
|
||||
|
||||
## 验证
|
||||
|
||||
- Spring 集成测试验证 execute starter 的完成、取消 Feign 入口能正确路由到任务处理器。
|
||||
- 编译 execute starter、task-api、task-server、lms-server、wms-server,验证跨模块接口一致。
|
||||
- 检查 MQ 成功回调和 Feign 成功返回最终都进入统一状态处理服务。
|
||||
|
||||
## 一致性边界
|
||||
|
||||
- MQ 和 Feign 均向 `AbstractTask` 透传稳定的 `eventId` 与 `eventType`,业务处理器应以该事件标识在业务数据库中实现幂等。
|
||||
- Redis 成功标识可以避免“业务已成功、Task 成功回调失败”时重复执行业务,但不能覆盖业务事务提交后、Redis 标识写入前进程崩溃的极小窗口,因此业务数据库幂等仍是最终保障。
|
||||
- Feign 调用与 Task 本地事务不构成分布式事务;远端成功后若本地事务回滚,重试仍依赖上述事件幂等。
|
||||
@@ -0,0 +1,53 @@
|
||||
package cn.code.nl.framework.common.pojo;
|
||||
|
||||
import cn.code.nl.framework.common.exception.enums.GlobalErrorCodeConstants;
|
||||
import lombok.Data;
|
||||
|
||||
import java.util.Objects;
|
||||
|
||||
/**
|
||||
* acs 固定返回对象
|
||||
* @Author: liyongde
|
||||
* @Date: 2026/7/30 9:05
|
||||
*/
|
||||
@Data
|
||||
public class AcsCommonResult<T> {
|
||||
/**
|
||||
* 总条数
|
||||
*/
|
||||
private Long totalElements;
|
||||
|
||||
/**
|
||||
* 业务数据体
|
||||
*/
|
||||
private T data;
|
||||
|
||||
/**
|
||||
* 时间戳
|
||||
*/
|
||||
private String timestamp;
|
||||
|
||||
/**
|
||||
* http状态码
|
||||
*/
|
||||
private Integer code;
|
||||
|
||||
/**
|
||||
* 消息
|
||||
*/
|
||||
private String message;
|
||||
|
||||
/**
|
||||
* 业务响应码 todo: 暂时不用
|
||||
*/
|
||||
private Integer respCode;
|
||||
|
||||
/**
|
||||
* 业务响应信息 todo: 暂时不用
|
||||
*/
|
||||
private String respMsg;
|
||||
|
||||
public static boolean isSuccess(Integer code) {
|
||||
return Objects.equals(code, GlobalErrorCodeConstants.SUCCESS.getCode()) || Objects.equals(code, 200);
|
||||
}
|
||||
}
|
||||
@@ -47,6 +47,13 @@
|
||||
<artifactId>springdoc-openapi-starter-webmvc-ui</artifactId>
|
||||
<scope>provided</scope> <!-- 设置为 provided,主要是 PageParam 使用到 -->
|
||||
</dependency>
|
||||
|
||||
<!-- Test 测试相关 -->
|
||||
<dependency>
|
||||
<groupId>org.springframework.boot</groupId>
|
||||
<artifactId>spring-boot-starter-test</artifactId>
|
||||
<scope>test</scope>
|
||||
</dependency>
|
||||
</dependencies>
|
||||
|
||||
</project>
|
||||
</project>
|
||||
|
||||
@@ -22,6 +22,29 @@ public abstract class AbstractTaskCommonApiImpl implements TaskCommonApi {
|
||||
@Resource
|
||||
private TaskFactory taskFactory;
|
||||
|
||||
/**
|
||||
* 处理任务完成
|
||||
*/
|
||||
@Override
|
||||
public CommonResult<AcsApplyActionRespVO> doHandleFinished(TaskStatusCallApiReqDTO taskStatusCallApiReqDTO) {
|
||||
return execute(taskStatusCallApiReqDTO, (task, taskExecuteDTO) -> {
|
||||
task.doHandleFinish(taskExecuteDTO);
|
||||
|
||||
return null;
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
* 处理任务取消
|
||||
*/
|
||||
@Override
|
||||
public CommonResult<AcsApplyActionRespVO> doHandleCancelled(TaskStatusCallApiReqDTO taskStatusCallApiReqDTO) {
|
||||
return execute(taskStatusCallApiReqDTO, (task, taskExecuteDTO) -> {
|
||||
task.doHandleCancel(taskExecuteDTO);
|
||||
return null;
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
* 处理取货完成
|
||||
*/
|
||||
@@ -97,6 +120,8 @@ public abstract class AbstractTaskCommonApiImpl implements TaskCommonApi {
|
||||
private TaskExecuteDTO buildTaskExecuteDTO(TaskStatusCallApiReqDTO taskStatusCallApiReqDTO) {
|
||||
return TaskExecuteDTO.builder()
|
||||
.taskId(taskStatusCallApiReqDTO.getTaskId())
|
||||
.eventId(taskStatusCallApiReqDTO.getEventId())
|
||||
.eventType(taskStatusCallApiReqDTO.getStatus())
|
||||
.payload(taskStatusCallApiReqDTO.getPayload())
|
||||
.build();
|
||||
}
|
||||
|
||||
@@ -14,6 +14,14 @@ import org.springframework.web.bind.annotation.RequestBody;
|
||||
*/
|
||||
public interface TaskCommonApi {
|
||||
|
||||
@Operation(summary = "处理任务完成")
|
||||
@PostMapping("/doHandleFinished")
|
||||
CommonResult<AcsApplyActionRespVO> doHandleFinished(@RequestBody TaskStatusCallApiReqDTO taskStatusCallApiReqDTO);
|
||||
|
||||
@Operation(summary = "处理任务取消")
|
||||
@PostMapping("/doHandleCancelled")
|
||||
CommonResult<AcsApplyActionRespVO> doHandleCancelled(@RequestBody TaskStatusCallApiReqDTO taskStatusCallApiReqDTO);
|
||||
|
||||
@Operation(summary = "请求取货")
|
||||
@PostMapping("/do-handle-picked")
|
||||
CommonResult<AcsApplyActionRespVO> doHandlePicked(@RequestBody TaskStatusCallApiReqDTO taskStatusCallApiReqDTO);
|
||||
|
||||
@@ -15,6 +15,9 @@ public class TaskStatusCallApiReqDTO implements Serializable {
|
||||
/** 任务ID */
|
||||
private Long taskId;
|
||||
|
||||
/** 任务事件唯一标识 */
|
||||
private String eventId;
|
||||
|
||||
/** 任务编码 */
|
||||
private String taskCode;
|
||||
|
||||
|
||||
@@ -23,6 +23,12 @@ import java.util.Map;
|
||||
@Component
|
||||
public class TaskCommonApiFactory {
|
||||
|
||||
/** LMS 业务归属编码 */
|
||||
private static final String LMS_OWNER_SERVICE = "LMS";
|
||||
|
||||
/** WMS 业务归属编码 */
|
||||
private static final String WMS_OWNER_SERVICE = "WMS";
|
||||
|
||||
@Autowired(required = false)
|
||||
private LmsTaskCommonApi lmsTaskCommonApi;
|
||||
|
||||
@@ -41,9 +47,11 @@ public class TaskCommonApiFactory {
|
||||
public void init() {
|
||||
if (lmsTaskCommonApi != null) {
|
||||
apiMap.put(RpcConstants.LMS_NAME, lmsTaskCommonApi);
|
||||
apiMap.put(LMS_OWNER_SERVICE, lmsTaskCommonApi);
|
||||
}
|
||||
if (wmsTaskCommonApi != null) {
|
||||
apiMap.put(RpcConstants.WMS_NAME, wmsTaskCommonApi);
|
||||
apiMap.put(WMS_OWNER_SERVICE, wmsTaskCommonApi);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -24,6 +24,12 @@ public class TaskExecuteDTO {
|
||||
@NotNull(message = "任务ID不能为空")
|
||||
private Long taskId;
|
||||
|
||||
@Schema(description = "任务事件唯一标识")
|
||||
private String eventId;
|
||||
|
||||
@Schema(description = "任务事件类型")
|
||||
private String eventType;
|
||||
|
||||
@Schema(description = "ACS反馈扩展数据")
|
||||
private Map<String, Object> payload;
|
||||
}
|
||||
|
||||
@@ -1,17 +1,22 @@
|
||||
package cn.code.nl.framework.execute.util.http;
|
||||
|
||||
import cn.code.nl.framework.common.exception.ServiceException;
|
||||
import cn.code.nl.framework.common.pojo.AcsBaseRespDTO;
|
||||
import cn.code.nl.framework.common.exception.enums.GlobalErrorCodeConstants;
|
||||
import cn.code.nl.framework.common.pojo.AcsCommonResult;
|
||||
import cn.code.nl.framework.common.util.http.HttpUtils;
|
||||
import cn.code.nl.framework.common.util.json.JsonUtils;
|
||||
import cn.code.nl.framework.common.util.spring.SpringUtils;
|
||||
import cn.code.nl.module.infra.api.config.ConfigApi;
|
||||
import cn.hutool.core.util.IdUtil;
|
||||
import cn.hutool.core.util.ObjectUtil;
|
||||
import cn.hutool.core.util.StrUtil;
|
||||
import com.alibaba.fastjson.JSON;
|
||||
import com.fasterxml.jackson.core.type.TypeReference;
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
|
||||
import java.net.*;
|
||||
import java.net.ConnectException;
|
||||
import java.net.NoRouteToHostException;
|
||||
import java.net.SocketException;
|
||||
import java.net.SocketTimeoutException;
|
||||
import java.net.UnknownHostException;
|
||||
import java.util.Collections;
|
||||
import java.util.Map;
|
||||
|
||||
@@ -32,29 +37,32 @@ public class AcsUtil {
|
||||
private static final String ACS_DISABLED = "0";
|
||||
|
||||
/**
|
||||
* 发送 POST 请求并解析响应
|
||||
* 发送 POST 请求并解析 ACS 通用响应
|
||||
*
|
||||
* @param serverAddress ACS 服务地址
|
||||
* @param api API 路径
|
||||
* @param request 请求参数
|
||||
* @param responseType 响应类型
|
||||
* @return 响应对象
|
||||
* @param api API 路径
|
||||
* @param request 请求参数
|
||||
* @param responseType data 响应类型
|
||||
* @return ACS 通用响应
|
||||
*/
|
||||
public static <T> T post(String serverAddress, String api, Object request, Class<T> responseType) {
|
||||
public static <T> AcsCommonResult<T> post(String serverAddress, String api, Object request, Class<T> responseType) {
|
||||
if (isAcsDisabled()) {
|
||||
return buildDefaultSuccessResponse(responseType);
|
||||
return buildDefaultSuccessResponse();
|
||||
}
|
||||
String url = buildUrl(serverAddress, api);
|
||||
String response;
|
||||
try {
|
||||
response = HttpUtils.post(url, headers(), JsonUtils.toJsonString(request));
|
||||
log.info("ACS 返回值:{}", response);
|
||||
} catch (RuntimeException ex) {
|
||||
if (isNetworkException(ex)) {
|
||||
throw new ServiceException(500, "ACS服务网络不通");
|
||||
log.error("ACS 服务网络不通,url={},request={}", url, JsonUtils.toJsonString(request), ex);
|
||||
return buildErrorResponse("ACS服务网络不通");
|
||||
}
|
||||
throw ex;
|
||||
log.error("ACS 请求失败,url={},request={}", url, JsonUtils.toJsonString(request), ex);
|
||||
return buildErrorResponse(ex.getMessage());
|
||||
}
|
||||
return JsonUtils.parseObject(response, responseType);
|
||||
return parseResponse(response, responseType);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -82,13 +90,40 @@ public class AcsUtil {
|
||||
/**
|
||||
* 构建默认成功响应
|
||||
*/
|
||||
private static <T> T buildDefaultSuccessResponse(Class<T> responseType) {
|
||||
AcsBaseRespDTO response = new AcsBaseRespDTO();
|
||||
response.setSuccess(true);
|
||||
response.setCode("200");
|
||||
response.setTraceId(IdUtil.simpleUUID());
|
||||
response.setTimestamp(System.currentTimeMillis());
|
||||
return JsonUtils.convertObject(response, responseType);
|
||||
private static <T> AcsCommonResult<T> buildDefaultSuccessResponse() {
|
||||
AcsCommonResult<T> response = new AcsCommonResult<>();
|
||||
response.setCode(GlobalErrorCodeConstants.SUCCESS.getCode());
|
||||
response.setMessage("ACS 已关闭,跳过请求");
|
||||
response.setTimestamp(String.valueOf(System.currentTimeMillis()));
|
||||
return response;
|
||||
}
|
||||
|
||||
/**
|
||||
* 构建失败响应,具体是否抛异常由业务层决定
|
||||
*/
|
||||
private static <T> AcsCommonResult<T> buildErrorResponse(String message) {
|
||||
AcsCommonResult<T> response = new AcsCommonResult<>();
|
||||
response.setCode(GlobalErrorCodeConstants.INTERNAL_SERVER_ERROR.getCode());
|
||||
response.setMessage(StrUtil.blankToDefault(message, "ACS 请求失败"));
|
||||
response.setTimestamp(String.valueOf(System.currentTimeMillis()));
|
||||
return response;
|
||||
}
|
||||
|
||||
/**
|
||||
* 解析 ACS 响应,并将 data 转为业务指定类型
|
||||
*/
|
||||
private static <T> AcsCommonResult<T> parseResponse(String response, Class<T> responseType) {
|
||||
AcsCommonResult<Object> rawResult = JsonUtils.parseObject(response, new TypeReference<AcsCommonResult<Object>>() {
|
||||
});
|
||||
AcsCommonResult<T> result = new AcsCommonResult<>();
|
||||
result.setTotalElements(rawResult.getTotalElements());
|
||||
result.setData(JsonUtils.convertObject(rawResult.getData(), responseType));
|
||||
result.setTimestamp(rawResult.getTimestamp());
|
||||
result.setCode(rawResult.getCode());
|
||||
result.setMessage(rawResult.getMessage());
|
||||
result.setRespCode(rawResult.getRespCode());
|
||||
result.setRespMsg(rawResult.getRespMsg());
|
||||
return result;
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -0,0 +1,126 @@
|
||||
package cn.code.nl.framework.execute.biz.api;
|
||||
|
||||
import cn.code.nl.framework.execute.biz.dto.TaskStatusCallApiReqDTO;
|
||||
import cn.code.nl.framework.execute.core.AbstractTask;
|
||||
import cn.code.nl.framework.execute.core.TaskFactory;
|
||||
import cn.code.nl.framework.execute.core.dto.TaskExecuteDTO;
|
||||
import jakarta.annotation.Resource;
|
||||
import org.junit.jupiter.api.Assertions;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.springframework.boot.test.context.SpringBootTest;
|
||||
import org.springframework.stereotype.Component;
|
||||
|
||||
import java.util.Map;
|
||||
|
||||
/**
|
||||
* 任务通用 Feign 接口本地 Spring 集成测试
|
||||
*/
|
||||
@SpringBootTest(classes = {
|
||||
TaskFactory.class,
|
||||
AbstractTaskCommonApiImplLocalSpringTest.TestTaskCommonApi.class,
|
||||
AbstractTaskCommonApiImplLocalSpringTest.TestTask.class
|
||||
})
|
||||
class AbstractTaskCommonApiImplLocalSpringTest {
|
||||
|
||||
@Resource
|
||||
private TestTaskCommonApi testTaskCommonApi;
|
||||
|
||||
@Resource
|
||||
private TestTask testTask;
|
||||
|
||||
/**
|
||||
* 验证完成接口按 handleCode 路由并传递任务参数
|
||||
*/
|
||||
@Test
|
||||
void doHandleFinishedShouldRouteTask() {
|
||||
TaskStatusCallApiReqDTO reqDTO = buildReqDTO();
|
||||
|
||||
testTaskCommonApi.doHandleFinished(reqDTO).getCheckedData();
|
||||
|
||||
Assertions.assertNotNull(testTask.getFinishedDTO());
|
||||
Assertions.assertEquals(1001L, testTask.getFinishedDTO().getTaskId());
|
||||
Assertions.assertEquals("1001_TASK_FINISHED", testTask.getFinishedDTO().getEventId());
|
||||
Assertions.assertEquals("TASK_FINISHED", testTask.getFinishedDTO().getEventType());
|
||||
Assertions.assertEquals("value", testTask.getFinishedDTO().getPayload().get("key"));
|
||||
}
|
||||
|
||||
/**
|
||||
* 验证取消接口按 handleCode 路由并传递任务参数
|
||||
*/
|
||||
@Test
|
||||
void doHandleCancelledShouldRouteTask() {
|
||||
TaskStatusCallApiReqDTO reqDTO = buildReqDTO();
|
||||
reqDTO.setEventId("1001_TASK_CANCELLED");
|
||||
reqDTO.setStatus("TASK_CANCELLED");
|
||||
|
||||
testTaskCommonApi.doHandleCancelled(reqDTO).getCheckedData();
|
||||
|
||||
Assertions.assertNotNull(testTask.getCancelledDTO());
|
||||
Assertions.assertEquals(1001L, testTask.getCancelledDTO().getTaskId());
|
||||
Assertions.assertEquals("1001_TASK_CANCELLED", testTask.getCancelledDTO().getEventId());
|
||||
Assertions.assertEquals("TASK_CANCELLED", testTask.getCancelledDTO().getEventType());
|
||||
Assertions.assertEquals("value", testTask.getCancelledDTO().getPayload().get("key"));
|
||||
}
|
||||
|
||||
/**
|
||||
* 构建任务状态调用参数
|
||||
*/
|
||||
private TaskStatusCallApiReqDTO buildReqDTO() {
|
||||
TaskStatusCallApiReqDTO reqDTO = new TaskStatusCallApiReqDTO();
|
||||
reqDTO.setTaskId(1001L);
|
||||
reqDTO.setTaskCode("TASK-1001");
|
||||
reqDTO.setEventId("1001_TASK_FINISHED");
|
||||
reqDTO.setStatus("TASK_FINISHED");
|
||||
reqDTO.setHandleCode("testTask");
|
||||
reqDTO.setPayload(Map.of("key", "value"));
|
||||
return reqDTO;
|
||||
}
|
||||
|
||||
/**
|
||||
* 测试任务通用接口实现
|
||||
*/
|
||||
@Component
|
||||
static class TestTaskCommonApi extends AbstractTaskCommonApiImpl {
|
||||
}
|
||||
|
||||
/**
|
||||
* 测试任务处理器
|
||||
*/
|
||||
@Component("testTask")
|
||||
static class TestTask extends AbstractTask {
|
||||
|
||||
private TaskExecuteDTO finishedDTO;
|
||||
|
||||
private TaskExecuteDTO cancelledDTO;
|
||||
|
||||
/**
|
||||
* 记录完成任务参数
|
||||
*/
|
||||
@Override
|
||||
public void doHandleFinish(TaskExecuteDTO taskExecuteDTO) {
|
||||
this.finishedDTO = taskExecuteDTO;
|
||||
}
|
||||
|
||||
/**
|
||||
* 记录取消任务参数
|
||||
*/
|
||||
@Override
|
||||
public void doHandleCancel(TaskExecuteDTO taskExecuteDTO) {
|
||||
this.cancelledDTO = taskExecuteDTO;
|
||||
}
|
||||
|
||||
/**
|
||||
* 获取完成任务参数
|
||||
*/
|
||||
public TaskExecuteDTO getFinishedDTO() {
|
||||
return finishedDTO;
|
||||
}
|
||||
|
||||
/**
|
||||
* 获取取消任务参数
|
||||
*/
|
||||
public TaskExecuteDTO getCancelledDTO() {
|
||||
return cancelledDTO;
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,61 @@
|
||||
package cn.code.nl.framework.execute.core;
|
||||
|
||||
import cn.code.nl.framework.execute.biz.api.AbstractTaskCommonApiImpl;
|
||||
import cn.code.nl.framework.execute.biz.api.lms.LmsTaskCommonApi;
|
||||
import cn.code.nl.framework.execute.biz.api.wms.WmsTaskCommonApi;
|
||||
import jakarta.annotation.Resource;
|
||||
import org.junit.jupiter.api.Assertions;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.springframework.boot.test.context.SpringBootTest;
|
||||
import org.springframework.stereotype.Component;
|
||||
|
||||
/**
|
||||
* 任务通用 API 工厂本地 Spring 集成测试
|
||||
*/
|
||||
@SpringBootTest(classes = {
|
||||
TaskFactory.class,
|
||||
TaskCommonApiFactory.class,
|
||||
TaskCommonApiFactoryLocalSpringTest.TestLmsTaskCommonApi.class,
|
||||
TaskCommonApiFactoryLocalSpringTest.TestWmsTaskCommonApi.class
|
||||
})
|
||||
class TaskCommonApiFactoryLocalSpringTest {
|
||||
|
||||
@Resource
|
||||
private TaskCommonApiFactory taskCommonApiFactory;
|
||||
|
||||
@Resource
|
||||
private TestLmsTaskCommonApi testLmsTaskCommonApi;
|
||||
|
||||
@Resource
|
||||
private TestWmsTaskCommonApi testWmsTaskCommonApi;
|
||||
|
||||
/**
|
||||
* 验证业务归属编码 LMS 能定位 Feign 接口
|
||||
*/
|
||||
@Test
|
||||
void getByServerNameShouldSupportLmsOwnerService() {
|
||||
Assertions.assertSame(testLmsTaskCommonApi, taskCommonApiFactory.getByServerName("LMS"));
|
||||
}
|
||||
|
||||
/**
|
||||
* 验证业务归属编码 WMS 能定位 Feign 接口
|
||||
*/
|
||||
@Test
|
||||
void getByServerNameShouldSupportWmsOwnerService() {
|
||||
Assertions.assertSame(testWmsTaskCommonApi, taskCommonApiFactory.getByServerName("WMS"));
|
||||
}
|
||||
|
||||
/**
|
||||
* 测试 LMS Feign 接口
|
||||
*/
|
||||
@Component
|
||||
static class TestLmsTaskCommonApi extends AbstractTaskCommonApiImpl implements LmsTaskCommonApi {
|
||||
}
|
||||
|
||||
/**
|
||||
* 测试 WMS Feign 接口
|
||||
*/
|
||||
@Component
|
||||
static class TestWmsTaskCommonApi extends AbstractTaskCommonApiImpl implements WmsTaskCommonApi {
|
||||
}
|
||||
}
|
||||
@@ -5,17 +5,22 @@ import cn.code.nl.framework.execute.core.AbstractTask;
|
||||
import cn.code.nl.framework.execute.core.TaskFactory;
|
||||
import cn.code.nl.framework.execute.core.dto.TaskExecuteDTO;
|
||||
import cn.code.nl.module.task.api.TransportTaskApi;
|
||||
import cn.code.nl.module.task.dto.TaskCallbackResultReqDTO;
|
||||
import cn.code.nl.module.task.dto.TaskInfoDTO;
|
||||
import cn.code.nl.module.task.enums.CallbackStatusEnum;
|
||||
import cn.code.nl.module.task.enums.TransportTaskStatusEnum;
|
||||
import cn.code.nl.module.task.message.TaskEventMessage;
|
||||
import jakarta.annotation.Resource;
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
|
||||
import org.apache.rocketmq.spring.core.RocketMQListener;
|
||||
import org.redisson.api.RBucket;
|
||||
import org.redisson.api.RLock;
|
||||
import org.redisson.api.RedissonClient;
|
||||
import org.springframework.stereotype.Component;
|
||||
|
||||
import java.time.Duration;
|
||||
|
||||
import static cn.code.nl.module.task.message.LockKeyConstants.TASK_STATUS_CHANGE_LOCK_KEY;
|
||||
|
||||
/**
|
||||
@@ -32,6 +37,12 @@ import static cn.code.nl.module.task.message.LockKeyConstants.TASK_STATUS_CHANGE
|
||||
)
|
||||
public class LmsTaskStatusChangeConsumer implements RocketMQListener<TaskEventMessage> {
|
||||
|
||||
/** 任务业务处理成功标识前缀 */
|
||||
private static final String TASK_BUSINESS_HANDLED_KEY = "task:business:handled:";
|
||||
|
||||
/** 任务业务处理成功标识有效期 */
|
||||
private static final Duration TASK_BUSINESS_HANDLED_TIMEOUT = Duration.ofDays(7);
|
||||
|
||||
@Resource
|
||||
private TaskFactory taskFactory;
|
||||
|
||||
@@ -56,7 +67,19 @@ public class LmsTaskStatusChangeConsumer implements RocketMQListener<TaskEventMe
|
||||
if (!isTaskCallbackPending(message)) {
|
||||
return;
|
||||
}
|
||||
executeTaskHandler(message, handleCode, eventType);
|
||||
RBucket<Boolean> handledBucket = redissonClient.getBucket(
|
||||
TASK_BUSINESS_HANDLED_KEY + message.getEventId());
|
||||
if (!Boolean.TRUE.equals(handledBucket.get())) {
|
||||
try {
|
||||
executeTaskHandler(message, handleCode, eventType);
|
||||
handledBucket.set(true, TASK_BUSINESS_HANDLED_TIMEOUT);
|
||||
} catch (RuntimeException exception) {
|
||||
callbackFailedResult(message, exception);
|
||||
throw exception;
|
||||
}
|
||||
}
|
||||
callbackTaskResult(message, CallbackStatusEnum.SUCCESS.getCode(), "任务业务处理成功");
|
||||
handledBucket.delete();
|
||||
} finally {
|
||||
if (lock.isHeldByCurrentThread()) {
|
||||
lock.unlock();
|
||||
@@ -99,12 +122,13 @@ public class LmsTaskStatusChangeConsumer implements RocketMQListener<TaskEventMe
|
||||
private void executeTaskHandler(TaskEventMessage message, String handleCode, String eventType) {
|
||||
AbstractTask task = taskFactory.getTask(handleCode);
|
||||
if (task == null) {
|
||||
log.warn("未找到对应任务处理器, handleCode={}", handleCode);
|
||||
return;
|
||||
throw new IllegalStateException("未找到对应任务处理器:" + handleCode);
|
||||
}
|
||||
|
||||
TaskExecuteDTO dto = new TaskExecuteDTO();
|
||||
dto.setTaskId(message.getTaskId());
|
||||
dto.setEventId(message.getEventId());
|
||||
dto.setEventType(message.getEventType());
|
||||
dto.setPayload(message.getPayload());
|
||||
|
||||
if (TaskEventMessage.EVENT_TYPE_FINISHED.equals(eventType)) {
|
||||
@@ -115,4 +139,31 @@ public class LmsTaskStatusChangeConsumer implements RocketMQListener<TaskEventMe
|
||||
log.warn("未知事件类型, eventType={}, handleCode={}", eventType, handleCode);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 回调 Task 服务记录业务处理结果
|
||||
*/
|
||||
private void callbackTaskResult(TaskEventMessage message, String result, String resultMessage) {
|
||||
TaskCallbackResultReqDTO reqDTO = new TaskCallbackResultReqDTO();
|
||||
reqDTO.setTaskId(message.getTaskId());
|
||||
reqDTO.setEventId(message.getEventId());
|
||||
reqDTO.setResult(result);
|
||||
reqDTO.setMessage(resultMessage);
|
||||
CommonResult<Boolean> callbackResult = transportTaskApi.receiveCallbackResult(reqDTO);
|
||||
if (!callbackResult.isSuccess()) {
|
||||
log.error("任务业务结果回调失败, reqDTO={}, callbackResult={}", reqDTO, callbackResult);
|
||||
callbackResult.checkError();
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 尝试回调 Task 服务记录业务处理失败
|
||||
*/
|
||||
private void callbackFailedResult(TaskEventMessage message, RuntimeException exception) {
|
||||
try {
|
||||
callbackTaskResult(message, CallbackStatusEnum.FAILED.getCode(), exception.getMessage());
|
||||
} catch (RuntimeException callbackException) {
|
||||
log.error("任务业务失败结果回调异常, message={}", message, callbackException);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -49,6 +49,13 @@
|
||||
<artifactId>nl-spring-boot-starter-execute</artifactId>
|
||||
</dependency>
|
||||
|
||||
<!-- Test 测试相关 -->
|
||||
<dependency>
|
||||
<groupId>org.springframework.boot</groupId>
|
||||
<artifactId>spring-boot-starter-test</artifactId>
|
||||
<scope>test</scope>
|
||||
</dependency>
|
||||
|
||||
</dependencies>
|
||||
|
||||
</project>
|
||||
</project>
|
||||
|
||||
@@ -0,0 +1,35 @@
|
||||
package cn.code.nl.module.task.enums;
|
||||
|
||||
import lombok.Getter;
|
||||
|
||||
/**
|
||||
* 任务业务回调方式
|
||||
*/
|
||||
@Getter
|
||||
public enum TaskCallbackTypeEnum {
|
||||
|
||||
MQ("mq"),
|
||||
|
||||
FEIGN("feign");
|
||||
|
||||
private final String code;
|
||||
|
||||
TaskCallbackTypeEnum(String code) {
|
||||
this.code = code;
|
||||
}
|
||||
|
||||
/**
|
||||
* 根据编码获取任务业务回调方式
|
||||
*
|
||||
* @param code 回调方式编码
|
||||
* @return 任务业务回调方式
|
||||
*/
|
||||
public static TaskCallbackTypeEnum getByCode(String code) {
|
||||
for (TaskCallbackTypeEnum value : values()) {
|
||||
if (value.getCode().equals(code)) {
|
||||
return value;
|
||||
}
|
||||
}
|
||||
return null;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,50 @@
|
||||
package cn.code.nl.module.task.enums;
|
||||
|
||||
import lombok.AllArgsConstructor;
|
||||
import lombok.Getter;
|
||||
|
||||
/**
|
||||
* 任务业务归属服务枚举
|
||||
*/
|
||||
@Getter
|
||||
@AllArgsConstructor
|
||||
public enum TaskOwnerServiceEnum {
|
||||
|
||||
LMS("lms", "lms-server"),
|
||||
WMS("wms", "wms-server");
|
||||
|
||||
/**
|
||||
* MQ Topic 前缀
|
||||
*/
|
||||
private final String topicPrefix;
|
||||
|
||||
/**
|
||||
* 服务注册名称
|
||||
*/
|
||||
private final String serviceName;
|
||||
|
||||
/**
|
||||
* 根据业务归属编码获取枚举
|
||||
*
|
||||
* @param code 业务归属编码
|
||||
* @return 业务归属服务枚举
|
||||
*/
|
||||
public static TaskOwnerServiceEnum getByCode(String code) {
|
||||
for (TaskOwnerServiceEnum value : values()) {
|
||||
if (value.name().equalsIgnoreCase(code) || value.serviceName.equalsIgnoreCase(code)) {
|
||||
return value;
|
||||
}
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
/**
|
||||
* 构建业务服务消费的 Topic
|
||||
*
|
||||
* @param topic 基础 Topic
|
||||
* @return 完整 Topic
|
||||
*/
|
||||
public String buildTopic(String topic) {
|
||||
return topicPrefix + "_" + topic;
|
||||
}
|
||||
}
|
||||
@@ -44,4 +44,15 @@ public class TaskEventMessage {
|
||||
|
||||
/** 扩展数据 */
|
||||
private Map<String, Object> payload;
|
||||
|
||||
/**
|
||||
* 构建任务事件唯一标识
|
||||
*
|
||||
* @param taskId 任务标识
|
||||
* @param eventType 事件类型
|
||||
* @return 任务事件唯一标识
|
||||
*/
|
||||
public static String buildEventId(Long taskId, String eventType) {
|
||||
return taskId + "_" + eventType;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,56 @@
|
||||
package cn.code.nl.module.task.enums;
|
||||
|
||||
import org.junit.jupiter.api.Assertions;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.springframework.boot.test.context.SpringBootTest;
|
||||
import org.springframework.context.annotation.Configuration;
|
||||
|
||||
/**
|
||||
* 任务业务归属服务本地 Spring 集成测试
|
||||
*/
|
||||
@SpringBootTest(classes = TaskOwnerServiceEnumLocalSpringTest.TestConfiguration.class)
|
||||
class TaskOwnerServiceEnumLocalSpringTest {
|
||||
|
||||
/**
|
||||
* 验证 LMS 业务归属编码生成消费者 Topic
|
||||
*/
|
||||
@Test
|
||||
void buildTopicShouldSupportLmsOwnerService() {
|
||||
TaskOwnerServiceEnum ownerService = TaskOwnerServiceEnum.getByCode("LMS");
|
||||
|
||||
Assertions.assertNotNull(ownerService);
|
||||
Assertions.assertEquals("lms_task_status_change_dev_topic",
|
||||
ownerService.buildTopic("task_status_change_dev_topic"));
|
||||
}
|
||||
|
||||
/**
|
||||
* 验证 WMS 业务归属编码生成消费者 Topic
|
||||
*/
|
||||
@Test
|
||||
void buildTopicShouldSupportWmsOwnerService() {
|
||||
TaskOwnerServiceEnum ownerService = TaskOwnerServiceEnum.getByCode("WMS");
|
||||
|
||||
Assertions.assertNotNull(ownerService);
|
||||
Assertions.assertEquals("wms_task_status_change_dev_topic",
|
||||
ownerService.buildTopic("task_status_change_dev_topic"));
|
||||
}
|
||||
|
||||
/**
|
||||
* 验证兼容服务注册名称并拒绝非法业务归属
|
||||
*/
|
||||
@Test
|
||||
void getByCodeShouldSupportServiceNameAndRejectInvalidCode() {
|
||||
Assertions.assertEquals(TaskOwnerServiceEnum.LMS,
|
||||
TaskOwnerServiceEnum.getByCode("lms-server"));
|
||||
Assertions.assertEquals(TaskOwnerServiceEnum.WMS,
|
||||
TaskOwnerServiceEnum.getByCode("wms-server"));
|
||||
Assertions.assertNull(TaskOwnerServiceEnum.getByCode("invalid-server"));
|
||||
}
|
||||
|
||||
/**
|
||||
* 测试 Spring 配置
|
||||
*/
|
||||
@Configuration
|
||||
static class TestConfiguration {
|
||||
}
|
||||
}
|
||||
@@ -8,9 +8,9 @@ package cn.code.nl.module.task.enums;
|
||||
public interface AcsApiConstants {
|
||||
|
||||
/** 下发任务 */
|
||||
String ACS_TASK_API = "/acs-api/wms/issue-task";
|
||||
String ACS_TASK_API = "/acs-api/wms-to-acs/issue-task";
|
||||
|
||||
/** 检测任务 */
|
||||
String ACS_OPERATE_CHECK_API = "/acs-api/wms/check-enable-operate";
|
||||
String ACS_OPERATE_CHECK_API = "/acs-api/wms-to-acs/check-enable-operate";
|
||||
|
||||
}
|
||||
|
||||
@@ -1,15 +1,14 @@
|
||||
package cn.code.nl.module.task.job.dto;
|
||||
|
||||
import cn.code.nl.framework.common.pojo.AcsBaseRespDTO;
|
||||
import lombok.Data;
|
||||
|
||||
import java.util.List;
|
||||
|
||||
/**
|
||||
* ACS 任务下发响应
|
||||
* ACS 任务下发 data 响应
|
||||
*/
|
||||
@Data
|
||||
public class AcsIssueResultRespDTO extends AcsBaseRespDTO {
|
||||
public class AcsIssueResultRespDTO {
|
||||
|
||||
/**
|
||||
* 下发失败的任务
|
||||
|
||||
@@ -1,12 +1,13 @@
|
||||
package cn.code.nl.module.task.job.dto;
|
||||
|
||||
import cn.code.nl.framework.common.pojo.AcsBaseReqDTO;
|
||||
import lombok.Data;
|
||||
|
||||
/**
|
||||
* ACS 任务操作校验请求
|
||||
*/
|
||||
@Data
|
||||
public class AcsOperateCheckReqDTO {
|
||||
public class AcsOperateCheckReqDTO extends AcsBaseReqDTO {
|
||||
|
||||
/**
|
||||
* 任务标识
|
||||
|
||||
@@ -1,29 +1,18 @@
|
||||
package cn.code.nl.module.task.job.dto;
|
||||
|
||||
import cn.code.nl.framework.common.pojo.AcsBaseRespDTO;
|
||||
import lombok.Data;
|
||||
|
||||
/**
|
||||
* ACS 任务操作校验响应
|
||||
* ACS 任务操作校验 data 响应
|
||||
*/
|
||||
@Data
|
||||
public class AcsOperateCheckRespDTO extends AcsBaseRespDTO {
|
||||
public class AcsOperateCheckRespDTO {
|
||||
|
||||
/**
|
||||
* 是否允许操作
|
||||
*/
|
||||
private Boolean enableOperate;
|
||||
|
||||
/**
|
||||
* 是否允许操作,兼容 ACS 字段
|
||||
*/
|
||||
private Boolean canOperate;
|
||||
|
||||
/**
|
||||
* 是否允许操作,兼容通用 data 字段
|
||||
*/
|
||||
private Boolean data;
|
||||
|
||||
/**
|
||||
* 不允许操作的原因
|
||||
*/
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
package cn.code.nl.module.task.manage;
|
||||
|
||||
import cn.code.nl.framework.common.exception.ServiceException;
|
||||
import cn.code.nl.framework.common.pojo.AcsCommonResult;
|
||||
import cn.code.nl.framework.common.pojo.CommonResult;
|
||||
import cn.code.nl.framework.common.util.json.JsonUtils;
|
||||
import cn.code.nl.framework.execute.util.http.AcsUtil;
|
||||
@@ -74,7 +75,8 @@ public class TransportTaskIssueManager {
|
||||
String requestJson = JsonUtils.toJsonString(acsTasks);
|
||||
try {
|
||||
String serverAddress = getAcsServerAddress(task.getProductArea());
|
||||
AcsIssueResultRespDTO result = AcsUtil.post(serverAddress, ACS_TASK_API, acsTasks, AcsIssueResultRespDTO.class);
|
||||
AcsCommonResult<AcsIssueResultRespDTO> result = AcsUtil.post(serverAddress, ACS_TASK_API, acsTasks,
|
||||
AcsIssueResultRespDTO.class);
|
||||
handleSingleIssueResult(task, result);
|
||||
} catch (ServiceException ex) {
|
||||
throw ex;
|
||||
@@ -95,7 +97,8 @@ public class TransportTaskIssueManager {
|
||||
String requestJson = JsonUtils.toJsonString(acsTasks);
|
||||
try {
|
||||
String serverAddress = getAcsServerAddress(productArea);
|
||||
AcsIssueResultRespDTO result = AcsUtil.post(serverAddress, ACS_TASK_API, acsTasks, AcsIssueResultRespDTO.class);
|
||||
AcsCommonResult<AcsIssueResultRespDTO> result = AcsUtil.post(serverAddress, ACS_TASK_API, acsTasks,
|
||||
AcsIssueResultRespDTO.class);
|
||||
handleIssueResult(tasks, result);
|
||||
} catch (Exception ex) {
|
||||
log.error("自动下发任务失败,productArea={},tasks={}", productArea, requestJson, ex);
|
||||
@@ -125,20 +128,18 @@ public class TransportTaskIssueManager {
|
||||
* @param tasks 本次下发任务
|
||||
* @param result ACS 响应
|
||||
*/
|
||||
private void handleIssueResult(List<TransportTaskDO> tasks, AcsIssueResultRespDTO result) {
|
||||
private void handleIssueResult(List<TransportTaskDO> tasks, AcsCommonResult<AcsIssueResultRespDTO> result) {
|
||||
if (result == null) {
|
||||
tasks.forEach(task -> updateIssueFailed(task.getTaskId(), "ACS 返回为空", null));
|
||||
return;
|
||||
}
|
||||
String resultJson = JsonUtils.toJsonString(result);
|
||||
List<AcsIssueResultRespDTO.FailedTask> failedTasks = result.getFailedTasks() == null
|
||||
? Collections.emptyList()
|
||||
: result.getFailedTasks();
|
||||
if (Boolean.FALSE.equals(result.getSuccess()) && CollUtil.isEmpty(failedTasks)) {
|
||||
String errorMessage = StrUtil.blankToDefault(result.getMsg(), "ACS 下发失败");
|
||||
if (!AcsCommonResult.isSuccess(result.getCode())) {
|
||||
String errorMessage = StrUtil.blankToDefault(result.getMessage(), "ACS 下发失败");
|
||||
tasks.forEach(task -> updateIssueFailed(task.getTaskId(), errorMessage, resultJson));
|
||||
return;
|
||||
}
|
||||
List<AcsIssueResultRespDTO.FailedTask> failedTasks = getFailedTasks(result);
|
||||
Map<Long, AcsIssueResultRespDTO.FailedTask> failedTaskMap = failedTasks.stream()
|
||||
.filter(failedTask -> failedTask.getTaskId() != null)
|
||||
.collect(Collectors.toMap(AcsIssueResultRespDTO.FailedTask::getTaskId, Function.identity(), (first, second) -> first));
|
||||
@@ -158,21 +159,34 @@ public class TransportTaskIssueManager {
|
||||
* @param task 本次下发任务
|
||||
* @param result ACS 响应
|
||||
*/
|
||||
private void handleSingleIssueResult(TransportTaskDO task, AcsIssueResultRespDTO result) {
|
||||
private void handleSingleIssueResult(TransportTaskDO task, AcsCommonResult<AcsIssueResultRespDTO> result) {
|
||||
if (result == null) {
|
||||
throw exception(TRANSPORT_TASK_ISSUE_FAILED, "ACS 返回为空");
|
||||
}
|
||||
String resultJson = JsonUtils.toJsonString(result);
|
||||
if (Boolean.TRUE.equals(result.getSuccess())) {
|
||||
if (!AcsCommonResult.isSuccess(result.getCode())) {
|
||||
throw exception(TRANSPORT_TASK_ISSUE_FAILED, StrUtil.blankToDefault(result.getMessage(), "ACS 下发失败"));
|
||||
}
|
||||
List<AcsIssueResultRespDTO.FailedTask> failedTasks = getFailedTasks(result);
|
||||
if (CollUtil.isEmpty(failedTasks)) {
|
||||
transportTaskMapper.updateIssueResult(task.getTaskId(), TransportTaskStatusEnum.ISSUED.getCode(), resultJson, null);
|
||||
return;
|
||||
}
|
||||
String errorMessage = CollUtil.isNotEmpty(result.getFailedTasks())
|
||||
? result.getFailedTasks().get(0).getErrorMessage()
|
||||
: result.getMsg();
|
||||
String errorMessage = failedTasks.get(0).getErrorMessage();
|
||||
throw exception(TRANSPORT_TASK_ISSUE_FAILED, StrUtil.blankToDefault(errorMessage, "ACS 下发失败"));
|
||||
}
|
||||
|
||||
/**
|
||||
* 获取下发失败任务
|
||||
*/
|
||||
private List<AcsIssueResultRespDTO.FailedTask> getFailedTasks(AcsCommonResult<AcsIssueResultRespDTO> result) {
|
||||
AcsIssueResultRespDTO data = result.getData();
|
||||
if (data == null || data.getFailedTasks() == null) {
|
||||
return Collections.emptyList();
|
||||
}
|
||||
return data.getFailedTasks();
|
||||
}
|
||||
|
||||
/**
|
||||
* 更新任务为下发失败
|
||||
*
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
package cn.code.nl.module.task.manage;
|
||||
|
||||
import cn.code.nl.framework.common.exception.ServiceException;
|
||||
import cn.code.nl.framework.common.pojo.AcsCommonResult;
|
||||
import cn.code.nl.framework.common.pojo.CommonResult;
|
||||
import cn.code.nl.framework.common.util.json.JsonUtils;
|
||||
import cn.code.nl.framework.execute.util.http.AcsUtil;
|
||||
@@ -48,7 +49,7 @@ public class TransportTaskOperateCheckManager {
|
||||
}
|
||||
String serverAddress = getAcsServerAddress(task.getProductArea());
|
||||
AcsOperateCheckReqDTO reqDTO = buildReqDTO(task, type);
|
||||
AcsOperateCheckRespDTO result;
|
||||
AcsCommonResult<AcsOperateCheckRespDTO> result;
|
||||
try {
|
||||
result = AcsUtil.post(serverAddress, ACS_OPERATE_CHECK_API, reqDTO, AcsOperateCheckRespDTO.class);
|
||||
} catch (ServiceException ex) {
|
||||
@@ -103,15 +104,23 @@ public class TransportTaskOperateCheckManager {
|
||||
/**
|
||||
* 处理 ACS 操作校验结果
|
||||
*/
|
||||
private void handleCheckResult(AcsOperateCheckReqDTO reqDTO, AcsOperateCheckRespDTO result) {
|
||||
private void handleCheckResult(AcsOperateCheckReqDTO reqDTO, AcsCommonResult<AcsOperateCheckRespDTO> result) {
|
||||
if (result == null) {
|
||||
throw exception(TRANSPORT_TASK_ACS_OPERATE_CHECK_RESULT_EMPTY);
|
||||
}
|
||||
Boolean enableOperate = getEnableOperate(result);
|
||||
if (!AcsCommonResult.isSuccess(result.getCode())) {
|
||||
log.warn("ACS 操作校验失败,reqDTO={},result={}", JsonUtils.toJsonString(reqDTO), JsonUtils.toJsonString(result));
|
||||
throw exception(TRANSPORT_TASK_ACS_OPERATE_CHECK_FAILED);
|
||||
}
|
||||
AcsOperateCheckRespDTO data = result.getData();
|
||||
if (data == null) {
|
||||
throw exception(TRANSPORT_TASK_ACS_OPERATE_CHECK_RESULT_EMPTY);
|
||||
}
|
||||
Boolean enableOperate = data.getEnableOperate();
|
||||
if (Boolean.TRUE.equals(enableOperate)) {
|
||||
return;
|
||||
}
|
||||
String message = getMessage(result);
|
||||
String message = getMessage(result, data);
|
||||
log.warn("ACS 拒绝 PC 端任务操作,reqDTO={},result={}",
|
||||
JsonUtils.toJsonString(reqDTO), JsonUtils.toJsonString(result));
|
||||
throw exception(TRANSPORT_TASK_ACS_OPERATE_NOT_ALLOW, message);
|
||||
@@ -120,27 +129,11 @@ public class TransportTaskOperateCheckManager {
|
||||
/**
|
||||
* 获取 ACS 是否允许操作
|
||||
*/
|
||||
private Boolean getEnableOperate(AcsOperateCheckRespDTO result) {
|
||||
if (result.getEnableOperate() != null) {
|
||||
return result.getEnableOperate();
|
||||
private String getMessage(AcsCommonResult<AcsOperateCheckRespDTO> result, AcsOperateCheckRespDTO data) {
|
||||
if (StrUtil.isNotBlank(data.getMessage())) {
|
||||
return data.getMessage();
|
||||
}
|
||||
if (result.getCanOperate() != null) {
|
||||
return result.getCanOperate();
|
||||
}
|
||||
if (result.getData() != null) {
|
||||
return result.getData();
|
||||
}
|
||||
return result.getSuccess();
|
||||
}
|
||||
|
||||
/**
|
||||
* 获取 ACS 拒绝原因
|
||||
*/
|
||||
private String getMessage(AcsOperateCheckRespDTO result) {
|
||||
if (StrUtil.isNotBlank(result.getMessage())) {
|
||||
return result.getMessage();
|
||||
}
|
||||
return StrUtil.blankToDefault(result.getMsg(), "ACS 不允许执行该操作");
|
||||
return StrUtil.blankToDefault(result.getMessage(), "ACS 不允许执行该操作");
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -1,22 +1,30 @@
|
||||
package cn.code.nl.module.task.manage;
|
||||
|
||||
import cn.code.nl.framework.common.exception.ServiceException;
|
||||
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.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.TaskCallbackResultReqDTO;
|
||||
import cn.code.nl.module.task.enums.CallbackStatusEnum;
|
||||
import cn.code.nl.module.task.enums.FinishedTypeEnum;
|
||||
import cn.code.nl.module.task.enums.TaskCallbackTypeEnum;
|
||||
import cn.code.nl.module.task.enums.TaskEventTypeEnum;
|
||||
import cn.code.nl.module.task.enums.TaskOperationTypeEnum;
|
||||
import cn.code.nl.module.task.enums.TaskOwnerServiceEnum;
|
||||
import cn.code.nl.module.task.enums.TransportTaskStatusEnum;
|
||||
import cn.code.nl.module.task.message.TaskEventMessage;
|
||||
import cn.code.nl.module.task.mq.producer.TaskEventProducer;
|
||||
import cn.code.nl.module.task.service.transporttask.TransportTaskCallbackService;
|
||||
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.beans.factory.annotation.Value;
|
||||
import org.springframework.stereotype.Component;
|
||||
import org.springframework.transaction.support.TransactionSynchronization;
|
||||
import org.springframework.transaction.support.TransactionSynchronizationManager;
|
||||
@@ -51,9 +59,21 @@ public class TransportTaskOperationManager {
|
||||
@Resource
|
||||
private TransportTaskIssueManager transportTaskIssueManager;
|
||||
|
||||
@Resource
|
||||
private TaskCommonApiFactory taskCommonApiFactory;
|
||||
|
||||
@Resource
|
||||
private TransportTaskCallbackService transportTaskCallbackService;
|
||||
|
||||
@Resource
|
||||
private RedissonClient redissonClient;
|
||||
|
||||
@Value("${nl.task.callback-type}")
|
||||
private String callbackType;
|
||||
|
||||
/** 任务业务回调方式 */
|
||||
private TaskCallbackTypeEnum callbackTypeEnum;
|
||||
|
||||
/** 操作类型路由表 */
|
||||
private final Map<TaskOperationTypeEnum, BiConsumer<TransportTaskDO, AcsFeedbackReqDTO>> operationHandlers =
|
||||
new EnumMap<>(TaskOperationTypeEnum.class);
|
||||
@@ -63,6 +83,10 @@ public class TransportTaskOperationManager {
|
||||
*/
|
||||
@PostConstruct
|
||||
public void initOperationHandlers() {
|
||||
callbackTypeEnum = TaskCallbackTypeEnum.getByCode(callbackType);
|
||||
if (callbackTypeEnum == null) {
|
||||
throw new IllegalArgumentException("任务业务回调方式配置错误:" + callbackType);
|
||||
}
|
||||
operationHandlers.put(TaskOperationTypeEnum.EXECUTING, this::handleExecuting);
|
||||
operationHandlers.put(TaskOperationTypeEnum.ISSUE, (task, reqDTO) -> handleIssue(task));
|
||||
operationHandlers.put(TaskOperationTypeEnum.FINISHED, this::handleFinished);
|
||||
@@ -91,7 +115,7 @@ public class TransportTaskOperationManager {
|
||||
throw exception(TRANSPORT_TASK_OPERATION_FAILED, ex.getMessage());
|
||||
} finally {
|
||||
if (lock.isHeldByCurrentThread()) {
|
||||
lock.unlock();
|
||||
unlockAfterTransaction(lock);
|
||||
}
|
||||
}
|
||||
} else {
|
||||
@@ -147,19 +171,7 @@ public class TransportTaskOperationManager {
|
||||
task.setCallbackStatus(CallbackStatusEnum.PENDING.getCode());
|
||||
transportTaskMapper.updateById(task);
|
||||
|
||||
// MQ 在事务提交后发送,避免 Consumer 读到未提交的数据
|
||||
Map<String, Object> payload = reqDTO.getPayload();
|
||||
if (TransactionSynchronizationManager.isSynchronizationActive()) {
|
||||
TransactionSynchronizationManager.registerSynchronization(
|
||||
new TransactionSynchronization() {
|
||||
@Override
|
||||
public void afterCommit() {
|
||||
publishEvent(task, TaskEventTypeEnum.TASK_FINISHED, payload);
|
||||
}
|
||||
});
|
||||
} else {
|
||||
publishEvent(task, TaskEventTypeEnum.TASK_FINISHED, payload);
|
||||
}
|
||||
executeBusinessCallback(task, TaskEventTypeEnum.TASK_FINISHED, reqDTO.getPayload());
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -182,19 +194,7 @@ public class TransportTaskOperationManager {
|
||||
task.setCallbackStatus(CallbackStatusEnum.PENDING.getCode());
|
||||
transportTaskMapper.updateById(task);
|
||||
|
||||
// MQ 在事务提交后发送,避免 Consumer 读到未提交的数据
|
||||
Map<String, Object> payload = reqDTO.getPayload();
|
||||
if (TransactionSynchronizationManager.isSynchronizationActive()) {
|
||||
TransactionSynchronizationManager.registerSynchronization(
|
||||
new TransactionSynchronization() {
|
||||
@Override
|
||||
public void afterCommit() {
|
||||
publishEvent(task, TaskEventTypeEnum.TASK_CANCELLED, payload);
|
||||
}
|
||||
});
|
||||
} else {
|
||||
publishEvent(task, TaskEventTypeEnum.TASK_CANCELLED, payload);
|
||||
}
|
||||
executeBusinessCallback(task, TaskEventTypeEnum.TASK_CANCELLED, reqDTO.getPayload());
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -207,6 +207,104 @@ public class TransportTaskOperationManager {
|
||||
log.info("任务强制完成, taskId={}", task.getTaskId());
|
||||
}
|
||||
|
||||
/**
|
||||
* 在事务完成后释放任务状态锁
|
||||
*/
|
||||
private void unlockAfterTransaction(RLock lock) {
|
||||
if (!TransactionSynchronizationManager.isSynchronizationActive()) {
|
||||
lock.unlock();
|
||||
return;
|
||||
}
|
||||
TransactionSynchronizationManager.registerSynchronization(new TransactionSynchronization() {
|
||||
@Override
|
||||
public void afterCompletion(int status) {
|
||||
if (lock.isHeldByCurrentThread()) {
|
||||
lock.unlock();
|
||||
}
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
* 根据配置执行任务业务回调
|
||||
*/
|
||||
private void executeBusinessCallback(TransportTaskDO task, TaskEventTypeEnum eventType,
|
||||
Map<String, Object> payload) {
|
||||
TaskOwnerServiceEnum ownerService = TaskOwnerServiceEnum.getByCode(task.getOwnerService());
|
||||
if (ownerService == null) {
|
||||
throw new ServiceException(500, "任务业务归属服务配置错误:" + task.getOwnerService());
|
||||
}
|
||||
if (callbackTypeEnum == TaskCallbackTypeEnum.MQ) {
|
||||
publishEventAfterCommit(task, eventType, payload);
|
||||
return;
|
||||
}
|
||||
executeFeignCallback(task, ownerService, eventType, payload);
|
||||
}
|
||||
|
||||
/**
|
||||
* 在事务提交后发布任务事件
|
||||
*/
|
||||
private void publishEventAfterCommit(TransportTaskDO task, TaskEventTypeEnum eventType,
|
||||
Map<String, Object> payload) {
|
||||
if (TransactionSynchronizationManager.isSynchronizationActive()) {
|
||||
TransactionSynchronizationManager.registerSynchronization(
|
||||
new TransactionSynchronization() {
|
||||
@Override
|
||||
public void afterCommit() {
|
||||
publishEvent(task, eventType, payload);
|
||||
}
|
||||
});
|
||||
return;
|
||||
}
|
||||
publishEvent(task, eventType, payload);
|
||||
}
|
||||
|
||||
/**
|
||||
* 通过 Feign 同步执行任务业务并更新最终状态
|
||||
*/
|
||||
private void executeFeignCallback(TransportTaskDO task, TaskOwnerServiceEnum ownerService,
|
||||
TaskEventTypeEnum eventType, Map<String, Object> payload) {
|
||||
TaskCommonApi taskCommonApi = taskCommonApiFactory.getByServerName(ownerService.name());
|
||||
TaskStatusCallApiReqDTO callApiReqDTO = buildCallApiReqDTO(task, eventType, payload);
|
||||
if (eventType == TaskEventTypeEnum.TASK_FINISHED) {
|
||||
taskCommonApi.doHandleFinished(callApiReqDTO).getCheckedData();
|
||||
} else {
|
||||
taskCommonApi.doHandleCancelled(callApiReqDTO).getCheckedData();
|
||||
}
|
||||
transportTaskCallbackService.receiveCallbackResult(buildSuccessCallbackReq(task, eventType));
|
||||
}
|
||||
|
||||
/**
|
||||
* 构建 Feign 任务业务调用参数
|
||||
*/
|
||||
private TaskStatusCallApiReqDTO buildCallApiReqDTO(TransportTaskDO task, TaskEventTypeEnum eventType,
|
||||
Map<String, Object> payload) {
|
||||
TaskStatusCallApiReqDTO reqDTO = new TaskStatusCallApiReqDTO();
|
||||
reqDTO.setTaskId(task.getTaskId());
|
||||
reqDTO.setEventId(TaskEventMessage.buildEventId(task.getTaskId(), eventType.getCode()));
|
||||
reqDTO.setTaskCode(task.getTaskCode());
|
||||
reqDTO.setStatus(eventType.getCode());
|
||||
reqDTO.setOwnerService(task.getOwnerService());
|
||||
reqDTO.setBizType(task.getBizType());
|
||||
reqDTO.setBizId(task.getBizId());
|
||||
reqDTO.setHandleCode(task.getHandleCode());
|
||||
reqDTO.setPayload(payload);
|
||||
return reqDTO;
|
||||
}
|
||||
|
||||
/**
|
||||
* 构建业务处理成功回调参数
|
||||
*/
|
||||
private TaskCallbackResultReqDTO buildSuccessCallbackReq(TransportTaskDO task,
|
||||
TaskEventTypeEnum eventType) {
|
||||
TaskCallbackResultReqDTO reqDTO = new TaskCallbackResultReqDTO();
|
||||
reqDTO.setTaskId(task.getTaskId());
|
||||
reqDTO.setEventId(TaskEventMessage.buildEventId(task.getTaskId(), eventType.getCode()));
|
||||
reqDTO.setResult(CallbackStatusEnum.SUCCESS.getCode());
|
||||
reqDTO.setMessage("任务业务处理成功");
|
||||
return reqDTO;
|
||||
}
|
||||
|
||||
/**
|
||||
* 发布 MQ 事件
|
||||
*/
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
package cn.code.nl.module.task.mq.producer;
|
||||
|
||||
import cn.code.nl.module.task.enums.TaskOwnerServiceEnum;
|
||||
import cn.code.nl.module.task.message.TaskEventMessage;
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
import org.apache.rocketmq.spring.core.RocketMQTemplate;
|
||||
@@ -7,7 +8,6 @@ import org.springframework.stereotype.Component;
|
||||
|
||||
import jakarta.annotation.Resource;
|
||||
import org.springframework.beans.factory.annotation.Value;
|
||||
import java.util.UUID;
|
||||
|
||||
/**
|
||||
* 任务事件 MQ 生产者
|
||||
@@ -31,9 +31,12 @@ public class TaskEventProducer {
|
||||
* @param message 事件消息
|
||||
*/
|
||||
public void publishEvent(TaskEventMessage message) {
|
||||
message.setEventId(message.getTaskId().toString());
|
||||
|
||||
String destination = message.getOwnerService() + "_" + topic;
|
||||
message.setEventId(TaskEventMessage.buildEventId(message.getTaskId(), message.getEventType()));
|
||||
TaskOwnerServiceEnum ownerService = TaskOwnerServiceEnum.getByCode(message.getOwnerService());
|
||||
if (ownerService == null) {
|
||||
throw new IllegalArgumentException("任务业务归属服务配置错误:" + message.getOwnerService());
|
||||
}
|
||||
String destination = ownerService.buildTopic(topic);
|
||||
try {
|
||||
rocketMQTemplate.syncSend(destination, message);
|
||||
log.info("MQ 发送成功, destination={}, taskId={}, eventType={}",
|
||||
|
||||
@@ -1,24 +1,26 @@
|
||||
package cn.code.nl.module.task.service.transporttask;
|
||||
|
||||
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.TaskCallbackResultReqDTO;
|
||||
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.message.TaskEventMessage;
|
||||
import com.mzt.logapi.context.LogRecordContext;
|
||||
import com.mzt.logapi.starter.annotation.LogRecord;
|
||||
import jakarta.annotation.Resource;
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
import org.springframework.stereotype.Service;
|
||||
import org.springframework.validation.annotation.Validated;
|
||||
|
||||
import jakarta.annotation.Resource;
|
||||
|
||||
import static cn.code.nl.module.system.enums.LogRecordConstants.*;
|
||||
import static cn.code.nl.module.system.enums.LogRecordConstants.TASK_INFO;
|
||||
import static cn.code.nl.module.system.enums.LogRecordConstants.TASK_INFO_CALL_BY_RPC;
|
||||
import static cn.code.nl.module.system.enums.LogRecordConstants.TASK_INFO_OPERATE_TYPE;
|
||||
|
||||
/**
|
||||
* 业务回调结果处理服务实现
|
||||
*
|
||||
* @author 诺力管理员
|
||||
*/
|
||||
@Slf4j
|
||||
@Service
|
||||
@@ -34,35 +36,37 @@ public class TransportTaskCallbackServiceImpl implements TransportTaskCallbackSe
|
||||
bizNo = "{{#reqDTO.taskId}}",
|
||||
success = TASK_INFO_CALL_BY_RPC)
|
||||
public void receiveCallbackResult(TaskCallbackResultReqDTO reqDTO) {
|
||||
// 1. 查询任务
|
||||
TransportTaskDO task = transportTaskMapper.selectById(reqDTO.getTaskId());
|
||||
if (task == null) {
|
||||
log.warn("回调 taskId 不存在, taskId={}", reqDTO.getTaskId());
|
||||
log.warn("回调任务不存在,taskId={}", reqDTO.getTaskId());
|
||||
return;
|
||||
}
|
||||
|
||||
// 2. 幂等:已是终态则直接返回
|
||||
String currentStatus = task.getTaskStatus();
|
||||
String finalFinished = statusCode(TransportTaskStatusEnum.FINISHED);
|
||||
String finalCancelled = statusCode(TransportTaskStatusEnum.CANCELLED);
|
||||
String finalFinished = TransportTaskStatusEnum.FINISHED.getCode();
|
||||
String finalCancelled = TransportTaskStatusEnum.CANCELLED.getCode();
|
||||
if (finalFinished.equals(currentStatus) || finalCancelled.equals(currentStatus)) {
|
||||
log.info("任务已处终态, taskId={}, status={}, 忽略回调",
|
||||
log.info("任务已处于终态,忽略重复回调,taskId={}, status={}",
|
||||
reqDTO.getTaskId(), currentStatus);
|
||||
return;
|
||||
}
|
||||
|
||||
// 3. 根据结果处理
|
||||
boolean success = "SUCCESS".equalsIgnoreCase(reqDTO.getResult());
|
||||
TaskEventTypeEnum eventType = getPendingEventType(currentStatus);
|
||||
if (eventType == null) {
|
||||
throw new ServiceException(500, "任务当前状态不允许处理业务回调:" + currentStatus);
|
||||
}
|
||||
String expectedEventId = TaskEventMessage.buildEventId(reqDTO.getTaskId(), eventType.getCode());
|
||||
if (!expectedEventId.equals(reqDTO.getEventId())) {
|
||||
throw new ServiceException(500, "任务业务回调事件标识不匹配:" + reqDTO.getEventId());
|
||||
}
|
||||
|
||||
boolean success = CallbackStatusEnum.SUCCESS.getCode().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.setTaskStatus(eventType == TaskEventTypeEnum.TASK_FINISHED
|
||||
? finalFinished : finalCancelled);
|
||||
task.setCallbackStatus(CallbackStatusEnum.SUCCESS.getCode());
|
||||
task.setCallbackErrorMsg(null);
|
||||
} else {
|
||||
// 业务处理失败:记录失败信息,增加重试次数
|
||||
task.setCallbackStatus(CallbackStatusEnum.FAILED.getCode());
|
||||
task.setCallbackErrorMsg(reqDTO.getMessage());
|
||||
task.setCallbackRetryCount(task.getCallbackRetryCount() != null
|
||||
@@ -70,15 +74,22 @@ public class TransportTaskCallbackServiceImpl implements TransportTaskCallbackSe
|
||||
}
|
||||
|
||||
transportTaskMapper.updateById(task);
|
||||
|
||||
LogRecordContext.putVariable("ownerService", task.getOwnerService());
|
||||
LogRecordContext.putVariable("taskId", task.getTaskId());
|
||||
|
||||
log.info("业务回调结果已处理, taskId={}, result={}, newStatus={}",
|
||||
log.info("业务回调结果已处理,taskId={}, result={}, newStatus={}",
|
||||
reqDTO.getTaskId(), reqDTO.getResult(), task.getTaskStatus());
|
||||
}
|
||||
|
||||
private String statusCode(TransportTaskStatusEnum statusEnum) {
|
||||
return String.format("%03d", statusEnum.getCode());
|
||||
/**
|
||||
* 根据任务待处理状态获取事件类型
|
||||
*/
|
||||
private TaskEventTypeEnum getPendingEventType(String currentStatus) {
|
||||
if (TransportTaskStatusEnum.FINISHED_CALLBACK_PENDING.getCode().equals(currentStatus)) {
|
||||
return TaskEventTypeEnum.TASK_FINISHED;
|
||||
}
|
||||
if (TransportTaskStatusEnum.CANCEL_CALLBACK_PENDING.getCode().equals(currentStatus)) {
|
||||
return TaskEventTypeEnum.TASK_CANCELLED;
|
||||
}
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -7,12 +7,12 @@ spring:
|
||||
username: nacos
|
||||
password: nacos
|
||||
discovery: # 【注册中心】配置项
|
||||
namespace: dev # 命名空间。这里使用 dev 开发环境
|
||||
namespace: 91c8ac41-7fb0-423b-947e-2caad4c87d41 # 命名空间。这里使用 dev 开发环境
|
||||
group: DEFAULT_GROUP # 使用的 Nacos 配置分组,默认为 DEFAULT_GROUP
|
||||
metadata:
|
||||
version: 1.0.0 # 服务实例的版本号,可用于灰度发布
|
||||
config: # 【配置中心】配置项
|
||||
namespace: dev # 命名空间。这里使用 dev 开发环境
|
||||
namespace: 91c8ac41-7fb0-423b-947e-2caad4c87d41 # 命名空间。这里使用 dev 开发环境
|
||||
group: DEFAULT_GROUP # 使用的 Nacos 配置分组,默认为 DEFAULT_GROUP
|
||||
|
||||
--- #################### 数据库相关配置 ####################
|
||||
@@ -59,12 +59,12 @@ spring:
|
||||
master:
|
||||
url: jdbc:mysql://192.168.81.193:3306/huachuang_lms_dev?useSSL=false&serverTimezone=Asia/Shanghai&allowPublicKeyRetrieval=true&nullCatalogMeansCurrent=true&rewriteBatchedStatements=true # MySQL Connector/J 8.X 连接的示例
|
||||
username: root
|
||||
password: root
|
||||
password: root123
|
||||
slave: # 模拟从库,可根据自己需要修改 # 模拟从库,可根据自己需要修改
|
||||
lazy: true # 开启懒加载,保证启动速度
|
||||
url: jdbc:mysql://192.168.81.193:3306/huachuang_lms_dev?useSSL=false&serverTimezone=Asia/Shanghai&allowPublicKeyRetrieval=true&nullCatalogMeansCurrent=true&rewriteBatchedStatements=true # MySQL Connector/J 8.X 连接的示例
|
||||
username: root
|
||||
password: root
|
||||
password: root123
|
||||
|
||||
# Redis 配置。Redisson 默认的配置足够使用,一般不需要进行调优
|
||||
data:
|
||||
@@ -83,7 +83,7 @@ rocketmq:
|
||||
group: task_producer_dev_group # 生产者分组
|
||||
consumer:
|
||||
task:
|
||||
topic: task-status-change-dev-topic
|
||||
topic: task_status_change_dev_topic
|
||||
|
||||
spring:
|
||||
# RabbitMQ 配置项,对应 RabbitProperties 配置类
|
||||
@@ -100,11 +100,11 @@ spring:
|
||||
xxl:
|
||||
job:
|
||||
admin:
|
||||
addresses: http://192.168.10.2:8080/xxl-job-admin # 调度中心部署跟地址
|
||||
accessToken: 123456 # 执行器通讯TOKEN
|
||||
addresses: http://192.168.81.193:8080/xxl-job-admin # 调度中心部署跟地址
|
||||
accessToken: default_token # 执行器通讯TOKEN
|
||||
executor:
|
||||
ip: 192.168.10.2
|
||||
port: 9994
|
||||
ip: localhost
|
||||
port: 8995
|
||||
|
||||
--- #################### 服务保障相关配置 ####################
|
||||
|
||||
@@ -145,4 +145,4 @@ logging:
|
||||
|
||||
# 芋道配置项,设置当前项目所有自定义的配置
|
||||
nl:
|
||||
demo: false # 开启演示模式
|
||||
demo: false # 开启演示模式
|
||||
|
||||
@@ -113,6 +113,8 @@ nl:
|
||||
info:
|
||||
version: 1.0.0
|
||||
base-package: cn.code.nl.module.task
|
||||
task:
|
||||
callback-type: mq # 任务完成/取消业务处理方式:mq 异步、feign 同步
|
||||
web:
|
||||
admin-ui:
|
||||
url: http://dashboard.nl.iocoder.cn # Admin 管理后台 UI 的地址
|
||||
|
||||
@@ -5,17 +5,22 @@ import cn.code.nl.framework.execute.core.AbstractTask;
|
||||
import cn.code.nl.framework.execute.core.TaskFactory;
|
||||
import cn.code.nl.framework.execute.core.dto.TaskExecuteDTO;
|
||||
import cn.code.nl.module.task.api.TransportTaskApi;
|
||||
import cn.code.nl.module.task.dto.TaskCallbackResultReqDTO;
|
||||
import cn.code.nl.module.task.dto.TaskInfoDTO;
|
||||
import cn.code.nl.module.task.enums.CallbackStatusEnum;
|
||||
import cn.code.nl.module.task.enums.TransportTaskStatusEnum;
|
||||
import cn.code.nl.module.task.message.TaskEventMessage;
|
||||
import jakarta.annotation.Resource;
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
|
||||
import org.apache.rocketmq.spring.core.RocketMQListener;
|
||||
import org.redisson.api.RBucket;
|
||||
import org.redisson.api.RLock;
|
||||
import org.redisson.api.RedissonClient;
|
||||
import org.springframework.stereotype.Component;
|
||||
|
||||
import java.time.Duration;
|
||||
|
||||
import static cn.code.nl.module.task.message.LockKeyConstants.TASK_STATUS_CHANGE_LOCK_KEY;
|
||||
|
||||
/**
|
||||
@@ -32,6 +37,12 @@ import static cn.code.nl.module.task.message.LockKeyConstants.TASK_STATUS_CHANGE
|
||||
)
|
||||
public class WmsTaskStatusChangeConsumer implements RocketMQListener<TaskEventMessage> {
|
||||
|
||||
/** 任务业务处理成功标识前缀 */
|
||||
private static final String TASK_BUSINESS_HANDLED_KEY = "task:business:handled:";
|
||||
|
||||
/** 任务业务处理成功标识有效期 */
|
||||
private static final Duration TASK_BUSINESS_HANDLED_TIMEOUT = Duration.ofDays(7);
|
||||
|
||||
@Resource
|
||||
private TaskFactory taskFactory;
|
||||
|
||||
@@ -56,7 +67,19 @@ public class WmsTaskStatusChangeConsumer implements RocketMQListener<TaskEventMe
|
||||
if (!isTaskCallbackPending(message)) {
|
||||
return;
|
||||
}
|
||||
executeTaskHandler(message, handleCode, eventType);
|
||||
RBucket<Boolean> handledBucket = redissonClient.getBucket(
|
||||
TASK_BUSINESS_HANDLED_KEY + message.getEventId());
|
||||
if (!Boolean.TRUE.equals(handledBucket.get())) {
|
||||
try {
|
||||
executeTaskHandler(message, handleCode, eventType);
|
||||
handledBucket.set(true, TASK_BUSINESS_HANDLED_TIMEOUT);
|
||||
} catch (RuntimeException exception) {
|
||||
callbackFailedResult(message, exception);
|
||||
throw exception;
|
||||
}
|
||||
}
|
||||
callbackTaskResult(message, CallbackStatusEnum.SUCCESS.getCode(), "任务业务处理成功");
|
||||
handledBucket.delete();
|
||||
} finally {
|
||||
if (lock.isHeldByCurrentThread()) {
|
||||
lock.unlock();
|
||||
@@ -99,12 +122,13 @@ public class WmsTaskStatusChangeConsumer implements RocketMQListener<TaskEventMe
|
||||
private void executeTaskHandler(TaskEventMessage message, String handleCode, String eventType) {
|
||||
AbstractTask task = taskFactory.getTask(handleCode);
|
||||
if (task == null) {
|
||||
log.warn("未找到对应任务处理器, handleCode={}", handleCode);
|
||||
return;
|
||||
throw new IllegalStateException("未找到对应任务处理器:" + handleCode);
|
||||
}
|
||||
|
||||
TaskExecuteDTO dto = new TaskExecuteDTO();
|
||||
dto.setTaskId(message.getTaskId());
|
||||
dto.setEventId(message.getEventId());
|
||||
dto.setEventType(message.getEventType());
|
||||
dto.setPayload(message.getPayload());
|
||||
|
||||
if (TaskEventMessage.EVENT_TYPE_FINISHED.equals(eventType)) {
|
||||
@@ -115,4 +139,31 @@ public class WmsTaskStatusChangeConsumer implements RocketMQListener<TaskEventMe
|
||||
log.warn("未知事件类型, eventType={}, handleCode={}", eventType, handleCode);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 回调 Task 服务记录业务处理结果
|
||||
*/
|
||||
private void callbackTaskResult(TaskEventMessage message, String result, String resultMessage) {
|
||||
TaskCallbackResultReqDTO reqDTO = new TaskCallbackResultReqDTO();
|
||||
reqDTO.setTaskId(message.getTaskId());
|
||||
reqDTO.setEventId(message.getEventId());
|
||||
reqDTO.setResult(result);
|
||||
reqDTO.setMessage(resultMessage);
|
||||
CommonResult<Boolean> callbackResult = transportTaskApi.receiveCallbackResult(reqDTO);
|
||||
if (!callbackResult.isSuccess()) {
|
||||
log.error("任务业务结果回调失败, reqDTO={}, callbackResult={}", reqDTO, callbackResult);
|
||||
callbackResult.checkError();
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 尝试回调 Task 服务记录业务处理失败
|
||||
*/
|
||||
private void callbackFailedResult(TaskEventMessage message, RuntimeException exception) {
|
||||
try {
|
||||
callbackTaskResult(message, CallbackStatusEnum.FAILED.getCode(), exception.getMessage());
|
||||
} catch (RuntimeException callbackException) {
|
||||
log.error("任务业务失败结果回调异常, message={}", message, callbackException);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -194,6 +194,12 @@
|
||||
</exclusions>
|
||||
</dependency>
|
||||
|
||||
<dependency>
|
||||
<groupId>org.springframework.boot</groupId>
|
||||
<artifactId>spring-boot-starter-test</artifactId>
|
||||
<scope>test</scope>
|
||||
</dependency>
|
||||
|
||||
</dependencies>
|
||||
|
||||
<build>
|
||||
|
||||
@@ -47,14 +47,14 @@ spring:
|
||||
primary: master
|
||||
datasource:
|
||||
master:
|
||||
url: jdbc:mysql://192.168.10.41:3306/huachuang_lms_dev?useSSL=false&serverTimezone=Asia/Shanghai&allowPublicKeyRetrieval=true&nullCatalogMeansCurrent=true&rewriteBatchedStatements=true # MySQL Connector/J 8.X 连接的示例
|
||||
url: jdbc:mysql://192.168.81.193:3306/huachuang_lms_dev?useSSL=false&serverTimezone=Asia/Shanghai&allowPublicKeyRetrieval=true&nullCatalogMeansCurrent=true&rewriteBatchedStatements=true # MySQL Connector/J 8.X 连接的示例
|
||||
username: root
|
||||
password: root
|
||||
password: root123
|
||||
slave: # 模拟从库,可根据自己需要修改 # 模拟从库,可根据自己需要修改
|
||||
lazy: true # 开启懒加载,保证启动速度
|
||||
url: jdbc:mysql://192.168.10.41:3306/huachuang_lms_dev?useSSL=false&serverTimezone=Asia/Shanghai&allowPublicKeyRetrieval=true&nullCatalogMeansCurrent=true&rewriteBatchedStatements=true # MySQL Connector/J 8.X 连接的示例
|
||||
url: jdbc:mysql://192.168.81.193:3306/huachuang_lms_dev?useSSL=false&serverTimezone=Asia/Shanghai&allowPublicKeyRetrieval=true&nullCatalogMeansCurrent=true&rewriteBatchedStatements=true # MySQL Connector/J 8.X 连接的示例
|
||||
username: root
|
||||
password: root
|
||||
password: root123
|
||||
|
||||
# Redis 配置。Redisson 默认的配置足够使用,一般不需要进行调优
|
||||
data:
|
||||
|
||||
@@ -146,13 +146,6 @@ aj:
|
||||
req-verify-minute-limit: 60 # verify 接口一分钟内请求数限制
|
||||
|
||||
--- #################### 消息队列相关 ####################
|
||||
|
||||
# rocketmq 配置项,对应 RocketMQProperties 配置类
|
||||
rocketmq:
|
||||
# Producer 配置项
|
||||
producer:
|
||||
group: ${spring.application.name}_PRODUCER # 生产者分组
|
||||
|
||||
spring:
|
||||
# Kafka 配置项,对应 KafkaProperties 配置类
|
||||
kafka:
|
||||
@@ -289,6 +282,8 @@ nl:
|
||||
info:
|
||||
version: 1.0.0
|
||||
base-package: cn.code.nl
|
||||
task:
|
||||
callback-type: mq # 任务完成/取消业务处理方式:mq 异步、feign 同步
|
||||
web:
|
||||
admin-ui:
|
||||
url: http://dashboard.nl.iocoder.cn # Admin 管理后台 UI 的地址
|
||||
@@ -379,4 +374,4 @@ nl:
|
||||
message-bus:
|
||||
type: redis # 消息总线的类型
|
||||
|
||||
debug: false
|
||||
debug: false
|
||||
|
||||
@@ -0,0 +1,76 @@
|
||||
package cn.code.nl.server.acs;
|
||||
|
||||
import cn.code.nl.framework.common.pojo.AcsCommonResult;
|
||||
import cn.code.nl.framework.common.pojo.CommonResult;
|
||||
import cn.code.nl.framework.common.util.spring.SpringUtils;
|
||||
import cn.code.nl.framework.execute.util.http.AcsUtil;
|
||||
import cn.code.nl.module.infra.api.config.ConfigApi;
|
||||
import cn.code.nl.module.task.job.dto.AcsTaskDTO;
|
||||
import lombok.Data;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.springframework.context.annotation.Bean;
|
||||
import org.springframework.context.annotation.Configuration;
|
||||
import org.springframework.test.context.ContextConfiguration;
|
||||
import org.springframework.test.context.junit.jupiter.SpringJUnitConfig;
|
||||
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
|
||||
import static cn.code.nl.module.task.enums.AcsApiConstants.ACS_TASK_API;
|
||||
import static org.junit.jupiter.api.Assertions.assertNotNull;
|
||||
|
||||
@SpringJUnitConfig
|
||||
@ContextConfiguration(classes = AcsUtilLocalSpringTest.TestConfiguration.class)
|
||||
class AcsUtilLocalSpringTest {
|
||||
|
||||
private static final String SERVER_ADDRESS = "http://127.0.0.1:8011";
|
||||
|
||||
@Test
|
||||
void shouldPostIssueTaskToLocalAcs() {
|
||||
AcsCommonResult<AcsIssueTaskResp> result = AcsUtil.post(SERVER_ADDRESS, ACS_TASK_API,
|
||||
List.of(buildAcsTask()), AcsIssueTaskResp.class);
|
||||
|
||||
assertNotNull(result);
|
||||
}
|
||||
|
||||
private AcsTaskDTO buildAcsTask() {
|
||||
AcsTaskDTO task = new AcsTaskDTO();
|
||||
// task.setTaskId(10001L);
|
||||
task.setTaskCode("LOCAL-ACS-TEST-10001");
|
||||
task.setStartDeviceCode("START-TEST-01");
|
||||
task.setNextDeviceCode("END-TEST-01");
|
||||
task.setPriority("1");
|
||||
task.setVehicleCode("BOX-TEST-01");
|
||||
task.setTaskType("TRANSFER");
|
||||
task.setAgvSystemType("ACS");
|
||||
task.setProductArea("TEST");
|
||||
task.setRemark("local acs issue task test");
|
||||
task.setPayload(Map.of("source", "AcsUtilLocalSpringTest"));
|
||||
return task;
|
||||
}
|
||||
|
||||
@Data
|
||||
static class AcsIssueTaskResp {
|
||||
private List<FailedTask> failedTasks;
|
||||
}
|
||||
|
||||
@Data
|
||||
static class FailedTask {
|
||||
private Long taskId;
|
||||
private String errorMessage;
|
||||
}
|
||||
|
||||
@Configuration(proxyBeanMethods = false)
|
||||
static class TestConfiguration {
|
||||
|
||||
@Bean
|
||||
SpringUtils springUtils() {
|
||||
return new SpringUtils();
|
||||
}
|
||||
|
||||
@Bean
|
||||
ConfigApi configApi() {
|
||||
return key -> CommonResult.success("1");
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user