Browse Source

单聊消息服务接口转移模块

master
tangmingyou 4 years ago
parent
commit
a8eb491825
  1. 2
      im-das-api/im-das-message-api/pom.xml
  2. 2
      im-entry/src/main/java/net/sopod/soim/entry/handlers/chat/ReqGroupMessageHandler.java
  3. 23
      im-entry/src/main/java/net/sopod/soim/entry/handlers/chat/ReqTextChatHandler.java
  4. 2
      im-entry/src/main/java/net/sopod/soim/entry/service/TextChatServiceImpl.java
  5. 2
      im-service-api/im-logic-common/src/main/java/net/sopod/soim/logic/common/model/TextChat.java
  6. 2
      im-service-api/im-logic-common/src/main/java/net/sopod/soim/logic/common/model/message/GroupMessage.java
  7. 30
      im-service-api/im-logic-common/src/main/java/net/sopod/soim/logic/common/model/message/UserMessage.java
  8. 7
      im-service-api/im-logic-message-api/pom.xml
  9. 2
      im-service-api/im-logic-message-api/src/main/java/net/sopod/soim/logic/api/message/service/ImGroupChatService.java
  10. 17
      im-service-api/im-logic-message-api/src/main/java/net/sopod/soim/logic/api/message/service/ImUserChatService.java
  11. 15
      im-service-api/im-logic-user-api/src/main/java/net/sopod/soim/logic/api/user/service/ChatService.java
  12. 9
      im-service-api/im-router-api/src/main/java/net/sopod/soim/router/api/service/MessageRouteService.java
  13. 6
      im-service-api/im-router-api/src/main/java/net/sopod/soim/router/api/service/UserRouteService.java
  14. 17
      im-service/im-logic-message/pom.xml
  15. 6
      im-service/im-logic-message/src/main/java/net/sopod/soim/logic/message/service/ImGroupChatServiceImpl.java
  16. 93
      im-service/im-logic-message/src/main/java/net/sopod/soim/logic/message/service/ImUserChatServiceImpl.java
  17. 60
      im-service/im-logic-user/src/main/java/net/sopod/soim/logic/user/service/ChatServiceImpl.java
  18. 37
      im-service/im-router/pom.xml
  19. 25
      im-service/im-router/src/main/java/net/sopod/soim/router/service/MessageRouteServiceImpl.java
  20. 16
      im-service/im-router/src/main/java/net/sopod/soim/router/service/RouterUserService.java
  21. 23
      im-service/im-router/src/main/java/net/sopod/soim/router/service/UserRouteServiceImpl.java

2
im-das-api/im-das-message-api/pom.xml

@ -31,12 +31,10 @@
<dependency> <dependency>
<groupId>org.springframework.boot</groupId> <groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-amqp</artifactId> <artifactId>spring-boot-starter-amqp</artifactId>
<scope>provided</scope>
</dependency> </dependency>
<dependency> <dependency>
<groupId>org.msgpack</groupId> <groupId>org.msgpack</groupId>
<artifactId>jackson-dataformat-msgpack</artifactId> <artifactId>jackson-dataformat-msgpack</artifactId>
<scope>provided</scope>
</dependency> </dependency>
<dependency> <dependency>
<groupId>org.xerial.snappy</groupId> <groupId>org.xerial.snappy</groupId>

2
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.AccountMessageHandler;
import net.sopod.soim.entry.server.handler.ImContext; import net.sopod.soim.entry.server.handler.ImContext;
import net.sopod.soim.entry.server.session.Account; 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 net.sopod.soim.logic.api.message.service.ImGroupChatService;
import org.apache.dubbo.config.annotation.DubboReference; import org.apache.dubbo.config.annotation.DubboReference;
import org.slf4j.Logger; import org.slf4j.Logger;

23
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.handler.ImContext;
import net.sopod.soim.entry.server.session.Account; import net.sopod.soim.entry.server.session.Account;
import net.sopod.soim.data.msg.chat.Chat; import net.sopod.soim.data.msg.chat.Chat;
import net.sopod.soim.logic.common.model.TextChat; import net.sopod.soim.entry.worker.FutureExecutor;
import net.sopod.soim.logic.api.user.service.ChatService; 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.apache.dubbo.config.annotation.DubboReference;
import org.slf4j.Logger; import org.slf4j.Logger;
import org.slf4j.LoggerFactory; import org.slf4j.LoggerFactory;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
import java.util.concurrent.CompletableFuture;
/** /**
* ReqTextChatHandler * ReqTextChatHandler
* *
@ -24,18 +27,24 @@ public class ReqTextChatHandler extends AccountMessageHandler<Chat.TextChat> {
private static final Logger logger = LoggerFactory.getLogger(ReqTextChatHandler.class); private static final Logger logger = LoggerFactory.getLogger(ReqTextChatHandler.class);
@DubboReference @DubboReference
private ChatService chatService; private ImUserChatService imUserChatService;
@Override @Override
public MessageLite handle(ImContext ctx, Account account, Chat.TextChat msg) { public MessageLite handle(ImContext ctx, Account account, Chat.TextChat msg) {
TextChat textChat = new TextChat() UserMessage textChat = new UserMessage()
.setUid(msg.getSender()) .setSenderUid(msg.getSender())
.setReceiverUid(msg.getReceiver()) .setReceiverUid(msg.getReceiver())
.setReceiverName(msg.getReceiverAccount()) .setReceiverName(msg.getReceiverAccount())
.setTime(msg.getTime()) .setTime(msg.getTime())
.setMessage(msg.getMessage()); .setMessage(msg.getMessage());
Boolean result = chatService.textChat(textChat); CompletableFuture<String> stringCompletableFuture = imUserChatService.userMessage(textChat);
logger.info("发送结果: {}", result); stringCompletableFuture.whenCompleteAsync((res, e) -> {
if (e != null) {
logger.error("发送失败:", e);
return;
}
logger.info("发送结果: {}", res);
}, FutureExecutor.getInstance());
return null; return null;
} }

2
im-entry/src/main/java/net/sopod/soim/entry/service/TextChatServiceImpl.java

@ -35,7 +35,7 @@ public class TextChatServiceImpl implements TextChatService {
return Boolean.FALSE; return Boolean.FALSE;
} }
Chat.TextChat resTextChat = Chat.TextChat.newBuilder() Chat.TextChat resTextChat = Chat.TextChat.newBuilder()
.setSender(chat.getUid()) .setSender(chat.getSenderUid())
.setReceiver(chat.getReceiverUid()) .setReceiver(chat.getReceiverUid())
.setReceiverAccount(chat.getReceiverName()) .setReceiverAccount(chat.getReceiverName())
.setMessage(chat.getMessage()) .setMessage(chat.getMessage())

2
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 static final long serialVersionUID = -3994628221046563769L;
private Long uid; private Long senderUid;
private Long receiverUid; private Long receiverUid;

2
im-service-api/im-logic-message-api/src/main/java/net/sopod/soim/logic/api/message/mode/GroupMessage.java → 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.Data;
import lombok.experimental.Accessors; import lombok.experimental.Accessors;

30
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;
}

7
im-service-api/im-logic-message-api/pom.xml

@ -10,6 +10,13 @@
<modelVersion>4.0.0</modelVersion> <modelVersion>4.0.0</modelVersion>
<artifactId>im-logic-message-api</artifactId> <artifactId>im-logic-message-api</artifactId>
<dependencies>
<dependency>
<groupId>net.sopod</groupId>
<artifactId>im-logic-common</artifactId>
<version>${soim.version}</version>
</dependency>
</dependencies>
</project> </project>

2
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; 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; import java.util.concurrent.CompletableFuture;

17
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<String> userMessage(UserMessage msg);
}

15
im-service-api/im-logic-user-api/src/main/java/net/sopod/soim/logic/api/user/service/ChatService.java

@ -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);
}

9
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; package net.sopod.soim.router.api.service;
import net.sopod.soim.logic.common.model.message.UserMessage;
import java.util.List; import java.util.List;
/** /**
@ -10,6 +12,11 @@ import java.util.List;
*/ */
public interface MessageRouteService { public interface MessageRouteService {
List<Boolean> routeGroupMessage(List<Long> uids, String message); List<Boolean> routeGroupMessage(Long sender, List<Long> uids, String message);
/**
* @param textChat 聊天内容
*/
Boolean routeUserMessage(UserMessage textChat);
} }

6
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<UserInfo> onlineUserList(String keyword); List<UserInfo> onlineUserList(String keyword);
/**
* @param relationId 好友关系id
* @param textChat 聊天内容
*/
Boolean routeTextChat(Long relationId, TextChat textChat);
List<Boolean> isOnlineUsers(List<Long> userIds); List<Boolean> isOnlineUsers(List<Long> userIds);
} }

17
im-service/im-logic-message/pom.xml

@ -27,6 +27,11 @@
<artifactId>im-das-group-api</artifactId> <artifactId>im-das-group-api</artifactId>
<version>${soim.version}</version> <version>${soim.version}</version>
</dependency> </dependency>
<dependency>
<groupId>net.sopod</groupId>
<artifactId>im-das-user-api</artifactId>
<version>${soim.version}</version>
</dependency>
<dependency> <dependency>
<groupId>net.sopod</groupId> <groupId>net.sopod</groupId>
<artifactId>im-segment-id-api</artifactId> <artifactId>im-segment-id-api</artifactId>
@ -74,18 +79,6 @@
<groupId>org.apache.dubbo</groupId> <groupId>org.apache.dubbo</groupId>
<artifactId>dubbo-registry-nacos</artifactId> <artifactId>dubbo-registry-nacos</artifactId>
</dependency> </dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-amqp</artifactId>
</dependency>
<dependency>
<groupId>org.msgpack</groupId>
<artifactId>jackson-dataformat-msgpack</artifactId>
</dependency>
<dependency>
<groupId>org.mybatis</groupId>
<artifactId>mybatis</artifactId>
</dependency>
<dependency> <dependency>
<groupId>org.springframework.boot</groupId> <groupId>org.springframework.boot</groupId>

6
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.group.api.service.DasGroupUserService;
import net.sopod.soim.das.message.api.entity.ImGroupMessage; import net.sopod.soim.das.message.api.entity.ImGroupMessage;
import net.sopod.soim.das.message.api.service.DasMQMessagePersistentService; 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.message.service.ImGroupChatService;
import net.sopod.soim.logic.api.segmentid.core.SegmentIdGenerator; 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.logic.common.util.RpcContextUtil;
import net.sopod.soim.router.api.route.UidConsistentHashSelector; import net.sopod.soim.router.api.route.UidConsistentHashSelector;
import net.sopod.soim.router.api.service.MessageRouteService; import net.sopod.soim.router.api.service.MessageRouteService;
@ -69,7 +69,7 @@ public class ImGroupChatServiceImpl implements ImGroupChatService {
UidConsistentHashSelector<?> uidSelector = UidConsistentHashSelector.getCurrent(); UidConsistentHashSelector<?> uidSelector = UidConsistentHashSelector.getCurrent();
if (uidSelector == null) { if (uidSelector == null) {
List<Long> uidList = groupUsers.stream().map(GroupUser_0::getUid).collect(Collectors.toList()); List<Long> uidList = groupUsers.stream().map(GroupUser_0::getUid).collect(Collectors.toList());
List<Boolean> results = messageRouteService.routeGroupMessage(uidList, msg.getMessage()); List<Boolean> results = messageRouteService.routeGroupMessage(msg.getSender(), uidList, msg.getMessage());
// TODO 更新未读消息数据偏移量,未读消息条数 // TODO 更新未读消息数据偏移量,未读消息条数
return CompletableFuture.completedFuture("OK"); return CompletableFuture.completedFuture("OK");
} }
@ -81,7 +81,7 @@ public class ImGroupChatServiceImpl implements ImGroupChatService {
// 批量路由消息到对应的 router // 批量路由消息到对应的 router
for (List<Long> uidGroup : uidGroups.values()) { for (List<Long> uidGroup : uidGroups.values()) {
RpcContextUtil.setContextUid(uidGroup.get(0)); RpcContextUtil.setContextUid(uidGroup.get(0));
List<Boolean> results = messageRouteService.routeGroupMessage(uidGroup, msg.getMessage()); List<Boolean> results = messageRouteService.routeGroupMessage(msg.getSender(), uidGroup, msg.getMessage());
Iterator<Long> iterator = uidGroup.iterator(); Iterator<Long> iterator = uidGroup.iterator();
for (int i = 0; iterator.hasNext(); i++) { for (int i = 0; iterator.hasNext(); i++) {
GroupUser_0 gUser = groupUserMap.get(iterator.next()); GroupUser_0 gUser = groupUserMap.get(iterator.next());

93
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<String> 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);
}
}
}

60
im-service/im-logic-user/src/main/java/net/sopod/soim/logic/user/service/ChatServiceImpl.java

@ -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);
}
}

37
im-service/im-router/pom.xml

@ -32,11 +32,6 @@
<artifactId>im-das-user-api</artifactId> <artifactId>im-das-user-api</artifactId>
<version>${soim.version}</version> <version>${soim.version}</version>
</dependency> </dependency>
<dependency>
<groupId>net.sopod</groupId>
<artifactId>im-das-message-api</artifactId>
<version>${soim.version}</version>
</dependency>
<dependency> <dependency>
<groupId>org.springframework.boot</groupId> <groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter</artifactId> <artifactId>spring-boot-starter</artifactId>
@ -90,22 +85,22 @@
<groupId>cglib</groupId> <groupId>cglib</groupId>
<artifactId>cglib</artifactId> <artifactId>cglib</artifactId>
</dependency> </dependency>
<dependency> <!-- <dependency>-->
<groupId>org.msgpack</groupId> <!-- <groupId>org.msgpack</groupId>-->
<artifactId>jackson-dataformat-msgpack</artifactId> <!-- <artifactId>jackson-dataformat-msgpack</artifactId>-->
</dependency> <!-- </dependency>-->
<dependency> <!-- <dependency>-->
<groupId>org.xerial.snappy</groupId> <!-- <groupId>org.xerial.snappy</groupId>-->
<artifactId>snappy-java</artifactId> <!-- <artifactId>snappy-java</artifactId>-->
</dependency> <!-- </dependency>-->
<dependency> <!-- <dependency>-->
<groupId>org.springframework.boot</groupId> <!-- <groupId>org.springframework.boot</groupId>-->
<artifactId>spring-boot-starter-amqp</artifactId> <!-- <artifactId>spring-boot-starter-amqp</artifactId>-->
</dependency> <!-- </dependency>-->
<dependency> <!-- <dependency>-->
<groupId>org.mybatis</groupId> <!-- <groupId>org.mybatis</groupId>-->
<artifactId>mybatis</artifactId> <!-- <artifactId>mybatis</artifactId>-->
</dependency> <!-- </dependency>-->
</dependencies> </dependencies>
</project> </project>

25
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.AllArgsConstructor;
import lombok.extern.slf4j.Slf4j; 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.api.service.MessageRouteService;
import net.sopod.soim.router.cache.RouterUser; import net.sopod.soim.router.cache.RouterUser;
import net.sopod.soim.router.cache.RouterUserStorage; import net.sopod.soim.router.cache.RouterUserStorage;
import org.apache.dubbo.config.annotation.DubboService; import org.apache.dubbo.config.annotation.DubboService;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.ArrayList; import java.util.ArrayList;
import java.util.Iterator;
import java.util.List; import java.util.List;
/** /**
@ -23,10 +25,12 @@ import java.util.List;
@Slf4j @Slf4j
public class MessageRouteServiceImpl implements MessageRouteService { public class MessageRouteServiceImpl implements MessageRouteService {
private static final Logger logger = LoggerFactory.getLogger(MessageRouteServiceImpl.class);
private RouterUserService routerUserService; private RouterUserService routerUserService;
@Override @Override
public List<Boolean> routeGroupMessage(List<Long> uids, String message) { public List<Boolean> routeGroupMessage(Long sender, List<Long> uids, String groupMessage) {
RouterUserStorage storage = RouterUserStorage.getInstance(); RouterUserStorage storage = RouterUserStorage.getInstance();
List<Boolean> results = new ArrayList<>(uids.size()); List<Boolean> results = new ArrayList<>(uids.size());
for (Long uid : uids) { for (Long uid : uids) {
@ -36,10 +40,25 @@ public class MessageRouteServiceImpl implements MessageRouteService {
results.add(false); results.add(false);
continue; continue;
} }
Boolean success = routerUserService.routeGroupMessage(routerUser, message); Boolean success = routerUserService.routeGroupMessage(sender, routerUser, groupMessage);
results.add(Boolean.TRUE.equals(success)); results.add(Boolean.TRUE.equals(success));
} }
return results; 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);
}
} }

16
im-service/im-router/src/main/java/net/sopod/soim/router/service/RouterUserService.java

@ -20,18 +20,26 @@ public class RouterUserService {
@DubboReference @DubboReference
private TextChatService textChatService; 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() TextChat textChat = new TextChat()
.setUid(receiverUser.getUid()) // TODO 群消息优化 .setSenderUid(sender)
.setReceiverUid(receiverUser.getUid()) .setReceiverUid(receiverUser.getUid())
.setMessage(message) .setMessage(message)
.setTime(ImClock.millis()) .setTime(ImClock.millis())
.setReceiverName(receiverUser.getAccount()); .setReceiverName(receiverUser.getAccount());
RpcContextUtil.setContextUid(receiverUser.getUid()); RpcContextUtil.setContextUid(receiverUser.getUid());
// TODO 单独群消息接口
return textChatService.sendTextChat(textChat); return textChatService.sendTextChat(textChat);
} }

23
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()); .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 @Override
public List<Boolean> isOnlineUsers(List<Long> userIds) { public List<Boolean> isOnlineUsers(List<Long> userIds) {

Loading…
Cancel
Save