节点中过期用户清除

This commit is contained in:
lys 2026-03-24 17:59:31 +08:00
parent a15f244317
commit 616b981d34
5 changed files with 627 additions and 0 deletions

View File

@ -0,0 +1,495 @@
package org.dromara.job.snailjob;
import cn.hutool.core.date.DateUtil;
import cn.hutool.json.JSONUtil;
import com.aizuda.snailjob.client.job.core.annotation.JobExecutor;
import com.aizuda.snailjob.client.job.core.dto.JobArgs;
import com.aizuda.snailjob.common.log.SnailJobLog;
import com.aizuda.snailjob.model.dto.ExecuteResult;
import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.dromara.common.core.constant.NetCacheNameConstants;
import org.dromara.common.core.enums.NetNodeSoftTypeEnum;
import org.dromara.common.core.enums.NetProtocolSubEnum;
import org.dromara.common.core.constant.NetUserStatusConstants;
import org.dromara.common.redis.utils.RedisUtils;
import org.dromara.common.tenant.helper.TenantHelper;
import org.dromara.net.domain.NetNode;
import org.dromara.net.domain.NetNodeConfig;
import org.dromara.net.domain.NetPlan;
import org.dromara.net.domain.NetUser;
import org.dromara.net.domain.NetUserFlowRecord;
import org.dromara.net.dto.GetInboundUsersRequest;
import org.dromara.net.dto.GetInboundUsersResponse;
import org.dromara.net.dto.RemoveUsersRequest;
import org.dromara.net.dto.RemoveUsersResponse;
import org.dromara.net.mapper.NetNodeConfigMapper;
import org.dromara.net.mapper.NetNodeMapper;
import org.dromara.net.mapper.NetPlanMapper;
import org.dromara.net.mapper.NetUserFlowRecordMapper;
import org.dromara.net.mapper.NetUserMapper;
import org.dromara.net.util.RemnaNodeHttpClient;
import org.springframework.stereotype.Component;
import java.util.ArrayList;
import java.util.Date;
import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.LinkedBlockingQueue;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.stream.Collectors;
/**
* Remna节点用户验证定时任务
* <p>
* 功能说明
* 1. 查询所有 online_status xray_status true remna 节点
* 2. 通过线程池为每个节点调用 /node/handler/get-inbound-users 接口获取节点的所有用户
* 3. 批量查询数据库 NetUser NetPlan
* 4. 判断用户是否应该被清理
* - 流量超过订阅计划的最大流量
* - 用户不在有效期
* - 对于符合条件的用户如果不存在 redis 信息且上个小时无流量记录
* 5. 汇总后批量调用 /node/handler/remove-users
* </p>
*
* @author lys
* @date 2025-03-24
*/
@Slf4j
@Component
@RequiredArgsConstructor
@JobExecutor(name = "remnaNodeVerifyUserJob")
public class RemnaNodeVerifyUserJob {
private final NetNodeMapper netNodeMapper;
private final NetNodeConfigMapper netNodeConfigMapper;
private final NetUserMapper netUserMapper;
private final NetPlanMapper netPlanMapper;
private final NetUserFlowRecordMapper netUserFlowRecordMapper;
/**
* 请求超时时间毫秒
*/
private static final int REQUEST_TIMEOUT = 10000;
/**
* 获取入站用户接口路径
*/
private static final String GET_INBOUND_USERS_PATH = "/node/handler/get-inbound-users";
/**
* 批量移除用户接口路径
*/
private static final String REMOVE_USERS_PATH = "/node/handler/remove-users";
/**
* 用于节点用户验证的线程池
*/
private static final ExecutorService NODE_VERIFY_EXECUTOR = new ThreadPoolExecutor(
8,
32,
60L,
TimeUnit.SECONDS,
new LinkedBlockingQueue<>(300),
r -> {
Thread t = new Thread(r, "node-verify-" + System.currentTimeMillis());
t.setDaemon(true);
return t;
},
new ThreadPoolExecutor.CallerRunsPolicy()
);
public ExecuteResult jobExecute(JobArgs jobArgs) {
SnailJobLog.LOCAL.info("remnaNodeVerifyUserJob 开始执行. jobArgs:{}", JSONUtil.toJsonStr(jobArgs));
SnailJobLog.REMOTE.info("remnaNodeVerifyUserJob 开始执行");
// 1. 查询所有可用的REMNAWAVE节点online_status和xray_status都为true
List<NetNode> nodes = queryAvailableNodes();
SnailJobLog.LOCAL.info("remnaNodeVerifyUserJob 找到 {} 个可用RemnaWave节点", nodes.size());
SnailJobLog.REMOTE.info("remnaNodeVerifyUserJob 找到 {} 个可用RemnaWave节点", nodes.size());
if (nodes.isEmpty()) {
return ExecuteResult.success("没有可用的RemnaWave节点");
}
// 获取上一小时的时间范围
Date lastHourStart = DateUtil.beginOfHour(DateUtil.offsetHour(new Date(), -1));
Date lastHourEnd = DateUtil.endOfHour(lastHourStart);
// 2. 并行处理每个节点 // 存储每个节点需要清理的用户
ConcurrentHashMap<Long, List<RemoveUsersRequest.UserRemoveInfo>> nodeUsersToRemove = new ConcurrentHashMap<>();
List<CompletableFuture<Void>> futures = nodes.stream()
.map(node -> CompletableFuture.runAsync(() -> {
try {
processNode(node, nodeUsersToRemove, lastHourStart, lastHourEnd);
} catch (Exception e) {
SnailJobLog.LOCAL.error("remnaNodeVerifyUserJob 处理节点异常: nodeId={}, nodeName={}, error={}",
node.getId(), node.getNodeName(), e.getMessage(), e);
}
}, NODE_VERIFY_EXECUTOR))
.toList();
// 等待所有节点处理完成
try {
CompletableFuture.allOf(futures.toArray(new CompletableFuture[0])).join();
} catch (Exception e) {
SnailJobLog.LOCAL.error("remnaNodeVerifyUserJob 等待节点处理完成异常", e);
}
// 3. 构建节点Map方便查找
Map<Long, NetNode> nodeMap = nodes.stream()
.collect(Collectors.toMap(NetNode::getId, n -> n));
// 4. 并行批量移除用户
AtomicInteger totalRemoved = new AtomicInteger(0);
AtomicInteger successNodeCount = new AtomicInteger(0);
AtomicInteger failNodeCount = new AtomicInteger(0);
List<CompletableFuture<Void>> removeFutures = nodeUsersToRemove.entrySet().stream()
.filter(entry -> !entry.getValue().isEmpty())
.map(entry -> CompletableFuture.runAsync(() -> {
Long nodeId = entry.getKey();
List<RemoveUsersRequest.UserRemoveInfo> usersToRemove = entry.getValue();
NetNode node = nodeMap.get(nodeId);
if (node == null) {
return;
}
try {
boolean success = removeUsersFromNode(node, usersToRemove);
if (success) {
totalRemoved.addAndGet(usersToRemove.size());
successNodeCount.incrementAndGet();
SnailJobLog.LOCAL.info("remnaNodeVerifyUserJob 节点 {} 移除 {} 个用户成功",
node.getNodeName(), usersToRemove.size());
} else {
failNodeCount.incrementAndGet();
}
} catch (Exception e) {
failNodeCount.incrementAndGet();
SnailJobLog.LOCAL.error("remnaNodeVerifyUserJob 节点 {} 移除用户异常: {}",
node.getNodeName(), e.getMessage(), e);
}
}, NODE_VERIFY_EXECUTOR))
.toList();
// 等待所有删除操作完成
try {
CompletableFuture.allOf(removeFutures.toArray(new CompletableFuture[0])).join();
} catch (Exception e) {
SnailJobLog.LOCAL.error("remnaNodeVerifyUserJob 等待删除操作完成异常", e);
}
String result = String.format("执行完成. 处理节点: %d, 成功移除节点数: %d, 失败节点数: %d, 总移除用户数: %d",
nodes.size(), successNodeCount.get(), failNodeCount.get(), totalRemoved.get());
SnailJobLog.LOCAL.info("remnaNodeVerifyUserJob {}", result);
SnailJobLog.REMOTE.info("remnaNodeVerifyUserJob {}", result);
return ExecuteResult.success(result);
}
/**
* 查询可用的REMNAWAVE节点
*/
private List<NetNode> queryAvailableNodes() {
return TenantHelper.ignore(() -> {
LambdaQueryWrapper<NetNode> queryWrapper = new LambdaQueryWrapper<>();
queryWrapper.eq(NetNode::getNodeSoftType, NetNodeSoftTypeEnum.REMNAWAVE.getNodeValue())
.eq(NetNode::getOnlineStatus, true)
.eq(NetNode::getXrayStatus, true);
return netNodeMapper.selectList(queryWrapper);
});
}
/**
* 处理单个节点
*/
private void processNode(NetNode node,
ConcurrentHashMap<Long, List<RemoveUsersRequest.UserRemoveInfo>> nodeUsersToRemove,
Date lastHourStart,
Date lastHourEnd) {
// 获取节点配置
NetNodeConfig config = TenantHelper.ignore(() -> {
LambdaQueryWrapper<NetNodeConfig> configQueryWrapper = new LambdaQueryWrapper<>();
configQueryWrapper.eq(NetNodeConfig::getNodeId, node.getId());
return netNodeConfigMapper.selectOne(configQueryWrapper);
});
if (config == null) {
SnailJobLog.LOCAL.warn("remnaNodeVerifyUserJob 节点 {} 配置不存在,跳过", node.getNodeName());
return;
}
// 获取节点的入站标签
String inboundTag = getInboundTag(node.getProtocolType());
if (inboundTag == null) {
SnailJobLog.LOCAL.warn("remnaNodeVerifyUserJob 节点 {} 协议类型 {} 无法确定入站标签,跳过",
node.getNodeName(), node.getProtocolType());
return;
}
// 调用节点接口获取用户列表
String nodeUrl = RemnaNodeHttpClient.buildNodeUrl(node.getNodeIp(), node.getServerPort());
String getUrl = nodeUrl + GET_INBOUND_USERS_PATH;
GetInboundUsersRequest request = new GetInboundUsersRequest();
request.setTag(inboundTag);
GetInboundUsersResponse response = RemnaNodeHttpClient.postAndParse(
config, getUrl, request, GetInboundUsersResponse.class, REQUEST_TIMEOUT
);
if (response == null || response.getResponse() == null || response.getResponse().getUsers() == null) {
SnailJobLog.LOCAL.warn("remnaNodeVerifyUserJob 节点 {} 获取用户列表失败或为空", node.getNodeName());
return;
}
List<GetInboundUsersResponse.UserInfo> nodeUsers = response.getResponse().getUsers();
if (nodeUsers.isEmpty()) {
SnailJobLog.LOCAL.debug("remnaNodeVerifyUserJob 节点 {} 无用户", node.getNodeName());
return;
}
SnailJobLog.LOCAL.info("remnaNodeVerifyUserJob 节点 {} 有 {} 个用户", node.getNodeName(), nodeUsers.size());
// 解析用户ID列表
Set<Long> userIds = nodeUsers.stream()
.map(GetInboundUsersResponse.UserInfo::getUsername)
.filter(username -> username != null && !username.isEmpty())
.map(this::parseUserId)
.filter(id -> id != null)
.collect(Collectors.toSet());
if (userIds.isEmpty()) {
SnailJobLog.LOCAL.debug("remnaNodeVerifyUserJob 节点 {} 无有效用户ID", node.getNodeName());
return;
}
// 批量查询数据库中的用户和订阅计划
List<NetUser> dbUsers = TenantHelper.ignore(() -> {
LambdaQueryWrapper<NetUser> userQueryWrapper = new LambdaQueryWrapper<>();
userQueryWrapper.in(NetUser::getId, userIds);
return netUserMapper.selectList(userQueryWrapper);
});
if (dbUsers.isEmpty()) {
SnailJobLog.LOCAL.debug("remnaNodeVerifyUserJob 节点 {} 的用户在数据库中不存在", node.getNodeName());
// 所有节点用户都不在数据库中全部清理
List<RemoveUsersRequest.UserRemoveInfo> usersToRemove = nodeUsers.stream()
.map(user -> {
RemoveUsersRequest.UserRemoveInfo info = new RemoveUsersRequest.UserRemoveInfo();
info.setUserId(user.getUsername());
info.setHashUuid(""); // 无法获取nodeToken设置为空
return info;
})
.toList();
nodeUsersToRemove.put(node.getId(), usersToRemove);
return;
}
// 构建用户Map
Map<Long, NetUser> dbUserMap = dbUsers.stream()
.collect(Collectors.toMap(NetUser::getId, u -> u));
// 获取所有用户的planId
Set<Long> planIds = dbUsers.stream()
.map(NetUser::getPlanId)
.filter(id -> id != null)
.collect(Collectors.toSet());
// 批量查询订阅计划
Map<Long, NetPlan> planMap = new java.util.HashMap<>();
if (!planIds.isEmpty()) {
List<NetPlan> plans = TenantHelper.ignore(() -> {
LambdaQueryWrapper<NetPlan> planQueryWrapper = new LambdaQueryWrapper<>();
planQueryWrapper.in(NetPlan::getId, planIds);
return netPlanMapper.selectList(planQueryWrapper);
});
planMap = plans.stream()
.collect(Collectors.toMap(NetPlan::getId, p -> p));
}
// 查询上一小时有流量记录的用户
Set<Long> usersWithFlowRecord = queryUsersWithFlowRecord(userIds, lastHourStart, lastHourEnd);
// 判断每个用户是否需要清理
List<RemoveUsersRequest.UserRemoveInfo> usersToRemove = new ArrayList<>();
Date now = new Date();
for (GetInboundUsersResponse.UserInfo nodeUser : nodeUsers) {
Long userId = parseUserId(nodeUser.getUsername());
if (userId == null) {
continue;
}
NetUser dbUser = dbUserMap.get(userId);
if (dbUser == null) {
// 用户不在数据库中需要清理
RemoveUsersRequest.UserRemoveInfo info = new RemoveUsersRequest.UserRemoveInfo();
info.setUserId(nodeUser.getUsername());
info.setHashUuid("");
usersToRemove.add(info);
continue;
}
// 检查用户状态
if (!NetUserStatusConstants.NORMAL.equals(dbUser.getStatus())) {
// 用户状态异常需要清理
RemoveUsersRequest.UserRemoveInfo info = new RemoveUsersRequest.UserRemoveInfo();
info.setUserId(nodeUser.getUsername());
info.setHashUuid(dbUser.getNodeToken());
usersToRemove.add(info);
continue;
}
// 检查有效期
boolean expired = dbUser.getValidityTime() != null && now.after(dbUser.getValidityTime());
// 检查流量是否超限
boolean flowExceeded = false;
NetPlan plan = planMap.get(dbUser.getPlanId());
if (plan != null && plan.getTotalFlow() != null && plan.getTotalFlow() > 0) {
Long monthUseFlow = dbUser.getMonthUseFlow() != null ? dbUser.getMonthUseFlow() : 0L;
if (monthUseFlow >= plan.getTotalFlow()) {
flowExceeded = true;
}
}
if (expired || flowExceeded) {
// 过期或流量超限需要清理
RemoveUsersRequest.UserRemoveInfo info = new RemoveUsersRequest.UserRemoveInfo();
info.setUserId(nodeUser.getUsername());
info.setHashUuid(dbUser.getNodeToken());
usersToRemove.add(info);
continue;
}
// 对于没有过期且没有超出流量限制的用户
// 检查 redis 信息和流量记录
String cacheKey = NetCacheNameConstants.NODE_SUB_ACTIVE_USER + node.getId() + ":" + userId;
boolean hasRedisCache = Boolean.TRUE.equals(RedisUtils.hasKey(cacheKey));
boolean hasFlowRecord = usersWithFlowRecord.contains(userId);
if (!hasRedisCache && !hasFlowRecord) {
// 没有 redis 信息且上个小时无流量记录需要清理
RemoveUsersRequest.UserRemoveInfo info = new RemoveUsersRequest.UserRemoveInfo();
info.setUserId(nodeUser.getUsername());
info.setHashUuid(dbUser.getNodeToken());
usersToRemove.add(info);
SnailJobLog.LOCAL.debug("remnaNodeVerifyUserJob 节点 {} 用户 {} 无redis信息且无流量记录加入清理列表",
node.getNodeName(), userId);
}
}
if (!usersToRemove.isEmpty()) {
nodeUsersToRemove.put(node.getId(), usersToRemove);
SnailJobLog.LOCAL.info("remnaNodeVerifyUserJob 节点 {} 需要清理 {} 个用户",
node.getNodeName(), usersToRemove.size());
}
}
/**
* 获取入站标签
*/
private String getInboundTag(String protocolType) {
if (protocolType == null) {
return null;
}
return protocolType + "-inbound";
}
/**
* 解析用户ID
*/
private Long parseUserId(String username) {
if (username == null || username.isEmpty()) {
return null;
}
try {
return Long.parseLong(username);
} catch (NumberFormatException e) {
return null;
}
}
/**
* 查询上一小时有流量记录的用户
*/
private Set<Long> queryUsersWithFlowRecord(Set<Long> userIds, Date startTime, Date endTime) {
if (userIds.isEmpty()) {
return new HashSet<>();
}
return TenantHelper.ignore(() -> {
LambdaQueryWrapper<NetUserFlowRecord> queryWrapper = new LambdaQueryWrapper<>();
queryWrapper.in(NetUserFlowRecord::getUserId, userIds)
.ge(NetUserFlowRecord::getCreateTime, startTime)
.le(NetUserFlowRecord::getCreateTime, endTime)
.select(NetUserFlowRecord::getUserId);
List<NetUserFlowRecord> records = netUserFlowRecordMapper.selectList(queryWrapper);
return records.stream()
.map(NetUserFlowRecord::getUserId)
.collect(Collectors.toSet());
});
}
/**
* 从节点批量移除用户
*/
private boolean removeUsersFromNode(NetNode node, List<RemoveUsersRequest.UserRemoveInfo> usersToRemove) {
// 获取节点配置
NetNodeConfig config = TenantHelper.ignore(() -> {
LambdaQueryWrapper<NetNodeConfig> configQueryWrapper = new LambdaQueryWrapper<>();
configQueryWrapper.eq(NetNodeConfig::getNodeId, node.getId());
return netNodeConfigMapper.selectOne(configQueryWrapper);
});
if (config == null) {
SnailJobLog.LOCAL.warn("remnaNodeVerifyUserJob 节点 {} 配置不存在,无法移除用户", node.getNodeName());
return false;
}
String nodeUrl = RemnaNodeHttpClient.buildNodeUrl(node.getNodeIp(), node.getServerPort());
String url = nodeUrl + REMOVE_USERS_PATH;
RemoveUsersRequest request = new RemoveUsersRequest();
request.setUsers(usersToRemove);
try {
RemoveUsersResponse response = RemnaNodeHttpClient.postAndParse(
config, url, request, RemoveUsersResponse.class, REQUEST_TIMEOUT
);
if (response != null && response.getResponse() != null) {
if (Boolean.TRUE.equals(response.getResponse().getSuccess())) {
return true;
} else {
SnailJobLog.LOCAL.warn("remnaNodeVerifyUserJob 节点 {} 移除用户失败: {}",
node.getNodeName(), response.getResponse().getError());
return false;
}
} else {
SnailJobLog.LOCAL.warn("remnaNodeVerifyUserJob 节点 {} 移除用户响应为空", node.getNodeName());
return false;
}
} catch (Exception e) {
SnailJobLog.LOCAL.error("remnaNodeVerifyUserJob 节点 {} 移除用户异常: {}",
node.getNodeName(), e.getMessage(), e);
return false;
}
}
}

View File

@ -0,0 +1,19 @@
package org.dromara.net.dto;
import lombok.Data;
/**
* 获取入站用户列表请求DTO
* 对应 POST /node/handler/get-inbound-users 接口请求
*
* @author lys
* @date 2025-03-24
*/
@Data
public class GetInboundUsersRequest {
/**
* 入站标签
*/
private String tag;
}

View File

@ -0,0 +1,47 @@
package org.dromara.net.dto;
import lombok.Data;
import java.util.List;
/**
* 获取入站用户列表响应DTO
* 对应 POST /node/handler/get-inbound-users 接口响应
*
* @author lys
* @date 2025-03-24
*/
@Data
public class GetInboundUsersResponse {
/**
* 响应数据
*/
private ResponseData response;
@Data
public static class ResponseData {
/**
* 用户列表
*/
private List<UserInfo> users;
}
@Data
public static class UserInfo {
/**
* 用户名在节点中是用户ID字符串
*/
private String username;
/**
* 邮箱
*/
private String email;
/**
* 等级
*/
private Integer level;
}
}

View File

@ -0,0 +1,34 @@
package org.dromara.net.dto;
import lombok.Data;
import java.util.List;
/**
* 批量移除用户请求DTO
* 对应 POST /node/handler/remove-users 接口请求
*
* @author lys
* @date 2025-03-24
*/
@Data
public class RemoveUsersRequest {
/**
* 要移除的用户列表
*/
private List<UserRemoveInfo> users;
@Data
public static class UserRemoveInfo {
/**
* 用户ID字符串形式
*/
private String userId;
/**
* 哈希UUID用户的nodeToken
*/
private String hashUuid;
}
}

View File

@ -0,0 +1,32 @@
package org.dromara.net.dto;
import lombok.Data;
/**
* 批量移除用户响应DTO
* 对应 POST /node/handler/remove-users 接口响应
*
* @author lys
* @date 2025-03-24
*/
@Data
public class RemoveUsersResponse {
/**
* 响应数据
*/
private ResponseData response;
@Data
public static class ResponseData {
/**
* 是否成功
*/
private Boolean success;
/**
* 错误信息
*/
private String error;
}
}