diff --git a/im-das-api/im-das-message-api/pom.xml b/im-das-api/im-das-message-api/pom.xml
index 04bcbff..79d9421 100644
--- a/im-das-api/im-das-message-api/pom.xml
+++ b/im-das-api/im-das-message-api/pom.xml
@@ -31,12 +31,10 @@
org.springframework.boot
spring-boot-starter-amqp
- provided
org.msgpack
jackson-dataformat-msgpack
- provided
org.xerial.snappy
diff --git a/im-entry/src/main/java/net/sopod/soim/entry/handlers/chat/ReqGroupMessageHandler.java b/im-entry/src/main/java/net/sopod/soim/entry/handlers/chat/ReqGroupMessageHandler.java
index be7f8d0..db0f49a 100644
--- a/im-entry/src/main/java/net/sopod/soim/entry/handlers/chat/ReqGroupMessageHandler.java
+++ b/im-entry/src/main/java/net/sopod/soim/entry/handlers/chat/ReqGroupMessageHandler.java
@@ -7,7 +7,7 @@ import net.sopod.soim.data.msg.group.Group;
import net.sopod.soim.entry.server.handler.AccountMessageHandler;
import net.sopod.soim.entry.server.handler.ImContext;
import net.sopod.soim.entry.server.session.Account;
-import net.sopod.soim.logic.api.message.mode.GroupMessage;
+import net.sopod.soim.logic.common.model.message.GroupMessage;
import net.sopod.soim.logic.api.message.service.ImGroupChatService;
import org.apache.dubbo.config.annotation.DubboReference;
import org.slf4j.Logger;
diff --git a/im-entry/src/main/java/net/sopod/soim/entry/handlers/chat/ReqTextChatHandler.java b/im-entry/src/main/java/net/sopod/soim/entry/handlers/chat/ReqTextChatHandler.java
index 7ee1cd0..4ad2094 100644
--- a/im-entry/src/main/java/net/sopod/soim/entry/handlers/chat/ReqTextChatHandler.java
+++ b/im-entry/src/main/java/net/sopod/soim/entry/handlers/chat/ReqTextChatHandler.java
@@ -5,13 +5,16 @@ import net.sopod.soim.entry.server.handler.AccountMessageHandler;
import net.sopod.soim.entry.server.handler.ImContext;
import net.sopod.soim.entry.server.session.Account;
import net.sopod.soim.data.msg.chat.Chat;
-import net.sopod.soim.logic.common.model.TextChat;
-import net.sopod.soim.logic.api.user.service.ChatService;
+import net.sopod.soim.entry.worker.FutureExecutor;
+import net.sopod.soim.logic.api.message.service.ImUserChatService;
+import net.sopod.soim.logic.common.model.message.UserMessage;
import org.apache.dubbo.config.annotation.DubboReference;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.stereotype.Service;
+import java.util.concurrent.CompletableFuture;
+
/**
* ReqTextChatHandler
*
@@ -24,18 +27,24 @@ public class ReqTextChatHandler extends AccountMessageHandler {
private static final Logger logger = LoggerFactory.getLogger(ReqTextChatHandler.class);
@DubboReference
- private ChatService chatService;
+ private ImUserChatService imUserChatService;
@Override
public MessageLite handle(ImContext ctx, Account account, Chat.TextChat msg) {
- TextChat textChat = new TextChat()
- .setUid(msg.getSender())
+ UserMessage textChat = new UserMessage()
+ .setSenderUid(msg.getSender())
.setReceiverUid(msg.getReceiver())
.setReceiverName(msg.getReceiverAccount())
.setTime(msg.getTime())
.setMessage(msg.getMessage());
- Boolean result = chatService.textChat(textChat);
- logger.info("发送结果: {}", result);
+ CompletableFuture stringCompletableFuture = imUserChatService.userMessage(textChat);
+ stringCompletableFuture.whenCompleteAsync((res, e) -> {
+ if (e != null) {
+ logger.error("发送失败:", e);
+ return;
+ }
+ logger.info("发送结果: {}", res);
+ }, FutureExecutor.getInstance());
return null;
}
diff --git a/im-entry/src/main/java/net/sopod/soim/entry/service/TextChatServiceImpl.java b/im-entry/src/main/java/net/sopod/soim/entry/service/TextChatServiceImpl.java
index dbcaeaf..0993a36 100644
--- a/im-entry/src/main/java/net/sopod/soim/entry/service/TextChatServiceImpl.java
+++ b/im-entry/src/main/java/net/sopod/soim/entry/service/TextChatServiceImpl.java
@@ -35,7 +35,7 @@ public class TextChatServiceImpl implements TextChatService {
return Boolean.FALSE;
}
Chat.TextChat resTextChat = Chat.TextChat.newBuilder()
- .setSender(chat.getUid())
+ .setSender(chat.getSenderUid())
.setReceiver(chat.getReceiverUid())
.setReceiverAccount(chat.getReceiverName())
.setMessage(chat.getMessage())
diff --git a/im-service-api/im-logic-common/src/main/java/net/sopod/soim/logic/common/model/TextChat.java b/im-service-api/im-logic-common/src/main/java/net/sopod/soim/logic/common/model/TextChat.java
index 33f419c..536c324 100644
--- a/im-service-api/im-logic-common/src/main/java/net/sopod/soim/logic/common/model/TextChat.java
+++ b/im-service-api/im-logic-common/src/main/java/net/sopod/soim/logic/common/model/TextChat.java
@@ -17,7 +17,7 @@ public class TextChat implements Serializable {
private static final long serialVersionUID = -3994628221046563769L;
- private Long uid;
+ private Long senderUid;
private Long receiverUid;
diff --git a/im-service-api/im-logic-message-api/src/main/java/net/sopod/soim/logic/api/message/mode/GroupMessage.java b/im-service-api/im-logic-common/src/main/java/net/sopod/soim/logic/common/model/message/GroupMessage.java
similarity index 89%
rename from im-service-api/im-logic-message-api/src/main/java/net/sopod/soim/logic/api/message/mode/GroupMessage.java
rename to im-service-api/im-logic-common/src/main/java/net/sopod/soim/logic/common/model/message/GroupMessage.java
index d8edd8d..2a49532 100644
--- a/im-service-api/im-logic-message-api/src/main/java/net/sopod/soim/logic/api/message/mode/GroupMessage.java
+++ b/im-service-api/im-logic-common/src/main/java/net/sopod/soim/logic/common/model/message/GroupMessage.java
@@ -1,4 +1,4 @@
-package net.sopod.soim.logic.api.message.mode;
+package net.sopod.soim.logic.common.model.message;
import lombok.Data;
import lombok.experimental.Accessors;
diff --git a/im-service-api/im-logic-common/src/main/java/net/sopod/soim/logic/common/model/message/UserMessage.java b/im-service-api/im-logic-common/src/main/java/net/sopod/soim/logic/common/model/message/UserMessage.java
new file mode 100644
index 0000000..fd8d269
--- /dev/null
+++ b/im-service-api/im-logic-common/src/main/java/net/sopod/soim/logic/common/model/message/UserMessage.java
@@ -0,0 +1,30 @@
+package net.sopod.soim.logic.common.model.message;
+
+import lombok.Data;
+import lombok.experimental.Accessors;
+
+import java.io.Serializable;
+
+/**
+ * UserMessage
+ *
+ * @author tmy
+ * @date 2022-06-08 16:53
+ */
+@Data
+@Accessors(chain = true)
+public class UserMessage implements Serializable {
+
+ private static final long serialVersionUID = -572323953445204089L;
+
+ private Long senderUid;
+
+ private Long receiverUid;
+
+ private String receiverName;
+
+ private String message;
+
+ private Long time;
+
+}
diff --git a/im-service-api/im-logic-message-api/pom.xml b/im-service-api/im-logic-message-api/pom.xml
index 7249be7..2149ba8 100644
--- a/im-service-api/im-logic-message-api/pom.xml
+++ b/im-service-api/im-logic-message-api/pom.xml
@@ -10,6 +10,13 @@
4.0.0
im-logic-message-api
+
+
+ net.sopod
+ im-logic-common
+ ${soim.version}
+
+
\ No newline at end of file
diff --git a/im-service-api/im-logic-message-api/src/main/java/net/sopod/soim/logic/api/message/service/ImGroupChatService.java b/im-service-api/im-logic-message-api/src/main/java/net/sopod/soim/logic/api/message/service/ImGroupChatService.java
index 4904a34..7273595 100644
--- a/im-service-api/im-logic-message-api/src/main/java/net/sopod/soim/logic/api/message/service/ImGroupChatService.java
+++ b/im-service-api/im-logic-message-api/src/main/java/net/sopod/soim/logic/api/message/service/ImGroupChatService.java
@@ -1,6 +1,6 @@
package net.sopod.soim.logic.api.message.service;
-import net.sopod.soim.logic.api.message.mode.GroupMessage;
+import net.sopod.soim.logic.common.model.message.GroupMessage;
import java.util.concurrent.CompletableFuture;
diff --git a/im-service-api/im-logic-message-api/src/main/java/net/sopod/soim/logic/api/message/service/ImUserChatService.java b/im-service-api/im-logic-message-api/src/main/java/net/sopod/soim/logic/api/message/service/ImUserChatService.java
new file mode 100644
index 0000000..731cb95
--- /dev/null
+++ b/im-service-api/im-logic-message-api/src/main/java/net/sopod/soim/logic/api/message/service/ImUserChatService.java
@@ -0,0 +1,17 @@
+package net.sopod.soim.logic.api.message.service;
+
+import net.sopod.soim.logic.common.model.message.UserMessage;
+
+import java.util.concurrent.CompletableFuture;
+
+/**
+ * ImUserChatService
+ *
+ * @author tmy
+ * @date 2022-06-08 17:19
+ */
+public interface ImUserChatService {
+
+ CompletableFuture userMessage(UserMessage msg);
+
+}
diff --git a/im-service-api/im-logic-user-api/src/main/java/net/sopod/soim/logic/api/user/service/ChatService.java b/im-service-api/im-logic-user-api/src/main/java/net/sopod/soim/logic/api/user/service/ChatService.java
deleted file mode 100644
index e55705d..0000000
--- a/im-service-api/im-logic-user-api/src/main/java/net/sopod/soim/logic/api/user/service/ChatService.java
+++ /dev/null
@@ -1,15 +0,0 @@
-package net.sopod.soim.logic.api.user.service;
-
-import net.sopod.soim.logic.common.model.TextChat;
-
-/**
- * ChatService
- *
- * @author tmy
- * @date 2022-04-28 14:19
- */
-public interface ChatService {
-
- Boolean textChat(TextChat textChat);
-
-}
diff --git a/im-service-api/im-router-api/src/main/java/net/sopod/soim/router/api/service/MessageRouteService.java b/im-service-api/im-router-api/src/main/java/net/sopod/soim/router/api/service/MessageRouteService.java
index f5f71ca..2af00ee 100644
--- a/im-service-api/im-router-api/src/main/java/net/sopod/soim/router/api/service/MessageRouteService.java
+++ b/im-service-api/im-router-api/src/main/java/net/sopod/soim/router/api/service/MessageRouteService.java
@@ -1,5 +1,7 @@
package net.sopod.soim.router.api.service;
+import net.sopod.soim.logic.common.model.message.UserMessage;
+
import java.util.List;
/**
@@ -10,6 +12,11 @@ import java.util.List;
*/
public interface MessageRouteService {
- List routeGroupMessage(List uids, String message);
+ List routeGroupMessage(Long sender, List uids, String message);
+
+ /**
+ * @param textChat 聊天内容
+ */
+ Boolean routeUserMessage(UserMessage textChat);
}
diff --git a/im-service-api/im-router-api/src/main/java/net/sopod/soim/router/api/service/UserRouteService.java b/im-service-api/im-router-api/src/main/java/net/sopod/soim/router/api/service/UserRouteService.java
index 546d6cd..faffbbb 100644
--- a/im-service-api/im-router-api/src/main/java/net/sopod/soim/router/api/service/UserRouteService.java
+++ b/im-service-api/im-router-api/src/main/java/net/sopod/soim/router/api/service/UserRouteService.java
@@ -25,12 +25,6 @@ public interface UserRouteService {
*/
List onlineUserList(String keyword);
- /**
- * @param relationId 好友关系id
- * @param textChat 聊天内容
- */
- Boolean routeTextChat(Long relationId, TextChat textChat);
-
List isOnlineUsers(List userIds);
}
diff --git a/im-service/im-logic-message/pom.xml b/im-service/im-logic-message/pom.xml
index d76abc9..284da44 100644
--- a/im-service/im-logic-message/pom.xml
+++ b/im-service/im-logic-message/pom.xml
@@ -27,6 +27,11 @@
im-das-group-api
${soim.version}
+
+ net.sopod
+ im-das-user-api
+ ${soim.version}
+
net.sopod
im-segment-id-api
@@ -74,18 +79,6 @@
org.apache.dubbo
dubbo-registry-nacos
-
- org.springframework.boot
- spring-boot-starter-amqp
-
-
- org.msgpack
- jackson-dataformat-msgpack
-
-
- org.mybatis
- mybatis
-
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 84b99f9..4a67141 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
@@ -8,9 +8,9 @@ import net.sopod.soim.das.group.api.model.dto.GroupUser_0;
import net.sopod.soim.das.group.api.service.DasGroupUserService;
import net.sopod.soim.das.message.api.entity.ImGroupMessage;
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.model.message.GroupMessage;
import net.sopod.soim.logic.common.util.RpcContextUtil;
import net.sopod.soim.router.api.route.UidConsistentHashSelector;
import net.sopod.soim.router.api.service.MessageRouteService;
@@ -69,7 +69,7 @@ public class ImGroupChatServiceImpl implements ImGroupChatService {
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());
+ List results = messageRouteService.routeGroupMessage(msg.getSender(), uidList, msg.getMessage());
// TODO 更新未读消息数据偏移量,未读消息条数
return CompletableFuture.completedFuture("OK");
}
@@ -81,7 +81,7 @@ public class ImGroupChatServiceImpl implements ImGroupChatService {
// 批量路由消息到对应的 router
for (List uidGroup : uidGroups.values()) {
RpcContextUtil.setContextUid(uidGroup.get(0));
- List results = messageRouteService.routeGroupMessage(uidGroup, msg.getMessage());
+ List results = messageRouteService.routeGroupMessage(msg.getSender(), uidGroup, msg.getMessage());
Iterator iterator = uidGroup.iterator();
for (int i = 0; iterator.hasNext(); i++) {
GroupUser_0 gUser = groupUserMap.get(iterator.next());
diff --git a/im-service/im-logic-message/src/main/java/net/sopod/soim/logic/message/service/ImUserChatServiceImpl.java b/im-service/im-logic-message/src/main/java/net/sopod/soim/logic/message/service/ImUserChatServiceImpl.java
new file mode 100644
index 0000000..98ae419
--- /dev/null
+++ b/im-service/im-logic-message/src/main/java/net/sopod/soim/logic/message/service/ImUserChatServiceImpl.java
@@ -0,0 +1,93 @@
+package net.sopod.soim.logic.message.service;
+
+import net.sopod.soim.common.dubbo.exception.LogicException;
+import net.sopod.soim.common.dubbo.exception.SoimException;
+import net.sopod.soim.common.util.ImClock;
+import net.sopod.soim.das.common.config.LogicTables;
+import net.sopod.soim.das.message.api.entity.ImMessage;
+import net.sopod.soim.das.message.api.service.DasMQMessagePersistentService;
+import net.sopod.soim.das.user.api.model.entity.ImUser;
+import net.sopod.soim.das.user.api.service.FriendDas;
+import net.sopod.soim.das.user.api.service.UserDas;
+import net.sopod.soim.logic.api.message.service.ImUserChatService;
+import net.sopod.soim.logic.api.segmentid.core.SegmentIdGenerator;
+import net.sopod.soim.logic.common.model.message.UserMessage;
+import net.sopod.soim.logic.common.util.RpcContextUtil;
+import net.sopod.soim.router.api.service.MessageRouteService;
+import net.sopod.soim.router.api.service.UserRouteService;
+import org.apache.dubbo.config.annotation.DubboReference;
+import org.apache.dubbo.config.annotation.DubboService;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import javax.annotation.Resource;
+import java.util.Objects;
+import java.util.concurrent.CompletableFuture;
+
+/**
+ * ImUserChatServiceImpl
+ * 单聊消息
+ *
+ * @author tmy
+ * @date 2022-06-08 14:59
+ */
+@DubboService
+public class ImUserChatServiceImpl implements ImUserChatService {
+
+ private static final Logger logger = LoggerFactory.getLogger(ImUserChatServiceImpl.class);
+
+ @DubboReference
+ private UserRouteService userRouteService;
+
+ @DubboReference
+ private MessageRouteService messageRouteService;
+
+ @DubboReference
+ private UserDas userDas;
+
+ @DubboReference
+ private FriendDas friendDas;
+
+ @Resource
+ private SegmentIdGenerator segmentIdGenerator;
+
+ @Resource
+ private DasMQMessagePersistentService dasMQMessagePersistentService;
+
+ @Override
+ public CompletableFuture userMessage(UserMessage msg) {
+ if (msg.getReceiverUid() == null
+ || Objects.equals(msg.getReceiverUid(), 0L)) {
+ ImUser receiverUser = userDas.getNormalUserByAccount(msg.getReceiverName());
+ if (receiverUser == null) {
+ // TODO receiverUid 由客户端传递
+ throw new LogicException("聊天对象不存在");
+ }
+ msg.setReceiverUid(receiverUser.getId());
+ }
+ // 查询好友关系
+ Long relationId = friendDas.getRelationId(msg.getSenderUid(), msg.getReceiverUid());
+ // 不是好友
+ if (relationId == null) {
+ throw new LogicException("请添加对方好友后发送消息");
+ }
+ // 通过消息队列持久化存储到db
+ ImMessage imMessage = new ImMessage()
+ .setRelationId(relationId)
+ .setId(segmentIdGenerator.nextId(LogicTables.IM_MESSAGE))
+ .setContent(msg.getMessage())
+ .setSender(msg.getSenderUid())
+ .setReceiver(msg.getReceiverUid())
+ .setCreateTime(ImClock.date());
+ dasMQMessagePersistentService.saveImMessage(imMessage);
+ // 设置调用 router 为消息接受者地址
+ RpcContextUtil.setContextUid(msg.getReceiverUid());
+ try {
+ Boolean success = messageRouteService.routeUserMessage(msg);
+ return CompletableFuture.completedFuture(Boolean.TRUE.equals(success) ? null : "消息未接受");
+ } catch (SoimException e) {
+ return CompletableFuture.failedFuture(e);
+ }
+ }
+
+}
diff --git a/im-service/im-logic-user/src/main/java/net/sopod/soim/logic/user/service/ChatServiceImpl.java b/im-service/im-logic-user/src/main/java/net/sopod/soim/logic/user/service/ChatServiceImpl.java
deleted file mode 100644
index 29eb87a..0000000
--- a/im-service/im-logic-user/src/main/java/net/sopod/soim/logic/user/service/ChatServiceImpl.java
+++ /dev/null
@@ -1,60 +0,0 @@
-package net.sopod.soim.logic.user.service;
-
-import net.sopod.soim.common.dubbo.exception.LogicException;
-import net.sopod.soim.das.user.api.model.entity.ImUser;
-import net.sopod.soim.das.user.api.service.FriendDas;
-import net.sopod.soim.das.user.api.service.UserDas;
-import net.sopod.soim.logic.api.user.service.ChatService;
-import net.sopod.soim.logic.common.model.TextChat;
-import net.sopod.soim.logic.common.util.RpcContextUtil;
-import net.sopod.soim.router.api.service.UserRouteService;
-import org.apache.dubbo.config.annotation.DubboReference;
-import org.apache.dubbo.config.annotation.DubboService;
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
-
-import java.util.Objects;
-
-/**
- * ChatServiceImpl
- *
- * @author tmy
- * @date 2022-04-28 14:22
- */
-@DubboService
-public class ChatServiceImpl implements ChatService {
-
- private static final Logger logger = LoggerFactory.getLogger(ChatServiceImpl.class);
-
- @DubboReference
- private UserRouteService userRouteService;
-
- @DubboReference
- private UserDas userDas;
-
- @DubboReference
- private FriendDas friendDas;
-
- @Override
- public Boolean textChat(TextChat textChat) {
- if (textChat.getReceiverUid() == null
- || Objects.equals(textChat.getReceiverUid(), 0L)) {
- ImUser receiverUser = userDas.getNormalUserByAccount(textChat.getReceiverName());
- if (receiverUser == null) {
- // TODO receiverUid 由客户端传递
- throw new LogicException("聊天对象不存在");
- }
- textChat.setReceiverUid(receiverUser.getId());
- }
- // 查询好友关系
- Long relationId = friendDas.getRelationId(textChat.getUid(), textChat.getReceiverUid());
- // 不是好友
- if (relationId == null) {
- throw new LogicException("请添加好友后发送消息");
- }
- // 设置调用 router 为消息接受者地址
- RpcContextUtil.setContextUid(textChat.getReceiverUid());
- return userRouteService.routeTextChat(relationId, textChat);
- }
-
-}
diff --git a/im-service/im-router/pom.xml b/im-service/im-router/pom.xml
index 42f76aa..342d88c 100644
--- a/im-service/im-router/pom.xml
+++ b/im-service/im-router/pom.xml
@@ -32,11 +32,6 @@
im-das-user-api
${soim.version}
-
- net.sopod
- im-das-message-api
- ${soim.version}
-
org.springframework.boot
spring-boot-starter
@@ -90,22 +85,22 @@
cglib
cglib
-
- org.msgpack
- jackson-dataformat-msgpack
-
-
- org.xerial.snappy
- snappy-java
-
-
- org.springframework.boot
- spring-boot-starter-amqp
-
-
- org.mybatis
- mybatis
-
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
\ No newline at end of file
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 377125e..d3c7cce 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
@@ -2,13 +2,15 @@ package net.sopod.soim.router.service;
import lombok.AllArgsConstructor;
import lombok.extern.slf4j.Slf4j;
+import net.sopod.soim.logic.common.model.message.UserMessage;
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 org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
import java.util.ArrayList;
-import java.util.Iterator;
import java.util.List;
/**
@@ -23,10 +25,12 @@ import java.util.List;
@Slf4j
public class MessageRouteServiceImpl implements MessageRouteService {
+ private static final Logger logger = LoggerFactory.getLogger(MessageRouteServiceImpl.class);
+
private RouterUserService routerUserService;
@Override
- public List routeGroupMessage(List uids, String message) {
+ public List routeGroupMessage(Long sender, List uids, String groupMessage) {
RouterUserStorage storage = RouterUserStorage.getInstance();
List results = new ArrayList<>(uids.size());
for (Long uid : uids) {
@@ -36,10 +40,25 @@ public class MessageRouteServiceImpl implements MessageRouteService {
results.add(false);
continue;
}
- Boolean success = routerUserService.routeGroupMessage(routerUser, message);
+ Boolean success = routerUserService.routeGroupMessage(sender, routerUser, groupMessage);
results.add(Boolean.TRUE.equals(success));
}
return results;
}
+ /**
+ * 调用该方法时,将到 im-router 服务的路由 uid 设置为消息接受者的 uid
+ */
+ @Override
+ public Boolean routeUserMessage(UserMessage userMessage) {
+ RouterUserStorage storage = RouterUserStorage.getInstance();
+ RouterUser receiverUser = storage.getRouterUserMap().get(userMessage.getReceiverUid());
+ Boolean send = routerUserService.routeChatMessage(userMessage.getSenderUid(), receiverUser, userMessage.getMessage());
+ if (!Boolean.TRUE.equals(send)) {
+ // 消息未送达
+ logger.info("消息未送达: {}: {}", userMessage.getReceiverName(), userMessage.getMessage());
+ }
+ return Boolean.TRUE.equals(send);
+ }
+
}
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
index 58b2aec..fde577d 100644
--- 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
@@ -20,18 +20,26 @@ public class RouterUserService {
@DubboReference
private TextChatService textChatService;
- public void routeChatMessage(RouterUser routerUser, String message) {
-
+ public Boolean routeChatMessage(Long sender, RouterUser receiverUser, String message) {
+ RpcContextUtil.setContextUid(receiverUser.getUid());
+ TextChat textChat = new TextChat()
+ .setSenderUid(sender)
+ .setReceiverUid(receiverUser.getUid())
+ .setMessage(message)
+ .setTime(ImClock.millis())
+ .setReceiverName(receiverUser.getAccount());
+ return textChatService.sendTextChat(textChat);
}
- public Boolean routeGroupMessage(RouterUser receiverUser, String message) {
+ public Boolean routeGroupMessage(Long sender, RouterUser receiverUser, String message) {
TextChat textChat = new TextChat()
- .setUid(receiverUser.getUid()) // TODO 群消息优化
+ .setSenderUid(sender)
.setReceiverUid(receiverUser.getUid())
.setMessage(message)
.setTime(ImClock.millis())
.setReceiverName(receiverUser.getAccount());
RpcContextUtil.setContextUid(receiverUser.getUid());
+ // TODO 单独群消息接口
return textChatService.sendTextChat(textChat);
}
diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/service/UserRouteServiceImpl.java b/im-service/im-router/src/main/java/net/sopod/soim/router/service/UserRouteServiceImpl.java
index 25f7f69..f0c5115 100644
--- a/im-service/im-router/src/main/java/net/sopod/soim/router/service/UserRouteServiceImpl.java
+++ b/im-service/im-router/src/main/java/net/sopod/soim/router/service/UserRouteServiceImpl.java
@@ -91,29 +91,6 @@ public class UserRouteServiceImpl implements UserRouteService {
.collect(Collectors.toList());
}
- /**
- * 调用该方法时,将到 im-router 服务的路由 uid 设置为消息接受者的 uid
- */
- @Override
- public Boolean routeTextChat(Long relationId, TextChat textChat) {
- // 通过消息队列持久化存储到db
- ImMessage imMessage = new ImMessage()
- .setRelationId(relationId)
- .setId(segmentIdGenerator.nextId(LogicTables.IM_MESSAGE))
- .setContent(textChat.getMessage())
- .setSender(textChat.getUid())
- .setReceiver(textChat.getReceiverUid())
- .setCreateTime(ImClock.date());
- dasMQMessagePersistentService.saveImMessage(imMessage);
-
- RpcContextUtil.setContextUid(textChat.getReceiverUid());
- Boolean send = textChatService.sendTextChat(textChat);
- if (!Boolean.TRUE.equals(send)) {
- // TODO 未送到,消息存储,重发...
- logger.info("消息未送达: {}: {}", textChat.getReceiverName(), textChat.getMessage());
- }
- return Boolean.TRUE;
- }
@Override
public List isOnlineUsers(List userIds) {