Browse Source

im-entry适配请求编号

master
tangmingyou 4 years ago
parent
commit
4af0eb34ea
  1. 32
      im-client/src/main/java/net/sopod/soim/client/config/exception/ClientException.java
  2. 9
      im-client/src/main/java/net/sopod/soim/client/protocol/ImMessageInboundHandler.java
  3. 47
      im-client/src/main/java/net/sopod/soim/client/protocol/MessageQueueHolder.java
  4. 53
      im-client/src/main/java/net/sopod/soim/client/session/SoImSession.java
  5. 31
      im-common/src/main/java/net/sopod/soim/common/dubbo/exception/ConvertException.java
  6. 13
      im-common/src/main/java/net/sopod/soim/common/exception/ConvertException.java
  7. 11
      im-common/src/main/java/net/sopod/soim/common/res/R.java
  8. 8
      im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/config/ChatPersistentRabbitMQConfiguration.java
  9. 8
      im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/mq/ChatMQMessageConverter.java
  10. 10
      im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/mq/ChatQueueType.java
  11. 30
      im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/service/DasMQPersistentService.java
  12. 5
      im-das/im-das-user/src/main/java/net/sopod/soim/das/user/amqp/Sender.java
  13. 1
      im-das/im-das-user/src/main/java/net/sopod/soim/das/user/service/FriendDasImpl.java
  14. 2
      im-entry/src/main/java/net/sopod/soim/entry/delay/NetUserDelayTaskManager.java
  15. 3
      im-entry/src/main/java/net/sopod/soim/entry/server/ImEntryInitializer.java
  16. 10
      im-entry/src/main/java/net/sopod/soim/entry/server/InboundImMessageHandler.java
  17. 10
      im-entry/src/main/java/net/sopod/soim/entry/server/ProtoMessageDispatcher.java
  18. 4
      im-entry/src/main/java/net/sopod/soim/entry/server/handler/AccountMessageHandler.java
  19. 22
      im-entry/src/main/java/net/sopod/soim/entry/server/handler/ImContext.java
  20. 2
      im-entry/src/main/java/net/sopod/soim/entry/server/handler/MessageHandler.java
  21. 4
      im-entry/src/main/java/net/sopod/soim/entry/server/handler/NetUserMessageHandler.java
  22. 28
      im-entry/src/main/java/net/sopod/soim/entry/server/session/NetUser.java
  23. 2
      im-service-api/im-entry-protocol/src/main/java/net/sopod/soim/data/proto/ProtoMessageManager.java
  24. 100
      im-service-api/im-entry-protocol/src/main/java/net/sopod/soim/data/serialize/ImMessageCodec.java
  25. 9
      im-service/im-router/src/main/java/net/sopod/soim/router/service/UserRouteServiceImpl.java

32
im-client/src/main/java/net/sopod/soim/client/config/exception/ClientException.java

@ -0,0 +1,32 @@
package net.sopod.soim.client.config.exception;
import net.sopod.soim.common.dubbo.exception.SoimException;
/**
* ClientException
*
* @author tmy
* @date 2022-06-02 22:46
*/
public class ClientException extends SoimException {
public ClientException() {
}
public ClientException(String message) {
super(message);
}
public ClientException(String message, Throwable cause) {
super(message, cause);
}
public ClientException(Throwable cause) {
super(cause);
}
public ClientException(String message, Throwable cause, boolean writableStackTrace) {
super(message, cause, writableStackTrace);
}
}

9
im-client/src/main/java/net/sopod/soim/client/protocol/ImMessageInboundHandler.java

@ -1,8 +1,10 @@
package net.sopod.soim.client.protocol; package net.sopod.soim.client.protocol;
import io.netty.buffer.ByteBuf;
import io.netty.channel.ChannelHandlerContext; import io.netty.channel.ChannelHandlerContext;
import io.netty.channel.SimpleChannelInboundHandler; import io.netty.channel.SimpleChannelInboundHandler;
import net.sopod.soim.data.serialize.ImMessage; import net.sopod.soim.data.serialize.ImMessage;
import net.sopod.soim.data.serialize.ImMessageCodec;
/** /**
* ImMessageInboundHandler * ImMessageInboundHandler
@ -10,13 +12,14 @@ import net.sopod.soim.data.serialize.ImMessage;
* @author tmy * @author tmy
* @date 2022-06-02 17:51 * @date 2022-06-02 17:51
*/ */
public class ImMessageInboundHandler extends SimpleChannelInboundHandler<ImMessage> { public class ImMessageInboundHandler extends SimpleChannelInboundHandler<ByteBuf> {
@Override @Override
protected void channelRead0(ChannelHandlerContext channelHandlerContext, ImMessage imMessage) throws Exception { protected void channelRead0(ChannelHandlerContext channelHandlerContext, ByteBuf bytebuf) throws Exception {
ImMessage imMessage = ImMessageCodec.decodeImMessage(bytebuf);
// 请求序列号,complete 对应 CompletableFuture // 请求序列号,complete 对应 CompletableFuture
int serialNo = imMessage.getSerialNo(); int serialNo = imMessage.getSerialNo();
MessageQueueHolder.getInstance().futureComplete(serialNo, imMessage.getDecodeBody());
} }
} }

47
im-client/src/main/java/net/sopod/soim/client/protocol/MessageQueueHolder.java

@ -1,5 +1,14 @@
package net.sopod.soim.client.protocol; package net.sopod.soim.client.protocol;
import com.google.protobuf.MessageLite;
import org.apache.commons.lang3.tuple.Pair;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.atomic.AtomicInteger;
/** /**
* MessageQueueHolder * MessageQueueHolder
* 发送一条消息产生一个序列号在队列中等待响应消息 * 发送一条消息产生一个序列号在队列中等待响应消息
@ -9,10 +18,46 @@ package net.sopod.soim.client.protocol;
*/ */
public class MessageQueueHolder { public class MessageQueueHolder {
private static final Logger logger = LoggerFactory.getLogger(MessageQueueHolder.class);
private static MessageQueueHolder INSTANCE;
public void a() { private final AtomicInteger serialNoGen;
// TODO 超时处理
private final ConcurrentHashMap<Integer, CompletableFuture<Object>> futureMap;
private MessageQueueHolder() {
this.serialNoGen = new AtomicInteger();
this.futureMap = new ConcurrentHashMap<>();
}
public static MessageQueueHolder getInstance() {
if (INSTANCE == null) {
synchronized (MessageQueueHolder.class) {
if (INSTANCE == null) {
INSTANCE = new MessageQueueHolder();
}
}
}
return INSTANCE;
}
public Pair<Integer, CompletableFuture<Object>> nextSerialNo() {
int serialNo = serialNoGen.incrementAndGet();
CompletableFuture<Object> future = new CompletableFuture<>();
futureMap.put(serialNo, future);
logger.info("put future: {}", serialNo);
return Pair.of(serialNo, future);
}
public void futureComplete(int serialNo, MessageLite body) {
CompletableFuture<Object> completableFuture = futureMap.remove(serialNo);
if (completableFuture == null) {
logger.error("无等待消息的future: {}, {}", serialNo, body);
} else {
completableFuture.complete(body);
}
} }
} }

53
im-client/src/main/java/net/sopod/soim/client/session/SoImSession.java

@ -9,12 +9,20 @@ import io.netty.channel.ChannelInitializer;
import io.netty.channel.nio.NioEventLoopGroup; import io.netty.channel.nio.NioEventLoopGroup;
import io.netty.channel.socket.SocketChannel; import io.netty.channel.socket.SocketChannel;
import io.netty.channel.socket.nio.NioSocketChannel; import io.netty.channel.socket.nio.NioSocketChannel;
import io.netty.handler.logging.LogLevel;
import io.netty.handler.logging.LoggingHandler;
import net.sopod.soim.client.config.exception.ClientException;
import net.sopod.soim.client.logger.Logger; import net.sopod.soim.client.logger.Logger;
import net.sopod.soim.common.res.R; import net.sopod.soim.client.protocol.ImMessageInboundHandler;
import net.sopod.soim.client.protocol.MessageQueueHolder;
import net.sopod.soim.common.util.ImClock;
import net.sopod.soim.common.util.netty.Varint32FrameCodec; import net.sopod.soim.common.util.netty.Varint32FrameCodec;
import net.sopod.soim.data.proto.ProtoMessageManager;
import net.sopod.soim.data.serialize.ImMessage;
import net.sopod.soim.data.serialize.ImMessageCodec; import net.sopod.soim.data.serialize.ImMessageCodec;
import net.sopod.soim.data.msg.auth.Auth; import net.sopod.soim.data.msg.auth.Auth;
import net.sopod.soim.data.msg.chat.Chat; import net.sopod.soim.data.msg.chat.Chat;
import org.apache.commons.lang3.tuple.Pair;
import org.slf4j.LoggerFactory; import org.slf4j.LoggerFactory;
import java.util.concurrent.CompletableFuture; import java.util.concurrent.CompletableFuture;
@ -60,18 +68,23 @@ public class SoImSession {
@Override @Override
protected void initChannel(SocketChannel ch) throws Exception { protected void initChannel(SocketChannel ch) throws Exception {
ch.pipeline() ch.pipeline()
.addLast(new LoggingHandler(LogLevel.INFO))
.addLast(new Varint32FrameCodec()) .addLast(new Varint32FrameCodec())
.addLast(new ImMessageCodec.ProtoMsg2ImMessageEncoder()) // .addLast(new ImMessageCodec())
.addLast(new ImMessageCodec.ImMessage2ByteEncoder()) .addLast(new ImMessageCodec.ImMessage2ByteEncoder())
.addLast(new ImMessageCodec.ImMessageDecoder()) // .addLast(new ImMessageCodec.ImMessageDecoder())
.addLast(messageDispatcher); .addLast(new ImMessageInboundHandler());
// .addLast(messageDispatcher);
} }
}); });
try { try {
Logger.info("连接中..."); Logger.info("连接中...");
clientChannel = b.connect(host, port).await().channel(); clientChannel = b.connect(host, port).await().channel();
// 连接后立即发送认证消息,10s未认证连接关闭 // 连接后立即发送认证消息,10s未认证连接关闭
clientChannel.writeAndFlush(tokenAuth); CompletableFuture<?> future = this.send0(tokenAuth);
long start = ImClock.millis();
Object res = future.join();
logger.info("future 响应: {}, {}ms", res, ImClock.millis() - start);
this.account = account; this.account = account;
} catch (Exception e) { } catch (Exception e) {
Logger.error("连接服务器失败: {}", e.getMessage()); Logger.error("连接服务器失败: {}", e.getMessage());
@ -92,21 +105,37 @@ public class SoImSession {
this.close(); this.close();
} }
private CompletableFuture<?> send0(MessageLite message) {
if (clientChannel == null
|| !clientChannel.isActive()) {
throw new ClientException("连接已关闭");
}
// 构建 ImMessage
Integer serialNo = ProtoMessageManager.getSerialNo(message.getClass());
if (serialNo == null) {
throw new ClientException("未知的消息类型:" + message.getClass());
}
ImMessage imMessage = new ImMessage()
.setServiceNo(serialNo)
.setBody(message.toByteArray());
// serialNo, Future
Pair<Integer, CompletableFuture<Object>> futurePair = MessageQueueHolder.getInstance().nextSerialNo();
imMessage.setSerialNo(futurePair.getLeft());
clientChannel.writeAndFlush(imMessage);
return futurePair.getRight();
}
/** /**
* 发送消息 * 发送消息
*/ */
@SuppressWarnings("unchecked")
public <T> CompletableFuture<T> send(MessageLite message) { public <T> CompletableFuture<T> send(MessageLite message) {
if (!auth.get()) { if (!auth.get()) {
Logger.error("请先登录"); Logger.error("请先登录");
return null; return null;
} }
if (clientChannel == null return (CompletableFuture<T>) send0(message);
|| !clientChannel.isActive()) {
Logger.error("连接已关闭");
return null;
}
clientChannel.writeAndFlush(message);
return new CompletableFuture<>();
} }
public void textChat(String receiverName, String message) { public void textChat(String receiverName, String message) {

31
im-common/src/main/java/net/sopod/soim/common/dubbo/exception/ConvertException.java

@ -0,0 +1,31 @@
package net.sopod.soim.common.dubbo.exception;
/**
* ConvertException
*
* @author tmy
* @date 2022-06-02 22:58
*/
public class ConvertException extends SoimException {
public ConvertException() {
super();
}
public ConvertException(String message) {
super(message);
}
public ConvertException(String message, Throwable cause) {
super(message, cause);
}
public ConvertException(Throwable cause) {
super(cause);
}
public ConvertException(String message, Throwable cause, boolean writableStackTrace) {
super(message, cause, writableStackTrace);
}
}

13
im-common/src/main/java/net/sopod/soim/common/exception/ConvertException.java

@ -1,13 +0,0 @@
package net.sopod.soim.common.exception;
/**
* ConvertException
*
* @author tmy
* @date 2022-03-27 17:18
*/
public class ConvertException extends Exception {
private Throwable originException;
}

11
im-common/src/main/java/net/sopod/soim/common/res/R.java

@ -1,11 +0,0 @@
package net.sopod.soim.common.res;
/**
* R
*
* @author tmy
* @date 2022-06-02 15:31
*/
public class R {
}

8
im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/config/ChatPersistentRabbitMQConfiguration.java

@ -2,16 +2,15 @@ package net.sopod.soim.das.user.api.config;
import net.sopod.soim.das.user.api.mq.ChatMQMessageConverter; import net.sopod.soim.das.user.api.mq.ChatMQMessageConverter;
import net.sopod.soim.das.user.api.mq.ChatQueueType; import net.sopod.soim.das.user.api.mq.ChatQueueType;
import net.sopod.soim.das.user.api.service.DasMQPersistentService;
import org.springframework.amqp.core.Queue; import org.springframework.amqp.core.Queue;
import org.springframework.amqp.rabbit.connection.CachingConnectionFactory; import org.springframework.amqp.rabbit.connection.CachingConnectionFactory;
import org.springframework.amqp.rabbit.core.RabbitTemplate; import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.amqp.support.converter.MessageConverter; import org.springframework.amqp.support.converter.MessageConverter;
import org.springframework.beans.factory.support.BeanDefinitionRegistry; import org.springframework.beans.factory.support.BeanDefinitionRegistry;
import org.springframework.beans.factory.support.RootBeanDefinition; import org.springframework.beans.factory.support.RootBeanDefinition;
import org.springframework.boot.autoconfigure.condition.ConditionalOnClass;
import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration; import org.springframework.context.annotation.Configuration;
import org.springframework.context.annotation.Import;
import org.springframework.context.annotation.ImportBeanDefinitionRegistrar; import org.springframework.context.annotation.ImportBeanDefinitionRegistrar;
import org.springframework.core.type.AnnotationMetadata; import org.springframework.core.type.AnnotationMetadata;
@ -53,4 +52,9 @@ public class ChatPersistentRabbitMQConfiguration implements ImportBeanDefinition
return rabbitTemplate; return rabbitTemplate;
} }
@Bean
public DasMQPersistentService dasPersistentService(RabbitTemplate rabbitTemplate) {
return new DasMQPersistentService(rabbitTemplate);
}
} }

8
im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/mq/ChatMQMessageConverter.java

@ -60,8 +60,12 @@ public class ChatMQMessageConverter extends AbstractMessageConverter {
MessageProperties messageProperties = message.getMessageProperties(); MessageProperties messageProperties = message.getMessageProperties();
Object ordinal = messageProperties.getHeaders().get(MQ_MSG_TYPE); Object ordinal = messageProperties.getHeaders().get(MQ_MSG_TYPE);
Object isSnappy = messageProperties.getHeaders().get(MQ_MSG_IS_SNAPPY); Object isSnappy = messageProperties.getHeaders().get(MQ_MSG_IS_SNAPPY);
ChatQueueType chatQueueType = ChatQueueType.getMQType(Converter.toInt(ordinal), String consumerQueue = messageProperties.getConsumerQueue();
ori -> new MessageConversionException(String.format("没有ordinal为%d的消息类型,在MQType添加", ori))); // 获取mq消息类型
ChatQueueType chatQueueType = ordinal == null
? ChatQueueType.getMQType(consumerQueue, queue -> new MessageConversionException(String.format("没有队列为%s的消息类型,在MQType添加", queue)))
: ChatQueueType.getMQType(Converter.toInt(ordinal), ori -> new MessageConversionException(String.format("没有ordinal为%d的消息类型,在MQType添加", ori)));
// 反序列化字节数据
byte[] body = message.getBody(); byte[] body = message.getBody();
if (Boolean.TRUE.equals(isSnappy)) { if (Boolean.TRUE.equals(isSnappy)) {
try { try {

10
im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/mq/ChatQueueType.java

@ -70,4 +70,14 @@ public enum ChatQueueType {
throw exceptionProvider.apply(ordinal); throw exceptionProvider.apply(ordinal);
} }
public static ChatQueueType getMQType(String queueName, Function<String, RuntimeException> exceptionProvider) {
ChatQueueType[] values = values();
for (ChatQueueType queueType : values) {
if (queueType.getQueueName().equals(queueName)) {
return queueType;
}
}
throw exceptionProvider.apply(queueName);
}
} }

30
im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/service/DasMQPersistentService.java

@ -0,0 +1,30 @@
package net.sopod.soim.das.user.api.service;
import net.sopod.soim.das.user.api.model.entity.ImGroupMessage;
import net.sopod.soim.das.user.api.model.entity.ImMessage;
import net.sopod.soim.das.user.api.mq.ChatQueueType;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
/**
* DasPersistentService
*
* @author tmy
* @date 2022-06-02 20:28
*/
public class DasMQPersistentService {
private final RabbitTemplate rabbitTemplate;
public DasMQPersistentService(RabbitTemplate rabbitTemplate) {
this.rabbitTemplate = rabbitTemplate;
}
public void saveImMessage(ImMessage imMessage) {
rabbitTemplate.convertAndSend(ChatQueueType.IM_MESSAGE.getQueueName(), imMessage);
}
public void saveImGroupMessage(ImGroupMessage imGroupMessage) {
rabbitTemplate.convertAndSend(ChatQueueType.IM_GROUP_MESSAGE.getQueueName(), imGroupMessage);
}
}

5
im-das/im-das-user/src/main/java/net/sopod/soim/das/user/amqp/Sender.java

@ -7,14 +7,11 @@ import net.sopod.soim.das.user.api.model.entity.ImGroupMessage;
import net.sopod.soim.das.user.api.model.entity.ImMessage; import net.sopod.soim.das.user.api.model.entity.ImMessage;
import net.sopod.soim.das.user.api.mq.ChatQueueType; import net.sopod.soim.das.user.api.mq.ChatQueueType;
import net.sopod.soim.logic.api.segmentid.core.SegmentIdGenerator; import net.sopod.soim.logic.api.segmentid.core.SegmentIdGenerator;
import org.apache.dubbo.config.annotation.DubboReference;
import org.springframework.amqp.core.AmqpTemplate; import org.springframework.amqp.core.AmqpTemplate;
import org.springframework.boot.context.event.ApplicationReadyEvent; import org.springframework.boot.context.event.ApplicationReadyEvent;
import org.springframework.context.ApplicationListener; import org.springframework.context.ApplicationListener;
import org.springframework.context.ConfigurableApplicationContext; import org.springframework.context.ConfigurableApplicationContext;
import org.springframework.stereotype.Component;
import java.util.Random;
import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeUnit;
/** /**
@ -23,7 +20,7 @@ import java.util.concurrent.TimeUnit;
* @author tmy * @author tmy
* @date 2022-05-28 16:23 * @date 2022-05-28 16:23
*/ */
@Component //@Component
@AllArgsConstructor @AllArgsConstructor
public class Sender implements ApplicationListener<ApplicationReadyEvent> { public class Sender implements ApplicationListener<ApplicationReadyEvent> {

1
im-das/im-das-user/src/main/java/net/sopod/soim/das/user/service/FriendDasImpl.java

@ -3,7 +3,6 @@ package net.sopod.soim.das.user.service;
import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper; import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
import com.baomidou.mybatisplus.core.conditions.query.QueryWrapper; import com.baomidou.mybatisplus.core.conditions.query.QueryWrapper;
import lombok.AllArgsConstructor; import lombok.AllArgsConstructor;
import net.sopod.soim.common.constant.LogicConsts;
import net.sopod.soim.common.util.Collects; import net.sopod.soim.common.util.Collects;
import net.sopod.soim.common.util.ImClock; import net.sopod.soim.common.util.ImClock;
import net.sopod.soim.das.user.api.config.LogicTables; import net.sopod.soim.das.user.api.config.LogicTables;

2
im-entry/src/main/java/net/sopod/soim/entry/delay/NetUserDelayTaskManager.java

@ -41,7 +41,7 @@ public class NetUserDelayTaskManager {
if (netUser.isActive()) { if (netUser.isActive()) {
// 延时任务执行时 NetUser 可能已升级为 Account // 延时任务执行时 NetUser 可能已升级为 Account
NetUser curNetUser = netUser.channel().attr(NetUser.NET_USER_KEY).get(); NetUser curNetUser = netUser.channel().attr(NetUser.NET_USER_KEY).get();
ProtoMessageDispatcher.dispatch(curNetUser, delayTask.getTaskMsg()); ProtoMessageDispatcher.dispatch(null, curNetUser, delayTask.getTaskMsg());
} else { } else {
logger.info("delayTask netUser inactive, task cancel!"); logger.info("delayTask netUser inactive, task cancel!");
} }

3
im-entry/src/main/java/net/sopod/soim/entry/server/ImEntryInitializer.java

@ -27,7 +27,8 @@ public class ImEntryInitializer extends ChannelInitializer<SocketChannel> {
ChannelPipeline pipeline = socketChannel.pipeline(); ChannelPipeline pipeline = socketChannel.pipeline();
pipeline.addLast(new LoggingHandler(logLevel)) pipeline.addLast(new LoggingHandler(logLevel))
.addLast(new Varint32FrameCodec()) .addLast(new Varint32FrameCodec())
.addLast(new ImMessageCodec()) .addLast(new ImMessageCodec.ImMessageDecoder())
.addLast(new ImMessageCodec.ImMessage2ByteEncoder())
.addLast(new InboundImMessageHandler()); .addLast(new InboundImMessageHandler());
} }

10
im-entry/src/main/java/net/sopod/soim/entry/server/InboundImMessageHandler.java

@ -1,11 +1,12 @@
package net.sopod.soim.entry.server; package net.sopod.soim.entry.server;
import com.google.protobuf.MessageLite;
import io.netty.channel.Channel; import io.netty.channel.Channel;
import io.netty.channel.ChannelHandlerContext; import io.netty.channel.ChannelHandlerContext;
import io.netty.channel.SimpleChannelInboundHandler; import io.netty.channel.SimpleChannelInboundHandler;
import io.netty.util.Attribute; import io.netty.util.Attribute;
import net.sopod.soim.common.util.ImClock; import net.sopod.soim.common.util.ImClock;
import net.sopod.soim.data.serialize.ImMessage;
import net.sopod.soim.entry.server.handler.ImContext;
import net.sopod.soim.entry.server.session.NetUser; import net.sopod.soim.entry.server.session.NetUser;
import net.sopod.soim.data.msg.task.Tasks; import net.sopod.soim.data.msg.task.Tasks;
import net.sopod.soim.entry.delay.NetUserDelayTaskManager; import net.sopod.soim.entry.delay.NetUserDelayTaskManager;
@ -19,7 +20,7 @@ import java.util.concurrent.atomic.AtomicInteger;
* @author tmy * @author tmy
* @date 2022-04-10 22:40 * @date 2022-04-10 22:40
*/ */
public class InboundImMessageHandler extends SimpleChannelInboundHandler<MessageLite> { public class InboundImMessageHandler extends SimpleChannelInboundHandler<ImMessage> {
/** /**
* channel 建立连接设置初始属性 * channel 建立连接设置初始属性
@ -50,11 +51,12 @@ public class InboundImMessageHandler extends SimpleChannelInboundHandler<Message
} }
@Override @Override
protected void channelRead0(ChannelHandlerContext ctx, MessageLite messageLite) throws Exception { protected void channelRead0(ChannelHandlerContext ctx, ImMessage imMessage) throws Exception {
Attribute<NetUser> netUserAttr = ctx.channel().attr(NetUser.NET_USER_KEY); Attribute<NetUser> netUserAttr = ctx.channel().attr(NetUser.NET_USER_KEY);
NetUser netUser = netUserAttr.get(); NetUser netUser = netUserAttr.get();
ImContext execCtx = new ImContext(imMessage.getSerialNo());
// TODO dispatcher // TODO dispatcher
ProtoMessageDispatcher.dispatch(netUser, messageLite); ProtoMessageDispatcher.dispatch(execCtx, netUser, imMessage.getDecodeBody());
} }
} }

10
im-entry/src/main/java/net/sopod/soim/entry/server/ProtoMessageDispatcher.java

@ -2,6 +2,8 @@ package net.sopod.soim.entry.server;
import com.google.protobuf.MessageLite; import com.google.protobuf.MessageLite;
import net.sopod.soim.common.constant.DubboConstant; import net.sopod.soim.common.constant.DubboConstant;
import net.sopod.soim.common.dubbo.exception.SoimException;
import net.sopod.soim.entry.server.handler.ImContext;
import net.sopod.soim.entry.server.handler.MessageHandler; import net.sopod.soim.entry.server.handler.MessageHandler;
import net.sopod.soim.entry.registry.ProtoMessageHandlerRegistry; import net.sopod.soim.entry.registry.ProtoMessageHandlerRegistry;
import net.sopod.soim.entry.server.session.Account; import net.sopod.soim.entry.server.session.Account;
@ -18,20 +20,20 @@ import net.sopod.soim.entry.worker.WorkerGroup;
*/ */
public class ProtoMessageDispatcher { public class ProtoMessageDispatcher {
public static void dispatch(NetUser netUser, MessageLite message) { public static void dispatch(ImContext ctx, NetUser netUser, MessageLite message) {
MessageHandler<MessageLite> typeHandler = ProtoMessageHandlerRegistry MessageHandler<MessageLite> typeHandler = ProtoMessageHandlerRegistry
.getTypeHandler(message.getClass()); .getTypeHandler(message.getClass());
if (typeHandler == null) { if (typeHandler == null) {
throw new IllegalCallerException("no handler for message : " + message.getClass()); throw new SoimException("no handler for message type: " + message.getClass());
} }
Worker worker = WorkerGroup.next(); Worker worker = WorkerGroup.next();
worker.execute(() -> { worker.execute(() -> {
if (netUser.isAccount()) { if (netUser.isAccount()) {
Account account = (Account)netUser; Account account = (Account) netUser;
MessageHandlerContext.setAttribute(DubboConstant.CTX_UID, String.valueOf(account.getUid())); MessageHandlerContext.setAttribute(DubboConstant.CTX_UID, String.valueOf(account.getUid()));
} }
try { try {
typeHandler.exec(netUser, message); typeHandler.exec(ctx, netUser, message);
} finally { } finally {
MessageHandlerContext.remove(); MessageHandlerContext.remove();
} }

4
im-entry/src/main/java/net/sopod/soim/entry/server/handler/AccountMessageHandler.java

@ -14,13 +14,13 @@ import net.sopod.soim.entry.server.session.NetUser;
public abstract class AccountMessageHandler<T> implements MessageHandler<T> { public abstract class AccountMessageHandler<T> implements MessageHandler<T> {
@Override @Override
public final void exec(NetUser netUser, T req) { public final void exec(ImContext ctx, NetUser netUser, T req) {
if (!netUser.isAccount()) { if (!netUser.isAccount()) {
throw new IllegalStateException("NetUser is not account!" + netUser); throw new IllegalStateException("NetUser is not account!" + netUser);
} }
MessageLite res = handle((Account) netUser, req); MessageLite res = handle((Account) netUser, req);
if (res != null) { if (res != null) {
netUser.writeNow(res); netUser.writeNow(ctx, res);
} }
} }

22
im-entry/src/main/java/net/sopod/soim/entry/server/handler/ImContext.java

@ -0,0 +1,22 @@
package net.sopod.soim.entry.server.handler;
/**
* ExecContext
*
* @author tmy
* @date 2022-06-02 23:31
*/
public class ImContext {
/** 请求id */
private final Integer serialNo;
public ImContext(Integer serialNo) {
this.serialNo = serialNo;
}
public Integer getSerialNo() {
return serialNo;
}
}

2
im-entry/src/main/java/net/sopod/soim/entry/server/handler/MessageHandler.java

@ -10,6 +10,6 @@ import net.sopod.soim.entry.server.session.NetUser;
*/ */
public interface MessageHandler<T> { public interface MessageHandler<T> {
void exec(NetUser netUser, T msg); void exec(ImContext ctx, NetUser netUser, T msg);
} }

4
im-entry/src/main/java/net/sopod/soim/entry/server/handler/NetUserMessageHandler.java

@ -12,10 +12,10 @@ import net.sopod.soim.entry.server.session.NetUser;
public abstract class NetUserMessageHandler<T> implements MessageHandler<T> { public abstract class NetUserMessageHandler<T> implements MessageHandler<T> {
@Override @Override
public final void exec(NetUser netUser, T msg) { public final void exec(ImContext ctx, NetUser netUser, T msg) {
MessageLite res = handle(netUser, msg); MessageLite res = handle(netUser, msg);
if (res != null) { if (res != null) {
netUser.writeNow(res); netUser.writeNow(ctx, res);
} }
} }

28
im-entry/src/main/java/net/sopod/soim/entry/server/session/NetUser.java

@ -1,12 +1,20 @@
package net.sopod.soim.entry.server.session; package net.sopod.soim.entry.server.session;
import com.google.protobuf.MessageLite;
import io.netty.channel.Channel; import io.netty.channel.Channel;
import io.netty.util.AttributeKey; import io.netty.util.AttributeKey;
import net.sopod.soim.data.serialize.ImMessage;
import net.sopod.soim.data.serialize.ImMessageCodec;
import net.sopod.soim.entry.server.handler.ImContext;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.lang.ref.WeakReference; import java.lang.ref.WeakReference;
public class NetUser { public class NetUser {
private static final Logger logger = LoggerFactory.getLogger(NetUser.class);
/** channel 绑定 netUser 对象 */ /** channel 绑定 netUser 对象 */
public static final AttributeKey<NetUser> NET_USER_KEY = AttributeKey.valueOf(NetUser.class, "NET_USER"); public static final AttributeKey<NetUser> NET_USER_KEY = AttributeKey.valueOf(NetUser.class, "NET_USER");
@ -39,24 +47,24 @@ public class NetUser {
} }
} }
public void write(Object message) { public void writeNow(MessageLite message) {
write(message, false); write0(null, message);
} }
public void writeNow(Object message) { public void writeNow(ImContext ctx, MessageLite message) {
write(message, true); write0(ctx.getSerialNo(), message);
} }
private void write(Object message, boolean now) { private void write0(Integer serialNo, MessageLite message) {
Channel channel = this.channel.get(); Channel channel = this.channel.get();
if (channel == null) { if (channel == null) {
logger.error("channel closed drop message: {}, {}", serialNo, message);
return; return;
} }
if (now) { // 构建 ImMessage
channel.writeAndFlush(message); ImMessage imMessage = ImMessageCodec.encodeImProto(message);
} else { imMessage.setSerialNo(serialNo == null ? -1 : serialNo);
channel.write(message); channel.writeAndFlush(imMessage);
}
} }
@Override @Override

2
im-service-api/im-entry-protocol/src/main/java/net/sopod/soim/data/proto/ProtoMessageManager.java

@ -60,7 +60,7 @@ public class ProtoMessageManager {
} }
@Nullable @Nullable
public static MessageLite getProtoInstance(Integer serialNo) { public static MessageLite getDefaultInstance(Integer serialNo) {
String clazz = serialNoTypeMap.get(serialNo); String clazz = serialNoTypeMap.get(serialNo);
if (clazz == null) { if (clazz == null) {
logger.error("protoMsgDict serialNo proto class not found: {}", serialNo); logger.error("protoMsgDict serialNo proto class not found: {}", serialNo);

100
im-service-api/im-entry-protocol/src/main/java/net/sopod/soim/data/serialize/ImMessageCodec.java

@ -1,18 +1,17 @@
package net.sopod.soim.data.serialize; package net.sopod.soim.data.serialize;
import com.google.protobuf.InvalidProtocolBufferException;
import com.google.protobuf.MessageLite; import com.google.protobuf.MessageLite;
import io.netty.buffer.ByteBuf; import io.netty.buffer.ByteBuf;
import io.netty.channel.ChannelHandlerContext; import io.netty.channel.ChannelHandlerContext;
import io.netty.channel.CombinedChannelDuplexHandler;
import io.netty.handler.codec.ByteToMessageDecoder; import io.netty.handler.codec.ByteToMessageDecoder;
import io.netty.handler.codec.MessageToByteEncoder; import io.netty.handler.codec.MessageToByteEncoder;
import io.netty.handler.codec.MessageToMessageEncoder; import net.sopod.soim.common.dubbo.exception.ConvertException;
import net.sopod.soim.common.dubbo.exception.SoimException;
import net.sopod.soim.data.proto.ProtoMessageManager; import net.sopod.soim.data.proto.ProtoMessageManager;
import org.slf4j.Logger; import org.slf4j.Logger;
import org.slf4j.LoggerFactory; import org.slf4j.LoggerFactory;
import java.lang.reflect.Type;
import java.util.Arrays;
import java.util.List; import java.util.List;
/** /**
@ -21,90 +20,57 @@ import java.util.List;
* @author tmy * @author tmy
* @date 2022-03-28 11:29 * @date 2022-03-28 11:29
*/ */
public class ImMessageCodec //extends MessageToMessageCodec<ByteBuf, MessageLite> { public class ImMessageCodec {
extends CombinedChannelDuplexHandler<ImMessageCodec.ProtoMsgDecoder, ImMessageCodec.ProtoMsgEncoder> {
private static final Logger logger = LoggerFactory.getLogger(ImMessageCodec.class); private static final Logger logger = LoggerFactory.getLogger(ImMessageCodec.class);
public ImMessageCodec() {
super(new ProtoMsgDecoder(), new ProtoMsgEncoder());
}
public static class ImMessageDecoder extends ByteToMessageDecoder { public static class ImMessageDecoder extends ByteToMessageDecoder {
@Override @Override
protected void decode(ChannelHandlerContext ctx, ByteBuf byteBuf, List<Object> out) throws Exception { protected void decode(ChannelHandlerContext ctx, ByteBuf byteBuf, List<Object> out) throws Exception {
ImMessage message = ImMessage.read(byteBuf); ImMessage imMessage = ImMessageCodec.decodeImMessage(byteBuf);
boolean isMagicError; out.add(imMessage);
if ((isMagicError = (message == ImMessage.MAGIC_ERROR))
|| message == ImMessage.PROTOCOL_ERROR) {
logger.warn("decode im message error: {}, remote={}, closing channel.",
isMagicError ? "MagicError" : "ProtocolError",
ctx.channel().remoteAddress());
ctx.channel().close();
return;
}
// 解码 protobuf 消息体
int serviceNo = message.getServiceNo();
byte[] protoByte = message.getBody();
MessageLite protoClass = ProtoMessageManager.getProtoInstance(serviceNo);
MessageLite protoMsg = protoClass.getParserForType().parseFrom(protoByte);
message.setDecodeBody(protoMsg);
out.add(message);
} }
} }
public static class ProtoMsgDecoder extends ByteToMessageDecoder { public static class ImMessage2ByteEncoder extends MessageToByteEncoder<ImMessage> {
@Override @Override
protected void decode(ChannelHandlerContext ctx, ByteBuf byteBuf, List<Object> list) throws Exception { protected void encode(ChannelHandlerContext ctx, ImMessage imMessage, ByteBuf byteBuf) throws Exception {
ImMessage message = ImMessage.read(byteBuf); imMessage.write(byteBuf);
boolean isMagicError;
if ((isMagicError = (message == ImMessage.MAGIC_ERROR))
|| message == ImMessage.PROTOCOL_ERROR) {
logger.warn("decode im message error: {}, remote={}, closing channel.",
isMagicError ? "MagicError" : "ProtocolError",
ctx.channel().remoteAddress());
ctx.channel().close();
return;
}
// 解码 protobuf 消息体
int serviceNo = message.getServiceNo();
// message.getSerialNo()
byte[] protoByte = message.getBody();
MessageLite protoClass = ProtoMessageManager.getProtoInstance(serviceNo);
MessageLite protoMsg = protoClass.getParserForType().parseFrom(protoByte);
list.add(protoMsg);
} }
} }
public static class ProtoMsg2ImMessageEncoder extends MessageToMessageEncoder<MessageLite> { public static ImMessage decodeImMessage(ByteBuf byteBuf) {
@Override ImMessage imMessage = ImMessage.read(byteBuf);
protected void encode(ChannelHandlerContext ctx, MessageLite message, List<Object> out) throws Exception { boolean isMagicError;
Integer serialNo = ProtoMessageManager.getSerialNo(message.getClass()); if ((isMagicError = (imMessage == ImMessage.MAGIC_ERROR))
// TODO unknow class serialNo || imMessage == ImMessage.PROTOCOL_ERROR) {
ImMessage imMessage = new ImMessage() throw new ConvertException("ImMessage解码错误");
.setServiceNo(serialNo)
.setBody(message.toByteArray());
out.add(imMessage);
} }
// 解码 protobuf 消息体
int serviceNo = imMessage.getServiceNo();
byte[] protoByte = imMessage.getBody();
MessageLite protoClass = ProtoMessageManager.getDefaultInstance(serviceNo);
if (protoClass == null) {
throw new SoimException(String.format("不支持的消息编号: %d", serviceNo));
} }
MessageLite protoMsg;
public static class ImMessage2ByteEncoder extends MessageToByteEncoder<ImMessage> { try {
@Override protoMsg = protoClass.getParserForType().parseFrom(protoByte);
protected void encode(ChannelHandlerContext ctx, ImMessage imMessage, ByteBuf out) throws Exception { } catch (InvalidProtocolBufferException e) {
imMessage.write(out); throw new ConvertException("ImMessage消息体Protobuf解码错误: " + e.getMessage());
} }
imMessage.setDecodeBody(protoMsg);
return imMessage;
} }
public static class ProtoMsgEncoder extends MessageToByteEncoder<MessageLite> { public static ImMessage encodeImProto(MessageLite message) {
@Override
protected void encode(ChannelHandlerContext ctx, MessageLite message, ByteBuf byteBuf) throws Exception {
Integer serialNo = ProtoMessageManager.getSerialNo(message.getClass()); Integer serialNo = ProtoMessageManager.getSerialNo(message.getClass());
// TODO unknow class serialNo if (serialNo == null) {
ImMessage imMessage = new ImMessage() throw new SoimException("未知的Protobuf消息类型:" + message.getClass());
}
return new ImMessage()
.setServiceNo(serialNo) .setServiceNo(serialNo)
.setBody(message.toByteArray()); .setBody(message.toByteArray());
imMessage.write(byteBuf);
}
} }
} }

9
im-service/im-router/src/main/java/net/sopod/soim/router/service/UserRouteServiceImpl.java

@ -1,14 +1,12 @@
package net.sopod.soim.router.service; package net.sopod.soim.router.service;
import net.sopod.soim.common.constant.DubboConstant; import net.sopod.soim.common.constant.DubboConstant;
import net.sopod.soim.common.dubbo.exception.ServiceException;
import net.sopod.soim.common.util.ImClock; import net.sopod.soim.common.util.ImClock;
import net.sopod.soim.common.util.StringUtil; import net.sopod.soim.common.util.StringUtil;
import net.sopod.soim.das.user.api.config.LogicTables; import net.sopod.soim.das.user.api.config.LogicTables;
import net.sopod.soim.das.user.api.model.entity.ImMessage; import net.sopod.soim.das.user.api.model.entity.ImMessage;
import net.sopod.soim.das.user.api.model.entity.ImUser; import net.sopod.soim.das.user.api.model.entity.ImUser;
import net.sopod.soim.das.user.api.mq.ChatQueue; import net.sopod.soim.das.user.api.service.DasMQPersistentService;
import net.sopod.soim.das.user.api.mq.ChatQueueType;
import net.sopod.soim.das.user.api.service.FriendDas; import net.sopod.soim.das.user.api.service.FriendDas;
import net.sopod.soim.das.user.api.service.UserDas; import net.sopod.soim.das.user.api.service.UserDas;
import net.sopod.soim.entry.api.service.OnlineUserService; import net.sopod.soim.entry.api.service.OnlineUserService;
@ -27,7 +25,6 @@ import org.apache.dubbo.config.annotation.DubboService;
import org.apache.dubbo.rpc.RpcContext; import org.apache.dubbo.rpc.RpcContext;
import org.slf4j.Logger; import org.slf4j.Logger;
import org.slf4j.LoggerFactory; import org.slf4j.LoggerFactory;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import javax.annotation.Resource; import javax.annotation.Resource;
import java.util.ArrayList; import java.util.ArrayList;
@ -62,7 +59,7 @@ public class UserRouteServiceImpl implements UserRouteService {
private SegmentIdGenerator segmentIdGenerator; private SegmentIdGenerator segmentIdGenerator;
@Resource @Resource
private RabbitTemplate rabbitTemplate; private DasMQPersistentService dasMQPersistentService;
@Override @Override
public RegistryRes registryUserEntry(Long uid, String imEntryAddr) { public RegistryRes registryUserEntry(Long uid, String imEntryAddr) {
@ -107,7 +104,7 @@ public class UserRouteServiceImpl implements UserRouteService {
.setSender(textChat.getUid()) .setSender(textChat.getUid())
.setReceiver(textChat.getReceiverUid()) .setReceiver(textChat.getReceiverUid())
.setCreateTime(ImClock.date()); .setCreateTime(ImClock.date());
rabbitTemplate.convertAndSend(ChatQueueType.IM_MESSAGE.getQueueName(), imMessage); dasMQPersistentService.saveImMessage(imMessage);
RpcContextUtil.setContextUid(textChat.getReceiverUid()); RpcContextUtil.setContextUid(textChat.getReceiverUid());
Boolean send = textChatService.sendTextChat(textChat); Boolean send = textChatService.sendTextChat(textChat);

Loading…
Cancel
Save