This is an automated email from the ASF dual-hosted git repository.

yinyijun pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/hertzbeat.git


The following commit(s) were added to refs/heads/master by this push:
     new 1050032a13 [improve] SSE exception handling improvements (#3775)
1050032a13 is described below

commit 1050032a13c55888596cb3589f65b407d258deae
Author: Duansg <[email protected]>
AuthorDate: Fri Sep 19 00:24:49 2025 +0800

    [improve] SSE exception handling improvements (#3775)
---
 .../hertzbeat/alert/config/AlertSseManager.java    | 18 ++++--
 .../alert/config/AlertSseManagerTest.java          | 69 ++++++++++++++++++++++
 .../manager/config/ManagerSseManager.java          | 18 ++++--
 3 files changed, 97 insertions(+), 8 deletions(-)

diff --git 
a/hertzbeat-alerter/src/main/java/org/apache/hertzbeat/alert/config/AlertSseManager.java
 
b/hertzbeat-alerter/src/main/java/org/apache/hertzbeat/alert/config/AlertSseManager.java
index 27f817dd4b..e3d245621e 100644
--- 
a/hertzbeat-alerter/src/main/java/org/apache/hertzbeat/alert/config/AlertSseManager.java
+++ 
b/hertzbeat-alerter/src/main/java/org/apache/hertzbeat/alert/config/AlertSseManager.java
@@ -22,10 +22,12 @@ package org.apache.hertzbeat.alert.config;
 import lombok.extern.slf4j.Slf4j;
 import org.springframework.scheduling.annotation.Async;
 import org.springframework.stereotype.Component;
+import 
org.springframework.web.servlet.mvc.method.annotation.ResponseBodyEmitter;
 import org.springframework.web.servlet.mvc.method.annotation.SseEmitter;
 
 import java.io.IOException;
 import java.util.Map;
+import java.util.Optional;
 import java.util.concurrent.ConcurrentHashMap;
 
 /**
@@ -54,16 +56,24 @@ public class AlertSseManager {
                         .name("ALERT_EVENT")
                         .data(data));
             } catch (IOException | IllegalStateException e) {
-                emitter.complete();
-                removeEmitter(clientId);
+                tryCompleteAndClean(clientId, emitter);
             } catch (Exception exception) {
                 log.error("Failed to broadcast alert data to client: {}", 
exception.getMessage());
-                emitter.complete();
-                removeEmitter(clientId);
+                tryCompleteAndClean(clientId, emitter);
             }
         });
     }
 
+    private void tryCompleteAndClean(Long clientId, SseEmitter emitter) {
+        try {
+            
Optional.ofNullable(emitter).ifPresent(ResponseBodyEmitter::complete);
+        } catch (Throwable e) {
+            log.debug("Failed to complete emitter for client {}: {}", 
clientId, e.getMessage());
+        }
+        // execute clear
+        removeEmitter(clientId);
+    }
+
     private void removeEmitter(Long clientId) {
         emitters.remove(clientId);
     }
diff --git 
a/hertzbeat-alerter/src/test/java/org/apache/hertzbeat/alert/config/AlertSseManagerTest.java
 
b/hertzbeat-alerter/src/test/java/org/apache/hertzbeat/alert/config/AlertSseManagerTest.java
new file mode 100644
index 0000000000..40e88ddef9
--- /dev/null
+++ 
b/hertzbeat-alerter/src/test/java/org/apache/hertzbeat/alert/config/AlertSseManagerTest.java
@@ -0,0 +1,69 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.hertzbeat.alert.config;
+
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.springframework.web.servlet.mvc.method.annotation.SseEmitter;
+
+import java.lang.reflect.Field;
+import java.util.HashMap;
+import java.util.Map;
+
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.Mockito.doThrow;
+import static org.mockito.Mockito.mock;
+
+/**
+ * alert sse manager test
+ */
+public class AlertSseManagerTest {
+
+    private AlertSseManager alertSseManager;
+
+    @BeforeEach
+    void setUp() {
+        alertSseManager = new AlertSseManager();
+    }
+
+    @Test
+    void testCompleteThrowsException() throws Exception {
+        SseEmitter emitter = alertSseManager.createEmitter(1L);
+        assertNotNull(emitter);
+
+        Map<Long, SseEmitter> emitters = new HashMap<>();
+        SseEmitter spyEmitter = mock(SseEmitter.class);
+        
+        doThrow(new IllegalStateException("Simulated output stream 
error")).when(spyEmitter).send(any(SseEmitter.SseEventBuilder.class));
+        doThrow(new RuntimeException("Complete 
failed")).when(spyEmitter).complete();
+        
+        emitters.put(1L, spyEmitter);
+
+        Field emittersField = 
AlertSseManager.class.getDeclaredField("emitters");
+        emittersField.setAccessible(true);
+        emittersField.set(alertSseManager, emitters);
+
+        assertThrows(RuntimeException.class, () -> 
alertSseManager.broadcast("{\"id\":1,\"content\":\"Test alert\"}"));
+        Map<Long, SseEmitter> currentEmitters = (Map<Long, SseEmitter>) 
emittersField.get(alertSseManager);
+        assertFalse(currentEmitters.containsKey(1L), "Emitter should still 
exist because complete() threw exception");
+    }
+
+}
\ No newline at end of file
diff --git 
a/hertzbeat-manager/src/main/java/org/apache/hertzbeat/manager/config/ManagerSseManager.java
 
b/hertzbeat-manager/src/main/java/org/apache/hertzbeat/manager/config/ManagerSseManager.java
index 5db74883a4..1fd681566c 100644
--- 
a/hertzbeat-manager/src/main/java/org/apache/hertzbeat/manager/config/ManagerSseManager.java
+++ 
b/hertzbeat-manager/src/main/java/org/apache/hertzbeat/manager/config/ManagerSseManager.java
@@ -26,10 +26,12 @@ import 
org.apache.hertzbeat.common.entity.dto.ManagerMessage;
 import org.apache.hertzbeat.common.util.JsonUtil;
 import org.springframework.scheduling.annotation.Async;
 import org.springframework.stereotype.Component;
+import 
org.springframework.web.servlet.mvc.method.annotation.ResponseBodyEmitter;
 import org.springframework.web.servlet.mvc.method.annotation.SseEmitter;
 
 import java.io.IOException;
 import java.util.Map;
+import java.util.Optional;
 import java.util.concurrent.ConcurrentHashMap;
 
 /**
@@ -57,16 +59,24 @@ public class ManagerSseManager {
                         .name(eventName)
                         .data(data));
             } catch (IOException | IllegalStateException e) {
-                emitter.complete();
-                removeEmitter(clientId);
+                tryCompleteAndClean(clientId, emitter);
             } catch (Exception exception) {
                 log.error("Failed to broadcast manager message data to client: 
{}", exception.getMessage());
-                emitter.complete();
-                removeEmitter(clientId);
+                tryCompleteAndClean(clientId, emitter);
             }
         });
     }
 
+    private void tryCompleteAndClean(Long clientId, SseEmitter emitter) {
+        try {
+            
Optional.ofNullable(emitter).ifPresent(ResponseBodyEmitter::complete);
+        } catch (Throwable e) {
+            log.debug("Failed to complete emitter for client {}: {}", 
clientId, e.getMessage());
+        }
+        // execute clear
+        removeEmitter(clientId);
+    }
+
     public void broadcastImportTaskInProgress(String taskName, Integer 
progress){
         ManagerMessage managerMessage = 
ImportTaskMessage.createInProgressMessage(taskName, progress);
         broadcast(ManagerEventTypeEnum.IMPORT_TASK_EVENT.getValue(), 
JsonUtil.toJson(managerMessage));


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to