diff --git a/im-client/src/main/java/net/sopod/soim/client/config/exception/ClientException.java b/im-client/src/main/java/net/sopod/soim/client/config/exception/ClientException.java new file mode 100644 index 0000000..c1a1991 --- /dev/null +++ b/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); + } + +} diff --git a/im-client/src/main/java/net/sopod/soim/client/protocol/ImMessageInboundHandler.java b/im-client/src/main/java/net/sopod/soim/client/protocol/ImMessageInboundHandler.java index 39e9e45..c6cdc58 100644 --- a/im-client/src/main/java/net/sopod/soim/client/protocol/ImMessageInboundHandler.java +++ b/im-client/src/main/java/net/sopod/soim/client/protocol/ImMessageInboundHandler.java @@ -1,8 +1,10 @@ package net.sopod.soim.client.protocol; +import io.netty.buffer.ByteBuf; import io.netty.channel.ChannelHandlerContext; import io.netty.channel.SimpleChannelInboundHandler; import net.sopod.soim.data.serialize.ImMessage; +import net.sopod.soim.data.serialize.ImMessageCodec; /** * ImMessageInboundHandler @@ -10,13 +12,14 @@ import net.sopod.soim.data.serialize.ImMessage; * @author tmy * @date 2022-06-02 17:51 */ -public class ImMessageInboundHandler extends SimpleChannelInboundHandler { +public class ImMessageInboundHandler extends SimpleChannelInboundHandler { @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 int serialNo = imMessage.getSerialNo(); - + MessageQueueHolder.getInstance().futureComplete(serialNo, imMessage.getDecodeBody()); } } diff --git a/im-client/src/main/java/net/sopod/soim/client/protocol/MessageQueueHolder.java b/im-client/src/main/java/net/sopod/soim/client/protocol/MessageQueueHolder.java index 0ee47cd..f8fd727 100644 --- a/im-client/src/main/java/net/sopod/soim/client/protocol/MessageQueueHolder.java +++ b/im-client/src/main/java/net/sopod/soim/client/protocol/MessageQueueHolder.java @@ -1,5 +1,14 @@ 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 * 发送一条消息产生一个序列号,在队列中等待响应消息 @@ -9,10 +18,46 @@ package net.sopod.soim.client.protocol; */ public class MessageQueueHolder { + private static final Logger logger = LoggerFactory.getLogger(MessageQueueHolder.class); + + private static MessageQueueHolder INSTANCE; + private final AtomicInteger serialNoGen; - public void a() { + // TODO 超时处理 + private final ConcurrentHashMap> 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> nextSerialNo() { + int serialNo = serialNoGen.incrementAndGet(); + CompletableFuture 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 completableFuture = futureMap.remove(serialNo); + if (completableFuture == null) { + logger.error("无等待消息的future: {}, {}", serialNo, body); + } else { + completableFuture.complete(body); + } } } diff --git a/im-client/src/main/java/net/sopod/soim/client/session/SoImSession.java b/im-client/src/main/java/net/sopod/soim/client/session/SoImSession.java index d7fca2b..09e97ab 100644 --- a/im-client/src/main/java/net/sopod/soim/client/session/SoImSession.java +++ b/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.socket.SocketChannel; 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.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.data.proto.ProtoMessageManager; +import net.sopod.soim.data.serialize.ImMessage; import net.sopod.soim.data.serialize.ImMessageCodec; import net.sopod.soim.data.msg.auth.Auth; import net.sopod.soim.data.msg.chat.Chat; +import org.apache.commons.lang3.tuple.Pair; import org.slf4j.LoggerFactory; import java.util.concurrent.CompletableFuture; @@ -60,18 +68,23 @@ public class SoImSession { @Override protected void initChannel(SocketChannel ch) throws Exception { ch.pipeline() + .addLast(new LoggingHandler(LogLevel.INFO)) .addLast(new Varint32FrameCodec()) - .addLast(new ImMessageCodec.ProtoMsg2ImMessageEncoder()) +// .addLast(new ImMessageCodec()) .addLast(new ImMessageCodec.ImMessage2ByteEncoder()) - .addLast(new ImMessageCodec.ImMessageDecoder()) - .addLast(messageDispatcher); +// .addLast(new ImMessageCodec.ImMessageDecoder()) + .addLast(new ImMessageInboundHandler()); + // .addLast(messageDispatcher); } }); try { Logger.info("连接中..."); clientChannel = b.connect(host, port).await().channel(); // 连接后立即发送认证消息,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; } catch (Exception e) { Logger.error("连接服务器失败: {}", e.getMessage()); @@ -92,21 +105,37 @@ public class SoImSession { 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> futurePair = MessageQueueHolder.getInstance().nextSerialNo(); + imMessage.setSerialNo(futurePair.getLeft()); + + clientChannel.writeAndFlush(imMessage); + return futurePair.getRight(); + } + /** * 发送消息 */ + @SuppressWarnings("unchecked") public CompletableFuture send(MessageLite message) { if (!auth.get()) { Logger.error("请先登录"); return null; } - if (clientChannel == null - || !clientChannel.isActive()) { - Logger.error("连接已关闭"); - return null; - } - clientChannel.writeAndFlush(message); - return new CompletableFuture<>(); + return (CompletableFuture) send0(message); } public void textChat(String receiverName, String message) { diff --git a/im-common/src/main/java/net/sopod/soim/common/dubbo/exception/ConvertException.java b/im-common/src/main/java/net/sopod/soim/common/dubbo/exception/ConvertException.java new file mode 100644 index 0000000..09f2b38 --- /dev/null +++ b/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); + } + +} diff --git a/im-common/src/main/java/net/sopod/soim/common/exception/ConvertException.java b/im-common/src/main/java/net/sopod/soim/common/exception/ConvertException.java deleted file mode 100644 index 7f0b9ea..0000000 --- a/im-common/src/main/java/net/sopod/soim/common/exception/ConvertException.java +++ /dev/null @@ -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; - -} diff --git a/im-common/src/main/java/net/sopod/soim/common/res/R.java b/im-common/src/main/java/net/sopod/soim/common/res/R.java deleted file mode 100644 index 7dd6b23..0000000 --- a/im-common/src/main/java/net/sopod/soim/common/res/R.java +++ /dev/null @@ -1,11 +0,0 @@ -package net.sopod.soim.common.res; - -/** - * R - * - * @author tmy - * @date 2022-06-02 15:31 - */ -public class R { - -} diff --git a/im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/config/ChatPersistentRabbitMQConfiguration.java b/im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/config/ChatPersistentRabbitMQConfiguration.java index 5e057bf..095c9d2 100644 --- a/im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/config/ChatPersistentRabbitMQConfiguration.java +++ b/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.ChatQueueType; +import net.sopod.soim.das.user.api.service.DasMQPersistentService; import org.springframework.amqp.core.Queue; import org.springframework.amqp.rabbit.connection.CachingConnectionFactory; import org.springframework.amqp.rabbit.core.RabbitTemplate; import org.springframework.amqp.support.converter.MessageConverter; import org.springframework.beans.factory.support.BeanDefinitionRegistry; 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.Configuration; -import org.springframework.context.annotation.Import; import org.springframework.context.annotation.ImportBeanDefinitionRegistrar; import org.springframework.core.type.AnnotationMetadata; @@ -53,4 +52,9 @@ public class ChatPersistentRabbitMQConfiguration implements ImportBeanDefinition return rabbitTemplate; } + @Bean + public DasMQPersistentService dasPersistentService(RabbitTemplate rabbitTemplate) { + return new DasMQPersistentService(rabbitTemplate); + } + } diff --git a/im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/mq/ChatMQMessageConverter.java b/im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/mq/ChatMQMessageConverter.java index 67a3ee9..106cd66 100644 --- a/im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/mq/ChatMQMessageConverter.java +++ b/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(); Object ordinal = messageProperties.getHeaders().get(MQ_MSG_TYPE); Object isSnappy = messageProperties.getHeaders().get(MQ_MSG_IS_SNAPPY); - ChatQueueType chatQueueType = ChatQueueType.getMQType(Converter.toInt(ordinal), - ori -> new MessageConversionException(String.format("没有ordinal为%d的消息类型,在MQType添加", ori))); + String consumerQueue = messageProperties.getConsumerQueue(); + // 获取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(); if (Boolean.TRUE.equals(isSnappy)) { try { diff --git a/im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/mq/ChatQueueType.java b/im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/mq/ChatQueueType.java index 0454a24..bfb4ad5 100644 --- a/im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/mq/ChatQueueType.java +++ b/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); } + public static ChatQueueType getMQType(String queueName, Function exceptionProvider) { + ChatQueueType[] values = values(); + for (ChatQueueType queueType : values) { + if (queueType.getQueueName().equals(queueName)) { + return queueType; + } + } + throw exceptionProvider.apply(queueName); + } + } diff --git a/im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/service/DasMQPersistentService.java b/im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/service/DasMQPersistentService.java new file mode 100644 index 0000000..90caedd --- /dev/null +++ b/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); + } + +} diff --git a/im-das/im-das-user/src/main/java/net/sopod/soim/das/user/amqp/Sender.java b/im-das/im-das-user/src/main/java/net/sopod/soim/das/user/amqp/Sender.java index 7dbff13..368881d 100644 --- a/im-das/im-das-user/src/main/java/net/sopod/soim/das/user/amqp/Sender.java +++ b/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.mq.ChatQueueType; 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.boot.context.event.ApplicationReadyEvent; import org.springframework.context.ApplicationListener; import org.springframework.context.ConfigurableApplicationContext; -import org.springframework.stereotype.Component; -import java.util.Random; import java.util.concurrent.TimeUnit; /** @@ -23,7 +20,7 @@ import java.util.concurrent.TimeUnit; * @author tmy * @date 2022-05-28 16:23 */ -@Component +//@Component @AllArgsConstructor public class Sender implements ApplicationListener { diff --git a/im-das/im-das-user/src/main/java/net/sopod/soim/das/user/service/FriendDasImpl.java b/im-das/im-das-user/src/main/java/net/sopod/soim/das/user/service/FriendDasImpl.java index 33d1efa..06f147b 100644 --- a/im-das/im-das-user/src/main/java/net/sopod/soim/das/user/service/FriendDasImpl.java +++ b/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.QueryWrapper; import lombok.AllArgsConstructor; -import net.sopod.soim.common.constant.LogicConsts; import net.sopod.soim.common.util.Collects; import net.sopod.soim.common.util.ImClock; import net.sopod.soim.das.user.api.config.LogicTables; diff --git a/im-entry/src/main/java/net/sopod/soim/entry/delay/NetUserDelayTaskManager.java b/im-entry/src/main/java/net/sopod/soim/entry/delay/NetUserDelayTaskManager.java index d51607e..e6daa95 100644 --- a/im-entry/src/main/java/net/sopod/soim/entry/delay/NetUserDelayTaskManager.java +++ b/im-entry/src/main/java/net/sopod/soim/entry/delay/NetUserDelayTaskManager.java @@ -41,7 +41,7 @@ public class NetUserDelayTaskManager { if (netUser.isActive()) { // 延时任务执行时 NetUser 可能已升级为 Account NetUser curNetUser = netUser.channel().attr(NetUser.NET_USER_KEY).get(); - ProtoMessageDispatcher.dispatch(curNetUser, delayTask.getTaskMsg()); + ProtoMessageDispatcher.dispatch(null, curNetUser, delayTask.getTaskMsg()); } else { logger.info("delayTask netUser inactive, task cancel!"); } diff --git a/im-entry/src/main/java/net/sopod/soim/entry/server/ImEntryInitializer.java b/im-entry/src/main/java/net/sopod/soim/entry/server/ImEntryInitializer.java index 14eb834..501fe9b 100644 --- a/im-entry/src/main/java/net/sopod/soim/entry/server/ImEntryInitializer.java +++ b/im-entry/src/main/java/net/sopod/soim/entry/server/ImEntryInitializer.java @@ -27,7 +27,8 @@ public class ImEntryInitializer extends ChannelInitializer { ChannelPipeline pipeline = socketChannel.pipeline(); pipeline.addLast(new LoggingHandler(logLevel)) .addLast(new Varint32FrameCodec()) - .addLast(new ImMessageCodec()) + .addLast(new ImMessageCodec.ImMessageDecoder()) + .addLast(new ImMessageCodec.ImMessage2ByteEncoder()) .addLast(new InboundImMessageHandler()); } diff --git a/im-entry/src/main/java/net/sopod/soim/entry/server/InboundImMessageHandler.java b/im-entry/src/main/java/net/sopod/soim/entry/server/InboundImMessageHandler.java index 25e1389..2af8b24 100644 --- a/im-entry/src/main/java/net/sopod/soim/entry/server/InboundImMessageHandler.java +++ b/im-entry/src/main/java/net/sopod/soim/entry/server/InboundImMessageHandler.java @@ -1,11 +1,12 @@ package net.sopod.soim.entry.server; -import com.google.protobuf.MessageLite; import io.netty.channel.Channel; import io.netty.channel.ChannelHandlerContext; import io.netty.channel.SimpleChannelInboundHandler; import io.netty.util.Attribute; 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.data.msg.task.Tasks; import net.sopod.soim.entry.delay.NetUserDelayTaskManager; @@ -19,7 +20,7 @@ import java.util.concurrent.atomic.AtomicInteger; * @author tmy * @date 2022-04-10 22:40 */ -public class InboundImMessageHandler extends SimpleChannelInboundHandler { +public class InboundImMessageHandler extends SimpleChannelInboundHandler { /** * channel 建立连接,设置初始属性 @@ -50,11 +51,12 @@ public class InboundImMessageHandler extends SimpleChannelInboundHandler netUserAttr = ctx.channel().attr(NetUser.NET_USER_KEY); NetUser netUser = netUserAttr.get(); + ImContext execCtx = new ImContext(imMessage.getSerialNo()); // TODO dispatcher - ProtoMessageDispatcher.dispatch(netUser, messageLite); + ProtoMessageDispatcher.dispatch(execCtx, netUser, imMessage.getDecodeBody()); } } diff --git a/im-entry/src/main/java/net/sopod/soim/entry/server/ProtoMessageDispatcher.java b/im-entry/src/main/java/net/sopod/soim/entry/server/ProtoMessageDispatcher.java index fbf54a1..b0ecfd7 100644 --- a/im-entry/src/main/java/net/sopod/soim/entry/server/ProtoMessageDispatcher.java +++ b/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 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.registry.ProtoMessageHandlerRegistry; import net.sopod.soim.entry.server.session.Account; @@ -18,20 +20,20 @@ import net.sopod.soim.entry.worker.WorkerGroup; */ public class ProtoMessageDispatcher { - public static void dispatch(NetUser netUser, MessageLite message) { + public static void dispatch(ImContext ctx, NetUser netUser, MessageLite message) { MessageHandler typeHandler = ProtoMessageHandlerRegistry .getTypeHandler(message.getClass()); 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.execute(() -> { if (netUser.isAccount()) { - Account account = (Account)netUser; + Account account = (Account) netUser; MessageHandlerContext.setAttribute(DubboConstant.CTX_UID, String.valueOf(account.getUid())); } try { - typeHandler.exec(netUser, message); + typeHandler.exec(ctx, netUser, message); } finally { MessageHandlerContext.remove(); } diff --git a/im-entry/src/main/java/net/sopod/soim/entry/server/handler/AccountMessageHandler.java b/im-entry/src/main/java/net/sopod/soim/entry/server/handler/AccountMessageHandler.java index 12e1d31..b77ca6b 100644 --- a/im-entry/src/main/java/net/sopod/soim/entry/server/handler/AccountMessageHandler.java +++ b/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 implements MessageHandler { @Override - public final void exec(NetUser netUser, T req) { + public final void exec(ImContext ctx, NetUser netUser, T req) { if (!netUser.isAccount()) { throw new IllegalStateException("NetUser is not account!" + netUser); } MessageLite res = handle((Account) netUser, req); if (res != null) { - netUser.writeNow(res); + netUser.writeNow(ctx, res); } } diff --git a/im-entry/src/main/java/net/sopod/soim/entry/server/handler/ImContext.java b/im-entry/src/main/java/net/sopod/soim/entry/server/handler/ImContext.java new file mode 100644 index 0000000..f9028d1 --- /dev/null +++ b/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; + } + +} diff --git a/im-entry/src/main/java/net/sopod/soim/entry/server/handler/MessageHandler.java b/im-entry/src/main/java/net/sopod/soim/entry/server/handler/MessageHandler.java index d2242e5..5610e7a 100644 --- a/im-entry/src/main/java/net/sopod/soim/entry/server/handler/MessageHandler.java +++ b/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 { - void exec(NetUser netUser, T msg); + void exec(ImContext ctx, NetUser netUser, T msg); } diff --git a/im-entry/src/main/java/net/sopod/soim/entry/server/handler/NetUserMessageHandler.java b/im-entry/src/main/java/net/sopod/soim/entry/server/handler/NetUserMessageHandler.java index cbef33a..0c4c8e5 100644 --- a/im-entry/src/main/java/net/sopod/soim/entry/server/handler/NetUserMessageHandler.java +++ b/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 implements MessageHandler { @Override - public final void exec(NetUser netUser, T msg) { + public final void exec(ImContext ctx, NetUser netUser, T msg) { MessageLite res = handle(netUser, msg); if (res != null) { - netUser.writeNow(res); + netUser.writeNow(ctx, res); } } diff --git a/im-entry/src/main/java/net/sopod/soim/entry/server/session/NetUser.java b/im-entry/src/main/java/net/sopod/soim/entry/server/session/NetUser.java index acb66ad..ed589cf 100644 --- a/im-entry/src/main/java/net/sopod/soim/entry/server/session/NetUser.java +++ b/im-entry/src/main/java/net/sopod/soim/entry/server/session/NetUser.java @@ -1,12 +1,20 @@ package net.sopod.soim.entry.server.session; +import com.google.protobuf.MessageLite; import io.netty.channel.Channel; 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; public class NetUser { + private static final Logger logger = LoggerFactory.getLogger(NetUser.class); + /** channel 绑定 netUser 对象 */ public static final AttributeKey NET_USER_KEY = AttributeKey.valueOf(NetUser.class, "NET_USER"); @@ -39,24 +47,24 @@ public class NetUser { } } - public void write(Object message) { - write(message, false); + public void writeNow(MessageLite message) { + write0(null, message); } - public void writeNow(Object message) { - write(message, true); + public void writeNow(ImContext ctx, MessageLite message) { + write0(ctx.getSerialNo(), message); } - private void write(Object message, boolean now) { + private void write0(Integer serialNo, MessageLite message) { Channel channel = this.channel.get(); if (channel == null) { + logger.error("channel closed drop message: {}, {}", serialNo, message); return; } - if (now) { - channel.writeAndFlush(message); - } else { - channel.write(message); - } + // 构建 ImMessage + ImMessage imMessage = ImMessageCodec.encodeImProto(message); + imMessage.setSerialNo(serialNo == null ? -1 : serialNo); + channel.writeAndFlush(imMessage); } @Override diff --git a/im-service-api/im-entry-protocol/src/main/java/net/sopod/soim/data/proto/ProtoMessageManager.java b/im-service-api/im-entry-protocol/src/main/java/net/sopod/soim/data/proto/ProtoMessageManager.java index abf6f22..ed79427 100644 --- a/im-service-api/im-entry-protocol/src/main/java/net/sopod/soim/data/proto/ProtoMessageManager.java +++ b/im-service-api/im-entry-protocol/src/main/java/net/sopod/soim/data/proto/ProtoMessageManager.java @@ -60,7 +60,7 @@ public class ProtoMessageManager { } @Nullable - public static MessageLite getProtoInstance(Integer serialNo) { + public static MessageLite getDefaultInstance(Integer serialNo) { String clazz = serialNoTypeMap.get(serialNo); if (clazz == null) { logger.error("protoMsgDict serialNo proto class not found: {}", serialNo); diff --git a/im-service-api/im-entry-protocol/src/main/java/net/sopod/soim/data/serialize/ImMessageCodec.java b/im-service-api/im-entry-protocol/src/main/java/net/sopod/soim/data/serialize/ImMessageCodec.java index 10110f4..29de4ee 100644 --- a/im-service-api/im-entry-protocol/src/main/java/net/sopod/soim/data/serialize/ImMessageCodec.java +++ b/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; +import com.google.protobuf.InvalidProtocolBufferException; import com.google.protobuf.MessageLite; import io.netty.buffer.ByteBuf; import io.netty.channel.ChannelHandlerContext; -import io.netty.channel.CombinedChannelDuplexHandler; import io.netty.handler.codec.ByteToMessageDecoder; 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 org.slf4j.Logger; import org.slf4j.LoggerFactory; -import java.lang.reflect.Type; -import java.util.Arrays; import java.util.List; /** @@ -21,90 +20,57 @@ import java.util.List; * @author tmy * @date 2022-03-28 11:29 */ -public class ImMessageCodec //extends MessageToMessageCodec { - extends CombinedChannelDuplexHandler { +public class ImMessageCodec { private static final Logger logger = LoggerFactory.getLogger(ImMessageCodec.class); - public ImMessageCodec() { - super(new ProtoMsgDecoder(), new ProtoMsgEncoder()); - } - public static class ImMessageDecoder extends ByteToMessageDecoder { @Override protected void decode(ChannelHandlerContext ctx, ByteBuf byteBuf, List out) throws Exception { - ImMessage message = ImMessage.read(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(); - byte[] protoByte = message.getBody(); - MessageLite protoClass = ProtoMessageManager.getProtoInstance(serviceNo); - MessageLite protoMsg = protoClass.getParserForType().parseFrom(protoByte); - message.setDecodeBody(protoMsg); - out.add(message); + ImMessage imMessage = ImMessageCodec.decodeImMessage(byteBuf); + out.add(imMessage); } } - public static class ProtoMsgDecoder extends ByteToMessageDecoder { + public static class ImMessage2ByteEncoder extends MessageToByteEncoder { @Override - protected void decode(ChannelHandlerContext ctx, ByteBuf byteBuf, List list) throws Exception { - ImMessage message = ImMessage.read(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); + protected void encode(ChannelHandlerContext ctx, ImMessage imMessage, ByteBuf byteBuf) throws Exception { + imMessage.write(byteBuf); } } - public static class ProtoMsg2ImMessageEncoder extends MessageToMessageEncoder { - @Override - protected void encode(ChannelHandlerContext ctx, MessageLite message, List out) throws Exception { - Integer serialNo = ProtoMessageManager.getSerialNo(message.getClass()); - // TODO unknow class serialNo - ImMessage imMessage = new ImMessage() - .setServiceNo(serialNo) - .setBody(message.toByteArray()); - out.add(imMessage); + public static ImMessage decodeImMessage(ByteBuf byteBuf) { + ImMessage imMessage = ImMessage.read(byteBuf); + boolean isMagicError; + if ((isMagicError = (imMessage == ImMessage.MAGIC_ERROR)) + || imMessage == ImMessage.PROTOCOL_ERROR) { + throw new ConvertException("ImMessage解码错误"); } - } - - public static class ImMessage2ByteEncoder extends MessageToByteEncoder { - @Override - protected void encode(ChannelHandlerContext ctx, ImMessage imMessage, ByteBuf out) throws Exception { - imMessage.write(out); + // 解码 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; + try { + protoMsg = protoClass.getParserForType().parseFrom(protoByte); + } catch (InvalidProtocolBufferException e) { + throw new ConvertException("ImMessage消息体Protobuf解码错误: " + e.getMessage()); + } + imMessage.setDecodeBody(protoMsg); + return imMessage; } - public static class ProtoMsgEncoder extends MessageToByteEncoder { - @Override - protected void encode(ChannelHandlerContext ctx, MessageLite message, ByteBuf byteBuf) throws Exception { - Integer serialNo = ProtoMessageManager.getSerialNo(message.getClass()); - // TODO unknow class serialNo - ImMessage imMessage = new ImMessage() - .setServiceNo(serialNo) - .setBody(message.toByteArray()); - imMessage.write(byteBuf); + public static ImMessage encodeImProto(MessageLite message) { + Integer serialNo = ProtoMessageManager.getSerialNo(message.getClass()); + if (serialNo == null) { + throw new SoimException("未知的Protobuf消息类型:" + message.getClass()); } + return new ImMessage() + .setServiceNo(serialNo) + .setBody(message.toByteArray()); } } 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 55751ee..a776a9b 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 @@ -1,14 +1,12 @@ package net.sopod.soim.router.service; 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.StringUtil; 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.ImUser; -import net.sopod.soim.das.user.api.mq.ChatQueue; -import net.sopod.soim.das.user.api.mq.ChatQueueType; +import net.sopod.soim.das.user.api.service.DasMQPersistentService; import net.sopod.soim.das.user.api.service.FriendDas; import net.sopod.soim.das.user.api.service.UserDas; 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.slf4j.Logger; import org.slf4j.LoggerFactory; -import org.springframework.amqp.rabbit.core.RabbitTemplate; import javax.annotation.Resource; import java.util.ArrayList; @@ -62,7 +59,7 @@ public class UserRouteServiceImpl implements UserRouteService { private SegmentIdGenerator segmentIdGenerator; @Resource - private RabbitTemplate rabbitTemplate; + private DasMQPersistentService dasMQPersistentService; @Override public RegistryRes registryUserEntry(Long uid, String imEntryAddr) { @@ -107,7 +104,7 @@ public class UserRouteServiceImpl implements UserRouteService { .setSender(textChat.getUid()) .setReceiver(textChat.getReceiverUid()) .setCreateTime(ImClock.date()); - rabbitTemplate.convertAndSend(ChatQueueType.IM_MESSAGE.getQueueName(), imMessage); + dasMQPersistentService.saveImMessage(imMessage); RpcContextUtil.setContextUid(textChat.getReceiverUid()); Boolean send = textChatService.sendTextChat(textChat);