diff --git a/im-service/im-logic-message/pom.xml b/im-service/im-logic-message/pom.xml
index dbd8323..bfc863b 100644
--- a/im-service/im-logic-message/pom.xml
+++ b/im-service/im-logic-message/pom.xml
@@ -74,6 +74,10 @@
org.apache.dubbo
dubbo-registry-nacos
+
+ org.springframework.boot
+ spring-boot-starter-amqp
+
org.springframework.boot
diff --git a/im-service/im-logic-message/src/main/java/net/sopod/soim/logic/message/service/ImGroupChatServiceImpl.java b/im-service/im-logic-message/src/main/java/net/sopod/soim/logic/message/service/ImGroupChatServiceImpl.java
index c7408a4..3b504bc 100644
--- a/im-service/im-logic-message/src/main/java/net/sopod/soim/logic/message/service/ImGroupChatServiceImpl.java
+++ b/im-service/im-logic-message/src/main/java/net/sopod/soim/logic/message/service/ImGroupChatServiceImpl.java
@@ -1,6 +1,8 @@
package net.sopod.soim.logic.message.service;
+import net.sopod.soim.common.util.Collects;
import net.sopod.soim.common.util.ImClock;
+import net.sopod.soim.common.util.StringUtil;
import net.sopod.soim.das.common.config.LogicTables;
import net.sopod.soim.das.group.api.model.dto.GroupUser_0;
import net.sopod.soim.das.group.api.service.DasGroupUserService;
@@ -9,13 +11,19 @@ import net.sopod.soim.das.message.api.service.DasMQMessagePersistentService;
import net.sopod.soim.logic.api.message.mode.GroupMessage;
import net.sopod.soim.logic.api.message.service.ImGroupChatService;
import net.sopod.soim.logic.api.segmentid.core.SegmentIdGenerator;
+import net.sopod.soim.logic.common.util.RpcContextUtil;
import net.sopod.soim.router.api.route.UidConsistentHashSelector;
+import net.sopod.soim.router.api.service.MessageRouteService;
import org.apache.dubbo.config.annotation.DubboReference;
import org.apache.dubbo.config.annotation.DubboService;
import javax.annotation.Resource;
+import java.util.Iterator;
import java.util.List;
+import java.util.Map;
+import java.util.Set;
import java.util.concurrent.CompletableFuture;
+import java.util.stream.Collectors;
/**
* ImGroupChatServiceImpl
@@ -35,6 +43,9 @@ public class ImGroupChatServiceImpl implements ImGroupChatService {
@DubboReference
private DasGroupUserService dasGroupUserService;
+ @DubboReference
+ private MessageRouteService messageRouteService;
+
/**
* 序列化群消息
* 查询群成员列表
@@ -54,11 +65,33 @@ public class ImGroupChatServiceImpl implements ImGroupChatService {
dasMQMessagePersistentService.saveImGroupMessage(imGroupMessage);
// 查询群成员列表
List groupUsers = dasGroupUserService.listGroupUsers(msg.getGid());
+ Map groupUserMap = Collects.collect2Map(groupUsers, GroupUser_0::getUid);
// 分批路由消息
UidConsistentHashSelector> uidSelector = UidConsistentHashSelector.getCurrent();
+ if (uidSelector == null) {
+ List uidList = groupUsers.stream().map(GroupUser_0::getUid).collect(Collectors.toList());
+ List results = messageRouteService.routeGroupMessage(uidList, msg.getMessage());
+ // TODO 更新未读消息数据偏移量,未读消息条数
+ return CompletableFuture.completedFuture("OK");
+ }
+ Map, List> uidGroups = Collects.group(
+ groupUsers,
+ user -> uidSelector.select(StringUtil.toString(user.getUid())),
+ GroupUser_0::getUid
+ );
+ // 批量路由消息到对应的 router
+ for (List uidGroup : uidGroups.values()) {
+ RpcContextUtil.setContextUid(uidGroup.get(0));
+ List results = messageRouteService.routeGroupMessage(uidGroup, msg.getMessage());
+ Iterator iterator = uidGroup.iterator();
+ for (int i = 0; iterator.hasNext(); i++) {
+ GroupUser_0 gUser = groupUserMap.get(iterator.next());
+ Boolean received = results.get(i);
+ // TODO 更新未读消息数据偏移量,未读消息条数
-
- return null;
+ }
+ }
+ return CompletableFuture.completedFuture("OK");
}
}
diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/service/MessageRouteServiceImpl.java b/im-service/im-router/src/main/java/net/sopod/soim/router/service/MessageRouteServiceImpl.java
index d6db1f0..18bb0e7 100644
--- a/im-service/im-router/src/main/java/net/sopod/soim/router/service/MessageRouteServiceImpl.java
+++ b/im-service/im-router/src/main/java/net/sopod/soim/router/service/MessageRouteServiceImpl.java
@@ -1,30 +1,43 @@
package net.sopod.soim.router.service;
+import lombok.AllArgsConstructor;
import net.sopod.soim.router.api.service.MessageRouteService;
import net.sopod.soim.router.cache.RouterUser;
import net.sopod.soim.router.cache.RouterUserStorage;
import org.apache.dubbo.config.annotation.DubboService;
+import java.util.ArrayList;
+import java.util.Iterator;
import java.util.List;
/**
* MessageRouterServiceImpl
+ * 路由聊天消息
*
* @author tmy
* @date 2022-06-05 11:30
*/
@DubboService
+@AllArgsConstructor
public class MessageRouteServiceImpl implements MessageRouteService {
+ private RouterUserService routerUserService;
+
@Override
public List routeGroupMessage(List uids, String message) {
- // TODO
RouterUserStorage storage = RouterUserStorage.getInstance();
- for (Long uid : uids) {
- RouterUser routerUser = storage.get(uid);
-
+ List results = new ArrayList<>(uids.size());
+ Iterator iterator = uids.iterator();
+ for (int i = 0; iterator.hasNext(); i++) {
+ RouterUser routerUser = storage.get(iterator.next());
+ if (routerUser == null) {
+ results.set(i, false);
+ continue;
+ }
+ Boolean success = routerUserService.routeGroupMessage(routerUser, message);
+ results.set(i, Boolean.TRUE.equals(success));
}
- return null;
+ return results;
}
}
diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/service/RouterUserService.java b/im-service/im-router/src/main/java/net/sopod/soim/router/service/RouterUserService.java
new file mode 100644
index 0000000..d8fd95d
--- /dev/null
+++ b/im-service/im-router/src/main/java/net/sopod/soim/router/service/RouterUserService.java
@@ -0,0 +1,37 @@
+package net.sopod.soim.router.service;
+
+import net.sopod.soim.common.util.ImClock;
+import net.sopod.soim.entry.api.service.TextChatService;
+import net.sopod.soim.logic.common.model.TextChat;
+import net.sopod.soim.logic.common.util.RpcContextUtil;
+import net.sopod.soim.router.cache.RouterUser;
+import org.apache.dubbo.config.annotation.DubboReference;
+import org.springframework.stereotype.Service;
+
+/**
+ * RouterUserService
+ *
+ * @author tmy
+ * @date 2022-06-06 9:34
+ */
+@Service
+public class RouterUserService {
+
+ @DubboReference
+ private TextChatService textChatService;
+
+ public void routeChatMessage(RouterUser routerUser, String message) {
+
+ }
+
+ public Boolean routeGroupMessage(RouterUser receiverUser, String message) {
+ RpcContextUtil.setContextUid(receiverUser.getUid());
+ TextChat textChat = new TextChat()
+ .setReceiverUid(receiverUser.getUid())
+ .setMessage(message)
+ .setTime(ImClock.millis())
+ .setReceiverName(receiverUser.getAccount());
+ return textChatService.sendTextChat(textChat);
+ }
+
+}
diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/util/RpcContextUtil.java b/im-service/im-router/src/main/java/net/sopod/soim/router/util/RpcContextUtil.java
index c13fa08..e3a0c09 100644
--- a/im-service/im-router/src/main/java/net/sopod/soim/router/util/RpcContextUtil.java
+++ b/im-service/im-router/src/main/java/net/sopod/soim/router/util/RpcContextUtil.java
@@ -13,6 +13,7 @@ import java.util.Objects;
* @author tmy
* @date 2022-05-02 22:27
*/
+@Deprecated
public class RpcContextUtil {
public static boolean setImEntryRouteServerAddrByUid(Long uid) {