From 7a5448174e5cb929e78926cce3783366557b7e88 Mon Sep 17 00:00:00 2001
From: Junjie <fallin.jie@qq.com>
Date: 星期六, 21 三月 2026 17:53:37 +0800
Subject: [PATCH] #
---
src/main/java/com/zy/core/trace/StationTaskTraceRegistry.java | 299 +++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
1 files changed, 297 insertions(+), 2 deletions(-)
diff --git a/src/main/java/com/zy/core/trace/StationTaskTraceRegistry.java b/src/main/java/com/zy/core/trace/StationTaskTraceRegistry.java
index 4c8e76b..5c8b962 100644
--- a/src/main/java/com/zy/core/trace/StationTaskTraceRegistry.java
+++ b/src/main/java/com/zy/core/trace/StationTaskTraceRegistry.java
@@ -1,10 +1,17 @@
package com.zy.core.trace;
+import com.alibaba.fastjson.JSON;
import com.zy.asrs.domain.vo.StationTaskTraceEventVo;
import com.zy.asrs.domain.vo.StationTaskTraceSegmentVo;
import com.zy.asrs.domain.vo.StationTaskTraceVo;
+import com.zy.asrs.entity.WrkMast;
+import com.zy.asrs.service.WrkMastService;
+import com.zy.common.utils.RedisUtil;
+import com.zy.core.enums.WrkStsType;
import com.zy.core.model.command.StationCommand;
+import com.zy.core.enums.RedisKeyType;
import lombok.Data;
+import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
import java.util.ArrayList;
@@ -20,6 +27,7 @@
private static final int MAX_EVENT_COUNT = 200;
private static final long TERMINAL_KEEP_MS = 30L * 60L * 1000L;
+ private static final String TRACE_CACHE_KEY = RedisKeyType.STATION_TASK_TRACE_REGISTRY.key;
public static final String STATUS_WAITING = "WAITING";
public static final String STATUS_RUNNING = "RUNNING";
@@ -29,7 +37,14 @@
public static final String STATUS_FINISHED = "FINISHED";
public static final String STATUS_REROUTED = "REROUTED";
+ @Autowired
+ private RedisUtil redisUtil;
+
+ @Autowired
+ private WrkMastService wrkMastService;
+
private final Map<Integer, TraceTaskState> taskStateMap = new ConcurrentHashMap<>();
+ private volatile boolean loadedFromRedis = false;
public TraceRegistration registerPlan(Integer taskNo,
String threadImpl,
@@ -42,10 +57,13 @@
return new TraceRegistration();
}
+ ensureCacheLoaded();
cleanupExpired();
TraceTaskState taskState = taskStateMap.computeIfAbsent(taskNo, TraceTaskState::new);
- return taskState.registerPlan(threadImpl, startStationId, currentStationId, finalTargetStationId,
+ TraceRegistration registration = taskState.registerPlan(threadImpl, startStationId, currentStationId, finalTargetStationId,
copyIntegerList(localPathStationIds), copySegmentList(localSegmentList));
+ persistState(taskState);
+ return registration;
}
public void markSegmentIssued(Integer taskNo,
@@ -54,11 +72,13 @@
String eventType,
String message,
Map<String, Object> details) {
+ ensureCacheLoaded();
TraceTaskState state = taskStateMap.get(taskNo);
if (state == null) {
return;
}
state.markSegmentIssued(traceVersion, command, eventType, message, details);
+ persistState(state);
}
public void updateProgress(Integer taskNo,
@@ -67,63 +87,75 @@
String eventType,
String message,
Map<String, Object> details) {
+ ensureCacheLoaded();
TraceTaskState state = taskStateMap.get(taskNo);
if (state == null) {
return;
}
state.updateProgress(traceVersion, currentStationId, eventType, message, details);
+ persistState(state);
}
public void markBlocked(Integer taskNo,
Integer traceVersion,
Integer currentStationId,
Map<String, Object> details) {
+ ensureCacheLoaded();
TraceTaskState state = taskStateMap.get(taskNo);
if (state == null) {
return;
}
state.markTerminal(traceVersion, STATUS_BLOCKED, currentStationId, currentStationId,
"BLOCKED", "杈撻�佷换鍔¤繍琛屽牭濉�", details);
+ persistState(state);
}
public void markCancelled(Integer taskNo,
Integer traceVersion,
Integer currentStationId,
Map<String, Object> details) {
+ ensureCacheLoaded();
TraceTaskState state = taskStateMap.get(taskNo);
if (state == null) {
return;
}
state.markTerminal(traceVersion, STATUS_CANCELLED, currentStationId, null,
"CANCELLED", "杈撻�佷换鍔℃敹鍒板彇娑堜俊鍙�", details);
+ persistState(state);
}
public void markTimeout(Integer taskNo,
Integer traceVersion,
Integer currentStationId,
Map<String, Object> details) {
+ ensureCacheLoaded();
TraceTaskState state = taskStateMap.get(taskNo);
if (state == null) {
return;
}
state.markTerminal(traceVersion, STATUS_TIMEOUT, currentStationId, null,
"TIMEOUT", "杈撻�佷换鍔¢暱鏃堕棿鏃犳硶瀹氫綅褰撳墠浣嶇疆", details);
+ persistState(state);
}
public void markFinished(Integer taskNo,
Integer traceVersion,
Integer currentStationId,
Map<String, Object> details) {
+ ensureCacheLoaded();
TraceTaskState state = taskStateMap.get(taskNo);
if (state == null) {
return;
}
state.markTerminal(traceVersion, STATUS_FINISHED, currentStationId, null,
"FINISHED", "杈撻�佷换鍔¤建杩瑰畬鎴�", details);
+ persistState(state);
}
public List<StationTaskTraceVo> listLatestTraces() {
+ ensureCacheLoaded();
cleanupExpired();
+ reconcileInactiveBusinessTasks();
List<StationTaskTraceVo> result = new ArrayList<>();
for (TraceTaskState state : taskStateMap.values()) {
if (state != null) {
@@ -141,13 +173,155 @@
return result;
}
+ public List<StationTaskTraceVo> listActiveTraces() {
+ List<StationTaskTraceVo> latestTraces = listLatestTraces();
+ List<StationTaskTraceVo> result = new ArrayList<>();
+ for (StationTaskTraceVo traceVo : latestTraces) {
+ if (traceVo == null) {
+ continue;
+ }
+ String status = traceVo.getStatus();
+ if (STATUS_WAITING.equals(status)
+ || STATUS_RUNNING.equals(status)
+ || STATUS_REROUTED.equals(status)) {
+ result.add(traceVo);
+ }
+ }
+ return result;
+ }
+
private void cleanupExpired() {
long now = System.currentTimeMillis();
for (Map.Entry<Integer, TraceTaskState> entry : taskStateMap.entrySet()) {
TraceTaskState state = entry.getValue();
if (state != null && state.shouldRemove(now)) {
taskStateMap.remove(entry.getKey(), state);
+ removePersistedState(entry.getKey());
}
+ }
+ }
+
+ private void reconcileInactiveBusinessTasks() {
+ if (wrkMastService == null) {
+ return;
+ }
+ for (TraceTaskState state : taskStateMap.values()) {
+ if (state == null || state.isTerminal()) {
+ continue;
+ }
+ WrkMast wrkMast;
+ try {
+ wrkMast = wrkMastService.selectByWorkNo(state.taskNo);
+ } catch (Exception ignore) {
+ continue;
+ }
+ if (wrkMast == null) {
+ Integer currentStationId = state.currentStationId != null ? state.currentStationId : state.finalTargetStationId;
+ Map<String, Object> details = new LinkedHashMap<>();
+ details.put("reason", "wrk_missing");
+ state.markTerminal(state.traceVersion, STATUS_FINISHED, currentStationId, null,
+ "AUTO_FINISHED", "杈撻�佷换鍔℃。宸蹭笉瀛樺湪锛岃建杩硅嚜鍔ㄧ粨鏉�", details);
+ persistState(state);
+ continue;
+ }
+ if (isStationTraceActiveWrkStatus(wrkMast.getWrkSts())) {
+ continue;
+ }
+ Integer currentStationId = state.currentStationId != null ? state.currentStationId : state.finalTargetStationId;
+ Map<String, Object> details = new LinkedHashMap<>();
+ details.put("reason", "wrk_status_transition");
+ details.put("wrkSts", wrkMast.getWrkSts());
+ details.put("wrkStsDesc", wrkMast.getWrkSts$());
+ if (isManualWrkStatus(wrkMast.getWrkSts())) {
+ state.markTerminal(state.traceVersion, STATUS_CANCELLED, currentStationId, null,
+ "AUTO_CANCELLED", "杈撻�佷换鍔″凡閫�鍑鸿繍琛岋紝杞ㄨ抗鑷姩缁撴潫", details);
+ } else {
+ state.markTerminal(state.traceVersion, STATUS_FINISHED, currentStationId, null,
+ "AUTO_FINISHED", "杈撻�佷换鍔″凡缁撴潫锛岃建杩硅嚜鍔ㄧ粨鏉�", details);
+ }
+ persistState(state);
+ }
+ }
+
+ private boolean isStationTraceActiveWrkStatus(Long wrkSts) {
+ return Objects.equals(wrkSts, WrkStsType.INBOUND_DEVICE_RUN.sts)
+ || Objects.equals(wrkSts, WrkStsType.STATION_RUN.sts);
+ }
+
+ private boolean isManualWrkStatus(Long wrkSts) {
+ return Objects.equals(wrkSts, WrkStsType.INBOUND_MANUAL.sts)
+ || Objects.equals(wrkSts, WrkStsType.OUTBOUND_MANUAL.sts)
+ || Objects.equals(wrkSts, WrkStsType.LOC_MOVE_MANUAL.sts);
+ }
+
+ private void ensureCacheLoaded() {
+ if (loadedFromRedis) {
+ return;
+ }
+ synchronized (this) {
+ if (loadedFromRedis) {
+ return;
+ }
+ Map<Object, Object> cacheMap = redisUtil == null ? null : redisUtil.hmget(TRACE_CACHE_KEY);
+ if (cacheMap != null && !cacheMap.isEmpty()) {
+ long now = System.currentTimeMillis();
+ for (Map.Entry<Object, Object> entry : cacheMap.entrySet()) {
+ Integer taskNo = parseTaskNo(entry.getKey());
+ TraceSnapshot snapshot = parseSnapshot(entry.getValue(), taskNo);
+ if (snapshot == null || snapshot.getTaskNo() == null) {
+ if (taskNo != null) {
+ removePersistedState(taskNo);
+ }
+ continue;
+ }
+ if (snapshot.isExpired(now)) {
+ removePersistedState(snapshot.getTaskNo());
+ continue;
+ }
+ taskStateMap.put(snapshot.getTaskNo(), TraceTaskState.fromSnapshot(snapshot));
+ }
+ }
+ loadedFromRedis = true;
+ }
+ }
+
+ private void persistState(TraceTaskState state) {
+ if (state == null || state.taskNo == null || redisUtil == null) {
+ return;
+ }
+ redisUtil.hset(TRACE_CACHE_KEY, String.valueOf(state.taskNo), JSON.toJSONString(state.toSnapshot()));
+ }
+
+ private void removePersistedState(Integer taskNo) {
+ if (taskNo == null || redisUtil == null) {
+ return;
+ }
+ redisUtil.hdel(TRACE_CACHE_KEY, String.valueOf(taskNo));
+ }
+
+ private Integer parseTaskNo(Object value) {
+ if (value == null) {
+ return null;
+ }
+ try {
+ return Integer.parseInt(String.valueOf(value));
+ } catch (Exception e) {
+ return null;
+ }
+ }
+
+ private TraceSnapshot parseSnapshot(Object cacheValue, Integer fallbackTaskNo) {
+ if (cacheValue == null) {
+ return null;
+ }
+ try {
+ TraceSnapshot snapshot = JSON.parseObject(String.valueOf(cacheValue), TraceSnapshot.class);
+ if (snapshot != null && snapshot.getTaskNo() == null) {
+ snapshot.setTaskNo(fallbackTaskNo);
+ }
+ return snapshot;
+ } catch (Exception e) {
+ return null;
}
}
@@ -205,6 +379,29 @@
return result;
}
+ private static List<StationTaskTraceEventVo> copyEventList(List<StationTaskTraceEventVo> source) {
+ List<StationTaskTraceEventVo> result = new ArrayList<>();
+ if (source == null) {
+ return result;
+ }
+ for (StationTaskTraceEventVo item : source) {
+ if (item == null) {
+ continue;
+ }
+ StationTaskTraceEventVo copy = new StationTaskTraceEventVo();
+ copy.setTimestamp(item.getTimestamp());
+ copy.setEventType(item.getEventType());
+ copy.setMessage(item.getMessage());
+ copy.setStatus(item.getStatus());
+ copy.setCurrentStationId(item.getCurrentStationId());
+ copy.setTargetStationId(item.getTargetStationId());
+ copy.setTraceVersion(item.getTraceVersion());
+ copy.setDetails(copyDetails(item.getDetails()));
+ result.add(copy);
+ }
+ return result;
+ }
+
private static boolean isTerminalStatus(String status) {
return STATUS_BLOCKED.equals(status)
|| STATUS_CANCELLED.equals(status)
@@ -218,6 +415,33 @@
private Integer pathOffset;
private List<Integer> fullPathStationIds = new ArrayList<>();
private boolean rerouted;
+ }
+
+ @Data
+ private static class TraceSnapshot {
+ private Integer taskNo;
+ private String threadImpl;
+ private String status;
+ private Integer traceVersion;
+ private Integer startStationId;
+ private Integer currentStationId;
+ private Integer finalTargetStationId;
+ private Integer blockedStationId;
+ private List<Integer> fullPathStationIds = new ArrayList<>();
+ private List<Integer> issuedStationIds = new ArrayList<>();
+ private List<Integer> passedStationIds = new ArrayList<>();
+ private List<Integer> pendingStationIds = new ArrayList<>();
+ private List<Integer> latestIssuedSegmentPath = new ArrayList<>();
+ private List<StationTaskTraceSegmentVo> segmentList = new ArrayList<>();
+ private Integer issuedSegmentCount;
+ private Integer totalSegmentCount;
+ private Long updatedAt;
+ private Long terminalExpireAt;
+ private List<StationTaskTraceEventVo> events = new ArrayList<>();
+
+ private boolean isExpired(long now) {
+ return terminalExpireAt != null && terminalExpireAt <= now;
+ }
}
private static class TraceTaskState {
@@ -367,6 +591,9 @@
}
this.status = terminalStatus;
this.blockedStationId = blockedStationId;
+ if (shouldClearPathOnTerminal(terminalStatus)) {
+ clearPathState();
+ }
this.updatedAt = System.currentTimeMillis();
this.terminalExpireAt = this.updatedAt + TERMINAL_KEEP_MS;
@@ -378,8 +605,28 @@
appendEvent(eventType, message, nextDetails);
}
+ private void clearPathState() {
+ this.fullPathStationIds = new ArrayList<>();
+ this.issuedStationIds = new ArrayList<>();
+ this.passedStationIds = new ArrayList<>();
+ this.pendingStationIds = new ArrayList<>();
+ this.latestIssuedSegmentPath = new ArrayList<>();
+ this.segmentList = new ArrayList<>();
+ this.issuedSegmentCount = 0;
+ this.totalSegmentCount = 0;
+ }
+
+ private boolean shouldClearPathOnTerminal(String terminalStatus) {
+ return STATUS_BLOCKED.equals(terminalStatus)
+ || STATUS_CANCELLED.equals(terminalStatus);
+ }
+
private synchronized boolean shouldRemove(long now) {
return terminalExpireAt != null && terminalExpireAt <= now;
+ }
+
+ private synchronized boolean isTerminal() {
+ return isTerminalStatus(this.status);
}
private synchronized StationTaskTraceVo toVo() {
@@ -401,10 +648,58 @@
vo.setIssuedSegmentCount(issuedSegmentCount);
vo.setTotalSegmentCount(totalSegmentCount);
vo.setUpdatedAt(updatedAt);
- vo.setEvents(new ArrayList<>(events));
+ vo.setEvents(copyEventList(events));
return vo;
}
+ private synchronized TraceSnapshot toSnapshot() {
+ TraceSnapshot snapshot = new TraceSnapshot();
+ snapshot.setTaskNo(taskNo);
+ snapshot.setThreadImpl(threadImpl);
+ snapshot.setStatus(status);
+ snapshot.setTraceVersion(traceVersion);
+ snapshot.setStartStationId(startStationId);
+ snapshot.setCurrentStationId(currentStationId);
+ snapshot.setFinalTargetStationId(finalTargetStationId);
+ snapshot.setBlockedStationId(blockedStationId);
+ snapshot.setFullPathStationIds(copyIntegerList(fullPathStationIds));
+ snapshot.setIssuedStationIds(copyIntegerList(issuedStationIds));
+ snapshot.setPassedStationIds(copyIntegerList(passedStationIds));
+ snapshot.setPendingStationIds(copyIntegerList(pendingStationIds));
+ snapshot.setLatestIssuedSegmentPath(copyIntegerList(latestIssuedSegmentPath));
+ snapshot.setSegmentList(copySegmentList(segmentList));
+ snapshot.setIssuedSegmentCount(issuedSegmentCount);
+ snapshot.setTotalSegmentCount(totalSegmentCount);
+ snapshot.setUpdatedAt(updatedAt);
+ snapshot.setTerminalExpireAt(terminalExpireAt);
+ snapshot.setEvents(copyEventList(events));
+ return snapshot;
+ }
+
+ private static TraceTaskState fromSnapshot(TraceSnapshot snapshot) {
+ TraceTaskState state = new TraceTaskState(snapshot.getTaskNo());
+ state.threadImpl = snapshot.getThreadImpl();
+ state.status = snapshot.getStatus() == null ? STATUS_WAITING : snapshot.getStatus();
+ state.traceVersion = snapshot.getTraceVersion() == null ? 0 : snapshot.getTraceVersion();
+ state.startStationId = snapshot.getStartStationId();
+ state.currentStationId = snapshot.getCurrentStationId();
+ state.finalTargetStationId = snapshot.getFinalTargetStationId();
+ state.blockedStationId = snapshot.getBlockedStationId();
+ state.fullPathStationIds = copyIntegerList(snapshot.getFullPathStationIds());
+ state.issuedStationIds = copyIntegerList(snapshot.getIssuedStationIds());
+ state.passedStationIds = copyIntegerList(snapshot.getPassedStationIds());
+ state.pendingStationIds = copyIntegerList(snapshot.getPendingStationIds());
+ state.latestIssuedSegmentPath = copyIntegerList(snapshot.getLatestIssuedSegmentPath());
+ state.segmentList = copySegmentList(snapshot.getSegmentList());
+ state.issuedSegmentCount = snapshot.getIssuedSegmentCount() == null ? 0 : snapshot.getIssuedSegmentCount();
+ state.totalSegmentCount = snapshot.getTotalSegmentCount() == null ? state.segmentList.size() : snapshot.getTotalSegmentCount();
+ state.updatedAt = snapshot.getUpdatedAt() == null ? System.currentTimeMillis() : snapshot.getUpdatedAt();
+ state.terminalExpireAt = snapshot.getTerminalExpireAt();
+ state.events.clear();
+ state.events.addAll(copyEventList(snapshot.getEvents()));
+ return state;
+ }
+
private List<Integer> buildFullPathForRegistration(boolean rerouted,
Integer planCurrentStationId,
List<Integer> localPathStationIds) {
--
Gitblit v1.9.1