Просмотр исходного кода

feat(websocket): 添加设备状态变更事件发布器和工单提醒推送功能

- 新增 DeviceStatusChangePublisher 类用于发布设备状态变更事件到 Redis
- 实现事务提交后发布机制,确保事件只在事务成功后发送
- 添加 WorkOrderPushVO 数据传输对象用于工单提醒
- 实现 WorkOrderWebSocketPushTask 定时推送待处理及超时工单提醒
- 支持工单可见性检查和用户权限验证
- 添加指纹机制避免重复推送相同的工单数据
- 集成 WebSocket 连接管理和会话清理功能
林仔 4 дней назад
Родитель
Сommit
0d61fff386

+ 153 - 0
pipe-network-service/zksy-admin/src/main/java/com/zksy/web/websocket/WorkOrderWebSocketPushTask.java

@@ -0,0 +1,153 @@
+package com.zksy.web.websocket;
+
+import com.alibaba.fastjson.JSON;
+import com.alibaba.fastjson.JSONObject;
+import com.alibaba.fastjson.serializer.SerializerFeature;
+import com.zksy.base.domain.vo.WorkOrderPushVO;
+import com.zksy.base.service.BusinessWorkOrderService;
+import com.zksy.common.core.domain.model.LoginUser;
+import com.zksy.common.utils.DateUtils;
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.beans.factory.annotation.Value;
+import org.springframework.scheduling.annotation.Scheduled;
+import org.springframework.stereotype.Component;
+import org.springframework.web.socket.WebSocketSession;
+
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.Date;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.Set;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.stream.Collectors;
+
+/**
+ * 待处理及超时工单 WebSocket 主动推送任务。
+ */
+@Slf4j
+@Component
+public class WorkOrderWebSocketPushTask {
+
+    public static final String PUSH_TYPE = "workOrderReminderPush";
+
+    @Autowired
+    private BusinessWorkOrderService businessWorkOrderService;
+
+    @Autowired
+    private DeviceStatusWebSocketHandler webSocketHandler;
+
+    @Value("${websocket.work-order.max-items:100}")
+    private int maxItems;
+
+    /** sessionId -> 上次已推送的业务快照指纹。 */
+    private final Map<String, String> lastFingerprints = new ConcurrentHashMap<>();
+
+    @Scheduled(
+            fixedDelayString = "${websocket.work-order.scan-interval-ms:60000}",
+            initialDelayString = "${websocket.work-order.initial-delay-ms:10000}")
+    public void pushPendingOrOverdueWorkOrders() {
+        Map<String, WebSocketSession> sessions = webSocketHandler.getSessionSnapshot();
+        cleanupDisconnectedSessions(sessions.keySet());
+        if (sessions.isEmpty()) {
+            return;
+        }
+
+        // 先完成会话鉴权,再查库;只有确实存在需要接收工单提醒的用户时才执行扫描。
+        // 这样设备状态专用连接不会触发无意义的工单全表查询,同时每次扫描只校验一次 Token。
+        Map<String, LoginUser> authenticatedUsers = new HashMap<>();
+        for (Map.Entry<String, WebSocketSession> entry : sessions.entrySet()) {
+            LoginUser loginUser = webSocketHandler.getAuthenticatedUser(entry.getValue());
+            if (loginUser == null) {
+                lastFingerprints.remove(entry.getKey());
+            } else {
+                authenticatedUsers.put(entry.getKey(), loginUser);
+            }
+        }
+        if (authenticatedUsers.isEmpty()) {
+            return;
+        }
+
+        Date now = new Date();
+        List<WorkOrderPushVO> candidates = businessWorkOrderService.selectPendingOrOverdueWorkOrders(now);
+        if (candidates == null) {
+            candidates = Collections.emptyList();
+        }
+        for (Map.Entry<String, LoginUser> userEntry : authenticatedUsers.entrySet()) {
+            String sessionId = userEntry.getKey();
+            WebSocketSession session = sessions.get(sessionId);
+            LoginUser loginUser = userEntry.getValue();
+
+            List<WorkOrderPushVO> visibleItems = candidates.stream()
+                    .filter(item -> isVisibleTo(item, loginUser))
+                    .collect(Collectors.toList());
+            pushSnapshot(sessionId, session, loginUser, visibleItems, now);
+        }
+    }
+
+    private void pushSnapshot(String sessionId, WebSocketSession session,
+                              LoginUser loginUser, List<WorkOrderPushVO> visibleItems, Date now) {
+        int total = visibleItems.size();
+        int limit = Math.max(1, maxItems);
+        List<WorkOrderPushVO> items = total > limit
+                ? new ArrayList<>(visibleItems.subList(0, limit)) : visibleItems;
+        String fingerprint = buildFingerprint(visibleItems, loginUser);
+        String previous = lastFingerprints.get(sessionId);
+
+        if (fingerprint.equals(previous) || (previous == null && visibleItems.isEmpty())) {
+            return;
+        }
+
+        JSONObject data = new JSONObject();
+        data.put("items", items);
+        data.put("total", total);
+        data.put("truncated", total > items.size());
+        data.put("serverTime", DateUtils.parseDateToStr(DateUtils.YYYY_MM_DD_HH_MM_SS, now));
+
+        JSONObject response = new JSONObject();
+        response.put("type", PUSH_TYPE);
+        response.put("success", true);
+        response.put("data", data);
+        String payload = JSON.toJSONString(response, SerializerFeature.WriteDateUseDateFormat);
+        if (webSocketHandler.sendText(session, payload)) {
+            lastFingerprints.put(sessionId, fingerprint);
+            log.info("工单提醒推送完成: sessionId={}, total={}, sent={}",
+                    sessionId, total, items.size());
+        }
+    }
+
+    boolean isVisibleTo(WorkOrderPushVO item, LoginUser loginUser) {
+        return DeviceStatusWebSocketHandler.isWorkOrderVisibleTo(item, loginUser);
+    }
+
+    private String buildFingerprint(List<WorkOrderPushVO> items, LoginUser loginUser) {
+        if (items.isEmpty()) {
+            return "EMPTY:" + userFingerprint(loginUser);
+        }
+        return userFingerprint(loginUser) + ":" + items.stream()
+                .map(item -> String.valueOf(item.getOrderId()) + ':'
+                        + item.getOrderStatus() + ':'
+                        + time(item.getEffectiveDeadline()) + ':'
+                        + time(item.getUpdateTime()) + ':'
+                        + item.getDeptId() + ':'
+                        + item.getReceiveUser())
+                .collect(Collectors.joining("|"));
+    }
+
+    private String userFingerprint(LoginUser loginUser) {
+        if (loginUser == null) {
+            return "ANONYMOUS";
+        }
+        return String.valueOf(loginUser.getUserId()) + ":" + String.valueOf(loginUser.getDeptId());
+    }
+
+    private long time(Date value) {
+        return value != null ? value.getTime() : 0L;
+    }
+
+    private void cleanupDisconnectedSessions(Set<String> activeSessionIds) {
+        lastFingerprints.keySet().retainAll(activeSessionIds);
+    }
+}

+ 74 - 0
pipe-network-service/zksy-system/src/main/java/com/zksy/base/domain/vo/WorkOrderPushVO.java

@@ -0,0 +1,74 @@
+package com.zksy.base.domain.vo;
+
+import com.fasterxml.jackson.annotation.JsonFormat;
+import io.swagger.annotations.ApiModel;
+import io.swagger.annotations.ApiModelProperty;
+import lombok.Data;
+
+import java.io.Serializable;
+import java.util.Date;
+
+/**
+ * WebSocket 工单提醒出参。
+ * 只返回提醒场景需要的字段,避免把工单完整维护字段广播给客户端。
+ */
+@Data
+@ApiModel(value = "WorkOrderPushVO", description = "WebSocket 工单提醒")
+public class WorkOrderPushVO implements Serializable {
+
+    private static final long serialVersionUID = 1L;
+
+    @ApiModelProperty("工单ID")
+    private Long orderId;
+
+    @ApiModelProperty("工单编号")
+    private String orderNo;
+
+    @ApiModelProperty("工单状态:1-待派单 2-待接单 3-已接单 4-处理中 5-待验收 8-延期审核 9-已延期")
+    private Integer orderStatus;
+
+    @ApiModelProperty("工单优先级")
+    private Integer orderLevel;
+
+    @ApiModelProperty("所属部门ID")
+    private Long deptId;
+
+    @ApiModelProperty("接单人ID,多个使用逗号分隔")
+    private String receiveUser;
+
+    @ApiModelProperty("设备ID")
+    private String deviceId;
+
+    @ApiModelProperty("设备编码")
+    private String deviceCode;
+
+    @ApiModelProperty("工单描述")
+    private String orderDesc;
+
+    @JsonFormat(locale = "zh", pattern = "yyyy-MM-dd HH:mm:ss", timezone = "GMT+8")
+    @ApiModelProperty("计划完成时间")
+    private Date planFinishTime;
+
+    @JsonFormat(locale = "zh", pattern = "yyyy-MM-dd HH:mm:ss", timezone = "GMT+8")
+    @ApiModelProperty("延期截止时间")
+    private Date delayedDeadline;
+
+    @JsonFormat(locale = "zh", pattern = "yyyy-MM-dd HH:mm:ss", timezone = "GMT+8")
+    @ApiModelProperty("当前生效的截止时间")
+    private Date effectiveDeadline;
+
+    @ApiModelProperty("超时类型:PENDING-待处理,PLAN_FINISH_TIME-计划完成时间,DELAYED_DEADLINE-延期截止时间")
+    private String overdueType;
+
+    @ApiModelProperty("是否已经超过生效截止时间")
+    private Boolean overdue;
+
+    @ApiModelProperty("超过截止时间的分钟数")
+    private Long overdueMinutes;
+
+    @JsonFormat(locale = "zh", pattern = "yyyy-MM-dd HH:mm:ss", timezone = "GMT+8")
+    private Date createTime;
+
+    @JsonFormat(locale = "zh", pattern = "yyyy-MM-dd HH:mm:ss", timezone = "GMT+8")
+    private Date updateTime;
+}

+ 80 - 0
pipe-network-service/zksy-system/src/main/java/com/zksy/base/event/DeviceStatusChangePublisher.java

@@ -0,0 +1,80 @@
+package com.zksy.base.event;
+
+import com.fasterxml.jackson.core.JsonProcessingException;
+import com.fasterxml.jackson.databind.ObjectMapper;
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.data.redis.core.StringRedisTemplate;
+import org.springframework.stereotype.Component;
+import org.springframework.transaction.support.TransactionSynchronization;
+import org.springframework.transaction.support.TransactionSynchronizationManager;
+
+import java.time.LocalDateTime;
+import java.util.LinkedHashMap;
+import java.util.Map;
+import java.util.UUID;
+
+/**
+ * 设备状态变更事件发布器。
+ *
+ * <p>状态写库成功后发布 Redis 事件,由管网主服务的 WebSocket 监听器
+ * 查询最新状态并推送给订阅该设备的客户端。事务存在时严格在提交成功后发布,
+ * 避免客户端收到已回滚的数据。</p>
+ */
+@Slf4j
+@Component
+public class DeviceStatusChangePublisher {
+
+    public static final String CHANNEL = "device:status:change";
+
+    @Autowired
+    private StringRedisTemplate redisTemplate;
+
+    @Autowired
+    private ObjectMapper objectMapper;
+
+    public void publishAfterCommit(String deviceCode, String equipmentId,
+                                   Integer currentStatus, Integer alarmStatus,
+                                   Integer onlineStatus, String source) {
+        if (deviceCode == null || deviceCode.trim().isEmpty()) {
+            log.warn("设备状态变更事件未发布:deviceCode为空, equipmentId={}, source={}", equipmentId, source);
+            return;
+        }
+
+        Runnable publishAction = () -> publish(deviceCode, equipmentId,
+                currentStatus, alarmStatus, onlineStatus, source);
+        if (TransactionSynchronizationManager.isSynchronizationActive()) {
+            TransactionSynchronizationManager.registerSynchronization(new TransactionSynchronization() {
+                @Override
+                public void afterCommit() {
+                    publishAction.run();
+                }
+            });
+        } else {
+            publishAction.run();
+        }
+    }
+
+    private void publish(String deviceCode, String equipmentId,
+                         Integer currentStatus, Integer alarmStatus,
+                         Integer onlineStatus, String source) {
+        Map<String, Object> event = new LinkedHashMap<>();
+        event.put("eventId", UUID.randomUUID().toString());
+        event.put("deviceCode", deviceCode);
+        event.put("equipmentId", equipmentId);
+        event.put("currentStatus", currentStatus);
+        event.put("alarmStatus", alarmStatus);
+        event.put("onlineStatus", onlineStatus);
+        event.put("source", source);
+        event.put("changedAt", LocalDateTime.now().toString());
+        try {
+            redisTemplate.convertAndSend(CHANNEL, objectMapper.writeValueAsString(event));
+            log.info("设备状态变更事件已发布: deviceCode={}, equipmentId={}, source={}",
+                    deviceCode, equipmentId, source);
+        } catch (JsonProcessingException e) {
+            log.error("设备状态变更事件序列化失败: deviceCode={}, source={}", deviceCode, source, e);
+        } catch (Exception e) {
+            log.error("设备状态变更事件发布失败: deviceCode={}, source={}", deviceCode, source, e);
+        }
+    }
+}

+ 73 - 0
zk-api-service/src/main/java/com/zksy/api/event/DeviceStatusChangePublisher.java

@@ -0,0 +1,73 @@
+package com.zksy.api.event;
+
+import com.fasterxml.jackson.core.JsonProcessingException;
+import com.fasterxml.jackson.databind.ObjectMapper;
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.data.redis.core.StringRedisTemplate;
+import org.springframework.stereotype.Component;
+import org.springframework.transaction.support.TransactionSynchronization;
+import org.springframework.transaction.support.TransactionSynchronizationManager;
+
+import java.time.LocalDateTime;
+import java.util.LinkedHashMap;
+import java.util.Map;
+import java.util.UUID;
+
+/** 发布设备状态变更事件,供 pipe-network-service WebSocket 转发。 */
+@Slf4j
+@Component
+public class DeviceStatusChangePublisher {
+
+    public static final String CHANNEL = "device:status:change";
+
+    @Autowired
+    private StringRedisTemplate redisTemplate;
+
+    @Autowired
+    private ObjectMapper objectMapper;
+
+    public void publishAfterCommit(String deviceCode, String equipmentId,
+                                   Integer currentStatus, Integer alarmStatus,
+                                   Integer onlineStatus, String source) {
+        if (deviceCode == null || deviceCode.trim().isEmpty()) {
+            log.warn("设备状态变更事件未发布:deviceCode为空, equipmentId={}, source={}", equipmentId, source);
+            return;
+        }
+        Runnable publishAction = () -> publish(deviceCode, equipmentId,
+                currentStatus, alarmStatus, onlineStatus, source);
+        if (TransactionSynchronizationManager.isSynchronizationActive()) {
+            TransactionSynchronizationManager.registerSynchronization(new TransactionSynchronization() {
+                @Override
+                public void afterCommit() {
+                    publishAction.run();
+                }
+            });
+        } else {
+            publishAction.run();
+        }
+    }
+
+    private void publish(String deviceCode, String equipmentId,
+                         Integer currentStatus, Integer alarmStatus,
+                         Integer onlineStatus, String source) {
+        Map<String, Object> event = new LinkedHashMap<>();
+        event.put("eventId", UUID.randomUUID().toString());
+        event.put("deviceCode", deviceCode);
+        event.put("equipmentId", equipmentId);
+        event.put("currentStatus", currentStatus);
+        event.put("alarmStatus", alarmStatus);
+        event.put("onlineStatus", onlineStatus);
+        event.put("source", source);
+        event.put("changedAt", LocalDateTime.now().toString());
+        try {
+            redisTemplate.convertAndSend(CHANNEL, objectMapper.writeValueAsString(event));
+            log.info("设备状态变更事件已发布: deviceCode={}, equipmentId={}, source={}",
+                    deviceCode, equipmentId, source);
+        } catch (JsonProcessingException e) {
+            log.error("设备状态变更事件序列化失败: deviceCode={}, source={}", deviceCode, source, e);
+        } catch (Exception e) {
+            log.error("设备状态变更事件发布失败: deviceCode={}, source={}", deviceCode, source, e);
+        }
+    }
+}