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

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


The following commit(s) were added to refs/heads/master by this push:
     new a37d4dfad8 fix(admin): load alert receivers per namespace when 
dispatching (#7350)
a37d4dfad8 is described below

commit a37d4dfad8d99780617cae468b582f82948d9403
Author: BobSong <[email protected]>
AuthorDate: Wed Sep 30 10:00:58 2026 +0800

    fix(admin): load alert receivers per namespace when dispatching (#7350)
    
    Co-authored-by: BobSong-dev <[email protected]>
    Co-authored-by: aias00 <[email protected]>
    Co-authored-by: shown <[email protected]>
---
 .../shenyu/admin/mapper/AlertReceiverMapper.java   |  9 +++++
 .../service/impl/AlertDispatchServiceImpl.java     | 30 ++++++++++-----
 .../resources/mappers/alert-receiver-sqlmap.xml    |  8 ++++
 .../admin/service/AlertDispatchServiceTest.java    | 45 +++++++++++-----------
 4 files changed, 60 insertions(+), 32 deletions(-)

diff --git 
a/shenyu-admin/src/main/java/org/apache/shenyu/admin/mapper/AlertReceiverMapper.java
 
b/shenyu-admin/src/main/java/org/apache/shenyu/admin/mapper/AlertReceiverMapper.java
index a6d14af010..39797beea7 100644
--- 
a/shenyu-admin/src/main/java/org/apache/shenyu/admin/mapper/AlertReceiverMapper.java
+++ 
b/shenyu-admin/src/main/java/org/apache/shenyu/admin/mapper/AlertReceiverMapper.java
@@ -41,6 +41,15 @@ public interface AlertReceiverMapper extends ExistProvider {
      */
     List<AlertReceiverDTO> selectAll();
     
+    /**
+     * select receivers scoped to a namespace, including receivers without a 
namespace
+     * that match every namespace.
+     *
+     * @param namespaceId the namespace id
+     * @return receiver list
+     */
+    List<AlertReceiverDTO> selectByNamespaceId(@Param("namespaceId") String 
namespaceId);
+    
     /**
      * insert record to table.
      *
diff --git 
a/shenyu-admin/src/main/java/org/apache/shenyu/admin/service/impl/AlertDispatchServiceImpl.java
 
b/shenyu-admin/src/main/java/org/apache/shenyu/admin/service/impl/AlertDispatchServiceImpl.java
index 839e895ec3..3098d4d855 100644
--- 
a/shenyu-admin/src/main/java/org/apache/shenyu/admin/service/impl/AlertDispatchServiceImpl.java
+++ 
b/shenyu-admin/src/main/java/org/apache/shenyu/admin/service/impl/AlertDispatchServiceImpl.java
@@ -35,11 +35,12 @@ import org.springframework.util.CollectionUtils;
 import java.util.List;
 import java.util.Map;
 import java.util.Objects;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.ConcurrentMap;
 import java.util.concurrent.LinkedBlockingQueue;
 import java.util.concurrent.ThreadFactory;
 import java.util.concurrent.ThreadPoolExecutor;
 import java.util.concurrent.TimeUnit;
-import java.util.concurrent.atomic.AtomicReference;
 import java.util.stream.Collectors;
 
 /**
@@ -54,13 +55,18 @@ public class AlertDispatchServiceImpl implements 
AlertDispatchService, Disposabl
     
     private final AlertReceiverMapper alertReceiverMapper;
     
-    private final AtomicReference<List<AlertReceiverDTO>> 
alertReceiverReference;
+    /**
+     * Receivers cached per alert namespace. Values are scoped queries 
(namespace-local
+     * receivers plus namespace-free receivers), so a refresh no longer loads 
the whole
+     * table across all namespaces.
+     */
+    private final ConcurrentMap<String, List<AlertReceiverDTO>> 
alertReceiverCache;
 
     private final ThreadPoolExecutor workerExecutor;
     
     public AlertDispatchServiceImpl(final List<AlertNotifyHandler> 
alertNotifyHandlerList, final AlertReceiverMapper alertReceiverMapper) {
         this.alertReceiverMapper = alertReceiverMapper;
-        this.alertReceiverReference = new AtomicReference<>();
+        this.alertReceiverCache = new ConcurrentHashMap<>();
         alertNotifyHandlerMap = 
Maps.newHashMapWithExpectedSize(alertNotifyHandlerList.size());
         ThreadFactory threadFactory = new ThreadFactoryBuilder()
                 .setUncaughtExceptionHandler((thread, throwable) -> {
@@ -88,7 +94,7 @@ public class AlertDispatchServiceImpl implements 
AlertDispatchService, Disposabl
     
     @Override
     public void clearCache() {
-        this.alertReceiverReference.set(null);
+        this.alertReceiverCache.clear();
     }
     
     @Override
@@ -146,11 +152,8 @@ public class AlertDispatchServiceImpl implements 
AlertDispatchService, Disposabl
         }
         
         private List<AlertReceiverDTO> matchReceiverByRules(final AlarmContent 
alert) {
-            List<AlertReceiverDTO> dtoList = alertReceiverReference.get();
-            if (Objects.isNull(dtoList)) {
-                dtoList = alertReceiverMapper.selectAll();
-                alertReceiverReference.set(dtoList);
-            }
+            final String namespaceId = 
StringUtils.defaultString(alert.getNamespaceId());
+            List<AlertReceiverDTO> dtoList = 
alertReceiverCache.computeIfAbsent(namespaceId, this::loadReceivers);
             return dtoList.stream().filter(item -> {
                 if (item.isEnable()) {
                     if (item.isMatchAll()) {
@@ -184,5 +187,14 @@ public class AlertDispatchServiceImpl implements 
AlertDispatchService, Disposabl
                 }
             }).collect(Collectors.toList());
         }
+        
+        private List<AlertReceiverDTO> loadReceivers(final String namespaceId) 
{
+            // alerts without a namespace can be matched by any 
namespace-scoped receiver,
+            // so they keep loading the full list and rely on the in-memory 
namespace filter
+            if (StringUtils.isBlank(namespaceId)) {
+                return alertReceiverMapper.selectAll();
+            }
+            return alertReceiverMapper.selectByNamespaceId(namespaceId);
+        }
     }
 }
diff --git a/shenyu-admin/src/main/resources/mappers/alert-receiver-sqlmap.xml 
b/shenyu-admin/src/main/resources/mappers/alert-receiver-sqlmap.xml
index 17c4accf38..998615a6a4 100644
--- a/shenyu-admin/src/main/resources/mappers/alert-receiver-sqlmap.xml
+++ b/shenyu-admin/src/main/resources/mappers/alert-receiver-sqlmap.xml
@@ -60,6 +60,14 @@
         <include refid="Base_Column_List"/>
         from alert_receiver
     </select>
+    <select id="selectByNamespaceId" 
resultType="org.apache.shenyu.alert.model.AlertReceiverDTO">
+        select
+        <include refid="Base_Column_List"/>
+        from alert_receiver
+        where namespace_id = #{namespaceId,jdbcType=VARCHAR}
+           or namespace_id is null
+           or namespace_id = ''
+    </select>
     <insert id="insert" 
parameterType="org.apache.shenyu.admin.model.entity.AlertReceiverDO">
         <[email protected]>
         insert into alert_receiver (id, `name`, enable, `type`, phone, email, 
hook_url, wechat_id, access_token, tg_bot_token, tg_user_id, slack_web_hook_url,
diff --git 
a/shenyu-admin/src/test/java/org/apache/shenyu/admin/service/AlertDispatchServiceTest.java
 
b/shenyu-admin/src/test/java/org/apache/shenyu/admin/service/AlertDispatchServiceTest.java
index 08db1d3b8e..82c529b776 100644
--- 
a/shenyu-admin/src/test/java/org/apache/shenyu/admin/service/AlertDispatchServiceTest.java
+++ 
b/shenyu-admin/src/test/java/org/apache/shenyu/admin/service/AlertDispatchServiceTest.java
@@ -40,14 +40,13 @@ import java.util.HashMap;
 import java.util.List;
 import java.util.Map;
 import java.util.Objects;
+import java.util.concurrent.ConcurrentMap;
 import java.util.concurrent.CountDownLatch;
 import java.util.concurrent.ThreadPoolExecutor;
 import java.util.concurrent.TimeUnit;
-import java.util.concurrent.atomic.AtomicReference;
 
 import static org.junit.jupiter.api.Assertions.assertFalse;
 import static org.junit.jupiter.api.Assertions.assertNotNull;
-import static org.junit.jupiter.api.Assertions.assertNull;
 import static org.junit.jupiter.api.Assertions.assertTrue;
 import static org.mockito.ArgumentMatchers.any;
 import static org.mockito.ArgumentMatchers.eq;
@@ -130,7 +129,7 @@ public class AlertDispatchServiceTest {
     void testDispatchAlertSuccess() throws InterruptedException {
         final AlertReceiverDTO receiver = createTestReceiver(EMAIL_TYPE, true, 
false);
         
-        
when(alertReceiverMapper.selectAll()).thenReturn(Arrays.asList(receiver));
+        
when(alertReceiverMapper.selectByNamespaceId(TEST_NAMESPACE_ID)).thenReturn(Arrays.asList(receiver));
         
         final CountDownLatch latch = new CountDownLatch(1);
         final AlarmContent alarmContent = createTestAlarmContent();
@@ -148,7 +147,7 @@ public class AlertDispatchServiceTest {
         
         // Verify handler was called
         Mockito.verify(emailHandler, times(1)).send(eq(receiver), 
eq(alarmContent));
-        verify(alertReceiverMapper, times(1)).selectAll();
+        verify(alertReceiverMapper, 
times(1)).selectByNamespaceId(TEST_NAMESPACE_ID);
     }
 
     @Test
@@ -164,7 +163,7 @@ public class AlertDispatchServiceTest {
         verify(emailHandler, never()).send(any(), any());
         verify(webhookHandler, never()).send(any(), any());
         verify(wechatHandler, never()).send(any(), any());
-        verify(alertReceiverMapper, never()).selectAll();
+        verify(alertReceiverMapper, 
never()).selectByNamespaceId(TEST_NAMESPACE_ID);
     }
 
     @Test
@@ -172,7 +171,7 @@ public class AlertDispatchServiceTest {
         final AlertReceiverDTO emailReceiver = createTestReceiver(EMAIL_TYPE, 
true, false);
         final AlertReceiverDTO webhookReceiver = 
createTestReceiver(WEBHOOK_TYPE, true, false);
         
-        
when(alertReceiverMapper.selectAll()).thenReturn(Arrays.asList(emailReceiver, 
webhookReceiver));
+        
when(alertReceiverMapper.selectByNamespaceId(TEST_NAMESPACE_ID)).thenReturn(Arrays.asList(emailReceiver,
 webhookReceiver));
         
         final CountDownLatch latch = new CountDownLatch(2);
         final AlarmContent alarmContent = createTestAlarmContent();
@@ -201,7 +200,7 @@ public class AlertDispatchServiceTest {
     void testDispatchAlertWithHandlerException() throws InterruptedException {
         final AlertReceiverDTO receiver = createTestReceiver(EMAIL_TYPE, true, 
false);
         
-        
when(alertReceiverMapper.selectAll()).thenReturn(Arrays.asList(receiver));
+        
when(alertReceiverMapper.selectByNamespaceId(TEST_NAMESPACE_ID)).thenReturn(Arrays.asList(receiver));
         
         final CountDownLatch latch = new CountDownLatch(1);
         final AlarmContent alarmContent = createTestAlarmContent();
@@ -221,7 +220,7 @@ public class AlertDispatchServiceTest {
     @Test
     void testClearCache() throws Exception {
         final AlertReceiverDTO receiver = createTestReceiver(EMAIL_TYPE, true, 
false);
-        
when(alertReceiverMapper.selectAll()).thenReturn(Arrays.asList(receiver));
+        
when(alertReceiverMapper.selectByNamespaceId(TEST_NAMESPACE_ID)).thenReturn(Arrays.asList(receiver));
         
         final CountDownLatch latch = new CountDownLatch(1);
         final AlarmContent alarmContent = createTestAlarmContent();
@@ -234,12 +233,12 @@ public class AlertDispatchServiceTest {
         alertDispatchService.dispatchAlert(alarmContent);
         assertTrue(latch.await(5, TimeUnit.SECONDS));
         
-        final AtomicReference<List<AlertReceiverDTO>> cacheRef = 
getAlertReceiverReference();
-        assertNotNull(cacheRef.get());
+        final ConcurrentMap<String, List<AlertReceiverDTO>> cache = 
getAlertReceiverCache();
+        assertNotNull(cache.get(TEST_NAMESPACE_ID));
         
         alertDispatchService.clearCache();
         
-        assertNull(cacheRef.get());
+        assertTrue(cache.isEmpty());
     }
 
     @Test
@@ -293,7 +292,7 @@ public class AlertDispatchServiceTest {
         final AlarmContent alarmContent = createTestAlarmContent();
         final AlertReceiverDTO matchAllReceiver = 
createTestReceiver(EMAIL_TYPE, true, true);
         
-        
when(alertReceiverMapper.selectAll()).thenReturn(Arrays.asList(matchAllReceiver));
+        
when(alertReceiverMapper.selectByNamespaceId(TEST_NAMESPACE_ID)).thenReturn(Arrays.asList(matchAllReceiver));
         
         final CountDownLatch latch = new CountDownLatch(1);
         doAnswer(invocation -> {
@@ -312,7 +311,7 @@ public class AlertDispatchServiceTest {
         final AlarmContent alarmContent = createTestAlarmContent();
         final AlertReceiverDTO disabledReceiver = 
createTestReceiver(EMAIL_TYPE, false, false);
         
-        
when(alertReceiverMapper.selectAll()).thenReturn(Arrays.asList(disabledReceiver));
+        
when(alertReceiverMapper.selectByNamespaceId(TEST_NAMESPACE_ID)).thenReturn(Arrays.asList(disabledReceiver));
         
         alertDispatchService.dispatchAlert(alarmContent);
         
@@ -332,7 +331,7 @@ public class AlertDispatchServiceTest {
         final AlertReceiverDTO nonMatchingReceiver = 
createTestReceiver(WEBHOOK_TYPE, true, false);
         nonMatchingReceiver.setNamespaceId("different-namespace");
         
-        
when(alertReceiverMapper.selectAll()).thenReturn(Arrays.asList(matchingReceiver,
 nonMatchingReceiver));
+        
when(alertReceiverMapper.selectByNamespaceId(TEST_NAMESPACE_ID)).thenReturn(Arrays.asList(matchingReceiver,
 nonMatchingReceiver));
         
         final CountDownLatch latch = new CountDownLatch(1);
         doAnswer(invocation -> {
@@ -359,7 +358,7 @@ public class AlertDispatchServiceTest {
         final AlertReceiverDTO nonMatchingReceiver = 
createTestReceiver(WEBHOOK_TYPE, true, false);
         nonMatchingReceiver.setLevels(Arrays.asList((byte) 0, (byte) 2));
         
-        
when(alertReceiverMapper.selectAll()).thenReturn(Arrays.asList(matchingReceiver,
 nonMatchingReceiver));
+        
when(alertReceiverMapper.selectByNamespaceId(TEST_NAMESPACE_ID)).thenReturn(Arrays.asList(matchingReceiver,
 nonMatchingReceiver));
         
         final CountDownLatch latch = new CountDownLatch(1);
         doAnswer(invocation -> {
@@ -393,7 +392,7 @@ public class AlertDispatchServiceTest {
         nonMatchingLabels.put("service", "api");
         nonMatchingReceiver.setLabels(nonMatchingLabels);
         
-        
when(alertReceiverMapper.selectAll()).thenReturn(Arrays.asList(matchingReceiver,
 nonMatchingReceiver));
+        
when(alertReceiverMapper.selectByNamespaceId(TEST_NAMESPACE_ID)).thenReturn(Arrays.asList(matchingReceiver,
 nonMatchingReceiver));
         
         final CountDownLatch latch = new CountDownLatch(1);
         doAnswer(invocation -> {
@@ -418,7 +417,7 @@ public class AlertDispatchServiceTest {
         requiredLabels.put("service", "gateway");
         receiverWithLabels.setLabels(requiredLabels);
 
-        
when(alertReceiverMapper.selectAll()).thenReturn(Arrays.asList(receiverWithLabels));
+        
when(alertReceiverMapper.selectByNamespaceId(TEST_NAMESPACE_ID)).thenReturn(Arrays.asList(receiverWithLabels));
 
         alertDispatchService.dispatchAlert(alarmContent);
         
@@ -433,7 +432,7 @@ public class AlertDispatchServiceTest {
         final AlarmContent alarmContent = createTestAlarmContent();
         final AlertReceiverDTO receiver = createTestReceiver(EMAIL_TYPE, true, 
false);
         
-        
when(alertReceiverMapper.selectAll()).thenReturn(Arrays.asList(receiver));
+        
when(alertReceiverMapper.selectByNamespaceId(TEST_NAMESPACE_ID)).thenReturn(Arrays.asList(receiver));
         
         final CountDownLatch latch1 = new CountDownLatch(1);
         final CountDownLatch latch2 = new CountDownLatch(1);
@@ -455,7 +454,7 @@ public class AlertDispatchServiceTest {
         assertTrue(latch2.await(5, TimeUnit.SECONDS));
         
         // verify mapper is called only once (cache is used for second call)
-        verify(alertReceiverMapper, times(1)).selectAll();
+        verify(alertReceiverMapper, 
times(1)).selectByNamespaceId(TEST_NAMESPACE_ID);
         verify(emailHandler, times(2)).send(eq(receiver), eq(alarmContent));
     }
 
@@ -503,13 +502,13 @@ public class AlertDispatchServiceTest {
     }
 
     @SuppressWarnings("unchecked")
-    private AtomicReference<List<AlertReceiverDTO>> 
getAlertReceiverReference() {
+    private ConcurrentMap<String, List<AlertReceiverDTO>> 
getAlertReceiverCache() {
         try {
-            Field field = 
AlertDispatchServiceImpl.class.getDeclaredField("alertReceiverReference");
+            Field field = 
AlertDispatchServiceImpl.class.getDeclaredField("alertReceiverCache");
             field.setAccessible(true);
-            return (AtomicReference<List<AlertReceiverDTO>>) 
field.get(alertDispatchService);
+            return (ConcurrentMap<String, List<AlertReceiverDTO>>) 
field.get(alertDispatchService);
         } catch (Exception e) {
-            throw new RuntimeException("Failed to get alertReceiverReference 
field", e);
+            throw new RuntimeException("Failed to get alertReceiverCache 
field", e);
         }
     }
 

Reply via email to