Browse Source

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

- 集成 DeviceStatusChangePublisher 实现设备状态变更事件发布
- 新增 WebSocket 工单提醒推送功能,支持鉴权和实时工单状态查询
- 实现 Redis 消息监听器统一处理设备状态变更推送
- 优化设备状态更新逻辑,添加事务管理和状态变更检测
- 移除工单验收对设备报警状态的直接重置逻辑
- 添加 WebSocket 连接认证和会话管理机制
- 实现工单超时和待处理状态的定时扫描推送
- 统一 WebSocket 消息发送机制防止并发写冲突
林仔 4 days ago
parent
commit
ea44940fc1

+ 16 - 6
manhole-service/src/main/java/com/zksy/manhole/service/impl/BaseDevicesManholeServiceImpl.java

@@ -3,7 +3,7 @@ package com.zksy.manhole.service.impl;
 import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
 import com.baomidou.mybatisplus.extension.service.impl.ServiceImpl;
 
-import com.zksy.base.domain.EquipmentStatus;
+import com.zksy.api.event.DeviceStatusChangePublisher;
 import com.zksy.common.exception.ServiceException;
 import com.zksy.manhole.domain.BaseDevicesManhole;
 import com.zksy.manhole.domain.EquipmentStatusManhole;
@@ -12,6 +12,9 @@ import com.zksy.manhole.mapper.EquipmentStatusManholeMapper;
 import com.zksy.manhole.service.BaseDevicesManholeService;
 import org.springframework.beans.factory.annotation.Autowired;
 import org.springframework.stereotype.Service;
+import org.springframework.transaction.annotation.Transactional;
+
+import java.util.Objects;
 
 /**
 * @author Administrator
@@ -27,7 +30,11 @@ public class BaseDevicesManholeServiceImpl extends ServiceImpl<BaseDevicesManhol
     @Autowired
     private EquipmentStatusManholeMapper equipmentStatusBaseDevicesMapper;
 
+    @Autowired
+    private DeviceStatusChangePublisher deviceStatusChangePublisher;
+
     @Override
+    @Transactional(rollbackFor = Exception.class)
     public void getByDeviceNumberStatus(String deviceNumber,Integer queryStatus,Integer updateStatus) {
         LambdaQueryWrapper<BaseDevicesManhole> wrapper = new LambdaQueryWrapper<>();
         //设备编号
@@ -47,12 +54,15 @@ public class BaseDevicesManholeServiceImpl extends ServiceImpl<BaseDevicesManhol
 
         //修改状态表
         if(equipmentStatus != null){
+            boolean changed = !Objects.equals(equipmentStatus.getOnlineStatus(), updateStatus);
             equipmentStatus.setOnlineStatus(updateStatus);
-            equipmentStatusBaseDevicesMapper.updateById(equipmentStatus);
+            int affected = equipmentStatusBaseDevicesMapper.updateById(equipmentStatus);
+            if (affected > 0 && changed) {
+                deviceStatusChangePublisher.publishAfterCommit(device.getEquipmentCode(),
+                        deviceEquipmentId, equipmentStatus.getCurrentStatus(),
+                        equipmentStatus.getAlarmStatus(), equipmentStatus.getOnlineStatus(),
+                        "manhole:updateOnlineStatus");
+            }
         }
     }
 }
-
-
-
-

+ 20 - 16
pipe-network-service/zksy-admin/src/main/java/com/zksy/web/websocket/DeviceStatusRedisListener.java

@@ -8,10 +8,9 @@ import org.springframework.beans.factory.annotation.Autowired;
 import org.springframework.data.redis.connection.Message;
 import org.springframework.data.redis.connection.MessageListener;
 import org.springframework.stereotype.Component;
-import org.springframework.web.socket.TextMessage;
 import org.springframework.web.socket.WebSocketSession;
 
-import java.io.IOException;
+import java.nio.charset.StandardCharsets;
 import java.util.Set;
 
 /**
@@ -32,6 +31,9 @@ public class DeviceStatusRedisListener implements MessageListener {
     @Autowired
     private EquipmentStatusService equipmentStatusService;
 
+    @Autowired
+    private DeviceStatusWebSocketHandler webSocketHandler;
+
     /**
      * Redis 频道名称
      */
@@ -39,7 +41,7 @@ public class DeviceStatusRedisListener implements MessageListener {
 
     @Override
     public void onMessage(Message message, byte[] pattern) {
-        String body = new String(message.getBody());
+        String body = new String(message.getBody(), StandardCharsets.UTF_8);
         log.debug("收到 Redis 设备状态变更消息: channel={}, body={}", CHANNEL, body);
 
         try {
@@ -66,31 +68,32 @@ public class DeviceStatusRedisListener implements MessageListener {
             }
 
             // 构建推送消息
+            JSONObject data = JSONObject.parseObject(JSONObject.toJSONString(status));
+            // 客户端可能同时订阅多个设备,必须在推送中携带设备编码以便路由状态。
+            data.put("deviceCode", deviceCode);
+            if (msgObj.containsKey("eventId")) {
+                data.put("eventId", msgObj.getString("eventId"));
+            }
             JSONObject pushMsg = new JSONObject();
             pushMsg.put("type", "deviceStatusPush");
-            pushMsg.put("data", status);
+            pushMsg.put("deviceCode", deviceCode);
+            pushMsg.put("data", data);
 
             String pushPayload = pushMsg.toJSONString();
-            TextMessage textMessage = new TextMessage(pushPayload);
 
             // 推送到所有订阅了该设备的会话
             int successCount = 0;
             int failCount = 0;
             for (String sessionId : sessionIds) {
-                try {
-                    WebSocketSession session = DeviceStatusWebSocketHandler.getSession(sessionId);
-                    if (session != null && session.isOpen()) {
-                        synchronized (session) {
-                            session.sendMessage(textMessage);
-                        }
+                WebSocketSession session = DeviceStatusWebSocketHandler.getSession(sessionId);
+                if (session != null && session.isOpen()) {
+                    if (webSocketHandler.sendText(session, pushPayload)) {
                         successCount++;
                     } else {
-                        // 会话已失效,清理订阅
-                        subscriptionManager.removeSession(sessionId);
                         failCount++;
                     }
-                } catch (IOException e) {
-                    log.error("推送设备状态失败: sessionId={}, deviceCode={}", sessionId, deviceCode, e);
+                } else {
+                    // 会话已失效,清理订阅
                     subscriptionManager.removeSession(sessionId);
                     failCount++;
                 }
@@ -103,4 +106,5 @@ public class DeviceStatusRedisListener implements MessageListener {
             log.error("处理 Redis 设备状态变更消息异常: body={}", body, e);
         }
     }
-}
+
+}

+ 173 - 8
pipe-network-service/zksy-admin/src/main/java/com/zksy/web/websocket/DeviceStatusWebSocketHandler.java

@@ -4,8 +4,14 @@ import com.alibaba.fastjson.JSONArray;
 import com.alibaba.fastjson.JSONObject;
 import com.zksy.base.domain.EquipmentStatus;
 import com.zksy.base.service.EquipmentStatusService;
+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 com.zksy.framework.web.service.TokenService;
 import lombok.extern.slf4j.Slf4j;
 import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.beans.factory.annotation.Value;
 import org.springframework.stereotype.Component;
 import org.springframework.web.socket.CloseStatus;
 import org.springframework.web.socket.TextMessage;
@@ -14,8 +20,11 @@ import org.springframework.web.socket.handler.TextWebSocketHandler;
 
 import java.io.IOException;
 import java.util.HashSet;
+import java.util.HashMap;
+import java.util.Collections;
 import java.util.List;
 import java.util.Map;
+import java.util.stream.Collectors;
 import java.util.Set;
 import java.util.concurrent.ConcurrentHashMap;
 
@@ -24,14 +33,20 @@ import java.util.concurrent.ConcurrentHashMap;
  * <p>
  * 支持的消息类型(通过 JSON 中的 type 字段区分):
  * <ul>
+ *   <li><b>authenticate</b> — 使用登录 JWT 鉴权(工单推送必需)</li>
  *   <li><b>subscribe</b> — 批量订阅设备,后端状态变更时主动推送</li>
  *   <li><b>unsubscribe</b> — 取消订阅设备</li>
  *   <li><b>deviceStatus</b> — 单次查询设备状态(兼容旧接口)</li>
+ *   <li><b>workOrderReminder</b> — 查询当前用户可见的待处理/超时工单</li>
  *   <li><b>ping</b> — 心跳检测</li>
  * </ul>
  *
  * <h3>协议示例</h3>
  * <pre>
+ * // 工单鉴权(连接后首帧发送)
+ * → {"type":"authenticate","token":"Bearer &lt;JWT&gt;"}
+ * ← {"type":"authenticated","success":true,"data":{"userId":10,"deptId":100}}
+ *
  * // 订阅
  * → {"type":"subscribe","deviceCodes":["DEV001","DEV002"]}
  * ← {"type":"subscribed","deviceCodes":["DEV001","DEV002"]}
@@ -49,7 +64,10 @@ import java.util.concurrent.ConcurrentHashMap;
  * ← {"type":"deviceStatus","success":true,"data":{...}}
  *
  * // 设备状态变更推送(服务端主动推送)
- * ← {"type":"deviceStatusPush","data":{...}}
+ * ← {"type":"deviceStatusPush","deviceCode":"DEV001","data":{...}}
+ *
+ * // 工单提醒推送(服务端每分钟扫描,有变化才推送)
+ * ← {"type":"workOrderReminderPush","success":true,"data":{"items":[...],"total":1}}
  * </pre>
  *
  * @author zksy
@@ -64,6 +82,18 @@ public class DeviceStatusWebSocketHandler extends TextWebSocketHandler {
     @Autowired
     private SubscriptionManager subscriptionManager;
 
+    @Autowired
+    private TokenService tokenService;
+
+    @Autowired
+    private BusinessWorkOrderService businessWorkOrderService;
+
+    @Value("${websocket.work-order.max-items:100}")
+    private int workOrderMaxItems;
+
+    static final String WORK_ORDER_LOGIN_USER_ATTRIBUTE = "workOrderLoginUser";
+    static final String WORK_ORDER_TOKEN_ATTRIBUTE = "workOrderToken";
+
     /**
      * 维护所有在线连接,key 为 sessionId
      * 供 RedisListener 跨线程查找会话
@@ -77,6 +107,42 @@ public class DeviceStatusWebSocketHandler extends TextWebSocketHandler {
         return SESSION_MAP.get(sessionId);
     }
 
+    /**
+     * 获取当前在线会话快照,供定时推送任务安全遍历。
+     */
+    public Map<String, WebSocketSession> getSessionSnapshot() {
+        return new HashMap<>(SESSION_MAP);
+    }
+
+    public LoginUser getAuthenticatedUser(WebSocketSession session) {
+        if (session == null) {
+            return null;
+        }
+        Map<String, Object> attributes = session.getAttributes();
+        if (attributes == null) {
+            return null;
+        }
+        Object token;
+        Object value;
+        synchronized (session) {
+            token = attributes.get(WORK_ORDER_TOKEN_ATTRIBUTE);
+            value = attributes.get(WORK_ORDER_LOGIN_USER_ATTRIBUTE);
+        }
+        if (token instanceof String && !((String) token).trim().isEmpty()) {
+            LoginUser loginUser = tokenService.getLoginUserByToken((String) token);
+            if (loginUser == null) {
+                clearAuthentication(session);
+                return null;
+            }
+            tokenService.verifyToken(loginUser);
+            synchronized (session) {
+                attributes.put(WORK_ORDER_LOGIN_USER_ATTRIBUTE, loginUser);
+            }
+            return loginUser;
+        }
+        return value instanceof LoginUser ? (LoginUser) value : null;
+    }
+
     @Override
     public void afterConnectionEstablished(WebSocketSession session) {
         String sessionId = session.getId();
@@ -87,11 +153,11 @@ public class DeviceStatusWebSocketHandler extends TextWebSocketHandler {
     @Override
     protected void handleTextMessage(WebSocketSession session, TextMessage message) throws IOException {
         String payload = message.getPayload();
-        log.debug("收到 WebSocket 消息: sessionId={}, payload={}", session.getId(), payload);
 
         try {
             JSONObject request = JSONObject.parseObject(payload);
             String type = request.getString("type");
+            log.debug("收到 WebSocket 消息: sessionId={}, type={}", session.getId(), type);
 
             if (type == null || type.isEmpty()) {
                 sendError(session, "消息缺少 type 字段");
@@ -100,6 +166,12 @@ public class DeviceStatusWebSocketHandler extends TextWebSocketHandler {
 
             // === 根据 type 分发到不同的业务处理 ===
             switch (type) {
+                case "authenticate":
+                    handleAuthenticate(session, request);
+                    break;
+                case "workOrderReminder":
+                    handleWorkOrderReminder(session);
+                    break;
                 //订阅
                 case "subscribe":
                     handleSubscribe(session, request);
@@ -126,6 +198,65 @@ public class DeviceStatusWebSocketHandler extends TextWebSocketHandler {
         }
     }
 
+    /**
+     * WebSocket 工单推送鉴权。浏览器 WebSocket 无法设置 Authorization 请求头,
+     * 因此客户端连接成功后需先发送 {"type":"authenticate","token":"..."}。
+     */
+    private void handleAuthenticate(WebSocketSession session, JSONObject request) throws IOException {
+        String token = request.getString("token");
+        LoginUser loginUser = tokenService.getLoginUserByToken(token);
+        if (loginUser == null) {
+            clearAuthentication(session);
+            sendResponse(session, "authenticated", false, null, "token 无效或已过期");
+            return;
+        }
+
+        tokenService.verifyToken(loginUser);
+        synchronized (session) {
+            session.getAttributes().put(WORK_ORDER_LOGIN_USER_ATTRIBUTE, loginUser);
+            session.getAttributes().put(WORK_ORDER_TOKEN_ATTRIBUTE, token);
+        }
+        JSONObject data = new JSONObject();
+        data.put("userId", loginUser.getUserId());
+        data.put("deptId", loginUser.getDeptId());
+        sendResponse(session, "authenticated", true, data, null);
+        log.info("WebSocket 工单推送鉴权成功: sessionId={}, userId={}, deptId={}",
+                session.getId(), loginUser.getUserId(), loginUser.getDeptId());
+    }
+
+    private void handleWorkOrderReminder(WebSocketSession session) throws IOException {
+        LoginUser loginUser = getAuthenticatedUser(session);
+        if (loginUser == null) {
+            sendResponse(session, "workOrderReminder", false, null, "请先完成 WebSocket 鉴权");
+            return;
+        }
+        List<WorkOrderPushVO> candidates = businessWorkOrderService.selectPendingOrOverdueWorkOrders(new java.util.Date());
+        if (candidates == null) {
+            candidates = Collections.emptyList();
+        }
+        List<WorkOrderPushVO> visibleItems = candidates
+                .stream()
+                .filter(item -> isWorkOrderVisibleTo(item, loginUser))
+                .collect(Collectors.toList());
+        int limit = Math.max(1, workOrderMaxItems);
+        List<WorkOrderPushVO> items = visibleItems.stream().limit(limit).collect(Collectors.toList());
+        JSONObject data = new JSONObject();
+        data.put("items", items);
+        data.put("total", visibleItems.size());
+        data.put("truncated", visibleItems.size() > items.size());
+        data.put("serverTime", DateUtils.parseDateToStr(DateUtils.YYYY_MM_DD_HH_MM_SS, new java.util.Date()));
+        sendResponse(session, "workOrderReminder", true, data, null);
+    }
+
+    static boolean isWorkOrderVisibleTo(WorkOrderPushVO item, LoginUser loginUser) {
+        if (item == null || loginUser == null) return false;
+        boolean sameDepartment = loginUser.getDeptId() != null && loginUser.getDeptId().equals(item.getDeptId());
+        boolean receiver = item.getReceiveUser() != null && loginUser.getUserId() != null
+                && java.util.Arrays.stream(item.getReceiveUser().split(","))
+                .map(String::trim).anyMatch(String.valueOf(loginUser.getUserId())::equals);
+        return sameDepartment || receiver;
+    }
+
     /**
      * 处理批量订阅
      */
@@ -154,7 +285,7 @@ public class DeviceStatusWebSocketHandler extends TextWebSocketHandler {
         JSONObject response = new JSONObject();
         response.put("type", "subscribed");
         response.put("deviceCodes", subscribed);
-        session.sendMessage(new TextMessage(response.toJSONString()));
+        sendText(session, response.toJSONString());
 
         log.info("订阅成功: sessionId={}, 订阅设备={}", session.getId(), subscribed);
     }
@@ -188,7 +319,7 @@ public class DeviceStatusWebSocketHandler extends TextWebSocketHandler {
         response.put("type", "unsubscribed");
         response.put("deviceCodes", deviceCodes);
         response.put("remaining", remaining);
-        session.sendMessage(new TextMessage(response.toJSONString()));
+        sendText(session, response.toJSONString());
 
         log.info("取消订阅: sessionId={}, 取消设备={}, 剩余订阅={}", session.getId(), deviceCodes, remaining);
     }
@@ -220,13 +351,14 @@ public class DeviceStatusWebSocketHandler extends TextWebSocketHandler {
         JSONObject response = new JSONObject();
         response.put("type", "pong");
         response.put("timestamp", System.currentTimeMillis());
-        session.sendMessage(new TextMessage(response.toJSONString()));
+        sendText(session, response.toJSONString());
     }
 
     @Override
     public void afterConnectionClosed(WebSocketSession session, CloseStatus status) {
         String sessionId = session.getId();
         SESSION_MAP.remove(sessionId);
+        clearAuthentication(session);
         subscriptionManager.removeSession(sessionId);
         log.info("WebSocket 连接关闭: sessionId={}, closeStatus={}", sessionId, status);
     }
@@ -235,6 +367,7 @@ public class DeviceStatusWebSocketHandler extends TextWebSocketHandler {
     public void handleTransportError(WebSocketSession session, Throwable exception) {
         String sessionId = session.getId();
         SESSION_MAP.remove(sessionId);
+        clearAuthentication(session);
         subscriptionManager.removeSession(sessionId);
         log.error("WebSocket 传输异常: sessionId={}", sessionId, exception);
     }
@@ -253,7 +386,7 @@ public class DeviceStatusWebSocketHandler extends TextWebSocketHandler {
         if (msg != null) {
             response.put("msg", msg);
         }
-        session.sendMessage(new TextMessage(response.toJSONString()));
+        sendText(session, response.toJSONString());
     }
 
     /**
@@ -264,6 +397,38 @@ public class DeviceStatusWebSocketHandler extends TextWebSocketHandler {
         response.put("type", "error");
         response.put("success", false);
         response.put("msg", msg);
-        session.sendMessage(new TextMessage(response.toJSONString()));
+        sendText(session, response.toJSONString());
+    }
+
+    /**
+     * 统一串行发送,防止心跳、Redis 通知和工单定时推送并发写同一会话。
+     */
+    public boolean sendText(WebSocketSession session, String payload) {
+        if (session == null || !session.isOpen()) {
+            return false;
+        }
+        try {
+            synchronized (session) {
+                session.sendMessage(new TextMessage(payload));
+            }
+            return true;
+        } catch (Exception e) {
+            log.error("WebSocket 消息发送失败: sessionId={}", session.getId(), e);
+            SESSION_MAP.remove(session.getId());
+            subscriptionManager.removeSession(session.getId());
+            return false;
+        }
+    }
+
+    private void clearAuthentication(WebSocketSession session) {
+        if (session != null) {
+            Map<String, Object> attributes = session.getAttributes();
+            if (attributes != null) {
+                synchronized (session) {
+                    attributes.remove(WORK_ORDER_LOGIN_USER_ATTRIBUTE);
+                    attributes.remove(WORK_ORDER_TOKEN_ATTRIBUTE);
+                }
+            }
+        }
     }
-}
+}

+ 3 - 1
pipe-network-service/zksy-admin/src/main/java/com/zksy/web/websocket/RedisPubSubConfig.java

@@ -3,6 +3,7 @@ package com.zksy.web.websocket;
 import org.springframework.beans.factory.annotation.Autowired;
 import org.springframework.context.annotation.Bean;
 import org.springframework.context.annotation.Configuration;
+import org.springframework.scheduling.annotation.EnableScheduling;
 import org.springframework.data.redis.connection.RedisConnectionFactory;
 import org.springframework.data.redis.listener.ChannelTopic;
 import org.springframework.data.redis.listener.RedisMessageListenerContainer;
@@ -16,6 +17,7 @@ import org.springframework.data.redis.listener.RedisMessageListenerContainer;
  * @author zksy
  */
 @Configuration
+@EnableScheduling
 public class RedisPubSubConfig {
 
     @Autowired
@@ -39,4 +41,4 @@ public class RedisPubSubConfig {
 
         return container;
     }
-}
+}

+ 26 - 0
pipe-network-service/zksy-framework/src/main/java/com/zksy/framework/web/service/TokenService.java

@@ -82,6 +82,32 @@ public class TokenService
         return null;
     }
 
+    /**
+     * 根据 WebSocket 首帧携带的 JWT 查询登录用户。
+     * WebSocket 客户端通常无法自定义 Authorization 请求头,因此由握手建立后的
+     * authenticate 消息调用此方法完成同一套 JWT/Redis 校验。
+     */
+    public LoginUser getLoginUserByToken(String token)
+    {
+        if (StringUtils.isEmpty(token))
+        {
+            return null;
+        }
+        try
+        {
+            String rawToken = token.startsWith(Constants.TOKEN_PREFIX)
+                    ? token.substring(Constants.TOKEN_PREFIX.length()) : token;
+            Claims claims = parseToken(rawToken);
+            String uuid = (String) claims.get(Constants.LOGIN_USER_KEY);
+            return redisCache.getCacheObject(getTokenKey(uuid));
+        }
+        catch (Exception e)
+        {
+            log.warn("WebSocket token 校验失败: {}", e.getMessage());
+            return null;
+        }
+    }
+
     /**
      * 设置用户身份信息
      */

+ 8 - 0
pipe-network-service/zksy-system/src/main/java/com/zksy/base/mapper/WorkOrderMapper.java

@@ -3,6 +3,10 @@ package com.zksy.base.mapper;
 import com.baomidou.mybatisplus.core.mapper.BaseMapper;
 import com.zksy.base.domain.WorkOrder;
 import org.apache.ibatis.annotations.Mapper;
+import org.apache.ibatis.annotations.Param;
+
+import java.util.Date;
+import java.util.List;
 
 /**
 * @author Administrator
@@ -13,6 +17,10 @@ import org.apache.ibatis.annotations.Mapper;
 @Mapper
 public interface WorkOrderMapper extends BaseMapper<WorkOrder> {
 
+    /**
+     * 查询待处理或已超时的未办结工单。
+     */
+    List<WorkOrder> selectPendingOrOverdue(@Param("now") Date now);
 }
 
 

+ 7 - 0
pipe-network-service/zksy-system/src/main/java/com/zksy/base/service/BusinessWorkOrderService.java

@@ -4,6 +4,7 @@ import com.baomidou.mybatisplus.core.metadata.IPage;
 import com.baomidou.mybatisplus.extension.service.IService;
 import com.zksy.base.domain.WorkOrder;
 import com.zksy.base.domain.WorkOrderLog;
+import com.zksy.base.domain.vo.WorkOrderPushVO;
 
 import java.util.Date;
 import java.util.List;
@@ -93,4 +94,10 @@ public interface BusinessWorkOrderService extends IService<WorkOrder> {
      * @return 工单ID
      */
     Long createWorkOrder(WorkOrder workOrder);
+
+    /**
+     * 查询需要通过 WebSocket 提醒的工单。
+     * 状态 1/2 始终纳入;其他未办结状态仅在计划完成时间或延期截止时间超时后纳入。
+     */
+    List<WorkOrderPushVO> selectPendingOrOverdueWorkOrders(Date now);
 }

+ 48 - 44
pipe-network-service/zksy-system/src/main/java/com/zksy/base/service/impl/BusinessWorkOrderServiceImpl.java

@@ -5,12 +5,11 @@ import com.baomidou.mybatisplus.core.metadata.IPage;
 import com.baomidou.mybatisplus.extension.plugins.pagination.Page;
 import com.baomidou.mybatisplus.extension.service.impl.ServiceImpl;
 import com.zksy.base.domain.EquipmentBase;
-import com.zksy.base.domain.EquipmentStatus;
 import com.zksy.base.domain.HazardHiddenAccount;
 import com.zksy.base.domain.WorkOrder;
 import com.zksy.base.domain.WorkOrderLog;
+import com.zksy.base.domain.vo.WorkOrderPushVO;
 import com.zksy.base.mapper.EquipmentBaseMapper;
-import com.zksy.base.mapper.EquipmentStatusMapper;
 import com.zksy.base.mapper.HazardHiddenAccountMapper;
 import com.zksy.base.mapper.WorkOrderLogMapper;
 import com.zksy.base.mapper.WorkOrderMapper;
@@ -25,7 +24,7 @@ import org.springframework.transaction.annotation.Transactional;
 import java.util.Arrays;
 import java.util.Date;
 import java.util.List;
-import java.time.LocalDateTime;
+import java.util.ArrayList;
 import java.util.stream.Collectors;
 
 /**
@@ -44,9 +43,6 @@ public class BusinessWorkOrderServiceImpl extends ServiceImpl<WorkOrderMapper, W
     @Autowired
     private HazardHiddenAccountMapper hazardHiddenAccountMapper;
 
-    @Autowired
-    private EquipmentStatusMapper equipmentStatusMapper;
-
     @Override
     public IPage<WorkOrder> selectWorkOrderList(Integer orderType, Integer orderLevel, String orderStatus, String deviceCode, String keyword, String equipmentType, Integer pageNum, Integer pageSize) {
         Page<WorkOrder> page = new Page<>(pageNum, pageSize);
@@ -336,6 +332,51 @@ public class BusinessWorkOrderServiceImpl extends ServiceImpl<WorkOrderMapper, W
         return workOrder.getOrderId();
     }
 
+    @Override
+    public List<WorkOrderPushVO> selectPendingOrOverdueWorkOrders(Date now) {
+        Date currentTime = now != null ? now : new Date();
+        List<WorkOrder> orders = baseMapper.selectPendingOrOverdue(currentTime);
+
+        List<WorkOrderPushVO> result = new ArrayList<>();
+        if (orders == null || orders.isEmpty()) {
+            return result;
+        }
+        for (WorkOrder order : orders) {
+            Integer status = order.getOrderStatus();
+            Date effectiveDeadline = order.getDelayedDeadline() != null
+                    ? order.getDelayedDeadline() : order.getPlanFinishTime();
+            boolean pending = Integer.valueOf(1).equals(status) || Integer.valueOf(2).equals(status);
+            boolean overdue = effectiveDeadline != null && effectiveDeadline.before(currentTime);
+            if (!pending && !overdue) {
+                continue;
+            }
+
+            WorkOrderPushVO item = new WorkOrderPushVO();
+            item.setOrderId(order.getOrderId());
+            item.setOrderNo(order.getOrderNo());
+            item.setOrderStatus(status);
+            item.setOrderLevel(order.getOrderLevel());
+            item.setDeptId(order.getDeptId());
+            item.setReceiveUser(order.getReceiveUser());
+            item.setDeviceId(order.getDeviceId());
+            item.setDeviceCode(order.getDeviceCode());
+            item.setOrderDesc(order.getOrderDesc());
+            item.setPlanFinishTime(order.getPlanFinishTime());
+            item.setDelayedDeadline(order.getDelayedDeadline());
+            item.setEffectiveDeadline(effectiveDeadline);
+            item.setOverdue(overdue);
+            item.setOverdueType(overdue
+                    ? order.getDelayedDeadline() != null ? "DELAYED_DEADLINE" : "PLAN_FINISH_TIME"
+                    : "PENDING");
+            item.setOverdueMinutes(overdue
+                    ? Math.max(0L, (currentTime.getTime() - effectiveDeadline.getTime()) / 60000L) : 0L);
+            item.setCreateTime(order.getCreateTime());
+            item.setUpdateTime(order.getUpdateTime());
+            result.add(item);
+        }
+        return result;
+    }
+
     /**
      * 校验设备是否存在进行中工单(状态 1-5/8/9),存在则禁止再次创建工单
      * 业务规则:同一设备同一时间只能有一个进行中的工单,结案(6)或驳回(7)后方可再次创建
@@ -423,8 +464,7 @@ public class BusinessWorkOrderServiceImpl extends ServiceImpl<WorkOrderMapper, W
         // 工单已结案 → 同步隐患状态为"已整改"
         syncHazardStatus(order, "已整改", result);
 
-        // 工单验收通过 → 重置关联设备报警状态为正常
-        resetEquipmentAlarmStatus(order);
+        // 设备报警状态由实时监测链路维护;工单结案不直接清除报警,避免覆盖仍存在的实时告警。
 
         WorkOrderLog logEntry = new WorkOrderLog();
         logEntry.setOrderId(orderId);
@@ -467,40 +507,4 @@ public class BusinessWorkOrderServiceImpl extends ServiceImpl<WorkOrderMapper, W
         }
     }
 
-    /**
-     * 工单验收通过后,重置关联设备的报警状态为正常(0)
-     * 参考 EquipmentStatusApiServiceImpl.updateAlarmStatus 逻辑
-     */
-    private void resetEquipmentAlarmStatus(WorkOrder order) {
-        if (StringUtils.isEmpty(order.getDeviceId())) {
-            return;
-        }
-        try {
-            LambdaQueryWrapper<EquipmentStatus> wrapper = new LambdaQueryWrapper<>();
-            wrapper.eq(EquipmentStatus::getEquipmentId, order.getDeviceId()).last("LIMIT 1");
-            EquipmentStatus existingStatus = equipmentStatusMapper.selectOne(wrapper);
-            if (existingStatus != null) {
-                // 更新已有记录
-                existingStatus.setAlarmStatus(0);
-                existingStatus.setStatusUpdateTime(LocalDateTime.now());
-                equipmentStatusMapper.updateById(existingStatus);
-                log.info("工单[{}]验收通过,设备[{}]报警状态已重置为正常",
-                        order.getOrderNo(), order.getDeviceCode());
-            } else {
-                // 不存在则创建新记录
-                EquipmentStatus newStatus = new EquipmentStatus();
-                newStatus.setEquipmentId(order.getDeviceId());
-                newStatus.setAlarmStatus(0);
-                newStatus.setCurrentStatus(1); // 默认在用
-                newStatus.setOnlineStatus(1);  // 默认在线
-                newStatus.setStatusUpdateTime(LocalDateTime.now());
-                newStatus.setCreateTime(LocalDateTime.now());
-                equipmentStatusMapper.insert(newStatus);
-                log.info("工单[{}]验收通过,创建设备状态记录,设备[{}]报警状态为正常",
-                        order.getOrderNo(), order.getDeviceCode());
-            }
-        } catch (Exception e) {
-            log.error("工单[{}]重置设备报警状态失败", order.getOrderNo(), e);
-        }
-    }
 }

+ 24 - 2
pipe-network-service/zksy-system/src/main/java/com/zksy/base/service/impl/EquipmentBaseServiceImpl.java

@@ -11,6 +11,7 @@ import com.zksy.base.alarm.domain.WarningThreshold;
 import com.zksy.base.alarm.mapper.WarningThresholdMapper;
 import com.zksy.base.domain.*;
 import com.zksy.base.domain.vo.EquipmentFullVO;
+import com.zksy.base.event.DeviceStatusChangePublisher;
 import com.zksy.base.manhole.domain.ManholeData;
 import com.zksy.base.manhole.mapper.ManholeDataMapper;
 import com.zksy.base.mapper.*;
@@ -50,6 +51,9 @@ public class EquipmentBaseServiceImpl extends ServiceImpl<EquipmentBaseMapper, E
     @Autowired
     private EquipmentStatusMapper equipmentStatusMapper;
 
+    @Autowired
+    private DeviceStatusChangePublisher deviceStatusChangePublisher;
+
     @Autowired
     private EquipmentPointRelMapper equipmentPointRelMapper;
 
@@ -775,12 +779,30 @@ public class EquipmentBaseServiceImpl extends ServiceImpl<EquipmentBaseMapper, E
             status.setAlarmStatus(0);
             status.setOnlineStatus(1);
             status.setStatusUpdateTime(LocalDateTime.now());
-            return equipmentStatusMapper.insert(status) > 0;
+            boolean saved = equipmentStatusMapper.insert(status) > 0;
+            if (saved) {
+                publishEquipmentStatusChange(equipment, status, "pipe-network:updateEquipmentStatus");
+            }
+            return saved;
         } else {
+            boolean changed = !Objects.equals(status.getCurrentStatus(), currentStatus);
             status.setCurrentStatus(currentStatus);
             status.setStatusUpdateTime(LocalDateTime.now());
-            return equipmentStatusMapper.updateById(status) > 0;
+            boolean updated = equipmentStatusMapper.updateById(status) > 0;
+            if (updated && changed) {
+                publishEquipmentStatusChange(equipment, status, "pipe-network:updateEquipmentStatus");
+            }
+            return updated;
+        }
+    }
+
+    private void publishEquipmentStatusChange(EquipmentBase equipment, EquipmentStatus status, String source) {
+        if (equipment == null || status == null || StringUtils.isBlank(equipment.getEquipmentCode())) {
+            return;
         }
+        deviceStatusChangePublisher.publishAfterCommit(equipment.getEquipmentCode(),
+                status.getEquipmentId(), status.getCurrentStatus(), status.getAlarmStatus(),
+                status.getOnlineStatus(), source);
     }
 
     @Override

+ 74 - 4
pipe-network-service/zksy-system/src/main/java/com/zksy/base/service/impl/EquipmentStatusServiceImpl.java

@@ -5,6 +5,7 @@ import com.baomidou.mybatisplus.extension.plugins.pagination.Page;
 import com.baomidou.mybatisplus.extension.service.impl.ServiceImpl;
 import com.zksy.base.domain.EquipmentBase;
 import com.zksy.base.domain.EquipmentStatus;
+import com.zksy.base.event.DeviceStatusChangePublisher;
 import com.zksy.base.mapper.EquipmentBaseMapper;
 import com.zksy.base.mapper.EquipmentStatusMapper;
 import com.zksy.base.service.EquipmentStatusService;
@@ -20,6 +21,7 @@ import org.springframework.transaction.annotation.Transactional;
 import java.util.Date;
 import java.util.HashSet;
 import java.util.List;
+import java.util.Objects;
 import java.util.Set;
 import java.util.stream.Collectors;
 
@@ -34,6 +36,9 @@ public class EquipmentStatusServiceImpl extends ServiceImpl<EquipmentStatusMappe
     @Autowired
     private EquipmentBaseMapper equipmentBaseMapper;
 
+    @Autowired
+    private DeviceStatusChangePublisher deviceStatusChangePublisher;
+
     @Override
     public Page<EquipmentStatus> findByPage(long pageNum, long pageSize,
                                             String statusId,String equipmentId,
@@ -55,6 +60,7 @@ public class EquipmentStatusServiceImpl extends ServiceImpl<EquipmentStatusMappe
      * @return
      */
     @Override
+    @Transactional
     public boolean saveWithCheck(EquipmentStatus entity) {
         //校验关联的设备是否存在
         String equipmentId = entity.getEquipmentId();
@@ -67,8 +73,13 @@ public class EquipmentStatusServiceImpl extends ServiceImpl<EquipmentStatusMappe
         wrapper.eq(EquipmentStatus::getEquipmentId, equipmentId);
         EquipmentStatus existing = this.getOne(wrapper);
         if (existing != null) {
+            boolean changed = hasStateChange(existing, entity);
             entity.setStatusId(existing.getStatusId());
-            return this.updateById(entity);
+            boolean updated = this.updateById(entity);
+            if (updated && changed) {
+                publishLatest(existing.getEquipmentId(), existing.getStatusId(), "pipe-network:saveStatus");
+            }
+            return updated;
         }
 
         //校验当前状态值是否合法(1-在用,2-闲置,3-维修,4-报废,5-待入库)
@@ -86,7 +97,11 @@ public class EquipmentStatusServiceImpl extends ServiceImpl<EquipmentStatusMappe
             }
         }
         //新增
-        return this.save(entity);
+        boolean saved = this.save(entity);
+        if (saved) {
+            publish(entity, "pipe-network:saveStatus");
+        }
+        return saved;
     }
 
     /**
@@ -154,7 +169,13 @@ public class EquipmentStatusServiceImpl extends ServiceImpl<EquipmentStatusMappe
             }
         }
         //批量新增
-        return this.saveBatch(entityList);
+        boolean saved = this.saveBatch(entityList);
+        if (saved) {
+            for (EquipmentStatus entity : entityList) {
+                publish(entity, "pipe-network:saveStatusBatch");
+            }
+        }
+        return saved;
     }
 
     /**
@@ -163,6 +184,7 @@ public class EquipmentStatusServiceImpl extends ServiceImpl<EquipmentStatusMappe
      * @return
      */
     @Override
+    @Transactional
     public boolean updateWithCheck(EquipmentStatus entity) {
         //校验当前状态记录是否存在
         Integer statusId = entity.getStatusId();
@@ -205,7 +227,12 @@ public class EquipmentStatusServiceImpl extends ServiceImpl<EquipmentStatusMappe
             }
         }
         //更新
-        return this.updateById(entity);
+        boolean changed = hasStateChange(oldStatus, entity);
+        boolean updated = this.updateById(entity);
+        if (updated && changed) {
+            publishLatest(oldStatus.getEquipmentId(), oldStatus.getStatusId(), "pipe-network:updateStatus");
+        }
+        return updated;
     }
 
     /**
@@ -214,6 +241,7 @@ public class EquipmentStatusServiceImpl extends ServiceImpl<EquipmentStatusMappe
      * @return
      */
     @Override
+    @Transactional
     public boolean updateDeviceStatus(EquipmentStatusInDTO inDTO) {
         log.info("修改设备状态-入参:{}", inDTO);
         //校验状态值合法性
@@ -222,11 +250,19 @@ public class EquipmentStatusServiceImpl extends ServiceImpl<EquipmentStatusMappe
             throw new ServiceException("设备状态值无效(必须为1-5):" + status);
         }
         //校验当前状态记录是否存在
+        EquipmentStatus oldStatus = this.lambdaQuery()
+                .eq(EquipmentStatus::getEquipmentId, inDTO.getEquipmentId())
+                .one();
+        boolean changed = oldStatus != null
+                && !Objects.equals(oldStatus.getCurrentStatus(), inDTO.getCurrentStatus());
         boolean update = this.lambdaUpdate()
                 .set(EquipmentStatus::getCurrentStatus, inDTO.getCurrentStatus())
                 .set(EquipmentStatus::getStatusUpdateTime, new Date())
                 .eq(EquipmentStatus::getEquipmentId, inDTO.getEquipmentId())
                 .update();
+        if (update && changed) {
+            publishLatest(inDTO.getEquipmentId(), oldStatus.getStatusId(), "pipe-network:updateDeviceStatus");
+        }
         return update;
     }
 
@@ -257,4 +293,38 @@ public class EquipmentStatusServiceImpl extends ServiceImpl<EquipmentStatusMappe
                 deviceCode, equipmentBase.getEquipmentId(), status != null);
         return status;
     }
+
+    private boolean hasStateChange(EquipmentStatus oldStatus, EquipmentStatus patch) {
+        if (oldStatus == null || patch == null) {
+            return false;
+        }
+        return (patch.getCurrentStatus() != null
+                && !Objects.equals(oldStatus.getCurrentStatus(), patch.getCurrentStatus()))
+                || (patch.getAlarmStatus() != null
+                && !Objects.equals(oldStatus.getAlarmStatus(), patch.getAlarmStatus()))
+                || (patch.getOnlineStatus() != null
+                && !Objects.equals(oldStatus.getOnlineStatus(), patch.getOnlineStatus()));
+    }
+
+    private void publish(EquipmentStatus status, String source) {
+        if (status == null) {
+            return;
+        }
+        publishLatest(status.getEquipmentId(), status.getStatusId(), source);
+    }
+
+    private void publishLatest(String equipmentId, Integer statusId, String source) {
+        if (equipmentId == null) {
+            return;
+        }
+        EquipmentStatus latest = statusId != null ? this.getById(statusId) : null;
+        if (latest == null) {
+            latest = this.lambdaQuery().eq(EquipmentStatus::getEquipmentId, equipmentId).one();
+        }
+        EquipmentBase equipment = equipmentBaseMapper.selectById(equipmentId);
+        if (latest != null && equipment != null) {
+            deviceStatusChangePublisher.publishAfterCommit(equipment.getEquipmentCode(), equipmentId,
+                    latest.getCurrentStatus(), latest.getAlarmStatus(), latest.getOnlineStatus(), source);
+        }
+    }
 }

+ 21 - 2
pipe-network-service/zksy-system/src/main/java/com/zksy/base/service/impl/WaterSupplyDeviceServiceImpl.java

@@ -14,6 +14,7 @@ import com.zksy.base.domain.EquipmentMaintain;
 import com.zksy.base.domain.EquipmentStatus;
 import com.zksy.base.domain.EquipmentType;
 import com.zksy.base.domain.vo.*;
+import com.zksy.base.event.DeviceStatusChangePublisher;
 import com.zksy.base.environment.domain.ERealTimeData;
 import com.zksy.base.environment.mapper.ERealTimeDataMapper;
 import com.zksy.base.gas.domain.GasMonitorData;
@@ -53,6 +54,9 @@ public class WaterSupplyDeviceServiceImpl extends ServiceImpl<EquipmentBaseMappe
     @Autowired
     private EquipmentStatusMapper equipmentStatusMapper;
 
+    @Autowired
+    private DeviceStatusChangePublisher deviceStatusChangePublisher;
+
     @Autowired
     private EquipmentMaintainMapper equipmentMaintainMapper;
 
@@ -288,9 +292,13 @@ public class WaterSupplyDeviceServiceImpl extends ServiceImpl<EquipmentBaseMappe
                     new LambdaQueryWrapper<EquipmentStatus>().eq(EquipmentStatus::getEquipmentId, vo.getEquipmentId())
             );
             if (status != null) {
+                boolean changed = !Objects.equals(status.getCurrentStatus(), 3);
                 status.setCurrentStatus(3);
                 status.setStatusUpdateTime(LocalDateTime.now());
-                equipmentStatusMapper.updateById(status);
+                boolean updated = equipmentStatusMapper.updateById(status) > 0;
+                if (updated && changed) {
+                    publishEquipmentStatusChange(device, status, "pipe-network:createMaintainOrder");
+                }
             } else {
                 EquipmentStatus newStatus = new EquipmentStatus();
                 newStatus.setEquipmentId(vo.getEquipmentId());
@@ -299,7 +307,9 @@ public class WaterSupplyDeviceServiceImpl extends ServiceImpl<EquipmentBaseMappe
                 newStatus.setOnlineStatus(1);
                 newStatus.setStatusUpdateTime(LocalDateTime.now());
                 newStatus.setCreateTime(LocalDateTime.now());
-                equipmentStatusMapper.insert(newStatus);
+                if (equipmentStatusMapper.insert(newStatus) > 0) {
+                    publishEquipmentStatusChange(device, newStatus, "pipe-network:createMaintainOrder");
+                }
             }
             log.info("设备[{}]派单成功,工单ID:{},维修人:{}",
                     vo.getEquipmentId(), maintain.getMaintainId(), vo.getMaintainPerson());
@@ -311,6 +321,15 @@ public class WaterSupplyDeviceServiceImpl extends ServiceImpl<EquipmentBaseMappe
         return insert > 0;
     }
 
+    private void publishEquipmentStatusChange(EquipmentBase equipment, EquipmentStatus status, String source) {
+        if (equipment == null || status == null || StrUtil.isBlank(equipment.getEquipmentCode())) {
+            return;
+        }
+        deviceStatusChangePublisher.publishAfterCommit(equipment.getEquipmentCode(),
+                status.getEquipmentId(), status.getCurrentStatus(), status.getAlarmStatus(),
+                status.getOnlineStatus(), source);
+    }
+
     private void sendMaintainSmsNotify(EquipmentBase device, EquipmentMaintain maintain) {
         try {
             if (smsTemplate == null) {

+ 15 - 0
pipe-network-service/zksy-system/src/main/resources/mapper/WorkOrderMapper.xml

@@ -43,4 +43,19 @@
         real_finish_time,repair_result,maintenance_before_img,maintenance_after_img,verify_user,
         verify_time,create_time,update_time
     </sql>
+
+    <select id="selectPendingOrOverdue" resultMap="BaseResultMap">
+        SELECT
+        <include refid="Base_Column_List"/>
+        FROM app_user.work_order
+        WHERE order_status IN (1, 2)
+           OR (
+                order_status IN (3, 4, 5, 8, 9)
+                AND COALESCE(delayed_deadline, plan_finish_time) &lt; #{now}
+           )
+        ORDER BY order_level ASC NULLS LAST,
+                 order_status ASC,
+                 create_time DESC NULLS LAST,
+                 order_id DESC
+    </select>
 </mapper>

+ 6 - 1
zk-api-service/pom.xml

@@ -17,6 +17,11 @@
     </properties>
 
     <dependencies>
+        <dependency>
+            <groupId>org.springframework.boot</groupId>
+            <artifactId>spring-boot-starter-test</artifactId>
+            <scope>test</scope>
+        </dependency>
         <!--common-->
         <dependency>
             <groupId>com.zksy</groupId>
@@ -102,4 +107,4 @@
             </plugin>
         </plugins>
     </build>
-</project>
+</project>

+ 23 - 3
zk-api-service/src/main/java/com/zksy/api/service/impl/EquipmentStatusApiServiceImpl.java

@@ -3,6 +3,7 @@ package com.zksy.api.service.impl;
 import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
 import com.zksy.api.domain.EquipmentBase;
 import com.zksy.api.domain.EquipmentStatus;
+import com.zksy.api.event.DeviceStatusChangePublisher;
 import com.zksy.api.mapper.EquipmentBaseApiMapper;
 import com.zksy.api.mapper.EquipmentStatusApiMapper;
 import com.zksy.api.service.EquipmentStatusApiService;
@@ -12,6 +13,7 @@ import org.springframework.stereotype.Service;
 import org.springframework.transaction.annotation.Transactional;
 
 import java.time.LocalDateTime;
+import java.util.Objects;
 
 /**
  * 设备报警状态服务实现
@@ -27,6 +29,9 @@ public class EquipmentStatusApiServiceImpl implements EquipmentStatusApiService
     @Autowired
     private EquipmentStatusApiMapper equipmentStatusApiMapper;
 
+    @Autowired
+    private DeviceStatusChangePublisher deviceStatusChangePublisher;
+
     @Override
     @Transactional
     public boolean updateAlarmStatus(String deviceCode, Integer alarmStatus) {
@@ -53,12 +58,21 @@ public class EquipmentStatusApiServiceImpl implements EquipmentStatusApiService
                         .eq(EquipmentStatus::getEquipmentId, equipmentId));
 
         if (existingStatus != null) {
+            Integer previousAlarmStatus = existingStatus.getAlarmStatus();
             // 更新已有记录
             existingStatus.setAlarmStatus(alarmStatus);
             existingStatus.setStatusUpdateTime(LocalDateTime.now());
-            equipmentStatusApiMapper.updateById(existingStatus);
+            int affected = equipmentStatusApiMapper.updateById(existingStatus);
+            if (affected <= 0) {
+                return false;
+            }
             log.info("更新设备报警状态成功: deviceCode={}, equipmentId={}, alarmStatus={}",
                     deviceCode, equipmentId, alarmStatus);
+            if (!Objects.equals(previousAlarmStatus, alarmStatus)) {
+                deviceStatusChangePublisher.publishAfterCommit(deviceCode, equipmentId,
+                        existingStatus.getCurrentStatus(), existingStatus.getAlarmStatus(),
+                        existingStatus.getOnlineStatus(), "zk-api:updateAlarmStatus");
+            }
         } else {
             // 不存在则创建新记录
             EquipmentStatus newStatus = new EquipmentStatus();
@@ -68,11 +82,17 @@ public class EquipmentStatusApiServiceImpl implements EquipmentStatusApiService
             newStatus.setOnlineStatus(1);  // 默认在线
             newStatus.setStatusUpdateTime(LocalDateTime.now());
             newStatus.setCreateTime(LocalDateTime.now());
-            equipmentStatusApiMapper.insert(newStatus);
+            int affected = equipmentStatusApiMapper.insert(newStatus);
+            if (affected <= 0) {
+                return false;
+            }
             log.info("创建设备报警状态记录成功: deviceCode={}, equipmentId={}, alarmStatus={}",
                     deviceCode, equipmentId, alarmStatus);
+            deviceStatusChangePublisher.publishAfterCommit(deviceCode, equipmentId,
+                    newStatus.getCurrentStatus(), newStatus.getAlarmStatus(),
+                    newStatus.getOnlineStatus(), "zk-api:updateAlarmStatus");
         }
 
         return true;
     }
-}
+}