From 96941abc0d306064a0a078c0f566554f17e4901b Mon Sep 17 00:00:00 2001 From: yuhaiming <40624989@qq.com> Date: Fri, 17 Jul 2026 16:28:57 +0800 Subject: [PATCH] fix: clear mqtt pending commands after startup --- .../MqttPendingCommandCleanupListener.java | 25 +++++++++++++ ...MqttPendingCommandCleanupListenerTest.java | 35 +++++++++++++++++++ 2 files changed, 60 insertions(+) create mode 100644 water-common/water-common-mqtt/src/main/java/org/dromara/mqtt/MqttPendingCommandCleanupListener.java create mode 100644 water-common/water-common-mqtt/src/test/java/org/dromara/mqtt/MqttPendingCommandCleanupListenerTest.java diff --git a/water-common/water-common-mqtt/src/main/java/org/dromara/mqtt/MqttPendingCommandCleanupListener.java b/water-common/water-common-mqtt/src/main/java/org/dromara/mqtt/MqttPendingCommandCleanupListener.java new file mode 100644 index 0000000..c587732 --- /dev/null +++ b/water-common/water-common-mqtt/src/main/java/org/dromara/mqtt/MqttPendingCommandCleanupListener.java @@ -0,0 +1,25 @@ +package org.dromara.mqtt; + +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.boot.context.event.ApplicationReadyEvent; +import org.springframework.context.event.EventListener; +import org.springframework.stereotype.Component; + +@Slf4j +@Component +@RequiredArgsConstructor +public class MqttPendingCommandCleanupListener { + + private final MqttCommandAckService mqttCommandAckService; + + @EventListener(ApplicationReadyEvent.class) + public void onApplicationReady(ApplicationReadyEvent event) { + try { + int clearedCount = mqttCommandAckService.clearPendingCommands(); + log.info("[MQTT] 应用启动清理待确认命令完成,清理数量: {}", clearedCount); + } catch (RuntimeException e) { + log.error("[MQTT] 应用启动清理待确认命令失败", e); + } + } +} diff --git a/water-common/water-common-mqtt/src/test/java/org/dromara/mqtt/MqttPendingCommandCleanupListenerTest.java b/water-common/water-common-mqtt/src/test/java/org/dromara/mqtt/MqttPendingCommandCleanupListenerTest.java new file mode 100644 index 0000000..f1e62a4 --- /dev/null +++ b/water-common/water-common-mqtt/src/test/java/org/dromara/mqtt/MqttPendingCommandCleanupListenerTest.java @@ -0,0 +1,35 @@ +package org.dromara.mqtt; + +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.Tag; +import org.springframework.boot.context.event.ApplicationReadyEvent; + +import static org.assertj.core.api.Assertions.assertThatCode; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +@Tag("dev") +class MqttPendingCommandCleanupListenerTest { + + @Test + void onApplicationReadyClearsPendingCommands() { + MqttCommandAckService service = mock(MqttCommandAckService.class); + when(service.clearPendingCommands()).thenReturn(3); + MqttPendingCommandCleanupListener listener = new MqttPendingCommandCleanupListener(service); + + listener.onApplicationReady(mock(ApplicationReadyEvent.class)); + + verify(service).clearPendingCommands(); + } + + @Test + void onApplicationReadyDoesNotThrowWhenCleanupFails() { + MqttCommandAckService service = mock(MqttCommandAckService.class); + when(service.clearPendingCommands()).thenThrow(new IllegalStateException("redis unavailable")); + MqttPendingCommandCleanupListener listener = new MqttPendingCommandCleanupListener(service); + + assertThatCode(() -> listener.onApplicationReady(mock(ApplicationReadyEvent.class))) + .doesNotThrowAnyException(); + } +}