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]