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);
}
}