Compare commits

...
Author SHA1 Message Date
tomsun28 849f0b4887 Merge branch 'master' into fix/alert-80 2025-05-18 16:55:10 +08:00
dc8ae844c1 [improve] improve url validation for serverChan (#3364)
Signed-off-by: aias00 <liuhongyu@apache.org>
Co-authored-by: Copilot Autofix powered by AI <62310815+github-advanced-security[bot]@users.noreply.github.com>
Co-authored-by: Calvin <naruse_shinji@163.com>
Co-authored-by: tomsun28 <tomsun28@outlook.com>
2025-05-18 16:48:31 +08:00
91a7593b87 [improve] improve url validation for SlackAlertNotifyHandlerImpl (#3363)
Signed-off-by: aias00 <liuhongyu@apache.org>
Co-authored-by: Copilot Autofix powered by AI <62310815+github-advanced-security[bot]@users.noreply.github.com>
Co-authored-by: Calvin <naruse_shinji@163.com>
Co-authored-by: tomsun28 <tomsun28@outlook.com>
2025-05-18 16:17:31 +08:00
4442a52adf [improve] improve url validation for TelegramBotAlertNotifyHandlerImpl (#3362)
Signed-off-by: aias00 <liuhongyu@apache.org>
Co-authored-by: Copilot Autofix powered by AI <62310815+github-advanced-security[bot]@users.noreply.github.com>
Co-authored-by: tomsun28 <tomsun28@outlook.com>
2025-05-18 14:09:09 +08:00
Calvin cf544ffd64 Merge branch 'master' into fix/alert-80 2025-05-18 13:48:35 +08:00
60e3437e82 [improve] improve url validation for WeComRobotAlertNotifyHandlerImpl (#3361)
Signed-off-by: aias00 <liuhongyu@apache.org>
Co-authored-by: Copilot Autofix powered by AI <62310815+github-advanced-security[bot]@users.noreply.github.com>
Co-authored-by: tomsun28 <tomsun28@outlook.com>
2025-05-18 13:31:39 +08:00
21d6ec2f1b [bugfix] Incorrect SD sub-monitor status (#3340)
Signed-off-by: Sherlock Yin <sherlock.yin1994@gmail.com>
Co-authored-by: yinyijun <yinyijun6@mgtv.com>
Co-authored-by: tomsun28 <tomsun28@outlook.com>
Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com>
Co-authored-by: Calvin <naruse_shinji@163.com>
Co-authored-by: aias00 <liuhongyu@apache.org>
2025-05-18 01:00:37 +08:00
aias00andCopilot Autofix powered by AI ffc7a6d1e0 Potential fix for code scanning alert no. 80: Server-side request forgery
Co-authored-by: Copilot Autofix powered by AI <62310815+github-advanced-security[bot]@users.noreply.github.com>
Signed-off-by: aias00 <liuhongyu@apache.org>
2025-05-17 14:30:25 +08:00
12 changed files with 162 additions and 94 deletions
@@ -54,7 +54,15 @@ final class DingTalkRobotAlertNotifyHandlerImpl extends AbstractAlertNotifyHandl
HttpHeaders headers = new HttpHeaders();
headers.setContentType(MediaType.APPLICATION_JSON);
HttpEntity<DingTalkWebHookDto> httpEntity = new HttpEntity<>(dingTalkWebHookDto, headers);
String webHookUrl = alerterProperties.getDingTalkWebhookUrl() + receiver.getAccessToken();
String baseUrl = alerterProperties.getDingTalkWebhookUrl();
String accessToken = receiver.getAccessToken();
if (StringUtils.isBlank(accessToken) || !accessToken.matches("^[a-zA-Z0-9_-]+$")) {
throw new AlertNoticeException("Invalid access token provided for DingTalk webhook.");
}
String webHookUrl = baseUrl + accessToken;
if (!webHookUrl.startsWith(baseUrl)) {
throw new AlertNoticeException("Constructed webhook URL does not match the trusted base URL.");
}
ResponseEntity<CommonRobotNotifyResp> responseEntity = restTemplate.postForEntity(webHookUrl,
httpEntity, CommonRobotNotifyResp.class);
if (responseEntity.getStatusCode() == HttpStatus.OK) {
@@ -31,6 +31,8 @@ import org.springframework.http.MediaType;
import org.springframework.http.ResponseEntity;
import org.springframework.stereotype.Component;
import java.util.List;
/**
* Send alarm information through Server
*/
@@ -54,7 +56,16 @@ public class ServerChanAlertNotifyHandlerImpl extends AbstractAlertNotifyHandler
HttpHeaders headers = new HttpHeaders();
headers.setContentType(MediaType.APPLICATION_JSON);
HttpEntity<ServerChanAlertNotifyHandlerImpl.ServerChanWebHookDto> httpEntity = new HttpEntity<>(serverChanWebHookDto, headers);
String webHookUrl = String.format(alerterProperties.getServerChanWebhookUrl(), receiver.getServerChanToken());
String sanitizedToken = receiver.getServerChanToken().replaceAll("[^a-zA-Z0-9_-]", "");
String webHookUrl = String.format(alerterProperties.getServerChanWebhookUrl(), sanitizedToken);
// Validate the constructed URL against a whitelist
List<String> allowedBaseUrls = List.of("https://api.serverchan.com", "https://serverchan.example.com");
boolean isValidUrl = allowedBaseUrls.stream().anyMatch(webHookUrl::startsWith);
if (!isValidUrl) {
throw new AlertNoticeException("Invalid webhook URL: " + webHookUrl);
}
ResponseEntity<CommonRobotNotifyResp> responseEntity = restTemplate.postForEntity(webHookUrl,
httpEntity, CommonRobotNotifyResp.class);
if (responseEntity.getStatusCode() == HttpStatus.OK) {
@@ -17,6 +17,7 @@
package org.apache.hertzbeat.alert.notice.impl;
import java.net.URI;
import java.util.Objects;
import lombok.Builder;
import lombok.Data;
@@ -52,7 +53,12 @@ final class SlackAlertNotifyHandlerImpl extends AbstractAlertNotifyHandlerImpl {
HttpHeaders headers = new HttpHeaders();
headers.setContentType(MediaType.APPLICATION_JSON);
HttpEntity<SlackNotifyDTO> slackNotifyEntity = new HttpEntity<>(slackNotify, headers);
var entity = restTemplate.postForEntity(receiver.getSlackWebHookUrl(), slackNotifyEntity, String.class);
String slackWebHookUrl = receiver.getSlackWebHookUrl();
if (!isValidSlackWebHookUrl(slackWebHookUrl)) {
log.warn("Invalid Slack Webhook URL: {}", slackWebHookUrl);
throw new AlertNoticeException("Invalid Slack Webhook URL");
}
var entity = restTemplate.postForEntity(slackWebHookUrl, slackNotifyEntity, String.class);
if (entity.getStatusCode() == HttpStatus.OK && entity.getBody() != null) {
var body = entity.getBody();
if (Objects.equals(SUCCESS, body)) {
@@ -81,4 +87,21 @@ final class SlackAlertNotifyHandlerImpl extends AbstractAlertNotifyHandlerImpl {
private String text;
}
/**
* Validate if the Slack Webhook URL belongs to an allowed domain.
*
* @param url the Slack Webhook URL to validate
* @return true if the URL is valid, false otherwise
*/
private boolean isValidSlackWebHookUrl(String url) {
try {
URI uri = new URI(url);
String host = uri.getHost();
return "hooks.slack.com".equals(host);
} catch (Exception e) {
log.warn("Error validating Slack Webhook URL: {}", url, e);
return false;
}
}
}
@@ -42,9 +42,14 @@ import org.springframework.stereotype.Component;
final class TelegramBotAlertNotifyHandlerImpl extends AbstractAlertNotifyHandlerImpl {
@Override
public void send(NoticeReceiver receiver, NoticeTemplate noticeTemplate, GroupAlert alert) throws AlertNoticeException {
public void send(NoticeReceiver receiver, NoticeTemplate noticeTemplate, GroupAlert alert)
throws AlertNoticeException {
try {
String url = String.format(alerterProperties.getTelegramWebhookUrl(), receiver.getTgBotToken());
String token = receiver.getTgBotToken();
if (!isValidTelegramToken(token)) {
throw new AlertNoticeException("Invalid Telegram Bot Token");
}
String url = String.format(alerterProperties.getTelegramWebhookUrl(), token);
TelegramBotNotifyDTO notifyBody = TelegramBotNotifyDTO.builder()
.chatId(receiver.getTgUserId())
.text(renderContent(noticeTemplate, alert))
@@ -54,7 +59,8 @@ final class TelegramBotAlertNotifyHandlerImpl extends AbstractAlertNotifyHandler
HttpHeaders headers = new HttpHeaders();
headers.setContentType(MediaType.APPLICATION_JSON);
HttpEntity<TelegramBotNotifyDTO> telegramEntity = new HttpEntity<>(notifyBody, headers);
ResponseEntity<TelegramBotNotifyResponse> entity = restTemplate.postForEntity(url, telegramEntity, TelegramBotNotifyResponse.class);
ResponseEntity<TelegramBotNotifyResponse> entity = restTemplate.postForEntity(url, telegramEntity,
TelegramBotNotifyResponse.class);
if (entity.getStatusCode() == HttpStatus.OK && entity.getBody() != null) {
TelegramBotNotifyResponse body = entity.getBody();
if (body.ok) {
@@ -99,4 +105,10 @@ final class TelegramBotAlertNotifyHandlerImpl extends AbstractAlertNotifyHandler
private String description;
}
private boolean isValidTelegramToken(String token) {
// Adjusted pattern to match real Telegram Bot tokens like
// 110201543:AAHdqTcvCH1vGWJxfSeofSAs0K5PALDsaw
String tokenPattern = "^[0-9]+:[a-zA-Z0-9_-]+$";
return token != null && token.matches(tokenPattern);
}
}
@@ -57,7 +57,12 @@ final class WeComRobotAlertNotifyHandlerImpl extends AbstractAlertNotifyHandlerI
HttpHeaders headers = new HttpHeaders();
headers.setContentType(MediaType.APPLICATION_JSON);
HttpEntity<WeWorkWebHookDto> httpEntity = new HttpEntity<>(weWorkWebHookDTO, headers);
String webHookUrl = alerterProperties.getWeWorkWebhookUrl() + receiver.getWechatId();
String wechatId = receiver.getWechatId();
if (!isValidWechatId(wechatId)) {
log.warn("Invalid WeChat ID: {}", wechatId);
throw new AlertNoticeException("Invalid WeChat ID provided.");
}
String webHookUrl = alerterProperties.getWeWorkWebhookUrl() + wechatId;
ResponseEntity<CommonRobotNotifyResp> entity = restTemplate.postForEntity(webHookUrl, httpEntity, CommonRobotNotifyResp.class);
if (entity.getStatusCode() == HttpStatus.OK) {
assert entity.getBody() != null;
@@ -170,4 +175,15 @@ final class WeComRobotAlertNotifyHandlerImpl extends AbstractAlertNotifyHandlerI
}
}
/**
* Validate the WeChat ID to ensure it meets the expected format.
*
* @param wechatId the WeChat ID to validate
* @return true if valid, false otherwise
*/
private boolean isValidWechatId(String wechatId) {
// Example validation: ensure the ID is alphanumeric and non-empty
return StringUtils.isNotBlank(wechatId) && wechatId.matches("^[a-zA-Z0-9_-]+$");
}
}
@@ -20,6 +20,7 @@ package org.apache.hertzbeat.alert.notice.impl;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.Mockito.lenient;
import static org.mockito.Mockito.when;
import org.apache.hertzbeat.alert.AlerterProperties;
@@ -48,7 +49,7 @@ import java.util.ResourceBundle;
*/
@ExtendWith(MockitoExtension.class)
class ServerChanAlertNotifyHandlerImplTest {
@Mock
private RestTemplate restTemplate;
@@ -57,20 +58,21 @@ class ServerChanAlertNotifyHandlerImplTest {
@Mock
private ResourceBundle bundle;
@InjectMocks
private ServerChanAlertNotifyHandlerImpl serverChanAlertNotifyHandler;
private NoticeReceiver receiver;
private GroupAlert groupAlert;
private NoticeTemplate template;
@BeforeEach
public void setUp() {
receiver = new NoticeReceiver();
receiver.setId(1L);
receiver.setName("test-receiver");
receiver.setAccessToken("test-token");
receiver.setServerChanToken("SCT193569TSNm6xIabdjqeZPtOGOWcvU1e");
groupAlert = new GroupAlert();
SingleAlert singleAlert = new SingleAlert();
@@ -87,43 +89,31 @@ class ServerChanAlertNotifyHandlerImplTest {
template.setName("test-template");
template.setContent("test content");
when(alerterProperties.getServerChanWebhookUrl()).thenReturn("http://test.url/");
when(bundle.getString("alerter.notify.title")).thenReturn("Alert Notification");
lenient().when(alerterProperties.getServerChanWebhookUrl())
.thenReturn("https://api.serverchan.com/send/%s");
lenient().when(bundle.getString("alerter.notify.title")).thenReturn("Alert Notification");
}
@Test
public void testNotifyAlertSuccess() {
CommonRobotNotifyResp successResp = new CommonRobotNotifyResp();
successResp.setErrCode(0);
successResp.setMsg("success");
ResponseEntity<CommonRobotNotifyResp> responseEntity =
new ResponseEntity<>(successResp, HttpStatus.OK);
ResponseEntity<CommonRobotNotifyResp> responseEntity = new ResponseEntity<>(successResp, HttpStatus.OK);
when(restTemplate.postForEntity(
any(String.class),
any(),
eq(CommonRobotNotifyResp.class)
)).thenReturn(responseEntity);
eq(CommonRobotNotifyResp.class))).thenReturn(responseEntity);
serverChanAlertNotifyHandler.send(receiver, template, groupAlert);
}
@Test
public void testNotifyAlertFailure() {
CommonRobotNotifyResp failResp = new CommonRobotNotifyResp();
failResp.setCode(1);
failResp.setErrMsg("Test Error");
ResponseEntity<CommonRobotNotifyResp> responseEntity =
new ResponseEntity<>(failResp, HttpStatus.BAD_REQUEST);
when(restTemplate.postForEntity(
any(String.class),
any(),
eq(CommonRobotNotifyResp.class)
)).thenReturn(responseEntity);
public void testNotifyAlertWithInvalidUrl() {
when(alerterProperties.getServerChanWebhookUrl()).thenReturn("http://invalid-url.com/%s");
assertThrows(AlertNoticeException.class,
assertThrows(AlertNoticeException.class,
() -> serverChanAlertNotifyHandler.send(receiver, template, groupAlert));
}
}
@@ -51,7 +51,7 @@ class SlackAlertNotifyHandlerImplTest {
@Mock
private RestTemplate restTemplate;
@Mock
private ResourceBundle bundle;
@@ -68,18 +68,18 @@ class SlackAlertNotifyHandlerImplTest {
receiver.setId(1L);
receiver.setName("test-receiver");
receiver.setAccessToken("test-token");
receiver.setSlackWebHookUrl("http://localhost:8080");
receiver.setSlackWebHookUrl("https://hooks.slack.com/services/ABCDEF/GHIJKL/mnopqrstuvwxyz");
groupAlert = new GroupAlert();
SingleAlert singleAlert = new SingleAlert();
singleAlert.setLabels(new HashMap<>());
singleAlert.getLabels().put("severity", "critical");
singleAlert.getLabels().put("alertname", "Test Alert");
List<SingleAlert> alerts = new ArrayList<>();
alerts.add(singleAlert);
groupAlert.setAlerts(alerts);
template = new NoticeTemplate();
template.setId(1L);
template.setName("test-template");
@@ -90,22 +90,18 @@ class SlackAlertNotifyHandlerImplTest {
@Test
public void testNotifyAlertSuccess() {
ResponseEntity<String> responseEntity =
new ResponseEntity<>("ok", HttpStatus.OK);
ResponseEntity<String> responseEntity = new ResponseEntity<>("ok", HttpStatus.OK);
when(restTemplate.postForEntity(any(String.class), any(), eq(String.class))).thenReturn(responseEntity);
slackAlertNotifyHandler.send(receiver, template, groupAlert);
}
@Test
public void testNotifyAlertFailure() {
ResponseEntity<String> responseEntity =
new ResponseEntity<>("invalid_payload", HttpStatus.BAD_REQUEST);
public void testNotifyAlertWithInvalidUrl() {
receiver.setSlackWebHookUrl("http://localhost:8080");
when(restTemplate.postForEntity(any(String.class), any(), eq(String.class))).thenReturn(responseEntity);
assertThrows(AlertNoticeException.class,
assertThrows(AlertNoticeException.class,
() -> slackAlertNotifyHandler.send(receiver, template, groupAlert));
}
}
@@ -20,6 +20,7 @@ package org.apache.hertzbeat.alert.notice.impl;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.Mockito.lenient;
import static org.mockito.Mockito.when;
import org.apache.hertzbeat.alert.AlerterProperties;
@@ -51,10 +52,10 @@ class TelegramBotAlertNotifyHandlerImplTest {
@Mock
private RestTemplate restTemplate;
@Mock
private AlerterProperties alerterProperties;
@Mock
private ResourceBundle bundle;
@@ -70,58 +71,57 @@ class TelegramBotAlertNotifyHandlerImplTest {
receiver = new NoticeReceiver();
receiver.setId(1L);
receiver.setName("test-receiver");
receiver.setAccessToken("test-token");
receiver.setTgBotToken("123456:ABC-DEF1234ghIkl-zyx57W2v1u123ew11");
receiver.setTgUserId("123456789"); // Telegram specific - chat ID
groupAlert = new GroupAlert();
SingleAlert singleAlert = new SingleAlert();
singleAlert.setLabels(new HashMap<>());
singleAlert.getLabels().put("severity", "critical");
singleAlert.getLabels().put("alertname", "Test Alert");
List<SingleAlert> alerts = new ArrayList<>();
alerts.add(singleAlert);
groupAlert.setAlerts(alerts);
template = new NoticeTemplate();
template.setId(1L);
template.setName("test-template");
template.setContent("test content");
when(alerterProperties.getTelegramWebhookUrl()).thenReturn("https://api.telegram.org/bot%s/sendMessage");
when(bundle.getString("alerter.notify.title")).thenReturn("Alert Notification");
lenient().when(alerterProperties.getTelegramWebhookUrl())
.thenReturn("https://api.telegram.org/bot%s/sendMessage");
lenient().when(bundle.getString("alerter.notify.title")).thenReturn("Alert Notification");
}
@Test
public void testNotifyAlertSuccess() {
TelegramBotAlertNotifyHandlerImpl.TelegramBotNotifyResponse successResp =
new TelegramBotAlertNotifyHandlerImpl.TelegramBotNotifyResponse();
TelegramBotAlertNotifyHandlerImpl.TelegramBotNotifyResponse successResp = new TelegramBotAlertNotifyHandlerImpl.TelegramBotNotifyResponse();
successResp.setOk(true);
successResp.setDescription("Test Success");
ResponseEntity<TelegramBotAlertNotifyHandlerImpl.TelegramBotNotifyResponse> responseEntity =
new ResponseEntity<>(successResp, HttpStatus.OK);
when(restTemplate.postForEntity(any(String.class), any(),
ResponseEntity<TelegramBotAlertNotifyHandlerImpl.TelegramBotNotifyResponse> responseEntity = new ResponseEntity<>(
successResp, HttpStatus.OK);
when(restTemplate.postForEntity(any(String.class), any(),
eq(TelegramBotAlertNotifyHandlerImpl.TelegramBotNotifyResponse.class))).thenReturn(responseEntity);
telegramBotAlertNotifyHandler.send(receiver, template, groupAlert);
}
@Test
public void testNotifyAlertFailure() {
TelegramBotAlertNotifyHandlerImpl.TelegramBotNotifyResponse successResp =
new TelegramBotAlertNotifyHandlerImpl.TelegramBotNotifyResponse();
successResp.setOk(false);
successResp.setDescription("Test failed");
TelegramBotAlertNotifyHandlerImpl.TelegramBotNotifyResponse failureResp = new TelegramBotAlertNotifyHandlerImpl.TelegramBotNotifyResponse();
failureResp.setOk(false);
failureResp.setDescription("Test failed");
ResponseEntity<TelegramBotAlertNotifyHandlerImpl.TelegramBotNotifyResponse> responseEntity =
new ResponseEntity<>(successResp, HttpStatus.BAD_REQUEST);
ResponseEntity<TelegramBotAlertNotifyHandlerImpl.TelegramBotNotifyResponse> responseEntity = new ResponseEntity<>(
failureResp, HttpStatus.OK);
when(restTemplate.postForEntity(any(String.class), any(),
eq(TelegramBotAlertNotifyHandlerImpl.TelegramBotNotifyResponse.class))).thenReturn(responseEntity);
assertThrows(AlertNoticeException.class,
assertThrows(AlertNoticeException.class,
() -> telegramBotAlertNotifyHandler.send(receiver, template, groupAlert));
}
}
@@ -48,7 +48,7 @@ import java.util.ResourceBundle;
*/
@ExtendWith(MockitoExtension.class)
class WeComRobotAlertNotifyHandlerImplTest {
@Mock
private RestTemplate restTemplate;
@@ -57,20 +57,21 @@ class WeComRobotAlertNotifyHandlerImplTest {
@Mock
private ResourceBundle bundle;
@InjectMocks
private WeComRobotAlertNotifyHandlerImpl weComRobotAlertNotifyHandler;
private NoticeReceiver receiver;
private GroupAlert groupAlert;
private NoticeTemplate template;
@BeforeEach
public void setUp() {
receiver = new NoticeReceiver();
receiver.setId(1L);
receiver.setName("test-receiver");
receiver.setAccessToken("test-token");
receiver.setWechatId("test-wechat-id");
groupAlert = new GroupAlert();
SingleAlert singleAlert = new SingleAlert();
@@ -95,33 +96,29 @@ class WeComRobotAlertNotifyHandlerImplTest {
public void testNotifyAlertSuccess() {
CommonRobotNotifyResp successResp = new CommonRobotNotifyResp();
successResp.setErrCode(0);
ResponseEntity<CommonRobotNotifyResp> responseEntity =
new ResponseEntity<>(successResp, HttpStatus.OK);
ResponseEntity<CommonRobotNotifyResp> responseEntity = new ResponseEntity<>(successResp, HttpStatus.OK);
when(restTemplate.postForEntity(
any(String.class),
any(),
eq(CommonRobotNotifyResp.class)
)).thenReturn(responseEntity);
eq(CommonRobotNotifyResp.class))).thenReturn(responseEntity);
weComRobotAlertNotifyHandler.send(receiver, template, groupAlert);
}
@Test
public void testNotifyAlertFailure() {
CommonRobotNotifyResp failResp = new CommonRobotNotifyResp();
failResp.setCode(1);
failResp.setErrMsg("Test Error");
ResponseEntity<CommonRobotNotifyResp> responseEntity =
new ResponseEntity<>(failResp, HttpStatus.OK);
ResponseEntity<CommonRobotNotifyResp> responseEntity = new ResponseEntity<>(failResp, HttpStatus.OK);
when(restTemplate.postForEntity(
any(String.class),
any(),
eq(CommonRobotNotifyResp.class)
)).thenReturn(responseEntity);
eq(CommonRobotNotifyResp.class))).thenReturn(responseEntity);
assertThrows(AlertNoticeException.class,
assertThrows(AlertNoticeException.class,
() -> weComRobotAlertNotifyHandler.send(receiver, template, groupAlert));
}
}
@@ -150,10 +150,7 @@ public class ServiceDiscoveryWorker implements InitializingBean {
// Thus, all monitors still in hostMonitorMap need to be cancelled.
final Set<Long> needCancelMonitorIdSet = subMonitorBindMap.values().stream()
.map(MonitorBind::getMonitorId).collect(Collectors.toSet());
monitorService.cancelManageMonitors(needCancelMonitorIdSet);
for (Long id : needCancelMonitorIdSet) {
monitorBindDao.deleteMonitorBindByBizIdAndMonitorId(monitorId, id);
}
monitorService.deleteMonitors(needCancelMonitorIdSet);
} catch (Exception exception) {
log.error(exception.getMessage(), exception);
}
@@ -18,6 +18,8 @@
package org.apache.hertzbeat.manager.dao;
import java.util.List;
import java.util.Set;
import org.apache.hertzbeat.common.entity.manager.MonitorBind;
import org.springframework.data.jpa.repository.JpaRepository;
import org.springframework.data.jpa.repository.JpaSpecificationExecutor;
@@ -31,9 +33,15 @@ public interface MonitorBindDao extends JpaRepository<MonitorBind, Long>, JpaSpe
List<MonitorBind> findMonitorBindsByBizId(Long bizId);
List<MonitorBind> findMonitorBindsByBizIdIn(Set<Long> bizIds);
void deleteByMonitorId(Long monitorId);
@Modifying
@Transactional
void deleteMonitorBindByBizIdAndMonitorId(Long bizId, Long monitorId);
@Modifying
void deleteMonitorBindByBizIdIn(Set<Long> bizIds);
}
@@ -37,6 +37,7 @@ import org.apache.hertzbeat.common.entity.manager.Collector;
import org.apache.hertzbeat.common.entity.manager.CollectorMonitorBind;
import org.apache.hertzbeat.common.entity.manager.Label;
import org.apache.hertzbeat.common.entity.manager.Monitor;
import org.apache.hertzbeat.common.entity.manager.MonitorBind;
import org.apache.hertzbeat.common.entity.manager.Param;
import org.apache.hertzbeat.common.entity.manager.ParamDefine;
import org.apache.hertzbeat.common.entity.message.CollectRep;
@@ -86,6 +87,7 @@ import java.nio.charset.StandardCharsets;
import java.time.LocalDateTime;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Objects;
@@ -529,12 +531,16 @@ public class MonitorServiceImpl implements MonitorService {
if (CollectionUtils.isEmpty(ids)) {
return;
}
List<Monitor> monitors = monitorDao.findMonitorsByIdIn(ids);
Set<Long> subMonitorIds = monitorBindDao.findMonitorBindsByBizIdIn(ids).stream().map(MonitorBind::getMonitorId).collect(Collectors.toSet());
Set<Long> allMonitorIds = new HashSet<>(ids);
allMonitorIds.addAll(subMonitorIds);
List<Monitor> monitors = monitorDao.findMonitorsByIdIn(allMonitorIds);
if (!monitors.isEmpty()) {
monitorDao.deleteAll(monitors);
paramDao.deleteParamsByMonitorIdIn(ids);
Set<Long> monitorIds = monitors.stream().map(Monitor::getId).collect(Collectors.toSet());
alertDefineBindDao.deleteAlertDefineMonitorBindsByMonitorIdIn(monitorIds);
monitorBindDao.deleteMonitorBindByBizIdIn(monitorIds);
for (Monitor monitor : monitors) {
monitorBindDao.deleteByMonitorId(monitor.getId());
collectorMonitorBindDao.deleteCollectorMonitorBindsByMonitorId(monitor.getId());
@@ -650,6 +656,8 @@ public class MonitorServiceImpl implements MonitorService {
}
// Update monitoring status Delete corresponding monitoring periodic task
// The jobId is not deleted, and the jobId is reused again after the management is started.
Set<Long> subMonitorIds = monitorBindDao.findMonitorBindsByBizIdIn(ids).stream().map(MonitorBind::getMonitorId).collect(Collectors.toSet());
ids.addAll(subMonitorIds);
List<Monitor> managedMonitors = monitorDao.findMonitorsByIdIn(ids)
.stream().filter(monitor ->
monitor.getStatus() != CommonConstants.MONITOR_PAUSED_CODE)
@@ -666,6 +674,8 @@ public class MonitorServiceImpl implements MonitorService {
@Override
public void enableManageMonitors(Set<Long> ids) {
// Update monitoring status Add corresponding monitoring periodic task
Set<Long> subMonitorIds = monitorBindDao.findMonitorBindsByBizIdIn(ids).stream().map(MonitorBind::getMonitorId).collect(Collectors.toSet());
ids.addAll(subMonitorIds);
List<Monitor> unManagedMonitors = monitorDao.findMonitorsByIdIn(ids)
.stream().filter(monitor ->
monitor.getStatus() == CommonConstants.MONITOR_PAUSED_CODE)