From 4b7e27083e4ebbd710db82d53c2fa3506c33ab62 Mon Sep 17 00:00:00 2001 From: tangmingyou <234767776@qq.com> Date: Mon, 28 Mar 2022 18:02:12 +0800 Subject: [PATCH] entry client --- README.md | 11 +- im-client/pom.xml | 15 ++ .../net/sopod/soim/client/ClientMain.java | 8 + .../sopod/soim/client/net/ImNetClient.java | 52 +++++ .../soim/client/server/ClientInitial.java | 13 -- im-common/pom.xml | 12 ++ .../sopod/soim/common/constant/Consts.java | 11 + .../net/sopod/soim/common/util/Jackson.java | 192 ++++++++++++++++++ im-core/pom.xml | 31 +++ .../net/sopod/soim/core/net/ImEntryCodec.java | 44 ++++ .../sopod/soim/core/net/ImMessageHandler.java | 26 +++ .../soim/core/registry/ServiceRegistry.java | 41 ++++ .../sopod/soim/core/service/ReqHandler.java | 13 ++ im-data/pom.xml | 4 + .../soim/data/constant/SerializeType.java | 22 +- .../soim/data/serialize/ByteSerializer.java | 15 ++ .../sopod/soim/data/serialize/ImMessage.java | 2 +- .../data/serialize/JacksonByteSerializer.java | 40 ++++ .../net/sopod/soim/data/transer/LoginReq.java | 18 ++ im-entry/pom.xml | 2 +- .../java/net/sopod/soim/entry/EntryMain.java | 13 +- .../sopod/soim/entry/server/EntryServer.java | 88 ++++++++ .../soim/entry/server/ImEntryInitializer.java | 32 +++ pom.xml | 1 + 24 files changed, 674 insertions(+), 32 deletions(-) create mode 100644 im-client/src/main/java/net/sopod/soim/client/net/ImNetClient.java delete mode 100644 im-client/src/main/java/net/sopod/soim/client/server/ClientInitial.java create mode 100644 im-common/src/main/java/net/sopod/soim/common/constant/Consts.java create mode 100644 im-common/src/main/java/net/sopod/soim/common/util/Jackson.java create mode 100644 im-core/pom.xml create mode 100644 im-core/src/main/java/net/sopod/soim/core/net/ImEntryCodec.java create mode 100644 im-core/src/main/java/net/sopod/soim/core/net/ImMessageHandler.java create mode 100644 im-core/src/main/java/net/sopod/soim/core/registry/ServiceRegistry.java create mode 100644 im-core/src/main/java/net/sopod/soim/core/service/ReqHandler.java create mode 100644 im-data/src/main/java/net/sopod/soim/data/serialize/ByteSerializer.java create mode 100644 im-data/src/main/java/net/sopod/soim/data/serialize/JacksonByteSerializer.java create mode 100644 im-data/src/main/java/net/sopod/soim/data/transer/LoginReq.java create mode 100644 im-entry/src/main/java/net/sopod/soim/entry/server/EntryServer.java create mode 100644 im-entry/src/main/java/net/sopod/soim/entry/server/ImEntryInitializer.java diff --git a/README.md b/README.md index 32116d0..2c726ff 100644 --- a/README.md +++ b/README.md @@ -30,14 +30,23 @@ 子网1和2不相通 +dubbo native image +https://dubbo.apache.org/zh/docs/references/graalvm/support-graalvm/ + + entry <--> client 通信 -entry <--> logic dubbo service +- serviceId --> paramClass --> serviceHandler + + +entry <--> logic dubbo service client jconsle cmd + router <--> cache das <--> db router <--> das + das shardingjdbc table struct, sharding roles logic biz diff --git a/im-client/pom.xml b/im-client/pom.xml index c4b3077..f8c7160 100644 --- a/im-client/pom.xml +++ b/im-client/pom.xml @@ -11,5 +11,20 @@ im-client + + + net.sopod + im-core + ${soim.version} + + + org.projectlombok + lombok + + + io.netty + netty-all + + \ No newline at end of file diff --git a/im-client/src/main/java/net/sopod/soim/client/ClientMain.java b/im-client/src/main/java/net/sopod/soim/client/ClientMain.java index c64286b..f0c2dd1 100644 --- a/im-client/src/main/java/net/sopod/soim/client/ClientMain.java +++ b/im-client/src/main/java/net/sopod/soim/client/ClientMain.java @@ -1,5 +1,7 @@ package net.sopod.soim.client; +import net.sopod.soim.client.net.ImNetClient; + /** * Main * @@ -8,4 +10,10 @@ package net.sopod.soim.client; */ public class ClientMain { + public static void main(String[] args) throws InterruptedException { + ImNetClient client = new ImNetClient(); + client.connect("127.0.0.1", 8088); + + } + } diff --git a/im-client/src/main/java/net/sopod/soim/client/net/ImNetClient.java b/im-client/src/main/java/net/sopod/soim/client/net/ImNetClient.java new file mode 100644 index 0000000..28e3ae2 --- /dev/null +++ b/im-client/src/main/java/net/sopod/soim/client/net/ImNetClient.java @@ -0,0 +1,52 @@ +package net.sopod.soim.client.net; + +import io.netty.bootstrap.Bootstrap; +import io.netty.channel.Channel; +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 java.nio.charset.StandardCharsets; +import java.util.HashMap; +import java.util.Map; + +/** + * ClientInitial + * + * @author tmy + * @date 2022-03-27 22:36 + */ +public class ImNetClient { + + public void connect(String host, int port) throws InterruptedException { + NioEventLoopGroup work = new NioEventLoopGroup(2); + Bootstrap b = new Bootstrap() + .group(work) + .channel(NioSocketChannel.class) + .handler(new ChannelInitializer() { + @Override + protected void initChannel(SocketChannel ch) throws Exception { + ch.pipeline() + .addLast(new ImEntryCodec()); + } + }); + Channel channel = b.connect(host, port).await().channel(); + Map body = new HashMap<>(); + body.put("name", "二狗子"); + body.put("age", 16); + ImMessage imMessage = new ImMessage() + .setServiceNo(1) + .setSerializeType(SerializeType.json.ordinal()) + .setBody(Jackson.json().serialize(body).getBytes(StandardCharsets.UTF_8)); + channel.writeAndFlush(imMessage); + channel.close(); + work.shutdownGracefully(); + } + + +} diff --git a/im-client/src/main/java/net/sopod/soim/client/server/ClientInitial.java b/im-client/src/main/java/net/sopod/soim/client/server/ClientInitial.java deleted file mode 100644 index 0425732..0000000 --- a/im-client/src/main/java/net/sopod/soim/client/server/ClientInitial.java +++ /dev/null @@ -1,13 +0,0 @@ -package net.sopod.soim.client.server; - -/** - * ClientInitial - * - * @author tmy - * @date 2022-03-27 22:36 - */ -public class ClientInitial { - - - -} diff --git a/im-common/pom.xml b/im-common/pom.xml index fe6825d..d2f1473 100644 --- a/im-common/pom.xml +++ b/im-common/pom.xml @@ -11,5 +11,17 @@ im-common + + + com.fasterxml.jackson.core + jackson-databind + ${jackson.version} + + + com.fasterxml.jackson.datatype + jackson-datatype-jsr310 + ${jackson.version} + + \ No newline at end of file diff --git a/im-common/src/main/java/net/sopod/soim/common/constant/Consts.java b/im-common/src/main/java/net/sopod/soim/common/constant/Consts.java new file mode 100644 index 0000000..90a8421 --- /dev/null +++ b/im-common/src/main/java/net/sopod/soim/common/constant/Consts.java @@ -0,0 +1,11 @@ +package net.sopod.soim.common.constant; + +public class Consts { + + public static int KB = 1024; + + public static int MB = KB * KB; + + public static int GB = KB * MB; + +} diff --git a/im-common/src/main/java/net/sopod/soim/common/util/Jackson.java b/im-common/src/main/java/net/sopod/soim/common/util/Jackson.java new file mode 100644 index 0000000..033bf02 --- /dev/null +++ b/im-common/src/main/java/net/sopod/soim/common/util/Jackson.java @@ -0,0 +1,192 @@ +package net.sopod.soim.common.util; + +import com.fasterxml.jackson.core.JsonFactory; +import com.fasterxml.jackson.core.JsonParser; +import com.fasterxml.jackson.core.JsonProcessingException; +import com.fasterxml.jackson.core.json.JsonReadFeature; +import com.fasterxml.jackson.core.json.PackageVersion; +import com.fasterxml.jackson.databind.DeserializationFeature; +import com.fasterxml.jackson.databind.ObjectMapper; +import com.fasterxml.jackson.databind.SerializationFeature; +import com.fasterxml.jackson.databind.module.SimpleModule; +import com.fasterxml.jackson.datatype.jsr310.deser.LocalDateDeserializer; +import com.fasterxml.jackson.datatype.jsr310.deser.LocalDateTimeDeserializer; +import com.fasterxml.jackson.datatype.jsr310.deser.LocalTimeDeserializer; +import com.fasterxml.jackson.datatype.jsr310.ser.LocalDateSerializer; +import com.fasterxml.jackson.datatype.jsr310.ser.LocalDateTimeSerializer; +import com.fasterxml.jackson.datatype.jsr310.ser.LocalTimeSerializer; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.text.SimpleDateFormat; +import java.time.LocalDate; +import java.time.LocalDateTime; +import java.time.LocalTime; +import java.time.ZoneId; +import java.time.format.DateTimeFormatter; +import java.util.Locale; +import java.util.Map; +import java.util.TimeZone; +import java.util.function.Consumer; +import java.util.function.Supplier; + +/** + * + * @author tmy + * @date 2022-03-04 16:37 + */ +public class Jackson { + + private static final Logger log = LoggerFactory.getLogger(Jackson.class); + + private static final String DATE_TIME_FORMAT = "yyyy-MM-dd HH:mm:ss"; + private static final String DATE_FORMAT = "yyyy-MM-dd"; + private static final String TIME_FORMAT = "HH:mm:ss"; + + private static volatile Jackson JSON_INSTANCE; + private static volatile Jackson YAML_INSTANCE; + private static volatile Jackson XML_INSTANCE; + + private final ObjectMapper objectMapper; + + private Jackson(ObjectMapper objectMapper) { + this.objectMapper = objectMapper; + + } + + private static void getFactoryInstance(String clazz, Supplier predicate, Consumer setter) { + getFactoryInstance(clazz, predicate, mapper->{}, setter); + } + + private static void getFactoryInstance(String clazz, Supplier predicate, Consumer setting, Consumer setter) { + if (predicate.get()) { + synchronized (Jackson.class) { + if (predicate.get()) { + try { + Class factoryClazz = Class.forName(clazz); + Object factoryInstance = factoryClazz.getDeclaredConstructor().newInstance(); + ObjectMapper mapper; + if (factoryInstance instanceof ObjectMapper) { + mapper = new JacksonObjectMapper((ObjectMapper) factoryInstance); + } else { + mapper = new JacksonObjectMapper((JsonFactory) factoryInstance); + } + setter.accept(new Jackson(mapper)); + } catch (ReflectiveOperationException e) { + throw new RuntimeException(clazz, e); + } + } + } + } + } + + private static class JacksonObjectMapper extends ObjectMapper { + + private static final long serialVersionUID = -1848248131995072046L; + + public JacksonObjectMapper(ObjectMapper src) { + super(src); + init(); + } + + public JacksonObjectMapper(JsonFactory factory) { + super(factory); + init(); + } + + @Override + public ObjectMapper copy() { + return new JacksonObjectMapper(this); + } + + + private void init() { + super.setLocale(Locale.CHINA); + super.configure(SerializationFeature.WRITE_DATES_AS_TIMESTAMPS, false); + super.setTimeZone(TimeZone.getTimeZone(ZoneId.systemDefault())); + super.setDateFormat(new SimpleDateFormat(DATE_TIME_FORMAT, Locale.CHINA)); + super.configure(JsonParser.Feature.ALLOW_SINGLE_QUOTES, true); + super.configure(JsonReadFeature.ALLOW_UNESCAPED_CONTROL_CHARS.mappedFeature(), true); + super.configure(JsonReadFeature.ALLOW_BACKSLASH_ESCAPING_ANY_CHARACTER.mappedFeature(), true); + + super.configure(SerializationFeature.FAIL_ON_EMPTY_BEANS, false); + super.configure(DeserializationFeature.FAIL_ON_UNKNOWN_PROPERTIES, false); + super.configure(JsonReadFeature.ALLOW_SINGLE_QUOTES.mappedFeature(), true); + super.getDeserializationConfig().withoutFeatures(DeserializationFeature.FAIL_ON_UNKNOWN_PROPERTIES); + + super.findAndRegisterModules(); + // 放到jackson-datatype-jsr310后注册,不然被覆盖 + super.registerModule(new MyJavaTimeModule()); + } + + } + + /** + * jsr310 默认 LocalDateTime 正反序列化格式 DateTimeFormatter.ISO_LOCAL_DATE_TIME + */ + private static class MyJavaTimeModule extends SimpleModule { + public MyJavaTimeModule() { + super(PackageVersion.VERSION); + DateTimeFormatter dateTimeFormatter = DateTimeFormatter.ofPattern(DATE_TIME_FORMAT); + DateTimeFormatter DateFormatter = DateTimeFormatter.ofPattern(DATE_FORMAT); + DateTimeFormatter timeFormatter = DateTimeFormatter.ofPattern(TIME_FORMAT); + this.addDeserializer(LocalDateTime.class, new LocalDateTimeDeserializer(dateTimeFormatter)); + this.addDeserializer(LocalDate.class, new LocalDateDeserializer(DateFormatter)); + this.addDeserializer(LocalTime.class, new LocalTimeDeserializer(timeFormatter)); + this.addSerializer(LocalDateTime.class, new LocalDateTimeSerializer(dateTimeFormatter)); + this.addSerializer(LocalDate.class, new LocalDateSerializer(DateFormatter)); + this.addSerializer(LocalTime.class, new LocalTimeSerializer(timeFormatter)); + } + } + + public static Jackson json() { + getFactoryInstance("com.fasterxml.jackson.core.JsonFactory", + () -> JSON_INSTANCE == null, + jackson -> JSON_INSTANCE = jackson); + return JSON_INSTANCE; + } + + /** + * 依赖 jackson-dataformat-yaml + */ + public static Jackson yaml() { + getFactoryInstance("com.fasterxml.jackson.dataformat.yaml.YAMLFactory", + () -> YAML_INSTANCE == null, + jackson -> YAML_INSTANCE = jackson); + return YAML_INSTANCE; + } + + /** + * 依赖 jackson-dataformat-xml + */ + public static Jackson xml() { + getFactoryInstance("com.fasterxml.jackson.dataformat.xml.XmlMapper", + () -> XML_INSTANCE == null, + jackson -> XML_INSTANCE = jackson); + return XML_INSTANCE; + } + + public T deserialize(String content, Class valueType) { + try { + return objectMapper.readValue(content, valueType); + } catch (JsonProcessingException e) { + log.error(e.getMessage(), e); + return null; + } + } + + public String serialize(T value) { + try { + return objectMapper.writeValueAsString(value); + } catch (JsonProcessingException e) { + e.printStackTrace(); + log.error(e.getMessage(), e); + return null; + } + } + + public T toPojo(Map fromValue, Class toValueType) { + return objectMapper.convertValue(fromValue, toValueType); + } + +} diff --git a/im-core/pom.xml b/im-core/pom.xml new file mode 100644 index 0000000..2c6929b --- /dev/null +++ b/im-core/pom.xml @@ -0,0 +1,31 @@ + + + + so-im + net.sopod + 1.0.0 + + 4.0.0 + + im-core + + + + net.sopod + im-common + ${soim.version} + + + net.sopod + im-data + ${soim.version} + + + io.netty + netty-all + + + + \ No newline at end of file diff --git a/im-core/src/main/java/net/sopod/soim/core/net/ImEntryCodec.java b/im-core/src/main/java/net/sopod/soim/core/net/ImEntryCodec.java new file mode 100644 index 0000000..3e02783 --- /dev/null +++ b/im-core/src/main/java/net/sopod/soim/core/net/ImEntryCodec.java @@ -0,0 +1,44 @@ +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 java.util.List; + +/** + * ImEntryCodec + * + * @author tmy + * @date 2022-03-28 11:29 + */ +public class ImEntryCodec extends CombinedChannelDuplexHandler { + + public ImEntryCodec() { + super(new ImDecoder(), new ImEncoder()); + } + + public static class ImDecoder extends ByteToMessageDecoder { + @Override + protected void decode(ChannelHandlerContext ctx, ByteBuf byteBuf, List list) throws Exception { + ImMessage message = ImMessage.read(byteBuf); + if (message == ImMessage.MAGIC_ERROR + || message == ImMessage.PROTOCOL_ERROR) { + ctx.channel().close(); + return; + } + list.add(message); + } + } + + public static class ImEncoder extends MessageToByteEncoder { + @Override + protected void encode(ChannelHandlerContext ctx, ImMessage imMessage, ByteBuf byteBuf) throws Exception { + imMessage.write(byteBuf); + } + } + +} diff --git a/im-core/src/main/java/net/sopod/soim/core/net/ImMessageHandler.java b/im-core/src/main/java/net/sopod/soim/core/net/ImMessageHandler.java new file mode 100644 index 0000000..79bb651 --- /dev/null +++ b/im-core/src/main/java/net/sopod/soim/core/net/ImMessageHandler.java @@ -0,0 +1,26 @@ +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 { + + @Override + protected void channelRead0(ChannelHandlerContext channelHandlerContext, ImMessage imMessage) throws Exception { + byte[] body = imMessage.getBody(); + SerializeType serialize = SerializeType.getSerializeByOrdinal(imMessage.getSerializeType()); + Map data = serialize.getSerializer().deserialize(body, Map.class); + System.out.println(data); + } + +} diff --git a/im-core/src/main/java/net/sopod/soim/core/registry/ServiceRegistry.java b/im-core/src/main/java/net/sopod/soim/core/registry/ServiceRegistry.java new file mode 100644 index 0000000..505555a --- /dev/null +++ b/im-core/src/main/java/net/sopod/soim/core/registry/ServiceRegistry.java @@ -0,0 +1,41 @@ +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> serviceIdParamTypeMap; + + private static final ConcurrentHashMap, Integer> paramTypeServiceIdMap; + + private static final ConcurrentHashMap> serviceIdHandlers; + + static { + serviceIdParamTypeMap = new ConcurrentHashMap<>(); + paramTypeServiceIdMap = new ConcurrentHashMap<>(); + serviceIdHandlers = new ConcurrentHashMap<>(); + } + + private static final AtomicInteger serviceIdGen = new AtomicInteger(10000); + + private static void registry(Class paramType, ReqHandler handler) { + int serviceId = serviceIdGen.getAndIncrement(); + serviceIdParamTypeMap.put(serviceId, paramType); + paramTypeServiceIdMap.put(paramType, serviceId); + serviceIdHandlers.put(serviceId, handler); + } + + public void aaa() { + + } + +} diff --git a/im-core/src/main/java/net/sopod/soim/core/service/ReqHandler.java b/im-core/src/main/java/net/sopod/soim/core/service/ReqHandler.java new file mode 100644 index 0000000..d1bb276 --- /dev/null +++ b/im-core/src/main/java/net/sopod/soim/core/service/ReqHandler.java @@ -0,0 +1,13 @@ +package net.sopod.soim.core.service; + +/** + * ReqHandler + * + * @author tmy + * @date 2022-03-28 14:33 + */ +public interface ReqHandler { + + Object handle(T param); + +} diff --git a/im-data/pom.xml b/im-data/pom.xml index c58ec6d..5a9631b 100644 --- a/im-data/pom.xml +++ b/im-data/pom.xml @@ -21,6 +21,10 @@ io.netty netty-buffer + + org.projectlombok + lombok + \ No newline at end of file diff --git a/im-data/src/main/java/net/sopod/soim/data/constant/SerializeType.java b/im-data/src/main/java/net/sopod/soim/data/constant/SerializeType.java index ced722d..4854645 100644 --- a/im-data/src/main/java/net/sopod/soim/data/constant/SerializeType.java +++ b/im-data/src/main/java/net/sopod/soim/data/constant/SerializeType.java @@ -1,5 +1,8 @@ package net.sopod.soim.data.constant; +import net.sopod.soim.data.serialize.ByteSerializer; +import net.sopod.soim.data.serialize.JacksonByteSerializer; + /** * SerializeType * @@ -8,21 +11,22 @@ package net.sopod.soim.data.constant; */ public enum SerializeType { - json, - - protobuf; - + json(new JacksonByteSerializer()), - public static interface Converter { + protobuf(null); - T deserialize(byte[] data, Class clazz) throws Exception; - - byte[] serialize(Object pojo); + private final ByteSerializer serializer; + SerializeType(ByteSerializer serializer) { + this.serializer = serializer; } - public static void main(String[] args) { + public ByteSerializer getSerializer() { + return serializer; + } + public static SerializeType getSerializeByOrdinal(int ordinal) { + return values()[ordinal]; } } diff --git a/im-data/src/main/java/net/sopod/soim/data/serialize/ByteSerializer.java b/im-data/src/main/java/net/sopod/soim/data/serialize/ByteSerializer.java new file mode 100644 index 0000000..d52fa3b --- /dev/null +++ b/im-data/src/main/java/net/sopod/soim/data/serialize/ByteSerializer.java @@ -0,0 +1,15 @@ +package net.sopod.soim.data.serialize; + +/** + * ByteSerializer + * + * @author tmy + * @date 2022-03-28 11:03 + */ +public interface ByteSerializer { + + T deserialize(byte[] data, Class clazz) throws Exception; + + byte[] serialize(Object pojo) throws Exception; + +} diff --git a/im-data/src/main/java/net/sopod/soim/data/serialize/ImMessage.java b/im-data/src/main/java/net/sopod/soim/data/serialize/ImMessage.java index 88162ee..781b7a1 100644 --- a/im-data/src/main/java/net/sopod/soim/data/serialize/ImMessage.java +++ b/im-data/src/main/java/net/sopod/soim/data/serialize/ImMessage.java @@ -15,7 +15,7 @@ public class ImMessage { public static final short MAGIC = 0x7a20; - public static final int MESSAGE_HEAD_LEN = 17; + private static final int MESSAGE_HEAD_LEN = 17; public static final ImMessage PROTOCOL_ERROR; diff --git a/im-data/src/main/java/net/sopod/soim/data/serialize/JacksonByteSerializer.java b/im-data/src/main/java/net/sopod/soim/data/serialize/JacksonByteSerializer.java new file mode 100644 index 0000000..c445efc --- /dev/null +++ b/im-data/src/main/java/net/sopod/soim/data/serialize/JacksonByteSerializer.java @@ -0,0 +1,40 @@ +package net.sopod.soim.data.serialize; + +import net.sopod.soim.common.util.Jackson; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.nio.charset.StandardCharsets; +import java.util.Arrays; +import java.util.Collections; +import java.util.HashMap; +import java.util.Map; + +/** + * JsonByteSerializer + * + * @author tmy + * @date 2022-03-28 11:05 + */ +public class JacksonByteSerializer implements ByteSerializer{ + + private static final Logger logger = LoggerFactory.getLogger(JacksonByteSerializer.class); + + @Override + public T deserialize(byte[] data, Class clazz) throws Exception { + String json = new String(data, StandardCharsets.UTF_8); + T obj = Jackson.json().deserialize(json, clazz); + if (obj == null) { + logger.error("json数据解析失败:{}", json); + throw new IllegalStateException("json数据解析失败"); + } + return obj; + } + + @Override + public byte[] serialize(Object pojo) throws Exception { + String json = Jackson.json().serialize(pojo); + return json.getBytes(StandardCharsets.UTF_8); + } + +} diff --git a/im-data/src/main/java/net/sopod/soim/data/transer/LoginReq.java b/im-data/src/main/java/net/sopod/soim/data/transer/LoginReq.java new file mode 100644 index 0000000..d518d12 --- /dev/null +++ b/im-data/src/main/java/net/sopod/soim/data/transer/LoginReq.java @@ -0,0 +1,18 @@ +package net.sopod.soim.data.transer; + +import lombok.Data; + +/** + * LoginReq + * + * @author tmy + * @date 2022-03-28 14:24 + */ +@Data +public class LoginReq { + + private String name; + + private String password; + +} diff --git a/im-entry/pom.xml b/im-entry/pom.xml index 9232150..1450daa 100644 --- a/im-entry/pom.xml +++ b/im-entry/pom.xml @@ -18,7 +18,7 @@ net.sopod - im-common + im-core ${soim.version} diff --git a/im-entry/src/main/java/net/sopod/soim/entry/EntryMain.java b/im-entry/src/main/java/net/sopod/soim/entry/EntryMain.java index e7b032e..8de8bca 100644 --- a/im-entry/src/main/java/net/sopod/soim/entry/EntryMain.java +++ b/im-entry/src/main/java/net/sopod/soim/entry/EntryMain.java @@ -1,6 +1,6 @@ package net.sopod.soim.entry; -import net.sopod.soim.entry.model.A; +import net.sopod.soim.entry.server.EntryServer; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -15,13 +15,12 @@ public class EntryMain { private static final Logger logger = LoggerFactory.getLogger(EntryMain.class); public static void main(String[] args) { - logger.info("info"); - logger.warn("warn"); - logger.error("error..."); + EntryServer entryServer = new EntryServer("entry server", 8088); + entryServer.startServer(err -> { + logger.error("EntryServer 启动失败:", err); + }); - A a = new A(); - a.setName("太极拳"); - System.out.println(a); + Runtime.getRuntime().addShutdownHook(new Thread(entryServer::shutdown)); } } diff --git a/im-entry/src/main/java/net/sopod/soim/entry/server/EntryServer.java b/im-entry/src/main/java/net/sopod/soim/entry/server/EntryServer.java new file mode 100644 index 0000000..e4d6c9f --- /dev/null +++ b/im-entry/src/main/java/net/sopod/soim/entry/server/EntryServer.java @@ -0,0 +1,88 @@ +package net.sopod.soim.entry.server; + +import io.netty.bootstrap.ServerBootstrap; +import io.netty.buffer.PooledByteBufAllocator; +import io.netty.channel.ChannelFutureListener; +import io.netty.channel.ChannelOption; +import io.netty.channel.WriteBufferWaterMark; +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 org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.util.function.Consumer; + +/** + * EntryServer + * + * @author tmy + * @date 2022-03-28 11:19 + */ +public class EntryServer { + + private static final Logger logger = LoggerFactory.getLogger(EntryServer.class); + + private NioEventLoopGroup boss; + private NioEventLoopGroup worker; + private ServerBootstrap serverBootstrap; + + private String name; + private int port; + + public static final int LOW_WATER_MARK = 32 * Consts.KB; + public static final int HIGH_WATER_MARK = 64 * Consts.KB; + + public EntryServer(String name, int port) { + this.name = name; + this.port = port; + this.boss = + new NioEventLoopGroup(1, new DefaultThreadFactory("entry-boss", Thread.MAX_PRIORITY)); + this.worker = + new NioEventLoopGroup(new DefaultThreadFactory("entry-worker", Thread.MAX_PRIORITY)); + this.bootstrap(); + } + + private void bootstrap() { + this.serverBootstrap = new ServerBootstrap() + .group(boss, worker) + .channel(NioServerSocketChannel.class) + .childHandler(new ImEntryInitializer()) + .option(ChannelOption.ALLOCATOR, PooledByteBufAllocator.DEFAULT) + .option(ChannelOption.SO_RCVBUF, 32 * Consts.KB) + .option(ChannelOption.SO_REUSEADDR, true) + .option(ChannelOption.WRITE_BUFFER_WATER_MARK, + new WriteBufferWaterMark(LOW_WATER_MARK, HIGH_WATER_MARK)) + .childOption(ChannelOption.ALLOCATOR, PooledByteBufAllocator.DEFAULT) + .childOption(ChannelOption.TCP_NODELAY, true) + .childOption(ChannelOption.SO_RCVBUF, 32 * Consts.KB) + .childOption(ChannelOption.SO_SNDBUF, 64 * Consts.KB) + .childOption(ChannelOption.SO_REUSEADDR, true) + .childOption( + ChannelOption.WRITE_BUFFER_WATER_MARK, + new WriteBufferWaterMark(LOW_WATER_MARK, HIGH_WATER_MARK)); + } + + public void startServer(Consumer onFail) { + this.serverBootstrap + .bind(port) + .addListener((ChannelFutureListener) future -> { + if (!future.isSuccess()) { + if (onFail != null) { + onFail.accept(future.cause()); + } + return; + } + System.out.println(String.format("%s %s listen...", name, port)); + }); + } + + public void shutdown() { + logger.info("netty reactor group shutting down..."); + boss.shutdownGracefully(); + worker.shutdownGracefully(); + logger.info("netty reactor group already shutdown!"); + } + +} 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 new file mode 100644 index 0000000..2d63caa --- /dev/null +++ b/im-entry/src/main/java/net/sopod/soim/entry/server/ImEntryInitializer.java @@ -0,0 +1,32 @@ +package net.sopod.soim.entry.server; + +import io.netty.channel.ChannelInitializer; +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 org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * ImEntryInitializer + * + * @author tmy + * @date 2022-03-28 11:25 + */ +public class ImEntryInitializer extends ChannelInitializer { + + private static final Logger logger = LoggerFactory.getLogger(ImEntryInitializer.class); + + @Override + protected void initChannel(SocketChannel socketChannel) throws Exception { + LogLevel logLevel = logger.isDebugEnabled() ? LogLevel.DEBUG : LogLevel.INFO; + ChannelPipeline pipeline = socketChannel.pipeline(); + pipeline.addLast(new LoggingHandler(logLevel)) + .addLast(new ImEntryCodec()) + .addLast(new ImMessageHandler()); + } + +} diff --git a/pom.xml b/pom.xml index fb3062b..e1edcee 100644 --- a/pom.xml +++ b/pom.xml @@ -16,6 +16,7 @@ im-client im-logic-api im-logic-api/im-user-api + im-core