package com.zy.core.plugin.store;
|
|
import com.alibaba.fastjson.JSON;
|
import com.alibaba.fastjson.JSONObject;
|
import com.baomidou.mybatisplus.core.conditions.query.QueryWrapper;
|
import com.core.common.Cools;
|
import com.zy.asrs.domain.param.CreateInTaskParam;
|
import com.zy.asrs.entity.BasDevp;
|
import com.zy.asrs.entity.WrkMast;
|
import com.zy.asrs.service.WrkMastService;
|
import com.zy.common.model.StartupDto;
|
import com.zy.common.service.CommonService;
|
import com.zy.common.utils.RedisUtil;
|
import com.zy.core.News;
|
import com.zy.core.cache.SlaveConnection;
|
import com.zy.core.dispatch.StationCommandDispatchResult;
|
import com.zy.core.dispatch.StationCommandDispatcher;
|
import com.zy.core.enums.RedisKeyType;
|
import com.zy.core.enums.SlaveType;
|
import com.zy.core.enums.StationCommandType;
|
import com.zy.core.enums.WrkIoType;
|
import com.zy.core.model.StationObjModel;
|
import com.zy.core.model.command.StationCommand;
|
import com.zy.core.model.protocol.StationProtocol;
|
import com.zy.core.task.MainProcessLane;
|
import com.zy.core.task.MainProcessTaskSubmitter;
|
import com.zy.core.thread.StationThread;
|
import com.zy.core.utils.StationOperateProcessUtils;
|
import com.zy.core.utils.WmsOperateUtils;
|
import org.springframework.beans.factory.annotation.Autowired;
|
import org.springframework.stereotype.Service;
|
|
import java.util.HashMap;
|
import java.util.Map;
|
|
@Service
|
public class StoreInTaskGenerationService {
|
private static final int APPLY_IN_TASK_TIMEOUT_SECONDS = 5;
|
private static final int APPLY_FAIL_STATION_BACK_LOCK_SECONDS = 30;
|
|
@Autowired
|
private WrkMastService wrkMastService;
|
@Autowired
|
private StationOperateProcessUtils stationOperateProcessUtils;
|
@Autowired
|
private RedisUtil redisUtil;
|
@Autowired
|
private WmsOperateUtils wmsOperateUtils;
|
@Autowired
|
private CommonService commonService;
|
@Autowired
|
private MainProcessTaskSubmitter mainProcessTaskSubmitter;
|
@Autowired
|
private StationCommandDispatcher stationCommandDispatcher;
|
|
/**
|
* 保留当前按站点 lane 并发的能力,同时用一个简单计数避免并发生成把站点任务数顶穿上限。
|
*/
|
private int inFlightGenerateCount = 0;
|
|
public void generate(StoreInTaskPolicy policy, BasDevp basDevp, StationObjModel stationObjModel) {
|
try {
|
if (!policy.isEnabled()) {
|
return;
|
}
|
|
HashMap<String, String> systemConfigMap = getSystemConfigMap();
|
if (systemConfigMap == null) {
|
return;
|
}
|
if (!hasAvailableStationTaskCapacity(systemConfigMap)) {
|
return;
|
}
|
|
generateByStation(policy, basDevp, stationObjModel, systemConfigMap);
|
} catch (Exception e) {
|
Integer stationId = stationObjModel == null ? null : stationObjModel.getStationId();
|
News.error("生成入库任务异常,policy={},stationId={}", policy.getPolicyName(), stationId, e);
|
}
|
}
|
|
public void submitGenerateStoreTask(StoreInTaskPolicy policy,
|
BasDevp basDevp,
|
StationObjModel stationObjModel,
|
long minIntervalMs,
|
Runnable task) {
|
submitGenerateStoreTask(policy, basDevp, stationObjModel, MainProcessLane.GENERATE_STORE, minIntervalMs, task);
|
}
|
|
public void submitGenerateStoreTask(StoreInTaskPolicy policy,
|
BasDevp basDevp,
|
StationObjModel stationObjModel,
|
MainProcessLane lane,
|
long minIntervalMs,
|
Runnable task) {
|
Integer stationId = stationObjModel == null ? null : stationObjModel.getStationId();
|
mainProcessTaskSubmitter.submitKeyedSerialTask(
|
lane,
|
stationId,
|
"generateStoreWrkFile",
|
minIntervalMs,
|
task
|
);
|
}
|
|
private void generateByStation(StoreInTaskPolicy policy, BasDevp basDevp, StationObjModel stationObjModel,
|
HashMap<String, String> systemConfigMap) {
|
StoreInTaskContext context = buildContext(basDevp, stationObjModel);
|
if (context == null) {
|
return;
|
}
|
|
StationProtocol stationProtocol = context.getStationProtocol();
|
if (stationProtocol == null) {
|
return;
|
}
|
|
if (!stationProtocol.isAutoing()) {
|
return;
|
}
|
|
if (!stationProtocol.isLoading()) {
|
return;
|
}
|
|
if (!stationProtocol.isInEnable()) {
|
return;
|
}
|
|
if (stationProtocol.getTaskNo() == 0) {
|
return;
|
}
|
|
if (Cools.isEmpty(stationProtocol.getBarcode())) {
|
return;
|
}
|
|
if (stationProtocol.getError() > 0) {
|
return;
|
}
|
|
if (stationProtocol.isInBarcodeError()) {
|
return;
|
}
|
|
if (!stationProtocol.getIoMode().equals(1)) {
|
policy.setSystemWarning(context, "当前站点不处于入库模式");
|
return;
|
}
|
|
String barcode = context.getStationProtocol().getBarcode();
|
long count = wrkMastService.count(new QueryWrapper<WrkMast>().eq("barcode", barcode));
|
if (count > 0) {
|
Object tipsLimit = redisUtil.get(RedisKeyType.GENERATE_IN_TASK_SUCCESS_REPEAT_WARNING_TIPS_LIMIT.key + barcode);
|
if (tipsLimit == null) {
|
policy.setSystemWarning(context, "系统任务已存在");
|
}
|
return;
|
}
|
|
if (redisUtil.get(policy.getGenerateLockKey(context)) != null) {
|
return;
|
}
|
|
if (!tryReserveGenerateCapacity(systemConfigMap)) {
|
return;
|
}
|
|
InTaskApplyRequest request = policy.buildApplyRequest(context);
|
try {
|
policy.onRequestPermitGranted(context);
|
policy.setSystemWarning(context, "请求WMS中");
|
News.info("发起同步WMS入库请求,barcode={},stationId={},timeout={}s",
|
request.getBarcode(), request.getSourceStaNo(), APPLY_IN_TASK_TIMEOUT_SECONDS);
|
String response = wmsOperateUtils.applyInTask(request);
|
handleSyncApplyResponse(policy, context, request, response);
|
} finally {
|
releaseGenerateCapacity();
|
}
|
}
|
|
private StoreInTaskContext buildContext(BasDevp basDevp, StationObjModel stationObjModel) {
|
if (basDevp == null || stationObjModel == null || stationObjModel.getStationId() == null) {
|
return null;
|
}
|
|
StationThread stationThread = (StationThread) SlaveConnection.get(SlaveType.Devp, basDevp.getDevpNo());
|
if (stationThread == null) {
|
return null;
|
}
|
|
Integer stationId = stationObjModel.getStationId();
|
Map<Integer, StationProtocol> stationMap = stationThread.getStatusMap();
|
if (stationMap == null || !stationMap.containsKey(stationId)) {
|
return null;
|
}
|
|
StationProtocol stationProtocol = stationMap.get(stationId);
|
if (stationProtocol == null) {
|
return null;
|
}
|
|
return new StoreInTaskContext(basDevp, stationThread, stationObjModel, stationProtocol);
|
}
|
|
private void handleSyncApplyResponse(StoreInTaskPolicy policy, StoreInTaskContext context, InTaskApplyRequest request,
|
String response) {
|
if (Cools.isEmpty(response)) {
|
markApplyFailed(policy, context, request, null, "FAILED");
|
return;
|
}
|
|
try {
|
JSONObject jsonObject = JSON.parseObject(response);
|
if (jsonObject == null || !Integer.valueOf(200).equals(jsonObject.getInteger("code"))) {
|
markApplyFailed(policy, context, request, response, "WMS返回非200");
|
return;
|
}
|
|
StartupDto dto = jsonObject.getObject("data", StartupDto.class);
|
if (dto == null) {
|
markApplyFailed(policy, context, request, response, "WMS返回data为空");
|
return;
|
}
|
|
CreateInTaskParam taskParam = policy.buildCreateInTaskParam(context, dto);
|
WrkMast wrkMast = commonService.createInTask(taskParam);
|
policy.afterTaskCreated(context, wrkMast);
|
policy.clearSystemWarning(context);
|
redisUtil.set(RedisKeyType.GENERATE_IN_TASK_SUCCESS_REPEAT_WARNING_TIPS_LIMIT.key + wrkMast.getBarcode(), "lock", 30);
|
} catch (Exception e) {
|
News.error("处理WMS入库响应异常,barcode={},stationId={}", request.getBarcode(),
|
request.getSourceStaNo(), e);
|
markApplyFailed(policy, context, request, response, e.getMessage());
|
}
|
}
|
|
private void markApplyFailed(StoreInTaskPolicy policy, StoreInTaskContext context, InTaskApplyRequest request,
|
String response, String message) {
|
InTaskApplyResult result = new InTaskApplyResult();
|
result.setStatus(InTaskApplyStatus.RETRYABLE_FAIL);
|
result.setResponse(response);
|
result.setMessage(message);
|
|
News.error("WMS入库请求失败,barcode={},stationId={},response={},WCS响应={}",
|
request.getBarcode(), request.getSourceStaNo(), result.getResponse(), result.getMessage());
|
redisUtil.set(policy.getGenerateLockKey(context), "lock", policy.getRetryLockSeconds(context));
|
policy.onApplyFailed(context, result);
|
triggerStationBackOnApplyFailed(context, request, result);
|
}
|
|
/**
|
* WMS 申请入库失败后,货物仍停留在扫码站,此时补发退回到入库站的输送命令,避免货物长期滞留。
|
*/
|
private void triggerStationBackOnApplyFailed(StoreInTaskContext context,
|
InTaskApplyRequest request,
|
InTaskApplyResult result) {
|
if (context == null || context.getStationThread() == null || context.getStationObjModel() == null) {
|
return;
|
}
|
|
StationProtocol stationProtocol = context.getStationProtocol();
|
if (stationProtocol == null || stationProtocol.getStationId() == null) {
|
return;
|
}
|
|
StationObjModel backStation = context.getStationObjModel().getBackStation();
|
if (backStation == null || backStation.getStationId() == null) {
|
News.warn("WMS入库失败后无法退回入库站,退回站未配置。barcode={},stationId={}",
|
request == null ? null : request.getBarcode(), stationProtocol.getStationId());
|
return;
|
}
|
|
Integer currentStationId = stationProtocol.getStationId();
|
Integer backStationId = backStation.getStationId();
|
if (backStationId.equals(currentStationId)) {
|
return;
|
}
|
|
if (stationProtocol.getTaskNo() != null
|
&& stationProtocol.getTaskNo() > 0
|
&& backStationId.equals(stationProtocol.getTargetStaNo())) {
|
return;
|
}
|
|
String lockKey = RedisKeyType.GENERATE_STATION_BACK_LIMIT.key
|
+ currentStationId + "_" + stationProtocol.getTaskNo();
|
if (redisUtil.get(lockKey) != null) {
|
return;
|
}
|
|
Integer stationBackTaskNo = commonService.getWorkNo(WrkIoType.STATION_BACK.id);
|
StationCommand command = context.getStationThread().getCommand(
|
StationCommandType.MOVE,
|
stationBackTaskNo,
|
currentStationId,
|
backStationId,
|
0
|
);
|
if (command == null) {
|
News.warn("WMS入库失败后生成退回入库站命令失败。barcode={},stationId={},backStationId={},warning={}",
|
request == null ? null : request.getBarcode(),
|
currentStationId,
|
backStationId,
|
result == null ? null : result.getMessage());
|
return;
|
}
|
|
StationCommandDispatchResult dispatchResult = stationCommandDispatcher.dispatch(
|
context.getBasDevp().getDevpNo(),
|
command,
|
"store-in-task-generation",
|
"apply-failed-station-back"
|
);
|
if (!dispatchResult.isAccepted()) {
|
News.warn("WMS入库失败后退回入库站命令入队失败。barcode={},stationId={},backStationId={},reason={},warning={}",
|
request == null ? null : request.getBarcode(),
|
currentStationId,
|
backStationId,
|
dispatchResult.getReason(),
|
result == null ? null : result.getMessage());
|
return;
|
}
|
|
redisUtil.set(lockKey, "lock", APPLY_FAIL_STATION_BACK_LOCK_SECONDS);
|
News.warn("WMS入库失败,已触发货物退回入库站。barcode={},stationId={},backStationId={},warning={}",
|
request == null ? null : request.getBarcode(),
|
currentStationId,
|
backStationId,
|
result == null ? null : result.getMessage());
|
}
|
|
private HashMap<String, String> getSystemConfigMap() {
|
Object systemConfigMapObj = redisUtil.get(RedisKeyType.SYSTEM_CONFIG_MAP.key);
|
if (systemConfigMapObj == null) {
|
return null;
|
}
|
return (HashMap<String, String>) systemConfigMapObj;
|
}
|
|
private boolean hasAvailableStationTaskCapacity(HashMap<String, String> systemConfigMap) {
|
int conveyorStationTaskLimit = getConveyorStationTaskLimit(systemConfigMap);
|
int currentStationTaskCount = stationOperateProcessUtils.getCurrentStationTaskCount();
|
if (currentStationTaskCount > conveyorStationTaskLimit) {
|
News.error("输送站点任务已达到上限,上限值:{},站点任务数:{}", conveyorStationTaskLimit, currentStationTaskCount);
|
return false;
|
}
|
return true;
|
}
|
|
private synchronized boolean tryReserveGenerateCapacity(HashMap<String, String> systemConfigMap) {
|
int conveyorStationTaskLimit = getConveyorStationTaskLimit(systemConfigMap);
|
int currentStationTaskCount = stationOperateProcessUtils.getCurrentStationTaskCount();
|
if (currentStationTaskCount + inFlightGenerateCount >= conveyorStationTaskLimit) {
|
News.error("输送站点任务已达到上限,上限值:{},站点任务数:{},生成中任务数:{}",
|
conveyorStationTaskLimit, currentStationTaskCount, inFlightGenerateCount);
|
return false;
|
}
|
inFlightGenerateCount++;
|
return true;
|
}
|
|
private synchronized void releaseGenerateCapacity() {
|
if (inFlightGenerateCount > 0) {
|
inFlightGenerateCount--;
|
}
|
}
|
|
private int getConveyorStationTaskLimit(HashMap<String, String> systemConfigMap) {
|
int conveyorStationTaskLimit = 30;
|
String conveyorStationTaskLimitStr = systemConfigMap.get("conveyorStationTaskLimit");
|
if (conveyorStationTaskLimitStr != null) {
|
conveyorStationTaskLimit = Integer.parseInt(conveyorStationTaskLimitStr);
|
}
|
return conveyorStationTaskLimit;
|
}
|
|
}
|