`
+- Produces: `AppVersionVo` fields `platform`, `currentVersion`, `latestVersion`, `versionCode`, `updateAvailable`, `forceUpdate`, `downloadUrl`, `releaseNotes`
+
+- [ ] **Step 1: Write failing tests**
+
+Add `IAppVersionService` as a mocked dependency and add tests for update available, no update, unsupported platform, missing current version, and missing version record.
+
+- [ ] **Step 2: Run tests and verify RED**
+
+Run:
+
+```bash
+mvn -pl water-modules/water-app -am "-DskipTests=false" "-Dmaven.test.skip=false" "-Dtest=AppControllerTest#checkVersion*" "-Dsurefire.failIfNoSpecifiedTests=false" test
+```
+
+Expected: compilation or test failure because `IAppVersionService` and table-backed implementation do not exist yet.
+
+- [ ] **Step 3: Implement minimal production code**
+
+Create `AppVersion`, `AppVersionVo`, `AppVersionMapper`, `IAppVersionService`, `AppVersionServiceImpl`, inject `IAppVersionService` into `AppController`, and keep `GET /checkVersion`.
+
+- [ ] **Step 4: Run targeted tests and verify GREEN**
+
+Run:
+
+```bash
+mvn -pl water-modules/water-app -am "-DskipTests=false" "-Dmaven.test.skip=false" "-Dtest=AppControllerTest#checkVersion*" "-Dsurefire.failIfNoSpecifiedTests=false" test
+```
+
+Expected: targeted version check tests pass.
+
+- [ ] **Step 5: Run compile/package verification**
+
+Run:
+
+```bash
+mvn -pl water-modules/water-app -am "-DskipTests" package
+```
+
+Expected: module and dependencies compile/package successfully.
diff --git a/docs/superpowers/plans/2026-07-13-device-status-lwt-mac.md b/docs/superpowers/plans/2026-07-13-device-status-lwt-mac.md
new file mode 100644
index 0000000..56a0fb7
--- /dev/null
+++ b/docs/superpowers/plans/2026-07-13-device-status-lwt-mac.md
@@ -0,0 +1,134 @@
+# Device Status LWT MAC Resolution Implementation Plan
+
+> **For agentic workers:** REQUIRED SUB-SKILL: Use superpowers:subagent-driven-development (recommended) or superpowers:executing-plans to implement this plan task-by-task. Steps use checkbox (`- [ ]`) syntax for tracking.
+
+**Goal:** Verify and document that `/MAC/publish/status` last-will messages resolve the MAC address to `deviceNo` before updating device status.
+
+**Architecture:** Keep `DeviceStatusHandler` focused on status parsing and delegate Topic identity resolution to the existing `DeviceIdentityResolver`. Add a focused integration-style unit test using the real resolver and a mocked mapper so the complete MAC-to-device-number path is covered without duplicating database access in the handler.
+
+**Tech Stack:** Java 17, Spring Boot, JUnit 5, Mockito, AssertJ, MyBatis-Plus
+
+## Global Constraints
+
+- Only `/MAC/publish/status` may use a MAC address for this change.
+- Other MQTT communication continues to use `deviceNo`.
+- A missing MAC mapping must not update device status.
+- Preserve all existing uncommitted workspace changes.
+
+---
+
+### Task 1: Cover MAC Last-Will Resolution
+
+**Files:**
+- Modify: `water-modules/water-app/src/test/java/org/dromara/app/handler/DeviceStatusHandlerTest.java`
+
+**Interfaces:**
+- Consumes: `DeviceIdentityResolver(AppDeviceMapper)` and `DeviceStatusHandler.handle(String, String, boolean)`
+- Produces: Regression coverage proving a Topic MAC is converted to `AppDevice.deviceNo`
+
+- [ ] **Step 1: Add the mapper mock and MAC last-will test**
+
+Add imports and a mapper mock:
+
+```java
+import org.dromara.app.domain.AppDevice;
+import org.dromara.app.mapper.AppDeviceMapper;
+
+@Mock
+private AppDeviceMapper appDeviceMapper;
+```
+
+Add this test:
+
+```java
+@Test
+void handleResolvesLastWillTopicMacToDeviceNo() {
+ String macAddress = "DC:DA:0C:FA:29:5E";
+ AppDevice device = new AppDevice();
+ device.setDeviceNo("D01");
+ when(appDeviceMapper.selectById(macAddress)).thenReturn(null);
+ when(appDeviceMapper.selectByMac("dc:da:0c:fa:29:5e")).thenReturn(device);
+ DeviceIdentityResolver resolver = new DeviceIdentityResolver(appDeviceMapper);
+ DeviceStatusHandler handler = new DeviceStatusHandler(resolver, deviceStatusService, new ObjectMapper());
+
+ handler.handle(macAddress, "{\"status\":\"offline\"}", false);
+
+ verify(appDeviceMapper).selectByMac("dc:da:0c:fa:29:5e");
+ verify(deviceStatusService).markOffline("D01", "设备 MQTT 状态离线");
+ verify(deviceStatusService, never()).markOffline(macAddress, "设备 MQTT 状态离线");
+}
+```
+
+Add the missing-device guard test:
+
+```java
+@Test
+void handleDoesNotUpdateStatusWhenLastWillMacIsUnknown() {
+ String macAddress = "DC:DA:0C:FA:29:5E";
+ when(appDeviceMapper.selectById(macAddress)).thenReturn(null);
+ when(appDeviceMapper.selectByMac("dc:da:0c:fa:29:5e")).thenReturn(null);
+ DeviceIdentityResolver resolver = new DeviceIdentityResolver(appDeviceMapper);
+ DeviceStatusHandler handler = new DeviceStatusHandler(resolver, deviceStatusService, new ObjectMapper());
+
+ handler.handle(macAddress, "{\"status\":\"offline\"}", false);
+
+ verify(deviceStatusService, never()).markOffline(anyString(), anyString());
+ verify(deviceStatusService, never()).markOnline(anyString());
+}
+```
+
+- [ ] **Step 2: Run the focused test**
+
+Run:
+
+```powershell
+mvn -pl water-modules/water-app -am "-DskipTests=false" "-Dmaven.test.skip=false" "-Dprofiles.active=dev" "-Dtest=DeviceStatusHandlerTest#handleResolvesLastWillTopicMacToDeviceNo+handleDoesNotUpdateStatusWhenLastWillMacIsUnknown" "-Dsurefire.failIfNoSpecifiedTests=false" test
+```
+
+Expected: both tests PASS. The requested production path already exists in `DeviceIdentityResolver`; these characterization tests make that behavior explicit and prevent regression.
+
+### Task 2: Clarify Device Status Topic Identity
+
+**Files:**
+- Modify: `water-modules/water-app/src/main/java/org/dromara/app/handler/DeviceStatusHandler.java:15-19`
+
+**Interfaces:**
+- Consumes: Existing `DeviceIdentityResolver.resolveDeviceNo(String)` behavior
+- Produces: Javadoc that accurately documents device-number and MAC Topic identities
+
+- [ ] **Step 1: Update handler Javadoc**
+
+Replace the class description with:
+
+```java
+/**
+ * 设备在线/离线状态处理器,匹配 /{deviceIdentity}/publish/status。
+ *
+ * deviceIdentity 支持设备编号;设备遗嘱消息允许使用 MAC 地址,处理前统一解析为设备编号。
+ * 设备通过 status=online 标记上线,通过 LWT status=offline 标记异常离线。
+ */
+```
+
+- [ ] **Step 2: Run all handler tests**
+
+Run:
+
+```powershell
+mvn -pl water-modules/water-app -am "-DskipTests=false" "-Dmaven.test.skip=false" "-Dprofiles.active=dev" "-Dtest=DeviceStatusHandlerTest" "-Dsurefire.failIfNoSpecifiedTests=false" test
+```
+
+Expected: all `DeviceStatusHandlerTest` tests PASS with zero failures and errors.
+
+- [ ] **Step 3: Check the final diff**
+
+Run:
+
+```powershell
+git diff --check -- water-modules/water-app/src/main/java/org/dromara/app/handler/DeviceStatusHandler.java water-modules/water-app/src/test/java/org/dromara/app/handler/DeviceStatusHandlerTest.java
+```
+
+Expected: exit code 0 and no whitespace errors.
+
+- [ ] **Step 4: Leave implementation changes uncommitted for review**
+
+Both implementation files already contain user changes. Do not create a commit that would mix those changes with this task.
diff --git a/docs/superpowers/specs/2026-06-24-bind-schedule-device-dispatch-design.md b/docs/superpowers/specs/2026-06-24-bind-schedule-device-dispatch-design.md
new file mode 100644
index 0000000..bd809d4
--- /dev/null
+++ b/docs/superpowers/specs/2026-06-24-bind-schedule-device-dispatch-design.md
@@ -0,0 +1,59 @@
+# 绑定排程时向设备下发排程信息设计
+
+## 背景
+
+当前 `AppController.addScheduleDevice` 只负责建立排程与设备的关联关系,不会把排程详情同步下发给设备。这样会导致设备绑定成功后仍然缺少最新排程配置。
+
+## 目标
+
+在 APP 端绑定排程与设备成功时,服务端同步向对应设备下发该排程详情信息,确保设备立即拿到当前排程配置。
+
+## 设计
+
+### 入口与改动范围
+
+- 入口保持在 `water-modules/water-app/src/main/java/org/dromara/app/controller/AppController.java` 的 `addScheduleDevice`
+- 复用现有 `IDeviceCommandService.sendScheduleBindCommand(...)`
+- 不修改现有 MQTT 发布器、ACK 处理器和设备主题路由
+
+### 下发时机
+
+对请求中的每个 `deviceNo`:
+
+1. 校验设备归属
+2. 写入排程设备关联
+3. 读取当前排程详情
+4. 同步下发排程命令
+
+### 命令协议
+
+- `commandType`: `bindSchedule`
+- `payload.deviceNo`: 当前设备编号
+- `payload.details`: 当前排程详情列表
+
+`payload.details` 每项沿用当前 APP 查询排程详情时对外返回的结构,包含:
+
+- `weekday`
+- `timeData`
+- `triggerType`
+- `status`
+
+如果排程详情里存在 `zones` 字段且当前查询对象可直接取到,则一并下发;否则不额外改造现有排程详情出参结构。
+
+### 失败策略
+
+绑定接口采用同步失败策略:
+
+- 只要某台设备的排程命令下发失败,接口直接返回失败
+- 不吞掉下发异常
+- 前端可以明确感知“排程绑定/同步下发”未完全成功
+
+本次不额外引入补偿、重试编排或异步任务。
+
+## 测试与验证
+
+- 优先补最小范围的自动化测试;如果模块当前没有现成测试基线,则至少执行模块编译或定向测试命令做回归验证
+- 手工验证重点:
+ - 绑定排程后关联表仍正常写入
+ - MQTT 下发 payload 包含 `deviceNo`、`details`
+ - 设备未登记 MAC 或命令下发异常时,接口返回失败
diff --git a/script/sql/update/add_app_version.sql b/script/sql/update/add_app_version.sql
new file mode 100644
index 0000000..3e52f8c
--- /dev/null
+++ b/script/sql/update/add_app_version.sql
@@ -0,0 +1,18 @@
+CREATE TABLE `app_version` (
+ `id` bigint NOT NULL COMMENT '主键',
+ `platform` varchar(20) NOT NULL COMMENT '平台 android/ios',
+ `latest_version` varchar(50) NOT NULL COMMENT '最新版本号',
+ `version_code` int NOT NULL COMMENT '版本序号,用于排序取最新版本',
+ `force_update` char(1) NOT NULL DEFAULT '0' COMMENT '是否强制更新 0否 1是',
+ `download_url` varchar(500) DEFAULT NULL COMMENT '下载地址',
+ `release_notes` varchar(1000) DEFAULT NULL COMMENT '更新说明',
+ `status` char(1) NOT NULL DEFAULT '1' COMMENT '状态 0停用 1启用',
+ `create_dept` bigint DEFAULT NULL COMMENT '创建部门',
+ `create_by` bigint DEFAULT NULL COMMENT '创建者',
+ `create_time` datetime DEFAULT NULL COMMENT '创建时间',
+ `update_by` bigint DEFAULT NULL COMMENT '更新者',
+ `update_time` datetime DEFAULT NULL COMMENT '更新时间',
+ `remark` varchar(500) DEFAULT NULL COMMENT '备注',
+ PRIMARY KEY (`id`),
+ KEY `idx_app_version_platform_status_code` (`platform`, `status`, `version_code`)
+) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_general_ci COMMENT='APP版本表';
diff --git a/water-admin/src/main/java/org/dromara/web/controller/AuthController.java b/water-admin/src/main/java/org/dromara/web/controller/AuthController.java
index d524752..5ed9d58 100644
--- a/water-admin/src/main/java/org/dromara/web/controller/AuthController.java
+++ b/water-admin/src/main/java/org/dromara/web/controller/AuthController.java
@@ -190,8 +190,9 @@ public class AuthController {
@DeleteMapping("/account")
public R cancelAccount(@Validated @RequestBody AccountCancelBody body) {
StpUtil.checkLogin();
+ Long userId = LoginHelper.getUserId();
accountCancellationService.cancelCurrentAccount(body.getCode());
- StpUtil.logout();
+ StpUtil.logout(userId);
return R.ok("注销成功");
}
diff --git a/water-admin/src/main/java/org/dromara/web/controller/CaptchaController.java b/water-admin/src/main/java/org/dromara/web/controller/CaptchaController.java
index 37950ba..62e69ea 100644
--- a/water-admin/src/main/java/org/dromara/web/controller/CaptchaController.java
+++ b/water-admin/src/main/java/org/dromara/web/controller/CaptchaController.java
@@ -110,6 +110,7 @@ public class CaptchaController {
/**
* 注销账号验证码
*/
+ @RateLimiter(key = "#{T(org.dromara.common.satoken.utils.LoginHelper).getUserId()}", time = 60, count = 1)
@GetMapping("/resource/account/cancel/code")
public R accountCancelCode() {
StpUtil.checkLogin();
@@ -129,7 +130,7 @@ public class CaptchaController {
if (!mailProperties.getEnabled()) {
return R.fail("当前系统没有开启邮箱功能!");
}
- emailCodeImpl(user.getEmail(), username);
+ SpringUtils.getAopProxy(this).emailCodeImpl(user.getEmail(), username);
return R.ok("操作成功");
}
return R.fail("当前账号未绑定手机号或邮箱");
diff --git a/water-admin/src/main/java/org/dromara/web/service/AccountCancelCodeService.java b/water-admin/src/main/java/org/dromara/web/service/AccountCancelCodeService.java
new file mode 100644
index 0000000..71eb650
--- /dev/null
+++ b/water-admin/src/main/java/org/dromara/web/service/AccountCancelCodeService.java
@@ -0,0 +1,36 @@
+package org.dromara.web.service;
+
+import lombok.RequiredArgsConstructor;
+import org.dromara.common.core.constant.GlobalConstants;
+import org.dromara.common.core.exception.user.CaptchaExpireException;
+import org.dromara.common.core.exception.user.UserException;
+import org.dromara.common.core.utils.StringUtils;
+import org.dromara.common.redis.utils.RedisUtils;
+import org.springframework.stereotype.Service;
+
+/**
+ * 注销账号验证码服务
+ */
+@RequiredArgsConstructor
+@Service
+public class AccountCancelCodeService {
+
+ public void validate(String username, String code) {
+ String cacheCode = RedisUtils.getCacheObject(buildKey(username));
+ if (StringUtils.isBlank(cacheCode)) {
+ throw new CaptchaExpireException();
+ }
+ if (!StringUtils.equals(cacheCode, code)) {
+ throw new UserException("验证码无效");
+ }
+ }
+
+ public void delete(String username) {
+ RedisUtils.deleteObject(buildKey(username));
+ }
+
+ private String buildKey(String username) {
+ return GlobalConstants.CAPTCHA_CODE_KEY + username;
+ }
+
+}
diff --git a/water-admin/src/main/java/org/dromara/web/service/AccountCancellationService.java b/water-admin/src/main/java/org/dromara/web/service/AccountCancellationService.java
new file mode 100644
index 0000000..532e74e
--- /dev/null
+++ b/water-admin/src/main/java/org/dromara/web/service/AccountCancellationService.java
@@ -0,0 +1,156 @@
+package org.dromara.web.service;
+
+import com.baomidou.mybatisplus.core.toolkit.Wrappers;
+import lombok.RequiredArgsConstructor;
+import org.dromara.app.domain.*;
+import org.dromara.app.mapper.*;
+import org.dromara.app.service.IDeviceCommandService;
+import org.dromara.common.core.exception.ServiceException;
+import org.dromara.common.core.utils.StringUtils;
+import org.dromara.common.satoken.utils.LoginHelper;
+import org.dromara.system.domain.SysSocial;
+import org.dromara.system.domain.SysUserPost;
+import org.dromara.system.domain.SysUserRole;
+import org.dromara.system.mapper.SysSocialMapper;
+import org.dromara.system.mapper.SysUserMapper;
+import org.dromara.system.mapper.SysUserPostMapper;
+import org.dromara.system.mapper.SysUserRoleMapper;
+import org.springframework.stereotype.Service;
+import org.springframework.transaction.annotation.Transactional;
+
+import java.util.List;
+import java.util.Objects;
+
+@RequiredArgsConstructor
+@Service
+public class AccountCancellationService {
+
+ private final SysUserMapper sysUserMapper;
+ private final SysUserRoleMapper sysUserRoleMapper;
+ private final SysUserPostMapper sysUserPostMapper;
+ private final SysSocialMapper sysSocialMapper;
+ private final AppDeviceMapper appDeviceMapper;
+ private final AppScheduleMapper appScheduleMapper;
+ private final AppScheduleDetailMapper appScheduleDetailMapper;
+ private final AppSchedulingDeviceMapper appSchedulingDeviceMapper;
+ private final AppWateringLogMapper appWateringLogMapper;
+ private final AccountCancelCodeService accountCancelCodeService;
+ private final IDeviceCommandService deviceCommandService;
+
+ @Transactional(rollbackFor = Exception.class)
+ public void cancelCurrentAccount(String code) {
+ Long userId = LoginHelper.getUserId();
+ if (userId == null) {
+ throw new ServiceException("用户未登录");
+ }
+ if (LoginHelper.isSuperAdmin(userId)) {
+ throw new ServiceException("超级管理员账号不允许注销");
+ }
+ String username = getCurrentUsername();
+ accountCancelCodeService.validate(username, code);
+
+ List deviceNos = queryUserDeviceNos(userId);
+ List scheduleIds = queryUserScheduleIds(userId);
+
+ sendInitDeviceCommands(deviceNos);
+ deleteAppScheduleData(userId, scheduleIds);
+ deleteAppDeviceData(userId, deviceNos);
+ deleteSystemUserData(userId);
+ accountCancelCodeService.delete(username);
+ }
+
+ private String getCurrentUsername() {
+ String username = LoginHelper.getUsername();
+ if (StringUtils.isBlank(username)) {
+ throw new ServiceException("用户未登录");
+ }
+ return username;
+ }
+
+ private List queryUserDeviceNos(Long userId) {
+ return appDeviceMapper.selectList(
+ Wrappers.lambdaQuery()
+ .select(AppDevice::getDeviceNo)
+ .eq(AppDevice::getUserId, userId)
+ ).stream()
+ .map(AppDevice::getDeviceNo)
+ .filter(StringUtils::isNotBlank)
+ .distinct()
+ .toList();
+ }
+
+ private List queryUserScheduleIds(Long userId) {
+ return appScheduleMapper.selectList(
+ Wrappers.lambdaQuery()
+ .select(AppSchedule::getId)
+ .eq(AppSchedule::getUserId, userId)
+ ).stream()
+ .map(AppSchedule::getId)
+ .filter(Objects::nonNull)
+ .distinct()
+ .toList();
+ }
+
+ private void sendInitDeviceCommands(List deviceNos) {
+ for (String deviceNo : deviceNos) {
+ deviceCommandService.sendInitDeviceCommand(deviceNo);
+ }
+ }
+
+ private void deleteAppScheduleData(Long userId, List scheduleIds) {
+ if (!scheduleIds.isEmpty()) {
+ appSchedulingDeviceMapper.delete(
+ Wrappers.lambdaQuery()
+ .in(AppSchedulingDevice::getScheduleId, scheduleIds)
+ );
+ appScheduleDetailMapper.delete(
+ Wrappers.lambdaQuery()
+ .in(AppScheduleDetail::getScheduleId, scheduleIds)
+ );
+ }
+ appScheduleMapper.delete(
+ Wrappers.lambdaQuery()
+ .eq(AppSchedule::getUserId, userId)
+ );
+ }
+
+ private void deleteAppDeviceData(Long userId, List deviceNos) {
+ if (!deviceNos.isEmpty()) {
+ appSchedulingDeviceMapper.delete(
+ Wrappers.lambdaQuery()
+ .in(AppSchedulingDevice::getDeviceNo, deviceNos)
+ );
+ appWateringLogMapper.delete(
+ Wrappers.lambdaQuery()
+ .in(AppWateringLog::getDeviceNo, deviceNos)
+ );
+ }
+ appWateringLogMapper.delete(
+ Wrappers.lambdaQuery()
+ .eq(AppWateringLog::getUserId, userId)
+ );
+ appDeviceMapper.delete(
+ Wrappers.lambdaQuery()
+ .eq(AppDevice::getUserId, userId)
+ );
+ }
+
+ private void deleteSystemUserData(Long userId) {
+ sysSocialMapper.delete(
+ Wrappers.lambdaQuery()
+ .eq(SysSocial::getUserId, userId)
+ );
+ sysUserRoleMapper.delete(
+ Wrappers.lambdaQuery()
+ .eq(SysUserRole::getUserId, userId)
+ );
+ sysUserPostMapper.delete(
+ Wrappers.lambdaQuery()
+ .eq(SysUserPost::getUserId, userId)
+ );
+ int rows = sysUserMapper.deleteById(userId);
+ if (rows < 1) {
+ throw new ServiceException("注销账号失败");
+ }
+ }
+}
diff --git a/water-admin/src/main/resources/application-prod.yml b/water-admin/src/main/resources/application-prod.yml
index 0b56e5c..234d5a2 100644
--- a/water-admin/src/main/resources/application-prod.yml
+++ b/water-admin/src/main/resources/application-prod.yml
@@ -4,7 +4,7 @@ spring.servlet.multipart.location: /water/server/temp
--- # 监控中心配置
spring.boot.admin.client:
# 增加客户端开关
- enabled: true
+ enabled: false
url: http://localhost:9090/admin
instance:
service-host-type: IP
@@ -16,7 +16,7 @@ spring.boot.admin.client:
--- # snail-job 配置
snail-job:
- enabled: true
+ enabled: false
# 需要在 SnailJob 后台组管理创建对应名称的组,然后创建任务的时候选择对应的组,才能正确分派任务
group: "water_group"
# SnailJob 接入验证令牌 详见 script/sql/ry_job.sql `sj_group_config`表
diff --git a/water-admin/src/main/resources/application.yml b/water-admin/src/main/resources/application.yml
index 613f655..1df150c 100644
--- a/water-admin/src/main/resources/application.yml
+++ b/water-admin/src/main/resources/application.yml
@@ -1,7 +1,7 @@
# 开发环境配置
server:
# 服务器的HTTP端口,默认为8080
- port: 8081
+ port: 8082
servlet:
# 应用的访问路径
context-path: /
@@ -26,9 +26,9 @@ mqtt:
broker-url: tcp://47.97.217.123
username: admin
password: 61e6062129e9
- client-id: water-server-${server.port}
+ client-id: water-server-1${server.port}
qos: 1
- keep-alive: 60
+ keep-alive: 300
connection-timeout: 30
max-inflight: 5000
clean-session: false
@@ -55,9 +55,10 @@ mqtt:
ack-key-prefix: "mqtt:command:ack:"
pending-set-key: "mqtt:command:pending:ids"
device-status-cache-prefix: "mqtt:device:status:"
- device-status-cache-ttl-seconds: 300
+ device-status-cache-ttl-seconds: 900
offline-check:
enabled: true
+ ttl-compat-enabled: true
interval-ms: 90000
topics:
# 订阅/下发主题第一段为设备 MAC(小写、去分隔符),上行业务侧由 DeviceIdentityResolver 解析为设备编号
@@ -65,6 +66,7 @@ mqtt:
- /+/publish/finish/schedule #排程任务完成上报
- /+/publish/register #设备注册
+ - /+/publish/status #设备在线/离线状态(online/offline,离线配合设备 LWT)
- /+/publish/power #电量
- /+/publish/ack #应答
- /+/publish/start #开始浇水上报
@@ -236,7 +238,7 @@ api-decrypt:
springdoc:
api-docs:
# 是否开启接口文档
- enabled: true
+ enabled: false
info:
# 标题
title: '标题:water-Vue-Plus多租户管理系统_接口文档'
@@ -306,7 +308,7 @@ websocket:
--- # warm-flow工作流配置
warm-flow:
# 是否开启工作流,默认true
- enabled: true
+ enabled: false
# 是否开启设计器ui
ui: true
# 是否显示流程图顶部文字
diff --git a/water-admin/src/test/java/org/dromara/web/config/MqttCommandAckConfigUnitTest.java b/water-admin/src/test/java/org/dromara/web/config/MqttCommandAckConfigUnitTest.java
new file mode 100644
index 0000000..a5e1b15
--- /dev/null
+++ b/water-admin/src/test/java/org/dromara/web/config/MqttCommandAckConfigUnitTest.java
@@ -0,0 +1,65 @@
+package org.dromara.web.config;
+
+import org.junit.jupiter.api.Tag;
+import org.junit.jupiter.api.Test;
+import org.springframework.core.io.ClassPathResource;
+import org.yaml.snakeyaml.Yaml;
+
+import java.io.InputStream;
+import java.util.Map;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+@Tag("dev")
+class MqttCommandAckConfigUnitTest {
+
+ @Test
+ void commandAckMaxRetryCountIsThree() throws Exception {
+ Integer maxRetryCount = null;
+ Yaml yaml = new Yaml();
+ try (InputStream inputStream = new ClassPathResource("application.yml").getInputStream()) {
+ for (Object document : yaml.loadAll(inputStream)) {
+ if (!(document instanceof Map, ?> root)) {
+ continue;
+ }
+ Object mqtt = root.get("mqtt");
+ if (!(mqtt instanceof Map, ?> mqttConfig)) {
+ continue;
+ }
+ Object commandAck = mqttConfig.get("command-ack");
+ if (!(commandAck instanceof Map, ?> commandAckConfig)) {
+ continue;
+ }
+ Object value = commandAckConfig.get("max-retry-count");
+ if (value instanceof Number number) {
+ maxRetryCount = number.intValue();
+ }
+ }
+ }
+
+ assertThat(maxRetryCount).isEqualTo(3);
+ }
+
+ @Test
+ void offlineCheckTtlCompatibilityIsEnabled() throws Exception {
+ Boolean ttlCompatEnabled = null;
+ Yaml yaml = new Yaml();
+ try (InputStream inputStream = new ClassPathResource("application.yml").getInputStream()) {
+ for (Object document : yaml.loadAll(inputStream)) {
+ if (!(document instanceof Map, ?> root)) {
+ continue;
+ }
+ Object mqtt = root.get("mqtt");
+ if (!(mqtt instanceof Map, ?> mqttConfig)) {
+ continue;
+ }
+ Object offlineCheck = mqttConfig.get("offline-check");
+ if (offlineCheck instanceof Map, ?> offlineCheckConfig) {
+ ttlCompatEnabled = (Boolean) offlineCheckConfig.get("ttl-compat-enabled");
+ }
+ }
+ }
+
+ assertThat(ttlCompatEnabled).isTrue();
+ }
+}
diff --git a/water-admin/src/test/java/org/dromara/web/controller/AuthControllerUnitTest.java b/water-admin/src/test/java/org/dromara/web/controller/AuthControllerUnitTest.java
new file mode 100644
index 0000000..3a257e6
--- /dev/null
+++ b/water-admin/src/test/java/org/dromara/web/controller/AuthControllerUnitTest.java
@@ -0,0 +1,55 @@
+package org.dromara.web.controller;
+
+import cn.dev33.satoken.stp.StpUtil;
+import org.dromara.common.core.domain.R;
+import org.dromara.common.core.domain.model.AccountCancelBody;
+import org.dromara.common.satoken.utils.LoginHelper;
+import org.dromara.web.service.AccountCancellationService;
+import org.junit.jupiter.api.Tag;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.extension.ExtendWith;
+import org.mockito.Mock;
+import org.mockito.MockedStatic;
+import org.mockito.junit.jupiter.MockitoExtension;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.mockito.Mockito.mockStatic;
+import static org.mockito.Mockito.verify;
+
+@ExtendWith(MockitoExtension.class)
+@Tag("dev")
+class AuthControllerUnitTest {
+
+ @Mock
+ private AccountCancellationService accountCancellationService;
+
+ @Test
+ void cancelAccount_checksLoginCancelsAccountAndLogsOut() {
+ AuthController controller = new AuthController(
+ null,
+ null,
+ null,
+ null,
+ null,
+ null,
+ null,
+ null,
+ null,
+ accountCancellationService
+ );
+
+ try (MockedStatic stpUtil = mockStatic(StpUtil.class);
+ MockedStatic loginHelper = mockStatic(LoginHelper.class)) {
+ AccountCancelBody body = new AccountCancelBody();
+ body.setCode("123456");
+ loginHelper.when(LoginHelper::getUserId).thenReturn(100L);
+
+ R result = controller.cancelAccount(body);
+
+ assertThat(result.getCode()).isEqualTo(200);
+ verify(accountCancellationService).cancelCurrentAccount("123456");
+ stpUtil.verify(StpUtil::checkLogin);
+ stpUtil.verify(() -> StpUtil.logout(100L));
+ }
+ }
+}
diff --git a/water-admin/src/test/java/org/dromara/web/controller/CaptchaControllerUnitTest.java b/water-admin/src/test/java/org/dromara/web/controller/CaptchaControllerUnitTest.java
new file mode 100644
index 0000000..b633b07
--- /dev/null
+++ b/water-admin/src/test/java/org/dromara/web/controller/CaptchaControllerUnitTest.java
@@ -0,0 +1,119 @@
+package org.dromara.web.controller;
+
+import cn.dev33.satoken.stp.StpUtil;
+import org.dromara.common.core.domain.R;
+import org.dromara.common.core.utils.SpringUtils;
+import org.dromara.common.mail.config.properties.MailProperties;
+import org.dromara.common.ratelimiter.annotation.RateLimiter;
+import org.dromara.common.satoken.utils.LoginHelper;
+import org.dromara.system.domain.vo.SysUserVo;
+import org.dromara.system.service.ISysUserService;
+import org.junit.jupiter.api.Tag;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.extension.ExtendWith;
+import org.mockito.Mock;
+import org.mockito.MockedStatic;
+import org.mockito.junit.jupiter.MockitoExtension;
+
+import java.lang.reflect.Method;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.mockito.Mockito.*;
+
+@ExtendWith(MockitoExtension.class)
+@Tag("dev")
+class CaptchaControllerUnitTest {
+
+ @Mock
+ private MailProperties mailProperties;
+ @Mock
+ private ISysUserService userService;
+
+ @Test
+ void accountCancelCode_sendsSmsCodeWhenCurrentUserHasPhone() {
+ CaptchaController controller = spy(new CaptchaController(null, mailProperties, userService));
+ SysUserVo user = new SysUserVo();
+ user.setPhonenumber("13305376054");
+ user.setEmail("demo@example.com");
+ when(userService.selectUserById(100L)).thenReturn(user);
+ doReturn(R.ok("操作成功")).when(controller).sendSmsCode("13305376054", "zhangsan");
+
+ try (MockedStatic stpUtil = mockStatic(StpUtil.class);
+ MockedStatic springUtils = mockStatic(SpringUtils.class);
+ MockedStatic loginHelper = mockStatic(LoginHelper.class)) {
+ loginHelper.when(LoginHelper::getUserId).thenReturn(100L);
+ loginHelper.when(LoginHelper::getUsername).thenReturn("zhangsan");
+ springUtils.when(() -> SpringUtils.getAopProxy(controller)).thenReturn(controller);
+
+ R result = controller.accountCancelCode();
+
+ assertThat(result.getCode()).isEqualTo(200);
+ stpUtil.verify(StpUtil::checkLogin);
+ verify(controller).sendSmsCode("13305376054", "zhangsan");
+ verify(controller, never()).emailCodeImpl("demo@example.com");
+ }
+ }
+
+ @Test
+ void accountCancelCode_sendsEmailCodeWhenCurrentUserHasNoPhone() {
+ CaptchaController controller = spy(new CaptchaController(null, mailProperties, userService));
+ SysUserVo user = new SysUserVo();
+ user.setEmail("demo@example.com");
+ when(userService.selectUserById(100L)).thenReturn(user);
+ when(mailProperties.getEnabled()).thenReturn(true);
+ doNothing().when(controller).emailCodeImpl("demo@example.com", "zhangsan");
+
+ try (MockedStatic stpUtil = mockStatic(StpUtil.class);
+ MockedStatic springUtils = mockStatic(SpringUtils.class);
+ MockedStatic loginHelper = mockStatic(LoginHelper.class)) {
+ loginHelper.when(LoginHelper::getUserId).thenReturn(100L);
+ loginHelper.when(LoginHelper::getUsername).thenReturn("zhangsan");
+ springUtils.when(() -> SpringUtils.getAopProxy(controller)).thenReturn(controller);
+
+ R result = controller.accountCancelCode();
+
+ assertThat(result.getCode()).isEqualTo(200);
+ stpUtil.verify(StpUtil::checkLogin);
+ springUtils.verify(() -> SpringUtils.getAopProxy(controller));
+ verify(controller).emailCodeImpl("demo@example.com", "zhangsan");
+ verify(controller, never()).sendSmsCode("demo@example.com", "zhangsan");
+ }
+ }
+
+
+ @Test
+ void accountCancelCode_hasPerUserRateLimiter() throws NoSuchMethodException {
+ Method method = CaptchaController.class.getDeclaredMethod("accountCancelCode");
+
+ RateLimiter rateLimiter = method.getAnnotation(RateLimiter.class);
+
+ assertThat(rateLimiter).isNotNull();
+ assertThat(rateLimiter.time()).isEqualTo(60);
+ assertThat(rateLimiter.count()).isEqualTo(1);
+ assertThat(rateLimiter.key())
+ .isEqualTo("#{T(org.dromara.common.satoken.utils.LoginHelper).getUserId()}");
+ }
+
+ @Test
+ void accountCancelCode_failsWhenCurrentUserHasNoPhoneOrEmail() {
+ CaptchaController controller = spy(new CaptchaController(null, mailProperties, userService));
+ SysUserVo user = new SysUserVo();
+ when(userService.selectUserById(100L)).thenReturn(user);
+
+ try (MockedStatic stpUtil = mockStatic(StpUtil.class);
+ MockedStatic springUtils = mockStatic(SpringUtils.class);
+ MockedStatic loginHelper = mockStatic(LoginHelper.class)) {
+ loginHelper.when(LoginHelper::getUserId).thenReturn(100L);
+ loginHelper.when(LoginHelper::getUsername).thenReturn("zhangsan");
+ springUtils.when(() -> SpringUtils.getAopProxy(controller)).thenReturn(controller);
+
+ R result = controller.accountCancelCode();
+
+ assertThat(result.getCode()).isEqualTo(500);
+ assertThat(result.getMsg()).isEqualTo("当前账号未绑定手机号或邮箱");
+ stpUtil.verify(StpUtil::checkLogin);
+ verify(controller, never()).sendSmsCode("13305376054", "zhangsan");
+ verify(controller, never()).emailCodeImpl("demo@example.com", "zhangsan");
+ }
+ }
+}
diff --git a/water-admin/src/test/java/org/dromara/web/service/AccountCancellationServiceUnitTest.java b/water-admin/src/test/java/org/dromara/web/service/AccountCancellationServiceUnitTest.java
new file mode 100644
index 0000000..a1d88e1
--- /dev/null
+++ b/water-admin/src/test/java/org/dromara/web/service/AccountCancellationServiceUnitTest.java
@@ -0,0 +1,215 @@
+package org.dromara.web.service;
+
+import com.baomidou.mybatisplus.core.MybatisConfiguration;
+import com.baomidou.mybatisplus.core.conditions.Wrapper;
+import com.baomidou.mybatisplus.core.metadata.TableInfoHelper;
+import org.apache.ibatis.builder.MapperBuilderAssistant;
+import org.dromara.app.domain.*;
+import org.dromara.app.mapper.*;
+import org.dromara.app.service.IDeviceCommandService;
+import org.dromara.common.core.exception.ServiceException;
+import org.dromara.common.core.exception.user.CaptchaExpireException;
+import org.dromara.common.core.exception.user.UserException;
+import org.dromara.common.satoken.utils.LoginHelper;
+import org.dromara.system.domain.SysSocial;
+import org.dromara.system.domain.SysUserPost;
+import org.dromara.system.domain.SysUserRole;
+import org.dromara.system.mapper.SysSocialMapper;
+import org.dromara.system.mapper.SysUserMapper;
+import org.dromara.system.mapper.SysUserPostMapper;
+import org.dromara.system.mapper.SysUserRoleMapper;
+import org.junit.jupiter.api.BeforeAll;
+import org.junit.jupiter.api.Tag;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.extension.ExtendWith;
+import org.mockito.Mock;
+import org.mockito.MockedStatic;
+import org.mockito.junit.jupiter.MockitoExtension;
+
+import java.util.List;
+
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.Mockito.*;
+
+@ExtendWith(MockitoExtension.class)
+@Tag("dev")
+class AccountCancellationServiceUnitTest {
+
+ @Mock
+ private SysUserMapper sysUserMapper;
+ @Mock
+ private SysUserRoleMapper sysUserRoleMapper;
+ @Mock
+ private SysUserPostMapper sysUserPostMapper;
+ @Mock
+ private SysSocialMapper sysSocialMapper;
+ @Mock
+ private AppDeviceMapper appDeviceMapper;
+ @Mock
+ private AppScheduleMapper appScheduleMapper;
+ @Mock
+ private AppScheduleDetailMapper appScheduleDetailMapper;
+ @Mock
+ private AppSchedulingDeviceMapper appSchedulingDeviceMapper;
+ @Mock
+ private AppWateringLogMapper appWateringLogMapper;
+ @Mock
+ private AccountCancelCodeService accountCancelCodeService;
+ @Mock
+ private IDeviceCommandService deviceCommandService;
+
+ @BeforeAll
+ static void initTableInfo() {
+ MapperBuilderAssistant assistant = new MapperBuilderAssistant(new MybatisConfiguration(), "");
+ initTableInfo(assistant, AppDevice.class);
+ initTableInfo(assistant, AppSchedule.class);
+ initTableInfo(assistant, AppScheduleDetail.class);
+ initTableInfo(assistant, AppSchedulingDevice.class);
+ initTableInfo(assistant, AppWateringLog.class);
+ initTableInfo(assistant, SysSocial.class);
+ initTableInfo(assistant, SysUserRole.class);
+ initTableInfo(assistant, SysUserPost.class);
+ }
+
+ private static void initTableInfo(MapperBuilderAssistant assistant, Class> entityClass) {
+ TableInfoHelper.remove(entityClass);
+ TableInfoHelper.initTableInfo(assistant, entityClass);
+ }
+
+ @Test
+ void cancelCurrentAccount_deletesCurrentUserAndRelatedData() {
+ AccountCancellationService service = newService();
+ AppDevice device = new AppDevice();
+ device.setDeviceNo("D01");
+ AppSchedule schedule = new AppSchedule();
+ schedule.setId(10L);
+ when(appDeviceMapper.selectList(any(Wrapper.class))).thenReturn(List.of(device));
+ when(appScheduleMapper.selectList(any(Wrapper.class))).thenReturn(List.of(schedule));
+ when(sysUserMapper.deleteById(100L)).thenReturn(1);
+
+ try (MockedStatic loginHelper = mockStatic(LoginHelper.class)) {
+ loginHelper.when(LoginHelper::getUserId).thenReturn(100L);
+ loginHelper.when(LoginHelper::getUsername).thenReturn("zhangsan");
+ loginHelper.when(() -> LoginHelper.isSuperAdmin(100L)).thenReturn(false);
+
+ service.cancelCurrentAccount("123456");
+
+ verify(accountCancelCodeService).validate("zhangsan", "123456");
+ verify(deviceCommandService).sendInitDeviceCommand("D01");
+ verify(appSchedulingDeviceMapper, atLeastOnce()).delete(any(Wrapper.class));
+ verify(appScheduleDetailMapper).delete(any(Wrapper.class));
+ verify(appScheduleMapper).delete(any(Wrapper.class));
+ verify(appWateringLogMapper, atLeastOnce()).delete(any(Wrapper.class));
+ verify(appDeviceMapper).delete(any(Wrapper.class));
+ verify(sysSocialMapper).delete(any(Wrapper.class));
+ verify(sysUserRoleMapper).delete(any(Wrapper.class));
+ verify(sysUserPostMapper).delete(any(Wrapper.class));
+ verify(sysUserMapper).deleteById(100L);
+ verify(accountCancelCodeService).delete("zhangsan");
+ }
+ }
+
+ @Test
+ void cancelCurrentAccount_rejectsSuperAdmin() {
+ AccountCancellationService service = newService();
+
+ try (MockedStatic loginHelper = mockStatic(LoginHelper.class)) {
+ loginHelper.when(LoginHelper::getUserId).thenReturn(1L);
+ loginHelper.when(() -> LoginHelper.isSuperAdmin(1L)).thenReturn(true);
+
+ assertThatThrownBy(() -> service.cancelCurrentAccount("123456"))
+ .isInstanceOf(ServiceException.class)
+ .hasMessageContaining("超级管理员");
+
+ verifyNoInteractions(
+ sysUserMapper,
+ sysUserRoleMapper,
+ sysUserPostMapper,
+ sysSocialMapper,
+ appDeviceMapper,
+ appScheduleMapper,
+ appScheduleDetailMapper,
+ appSchedulingDeviceMapper,
+ appWateringLogMapper,
+ accountCancelCodeService
+ );
+ }
+ }
+
+ @Test
+ void cancelCurrentAccount_rejectsInvalidCodeWithoutDeletingData() {
+ AccountCancellationService service = newService();
+
+ try (MockedStatic loginHelper = mockStatic(LoginHelper.class)) {
+ loginHelper.when(LoginHelper::getUserId).thenReturn(100L);
+ loginHelper.when(LoginHelper::getUsername).thenReturn("zhangsan");
+ loginHelper.when(() -> LoginHelper.isSuperAdmin(100L)).thenReturn(false);
+ doThrow(new UserException("验证码无效"))
+ .when(accountCancelCodeService).validate("zhangsan", "654321");
+
+ assertThatThrownBy(() -> service.cancelCurrentAccount("654321"))
+ .isInstanceOf(UserException.class);
+
+ verifyNoInteractions(
+ sysUserMapper,
+ sysUserRoleMapper,
+ sysUserPostMapper,
+ sysSocialMapper,
+ appDeviceMapper,
+ appScheduleMapper,
+ appScheduleDetailMapper,
+ appSchedulingDeviceMapper,
+ appWateringLogMapper
+ );
+ verify(accountCancelCodeService).validate("zhangsan", "654321");
+ verify(accountCancelCodeService, never()).delete("zhangsan");
+ }
+ }
+
+ @Test
+ void cancelCurrentAccount_rejectsExpiredCodeWithoutDeletingData() {
+ AccountCancellationService service = newService();
+
+ try (MockedStatic loginHelper = mockStatic(LoginHelper.class)) {
+ loginHelper.when(LoginHelper::getUserId).thenReturn(100L);
+ loginHelper.when(LoginHelper::getUsername).thenReturn("zhangsan");
+ loginHelper.when(() -> LoginHelper.isSuperAdmin(100L)).thenReturn(false);
+ doThrow(new CaptchaExpireException())
+ .when(accountCancelCodeService).validate("zhangsan", "123456");
+
+ assertThatThrownBy(() -> service.cancelCurrentAccount("123456"))
+ .isInstanceOf(CaptchaExpireException.class);
+
+ verifyNoInteractions(
+ sysUserMapper,
+ sysUserRoleMapper,
+ sysUserPostMapper,
+ sysSocialMapper,
+ appDeviceMapper,
+ appScheduleMapper,
+ appScheduleDetailMapper,
+ appSchedulingDeviceMapper,
+ appWateringLogMapper
+ );
+ verify(accountCancelCodeService).validate("zhangsan", "123456");
+ verify(accountCancelCodeService, never()).delete("zhangsan");
+ }
+ }
+
+ private AccountCancellationService newService() {
+ return new AccountCancellationService(
+ sysUserMapper,
+ sysUserRoleMapper,
+ sysUserPostMapper,
+ sysSocialMapper,
+ appDeviceMapper,
+ appScheduleMapper,
+ appScheduleDetailMapper,
+ appSchedulingDeviceMapper,
+ appWateringLogMapper,
+ accountCancelCodeService,
+ deviceCommandService
+ );
+ }
+}
diff --git a/water-common/water-common-core/src/main/java/org/dromara/common/core/domain/model/AccountCancelBody.java b/water-common/water-common-core/src/main/java/org/dromara/common/core/domain/model/AccountCancelBody.java
new file mode 100644
index 0000000..f5bb096
--- /dev/null
+++ b/water-common/water-common-core/src/main/java/org/dromara/common/core/domain/model/AccountCancelBody.java
@@ -0,0 +1,24 @@
+package org.dromara.common.core.domain.model;
+
+import jakarta.validation.constraints.NotBlank;
+import lombok.Data;
+
+import java.io.Serial;
+import java.io.Serializable;
+
+/**
+ * 注销账号请求体
+ */
+@Data
+public class AccountCancelBody implements Serializable {
+
+ @Serial
+ private static final long serialVersionUID = 1L;
+
+ /**
+ * 验证码
+ */
+ @NotBlank(message = "{sms.code.not.blank}")
+ private String code;
+
+}
diff --git a/water-common/water-common-mqtt/src/main/java/org/dromara/mqtt/MqttClientManager.java b/water-common/water-common-mqtt/src/main/java/org/dromara/mqtt/MqttClientManager.java
index f93a19e..0530854 100644
--- a/water-common/water-common-mqtt/src/main/java/org/dromara/mqtt/MqttClientManager.java
+++ b/water-common/water-common-mqtt/src/main/java/org/dromara/mqtt/MqttClientManager.java
@@ -37,65 +37,18 @@ public class MqttClientManager implements DisposableBean {
private final MqttMessageDispatcher dispatcher;
private MqttAsyncClient client;
private ThreadPoolExecutor consumerExecutor;
- private BlockingQueue messageQueue;
+ private List> messageQueues;
private volatile boolean running;
+ private int consumerCount;
- @Bean
- @ConditionalOnProperty(prefix = "mqtt", name = "enabled", havingValue = "true")
- public MqttAsyncClient mqttConnect() throws MqttException {
- if (StringUtils.isBlank(props.getBrokerUrl())) {
- throw new ServiceException("启用 MQTT 时必须配置 broker-url");
+ static int deviceIdentityHash(String topic) {
+ if (topic == null) {
+ return 0;
}
- if (StringUtils.isBlank(props.getClientId())) {
- throw new ServiceException("启用 MQTT 时必须配置 client-id");
- }
- messageQueue = new ArrayBlockingQueue<>(Math.max(1, props.getAsync().getQueueCapacity()));
- consumerExecutor = createConsumerExecutor();
- running = true;
- startConsumers();
- client = new MqttAsyncClient(props.getBrokerUrl(), props.getClientId(), new MemoryPersistence());
-
- MqttConnectOptions options = new MqttConnectOptions();
- options.setUserName(props.getUsername());
- if (props.getPassword() != null) {
- options.setPassword(props.getPassword().toCharArray());
- }
- options.setKeepAliveInterval(props.getKeepAlive());
- options.setConnectionTimeout(props.getConnectionTimeout());
- options.setAutomaticReconnect(props.isAutomaticReconnect());
- options.setCleanSession(props.isCleanSession());
- options.setMaxInflight(props.getMaxInflight());
- configureTls(options);
-
- client.setCallback(new MqttCallbackExtended() {
- @Override
- public void connectionLost(Throwable cause) {
- log.warn("[MQTT] 连接已断开:{}", cause == null ? "" : cause.getMessage());
- }
-
- @Override
- public void connectComplete(boolean reconnect, String serverURI) {
- if (reconnect) {
- log.info("[MQTT] 已重新连接:{}", serverURI);
- subscribeTopics();
- }
- }
-
- @Override
- public void messageArrived(String topic, MqttMessage message) {
- String payload = new String(message.getPayload(), StandardCharsets.UTF_8);
- enqueue(topic, payload);
- }
-
- @Override
- public void deliveryComplete(IMqttDeliveryToken token) {
- }
- });
-
- client.connect(options).waitForCompletion();
- log.info("[MQTT] 已连接:{}", props.getBrokerUrl());
- subscribeTopics();
- return client;
+ int start = topic.startsWith("/") ? 1 : 0;
+ int end = topic.indexOf('/', start);
+ String identity = end < 0 ? topic.substring(start) : topic.substring(start, end);
+ return identity.hashCode();
}
private void configureTls(MqttConnectOptions options) throws MqttException {
@@ -171,39 +124,97 @@ public class MqttClientManager implements DisposableBean {
};
}
+ @Bean
+ @ConditionalOnProperty(prefix = "mqtt", name = "enabled", havingValue = "true")
+ public MqttAsyncClient mqttConnect() throws MqttException {
+ if (StringUtils.isBlank(props.getBrokerUrl())) {
+ throw new ServiceException("启用 MQTT 时必须配置 broker-url");
+ }
+ if (StringUtils.isBlank(props.getClientId())) {
+ throw new ServiceException("启用 MQTT 时必须配置 client-id");
+ }
+ consumerCount = Math.max(1, props.getAsync().getConsumerCount());
+ messageQueues = createMessageQueues();
+ consumerExecutor = createConsumerExecutor();
+ running = true;
+ startConsumers();
+ client = new MqttAsyncClient(props.getBrokerUrl(), props.getClientId(), new MemoryPersistence());
+
+ MqttConnectOptions options = new MqttConnectOptions();
+ options.setUserName(props.getUsername());
+ if (props.getPassword() != null) {
+ options.setPassword(props.getPassword().toCharArray());
+ }
+ options.setKeepAliveInterval(props.getKeepAlive());
+ options.setConnectionTimeout(props.getConnectionTimeout());
+ options.setAutomaticReconnect(props.isAutomaticReconnect());
+ options.setCleanSession(props.isCleanSession());
+ options.setMaxInflight(props.getMaxInflight());
+ configureTls(options);
+
+ client.setCallback(new MqttCallbackExtended() {
+ @Override
+ public void connectionLost(Throwable cause) {
+ log.warn("[MQTT] 连接已断开:{}", cause == null ? "" : cause.getMessage());
+ }
+
+ @Override
+ public void connectComplete(boolean reconnect, String serverURI) {
+ if (reconnect) {
+ log.info("[MQTT] 已重新连接:{}", serverURI);
+ subscribeTopics();
+ }
+ }
+
+ @Override
+ public void messageArrived(String topic, MqttMessage message) {
+ String payload = new String(message.getPayload(), StandardCharsets.UTF_8);
+ enqueue(topic, payload, message.isRetained());
+ }
+
+ @Override
+ public void deliveryComplete(IMqttDeliveryToken token) {
+ }
+ });
+
+ client.connect(options).waitForCompletion();
+ log.info("[MQTT] 已连接:{}", props.getBrokerUrl());
+ subscribeTopics();
+ return client;
+ }
+
private void startConsumers() {
- int consumerCount = Math.max(1, props.getAsync().getConsumerCount());
for (int i = 0; i < consumerCount; i++) {
- consumerExecutor.execute(this::consumeLoop);
+ final int consumerIndex = i;
+ consumerExecutor.execute(() -> consumeLoop(consumerIndex));
}
}
- private void enqueue(String topic, String payload) {
- InboundMessage inboundMessage = new InboundMessage(topic, payload);
+ private void enqueue(String topic, String payload, boolean retained) {
+ InboundMessage inboundMessage = new InboundMessage(topic, payload, retained);
try {
- if (!messageQueue.offer(inboundMessage, props.getAsync().getOfferTimeoutMs(), TimeUnit.MILLISECONDS)) {
- log.warn("[MQTT] 上行消息队列已满,丢弃 主题={}", topic);
- }
+ messageQueues.get(queueIndex(topic)).put(inboundMessage);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
log.warn("[MQTT] 上行消息入队被中断 主题={}", topic);
}
}
- private void consumeLoop() {
+ private void consumeLoop(int consumerIndex) {
MqttProperties.Async async = props.getAsync();
int batchSize = Math.max(1, async.getBatchSize());
long pollTimeoutMs = Math.max(1, async.getPollTimeoutMs());
List batch = new java.util.ArrayList<>(batchSize);
+ BlockingQueue queue = messageQueues.get(consumerIndex);
- while (running || !messageQueue.isEmpty()) {
+ while (running || !queue.isEmpty()) {
try {
- InboundMessage first = messageQueue.poll(pollTimeoutMs, TimeUnit.MILLISECONDS);
+ InboundMessage first = queue.poll(pollTimeoutMs, TimeUnit.MILLISECONDS);
if (first == null) {
continue;
}
batch.add(first);
- messageQueue.drainTo(batch, batchSize - 1);
+ queue.drainTo(batch, batchSize - 1);
for (InboundMessage message : batch) {
dispatch(message);
}
@@ -218,12 +229,26 @@ public class MqttClientManager implements DisposableBean {
private void dispatch(InboundMessage message) {
try {
- dispatcher.dispatch(message.topic(), message.payload());
+ dispatcher.dispatch(message.topic(), message.payload(), message.retained());
} catch (Exception e) {
log.error("[MQTT] 消息分发失败 主题={}", message.topic(), e);
}
}
+ private List> createMessageQueues() {
+ int capacity = Math.max(1, props.getAsync().getQueueCapacity());
+ int queueCapacity = Math.max(1, (capacity + Math.max(1, consumerCount) - 1) / Math.max(1, consumerCount));
+ List> queues = new java.util.ArrayList<>(consumerCount);
+ for (int i = 0; i < consumerCount; i++) {
+ queues.add(new ArrayBlockingQueue<>(queueCapacity));
+ }
+ return queues;
+ }
+
+ private int queueIndex(String topic) {
+ return Math.floorMod(deviceIdentityHash(topic), Math.max(1, consumerCount));
+ }
+
private void subscribeTopics() {
try {
if (client == null || !client.isConnected()) {
@@ -276,6 +301,6 @@ public class MqttClientManager implements DisposableBean {
}
}
- private record InboundMessage(String topic, String payload) {
+ private record InboundMessage(String topic, String payload, boolean retained) {
}
}
diff --git a/water-common/water-common-mqtt/src/main/java/org/dromara/mqtt/MqttCommandAckService.java b/water-common/water-common-mqtt/src/main/java/org/dromara/mqtt/MqttCommandAckService.java
index 64a9b1a..70b7ab3 100644
--- a/water-common/water-common-mqtt/src/main/java/org/dromara/mqtt/MqttCommandAckService.java
+++ b/water-common/water-common-mqtt/src/main/java/org/dromara/mqtt/MqttCommandAckService.java
@@ -22,6 +22,7 @@ import java.time.Duration;
import java.time.Instant;
import java.util.*;
import java.util.concurrent.TimeUnit;
+import java.util.function.BooleanSupplier;
@Slf4j
@Service
@@ -43,11 +44,10 @@ public class MqttCommandAckService implements IDeviceCommandAckHandler {
pendingIds().add(command.getCommandId());
}
- public void removePending(String commandId) {
- if (StringUtils.isBlank(commandId)) {
- return;
- }
- withCommandLock(commandId, () -> deletePending(commandId));
+ static boolean ackMatchesPendingCommand(String topicDeviceNo, DeviceCommand pendingCommand) {
+ return StringUtils.isNotBlank(topicDeviceNo)
+ && pendingCommand != null
+ && topicDeviceNo.equals(pendingCommand.getDeviceNo());
}
static boolean resolveMissingCommandId(DeviceCommandAck ack, List pendingCommands) {
@@ -68,6 +68,16 @@ public class MqttCommandAckService implements IDeviceCommandAckHandler {
return true;
}
+ public void removePending(String commandId) {
+ if (StringUtils.isBlank(commandId)) {
+ return;
+ }
+ withCommandLock(commandId, () -> {
+ deletePending(commandId);
+ return true;
+ });
+ }
+
@Override
public void handleAck(String deviceNo, String payload) {
if (StringUtils.isBlank(payload)) {
@@ -102,8 +112,19 @@ public class MqttCommandAckService implements IDeviceCommandAckHandler {
}
}
boolean ackSaved = withCommandLock(ack.getCommandId(), () -> {
+ DeviceCommand pending = RedisUtils.getCacheObject(pendingKey(ack.getCommandId()));
+ if (pending == null) {
+ log.warn("[MQTT] ACK 对应命令不存在或已过期 设备编号={} 命令编号={}", deviceNo, ack.getCommandId());
+ return false;
+ }
+ if (!ackMatchesPendingCommand(deviceNo, pending)) {
+ log.warn("[MQTT] ACK 设备编号与待确认命令不匹配,已拒绝 设备编号={} 命令编号={} 命令设备={}",
+ deviceNo, ack.getCommandId(), pending.getDeviceNo());
+ return false;
+ }
deletePending(ack.getCommandId());
RedisUtils.setCacheObject(ackKey(ack.getCommandId()), ack, Duration.ofSeconds(mqttProperties.getCommandAck().getAckTtlSeconds()));
+ return true;
});
if (!ackSaved) {
log.warn("[MQTT] ACK 处理未获得命令锁,本次跳过 设备编号={} 命令编号={}", deviceNo, ack.getCommandId());
@@ -130,8 +151,19 @@ public class MqttCommandAckService implements IDeviceCommandAckHandler {
ack.setMessage(payload);
boolean ackSaved = withCommandLock(commandId, () -> {
+ DeviceCommand pending = RedisUtils.getCacheObject(pendingKey(commandId));
+ if (pending == null) {
+ log.warn("[MQTT] 非 JSON ACK 对应命令不存在或已过期 设备编号={} 命令编号={}", deviceNo, commandId);
+ return false;
+ }
+ if (!ackMatchesPendingCommand(deviceNo, pending)) {
+ log.warn("[MQTT] 非 JSON ACK 设备编号与待确认命令不匹配,已拒绝 设备编号={} 命令编号={} 命令设备={}",
+ deviceNo, commandId, pending.getDeviceNo());
+ return false;
+ }
deletePending(commandId);
RedisUtils.setCacheObject(ackKey(commandId), ack, Duration.ofSeconds(mqttProperties.getCommandAck().getAckTtlSeconds()));
+ return true;
});
if (!ackSaved) {
log.warn("[MQTT] 非 JSON ACK 处理未获得命令锁,本次跳过 设备编号={} 命令编号={}", deviceNo, commandId);
@@ -160,7 +192,10 @@ public class MqttCommandAckService implements IDeviceCommandAckHandler {
}
private void retryCommandIfLocked(String commandId, long now) {
- withCommandLock(commandId, () -> retryCommand(commandId, now));
+ withCommandLock(commandId, () -> {
+ retryCommand(commandId, now);
+ return true;
+ });
}
private void retryCommand(String commandId, long now) {
@@ -178,24 +213,24 @@ public class MqttCommandAckService implements IDeviceCommandAckHandler {
return;
}
if (!isDeviceOnline(command.getDeviceNo())) {
- command.setRetryCount(command.getRetryCount() + 1);
- if (command.getRetryCount() >= mqttProperties.getCommandAck().getMaxRetryCount()) {
- deletePending(commandId);
- log.warn("[MQTT] 设备离线,命令达到最大等待次数,已删除待确认命令 设备编号={} 命令编号={}", command.getDeviceNo(), commandId);
- return;
- }
command.setNextRetryAt(now + mqttProperties.getCommandAck().getRetryIntervalMs());
savePending(command);
- log.debug("[MQTT] 设备离线,命令延后重试 设备编号={} 命令编号={} 等待次数={}", command.getDeviceNo(), commandId, command.getRetryCount());
+ log.debug("[MQTT] 设备离线,命令延后重试 设备编号={} 命令编号={}", command.getDeviceNo(), commandId);
return;
}
command.setRetryCount(command.getRetryCount() + 1);
command.setLastSentAt(now);
command.setNextRetryAt(now + mqttProperties.getCommandAck().getRetryIntervalMs());
- mqttClientManager.publish(command.getTopic(), JsonUtils.toJsonString(command.getPayload()));
- savePending(command);
- log.info("[MQTT] 命令已重新下发 设备编号={} 命令编号={} 重试次数={}", command.getDeviceNo(), commandId, command.getRetryCount());
+ try {
+ mqttClientManager.publish(command.getTopic(), JsonUtils.toJsonString(command.getPayload()));
+ savePending(command);
+ log.info("[MQTT] 命令已重新下发 设备编号={} 命令编号={} 重试次数={}", command.getDeviceNo(), commandId, command.getRetryCount());
+ } catch (RuntimeException e) {
+ savePending(command);
+ log.warn("[MQTT] 命令重新下发失败,已延后重试 设备编号={} 命令编号={} 重试次数={}",
+ command.getDeviceNo(), commandId, command.getRetryCount(), e);
+ }
}
private boolean isDeviceOnline(String deviceNo) {
@@ -249,7 +284,7 @@ public class MqttCommandAckService implements IDeviceCommandAckHandler {
pendingIds().remove(commandId);
}
- private boolean withCommandLock(String commandId, Runnable action) {
+ private boolean withCommandLock(String commandId, BooleanSupplier action) {
RLock lock = RedisUtils.getClient().getLock(mqttProperties.getCommandAck().getRetryLockKeyPrefix() + commandId);
boolean locked = false;
try {
@@ -257,8 +292,7 @@ public class MqttCommandAckService implements IDeviceCommandAckHandler {
if (!locked) {
return false;
}
- action.run();
- return true;
+ return action.getAsBoolean();
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
log.warn("[MQTT] 命令锁等待被中断 命令编号={}", commandId);
diff --git a/water-common/water-common-mqtt/src/main/java/org/dromara/mqtt/config/properties/MqttProperties.java b/water-common/water-common-mqtt/src/main/java/org/dromara/mqtt/config/properties/MqttProperties.java
index 916ad85..3f0569a 100644
--- a/water-common/water-common-mqtt/src/main/java/org/dromara/mqtt/config/properties/MqttProperties.java
+++ b/water-common/water-common-mqtt/src/main/java/org/dromara/mqtt/config/properties/MqttProperties.java
@@ -18,7 +18,7 @@ public class MqttProperties {
private String password;
private String clientId;
private int qos = 1;
- private int keepAlive = 60;
+ private int keepAlive = 300;
private int connectionTimeout = 30;
private int maxInflight = 1000;
private boolean cleanSession = false;
@@ -65,7 +65,7 @@ public class MqttProperties {
private String retryLockKeyPrefix = "lock:mqtt:command:retry:";
private long retryLockTtlMs = 30000;
private String deviceStatusCachePrefix = "mqtt:device:status:";
- private long deviceStatusCacheTtlSeconds = 300;
+ private long deviceStatusCacheTtlSeconds = 900;
}
}
diff --git a/water-common/water-common-mqtt/src/test/java/org/dromara/mqtt/DeviceMqttCommandPublisherTest.java b/water-common/water-common-mqtt/src/test/java/org/dromara/mqtt/DeviceMqttCommandPublisherTest.java
new file mode 100644
index 0000000..e2799ce
--- /dev/null
+++ b/water-common/water-common-mqtt/src/test/java/org/dromara/mqtt/DeviceMqttCommandPublisherTest.java
@@ -0,0 +1,26 @@
+package org.dromara.mqtt;
+
+import org.dromara.mqtt.config.properties.MqttProperties;
+import org.junit.jupiter.api.Tag;
+import org.junit.jupiter.api.Test;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.springframework.test.util.ReflectionTestUtils.invokeMethod;
+
+@Tag("dev")
+class DeviceMqttCommandPublisherTest {
+
+ @Test
+ void buildCommandTopicDefaultsToDeviceNo() {
+ MqttProperties properties = new MqttProperties();
+ DeviceMqttCommandPublisher publisher = new DeviceMqttCommandPublisher(
+ null,
+ null,
+ properties
+ );
+
+ String topic = invokeMethod(publisher, "buildCommandTopic", "D01");
+
+ assertThat(topic).isEqualTo("/d01/subscriber/cmd");
+ }
+}
diff --git a/water-common/water-common-mqtt/src/test/java/org/dromara/mqtt/MqttCommandAckServiceTest.java b/water-common/water-common-mqtt/src/test/java/org/dromara/mqtt/MqttCommandAckServiceTest.java
new file mode 100644
index 0000000..147c0c8
--- /dev/null
+++ b/water-common/water-common-mqtt/src/test/java/org/dromara/mqtt/MqttCommandAckServiceTest.java
@@ -0,0 +1,122 @@
+package org.dromara.mqtt;
+
+import cn.hutool.extra.spring.SpringUtil;
+import com.baomidou.mybatisplus.core.metadata.TableInfoHelper;
+import org.apache.ibatis.builder.MapperBuilderAssistant;
+import org.apache.ibatis.session.Configuration;
+import org.dromara.app.domain.AppDevice;
+import org.dromara.app.domain.mqtt.DeviceCommand;
+import org.dromara.app.domain.mqtt.DeviceCommandAck;
+import org.dromara.app.mapper.AppDeviceMapper;
+import org.dromara.common.redis.utils.RedisUtils;
+import org.dromara.mqtt.config.properties.MqttProperties;
+import org.junit.jupiter.api.AfterAll;
+import org.junit.jupiter.api.BeforeAll;
+import org.junit.jupiter.api.Tag;
+import org.junit.jupiter.api.Test;
+import org.mockito.MockedStatic;
+import org.redisson.api.RLock;
+import org.redisson.api.RedissonClient;
+import org.springframework.context.support.GenericApplicationContext;
+import org.springframework.test.util.ReflectionTestUtils;
+
+import java.time.Duration;
+import java.util.List;
+import java.util.Map;
+import java.util.concurrent.TimeUnit;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.Mockito.*;
+
+@Tag("dev")
+class MqttCommandAckServiceTest {
+
+ private static GenericApplicationContext applicationContext;
+
+ @BeforeAll
+ static void initializeRedisUtils() {
+ if (TableInfoHelper.getTableInfo(AppDevice.class) == null) {
+ TableInfoHelper.initTableInfo(new MapperBuilderAssistant(new Configuration(), ""), AppDevice.class);
+ }
+ applicationContext = new GenericApplicationContext();
+ applicationContext.registerBean(RedissonClient.class, () -> mock(RedissonClient.class));
+ applicationContext.refresh();
+ new SpringUtil().setApplicationContext(applicationContext);
+ }
+
+ @AfterAll
+ static void closeApplicationContext() {
+ applicationContext.close();
+ }
+
+ @Test
+ void resolveMissingCommandId_usesSinglePendingCommandForDevice() {
+ DeviceCommandAck ack = new DeviceCommandAck();
+ ack.setCommandId("");
+
+ DeviceCommand pending = new DeviceCommand();
+ pending.setCommandId("cmd-1");
+ pending.setDeviceNo("D01");
+
+ boolean resolved = MqttCommandAckService.resolveMissingCommandId(ack, List.of(pending));
+
+ assertThat(resolved).isTrue();
+ assertThat(ack.getCommandId()).isEqualTo("cmd-1");
+ assertThat(ack.getStatus()).isEqualTo("1");
+ }
+
+ @Test
+ void resolveMissingCommandId_refusesAmbiguousPendingCommands() {
+ DeviceCommandAck ack = new DeviceCommandAck();
+ ack.setCommandId("");
+
+ DeviceCommand first = new DeviceCommand();
+ first.setCommandId("cmd-1");
+ DeviceCommand second = new DeviceCommand();
+ second.setCommandId("cmd-2");
+
+ boolean resolved = MqttCommandAckService.resolveMissingCommandId(ack, List.of(first, second));
+
+ assertThat(resolved).isFalse();
+ assertThat(ack.getCommandId()).isEmpty();
+ }
+
+ @Test
+ void ackMatchesPendingCommandRequiresSameDevice() {
+ DeviceCommand pending = new DeviceCommand();
+ pending.setDeviceNo("D01");
+
+ assertThat(MqttCommandAckService.ackMatchesPendingCommand("D01", pending)).isTrue();
+ assertThat(MqttCommandAckService.ackMatchesPendingCommand("D02", pending)).isFalse();
+ assertThat(MqttCommandAckService.ackMatchesPendingCommand("D01", null)).isFalse();
+ }
+
+ @Test
+ void refreshDeviceOnlineRenewsStatusCacheTtl() throws Exception {
+ MqttProperties properties = new MqttProperties();
+ properties.getCommandAck().setDeviceStatusCacheTtlSeconds(900);
+ AppDeviceMapper appDeviceMapper = mock(AppDeviceMapper.class);
+ RedissonClient redissonClient = mock(RedissonClient.class);
+ RLock lock = mock(RLock.class);
+ when(redissonClient.getLock("lock:mqtt:device:status:D01")).thenReturn(lock);
+ when(lock.tryLock(0, 10, TimeUnit.SECONDS)).thenReturn(true);
+ when(lock.isHeldByCurrentThread()).thenReturn(true);
+
+ MqttCommandAckService service = new MqttCommandAckService(properties, appDeviceMapper);
+ try (MockedStatic redis = mockStatic(RedisUtils.class)) {
+ redis.when(RedisUtils::getClient).thenReturn(redissonClient);
+
+ ReflectionTestUtils.invokeMethod(service, "refreshDeviceOnline", "D01");
+
+ redis.verify(() -> RedisUtils.setCacheObject(
+ eq("mqtt:device:status:D01"),
+ any(Map.class),
+ eq(Duration.ofSeconds(900))
+ ));
+ verify(appDeviceMapper).update(eq(null), any());
+ verify(lock).unlock();
+ }
+ }
+}
diff --git a/water-modules/water-app/src/main/java/org/dromara/app/controller/AppVersionController.java b/water-modules/water-app/src/main/java/org/dromara/app/controller/AppVersionController.java
new file mode 100644
index 0000000..275e4e8
--- /dev/null
+++ b/water-modules/water-app/src/main/java/org/dromara/app/controller/AppVersionController.java
@@ -0,0 +1,106 @@
+package org.dromara.app.controller;
+
+import cn.dev33.satoken.annotation.SaCheckPermission;
+import jakarta.servlet.http.HttpServletResponse;
+import jakarta.validation.constraints.NotEmpty;
+import jakarta.validation.constraints.NotNull;
+import lombok.RequiredArgsConstructor;
+import org.dromara.app.domain.bo.AppVersionBo;
+import org.dromara.app.domain.vo.AppVersionVo;
+import org.dromara.app.service.IAppVersionService;
+import org.dromara.common.core.domain.R;
+import org.dromara.common.core.validate.AddGroup;
+import org.dromara.common.core.validate.EditGroup;
+import org.dromara.common.excel.utils.ExcelUtil;
+import org.dromara.common.idempotent.annotation.RepeatSubmit;
+import org.dromara.common.log.annotation.Log;
+import org.dromara.common.log.enums.BusinessType;
+import org.dromara.common.mybatis.core.page.PageQuery;
+import org.dromara.common.mybatis.core.page.TableDataInfo;
+import org.dromara.common.web.core.BaseController;
+import org.springframework.validation.annotation.Validated;
+import org.springframework.web.bind.annotation.*;
+
+import java.util.List;
+
+/**
+ * APP版本
+ *
+ * @author Lion Li
+ * @date 2026-07-01
+ */
+@Validated
+@RequiredArgsConstructor
+@RestController
+@RequestMapping("/app/version")
+public class AppVersionController extends BaseController {
+
+ private final IAppVersionService appVersionService;
+
+ /**
+ * 查询APP版本列表
+ */
+ @SaCheckPermission("app:version:list")
+ @GetMapping("/list")
+ public TableDataInfo list(AppVersionBo bo, PageQuery pageQuery) {
+ return appVersionService.queryPageList(bo, pageQuery);
+ }
+
+ /**
+ * 导出APP版本列表
+ */
+ @SaCheckPermission("app:version:export")
+ @Log(title = "APP版本", businessType = BusinessType.EXPORT)
+ @PostMapping("/export")
+ public void export(AppVersionBo bo, HttpServletResponse response) {
+ List list = appVersionService.queryList(bo);
+ ExcelUtil.exportExcel(list, "APP版本", AppVersionVo.class, response);
+ }
+
+ /**
+ * 获取APP版本详细信息
+ *
+ * @param id 主键
+ */
+ @SaCheckPermission("app:version:query")
+ @GetMapping("/{id}")
+ public R getInfo(@NotNull(message = "主键不能为空")
+ @PathVariable Long id) {
+ return R.ok(appVersionService.queryById(id));
+ }
+
+ /**
+ * 新增APP版本
+ */
+ @SaCheckPermission("app:version:add")
+ @Log(title = "APP版本", businessType = BusinessType.INSERT)
+ @RepeatSubmit()
+ @PostMapping()
+ public R add(@Validated(AddGroup.class) @RequestBody AppVersionBo bo) {
+ return toAjax(appVersionService.insertByBo(bo));
+ }
+
+ /**
+ * 修改APP版本
+ */
+ @SaCheckPermission("app:version:edit")
+ @Log(title = "APP版本", businessType = BusinessType.UPDATE)
+ @RepeatSubmit()
+ @PutMapping()
+ public R edit(@Validated(EditGroup.class) @RequestBody AppVersionBo bo) {
+ return toAjax(appVersionService.updateByBo(bo));
+ }
+
+ /**
+ * 删除APP版本
+ *
+ * @param ids 主键串
+ */
+ @SaCheckPermission("app:version:remove")
+ @Log(title = "APP版本", businessType = BusinessType.DELETE)
+ @DeleteMapping("/{ids}")
+ public R remove(@NotEmpty(message = "主键不能为空")
+ @PathVariable Long[] ids) {
+ return toAjax(appVersionService.deleteWithValidByIds(List.of(ids), true));
+ }
+}
diff --git a/water-modules/water-app/src/main/java/org/dromara/app/domain/AppVersion.java b/water-modules/water-app/src/main/java/org/dromara/app/domain/AppVersion.java
new file mode 100644
index 0000000..8bb8512
--- /dev/null
+++ b/water-modules/water-app/src/main/java/org/dromara/app/domain/AppVersion.java
@@ -0,0 +1,59 @@
+package org.dromara.app.domain;
+
+import com.baomidou.mybatisplus.annotation.TableId;
+import com.baomidou.mybatisplus.annotation.TableName;
+import lombok.Data;
+import lombok.EqualsAndHashCode;
+import org.dromara.common.mybatis.core.domain.BaseEntity;
+
+import java.io.Serial;
+
+/**
+ * APP版本对象 app_version
+ */
+@Data
+@EqualsAndHashCode(callSuper = true)
+@TableName("app_version")
+public class AppVersion extends BaseEntity {
+
+ @Serial
+ private static final long serialVersionUID = 1L;
+
+ @TableId(value = "id")
+ private Long id;
+
+ /**
+ * 平台 android/ios
+ */
+ private String platform;
+
+ /**
+ * 最新版本号
+ */
+ private String latestVersion;
+
+ /**
+ * 版本序号
+ */
+ private Integer versionCode;
+
+ /**
+ * 是否强制更新 0否 1是
+ */
+ private String forceUpdate;
+
+ /**
+ * 下载地址
+ */
+ private String downloadUrl;
+
+ /**
+ * 更新说明
+ */
+ private String releaseNotes;
+
+ /**
+ * 状态 0停用 1启用
+ */
+ private String status;
+}
diff --git a/water-modules/water-app/src/main/java/org/dromara/app/domain/bo/AppVersionBo.java b/water-modules/water-app/src/main/java/org/dromara/app/domain/bo/AppVersionBo.java
new file mode 100644
index 0000000..9d2f82d
--- /dev/null
+++ b/water-modules/water-app/src/main/java/org/dromara/app/domain/bo/AppVersionBo.java
@@ -0,0 +1,76 @@
+package org.dromara.app.domain.bo;
+
+import io.github.linpeilie.annotations.AutoMapper;
+import jakarta.validation.constraints.NotBlank;
+import jakarta.validation.constraints.NotNull;
+import lombok.Data;
+import lombok.EqualsAndHashCode;
+import org.dromara.app.domain.AppVersion;
+import org.dromara.common.core.validate.AddGroup;
+import org.dromara.common.core.validate.EditGroup;
+import org.dromara.common.mybatis.core.domain.BaseEntity;
+
+/**
+ * APP版本业务对象 app_version
+ *
+ * @author Lion Li
+ * @date 2026-07-01
+ */
+@Data
+@EqualsAndHashCode(callSuper = true)
+@AutoMapper(target = AppVersion.class, reverseConvertGenerate = false)
+public class AppVersionBo extends BaseEntity {
+
+ /**
+ * 主键
+ */
+ @NotNull(message = "主键不能为空", groups = { EditGroup.class })
+ private Long id;
+
+ /**
+ * 平台 android/ios
+ */
+ @NotBlank(message = "平台 android/ios不能为空", groups = { AddGroup.class, EditGroup.class })
+ private String platform;
+
+ /**
+ * 最新版本号
+ */
+ @NotBlank(message = "最新版本号不能为空", groups = { AddGroup.class, EditGroup.class })
+ private String latestVersion;
+
+ /**
+ * 版本序号,用于排序取最新版本
+ */
+ @NotNull(message = "版本序号,用于排序取最新版本不能为空", groups = { AddGroup.class, EditGroup.class })
+ private Long versionCode;
+
+ /**
+ * 是否强制更新 0否 1是
+ */
+ @NotBlank(message = "是否强制更新 0否 1是不能为空", groups = { AddGroup.class, EditGroup.class })
+ private String forceUpdate;
+
+ /**
+ * 下载地址
+ */
+ private String downloadUrl;
+
+ /**
+ * 更新说明
+ */
+ private String releaseNotes;
+
+ /**
+ * 状态 0停用 1启用
+ */
+ @NotBlank(message = "状态 0停用 1启用不能为空", groups = { AddGroup.class, EditGroup.class })
+ private String status;
+
+ /**
+ * 备注
+ */
+ private String remark;
+
+
+}
diff --git a/water-modules/water-app/src/main/java/org/dromara/app/domain/vo/AppVersionCheckVo.java b/water-modules/water-app/src/main/java/org/dromara/app/domain/vo/AppVersionCheckVo.java
new file mode 100644
index 0000000..1f197a2
--- /dev/null
+++ b/water-modules/water-app/src/main/java/org/dromara/app/domain/vo/AppVersionCheckVo.java
@@ -0,0 +1,53 @@
+package org.dromara.app.domain.vo;
+
+import lombok.Data;
+
+import java.io.Serial;
+import java.io.Serializable;
+
+@Data
+public class AppVersionCheckVo implements Serializable {
+
+ @Serial
+ private static final long serialVersionUID = 1L;
+
+ /**
+ * 平台类型:android 或 ios
+ */
+ private String platform;
+
+ /**
+ * APP当前版本号,由客户端上传
+ */
+ private String currentVersion;
+
+ /**
+ * 后台配置的最新版本号
+ */
+ private String latestVersion;
+
+ /**
+ * 最新版本序号,用于版本排序
+ */
+ private Integer versionCode;
+
+ /**
+ * 是否存在可更新版本
+ */
+ private Boolean updateAvailable;
+
+ /**
+ * 是否强制更新
+ */
+ private Boolean forceUpdate;
+
+ /**
+ * APP安装包下载地址
+ */
+ private String downloadUrl;
+
+ /**
+ * 更新说明
+ */
+ private String releaseNotes;
+}
diff --git a/water-modules/water-app/src/main/java/org/dromara/app/domain/vo/AppVersionVo.java b/water-modules/water-app/src/main/java/org/dromara/app/domain/vo/AppVersionVo.java
new file mode 100644
index 0000000..464794b
--- /dev/null
+++ b/water-modules/water-app/src/main/java/org/dromara/app/domain/vo/AppVersionVo.java
@@ -0,0 +1,49 @@
+package org.dromara.app.domain.vo;
+
+import cn.idev.excel.annotation.ExcelIgnoreUnannotated;
+import cn.idev.excel.annotation.ExcelProperty;
+import io.github.linpeilie.annotations.AutoMapper;
+import lombok.Data;
+import org.dromara.app.domain.AppVersion;
+
+import java.io.Serial;
+import java.io.Serializable;
+import java.util.Date;
+
+@Data
+@ExcelIgnoreUnannotated
+@AutoMapper(target = AppVersion.class)
+public class AppVersionVo implements Serializable {
+
+ @Serial
+ private static final long serialVersionUID = 1L;
+
+ private Long id;
+
+ @ExcelProperty(value = "平台 android/ios")
+ private String platform;
+
+ @ExcelProperty(value = "最新版本号")
+ private String latestVersion;
+
+ @ExcelProperty(value = "版本序号,用于排序取最新版本")
+ private Integer versionCode;
+
+ @ExcelProperty(value = "是否强制更新 0否 1是")
+ private String forceUpdate;
+
+ @ExcelProperty(value = "下载地址")
+ private String downloadUrl;
+
+ @ExcelProperty(value = "更新说明")
+ private String releaseNotes;
+
+ @ExcelProperty(value = "状态 0停用 1启用")
+ private String status;
+
+ @ExcelProperty(value = "备注")
+ private String remark;
+
+ @ExcelProperty(value = "创建时间")
+ private Date createTime;
+}
diff --git a/water-modules/water-app/src/main/java/org/dromara/app/handler/DeviceDataHandler.java b/water-modules/water-app/src/main/java/org/dromara/app/handler/DeviceDataHandler.java
index ba253a4..6f81e0c 100644
--- a/water-modules/water-app/src/main/java/org/dromara/app/handler/DeviceDataHandler.java
+++ b/water-modules/water-app/src/main/java/org/dromara/app/handler/DeviceDataHandler.java
@@ -1,25 +1,15 @@
package org.dromara.app.handler;
-import com.baomidou.mybatisplus.core.conditions.update.LambdaUpdateWrapper;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
-import org.dromara.app.domain.AppDevice;
import org.dromara.app.domain.bo.AppDeviceBo;
-import org.dromara.app.mapper.AppDeviceMapper;
import org.dromara.app.mqtt.MqttTopicHandler;
import org.dromara.app.service.IAppDeviceService;
import org.dromara.common.json.utils.JsonUtils;
-import org.dromara.common.redis.utils.RedisUtils;
-import org.redisson.api.RLock;
-import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Component;
-import java.time.Duration;
-import java.time.Instant;
import java.util.Date;
-import java.util.HashMap;
import java.util.Map;
-import java.util.concurrent.TimeUnit;
import java.util.regex.Pattern;
/**
@@ -32,14 +22,6 @@ public class DeviceDataHandler implements MqttTopicHandler {
private static final Pattern PATTERN = Pattern.compile("^/([^/]+)/publish/power$");
private final IAppDeviceService appDeviceService;
- private final AppDeviceMapper appDeviceMapper;
- private static final String DEVICE_STATUS_LOCK_PREFIX = "lock:mqtt:device:status:";
-
- @Value("${mqtt.command-ack.device-status-cache-prefix}")
- private String deviceStatusCachePrefix;
-
- @Value("${mqtt.command-ack.device-status-cache-ttl-seconds}")
- private int deviceStatusCacheTtlSeconds;
private final DeviceIdentityResolver deviceIdentityResolver;
@Override
@@ -60,7 +42,6 @@ public class DeviceDataHandler implements MqttTopicHandler {
Map dto = JsonUtils.parseObject(payload, Map.class);
if (dto == null || dto.get("powerLevel") == null) {
log.warn("[MQTT] 设备电量上报缺少电量字段 时间={} 设备编号={} 消息体={}", HandlerLogTime.now(), deviceNo, payload);
- refreshDeviceOnline(deviceNo);
return;
}
// 更新设备电量 + 同步在线状态到数据库
@@ -68,50 +49,11 @@ public class DeviceDataHandler implements MqttTopicHandler {
appDeviceBo.setDeviceNo(deviceNo);
appDeviceBo.setPowerLevel(dto.get("powerLevel").toString());
appDeviceBo.setPowerLevelUpdatatime(new Date());
- appDeviceBo.setStatus("1"); // 收到数据 = 在线
appDeviceService.updateByBo(appDeviceBo);
- // 同步刷新 Redis 在线缓存(定时任务离线检测依赖此 Key)
- refreshDeviceOnline(deviceNo);
log.info("[MQTT] 设备电量更新 时间={} 设备编号={} 消息体={}", HandlerLogTime.now(), deviceNo, payload);
} catch (Exception e) {
log.error("[MQTT] 设备电量更新失败 时间={} 设备标识={} 设备编号={} 消息体={}",
HandlerLogTime.now(), deviceIdentity, deviceNo, payload, e);
}
}
-
- /**
- * 刷新设备在线状态到 Redis(与 MqttCommandAckService 逻辑一致)
- * 写入幂等,未拿到锁则跳过由后续上报兜底
- */
- private void refreshDeviceOnline(String deviceNo) {
- RLock lock = RedisUtils.getClient().getLock(DEVICE_STATUS_LOCK_PREFIX + deviceNo);
- boolean locked = false;
- try {
- locked = lock.tryLock(0, 10, TimeUnit.SECONDS);
- if (!locked) {
- return;
- }
- Map statusCache = new HashMap<>();
- statusCache.put("deviceNo", deviceNo);
- statusCache.put("status", "1");
- statusCache.put("lastReportTime", Instant.now().toString());
- RedisUtils.setCacheObject(
- deviceStatusCachePrefix + deviceNo,
- statusCache,
- Duration.ofSeconds(deviceStatusCacheTtlSeconds)
- );
- appDeviceMapper.update(null,
- new LambdaUpdateWrapper()
- .set(AppDevice::getStatus, "1")
- .eq(AppDevice::getDeviceNo, deviceNo)
- );
- } catch (InterruptedException e) {
- Thread.currentThread().interrupt();
- log.warn("[MQTT] 刷新设备在线状态被中断 设备编号={}", deviceNo);
- } finally {
- if (locked && lock.isHeldByCurrentThread()) {
- lock.unlock();
- }
- }
- }
}
diff --git a/water-modules/water-app/src/main/java/org/dromara/app/handler/DeviceRegisterHandler.java b/water-modules/water-app/src/main/java/org/dromara/app/handler/DeviceRegisterHandler.java
index 2cbaf7f..dd64f82 100644
--- a/water-modules/water-app/src/main/java/org/dromara/app/handler/DeviceRegisterHandler.java
+++ b/water-modules/water-app/src/main/java/org/dromara/app/handler/DeviceRegisterHandler.java
@@ -1,6 +1,5 @@
package org.dromara.app.handler;
-import com.baomidou.mybatisplus.core.conditions.update.LambdaUpdateWrapper;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.dromara.app.domain.AppDevice;
@@ -11,19 +10,12 @@ import org.dromara.app.mqtt.MqttTopicHandler;
import org.dromara.app.service.IAppDeviceService;
import org.dromara.app.service.IDeviceCommandPublisher;
import org.dromara.common.json.utils.JsonUtils;
-import org.dromara.common.redis.utils.RedisUtils;
import org.jspecify.annotations.Nullable;
-import org.redisson.api.RLock;
import org.springframework.beans.factory.ObjectProvider;
-import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Component;
-import java.time.Duration;
-import java.time.Instant;
-import java.util.HashMap;
import java.util.Locale;
import java.util.Map;
-import java.util.concurrent.TimeUnit;
import java.util.regex.Pattern;
/**
@@ -39,14 +31,7 @@ DeviceRegisterHandler implements MqttTopicHandler {
private final IAppDeviceService appDeviceService;
private final AppDeviceMapper appDeviceMapper;
- private static final String DEVICE_STATUS_LOCK_PREFIX = "lock:mqtt:device:status:";
private final ObjectProvider commandPublisherProvider;
-
- @Value("${mqtt.command-ack.device-status-cache-prefix}")
- private String deviceStatusCachePrefix;
-
- @Value("${mqtt.command-ack.device-status-cache-ttl-seconds}")
- private int deviceStatusCacheTtlSeconds;
private final DeviceIdentityResolver deviceIdentityResolver;
@Override
@@ -59,7 +44,7 @@ DeviceRegisterHandler implements MqttTopicHandler {
try {
Map dto = JsonUtils.parseObject(payload, Map.class);
AppDevice exists = deviceIdentityResolver.resolve(deviceIdentity);
- String normalizedDeviceMac = resolveDeviceMac(deviceIdentity, dto);
+ String normalizedDeviceMac = resolveDeviceMac(deviceIdentity, dto, exists);
if (exists == null && normalizedDeviceMac != null) {
exists = appDeviceMapper.selectByMac(normalizedDeviceMac);
}
@@ -75,7 +60,6 @@ DeviceRegisterHandler implements MqttTopicHandler {
}
AppDeviceBo device = buildRegisterDevice(deviceNo, normalizedDeviceMac, dto);
appDeviceService.registerByMqtt(device);
- refreshDeviceOnline(deviceNo);
sendDeviceNoToDevice(deviceNo, normalizedDeviceMac);
log.info("[MQTT] 设备注册 时间={} MAC={} 设备编号={} 消息体={}",
HandlerLogTime.now(), normalizedDeviceMac, deviceNo, payload);
@@ -88,11 +72,9 @@ DeviceRegisterHandler implements MqttTopicHandler {
AppDeviceBo device = new AppDeviceBo();
device.setDeviceNo(deviceNo);
device.setMacAddress(deviceMac);
- device.setStatus("1"); // 注册即在线
if (dto == null) {
return device;
}
- device.setDeviceName(valueAsString(dto.get("deviceName")));
device.setDeviceInitName(valueAsString(dto.get("deviceName")));
device.setPowerLevel(valueAsString(dto.get("powerLevel")));
device.setDeviceEm(valueAsString(dto.get("deviceEm")));
@@ -101,9 +83,15 @@ DeviceRegisterHandler implements MqttTopicHandler {
return device;
}
- private String resolveDeviceMac(String topicDeviceMac, Map dto) {
+ private String resolveDeviceMac(String topicDeviceMac, Map dto, AppDevice exists) {
String payloadMac = dto == null ? null : firstNotBlank(dto.get("deviceMac"), dto.get("macAddress"));
- return normalizeMacAddress(firstNotBlank(payloadMac, topicDeviceMac));
+ if (payloadMac != null && !payloadMac.isBlank()) {
+ return normalizeMacAddress(payloadMac);
+ }
+ if (exists != null && exists.getMacAddress() != null && !exists.getMacAddress().isBlank()) {
+ return normalizeMacAddress(exists.getMacAddress());
+ }
+ return normalizeMacAddress(topicDeviceMac);
}
private String normalizeMacAddress(String macAddress) {
@@ -157,39 +145,4 @@ DeviceRegisterHandler implements MqttTopicHandler {
private String valueAsString(Object value) {
return value == null ? null : String.valueOf(value);
}
-
- /**
- * 刷新设备在线状态到 Redis;状态写入幂等,未拿到锁则跳过由后续上报兜底
- */
- private void refreshDeviceOnline(String deviceNo) {
- RLock lock = RedisUtils.getClient().getLock(DEVICE_STATUS_LOCK_PREFIX + deviceNo);
- boolean locked = false;
- try {
- locked = lock.tryLock(0, 10, TimeUnit.SECONDS);
- if (!locked) {
- return;
- }
- Map statusCache = new HashMap<>();
- statusCache.put("deviceNo", deviceNo);
- statusCache.put("status", "1");
- statusCache.put("lastReportTime", Instant.now().toString());
- RedisUtils.setCacheObject(
- deviceStatusCachePrefix + deviceNo,
- statusCache,
- Duration.ofSeconds(deviceStatusCacheTtlSeconds)
- );
- appDeviceMapper.update(null,
- new LambdaUpdateWrapper()
- .set(AppDevice::getStatus, "1")
- .eq(AppDevice::getDeviceNo, deviceNo)
- );
- } catch (InterruptedException e) {
- Thread.currentThread().interrupt();
- log.warn("[MQTT] 刷新设备在线状态被中断 设备编号={}", deviceNo);
- } finally {
- if (locked && lock.isHeldByCurrentThread()) {
- lock.unlock();
- }
- }
- }
}
diff --git a/water-modules/water-app/src/main/java/org/dromara/app/handler/DeviceStatusHandler.java b/water-modules/water-app/src/main/java/org/dromara/app/handler/DeviceStatusHandler.java
new file mode 100644
index 0000000..67cf924
--- /dev/null
+++ b/water-modules/water-app/src/main/java/org/dromara/app/handler/DeviceStatusHandler.java
@@ -0,0 +1,90 @@
+package org.dromara.app.handler;
+
+import com.fasterxml.jackson.core.type.TypeReference;
+import com.fasterxml.jackson.databind.ObjectMapper;
+import lombok.RequiredArgsConstructor;
+import lombok.extern.slf4j.Slf4j;
+import org.dromara.app.mqtt.MqttTopicHandler;
+import org.dromara.common.core.utils.StringUtils;
+import org.springframework.stereotype.Component;
+
+import java.util.Locale;
+import java.util.Map;
+import java.util.regex.Pattern;
+
+/**
+ * 设备在线/离线状态处理器,匹配 /{deviceIdentity}/publish/status。
+ *
+ * deviceIdentity 支持设备编号;设备遗嘱消息允许使用 MAC 地址,处理前统一解析为设备编号。
+ * 设备通过 status=online 标记上线,通过 LWT status=offline 标记异常离线。
+ */
+@Slf4j
+@Component
+@RequiredArgsConstructor
+public class DeviceStatusHandler implements MqttTopicHandler {
+
+ private static final Pattern PATTERN = Pattern.compile("^/([^/]+)/publish/status$");
+ private static final String OFFLINE_REASON = "设备 MQTT 状态离线";
+
+ private final DeviceIdentityResolver deviceIdentityResolver;
+ private final MqttDeviceStatusService deviceStatusService;
+ private final ObjectMapper objectMapper;
+
+ @Override
+ public Pattern topicPattern() {
+ return PATTERN;
+ }
+
+ @Override
+ public void handle(String deviceIdentity, String payload) {
+ handle(deviceIdentity, payload, false);
+ }
+
+ @Override
+ public void handle(String deviceIdentity, String payload, boolean retained) {
+ String deviceNo = deviceIdentityResolver.resolveDeviceNo(deviceIdentity);
+ if (deviceNo == null) {
+ log.warn("[MQTT] 设备状态上报未找到设备 时间={} 设备标识={} 消息体={}",
+ HandlerLogTime.now(), deviceIdentity, payload);
+ return;
+ }
+
+ String status = parseStatus(payload);
+ if ("online".equals(status) || "1".equals(status)) {
+ deviceStatusService.markOnline(deviceNo);
+ log.info("[MQTT] 设备状态在线 时间={} 设备编号={} 消息体={}", HandlerLogTime.now(), deviceNo, payload);
+ return;
+ }
+ if ("offline".equals(status) || "0".equals(status)) {
+ if (retained) {
+ log.warn("[MQTT] 忽略 retained 离线状态 时间={} 设备编号={} 消息体={}", HandlerLogTime.now(), deviceNo, payload);
+ return;
+ }
+ deviceStatusService.markOffline(deviceNo, OFFLINE_REASON);
+ log.info("[MQTT] 设备状态离线 时间={} 设备编号={} 消息体={}", HandlerLogTime.now(), deviceNo, payload);
+ return;
+ }
+
+ log.warn("[MQTT] 设备状态上报状态未知 时间={} 设备编号={} 消息体={}",
+ HandlerLogTime.now(), deviceNo, payload);
+ }
+
+ private String parseStatus(String payload) {
+ if (StringUtils.isBlank(payload)) {
+ return null;
+ }
+ String trimmed = payload.trim();
+ if (!trimmed.startsWith("{")) {
+ return trimmed.toLowerCase(Locale.ROOT);
+ }
+ try {
+ Map body = objectMapper.readValue(trimmed, new TypeReference