From 41aeff86351d1dd94fe2408175f96475f227c1b9 Mon Sep 17 00:00:00 2001
From: Junjie <fallin.jie@qq.com>
Date: 星期四, 02 四月 2026 17:15:27 +0800
Subject: [PATCH] #执行优化

---
 src/main/java/com/zy/core/utils/StationOperateProcessUtils.java |  727 +++++++++++++++++++++++++++++++------------------------
 1 files changed, 406 insertions(+), 321 deletions(-)

diff --git a/src/main/java/com/zy/core/utils/StationOperateProcessUtils.java b/src/main/java/com/zy/core/utils/StationOperateProcessUtils.java
index 8237bbe..247d6ee 100644
--- a/src/main/java/com/zy/core/utils/StationOperateProcessUtils.java
+++ b/src/main/java/com/zy/core/utils/StationOperateProcessUtils.java
@@ -1,379 +1,464 @@
 package com.zy.core.utils;
 
-import com.alibaba.fastjson.JSON;
-import com.alibaba.fastjson.JSONObject;
-import com.baomidou.mybatisplus.mapper.EntityWrapper;
-import com.core.exception.CoolException;
-import com.zy.asrs.entity.*;
+import com.baomidou.mybatisplus.core.conditions.query.QueryWrapper;
+import com.zy.asrs.entity.BasCrnp;
+import com.zy.asrs.entity.BasDevp;
+import com.zy.asrs.entity.WrkMast;
 import com.zy.asrs.service.*;
-import com.zy.common.entity.FindCrnNoResult;
-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.MessageQueue;
 import com.zy.core.cache.SlaveConnection;
-import com.zy.core.enums.*;
+import com.zy.core.enums.SlaveType;
+import com.zy.core.enums.WrkIoType;
+import com.zy.core.enums.WrkStsType;
 import com.zy.core.model.StationObjModel;
-import com.zy.core.model.Task;
-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.station.*;
+import com.zy.core.utils.station.model.RerouteCommandPlan;
+import com.zy.core.utils.station.model.RerouteContext;
+import com.zy.core.utils.station.model.RerouteDecision;
+import com.zy.core.utils.station.model.RerouteExecutionResult;
 import org.springframework.beans.factory.annotation.Autowired;
 import org.springframework.stereotype.Component;
 
-import java.util.*;
+import java.util.Date;
+import java.util.List;
+import java.util.Map;
 
 @Component
 public class StationOperateProcessUtils {
-
-    @Autowired
-    private BasDevpService basDevpService;
     @Autowired
     private WrkMastService wrkMastService;
     @Autowired
-    private CommonService commonService;
+    private BasDevpService basDevpService;
+    @Autowired
+    private BasStationService basStationService;
+    @Autowired
+    private WrkAnalysisService wrkAnalysisService;
+    @Autowired
+    private StationRegularDispatchProcessor stationRegularDispatchProcessor;
+    @Autowired
+    private StationDispatchLoadSupport stationDispatchLoadSupport;
+    @Autowired
+    private StationOutboundDispatchProcessor stationOutboundDispatchProcessor;
+    @Autowired
+    private StationRerouteProcessor stationRerouteProcessor;
+    @Autowired
+    private MainProcessTaskSubmitter mainProcessTaskSubmitter;
+    @Autowired
+    private StationOutboundDecisionSupport stationOutboundDecisionSupport;
     @Autowired
     private BasCrnpService basCrnpService;
-    @Autowired
-    private BasDualCrnpService basDualCrnpService;
-    @Autowired
-    private RedisUtil redisUtil;
-    @Autowired
-    private LocMastService locMastService;
-    @Autowired
-    private WmsOperateUtils wmsOperateUtils;
 
     //鎵ц杈撻�佺珯鐐瑰叆搴撲换鍔�
     public synchronized void stationInExecute() {
-        List<BasDevp> basDevps = basDevpService.selectList(new EntityWrapper<>());
-        for (BasDevp basDevp : basDevps) {
-            StationThread stationThread = (StationThread) SlaveConnection.get(SlaveType.Devp, basDevp.getDevpNo());
-            if(stationThread == null){
-                continue;
-            }
+        stationRegularDispatchProcessor.stationInExecute();
+    }
 
-            Map<Integer, StationProtocol> stationMap = stationThread.getStatusMap();
+    // 鎵ц鍗曚釜绔欑偣鐨勫叆搴撲换鍔′笅鍙�
+    public void stationInExecute(BasDevp basDevp, StationObjModel stationObjModel) {
+        stationRegularDispatchProcessor.stationInExecute(basDevp, stationObjModel);
+    }
 
-            List<StationObjModel> list = basDevp.getBarcodeStationList$();
-            for (StationObjModel entity : list) {
-                Integer stationId = entity.getStationId();
-                if(!stationMap.containsKey(stationId)){
-                    continue;
-                }
+    //鎵ц鍫嗗灈鏈鸿緭閫佺珯鐐瑰嚭搴撲换鍔�
+    public synchronized void crnStationOutExecute() {
+        stationOutboundDispatchProcessor.crnStationOutExecute();
+    }
 
-                StationProtocol stationProtocol = stationMap.get(stationId);
-                if (stationProtocol == null) {
-                    continue;
-                }
+    // 鎵ц鍗曚釜鍑哄簱浠诲姟瀵瑰簲鐨勮緭閫佺珯鐐逛笅鍙�
+    public void crnStationOutExecute(WrkMast wrkMast) {
+        stationOutboundDispatchProcessor.crnStationOutExecute(wrkMast);
+    }
 
-                Object lock = redisUtil.get(RedisKeyType.STATION_IN_EXECUTE_LIMIT.key + stationId);
-                if(lock != null){
-                    continue;
-                }
+    //鎵ц鍙屽伐浣嶅爢鍨涙満杈撻�佺珯鐐瑰嚭搴撲换鍔�
+    public synchronized void dualCrnStationOutExecute() {
+        stationOutboundDispatchProcessor.dualCrnStationOutExecute();
+    }
 
-                //婊¤冻鑷姩銆佹湁鐗┿�佹湁宸ヤ綔鍙�
-                if (stationProtocol.isAutoing()
-                        && stationProtocol.isLoading()
-                        && stationProtocol.getTaskNo() > 0
-                ) {
-                    //妫�娴嬩换鍔℃槸鍚︾敓鎴�
-                    WrkMast wrkMast = wrkMastService.selectOne(new EntityWrapper<WrkMast>().eq("barcode", stationProtocol.getBarcode()));
-                    if (wrkMast == null) {
-                        continue;
-                    }
+    // 妫�娴嬪崟涓嚭搴撲换鍔℃槸鍚﹀埌杈剧洰鏍囩珯鍙�
+    public void stationOutExecuteFinish(StationObjModel stationObjModel) {
+        stationRegularDispatchProcessor.stationOutExecuteFinish(stationObjModel);
+    }
 
-                    if (wrkMast.getWrkSts() == WrkStsType.INBOUND_DEVICE_RUN.sts) {
-                        continue;
-                    }
-
-                    String locNo = wrkMast.getLocNo();
-                    FindCrnNoResult findCrnNoResult = commonService.findCrnNoByLocNo(locNo);
-                    if (findCrnNoResult == null) {
-                        News.taskInfo(wrkMast.getWrkNo(), "{}宸ヤ綔,鏈尮閰嶅埌鍫嗗灈鏈�", wrkMast.getWrkNo());
-                        continue;
-                    }
-
-                    Integer targetStationId = commonService.findInStationId(findCrnNoResult, stationId);
-                    if (targetStationId == null) {
-                        News.taskInfo(wrkMast.getWrkNo(), "{}绔欑偣,鎼滅储鍏ュ簱绔欑偣澶辫触", stationId);
-                        continue;
-                    }
-
-                    StationCommand command = stationThread.getCommand(StationCommandType.MOVE, wrkMast.getWrkNo(), stationId, targetStationId, 0);
-                    if(command == null){
-                        News.taskInfo(wrkMast.getWrkNo(), "{}宸ヤ綔,鑾峰彇杈撻�佺嚎鍛戒护澶辫触", wrkMast.getWrkNo());
-                        continue;
-                    }
-
-                    wrkMast.setWrkSts(WrkStsType.INBOUND_DEVICE_RUN.sts);
-                    wrkMast.setSourceStaNo(stationProtocol.getStationId());
-                    wrkMast.setStaNo(targetStationId);
-                    wrkMast.setSystemMsg("");
-                    wrkMast.setIoTime(new Date());
-                    if (wrkMastService.updateById(wrkMast)) {
-                        MessageQueue.offer(SlaveType.Devp, basDevp.getDevpNo(), new Task(2, command));
-                        News.info("杈撻�佺珯鐐瑰叆搴撳懡浠や笅鍙戞垚鍔燂紝绔欑偣鍙�={}锛屽伐浣滃彿={}锛屽懡浠ゆ暟鎹�={}", stationId, wrkMast.getWrkNo(), JSON.toJSONString(command));
-                        redisUtil.set(RedisKeyType.STATION_IN_EXECUTE_LIMIT.key + stationId, "lock", 5);
-                    }
-                }
-            }
+    // 妫�娴嬪崟涓叆搴撲换鍔℃槸鍚﹀埌杈剧洰鏍囩珯鍙�
+    public void scanInboundStationArrival(StationObjModel stationObjModel) {
+        if (stationObjModel == null) {
+            return;
+        }
+        StationThread stationThread = (StationThread) SlaveConnection.get(SlaveType.Devp, stationObjModel.getDeviceNo());
+        if (stationThread == null) {
+            return;
+        }
+        Map<Integer, StationProtocol> statusMap = stationThread.getStatusMap();
+        StationProtocol stationProtocol = statusMap == null ? null : statusMap.get(stationObjModel.getStationId());
+        if (stationProtocol == null) {
+            return;
+        }
+        if (stationProtocol.getTaskNo() <= 0) {
+            return;
+        }
+        WrkMast wrkMast = wrkMastService.selectByWorkNo(stationProtocol.getTaskNo());
+        if (wrkMast == null) {
+            return;
+        }
+        boolean updated = wrkAnalysisService.completeInboundStationRun(wrkMast, new Date());
+        if (updated) {
+            News.info("鍏ュ簱绔欑偣鍒拌揪鎵弿鍛戒腑锛屽伐浣滃彿={}锛岀洰鏍囩珯={}", wrkMast.getWrkNo(), wrkMast.getStaNo());
         }
     }
 
-    //鎵ц杈撻�佺珯鐐瑰嚭搴撲换鍔�
-    public synchronized void stationOutExecute() {
-        List<WrkMast> wrkMasts = wrkMastService.selectList(new EntityWrapper<WrkMast>().eq("wrk_sts", WrkStsType.OUTBOUND_RUN_COMPLETE.sts));
-        for (WrkMast wrkMast : wrkMasts) {
-            List<StationObjModel> outStationList = new ArrayList<>();
-
-            BasCrnp basCrnp = basCrnpService.selectOne(new EntityWrapper<BasCrnp>().eq("crn_no", wrkMast.getCrnNo()));
-            if (basCrnp != null) {
-                outStationList = basCrnp.getOutStationList$();
-                if(outStationList.isEmpty()){
-                    News.info("鍫嗗灈鏈�:{} 鍑哄簱绔欑偣鏈缃�", basCrnp.getCrnNo());
-                    continue;
-                }
-            }
-
-            BasDualCrnp basDualCrnp = basDualCrnpService.selectOne(new EntityWrapper<BasDualCrnp>().eq("crn_no", wrkMast.getDualCrnNo()));
-            if (basDualCrnp != null) {
-                outStationList = basDualCrnp.getOutStationList$();
-                if(outStationList.isEmpty()){
-                    News.info("鍙屽伐浣嶅爢鍨涙満:{} 鍑哄簱绔欑偣鏈缃�", basDualCrnp.getCrnNo());
-                    continue;
-                }
-            }
-
-            for (StationObjModel stationObjModel : outStationList) {
-                StationThread stationThread = (StationThread) SlaveConnection.get(SlaveType.Devp, stationObjModel.getDeviceNo());
-                if(stationThread == null){
-                    continue;
-                }
-
-                Map<Integer, StationProtocol> stationMap = stationThread.getStatusMap();
-                StationProtocol stationProtocol = stationMap.get(stationObjModel.getStationId());
-                if (stationProtocol == null) {
-                    continue;
-                }
-
-                Object lock = redisUtil.get(RedisKeyType.STATION_OUT_EXECUTE_LIMIT.key + stationProtocol.getStationId());
-                if (lock != null) {
-                    continue;
-                }
-
-                //婊¤冻鑷姩銆佹湁鐗┿�佸伐浣滃彿0
-                if (stationProtocol.isAutoing()
-                        && stationProtocol.isLoading()
-                        && stationProtocol.getTaskNo() == 0
-                ) {
-                    StationCommand command = stationThread.getCommand(StationCommandType.MOVE, wrkMast.getWrkNo(), stationProtocol.getStationId(), wrkMast.getStaNo(), 0);
-                    if(command == null){
-                        News.taskInfo(wrkMast.getWrkNo(), "鑾峰彇杈撻�佺嚎鍛戒护澶辫触");
-                        continue;
-                    }
-
-                    wrkMast.setWrkSts(WrkStsType.STATION_RUN.sts);
-                    wrkMast.setSystemMsg("");
-                    wrkMast.setIoTime(new Date());
-                    if (wrkMastService.updateById(wrkMast)) {
-                        MessageQueue.offer(SlaveType.Devp, stationObjModel.getDeviceNo(), new Task(2, command));
-                        News.info("杈撻�佺珯鐐瑰嚭搴撳懡浠や笅鍙戞垚鍔燂紝绔欑偣鍙�={}锛屽伐浣滃彿={}锛屽懡浠ゆ暟鎹�={}", stationProtocol.getStationId(), wrkMast.getWrkNo(), JSON.toJSONString(command));
-                        redisUtil.set(RedisKeyType.STATION_OUT_EXECUTE_LIMIT.key + stationProtocol.getStationId(), "lock", 5);
-                        redisUtil.set(RedisKeyType.STATION_OUT_EXECUTE_COMPLETE_LIMIT.key + wrkMast.getWrkNo(), "lock", 60 * 5);
-                    }
-                }
-            }
-        }
+    // 妫�娴嬩换鍔¤浆瀹屾垚
+    public void checkTaskToComplete() {
+        stationRegularDispatchProcessor.checkTaskToComplete();
     }
 
-    //妫�娴嬭緭閫佺珯鐐瑰嚭搴撲换鍔℃墽琛屽畬鎴�
-    public synchronized void stationOutExecuteFinish() {
-        List<WrkMast> wrkMasts = wrkMastService.selectList(new EntityWrapper<WrkMast>().eq("wrk_sts", WrkStsType.STATION_RUN.sts));
-        for (WrkMast wrkMast : wrkMasts) {
-            Integer wrkNo = wrkMast.getWrkNo();
-
-            Object lock = redisUtil.get(RedisKeyType.STATION_OUT_EXECUTE_COMPLETE_LIMIT.key + wrkNo);
-            if (lock != null) {
-                continue;
-            }
-
-            boolean complete = true;
-            List<BasDevp> basDevps = basDevpService.selectList(new EntityWrapper<>());
-            for (BasDevp basDevp : basDevps) {
-                StationThread stationThread = (StationThread) SlaveConnection.get(SlaveType.Devp, basDevp.getDevpNo());
-                if (stationThread == null) {
-                    continue;
-                }
-
-                List<StationProtocol> list = stationThread.getStatus();
-                for (StationProtocol stationProtocol : list) {
-                    if (stationProtocol.getTaskNo().equals(wrkNo)) {
-                        complete = false;
-                    }
-                }
-            }
-
-            if (complete) {
-                wrkMast.setWrkSts(WrkStsType.COMPLETE_OUTBOUND.sts);
-                wrkMast.setIoTime(new Date());
-                wrkMastService.updateById(wrkMast);
-            }
-        }
+    // 妫�娴嬪崟涓嚭搴撲换鍔℃槸鍚﹀彲浠ヨ浆瀹屾垚
+    public void checkTaskToComplete(WrkMast wrkMast) {
+        stationRegularDispatchProcessor.checkTaskToComplete(wrkMast);
     }
 
     //妫�娴嬭緭閫佺珯鐐规槸鍚﹁繍琛屽牭濉�
-    public synchronized void checkStationRunBlock() {
-        List<BasDevp> basDevps = basDevpService.selectList(new EntityWrapper<>());
+    public void checkStationRunBlock() {
+        stationRerouteProcessor.checkStationRunBlock();
+    }
+
+    // 妫�娴嬪崟涓珯鐐规槸鍚﹁繍琛屽牭濉�
+    public void checkStationRunBlock(BasDevp basDevp, Integer stationId) {
+        stationRerouteProcessor.checkStationRunBlock(basDevp, stationId);
+    }
+
+    //妫�娴嬭緭閫佺珯鐐逛换鍔″仠鐣欒秴鏃跺悗閲嶆柊璁$畻璺緞
+    public void checkStationIdleRecover() {
+        stationRerouteProcessor.checkStationIdleRecover();
+    }
+
+    // 妫�娴嬪崟涓珯鐐逛换鍔″仠鐣欒秴鏃跺悗鐨勬仮澶嶅鐞�
+    public void checkStationIdleRecover(BasDevp basDevp, Integer stationId) {
+        stationRerouteProcessor.checkStationIdleRecover(basDevp, stationId);
+    }
+
+    //鑾峰彇杈撻�佺嚎浠诲姟鏁伴噺
+    public int getCurrentStationTaskCount() {
+        return stationDispatchLoadSupport.countCurrentStationTask();
+    }
+
+    public int getCurrentOutboundTaskCountByTargetStation(Integer stationId) {
+        if (stationId == null) {
+            return 0;
+        }
+        return (int) wrkMastService.count(new QueryWrapper<WrkMast>()
+                .eq("io_type", WrkIoType.OUT.id)
+                .eq("sta_no", stationId)
+                .in("wrk_sts",
+                        WrkStsType.OUTBOUND_RUN.sts,
+                        WrkStsType.OUTBOUND_RUN_COMPLETE.sts,
+                        WrkStsType.STATION_RUN.sts));
+    }
+
+    // 妫�娴嬪嚭搴撴帓搴�
+    public synchronized void checkStationOutOrder() {
+        stationRerouteProcessor.checkStationOutOrder();
+    }
+
+    // 妫�娴嬪崟涓珯鐐圭殑鍑哄簱鎺掑簭
+    public void checkStationOutOrder(BasDevp basDevp, StationObjModel stationObjModel) {
+        stationRerouteProcessor.checkStationOutOrder(basDevp, stationObjModel);
+    }
+
+    // 鐩戞帶缁曞湀绔欑偣
+    public synchronized void watchCircleStation() {
+        stationRerouteProcessor.watchCircleStation();
+    }
+
+    // 鐩戞帶鍗曚釜缁曞湀绔欑偣
+    public void watchCircleStation(BasDevp basDevp, Integer stationId) {
+        stationRerouteProcessor.watchCircleStation(basDevp, stationId);
+    }
+
+    public void submitStationInTasks(long minIntervalMs) {
+        submitStationInTasks(MainProcessLane.STATION_IN, minIntervalMs);
+    }
+
+    public void submitStationInTasks(MainProcessLane lane, long minIntervalMs) {
+        List<BasDevp> basDevps = basDevpService.list(new QueryWrapper<>());
         for (BasDevp basDevp : basDevps) {
             StationThread stationThread = (StationThread) SlaveConnection.get(SlaveType.Devp, basDevp.getDevpNo());
-            if(stationThread == null){
+            if (stationThread == null) {
                 continue;
             }
-
-            List<Integer> runBlockReassignLocStationList = new ArrayList<>();
-            for (StationObjModel stationObjModel : basDevp.getRunBlockReassignLocStationList$()) {
-                runBlockReassignLocStationList.add(stationObjModel.getStationId());
+            Map<Integer, StationProtocol> stationMap = stationThread.getStatusMap();
+            if (stationMap == null || stationMap.isEmpty()) {
+                continue;
             }
-
-            List<StationProtocol> list = stationThread.getStatus();
-            for (StationProtocol stationProtocol : list) {
-                if(stationProtocol.isAutoing()
-                    && stationProtocol.isLoading()
-                    && stationProtocol.getTaskNo() > 0
-                    && stationProtocol.isRunBlock()
-                ) {
-                    WrkMast wrkMast = wrkMastService.selectByWorkNo(stationProtocol.getTaskNo());
-                    if (wrkMast == null) {
-                        News.info("杈撻�佺珯鐐瑰彿={} 杩愯闃诲锛屼絾鏃犳硶鎵惧埌瀵瑰簲浠诲姟锛屽伐浣滃彿={}", stationProtocol.getStationId(), stationProtocol.getTaskNo());
-                        continue;
-                    }
-
-                    Object lock = redisUtil.get(RedisKeyType.CHECK_STATION_RUN_BLOCK_LIMIT_.key + stationProtocol.getTaskNo());
-                    if (lock != null) {
-                        continue;
-                    }
-                    redisUtil.set(RedisKeyType.CHECK_STATION_RUN_BLOCK_LIMIT_.key + stationProtocol.getTaskNo(), "lock", 15);
-
-                    if (wrkMast.getIoType() == WrkIoType.IN.id && runBlockReassignLocStationList.contains(stationProtocol.getStationId())) {
-                        //绔欑偣澶勪簬閲嶆柊鍒嗛厤搴撲綅鍖哄煙
-                        //杩愯鍫靛锛岄噸鏂扮敵璇蜂换鍔�
-                        String response = wmsOperateUtils.applyReassignTaskLocNo(wrkMast.getWrkNo(), stationProtocol.getStationId());
-                        if (response == null) {
-                            News.taskError(wrkMast.getWrkNo(), "璇锋眰WMS閲嶆柊鍒嗛厤搴撲綅鎺ュ彛澶辫触锛屾帴鍙f湭鍝嶅簲锛侊紒锛乺esponse锛歿}", response);
-                            continue;
-                        }
-                        JSONObject jsonObject = JSON.parseObject(response);
-                        if (jsonObject.getInteger("code").equals(200)) {
-                            StartupDto dto = jsonObject.getObject("data", StartupDto.class);
-
-                            String sourceLocNo = wrkMast.getLocNo();
-                            String locNo = dto.getLocNo();
-
-                            LocMast sourceLocMast = locMastService.queryByLoc(sourceLocNo);
-                            if (sourceLocMast == null) {
-                                News.taskInfo(wrkMast.getWrkNo(), "搴撲綅鍙�:{} 婧愬簱浣嶄俊鎭笉瀛樺湪", sourceLocNo);
-                                continue;
-                            }
-
-                            if (!sourceLocMast.getLocSts().equals("S")) {
-                                News.taskInfo(wrkMast.getWrkNo(), "搴撲綅鍙�:{} 婧愬簱浣嶇姸鎬佷笉澶勪簬鍏ュ簱棰勭害", sourceLocNo);
-                                continue;
-                            }
-
-                            LocMast locMast = locMastService.queryByLoc(locNo);
-                            if (locMast == null) {
-                                News.taskInfo(wrkMast.getWrkNo(), "搴撲綅鍙�:{} 鐩爣搴撲綅淇℃伅涓嶅瓨鍦�", locNo);
-                                continue;
-                            }
-
-                            if (!locMast.getLocSts().equals("O")) {
-                                News.taskInfo(wrkMast.getWrkNo(), "搴撲綅鍙�:{} 鐩爣搴撲綅鐘舵�佷笉澶勪簬绌哄簱浣�", locNo);
-                                continue;
-                            }
-
-                            FindCrnNoResult findCrnNoResult = commonService.findCrnNoByLocNo(locNo);
-                            if (findCrnNoResult == null) {
-                                News.taskInfo(wrkMast.getWrkNo(), "{}宸ヤ綔,鏈尮閰嶅埌鍫嗗灈鏈�", wrkMast.getWrkNo());
-                                continue;
-                            }
-                            Integer crnNo = findCrnNoResult.getCrnNo();
-
-                            Integer targetStationId = commonService.findInStationId(findCrnNoResult, stationProtocol.getStationId());
-                            if (targetStationId == null) {
-                                News.taskInfo(wrkMast.getWrkNo(), "{}绔欑偣,鎼滅储鍏ュ簱绔欑偣澶辫触", stationProtocol.getStationId());
-                                continue;
-                            }
-
-                            StationCommand command = stationThread.getCommand(StationCommandType.MOVE, wrkMast.getWrkNo(), stationProtocol.getStationId(), targetStationId, 0);
-                            if(command == null){
-                                News.taskInfo(wrkMast.getWrkNo(), "{}宸ヤ綔,鑾峰彇杈撻�佺嚎鍛戒护澶辫触", wrkMast.getWrkNo());
-                                continue;
-                            }
-
-                            //鏇存柊婧愬簱浣�
-                            sourceLocMast.setLocSts("O");
-                            sourceLocMast.setModiTime(new Date());
-                            locMastService.updateById(sourceLocMast);
-
-                            //鏇存柊鐩爣搴撲綅
-                            locMast.setLocSts("S");
-                            locMast.setModiTime(new Date());
-                            locMastService.updateById(locMast);
-
-                            //鏇存柊宸ヤ綔妗f暟鎹�
-                            wrkMast.setLocNo(locNo);
-                            wrkMast.setStaNo(targetStationId);
-
-                            if (findCrnNoResult.getCrnType().equals(SlaveType.Crn)) {
-                                wrkMast.setCrnNo(crnNo);
-                            } else if (findCrnNoResult.getCrnType().equals(SlaveType.DualCrn)) {
-                                wrkMast.setDualCrnNo(crnNo);
-                            }else {
-                                throw new CoolException("鏈煡璁惧绫诲瀷");
-                            }
-
-                            if (wrkMastService.updateById(wrkMast)) {
-                                MessageQueue.offer(SlaveType.Devp, basDevp.getDevpNo(), new Task(2, command));
-                            }
-                        } else {
-                            News.error("璇锋眰WMS鎺ュ彛澶辫触锛侊紒锛乺esponse锛歿}", response);
-                        }
-                    }else {
-                        //杩愯鍫靛锛岄噸鏂拌绠楄矾绾�
-                        StationCommand command = stationThread.getCommand(StationCommandType.MOVE, wrkMast.getWrkNo(), stationProtocol.getStationId(), wrkMast.getStaNo(), 0);
-                        if(command == null){
-                            News.taskInfo(wrkMast.getWrkNo(), "鑾峰彇杈撻�佺嚎鍛戒护澶辫触");
-                            continue;
-                        }
-
-                        MessageQueue.offer(SlaveType.Devp, basDevp.getDevpNo(), new Task(2, command));
-                        News.info("杈撻�佺珯鐐瑰牭濉炲悗閲嶆柊璁$畻璺緞鍛戒护涓嬪彂鎴愬姛锛岀珯鐐瑰彿={}锛屽伐浣滃彿={}锛屽懡浠ゆ暟鎹�={}", stationProtocol.getStationId(), wrkMast.getWrkNo(), JSON.toJSONString(command));
-                    }
+            for (StationObjModel stationObjModel : basDevp.getBarcodeStationList$()) {
+                Integer stationId = stationObjModel == null ? null : stationObjModel.getStationId();
+                if (stationId == null || !stationMap.containsKey(stationId)) {
+                    continue;
                 }
+                StationProtocol stationProtocol = stationMap.get(stationId);
+                if (stationProtocol == null
+                        || !stationProtocol.isAutoing()
+                        || !stationProtocol.isLoading()
+                        || stationProtocol.getTaskNo() <= 0) {
+                    continue;
+                }
+                mainProcessTaskSubmitter.submitKeyedSerialTask(
+                        lane,
+                        stationId,
+                        "stationInExecute",
+                        minIntervalMs,
+                        () -> stationInExecute(basDevp, stationObjModel)
+                );
             }
         }
     }
 
-    //鑾峰彇杈撻�佺嚎浠诲姟鏁伴噺
-    public synchronized int getCurrentStationTaskCount() {
-        int currentStationTaskCount = 0;
-        List<BasDevp> basDevps = basDevpService.selectList(new EntityWrapper<BasDevp>());
+    public void submitCrnStationOutTasks(long minIntervalMs) {
+        submitCrnStationOutTasks(MainProcessLane.STATION_OUT, minIntervalMs);
+    }
+
+    public void submitCrnStationOutTasks(MainProcessLane lane, long minIntervalMs) {
+        List<WrkMast> wrkMasts = wrkMastService.list(new QueryWrapper<WrkMast>()
+                .eq("wrk_sts", WrkStsType.OUTBOUND_RUN_COMPLETE.sts)
+                .isNotNull("crn_no"));
+        for (WrkMast wrkMast : wrkMasts) {
+            Integer laneKey = wrkMast == null ? null : wrkMast.getSourceStaNo();
+            if (laneKey == null) {
+                laneKey = wrkMast == null ? null : wrkMast.getWrkNo();
+            }
+            mainProcessTaskSubmitter.submitKeyedSerialTask(
+                    lane,
+                    laneKey,
+                    "crnStationOutExecute",
+                    minIntervalMs,
+                    () -> crnStationOutExecute(wrkMast)
+            );
+        }
+    }
+
+    public void submitStationOutExecuteFinishTasks(long minIntervalMs) {
+        submitStationOutExecuteFinishTasks(MainProcessLane.STATION_OUT_FINISH, minIntervalMs);
+    }
+
+    public void submitStationOutExecuteFinishTasks(MainProcessLane lane, long minIntervalMs) {
+        List<BasDevp> basDevps = basDevpService.list(new QueryWrapper<>());
         for (BasDevp basDevp : basDevps) {
-            StationThread stationThread = (StationThread) SlaveConnection.get(SlaveType.Devp, basDevp.getId());
+            for (StationObjModel stationObjModel : basDevp.getOutStationList$()) {
+                mainProcessTaskSubmitter.submitKeyedSerialTask(
+                        lane,
+                        stationObjModel.getStationId(),
+                        "stationOutExecuteFinish",
+                        minIntervalMs,
+                        () -> stationOutExecuteFinish(stationObjModel)
+                );
+            }
+        }
+    }
+
+    public void submitInboundStationArrivalTasks(long minIntervalMs) {
+        submitInboundStationArrivalTasks(MainProcessLane.STATION_IN_ARRIVAL, minIntervalMs);
+    }
+
+    public void submitInboundStationArrivalTasks(MainProcessLane lane, long minIntervalMs) {
+        List<BasCrnp> basCrnps = basCrnpService.list(new QueryWrapper<>());
+        for (BasCrnp basCrnp : basCrnps) {
+            Integer crnNo = basCrnp == null ? null : basCrnp.getCrnNo();
+            if (crnNo == null) {
+                continue;
+            }
+
+            for (StationObjModel stationObjModel : basCrnp.getInStationList$()) {
+                mainProcessTaskSubmitter.submitKeyedSerialTask(
+                        lane,
+                        stationObjModel.getStationId(),
+                        "scanInboundStationArrival",
+                        minIntervalMs,
+                        () -> scanInboundStationArrival(stationObjModel)
+                );
+            }
+        }
+    }
+
+    public void submitCheckTaskToCompleteTasks(long minIntervalMs) {
+        submitCheckTaskToCompleteTasks(MainProcessLane.STATION_COMPLETE, minIntervalMs);
+    }
+
+    public void submitCheckTaskToCompleteTasks(MainProcessLane lane, long minIntervalMs) {
+        List<WrkMast> wrkMasts = wrkMastService.list(new QueryWrapper<WrkMast>()
+                .eq("wrk_sts", WrkStsType.STATION_RUN_COMPLETE.sts)
+                .isNotNull("sta_no"));
+        for (WrkMast wrkMast : wrkMasts) {
+            mainProcessTaskSubmitter.submitKeyedSerialTask(
+                    lane,
+                    wrkMast.getStaNo(),
+                    "checkTaskToComplete",
+                    minIntervalMs,
+                    () -> checkTaskToComplete(wrkMast)
+            );
+        }
+    }
+
+    public void submitCheckStationOutOrderTasks(long minIntervalMs) {
+        submitCheckStationOutOrderTasks(MainProcessLane.STATION_OUT_ORDER, minIntervalMs);
+    }
+
+    public void submitCheckStationOutOrderTasks(MainProcessLane lane, long minIntervalMs) {
+        List<BasDevp> basDevps = basDevpService.list(new QueryWrapper<>());
+        for (BasDevp basDevp : basDevps) {
+            StationThread stationThread = (StationThread) SlaveConnection.get(SlaveType.Devp, basDevp.getDevpNo());
             if (stationThread == null) {
                 continue;
             }
 
-            for (StationProtocol stationProtocol : stationThread.getStatus()) {
-                if (stationProtocol.getTaskNo() > 0) {
-                    currentStationTaskCount++;
+            for (StationObjModel stationObjModel : basDevp.getOutOrderList$()) {
+                Integer stationId = stationObjModel == null ? null : stationObjModel.getStationId();
+                if (stationId == null) {
+                    continue;
                 }
+                Map<Integer, StationProtocol> statusMap = stationThread.getStatusMap();
+                StationProtocol stationProtocol = statusMap.get(stationId);
+                if (stationProtocol == null
+                        || !stationProtocol.isAutoing()
+                        || !stationProtocol.isLoading()
+                        || stationProtocol.getTaskNo() <= 0
+                        || stationProtocol.isRunBlock()
+                        || !stationProtocol.getStationId().equals(stationProtocol.getTargetStaNo())) {
+                    continue;
+                }
+                mainProcessTaskSubmitter.submitKeyedSerialTask(
+                        lane,
+                        stationId,
+                        "checkStationOutOrder",
+                        minIntervalMs,
+                        () -> checkStationOutOrder(basDevp, stationObjModel)
+                );
             }
         }
-
-        return currentStationTaskCount;
     }
 
+    public void submitWatchCircleStationTasks(long minIntervalMs) {
+        submitWatchCircleStationTasks(MainProcessLane.STATION_WATCH_CIRCLE, minIntervalMs);
+    }
 
+    public void submitWatchCircleStationTasks(MainProcessLane lane, long minIntervalMs) {
+        List<BasDevp> basDevps = basDevpService.list(new QueryWrapper<>());
+        for (BasDevp basDevp : basDevps) {
+            StationThread stationThread = (StationThread) SlaveConnection.get(SlaveType.Devp, basDevp.getDevpNo());
+            if (stationThread == null) {
+                continue;
+            }
+            for (StationProtocol stationProtocol : stationThread.getStatus()) {
+                Integer stationId = stationProtocol == null ? null : stationProtocol.getStationId();
+                if (stationId == null) {
+                    continue;
+                }
+                if (!stationProtocol.isAutoing()
+                        || !stationProtocol.isLoading()
+                        || stationProtocol.getTaskNo() <= 0
+                        || !stationOutboundDecisionSupport.isWatchingCircleArrival(stationProtocol.getTaskNo(), stationProtocol.getStationId())) {
+                    continue;
+                }
+
+                mainProcessTaskSubmitter.submitKeyedSerialTask(
+                        lane,
+                        stationId,
+                        "watchCircleStation",
+                        minIntervalMs,
+                        () -> watchCircleStation(basDevp, stationId)
+                );
+            }
+        }
+    }
+
+    public void submitCheckStationRunBlockTasks(long minIntervalMs) {
+        submitCheckStationRunBlockTasks(MainProcessLane.STATION_RUN_BLOCK, minIntervalMs);
+    }
+
+    public void submitCheckStationRunBlockTasks(MainProcessLane lane, long minIntervalMs) {
+        List<BasDevp> basDevps = basDevpService.list(new QueryWrapper<>());
+        for (BasDevp basDevp : basDevps) {
+            StationThread stationThread = (StationThread) SlaveConnection.get(SlaveType.Devp, basDevp.getDevpNo());
+            if (stationThread == null) {
+                continue;
+            }
+            for (StationProtocol stationProtocol : stationThread.getStatus()) {
+                Integer stationId = stationProtocol == null ? null : stationProtocol.getStationId();
+                if (stationId == null) {
+                    continue;
+                }
+                if (!stationProtocol.isAutoing()
+                        || !stationProtocol.isLoading()
+                        || stationProtocol.getTaskNo() <= 0
+                        || !stationProtocol.isRunBlock()) {
+                    continue;
+                }
+                mainProcessTaskSubmitter.submitKeyedSerialTask(
+                        lane,
+                        stationId,
+                        "checkStationRunBlock",
+                        minIntervalMs,
+                        () -> checkStationRunBlock(basDevp, stationId)
+                );
+            }
+        }
+    }
+
+    public void submitCheckStationIdleRecoverTasks(long minIntervalMs) {
+        submitCheckStationIdleRecoverTasks(MainProcessLane.STATION_IDLE_RECOVER, minIntervalMs);
+    }
+
+    public void submitCheckStationIdleRecoverTasks(MainProcessLane lane, long minIntervalMs) {
+        List<BasDevp> basDevps = basDevpService.list(new QueryWrapper<>());
+        for (BasDevp basDevp : basDevps) {
+            StationThread stationThread = (StationThread) SlaveConnection.get(SlaveType.Devp, basDevp.getDevpNo());
+            if (stationThread == null) {
+                continue;
+            }
+            for (StationProtocol stationProtocol : stationThread.getStatus()) {
+                Integer stationId = stationProtocol == null ? null : stationProtocol.getStationId();
+                if (stationId == null) {
+                    continue;
+                }
+                if (!stationProtocol.isAutoing()
+                        || !stationProtocol.isLoading()
+                        || stationProtocol.getTaskNo() <= 0
+                        || stationProtocol.isRunBlock()) {
+                    continue;
+                }
+                mainProcessTaskSubmitter.submitKeyedSerialTask(
+                        lane,
+                        stationId,
+                        "checkStationIdleRecover",
+                        minIntervalMs,
+                        () -> checkStationIdleRecover(basDevp, stationId)
+                );
+            }
+        }
+    }
+
+    RerouteCommandPlan buildRerouteCommandPlan(RerouteContext context,
+                                               RerouteDecision decision) {
+        return stationRerouteProcessor.buildRerouteCommandPlan(context, decision);
+    }
+
+    RerouteExecutionResult executeReroutePlan(RerouteContext context,
+                                              RerouteCommandPlan plan) {
+        return stationRerouteProcessor.executeReroutePlan(context, plan);
+    }
+
+    RerouteDecision resolveSharedRerouteDecision(RerouteContext context) {
+        return stationRerouteProcessor.resolveSharedRerouteDecision(context);
+    }
+
+    boolean shouldUseRunBlockDirectReassign(WrkMast wrkMast,
+                                            Integer stationId,
+                                            List<Integer> runBlockReassignLocStationList) {
+        return stationRerouteProcessor.shouldUseRunBlockDirectReassign(wrkMast, stationId, runBlockReassignLocStationList);
+    }
+
+    private boolean shouldSkipIdleRecoverForRecentDispatch(Integer taskNo, Integer stationId) {
+        return stationRerouteProcessor.shouldSkipIdleRecoverForRecentDispatch(taskNo, stationId);
+    }
 }

--
Gitblit v1.9.1