This is an automated email from the ASF dual-hosted git repository.
lizhimins pushed a commit to branch rocketmq-studio
in repository https://gitbox.apache.org/repos/asf/rocketmq-dashboard.git
The following commit(s) were added to refs/heads/rocketmq-studio by this push:
new 2fe074643 fix(alert): clear stale notification deliveries and stop
double-encoding DingTalk signatures (#5048)
2fe074643 is described below
commit 2fe0746438675b74278375cd33d57174f28e05d9
Author: zmuxuny <[email protected]>
AuthorDate: Fri Oct 9 23:37:57 2026 -0700
fix(alert): clear stale notification deliveries and stop double-encoding
DingTalk signatures (#5048)
Two defects on the notification path.
- The deliveries page only toasted when a filtered list load failed,
leaving the previous filter's rows
and their retry buttons on screen; all eleven `loadFailed` consumers on
the page now clear the stale
rows, matching the existing `AlertRuleAssetList` pattern. (#5048)
- `NotificationOutboxService` handed an already URL-escaped DingTalk
signature to
`postForEntity(String, ...)`, whose `TEMPLATE_AND_VALUES` encoding turned
`%3D` into `%253D` and so
broke the webhook signature. (#5493)
`NotificationDeliveriesPage` 7 tests and `NotificationWebhookEncodingTest`
10 tests green,
0 checkstyle violations. The new delivery-page case was rebased alongside
#5034's filter test rather
than replacing it.
Folded in #5493 (same author).
---
.../ops/alert/NotificationOutboxService.java | 5 +-
.../ops/alert/NotificationWebhookEncodingTest.java | 167 +++++++++++++++++++++
.../__tests__/NotificationDeliveriesPage.test.tsx | 62 +++++++-
web/src/pages/ops/notificationDeliveries.tsx | 39 ++++-
4 files changed, 266 insertions(+), 7 deletions(-)
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/ops/alert/NotificationOutboxService.java
b/server/src/main/java/org/apache/rocketmq/studio/ops/alert/NotificationOutboxService.java
index 760a5c974..955eb3178 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/ops/alert/NotificationOutboxService.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/ops/alert/NotificationOutboxService.java
@@ -50,6 +50,7 @@ import java.util.concurrent.ScheduledFuture;
import java.util.concurrent.ThreadFactory;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
+import java.net.URI;
import java.net.URLEncoder;
import java.nio.charset.StandardCharsets;
import java.util.Base64;
@@ -421,7 +422,9 @@ public class NotificationOutboxService {
throw new IllegalStateException("No configured " + channel + "
webhook");
}
UrlHostGuard.check(webhook, false);
- ResponseEntity<String> response =
restTemplate.postForEntity(dingTalkWebhook(webhook, settings, channel),
+ // The URL already contains escaped query values, including the
DingTalk signature.
+ // The String overload would treat it as a URI template and encode
those escapes again.
+ ResponseEntity<String> response =
restTemplate.postForEntity(URI.create(dingTalkWebhook(webhook, settings,
channel)),
payload(alert, channel, content), String.class);
if (!response.getStatusCode().is2xxSuccessful()) {
throw new IllegalStateException("Webhook returned " +
response.getStatusCode());
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/ops/alert/NotificationWebhookEncodingTest.java
b/server/src/test/java/org/apache/rocketmq/studio/ops/alert/NotificationWebhookEncodingTest.java
new file mode 100644
index 000000000..2fb21e751
--- /dev/null
+++
b/server/src/test/java/org/apache/rocketmq/studio/ops/alert/NotificationWebhookEncodingTest.java
@@ -0,0 +1,167 @@
+/*
+ * 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.rocketmq.studio.ops.alert;
+
+import java.net.URI;
+import java.net.URLDecoder;
+import java.net.URLEncoder;
+import java.nio.charset.StandardCharsets;
+import java.time.LocalDateTime;
+import java.util.Arrays;
+import java.util.Base64;
+import java.util.List;
+import java.util.Map;
+import java.util.Optional;
+import java.util.concurrent.atomic.AtomicReference;
+import java.util.stream.Collectors;
+import java.util.stream.Stream;
+import javax.crypto.Mac;
+import javax.crypto.spec.SecretKeySpec;
+import org.apache.rocketmq.studio.audit.OperationAuditService;
+import org.apache.rocketmq.studio.common.domain.enums.AlertLevel;
+import
org.apache.rocketmq.studio.persistence.entity.RmqAlertNotificationOutbox;
+import
org.apache.rocketmq.studio.persistence.mapper.RmqAlertNotificationOutboxMapper;
+import org.apache.rocketmq.studio.settings.GeneralSettingsVO;
+import org.apache.rocketmq.studio.settings.SettingsRepository;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.Arguments;
+import org.junit.jupiter.params.provider.MethodSource;
+import org.junit.jupiter.params.provider.ValueSource;
+import org.springframework.http.HttpMethod;
+import org.springframework.http.MediaType;
+import org.springframework.test.web.client.MockRestServiceServer;
+import org.springframework.web.client.RestTemplate;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+import static
org.springframework.test.web.client.match.MockRestRequestMatchers.method;
+import static
org.springframework.test.web.client.match.MockRestRequestMatchers.requestTo;
+import static
org.springframework.test.web.client.response.MockRestResponseCreators.withSuccess;
+
+class NotificationWebhookEncodingTest {
+ private static final String WEBHOOK = "https://192.0.2.1/robot/send";
+ private static final String SYNTHETIC_SECRET = "synthetic-test-secret";
+
+ private final RmqAlertNotificationOutboxMapper mapper =
mock(RmqAlertNotificationOutboxMapper.class);
+ private final SettingsRepository settings = mock(SettingsRepository.class);
+ private final AlertRepository alerts = mock(AlertRepository.class);
+ private final OperationAuditService audit =
mock(OperationAuditService.class);
+ private final RestTemplate client = new RestTemplate();
+ private final MockRestServiceServer server =
MockRestServiceServer.bindTo(client).build();
+ private final NotificationOutboxService service = new
NotificationOutboxService(mapper, settings,
+ mock(AlertSilenceService.class), alerts, audit, client);
+
+ @AfterEach
+ void close() {
+ service.closeHeartbeatExecutor();
+ }
+
+ @ParameterizedTest
+ @ValueSource(strings = {"", "?access_token=synthetic",
"?access_token=a%2Bb%2Fc%3D"})
+ void signsTestNotificationsWithExactlyOneLayerOfEncodingTest(String query)
throws Exception {
+
when(settings.loadGeneralSettings()).thenReturn(GeneralSettingsVO.builder()
+ .dingtalkWebhook(WEBHOOK +
query).dingtalkSigningSecret(SYNTHETIC_SECRET).build());
+ AtomicReference<URI> requestUri = new AtomicReference<>();
+ server.expect(request -> requestUri.set(request.getURI()))
+ .andExpect(method(HttpMethod.POST))
+ .andRespond(withSuccess("{\"errcode\":0}",
MediaType.APPLICATION_JSON));
+
+ service.sendTestMessage("dingtalk");
+
+ server.verify();
+ assertValidSignature(requestUri.get());
+ if (!query.isEmpty()) {
+
assertThat(requestUri.get().getRawQuery()).startsWith(query.substring(1) +
"×tamp=");
+ }
+ }
+
+ @Test
+ void signsQueuedDeliveriesWithExactlyOneLayerOfEncodingTest() throws
Exception {
+
when(settings.loadGeneralSettings()).thenReturn(GeneralSettingsVO.builder()
+ .dingtalkWebhook(WEBHOOK + "?access_token=synthetic")
+ .dingtalkSigningSecret(SYNTHETIC_SECRET).build());
+ RmqAlertNotificationOutbox row = new RmqAlertNotificationOutbox();
+ row.setId(8L);
+ row.setAlertId(9L);
+ row.setChannel("dingtalk");
+ row.setStatus("PENDING");
+ row.setAttemptCount(0);
+ when(mapper.findDispatchable(any(LocalDateTime.class),
any(LocalDateTime.class), any(Integer.class)))
+ .thenReturn(List.of(row));
+ when(mapper.claimForDispatch(any(), any(LocalDateTime.class),
any(LocalDateTime.class),
+ any(LocalDateTime.class), anyString())).thenReturn(1);
+ when(mapper.update(any(), any())).thenReturn(1);
+
when(alerts.findAlertById(9L)).thenReturn(Optional.of(SystemAlertVO.builder().id(9L)
+
.level(AlertLevel.warning).title("Lag").description("high").instanceId("local").build()));
+ AtomicReference<URI> requestUri = new AtomicReference<>();
+ server.expect(request -> requestUri.set(request.getURI()))
+ .andExpect(method(HttpMethod.POST))
+ .andRespond(withSuccess("{\"errcode\":0}",
MediaType.APPLICATION_JSON));
+
+ service.dispatch();
+
+ server.verify();
+ assertValidSignature(requestUri.get());
+ verify(mapper).update(any(), any());
+ verify(audit).record("DELIVER_ALERT_NOTIFICATION",
"ALERT_NOTIFICATION", "8", null,
+ "alertId=9, channel=dingtalk", "SUCCESS", null);
+ }
+
+ @ParameterizedTest
+ @MethodSource("unsignedWebhooks")
+ void preservesUnsignedWebhookUrisTest(String channel, String query) {
+ // A configured DingTalk secret must not add signing parameters to SMS
requests.
+
when(settings.loadGeneralSettings()).thenReturn(GeneralSettingsVO.builder()
+ .dingtalkWebhook(WEBHOOK + query).smsWebhook(WEBHOOK + query)
+ .dingtalkSigningSecret("sms".equals(channel) ?
SYNTHETIC_SECRET : null).build());
+ server.expect(requestTo(WEBHOOK + query))
+ .andExpect(method(HttpMethod.POST))
+ .andRespond(withSuccess("dingtalk".equals(channel) ?
"{\"errcode\":0}" : "accepted",
+ "dingtalk".equals(channel) ?
MediaType.APPLICATION_JSON : MediaType.TEXT_PLAIN));
+
+ service.sendTestMessage(channel);
+
+ server.verify();
+ }
+
+ private static Stream<Arguments> unsignedWebhooks() {
+ return Stream.of("dingtalk", "sms").flatMap(channel -> Stream.of("",
"?token=synthetic",
+ "?token=a%2Bb%2Fc%3D&label=hello%20world").map(query ->
Arguments.of(channel, query)));
+ }
+
+ private static void assertValidSignature(URI uri) throws Exception {
+ Map<String, String> rawQuery =
Arrays.stream(uri.getRawQuery().split("&"))
+ .map(part -> part.split("=", 2))
+ .collect(Collectors.toMap(parts -> parts[0], parts ->
parts[1]));
+ String timestamp = rawQuery.get("timestamp");
+ assertThat(timestamp).matches("[0-9]+");
+ Mac mac = Mac.getInstance("HmacSHA256");
+ mac.init(new
SecretKeySpec(SYNTHETIC_SECRET.getBytes(StandardCharsets.UTF_8), "HmacSHA256"));
+ String expected = Base64.getEncoder().encodeToString(
+ mac.doFinal((timestamp + "\n" +
SYNTHETIC_SECRET).getBytes(StandardCharsets.UTF_8)));
+ // Inspect the real RestTemplate request, not the URL before Spring
processes it.
+ assertThat(URLDecoder.decode(rawQuery.get("sign"),
StandardCharsets.UTF_8)).isEqualTo(expected);
+ assertThat(rawQuery.get("sign")).isEqualTo(URLEncoder.encode(expected,
StandardCharsets.UTF_8));
+ assertThat(rawQuery.get("sign")).endsWith("%3D");
+ }
+}
diff --git a/web/src/pages/ops/__tests__/NotificationDeliveriesPage.test.tsx
b/web/src/pages/ops/__tests__/NotificationDeliveriesPage.test.tsx
index adff33d95..8049d8ec9 100644
--- a/web/src/pages/ops/__tests__/NotificationDeliveriesPage.test.tsx
+++ b/web/src/pages/ops/__tests__/NotificationDeliveriesPage.test.tsx
@@ -24,10 +24,12 @@ vi.mock('../../../services/opsService', () => ({
const deferred = <T,>() => {
let resolve!: (value: T) => void;
- const promise = new Promise<T>((complete) => {
+ let reject!: (reason?: unknown) => void;
+ const promise = new Promise<T>((complete, fail) => {
resolve = complete;
+ reject = fail;
});
- return { promise, resolve };
+ return { promise, resolve, reject };
};
beforeAll(() => {
@@ -227,7 +229,9 @@ describe('NotificationDeliveriesPage', () => {
const secondPage = document.querySelector('.ant-pagination-item-2') as
HTMLElement;
await user.click(secondPage);
await waitFor(() =>
-
expect(listAlertDeliveriesPage).toHaveBeenLastCalledWith(expect.objectContaining({
page: 2 })),
+ expect(listAlertDeliveriesPage).toHaveBeenLastCalledWith(
+ expect.objectContaining({ page: 2 }),
+ ),
);
await user.type(screen.getByPlaceholderText('搜索告警标题或失败原因'),
'WebHook{Enter}');
@@ -277,4 +281,56 @@ describe('NotificationDeliveriesPage', () => {
vi.unstubAllEnvs();
}
});
+
+ it('does not expose stale delivery actions after a filtered list load
fails', async () => {
+ const user = userEvent.setup({ pointerEventsCheck: 0 });
+ render(
+ <App>
+ <LangProvider>
+ <NotificationDeliveriesPage />
+ </LangProvider>
+ </App>,
+ );
+
+ await screen.findByText('Broker disk usage');
+ const filteredLoad = deferred<Awaited<ReturnType<typeof
listAlertDeliveriesPage>>>();
+ vi.mocked(listAlertDeliveriesPage)
+ .mockImplementationOnce(() => filteredLoad.promise)
+ .mockResolvedValueOnce({
+ items: [
+ {
+ id: 8,
+ alertId: 4,
+ alertTitle: 'Delivered notification',
+ channel: 'email',
+ status: 'DELIVERED',
+ attemptCount: 1,
+ createdAt: '2026-08-23T10:00:00',
+ deliveredAt: '2026-08-23T10:01:00',
+ },
+ ],
+ total: 1,
+ page: 1,
+ size: 20,
+ });
+ await user.click(screen.getAllByRole('combobox')[1]);
+ await user.click(await screen.findByText('DELIVERED'));
+
+ await waitFor(() =>
+ expect(listAlertDeliveriesPage).toHaveBeenCalledWith(
+ expect.objectContaining({ status: 'DELIVERED' }),
+ ),
+ );
+ expect(screen.getByRole('button', { name: '重新投递' })).toBeDisabled();
+ await act(async () => filteredLoad.reject(new Error('delivery list
unavailable')));
+ expect(await screen.findByText('告警投递记录加载失败,请稍后重试')).toBeInTheDocument();
+ expect(screen.queryByText('Broker disk usage')).not.toBeInTheDocument();
+ expect(screen.queryByRole('button', { name: '重新投递'
})).not.toBeInTheDocument();
+
+ await user.click(screen.getByRole('button', { name: /^重\s*试$/ }));
+ expect(await screen.findByText('Delivered
notification')).toBeInTheDocument();
+ expect(listAlertDeliveriesPage).toHaveBeenLastCalledWith(
+ expect.objectContaining({ status: 'DELIVERED' }),
+ );
+ });
});
diff --git a/web/src/pages/ops/notificationDeliveries.tsx
b/web/src/pages/ops/notificationDeliveries.tsx
index 7bc3e1d07..1e8fdf0ae 100644
--- a/web/src/pages/ops/notificationDeliveries.tsx
+++ b/web/src/pages/ops/notificationDeliveries.tsx
@@ -6,6 +6,7 @@
*/
import { useEffect, useRef, useState } from 'react';
import {
+ Alert,
Button,
Card,
DatePicker,
@@ -58,6 +59,7 @@ const NotificationDeliveriesPage = () => {
const [search, setSearch] = useState<string>();
const [timeRange, setTimeRange] = useState<{ from?: string; to?: string
}>({});
const [instanceLoadFailed, setInstanceLoadFailed] = useState(false);
+ const [loadFailed, setLoadFailed] = useState(false);
const [instanceReloadNonce, setInstanceReloadNonce] = useState(0);
const [selectedDelivery, setSelectedDelivery] =
useState<NotificationDeliveryRecord>();
const [retryingIds, setRetryingIds] = useState<Set<number>>(() => new Set());
@@ -68,10 +70,12 @@ const NotificationDeliveriesPage = () => {
const refresh = () => {
setLoading(true);
+ setLoadFailed(false);
setRefreshNonce((current) => current + 1);
};
const retryDelivery = async (record: NotificationDeliveryRecord) => {
+ if (loading || loadFailed) return;
if (retryingVisibleInFlight.current ||
retryingIdsInFlight.current.has(record.id)) return;
retryingIdsInFlight.current.add(record.id);
setRetryingIds((current) => new Set(current).add(record.id));
@@ -97,6 +101,7 @@ const NotificationDeliveriesPage = () => {
};
const retryVisibleFailures = async () => {
+ if (loading || loadFailed) return;
const ids = items.filter((item) => item.status === 'FAILED').map((item) =>
item.id);
if (ids.length === 0) return;
if (retryingVisibleInFlight.current || ids.some((id) =>
retryingIdsInFlight.current.has(id)))
@@ -141,14 +146,27 @@ const NotificationDeliveriesPage = () => {
useEffect(() => {
let cancelled = false;
- void listAlertDeliveriesPage({ channel, status, instanceId, search,
...timeRange, page, pageSize })
+ void listAlertDeliveriesPage({
+ channel,
+ status,
+ instanceId,
+ search,
+ ...timeRange,
+ page,
+ pageSize,
+ })
.then((result) => {
if (cancelled) return;
setItems(result.items);
setTotal(result.total);
+ setLoadFailed(false);
})
.catch(() => {
- if (!cancelled) message.error(t('deliveries.loadFailed'));
+ if (cancelled) return;
+ setItems([]);
+ setTotal(0);
+ setSelectedDelivery(undefined);
+ setLoadFailed(true);
})
.finally(() => {
if (!cancelled) setLoading(false);
@@ -160,6 +178,8 @@ const NotificationDeliveriesPage = () => {
const resetPage = (change: () => void) => {
setLoading(true);
+ setLoadFailed(false);
+ setSelectedDelivery(undefined);
change();
setPage(1);
};
@@ -232,6 +252,7 @@ const NotificationDeliveriesPage = () => {
size="small"
icon={<Eye size={18} />}
aria-label={t('deliveries.viewDetails')}
+ disabled={loading || loadFailed}
onClick={() => setSelectedDelivery(record)}
/>
{record.status === 'FAILED' && (
@@ -242,6 +263,7 @@ const NotificationDeliveriesPage = () => {
icon={<ArrowClockwise size={18} />}
aria-label={t('deliveries.retry')}
loading={retryingIds.has(record.id)}
+ disabled={loading || loadFailed}
onClick={() => void retryDelivery(record)}
/>
</Tooltip>
@@ -259,7 +281,7 @@ const NotificationDeliveriesPage = () => {
<Flex gap={12} wrap="wrap" style={{ marginBottom: 20 }}>
<Button
icon={<ArrowClockwise size={18} />}
- disabled={!items.some((item) => item.status === 'FAILED')}
+ disabled={loading || loadFailed || !items.some((item) =>
item.status === 'FAILED')}
loading={retryingVisible}
onClick={() => void retryVisibleFailures()}
>
@@ -339,6 +361,15 @@ const NotificationDeliveriesPage = () => {
</Tooltip>
)}
</Flex>
+ {loadFailed && (
+ <Alert
+ type="error"
+ showIcon
+ message={t('deliveries.loadFailed')}
+ action={<Button onClick={refresh}>{t('common.retry')}</Button>}
+ style={{ marginBottom: 16 }}
+ />
+ )}
<Table
rowKey="id"
columns={columns}
@@ -354,6 +385,7 @@ const NotificationDeliveriesPage = () => {
showTotal: (count) => `${t('common.total')} ${count}`,
onChange: (nextPage, nextPageSize) => {
setLoading(true);
+ setLoadFailed(false);
setPage(nextPage);
setPageSize(nextPageSize);
},
@@ -409,6 +441,7 @@ const NotificationDeliveriesPage = () => {
<Button
icon={<ArrowClockwise size={18} />}
loading={retryingIds.has(selectedDelivery.id)}
+ disabled={loading || loadFailed}
onClick={() => void retryDelivery(selectedDelivery)}
>
{t('deliveries.retry')}