Browse Source

msg dispatcher

master
tangmingyou 4 years ago
parent
commit
7152d7ae9e
  1. 25
      im-client/src/main/java/net/sopod/soim/client/net/ImNetClient.java
  2. 39
      im-common/src/main/java/net/sopod/soim/common/util/Reflects.java
  3. 5
      im-core/pom.xml
  4. 28
      im-core/src/main/java/net/sopod/soim/core/handler/AccountMessageHandler.java
  5. 15
      im-core/src/main/java/net/sopod/soim/core/handler/MessageHandler.java
  6. 24
      im-core/src/main/java/net/sopod/soim/core/handler/NetUserMessageHandler.java
  7. 26
      im-core/src/main/java/net/sopod/soim/core/handler/ProtoMessageHandler.java
  8. 21
      im-core/src/main/java/net/sopod/soim/core/net/AttributeKeys.java
  9. 52
      im-core/src/main/java/net/sopod/soim/core/net/ImEntryCodec.java
  10. 76
      im-core/src/main/java/net/sopod/soim/core/net/ImMessageCodec.java
  11. 28
      im-core/src/main/java/net/sopod/soim/core/net/ImMessageHandler.java
  12. 75
      im-core/src/main/java/net/sopod/soim/core/registry/ProtoMessageHandlerRegistry.java
  13. 41
      im-core/src/main/java/net/sopod/soim/core/registry/ServiceRegistry.java
  14. 13
      im-core/src/main/java/net/sopod/soim/core/service/ReqHandler.java
  15. 23
      im-core/src/main/java/net/sopod/soim/core/session/Account.java
  16. 50
      im-core/src/main/java/net/sopod/soim/core/session/NetUser.java
  17. 17
      im-data/src/main/java/net/sopod/soim/data/GenProtobuf.java
  18. 12
      im-data/src/main/java/net/sopod/soim/data/proto/ProtoMessageHolder.java
  19. 100
      im-data/src/main/java/net/sopod/soim/data/proto/ProtoMessageManager.java
  20. 25
      im-data/src/main/java/net/sopod/soim/data/serialize/ImMessage.java
  21. 4
      im-data/src/main/resources/protoSerialNoTable.txt
  22. 24
      im-entry/src/main/java/net/sopod/soim/entry/config/ApplicationContextInitialed.java
  23. 26
      im-entry/src/main/java/net/sopod/soim/entry/handler/HelloHandler.java
  24. 16
      im-entry/src/main/java/net/sopod/soim/entry/model/A.java
  25. 5
      im-entry/src/main/java/net/sopod/soim/entry/server/EntryServer.java
  26. 7
      im-entry/src/main/java/net/sopod/soim/entry/server/EntryServerRunner.java
  27. 7
      im-entry/src/main/java/net/sopod/soim/entry/server/ImEntryInitializer.java
  28. 53
      im-entry/src/main/java/net/sopod/soim/entry/server/InboundImMessageHandler.java
  29. 28
      im-entry/src/main/java/net/sopod/soim/entry/util/FastThreadLocalThreadFactory.java

25
im-client/src/main/java/net/sopod/soim/client/net/ImNetClient.java

@ -6,12 +6,9 @@ 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 net.sopod.soim.common.util.Jackson;
import net.sopod.soim.core.net.ImEntryCodec;
import net.sopod.soim.data.constant.SerializeType;
import net.sopod.soim.data.serialize.ImMessage;
import net.sopod.soim.core.net.ImMessageCodec;
import net.sopod.soim.data.msg.hello.HelloPB;
import java.nio.charset.StandardCharsets;
import java.util.HashMap;
import java.util.Map;
@ -32,7 +29,8 @@ public class ImNetClient {
@Override
protected void initChannel(SocketChannel ch) throws Exception {
ch.pipeline()
.addLast(new ImEntryCodec());
.addLast(new ImMessageCodec());
// .addLast(new ProtoMessageCodec());
}
});
Channel channel = b.connect(host, port).await().channel();
@ -40,11 +38,16 @@ public class ImNetClient {
body.put("name", "二狗子");
body.put("age", 16);
body.put("birthday", "2002-04-19");
ImMessage imMessage = new ImMessage()
.setServiceNo(1)
.setSerializeType(SerializeType.json.ordinal())
.setBody(Jackson.json().serialize(body).getBytes(StandardCharsets.UTF_8));
channel.writeAndFlush(imMessage);
HelloPB.Hello hello = HelloPB.Hello.newBuilder()
.setId(1)
.setStr("手")
.build();
// ImMessage imMessage = new ImMessage()
// .setServiceNo(ProtoMessageManager.getSerialNo(hello.getClass()))
// .setSerializeType(SerializeType.json.ordinal())
// //.setBody(Jackson.json().serialize(body).getBytes(StandardCharsets.UTF_8));
// .setBody(hello.toByteArray());
channel.writeAndFlush(hello);
channel.close();
work.shutdownGracefully();
}

39
im-common/src/main/java/net/sopod/soim/common/util/Reflects.java

@ -0,0 +1,39 @@
package net.sopod.soim.common.util;
import java.lang.reflect.Type;
import java.util.Arrays;
import java.util.Collections;
import java.util.List;
/**
* Reflects
*
* @author tmy
* @date 2022-04-10 23:54
*/
public class Reflects {
/**
* 获取父类上的泛型
* @return 父类上的泛型
*/
public static List<String> getSuperclassGenericTypes(Class<?> clazz) {
// 获取 handler 的泛型消息
Type superType = clazz.getGenericSuperclass();
String typeName = superType.getTypeName();
int idx = typeName.indexOf('<');
if (idx == -1) {
// 父类没有泛型
return Collections.emptyList();
}
String genericName = typeName.substring(idx + 1, typeName.length() - 1);
// 父类只有一个泛型
if (!genericName.contains(",")) {
return Collections.singletonList(genericName);
}
// 父类有多个泛型
String[] genericNames = genericName.split(", ");
return Arrays.asList(genericNames);
}
}

5
im-core/pom.xml

@ -26,6 +26,11 @@
<groupId>io.netty</groupId>
<artifactId>netty-all</artifactId>
</dependency>
<dependency>
<groupId>org.springframework</groupId>
<artifactId>spring-context</artifactId>
<scope>provided</scope>
</dependency>
</dependencies>
</project>

28
im-core/src/main/java/net/sopod/soim/core/handler/AccountMessageHandler.java

@ -0,0 +1,28 @@
package net.sopod.soim.core.handler;
import com.google.protobuf.MessageLite;
import net.sopod.soim.core.session.Account;
import net.sopod.soim.core.session.NetUser;
/**
* AccountMessageHandler
*
* @author tmy
* @date 2022-04-10 23:41
*/
public abstract class AccountMessageHandler<T> implements MessageHandler<T> {
@Override
public final void exec(NetUser netUser, T msg) {
if (!netUser.isAccount()) {
throw new IllegalStateException("NetUser is not account!" + netUser);
}
MessageLite res = handle((Account) netUser, msg);
if (res != null) {
netUser.write(res);
}
}
public abstract MessageLite handle(Account account, T msg);
}

15
im-core/src/main/java/net/sopod/soim/core/handler/MessageHandler.java

@ -0,0 +1,15 @@
package net.sopod.soim.core.handler;
import net.sopod.soim.core.session.NetUser;
/**
* MessageHandler
*
* @author tmy
* @date 2022-04-10 23:40
*/
public interface MessageHandler<T> {
void exec(NetUser netUser, T msg);
}

24
im-core/src/main/java/net/sopod/soim/core/handler/NetUserMessageHandler.java

@ -0,0 +1,24 @@
package net.sopod.soim.core.handler;
import com.google.protobuf.MessageLite;
import net.sopod.soim.core.session.NetUser;
/**
* NetUserMessageHandler
*
* @author tmy
* @date 2022-04-10 23:40
*/
public abstract class NetUserMessageHandler<T> implements MessageHandler<T> {
@Override
public final void exec(NetUser netUser, T msg) {
MessageLite res = handle(netUser, msg);
if (res != null) {
netUser.write(res);
}
}
public abstract MessageLite handle(NetUser netUser, T msg);
}

26
im-core/src/main/java/net/sopod/soim/core/handler/ProtoMessageHandler.java

@ -0,0 +1,26 @@
package net.sopod.soim.core.handler;
import com.google.protobuf.MessageLite;
import net.sopod.soim.core.session.NetUser;
/**
* ProtoMessageHandler
* <T extends MessageLite>
*
* @author tmy
* @date 2022-04-10 19:19
*/
public abstract class ProtoMessageHandler<T> {
public final void exec(NetUser netUser, T msg) {
MessageLite res = handle(msg);
if (res != null) {
netUser.write(res);
}
}
public abstract Class<T> type();
public abstract MessageLite handle(T msg);
}

21
im-core/src/main/java/net/sopod/soim/core/net/AttributeKeys.java

@ -0,0 +1,21 @@
package net.sopod.soim.core.net;
import io.netty.util.AttributeKey;
import java.util.concurrent.atomic.AtomicInteger;
/**
* AttributeKeys
*
* @author tmy
* @date 2022-04-10 23:26
*/
public interface AttributeKeys {
/** channel 写失败次数 */
AttributeKey<AtomicInteger> WRITE_FAIL_TIMES = AttributeKey.valueOf("WRITE_FAIL_TIMES");
/** channel 登录失败次数 */
AttributeKey<AtomicInteger> LOGIN_FAIL_TIMES = AttributeKey.valueOf("LOGIN_FAIL_TIMES");
}

52
im-core/src/main/java/net/sopod/soim/core/net/ImEntryCodec.java

@ -1,52 +0,0 @@
package net.sopod.soim.core.net;
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 net.sopod.soim.data.serialize.ImMessage;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.List;
/**
* ImEntryCodec
*
* @author tmy
* @date 2022-03-28 11:29
*/
public class ImEntryCodec extends CombinedChannelDuplexHandler<ImEntryCodec.ImDecoder, ImEntryCodec.ImEncoder> {
private static final Logger logger = LoggerFactory.getLogger(ImEntryCodec.class);
public ImEntryCodec() {
super(new ImDecoder(), new ImEncoder());
}
public static class ImDecoder extends ByteToMessageDecoder {
@Override
protected void decode(ChannelHandlerContext ctx, ByteBuf byteBuf, List<Object> 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;
}
list.add(message);
}
}
public static class ImEncoder extends MessageToByteEncoder<ImMessage> {
@Override
protected void encode(ChannelHandlerContext ctx, ImMessage imMessage, ByteBuf byteBuf) throws Exception {
imMessage.write(byteBuf);
}
}
}

76
im-core/src/main/java/net/sopod/soim/core/net/ImMessageCodec.java

@ -0,0 +1,76 @@
package net.sopod.soim.core.net;
import com.google.protobuf.MessageLite;
import io.netty.buffer.ByteBuf;
import io.netty.buffer.Unpooled;
import io.netty.channel.ChannelHandlerContext;
import io.netty.channel.CombinedChannelDuplexHandler;
import io.netty.handler.codec.ByteToMessageDecoder;
import io.netty.handler.codec.MessageToByteEncoder;
import net.sopod.soim.data.proto.ProtoMessageManager;
import net.sopod.soim.data.serialize.ImMessage;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.lang.reflect.Type;
import java.util.Arrays;
import java.util.List;
/**
* ImEntryCodec
*
* @author tmy
* @date 2022-03-28 11:29
*/
public class ImMessageCodec //extends MessageToMessageCodec<ByteBuf, MessageLite> {
extends CombinedChannelDuplexHandler<ImMessageCodec.ProtoMsgDecoder, ImMessageCodec.ProtoMsgEncoder> {
public static void main(String[] args) throws ClassNotFoundException {
Type superType = ImMessageCodec.class.getGenericSuperclass();
String typeName = superType.getTypeName();
int idx = typeName.indexOf('<');
String genericName = typeName.substring(idx + 1, typeName.length() - 1);
System.out.println(genericName.trim());
System.out.println(Arrays.toString(genericName.split(", ")));
}
private static final Logger logger = LoggerFactory.getLogger(ImMessageCodec.class);
public ImMessageCodec() {
super(new ProtoMsgDecoder(), new ProtoMsgEncoder());
}
public static class ProtoMsgDecoder extends ByteToMessageDecoder {
@Override
protected void decode(ChannelHandlerContext ctx, ByteBuf byteBuf, List<Object> 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();
byte[] protoByte = message.getBody();
MessageLite protoClass = ProtoMessageManager.getProtoInstance(serviceNo);
MessageLite protoMsg = protoClass.getParserForType().parseFrom(protoByte);
list.add(protoMsg);
}
}
public static class ProtoMsgEncoder extends MessageToByteEncoder<MessageLite> {
@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);
}
}
}

28
im-core/src/main/java/net/sopod/soim/core/net/ImMessageHandler.java

@ -1,28 +0,0 @@
package net.sopod.soim.core.net;
import io.netty.channel.ChannelHandlerContext;
import io.netty.channel.SimpleChannelInboundHandler;
import net.sopod.soim.data.constant.SerializeType;
import net.sopod.soim.data.serialize.ImMessage;
import java.util.Map;
/**
* ImMessageHandler
*
* @author tmy
* @date 2022-03-28 13:27
*/
public class ImMessageHandler extends SimpleChannelInboundHandler<ImMessage> {
@Override
protected void channelRead0(ChannelHandlerContext channelHandlerContext, ImMessage imMessage) throws Exception {
byte[] body = imMessage.getBody();
int serviceNo = imMessage.getServiceNo();
SerializeType serialize = SerializeType.getSerializeByOrdinal(imMessage.getSerializeType());
Map data = serialize.getSerializer().deserialize(body, Map.class);
System.out.println(data);
}
}

75
im-core/src/main/java/net/sopod/soim/core/registry/ProtoMessageHandlerRegistry.java

@ -0,0 +1,75 @@
package net.sopod.soim.core.registry;
import net.sopod.soim.common.util.ImClock;
import net.sopod.soim.common.util.Reflects;
import net.sopod.soim.core.handler.MessageHandler;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.context.ApplicationContext;
import javax.annotation.Nullable;
import java.util.Collection;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.concurrent.CountDownLatch;
/**
* ProtoMessageDispatcher
* implements ApplicationContextAware
*
* @author tmy
* @date 2022-04-10 19:15
*/
public class ProtoMessageHandlerRegistry {
private static final Logger logger = LoggerFactory.getLogger(ProtoMessageHandlerRegistry.class);
private static final Map<Class<?>, MessageHandler<?>> TYPE_HANDLER_MAP = new HashMap<>(32);
private static final CountDownLatch CONTEXT_AWARE_AWAIT = new CountDownLatch(1);
/**
* spring ioc 容器中获取 msgType handler
* @param context spring ioc 上下文
*/
public static synchronized void registerHandlerWithApplicationContext(ApplicationContext context) {
if (CONTEXT_AWARE_AWAIT.getCount() <= 0) {
throw new IllegalStateException("proto message registry already initialed!");
}
logger.info("proto message registry initial...");
long start = ImClock.millis();
Map<String, MessageHandler> beansOfType = context.getBeansOfType(MessageHandler.class);
Collection<MessageHandler> handlers = beansOfType.values();
for (MessageHandler<?> handler : handlers) {
// 获取 handler 泛型
List<String> genericTypes = Reflects.getSuperclassGenericTypes(handler.getClass());
try {
Class<?> type = genericTypes.size() == 0 ? Object.class : Class.forName(genericTypes.get(0));
MessageHandler<?> existHandler = TYPE_HANDLER_MAP.putIfAbsent(type, handler);
if (existHandler != null) {
// 消息类型有重复的 handler!
throw new IllegalStateException("msg type " + type + " handler duplicate; " +
"[" + existHandler.getClass() + "] and [" + handler.getClass() + "]");
}
} catch (ClassNotFoundException e) {
logger.warn("handler msgType class not found!", e);
}
}
CONTEXT_AWARE_AWAIT.countDown();
logger.info("proto message registry complete, {} handlers at {}ms.", handlers.size(), ImClock.millis() - start);
}
@Nullable
public static <T> MessageHandler<T> getTypeHandler(Class<T> type) {
if (CONTEXT_AWARE_AWAIT.getCount() > 0) {
try {
CONTEXT_AWARE_AWAIT.await();
} catch (InterruptedException e) {
logger.error("proto message type dispatcher, wait context ready error!", e);
}
}
return (MessageHandler<T>) TYPE_HANDLER_MAP.get(type);
}
}

41
im-core/src/main/java/net/sopod/soim/core/registry/ServiceRegistry.java

@ -1,41 +0,0 @@
package net.sopod.soim.core.registry;
import net.sopod.soim.core.service.ReqHandler;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.atomic.AtomicInteger;
/**
* ServiceRegistry
*
* @author tmy
* @date 2022-03-28 14:30
*/
public class ServiceRegistry {
private static final ConcurrentHashMap<Integer, Class<?>> serviceIdParamTypeMap;
private static final ConcurrentHashMap<Class<?>, Integer> paramTypeServiceIdMap;
private static final ConcurrentHashMap<Integer, ReqHandler<?>> serviceIdHandlers;
static {
serviceIdParamTypeMap = new ConcurrentHashMap<>();
paramTypeServiceIdMap = new ConcurrentHashMap<>();
serviceIdHandlers = new ConcurrentHashMap<>();
}
private static final AtomicInteger serviceIdGen = new AtomicInteger(10000);
private static <T> void registry(Class<T> paramType, ReqHandler<T> handler) {
int serviceId = serviceIdGen.getAndIncrement();
serviceIdParamTypeMap.put(serviceId, paramType);
paramTypeServiceIdMap.put(paramType, serviceId);
serviceIdHandlers.put(serviceId, handler);
}
public void aaa() {
}
}

13
im-core/src/main/java/net/sopod/soim/core/service/ReqHandler.java

@ -1,13 +0,0 @@
package net.sopod.soim.core.service;
/**
* ReqHandler
*
* @author tmy
* @date 2022-03-28 14:33
*/
public interface ReqHandler<T> {
Object handle(T param);
}

23
im-core/src/main/java/net/sopod/soim/core/session/Account.java

@ -0,0 +1,23 @@
package net.sopod.soim.core.session;
import io.netty.channel.Channel;
import io.netty.util.AttributeKey;
public class Account extends NetUser {
public static final AttributeKey<Account> ACCOUNT_KEY = AttributeKey.valueOf(Account.class, "ACCOUNT");
private long accountId;
private String name;
public Account(Channel channel) {
super(channel);
}
@Override
public boolean isAccount() {
return true;
}
}

50
im-core/src/main/java/net/sopod/soim/core/session/NetUser.java

@ -0,0 +1,50 @@
package net.sopod.soim.core.session;
import io.netty.channel.Channel;
import io.netty.util.AttributeKey;
import java.lang.ref.WeakReference;
public class NetUser {
/** channel 绑定 netUser 对象 */
public static final AttributeKey<NetUser> NET_USER_KEY = AttributeKey.valueOf(NetUser.class, "NET_USER");
private final WeakReference<Channel> channel;
public NetUser(Channel channel) {
this.channel = new WeakReference<>(channel);
}
public boolean isAccount() {
return false;
}
public void write(Object message) {
write(message, false);
}
public void writeNow(Object message) {
write(message, true);
}
private void write(Object message, boolean now) {
Channel channel = this.channel.get();
if (channel == null) {
return;
}
if (now) {
channel.writeAndFlush(message);
} else {
channel.write(message);
}
}
@Override
public String toString() {
return "NetUser{" +
"channel=" + channel +
'}';
}
}

17
im-data/src/main/java/net/sopod/soim/data/GenProtobuf.java

@ -4,15 +4,19 @@ import com.google.protobuf.GeneratedMessageV3;
import java.io.File;
import java.io.FileWriter;
import java.io.IOException;
import java.util.Iterator;
import java.util.Set;
import net.sopod.soim.common.util.ExecUtil;
import net.sopod.soim.data.proto.ProtoMessageHolder;
import net.sopod.soim.data.proto.ProtoMessageManager;
import org.reflections.Reflections;
import org.reflections.scanners.Scanner;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
/**
* TODO DB存储消息版本号唯一
* client 连接时从服务端获取
* TODO 消息类型分组 (user, group, ...) 避免大数据量传输
*/
public class GenProtobuf {
private static final Logger logger = LoggerFactory.getLogger(GenProtobuf.class);
@ -26,7 +30,7 @@ public class GenProtobuf {
/** !!生成新class删除.java文件目录 */
private static final String javaClassDir = "im-data/src/main/java/net/sopod/soim/data/msg";
private static final String msgSerialNoTableFilePath = "./im-data/src/main/resources/" + ProtoMessageHolder.protoSerialNoTableName;
private static final String msgSerialNoTableFilePath = "./im-data/src/main/resources/" + ProtoMessageManager.protoSerialNoTableName;
private static void genProtobufMsgClasses() {
removeOldClass(new File(javaClassDir));
@ -40,9 +44,12 @@ public class GenProtobuf {
String pack = "net.sopod.soim.data.msg";
Reflections collect = new Reflections(pack);
Set<Class<? extends GeneratedMessageV3>> types = collect.getSubTypesOf(GeneratedMessageV3.class);
// TreeSet<String> classNames = new TreeSet<>();
StringBuilder msgTableBuilder = new StringBuilder();
// 已存在不修改, TODO classDict id生成
int idx = 10000;
for (Class<? extends GeneratedMessageV3> type : types) {
// classNames.add(type.getName());
msgTableBuilder.append(type.getName()).append('=').append(idx++).append("\n");
}
return msgTableBuilder;
@ -66,7 +73,7 @@ public class GenProtobuf {
}
}
public static void main(String[] args) throws IOException, ClassNotFoundException {
public static void main(String[] args) throws IOException {
// 从.proto生成.java
genProtobufMsgClasses();
// 将 protobuf java class 编码 id

12
im-data/src/main/java/net/sopod/soim/data/proto/ProtoMessageHolder.java

@ -1,12 +0,0 @@
package net.sopod.soim.data.proto;
/**
* MessageHolder
*
* @author tmy
* @date 2022-04-08 18:02
*/
public class ProtoMessageHolder {
public static final String protoSerialNoTableName = "protoSerialNoTable.txt";
}

100
im-data/src/main/java/net/sopod/soim/data/proto/ProtoMessageManager.java

@ -0,0 +1,100 @@
package net.sopod.soim.data.proto;
import com.google.protobuf.GeneratedMessageV3;
import com.google.protobuf.MessageLite;
import net.sopod.soim.data.msg.hello.HelloPB;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import javax.annotation.Nullable;
import java.io.BufferedReader;
import java.io.IOException;
import java.io.InputStream;
import java.io.InputStreamReader;
import java.lang.reflect.Method;
import java.nio.charset.StandardCharsets;
import java.util.HashMap;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
/**
* MessageHolder
*
* @author tmy
* @date 2022-04-08 18:02
*/
public class ProtoMessageManager {
private static final Logger logger = LoggerFactory.getLogger(ProtoMessageManager.class);
public static final String protoSerialNoTableName = "protoSerialNoTable.txt";
private static final Map<Integer, String> serialNoTypeMap = new HashMap<>(32);
private static final Map<String, Integer> typeSerialNoMap = new HashMap<>(32);
private static Map<String, MessageLite> typeNameClazzMap = new ConcurrentHashMap<>();
static {
try {
init();
} catch (IOException e) {
throw new IllegalStateException("protobuf序列号列表初始化失败", e);
}
}
private static void init() throws IOException {
InputStream in = ProtoMessageManager.class.getClassLoader().getResourceAsStream(protoSerialNoTableName);
if (in == null) {
throw new IllegalStateException("classpath:" + protoSerialNoTableName + " 文件未找到");
}
InputStreamReader reader = new InputStreamReader(in, StandardCharsets.UTF_8);
BufferedReader bufReader = new BufferedReader(reader);
String line;
while (null != (line = bufReader.readLine())) {
int idx = line.indexOf('=');
String clazz = line.substring(0, idx);
Integer num = Integer.valueOf(line.substring(idx + 1));
serialNoTypeMap.put(num, clazz);
typeSerialNoMap.put(clazz, num);
}
bufReader.close();
reader.close();
in.close();
}
@Nullable
public static MessageLite getProtoInstance(Integer serialNo) {
String clazz = serialNoTypeMap.get(serialNo);
if (clazz == null) {
logger.error("protoMsgDict serialNo proto class not found: {}", serialNo);
return null;
}
return typeNameClazzMap.computeIfAbsent(clazz, c -> {
try {
Class<?> type = Class.forName(c);
if (!MessageLite.class.isAssignableFrom(type)) {
return null;
}
return getDefaultInstance((Class<? extends MessageLite>) type);
} catch (ClassNotFoundException e) {
logger.error("protoMsgDict proto class not found: {}, {}", serialNo, c);
return null;
}
});
}
@Nullable
public static Integer getSerialNo(Class<? extends MessageLite> type) {
return typeSerialNoMap.get(type.getName());
}
private static MessageLite getDefaultInstance(Class<? extends MessageLite> clazz) {
try {
Method getDefaultInstance = clazz.getDeclaredMethod("getDefaultInstance");
return (MessageLite)getDefaultInstance.invoke(null);
} catch (Exception e) {
System.out.println("get instance exception:" + clazz.getName());
}
return null;
}
}

25
im-data/src/main/java/net/sopod/soim/data/serialize/ImMessage.java

@ -15,7 +15,12 @@ public class ImMessage {
public static final short MAGIC = 0x7a20;
private static final int MESSAGE_HEAD_LEN = 17;
/**
* {@link ImMessage#write(ByteBuf)}
* (short)magic + (int)serialNo + (int)serviceNo + (byte)serializeType + (byte)zipType + (byte)platformNo + (int)body.length
* = 17 byte
*/
private static final int MESSAGE_HEAD_LEN = 2 + 4 + 4 + 1 + 1 + 1 + 4;
public static final ImMessage PROTOCOL_ERROR;
@ -48,9 +53,6 @@ public class ImMessage {
/** 平台号 */
private int platformNo;
/** 请求体长度 */
private int bodyLength;
private byte[] body;
/**
@ -71,11 +73,11 @@ public class ImMessage {
.setSerializeType(buf.readByte())
.setZipType(buf.readByte())
.setPlatformNo(buf.readByte());
msg.bodyLength = buf.readInt();
if (bufLen != MESSAGE_HEAD_LEN + msg.getBodyLength()) {
int bodyLength = buf.readInt();
if (bufLen != MESSAGE_HEAD_LEN + bodyLength) {
return PROTOCOL_ERROR;
}
byte[] body = new byte[msg.getBodyLength()];
byte[] body = new byte[bodyLength];
buf.readBytes(body);
msg.setBody(body);
return msg;
@ -144,10 +146,6 @@ public class ImMessage {
return this;
}
public int getBodyLength() {
return bodyLength;
}
public byte[] getBody() {
return body;
}
@ -156,4 +154,9 @@ public class ImMessage {
this.body = body;
return this;
}
public int byteLength() {
return MESSAGE_HEAD_LEN + body.length;
}
}

4
im-data/src/main/resources/protoSerialNoTable.txt

@ -1,2 +1,2 @@
net.sopod.soim.data.msg.hello.HelloPB$World=10000
net.sopod.soim.data.msg.hello.HelloPB$Hello=10001
net.sopod.soim.data.msg.hello.HelloPB$Hello=10000
net.sopod.soim.data.msg.hello.HelloPB$World=10001

24
im-entry/src/main/java/net/sopod/soim/entry/config/ApplicationContextInitialed.java

@ -0,0 +1,24 @@
package net.sopod.soim.entry.config;
import net.sopod.soim.core.registry.ProtoMessageHandlerRegistry;
import org.springframework.beans.BeansException;
import org.springframework.context.ApplicationContext;
import org.springframework.context.ApplicationContextAware;
import org.springframework.context.annotation.Configuration;
/**
* ApplicationContextInitialed
*
* @author tmy
* @date 2022-04-10 22:20
*/
@Configuration
public class ApplicationContextInitialed implements ApplicationContextAware {
@Override
public void setApplicationContext(ApplicationContext applicationContext) throws BeansException {
// 注册 protobuf 消息 handler
ProtoMessageHandlerRegistry.registerHandlerWithApplicationContext(applicationContext);
}
}

26
im-entry/src/main/java/net/sopod/soim/entry/handler/HelloHandler.java

@ -0,0 +1,26 @@
package net.sopod.soim.entry.handler;
import com.google.protobuf.MessageLite;
import net.sopod.soim.core.handler.NetUserMessageHandler;
import net.sopod.soim.core.session.NetUser;
import net.sopod.soim.data.msg.hello.HelloPB;
import org.springframework.stereotype.Service;
/**
* HelloHandler
*
* @author tmy
* @date 2022-04-10 19:19
*/
@Service
public class HelloHandler extends NetUserMessageHandler<HelloPB.Hello> {
@Override
public MessageLite handle(NetUser netUser, HelloPB.Hello msg) {
System.out.println("get hello message");
System.out.println(msg);
System.out.println(msg.getStr());
return null;
}
}

16
im-entry/src/main/java/net/sopod/soim/entry/model/A.java

@ -1,16 +0,0 @@
package net.sopod.soim.entry.model;
import lombok.Data;
/**
* A
*
* @author tmy
* @date 2022-03-27 17:01
*/
@Data
public class A {
private String name;
}

5
im-entry/src/main/java/net/sopod/soim/entry/server/EntryServer.java

@ -9,6 +9,7 @@ import io.netty.channel.nio.NioEventLoopGroup;
import io.netty.channel.socket.nio.NioServerSocketChannel;
import io.netty.util.concurrent.DefaultThreadFactory;
import net.sopod.soim.common.constant.Consts;
import net.sopod.soim.entry.util.FastThreadLocalThreadFactory;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@ -38,9 +39,9 @@ public class EntryServer {
this.name = name;
this.port = port;
this.boss =
new NioEventLoopGroup(1, new DefaultThreadFactory("entry-b", Thread.MAX_PRIORITY));
new NioEventLoopGroup(1, new FastThreadLocalThreadFactory("entry-b-%d", Thread.MAX_PRIORITY));
this.worker =
new NioEventLoopGroup(new DefaultThreadFactory("entry-w", Thread.MAX_PRIORITY));
new NioEventLoopGroup(new FastThreadLocalThreadFactory("entry-w-%d", Thread.MAX_PRIORITY));
this.bootstrap();
}

7
im-entry/src/main/java/net/sopod/soim/entry/EntryServerStarter.java → im-entry/src/main/java/net/sopod/soim/entry/server/EntryServerRunner.java

@ -1,6 +1,5 @@
package net.sopod.soim.entry;
package net.sopod.soim.entry.server;
import net.sopod.soim.entry.server.EntryServer;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.boot.ApplicationArguments;
@ -14,9 +13,9 @@ import org.springframework.stereotype.Component;
* @date 2022-03-29 00:22
*/
@Component
public class EntryServerStarter implements ApplicationRunner {
public class EntryServerRunner implements ApplicationRunner {
private static final Logger logger = LoggerFactory.getLogger(EntryServerStarter.class);
private static final Logger logger = LoggerFactory.getLogger(EntryServerRunner.class);
@Override
public void run(ApplicationArguments args) {

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

@ -5,8 +5,7 @@ import io.netty.channel.ChannelPipeline;
import io.netty.channel.socket.SocketChannel;
import io.netty.handler.logging.LogLevel;
import io.netty.handler.logging.LoggingHandler;
import net.sopod.soim.core.net.ImEntryCodec;
import net.sopod.soim.core.net.ImMessageHandler;
import net.sopod.soim.core.net.ImMessageCodec;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@ -25,8 +24,8 @@ public class ImEntryInitializer extends ChannelInitializer<SocketChannel> {
LogLevel logLevel = logger.isDebugEnabled() ? LogLevel.DEBUG : LogLevel.INFO;
ChannelPipeline pipeline = socketChannel.pipeline();
pipeline.addLast(new LoggingHandler(logLevel))
.addLast(new ImEntryCodec())
.addLast(new ImMessageHandler());
.addLast(new ImMessageCodec())
.addLast(new InboundImMessageHandler());
}
}

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

@ -0,0 +1,53 @@
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.core.handler.MessageHandler;
import net.sopod.soim.core.net.AttributeKeys;
import net.sopod.soim.core.registry.ProtoMessageHandlerRegistry;
import net.sopod.soim.core.session.NetUser;
import java.util.concurrent.atomic.AtomicInteger;
/**
* MessageLiteHandler
*
* @author tmy
* @date 2022-04-10 22:40
*/
public class InboundImMessageHandler extends SimpleChannelInboundHandler<MessageLite> {
/**
* channel 建立连接设置初始属性
* TODO 登录倒计时 5s 断开连接, 登录失败次数
*/
@Override
public void channelActive(ChannelHandlerContext ctx) throws Exception {
Channel channel = ctx.channel();
channel.attr(NetUser.NET_USER_KEY).set(new NetUser(channel));
channel.attr(AttributeKeys.WRITE_FAIL_TIMES).set(new AtomicInteger());
channel.attr(AttributeKeys.LOGIN_FAIL_TIMES).set(new AtomicInteger());
ctx.fireChannelActive();
}
@Override
public void channelInactive(ChannelHandlerContext ctx) {
ctx.fireChannelInactive();
}
@Override
protected void channelRead0(ChannelHandlerContext ctx, MessageLite messageLite) throws Exception {
Attribute<NetUser> netUserAttr = ctx.channel().attr(NetUser.NET_USER_KEY);
NetUser netUser = netUserAttr.get();
// TODO dispatcher
MessageHandler<MessageLite> typeHandler = (MessageHandler<MessageLite>) ProtoMessageHandlerRegistry
.getTypeHandler(messageLite.getClass());
typeHandler.exec(netUser, messageLite);
}
}

28
im-entry/src/main/java/net/sopod/soim/entry/util/FastThreadLocalThreadFactory.java

@ -0,0 +1,28 @@
package net.sopod.soim.entry.util;
import io.netty.util.concurrent.FastThreadLocalThread;
import java.util.concurrent.ThreadFactory;
import java.util.concurrent.atomic.AtomicInteger;
public class FastThreadLocalThreadFactory implements ThreadFactory {
private String name;
private int priority;
private AtomicInteger counter;
public FastThreadLocalThreadFactory(String name, int priority) {
this.name = name;
this.priority = priority;
this.counter = new AtomicInteger();
}
@Override
public Thread newThread(Runnable runnable) {
FastThreadLocalThread thread = new FastThreadLocalThread(
runnable,
String.format(name, counter.incrementAndGet())
);
thread.setPriority(priority);
return thread;
}
}
Loading…
Cancel
Save