u
This commit is contained in:
1 parent
63f1d4d1dc
commit
655581b217
17 files changed
+188
-146
No files matched your search
+17
-13
@@ -1,21 +1,12 @@
|
||||
package com.link.im.controller;
|
||||
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import com.link.api.util.DruidDataSourceManager;
|
||||
import com.link.api.util.PaginationUtils;
|
||||
import com.link.api.util.Result;
|
||||
import com.link.api.util.UserContext;
|
||||
import com.link.im.manager.UserChannelManager;
|
||||
import com.link.im.utils.DbUtils;
|
||||
import org.springframework.util.StringUtils;
|
||||
import org.springframework.web.bind.annotation.*;
|
||||
|
||||
import java.sql.Connection;
|
||||
import java.sql.PreparedStatement;
|
||||
import java.util.HashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.stream.Collectors;
|
||||
|
||||
@RestController
|
||||
@CrossOrigin
|
||||
@@ -29,15 +20,28 @@ public class CompanyController {
|
||||
|
||||
List<Map<String, Object>> maps = DbUtils.queryForList(orgId, sql);
|
||||
|
||||
for (Map<String, Object> map : maps) {
|
||||
String bId = map.get("b_id").toString();
|
||||
return Result.ok("成功",maps);
|
||||
} catch (Exception e) {
|
||||
return Result.fail("公司列表查询异常:" + e.getMessage());
|
||||
}
|
||||
}
|
||||
|
||||
map.put("_online",UserChannelManager.isOnline(bId));
|
||||
}
|
||||
@GetMapping("/api/loadFriendList")
|
||||
public Result<?> loadFriendList(){
|
||||
try {
|
||||
String orgId = UserContext.getOrgId();
|
||||
String sql = "select b_id,b_name from b_user where b_canuse = '1' ";
|
||||
|
||||
List<Map<String, Object>> maps = DbUtils.queryForList(orgId, sql);
|
||||
|
||||
return Result.ok("成功",maps);
|
||||
} catch (Exception e) {
|
||||
return Result.fail("公司列表查询异常:" + e.getMessage());
|
||||
}
|
||||
}
|
||||
|
||||
@GetMapping("/api/loadGroupList")
|
||||
public Result<?> loadGroupList(){
|
||||
return Result.ok("");
|
||||
}
|
||||
}
|
||||
@@ -1,6 +1,5 @@
|
||||
package com.link.im.handler;
|
||||
|
||||
import com.link.api.util.UserContext;
|
||||
import com.link.im.manager.UserChannelManager;
|
||||
import com.link.im.manager.WriteManager;
|
||||
import com.link.im.model.ImRequest;
|
||||
@@ -9,7 +8,6 @@ import com.link.im.utils.CommonUtils;
|
||||
import com.link.im.utils.DbUtils;
|
||||
import io.netty.channel.ChannelHandlerContext;
|
||||
import org.springframework.stereotype.Component;
|
||||
import org.springframework.util.StringUtils;
|
||||
|
||||
import java.sql.SQLException;
|
||||
import java.util.List;
|
||||
@@ -38,35 +36,68 @@ public class ImMessageProcessor {
|
||||
}
|
||||
|
||||
private void handlePing(ChannelHandlerContext ctx) {
|
||||
WriteManager.writeToChannel(ctx.channel(), ImResponse.success("pong",null));
|
||||
WriteManager.writeToChannel(ctx.channel(), ImResponse.success("pong",null,null));
|
||||
}
|
||||
|
||||
private void handleLogin(ChannelHandlerContext ctx,ImRequest request){
|
||||
String fromUserid = request.getFrom();
|
||||
UserChannelManager.addUser(fromUserid,ctx.channel());
|
||||
String orgId = request.getOrgId();
|
||||
|
||||
System.out.println(orgId);
|
||||
String fromUserid = request.getFrom();
|
||||
// 单个上线
|
||||
UserChannelManager.addUser(fromUserid,ctx.channel());
|
||||
// 群组上线
|
||||
|
||||
// 1.通知给好友
|
||||
String sql1 = "select b_frienduser_Id from IM_Friends where b_user_id = '" + fromUserid +"'";
|
||||
try {
|
||||
List<Map<String, Object>> maps = DbUtils.queryForList(orgId, sql1);
|
||||
String friendSql = "select b_frienduser_Id from IM_Friends where b_user_id = '" + fromUserid +"'";
|
||||
|
||||
List<Map<String, Object>> maps = DbUtils.queryForList(orgId, friendSql);
|
||||
|
||||
ImResponse response = ImResponse.success();
|
||||
response.setType("online");
|
||||
response.setFrom(fromUserid);
|
||||
|
||||
for (Map<String, Object> map : maps) {
|
||||
String bFrienduserId = map.get("b_frienduser_Id").toString();
|
||||
if (CommonUtils.isNotEmpty(bFrienduserId)) {
|
||||
boolean friendOnline = UserChannelManager.isOnline(bFrienduserId);
|
||||
if(friendOnline){
|
||||
|
||||
if (CommonUtils.isNotEmpty(bFrienduserId) && UserChannelManager.isOnline(bFrienduserId)) {
|
||||
WriteManager.writeToUser(bFrienduserId,response);
|
||||
}
|
||||
}
|
||||
} catch (SQLException e) {
|
||||
throw new RuntimeException(e);
|
||||
}
|
||||
|
||||
try {
|
||||
// 1. 查询用户所属群组
|
||||
String sqlGroups = "SELECT b_group_id FROM IM_GroupMembers WHERE b_member_id = '" + fromUserid + "'";
|
||||
List<Map<String, Object>> groupList = DbUtils.queryForList(orgId, sqlGroups);
|
||||
|
||||
for (Map<String, Object> g : groupList) {
|
||||
String groupId = (String) g.get("b_group_id");
|
||||
|
||||
if (CommonUtils.isEmpty(groupId)) continue;
|
||||
|
||||
// 2. 查询群组成员
|
||||
String sqlMembers = "SELECT b_member_id FROM IM_GroupMembers WHERE b_group_id = '" + groupId + "'";
|
||||
List<Map<String, Object>> memberList = DbUtils.queryForList(orgId, sqlMembers);
|
||||
|
||||
// 3. 构造消息
|
||||
ImResponse response = ImResponse.success();
|
||||
response.setFrom(fromUserid);
|
||||
response.setData(CommonUtils.mapOf("groupId", groupId));
|
||||
|
||||
// 4. 遍历群成员发送消息
|
||||
for (Map<String, Object> member : memberList) {
|
||||
String memberId = (String) member.get("b_member_id");
|
||||
if (!fromUserid.equals(memberId)) {
|
||||
WriteManager.writeToUser(memberId, response);
|
||||
}
|
||||
}
|
||||
}
|
||||
} catch (SQLException e) {
|
||||
throw new RuntimeException(e);
|
||||
}
|
||||
// 2.通知给群组
|
||||
|
||||
// 3.通知给全部(测试用,先留,后面会删)
|
||||
}
|
||||
|
||||
}
|
||||
@@ -1,10 +1,8 @@
|
||||
package com.link.im.handler;
|
||||
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import com.link.im.manager.GroupManager;
|
||||
import com.link.im.manager.UserChannelManager;
|
||||
import com.link.im.manager.WriteManager;
|
||||
import com.link.im.model.ImMessage;
|
||||
import com.link.im.model.ImRequest;
|
||||
import com.link.im.model.ImResponse;
|
||||
import io.netty.channel.ChannelHandler;
|
||||
@@ -31,9 +29,6 @@ public class ImWebSocketHandler extends SimpleChannelInboundHandler<TextWebSocke
|
||||
// 下线时通过 UserChannelManager 移除
|
||||
UserChannelManager.removeUserByChannel(ctx.channel());
|
||||
|
||||
// 同时从所有群中移除
|
||||
GroupManager.removeUserFromAllGroups(ctx.channel());
|
||||
|
||||
System.out.println("❌ 客户端断开: " + ctx.channel().id().asShortText());
|
||||
}
|
||||
|
||||
|
||||
@@ -1,71 +0,0 @@
|
||||
package com.link.im.manager;
|
||||
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import io.netty.channel.Channel;
|
||||
import io.netty.handler.codec.http.websocketx.TextWebSocketFrame;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.util.Collections;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
|
||||
public class GroupManager {
|
||||
private static final Logger logger = LoggerFactory.getLogger(GroupManager.class);
|
||||
|
||||
// groupId -> Set<userId>
|
||||
private static final ConcurrentHashMap<String, Set<String>> GROUP_MAP = new ConcurrentHashMap<>();
|
||||
|
||||
/**
|
||||
* 创建或获取群成员集合
|
||||
*/
|
||||
private static Set<String> getGroup(String groupId) {
|
||||
return GROUP_MAP.computeIfAbsent(groupId, k -> Collections.newSetFromMap(new ConcurrentHashMap<>()));
|
||||
}
|
||||
|
||||
/**
|
||||
* 添加成员到群
|
||||
*/
|
||||
public static void addUser(String groupId, String userId) {
|
||||
if (groupId == null || userId == null) return;
|
||||
getGroup(groupId).add(userId);
|
||||
}
|
||||
|
||||
/**
|
||||
* 从群中移除成员
|
||||
*/
|
||||
public static void removeUser(String groupId, String userId) {
|
||||
Set<String> members = GROUP_MAP.get(groupId);
|
||||
if (members != null) {
|
||||
members.remove(userId);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 根据 Channel 移除该用户在所有群的成员记录
|
||||
*/
|
||||
public static void removeUserFromAllGroups(Channel channel) {
|
||||
if (channel == null) return;
|
||||
|
||||
// 遍历所有群
|
||||
GROUP_MAP.forEach((groupId, members) -> {
|
||||
members.removeIf(userId -> {
|
||||
Channel ch = UserChannelManager.getChannel(userId);
|
||||
boolean match = ch != null && ch.equals(channel);
|
||||
if (match) {
|
||||
logger.info("用户从群移除: groupId={}, userId={}", groupId, userId);
|
||||
}
|
||||
return match;
|
||||
});
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
* 获取群成员列表(不可修改)
|
||||
*/
|
||||
public static Set<String> getMembers(String groupId) {
|
||||
Set<String> members = GROUP_MAP.get(groupId);
|
||||
if (members == null) return Collections.emptySet();
|
||||
return Collections.unmodifiableSet(members);
|
||||
}
|
||||
}
|
||||
@@ -1,13 +1,12 @@
|
||||
package com.link.im.manager;
|
||||
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import com.link.im.handler.ImWebSocketHandler;
|
||||
import io.netty.channel.Channel;
|
||||
import io.netty.handler.codec.http.websocketx.TextWebSocketFrame;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.util.Set;
|
||||
import java.util.List;
|
||||
|
||||
public class WriteManager {
|
||||
private static final ObjectMapper MAPPER = new ObjectMapper();
|
||||
@@ -15,8 +14,16 @@ public class WriteManager {
|
||||
|
||||
/**
|
||||
* 发送消息 - 根据 userId
|
||||
* 内部会自动判断用户是否在线
|
||||
*/
|
||||
public static void writeToUser(String userId, Object message) {
|
||||
if (userId == null || message == null) return;
|
||||
|
||||
if (!UserChannelManager.isOnline(userId)) {
|
||||
// 用户不在线,直接跳过
|
||||
return;
|
||||
}
|
||||
|
||||
Channel ch = UserChannelManager.getChannel(userId);
|
||||
writeToChannel(ch, message);
|
||||
}
|
||||
@@ -36,22 +43,25 @@ public class WriteManager {
|
||||
}
|
||||
|
||||
/**
|
||||
* 向群内所有在线成员广播消息
|
||||
* 只在异常时记录日志
|
||||
* 广播消息给指定用户列表
|
||||
* 内部会自动过滤离线用户
|
||||
*/
|
||||
public static void writeToGroup(String groupId, Object message) {
|
||||
Set<String> members = GroupManager.getMembers(groupId);
|
||||
for (String userId : members) {
|
||||
try {
|
||||
Channel ch = UserChannelManager.getChannel(userId);
|
||||
if (ch != null && ch.isActive() && message != null) {
|
||||
String json = message instanceof String ? (String) message : MAPPER.writeValueAsString(message);
|
||||
public static void broadcastToUsers(List<String> userIds, Object message) {
|
||||
if (message == null || userIds == null || userIds.isEmpty()) return;
|
||||
|
||||
ch.writeAndFlush(new TextWebSocketFrame(json));
|
||||
}
|
||||
} catch (Exception e) {
|
||||
logger.error("❌ 群消息发送失败 groupId={} userId={}", groupId, userId, e);
|
||||
}
|
||||
for (String userId : userIds) {
|
||||
writeToUser(userId, message); // 内部自动判断在线
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 广播消息给指定 Channel 列表
|
||||
*/
|
||||
public static void broadcastToChannels(List<Channel> channels, Object message) {
|
||||
if (message == null || channels == null || channels.isEmpty()) return;
|
||||
|
||||
for (Channel ch : channels) {
|
||||
writeToChannel(ch, message);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -5,37 +5,45 @@ public class ImResponse {
|
||||
private int code;
|
||||
private String message;
|
||||
private long timestamp;
|
||||
private String from;
|
||||
private Object data;
|
||||
|
||||
// 无参构造
|
||||
public ImResponse() {}
|
||||
public ImResponse() {
|
||||
}
|
||||
|
||||
// 全参构造
|
||||
public ImResponse(String type, int code, String message, long timestamp, Object data) {
|
||||
public ImResponse(String type, int code, String message, long timestamp, String from,Object data) {
|
||||
this.type = type;
|
||||
this.code = code;
|
||||
this.message = message;
|
||||
this.timestamp = timestamp;
|
||||
this.from = from;
|
||||
this.data = data;
|
||||
}
|
||||
|
||||
public static ImResponse success(String type, Object data) {
|
||||
return new ImResponse(type, 200, "success", System.currentTimeMillis(), data);
|
||||
public static ImResponse success() {
|
||||
ImResponse response = new ImResponse();
|
||||
response.setCode(200);
|
||||
response.setMessage("success");
|
||||
response.setTimestamp( System.currentTimeMillis());
|
||||
return response;
|
||||
}
|
||||
|
||||
public static ImResponse success(String type, String message, Object data) {
|
||||
return new ImResponse(type, 200, message, System.currentTimeMillis(), data);
|
||||
public static ImResponse success(String type, String from,Object data) {
|
||||
return new ImResponse(type, 200, "success", System.currentTimeMillis(), from,data);
|
||||
}
|
||||
|
||||
public static ImResponse success(String type, String message,String from,Object data) {
|
||||
return new ImResponse(type, 200, message, System.currentTimeMillis(), from,data);
|
||||
}
|
||||
|
||||
public static ImResponse error(String type, String message) {
|
||||
return new ImResponse(type, 500, message, System.currentTimeMillis(), null);
|
||||
return new ImResponse(type, 500, message, System.currentTimeMillis(), null,null);
|
||||
}
|
||||
|
||||
public static ImResponse error(String type, int code, String message) {
|
||||
return new ImResponse(type, code, message, System.currentTimeMillis(), null);
|
||||
return new ImResponse(type, code, message, System.currentTimeMillis(), null,null);
|
||||
}
|
||||
|
||||
// Getter & Setter
|
||||
public String getType() {
|
||||
return type;
|
||||
}
|
||||
@@ -68,6 +76,14 @@ public class ImResponse {
|
||||
this.timestamp = timestamp;
|
||||
}
|
||||
|
||||
public String getFrom() {
|
||||
return from;
|
||||
}
|
||||
|
||||
public void setFrom(String from) {
|
||||
this.from = from;
|
||||
}
|
||||
|
||||
public Object getData() {
|
||||
return data;
|
||||
}
|
||||
@@ -75,15 +91,4 @@ public class ImResponse {
|
||||
public void setData(Object data) {
|
||||
this.data = data;
|
||||
}
|
||||
|
||||
@Override
|
||||
public String toString() {
|
||||
return "ImResponse{" +
|
||||
"type='" + type + '\'' +
|
||||
", code=" + code +
|
||||
", message='" + message + '\'' +
|
||||
", timestamp=" + timestamp +
|
||||
", data=" + data +
|
||||
'}';
|
||||
}
|
||||
}
|
||||
@@ -1,5 +1,8 @@
|
||||
package com.link.im.utils;
|
||||
|
||||
import java.util.HashMap;
|
||||
import java.util.Map;
|
||||
|
||||
public class CommonUtils {
|
||||
/**
|
||||
* 判断字符串是否为空(null 或 "")
|
||||
@@ -18,4 +21,29 @@ public class CommonUtils {
|
||||
public static boolean isNotEmpty(String str) {
|
||||
return !isEmpty(str);
|
||||
}
|
||||
|
||||
|
||||
/**
|
||||
* 快速创建 Map,类似 JDK9 的 Map.of()
|
||||
* 用法示例:
|
||||
* MapUtils.of("groupId", 1, "userId", 2, "status", "online")
|
||||
*/
|
||||
public static <K, V> Map<K, V> mapOf(Object... keyValues) {
|
||||
if (keyValues == null || keyValues.length == 0) {
|
||||
return new HashMap<>();
|
||||
}
|
||||
if (keyValues.length % 2 != 0) {
|
||||
throw new IllegalArgumentException("参数数量必须为偶数(键值成对出现)");
|
||||
}
|
||||
|
||||
Map<K, V> map = new HashMap<>(keyValues.length / 2);
|
||||
for (int i = 0; i < keyValues.length; i += 2) {
|
||||
@SuppressWarnings("unchecked")
|
||||
K key = (K) keyValues[i];
|
||||
@SuppressWarnings("unchecked")
|
||||
V value = (V) keyValues[i + 1];
|
||||
map.put(key, value);
|
||||
}
|
||||
return map;
|
||||
}
|
||||
}
|
||||
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
Reference in new issue
Block a user