diff --git a/README.md b/README.md index b2a922e..3677326 100644 --- a/README.md +++ b/README.md @@ -1,4 +1,48 @@ -TODO + + +### 开发计划 + +- (05-08~05-11)router 一致性 hash,新增、删除、备份处理 + - 数据同步删除后,保留id一段时间,有请求进行重定向/转发 + - 新增:重算hash,发起数据同步,接收其他节点推送数据及数据更改日志,注册服务,其他节点删除数据 + - 正常删除节点:重算 hash,数据和更改日志推送给其他节点,取消注册,关闭服务 + - 宕机:备份服务接收全量同步数据和更改日志,服务探测到主服务不可用,注册到注册中心提供服务 +- (05-12~05-12) 备用计划:router 层 id 号段负载均衡 +- (05-13~05-14) dubbo 服务异步处理 +- 功能开发: + - (05-13~05-13) 用户注册功能 + - (05-14~05-14) 好友列表(在线状态:批量uid一致性hash, router查询) + - 用户查询 + - 添加好友 + - 好友在线状态列表 + - (05-15~05-15) 消息群发(im-router 群消息路由,批量uid一致性哈希路由) + - 创群: + - 加群:发送申请,推送申请,推送回复 + - 删群 + - 群列表 + - (05-16~05-16) 聊天记录查询 +- (05-17~05-18) das 部分接口消息队列异步写 +- (05-19~05-20) 集群部署:docker swarm +- 监控:log4j2完善 + prometheus + grafana + loki + arthas + - (05-21~05-21) prometheus 服务发现,exporter(服务发现) 开发 + - (05-22~05-22) 容器监控,主机监控,数据库监控,日志收集 + - (05-23~05-24) 应用监控 + - 通讯指标监控,请求量,吞吐量,延时,请求节点分布,链路追踪 + - 数据库连接池监控 +- (05-25~05-25) 压测:压测开发 +- (05-26~05-26) client: 控制台完善,grallvm 打包 +- 后续: + - 通讯加密 + - websocket 网关 + - web 页面开发 + - 异/同设备,多地登录 + + + + + +### TODO + - router -> entry 负载均衡 - router 新增节点顺时针相邻节点数据一致性哈希迁移 - router 冗余节点存储数据不提供服务 @@ -22,9 +66,9 @@ TODO - 内存存储层 router - 固化存储层 das - 功能组件: 1.0 + - IOC: guice - 通信层: netty - 序列化: protobuf @@ -43,14 +87,13 @@ TODO - ID生成器 2.0 -- quarkus +- quarkus 集群负载均衡: 子网1,子网2,公网 子网1和2不相通 - dubbo native image https://dubbo.apache.org/zh/docs/references/graalvm/support-graalvm/ @@ -58,11 +101,9 @@ https://dubbo.apache.org/zh/docs/references/graalvm/support-graalvm/ entry <--> client 通信 - serviceId --> paramClass --> serviceHandler - entry <--> logic dubbo service client jconsle cmd - router <--> cache das <--> db router <--> das @@ -71,5 +112,5 @@ router <--> das das shardingjdbc table struct, sharding roles logic biz - + 请求响应消息队列异步处理,减少线程 cpu 占用, dubbo async \ No newline at end of file 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 index 9845667..b1fb6ea 100644 --- 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 @@ -8,6 +8,8 @@ import com.fasterxml.jackson.core.json.JsonReadFeature; import com.fasterxml.jackson.core.json.PackageVersion; import com.fasterxml.jackson.databind.*; import com.fasterxml.jackson.databind.module.SimpleModule; +import com.fasterxml.jackson.databind.type.CollectionLikeType; +import com.fasterxml.jackson.databind.type.MapType; import com.fasterxml.jackson.datatype.jsr310.deser.LocalDateDeserializer; import com.fasterxml.jackson.datatype.jsr310.deser.LocalDateTimeDeserializer; import com.fasterxml.jackson.datatype.jsr310.deser.LocalTimeDeserializer; @@ -24,9 +26,7 @@ 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.*; import java.util.function.Consumer; import java.util.function.Supplier; @@ -46,6 +46,7 @@ public class Jackson { private static volatile Jackson JSON_INSTANCE; private static volatile Jackson YAML_INSTANCE; private static volatile Jackson XML_INSTANCE; + private static volatile Jackson MSGPACK_INSTANCE; private final ObjectMapper objectMapper; @@ -176,6 +177,19 @@ public class Jackson { return XML_INSTANCE; } + /** + * 序列化的时候,不写入字段名字,会按字段顺序写入值 + * 如果在bean中要增加新字段,请务必保证新字段加在字段序的最后! + * 对象新增字段,放在中间位置,会导致序列化失败! + */ + public static Jackson msgpack() { + getFactoryInstance("org.msgpack.jackson.dataformat.MessagePackFactory", + () -> MSGPACK_INSTANCE == null, + jackson -> MSGPACK_INSTANCE = jackson + ); + return MSGPACK_INSTANCE; + } + public boolean canSerialize(Class clazz) { return objectMapper.canSerialize(clazz); } @@ -199,8 +213,107 @@ public class Jackson { } } + public byte[] serializeBytes(T value) { + try { + return objectMapper.writeValueAsBytes(value); + } catch (JsonProcessingException e) { + e.printStackTrace(); + return null; + } + } + + public T deserializeBytes(byte[] bytes, Class valueType) { + try { + return objectMapper.readValue(bytes, valueType); + } catch (IOException e) { + e.printStackTrace(); + return null; + } + } + public T toPojo(Map fromValue, Class toValueType) { return objectMapper.convertValue(fromValue, toValueType); } + public Map readMap(String content, Class keyClass, Class valueClass) { + if (content == null || content.length() == 0) { + return Collections.emptyMap(); + } else { + try { + return objectMapper.readValue(content, getMapType(keyClass, valueClass)); + } catch (IOException e) { + e.printStackTrace(); + return null; + } + } + } + + public Map readMap(byte[] content, Class keyClass, Class valueClass) { + if (content == null || content.length == 0) { + return Collections.emptyMap(); + } else { + try { + return objectMapper.readValue(content, getMapType(keyClass, valueClass)); + } catch (IOException e) { + e.printStackTrace(); + return null; + } + } + } + + private MapType getMapType(Class keyClass, Class valueClass) { + return objectMapper.getTypeFactory().constructMapType(Map.class, keyClass, valueClass); + } + + + public List> readListMap(String content) { + return readListMap(content, Object.class); + } + + public List> readListMap(String content, Class valueType) { + if (content == null || content.length() == 0) { + return Collections.emptyList(); + } else { + try { + return objectMapper.readValue(content, objectMapper.getTypeFactory() + .constructCollectionLikeType(List.class, getMapType(String.class, valueType))); + } catch (IOException e) { + e.printStackTrace(); + return null; + } + } + } + + + public List readList(byte[] content, Class elementClass) { + if (content == null || content.length == 0) { + return Collections.emptyList(); + } else { + try { + return objectMapper.readValue(content, getListType(elementClass)); + } catch (IOException e) { + e.printStackTrace(); + return null; + } + } + } + + public List readList(String content, Class elementClass) { + if (content == null || content.length() == 0) { + return Collections.emptyList(); + } else { + try { + return objectMapper.readValue(content, getListType(elementClass)); + } catch (IOException e) { + e.printStackTrace(); + return null; + } + } + } + + private CollectionLikeType getListType(Class elementClass) { + objectMapper.getTypeFactory().constructCollectionLikeType(List.class, getMapType(String.class, Object.class)); + return objectMapper.getTypeFactory().constructCollectionLikeType(List.class, elementClass); + } + } diff --git a/im-common/src/main/resources/banner.txt b/im-common/src/main/resources/banner.txt new file mode 100644 index 0000000..be8e9c1 --- /dev/null +++ b/im-common/src/main/resources/banner.txt @@ -0,0 +1,18 @@ +⣿⣿⣿⣿⣿⣿⣿⣿⣿⠻⢿⣿⣿⣿⣿⣿⣿⣿⣿⣿⣿⣿⣿⣿⣿ +⣿⣿⣿⣿⣿⣿⣿⣿⣿⠀⠀⠹⢿⣿⣿⣿⣿⣿⡿⠋⣿⣿⣿⣿⣿ +⣿⣿⣿⣿⣿⣿⣿⣿⣿⠀⠀⠀⠈⠻⣿⣿⡿⠏⠀⠀⣿⣿⣿⣿⣿ +⣿⣿⣿⣿⣿⣿⣿⣿⣿⠀⠀⠀⠀⠀⠙⠋⠀⠀⠀⢀⣿⣿⣿⣿⣿ +⣿⣿⣷⡈⠉⠉⠉⠉⠉⠀⣀⡤⠴⠶⠶⠶⠤⣄⡀⠸⠿⠿⠛⢛⣿ +⣿⣿⣿⣿⡄⠀⠀⠀⣰⠚⠁⠀⠀⠀⠀⠀⠀⠀⠳⣄⠀⠀⢠⣾⣿ +⣿⣿⣿⣿⣷⠀⠀⡜⠁⠀⠀⠀⠀⣀⣀⣀⣀⣀⣀⠘⡆⢠⣿⣿⣿ +⣿⣿⠿⠛⠉⠀⢰⠇⠀⠰⣾⡯⠭⠭⣭⠭⡭⠭⠭⠭⢿⡀⠙⠿⣿ +⣿⣤⣀⠀⠀⠀⢸⠀⢠⣧⣤⣤⣤⠤⢼⣾⠥⣤⡤⠤⠬⡇⣠⣴⣾ +⣿⣿⣿⣿⡦⠀⢸⠀⠈⢧⡀⠁⠀⠀⡠⠛⣄⠈⠀⢀⣠⠇⣿⣿⣿ +⣿⣿⠿⠋⠀⠀⣸⡄⠀⢤⣉⠓⠚⠋⠁⡀⢸⡩⣯⡭⢿⡀⠙⠿⣿ +⣿⣷⣤⣄⡀⢼⠋⠀⠀⠀⠉⠉⠁⠀⠀⠳⣼⠇⠀⠀⣼⢹⣾⣿⣿ +⣿⣿⣿⣿⡟⠈⠳⣤⠀⠀⡀⠀⠀⣀⣀⣀⣀⣀⡀⠀⣿⡻⣿⣿⣿ +⣿⣿⣿⣿⣾⣿⣿⠞⣆⠸⡇⠚⠉⠁⠀⣿⣿⡟⠉⢁⣿⣿⣿⣿⣿ +⣿⣿⣿⣿⣿⣿⣿⣶⣿⣆⠈⠁⠀⠰⡖⠙⠛⠃⣠⣿⣿⣿⣿⣿⣿ +⣿⣿⣿⣿⣿⣿⣿⣿⣿⣿⣿⡦⣄⣀⣀⣠⣴⣾⣿⣿⣿⣿⣿⣿⣿ +⣿⣿⣿⣿⣿⣿⡿⠿⠟⠛⡛⢷⣤⣀⣠⢾⣛⠛⠿⢿⣿⣿⣿⣿⣿ +⣿⣿⣿⣿⣿⠃⠀⠀⠀⢰⠃⢸⠀⠀⠀⠘⡟⣆⠀⠀⢹⣿⣿⣿⣿ \ No newline at end of file diff --git a/im-service/im-router/pom.xml b/im-service/im-router/pom.xml index 86d1699..d53a16c 100644 --- a/im-service/im-router/pom.xml +++ b/im-service/im-router/pom.xml @@ -85,6 +85,14 @@ cglib cglib + + org.msgpack + jackson-dataformat-msgpack + + + org.xerial.snappy + snappy-java + \ No newline at end of file diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/DataChangeTrigger.java b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/DataChangeTrigger.java index f23bba1..49bece5 100644 --- a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/DataChangeTrigger.java +++ b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/DataChangeTrigger.java @@ -1,6 +1,6 @@ package net.sopod.soim.router.datasync; -import net.sopod.soim.router.datasync.server.SyncLog; +import net.sopod.soim.router.datasync.server.data.SyncLog; import java.util.Queue; import java.util.concurrent.ConcurrentHashMap; diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/SyncService.java b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/SyncService.java new file mode 100644 index 0000000..cbab3d8 --- /dev/null +++ b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/SyncService.java @@ -0,0 +1,31 @@ +package net.sopod.soim.router.datasync; + +import net.sopod.soim.common.util.Jackson; +import net.sopod.soim.router.cache.RouterUser; +import net.sopod.soim.router.cache.RouterUserStorage; + +import java.util.Map; + +/** + * SyncService + * + * @author tmy + * @date 2022-05-08 09:42 + */ +public class SyncService { + + public void fullSync() { + Map routerUserMap = RouterUserStorage.getInstance().getRouterUserMap(); + for (Map.Entry entry : routerUserMap.entrySet()) { + Long key = entry.getKey(); + RouterUser value = entry.getValue(); + Jackson.json().serialize(value); + + } + } + + public static void main(String[] args) { + + } + +} diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncClient.java b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncClient.java index cbb6acc..20838be 100644 --- a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncClient.java +++ b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncClient.java @@ -1,9 +1,17 @@ package net.sopod.soim.router.datasync.server; 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.util.AttributeKey; +import net.sopod.soim.router.datasync.server.codec.SyncCmdCodec; +import net.sopod.soim.router.datasync.server.codec.SyncLogEncoder; +import net.sopod.soim.router.datasync.server.data.SyncCmd; +import net.sopod.soim.router.datasync.server.data.SyncLog; +import net.sopod.soim.router.datasync.server.handler.SyncCmdClientHandler; +import net.sopod.soim.router.datasync.server.handler.SyncLogClientHandler; /** * SyncClient @@ -13,18 +21,52 @@ import io.netty.channel.socket.SocketChannel; */ public class SyncClient { + private NioEventLoopGroup group; + + private final Bootstrap bootstrap; + + private Channel clientChannel; + public SyncClient() { - NioEventLoopGroup group = new NioEventLoopGroup(2); - new Bootstrap() - .group(group) + this.bootstrap = new Bootstrap() .handler(new ChannelInitializer() { @Override protected void initChannel(SocketChannel channel) throws Exception { channel.pipeline() - .addLast(new SyncLogDataCodec()); + .addLast(new SyncCmdCodec()) + .addLast(new SyncLogEncoder()) + .addLast(new SyncCmdClientHandler()) + .addLast(new SyncLogClientHandler()); } }); } + public void connect(String host, int port) throws InterruptedException { + this.group = new NioEventLoopGroup(2); + this.clientChannel = bootstrap + .group(group) + .connect(host, port) + .await().channel(); + } + + public void write(SyncCmd syncCmd) { + clientChannel.writeAndFlush(syncCmd); + } + + public void write(SyncLog syncLog) { + clientChannel.write(syncLog); + } + + public void setAttr(AttributeKey key, T val) { + clientChannel.attr(key).set(val); + } + + public void getAttr(AttributeKey key) { + clientChannel.attr(key).get(); + } + + public void close() { + this.group.shutdownGracefully(); + } } diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncLogDataCodec.java b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncLogDataCodec.java deleted file mode 100644 index 3d86700..0000000 --- a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncLogDataCodec.java +++ /dev/null @@ -1,50 +0,0 @@ -package net.sopod.soim.router.datasync.server; - -import io.netty.buffer.ByteBuf; -import io.netty.channel.*; -import io.netty.handler.codec.ByteToMessageDecoder; -import io.netty.handler.codec.MessageToByteEncoder; - -import java.util.List; - -/** - * SyncDataInboundHandler - * - * @author tmy - * @date 2022-05-05 10:20 - */ -public class SyncLogDataCodec - extends CombinedChannelDuplexHandler { - - public SyncLogDataCodec() { - super(new SyncDataDecoder(), new SyncDataEncoder()); - } - - /** - * 解码器 - */ - public static class SyncDataDecoder extends ByteToMessageDecoder { - - @Override - protected void decode(ChannelHandlerContext channelHandlerContext, ByteBuf byteBuf, List list) throws Exception { - SyncLog syncLog = SyncLog.read(byteBuf); - list.add(syncLog); - } - - } - - /** - * 编码器 - */ - public static class SyncDataEncoder extends MessageToByteEncoder { - - @Override - protected void encode(ChannelHandlerContext channelHandlerContext, SyncLog syncLog, ByteBuf byteBuf) throws Exception { - byte[] bytes = syncLog.toBytes(); - byteBuf.writeBytes(bytes); - } - - } - -} diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncServer.java b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncServer.java index 2404790..979e7f1 100644 --- a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncServer.java +++ b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncServer.java @@ -10,6 +10,10 @@ import io.netty.channel.socket.nio.NioServerSocketChannel; import io.netty.handler.logging.LogLevel; import io.netty.handler.logging.LoggingHandler; import net.sopod.soim.common.util.netty.FastThreadLocalThreadFactory; +import net.sopod.soim.router.datasync.server.codec.SyncCmdCodec; +import net.sopod.soim.router.datasync.server.codec.SyncLogEncoder; +import net.sopod.soim.router.datasync.server.handler.SyncCmdServerHandler; +import net.sopod.soim.router.datasync.server.handler.SyncLogServerHandler; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -25,16 +29,14 @@ public class SyncServer { private static final Logger logger = LoggerFactory.getLogger(SyncServer.class); - private final int port; - private NioEventLoopGroup boss; private NioEventLoopGroup worker; - public SyncServer(int port) { - this.port = port; + public SyncServer() { + } - public void start(Consumer onFail) throws InterruptedException { + public void start(int port, Consumer onFail) { this.boss = new NioEventLoopGroup(1, new FastThreadLocalThreadFactory("sync-server-boss-%d", Thread.NORM_PRIORITY)); this.worker = new NioEventLoopGroup(2, new FastThreadLocalThreadFactory("sync-server-worker-%d", Thread.NORM_PRIORITY)); ServerBootstrap serverBoot = new ServerBootstrap() @@ -49,8 +51,10 @@ public class SyncServer { : logger.isErrorEnabled() ? LogLevel.ERROR : LogLevel.INFO; ChannelPipeline pipeline = channel.pipeline(); pipeline.addLast(new LoggingHandler(logLevel)) - .addLast(new SyncLogDataCodec()) - .addLast(new SyncServerInboundHandler()); + .addLast(new SyncCmdCodec()) + .addLast(new SyncLogEncoder()) + .addLast(new SyncCmdServerHandler()) + .addLast(new SyncLogServerHandler()); } }); serverBoot.bind(port).addListener((ChannelFutureListener) future -> { @@ -73,4 +77,13 @@ public class SyncServer { } } + public static void main(String[] args) { + new SyncServer() + .start(9999, err -> { + logger.error("服务启动失败"); + }); + + + } + } diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncServerInboundHandler.java b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncServerInboundHandler.java deleted file mode 100644 index bfc1f1c..0000000 --- a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncServerInboundHandler.java +++ /dev/null @@ -1,33 +0,0 @@ -package net.sopod.soim.router.datasync.server; - -import io.netty.channel.ChannelHandlerContext; -import io.netty.channel.ChannelInboundHandlerAdapter; - -/** - * SyncLogInboundHandler - * - * @author tmy - * @date 2022-05-05 14:57 - */ -public class SyncServerInboundHandler extends ChannelInboundHandlerAdapter { - - public SyncServerInboundHandler() { - - } - - @Override - public void channelActive(ChannelHandlerContext ctx) throws Exception { - super.channelActive(ctx); - } - - @Override - public void channelInactive(ChannelHandlerContext ctx) throws Exception { - super.channelInactive(ctx); - } - - @Override - public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) throws Exception { - super.exceptionCaught(ctx, cause); - } - -} diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/codec/CodecUtil.java b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/codec/CodecUtil.java new file mode 100644 index 0000000..5183d8b --- /dev/null +++ b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/codec/CodecUtil.java @@ -0,0 +1,50 @@ +package net.sopod.soim.router.datasync.server.codec; + +import net.sopod.soim.common.util.Jackson; +import org.xerial.snappy.Snappy; + +import java.io.IOException; + +/** + * CodecUtil + * + * @author tmy + * @date 2022-05-08 16:54 + */ +public class CodecUtil { + + public static byte[] encode(Object data) { + return Jackson.msgpack().serializeBytes(data); + } + + public static byte[] compress(byte[] bytes) { + try { + return Snappy.compress(bytes); + } catch (IOException e) { + throw new RuntimeException(e); + } + } + + public static byte[] codecAndCompress(Object data) { + byte[] bytes = encode(data); + return compress(bytes); + } + + public static T decode(byte[] bytes, Class type) { + return Jackson.msgpack().deserializeBytes(bytes, type); + } + + public static byte[] uncompress(byte[] bytes) { + try { + return Snappy.uncompress(bytes); + } catch (IOException e) { + throw new RuntimeException(e); + } + } + + public static T uncompressAndDecode(byte[] bytes, Class type) { + byte[] uncompress = uncompress(bytes); + return decode(uncompress, type); + } + +} diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/codec/SyncCmdCodec.java b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/codec/SyncCmdCodec.java new file mode 100644 index 0000000..ad70035 --- /dev/null +++ b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/codec/SyncCmdCodec.java @@ -0,0 +1,61 @@ +package net.sopod.soim.router.datasync.server.codec; + +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.router.datasync.server.data.SyncCmd; +import net.sopod.soim.router.datasync.server.data.SyncLog; + +import java.util.List; + +/** + * SyncChannelHandler + * SyncCmd 编解码器 + * + * @author tmy + * @date 2022-05-07 21:37 + */ +public class SyncCmdCodec extends CombinedChannelDuplexHandler { + + public SyncCmdCodec() { + super(new SyncCmdDecoder(), new SyncCmdEncoder()); + } + + public static class SyncCmdDecoder extends ByteToMessageDecoder { + @Override + protected void decode(ChannelHandlerContext channelHandlerContext, ByteBuf byteBuf, List list) throws Exception { + SyncCmd syncCmd = SyncCmd.read(byteBuf); + if (!SyncCmd.SYNC_LOG.equals(syncCmd.getCmdType())) { + list.add(syncCmd); + return; + } + // 数据同步请求命令 + // 同步数据请求体长度 + int syncLogBytesLen = byteBuf.readInt(); + // 读取、解压、解析字节数据 + byte[] syncLogCompressBytes = new byte[syncLogBytesLen]; + byteBuf.readBytes(syncLogCompressBytes); + byte[] syncLogBytes = CodecUtil.uncompress(syncLogCompressBytes); + ByteBuf syncLogByteBuf = Unpooled.wrappedBuffer(syncLogBytes); + try { + SyncLog syncLog = SyncLog.read(syncLogByteBuf); + list.add(syncLog); + } finally { + syncLogByteBuf.release(); + } + } + } + + public static class SyncCmdEncoder extends MessageToByteEncoder { + @Override + protected void encode(ChannelHandlerContext channelHandlerContext, SyncCmd syncCmd, ByteBuf byteBuf) throws Exception { + byte[] bytes = syncCmd.toByte(); + byteBuf.writeBytes(bytes); + } + } + +} diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/codec/SyncLogEncoder.java b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/codec/SyncLogEncoder.java new file mode 100644 index 0000000..d505611 --- /dev/null +++ b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/codec/SyncLogEncoder.java @@ -0,0 +1,34 @@ +package net.sopod.soim.router.datasync.server.codec; + +import io.netty.buffer.ByteBuf; +import io.netty.channel.ChannelHandlerContext; +import io.netty.handler.codec.MessageToByteEncoder; +import net.sopod.soim.router.datasync.server.data.SyncCmd; +import net.sopod.soim.router.datasync.server.data.SyncLog; + +/** + * SyncDataInboundHandler + * SyncLog 编码器 + * + * @author tmy + * @date 2022-05-05 10:20 + */ +public class SyncLogEncoder extends MessageToByteEncoder { + + @Override + protected void encode(ChannelHandlerContext channelHandlerContext, SyncLog syncLog, ByteBuf byteBuf) throws Exception { + // 同步数据指令 + SyncCmd syncCmd = new SyncCmd() + .setCmdType(SyncCmd.SYNC_LOG); + byte[] syncCmdBytes = syncCmd.toByte(); + byteBuf.writeBytes(syncCmdBytes); + + // 同步数据内容 + byte[] bytes = syncLog.toBytes(); + // snappy 压缩一下 + byte[] compressLogBytes = CodecUtil.compress(bytes); + byteBuf.writeInt(compressLogBytes.length); + byteBuf.writeBytes(compressLogBytes); + } + +} diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/data/SyncCmd.java b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/data/SyncCmd.java new file mode 100644 index 0000000..21a3105 --- /dev/null +++ b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/data/SyncCmd.java @@ -0,0 +1,113 @@ +package net.sopod.soim.router.datasync.server.data; + +import io.netty.buffer.ByteBuf; +import lombok.Data; +import lombok.experimental.Accessors; +import org.apache.dubbo.common.io.Bytes; + +import java.nio.charset.StandardCharsets; +import java.util.concurrent.atomic.AtomicBoolean; + +/** + * SyncCmd + * + * @author tmy + * @date 2022-05-07 21:41 + */ +@Data +@Accessors(chain = true) +public class SyncCmd { + + private static final short MAGIC = 0x7a21; + + public static final int PING = 1; + + public static final int PONG = 2; + + /** + * 同步命令 + */ + public static final Integer SYNC = 3; + + /** + * 同步结束命令 + */ + public static final Integer SYNC_END = 4; + + /** + * SyncLog 推送命令 + */ + public static final Integer SYNC_LOG = 10; + + private Integer cmdType; + + private String param1; + + private String param2; + + private String param3; + + public byte[] toByte() { + byte[] param1Bytes = param1 == null ? new byte[0] : param1.getBytes(StandardCharsets.UTF_8); + byte[] param2Bytes = param2 == null ? new byte[0] : param2.getBytes(StandardCharsets.UTF_8); + byte[] param3Bytes = param3 == null ? new byte[0] : param3.getBytes(StandardCharsets.UTF_8); + byte[] bytes = new byte[2 + 4 + 4 + param1Bytes.length + 4 + param2Bytes.length + 4 + param3Bytes.length]; + + int offset = 0; + Bytes.short2bytes(MAGIC, bytes, offset); + offset += 2; + Bytes.int2bytes(cmdType, bytes, offset); + offset += 4; + Bytes.int2bytes(param1Bytes.length, bytes, offset); + offset += 4; + System.arraycopy(param1Bytes, 0, bytes, offset, param1Bytes.length); + offset += param1Bytes.length; + + Bytes.int2bytes(param2Bytes.length, bytes, offset); + offset += 4; + System.arraycopy(param2Bytes, 0, bytes, offset, param2Bytes.length); + offset += param2Bytes.length; + + Bytes.int2bytes(param3Bytes.length, bytes, offset); + offset += 4; + System.arraycopy(param3Bytes, 0, bytes, offset, param3Bytes.length); + offset += param3Bytes.length; + return bytes; + } + + public static SyncCmd read(ByteBuf buf) { + short magic = buf.readShort(); + if (magic != MAGIC) { + throw new IllegalStateException("unknown bytes magic error"); + } + SyncCmd cmd = new SyncCmd(); + cmd.cmdType = buf.readInt(); + int param1Len = buf.readInt(); + byte[] paramBytes = new byte[param1Len]; + if (param1Len > 0) { + buf.readBytes(paramBytes); + cmd.param1 = new String(paramBytes, StandardCharsets.UTF_8); + } + int param2Len = buf.readInt(); + if (param2Len > 0) { + paramBytes = paramBytes.length >= param2Len ? paramBytes : new byte[param2Len]; + buf.readBytes(paramBytes, 0, param2Len); + cmd.param2 = new String(paramBytes, 0, param2Len, StandardCharsets.UTF_8); + } + int param3Len = buf.readInt(); + if (param3Len > 0) { + paramBytes = paramBytes.length >= param3Len ? paramBytes : new byte[param3Len]; + buf.readBytes(paramBytes, 0, param3Len); + cmd.param3 = new String(paramBytes, 0, param3Len, StandardCharsets.UTF_8); + } + return cmd; + } + + public static void main(String[] args) { + AtomicBoolean locked = new AtomicBoolean(false); + System.out.println(locked.compareAndSet(false, false)); + System.out.println(locked.compareAndExchange(false, true)); + + } + +} diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncLog.java b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/data/SyncLog.java similarity index 92% rename from im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncLog.java rename to im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/data/SyncLog.java index 260786f..ed6f69b 100644 --- a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncLog.java +++ b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/data/SyncLog.java @@ -1,4 +1,4 @@ -package net.sopod.soim.router.datasync.server; +package net.sopod.soim.router.datasync.server.data; import io.netty.buffer.ByteBuf; import io.netty.buffer.Unpooled; @@ -12,8 +12,10 @@ import net.sopod.soim.router.datasync.SyncTypes; import org.apache.dubbo.common.io.Bytes; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import org.xerial.snappy.Snappy; import javax.annotation.Nullable; +import java.io.IOException; import java.io.Serializable; import java.lang.reflect.Method; import java.nio.charset.StandardCharsets; @@ -37,7 +39,7 @@ public class SyncLog implements Serializable { private static final long serialVersionUID = -8925525709590423526L; - private static final short MAGIC = 0x7a21; + // private static final short MAGIC = 0x7a21; public static final int OPT_ADD = 1; public static final int OPT_REMOVE = 2; @@ -87,7 +89,7 @@ public class SyncLog implements Serializable { protected String[] args; /** ================ 新增数据:序列化后的数据(避免修改) ===================== */ - protected List serializeDataCollect; + protected List serializeDataCollect; /** 数据可能同步给多个订阅者,缓存一下 */ private transient byte[] toBytesCache; @@ -125,8 +127,9 @@ public class SyncLog implements Serializable { int dataCollectByteLen = 0; if (dataCollectSize > 0) { int i = 0; - for (String dataCollect : serializeDataCollect) { - byteData[i] = dataCollect.getBytes(StandardCharsets.UTF_8); + for (byte[] dataCollect : serializeDataCollect) { + // byteData[i] = dataCollect.getBytes(StandardCharsets.UTF_8); + byteData[i] = dataCollect; dataCollectByteLen += 4; dataCollectByteLen += byteData[i].length; i++; @@ -134,8 +137,8 @@ public class SyncLog implements Serializable { } // 总长 + 同步数据类型(byte) + byte[] bytes = new byte[ - 2 // 魔术 - + 4 // 请求体总长度 + // 2 + // 魔术 + 4 // 请求体总长度 + 1 // 同步数据类型 SyncType + 1 // 新增/删除/更新 + 4 // logSeq 序列号 @@ -151,8 +154,8 @@ public class SyncLog implements Serializable { + dataCollectByteLen // 序列化数据字节(len,dataBytes,len,dataBytes...) ]; int offset = 0; - Bytes.short2bytes(MAGIC, bytes, offset); - offset += 2; +// Bytes.short2bytes(MAGIC, bytes, offset); +// offset += 2; Bytes.int2bytes(bytes.length - 6, bytes, offset); offset += 4; @@ -204,10 +207,6 @@ public class SyncLog implements Serializable { } public static SyncLog read(ByteBuf buf) { - short magic = buf.readShort(); - if (magic != MAGIC) { - throw new IllegalStateException("unknown bytes magic error"); - } // 后续bytes长度 int dataLen = buf.readInt(); int syncDataType = buf.readByte(); @@ -258,18 +257,17 @@ public class SyncLog implements Serializable { int dataSize = buf.readInt(); log.serializeDataCollect = dataSize > 0 ? new ArrayList<>(dataSize) : Collections.emptyList(); if (dataSize > 0) { - byte[] dataByte = new byte[0]; for (int i = 0; i < dataSize; i++) { int dataByteLen = buf.readInt(); - dataByte = dataByte.length >= dataByteLen ? dataByte : new byte[dataByteLen]; + byte[] dataByte = new byte[dataByteLen]; buf.readBytes(dataByte, 0, dataByteLen); - log.serializeDataCollect.add(new String(dataByte, 0, dataByteLen, StandardCharsets.UTF_8)); + log.serializeDataCollect.add(dataByte); } } return log; } - public static void main(String[] args) { + public static void main(String[] args) throws IOException { RouterUser user1 = new RouterUser() .setUid(10001L) .setAccount("前线") @@ -288,6 +286,7 @@ public class SyncLog implements Serializable { // System.out.println(json.getBytes(StandardCharsets.UTF_8).length); byte[] bytes = addLog.toBytes(); + byte[] compress = Snappy.compress(bytes); ByteBuf buf = Unpooled.wrappedBuffer(bytes); SyncLog newLog = SyncLog.read(buf); buf.release(); @@ -328,7 +327,7 @@ public class SyncLog implements Serializable { this.logSeq = logSeq; } - public AddLog setSerializeDataCollect(List serializeDataCollect) { + public AddLog setSerializeDataCollect(List serializeDataCollect) { this.serializeDataCollect = serializeDataCollect; return this; } @@ -337,7 +336,8 @@ public class SyncLog implements Serializable { if (this.serializeDataCollect == null) { this.serializeDataCollect = new ArrayList<>(); } - this.serializeDataCollect.add(Jackson.json().serialize(data)); + byte[] dataBytes = Jackson.msgpack().serializeBytes(data); + this.serializeDataCollect.add(dataBytes); return this; } } diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/handler/SyncCmdClientHandler.java b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/handler/SyncCmdClientHandler.java new file mode 100644 index 0000000..dc1a4e5 --- /dev/null +++ b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/handler/SyncCmdClientHandler.java @@ -0,0 +1,41 @@ +package net.sopod.soim.router.datasync.server.handler; + +import io.netty.channel.ChannelHandlerContext; +import io.netty.channel.SimpleChannelInboundHandler; +import net.sopod.soim.router.datasync.server.data.SyncCmd; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * SyncHandler + * + * @author tmy + * @date 2022-05-07 21:43 + */ +public class SyncCmdClientHandler extends SimpleChannelInboundHandler { + + private static final Logger logger = LoggerFactory.getLogger(SyncCmdClientHandler.class); + + @Override + protected void channelRead0(ChannelHandlerContext ctx, SyncCmd syncCmd) { + switch (syncCmd.getCmdType()) { + case SyncCmd.PING: + this.handlePing(ctx, syncCmd); + break; + case SyncCmd.PONG: + this.handlePong(ctx, syncCmd); + break; + } + } + + private void handlePing(ChannelHandlerContext ctx, SyncCmd syncCmd) { + logger.info("ping: {}", ctx.channel()); + SyncCmd pong = new SyncCmd().setCmdType(SyncCmd.PONG); + ctx.writeAndFlush(pong); + } + + private void handlePong(ChannelHandlerContext ctx, SyncCmd syncCmd) { + logger.info("pong: {}", ctx.channel()); + } + +} diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/handler/SyncCmdServerHandler.java b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/handler/SyncCmdServerHandler.java new file mode 100644 index 0000000..38d4396 --- /dev/null +++ b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/handler/SyncCmdServerHandler.java @@ -0,0 +1,20 @@ +package net.sopod.soim.router.datasync.server.handler; + +import io.netty.channel.ChannelHandlerContext; +import io.netty.channel.SimpleChannelInboundHandler; +import net.sopod.soim.router.datasync.server.data.SyncCmd; + +/** + * SyncCmdServerHandler + * + * @author tmy + * @date 2022-05-08 18:48 + */ +public class SyncCmdServerHandler extends SimpleChannelInboundHandler { + + @Override + protected void channelRead0(ChannelHandlerContext channelHandlerContext, SyncCmd syncCmd) throws Exception { + + } + +} diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/handler/SyncLogClientHandler.java b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/handler/SyncLogClientHandler.java new file mode 100644 index 0000000..ec92d7e --- /dev/null +++ b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/handler/SyncLogClientHandler.java @@ -0,0 +1,21 @@ +package net.sopod.soim.router.datasync.server.handler; + +import io.netty.channel.ChannelHandlerContext; +import io.netty.channel.SimpleChannelInboundHandler; +import net.sopod.soim.router.datasync.server.data.SyncLog; + +/** + * SyncLogHandler + * 其他节点会推送同步数据到当前新增的节点,这里接收数据 + * + * @author tmy + * @date 2022-05-08 17:20 + */ +public class SyncLogClientHandler extends SimpleChannelInboundHandler { + + @Override + protected void channelRead0(ChannelHandlerContext channelHandlerContext, SyncLog syncLog) throws Exception { + + } + +} diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/handler/SyncLogServerHandler.java b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/handler/SyncLogServerHandler.java new file mode 100644 index 0000000..ff02c64 --- /dev/null +++ b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/handler/SyncLogServerHandler.java @@ -0,0 +1,21 @@ +package net.sopod.soim.router.datasync.server.handler; + +import io.netty.channel.ChannelHandlerContext; +import io.netty.channel.SimpleChannelInboundHandler; +import net.sopod.soim.router.datasync.server.data.SyncLog; + +/** + * SyncLogServerHandler + * 即将删除的节点:会主动推送同步数据,这里进行接收 + * + * @author tmy + * @date 2022-05-08 18:48 + */ +public class SyncLogServerHandler extends SimpleChannelInboundHandler { + + @Override + protected void channelRead0(ChannelHandlerContext channelHandlerContext, SyncLog syncLog) throws Exception { + + } + +} diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/session/SyncSession.java b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/session/SyncSession.java new file mode 100644 index 0000000..a65d28a --- /dev/null +++ b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/session/SyncSession.java @@ -0,0 +1,56 @@ +package net.sopod.soim.router.datasync.server.session; + +import net.sopod.soim.router.datasync.server.data.SyncCmd; +import net.sopod.soim.router.datasync.server.SyncServer; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * SyncSession + * 作为客户端或服务端主动操作连接 + * + * @author tmy + * @date 2022-05-08 17:34 + */ +public class SyncSession { + private static final Logger logger = LoggerFactory.getLogger(SyncSession.class); + + public static SyncSession client() { + return null; + } + + public static SyncSession server() { + + return null; + } + + public static class SyncSessionServer { + + private final SyncServer syncServer; + + public SyncSessionServer(int port) { + this.syncServer = new SyncServer(); + this.syncServer.start(9999, err -> { + logger.error("sync server start fail!", err); + }); + } + + public void writeTo() { + + } + + } + + public static class SyncSessionClient { + + public SyncSessionClient() { + + } + + public void write(SyncCmd syncCmd) { + + } + + } + +} diff --git a/pom.xml b/pom.xml index 5ea5c47..bdb08c8 100644 --- a/pom.xml +++ b/pom.xml @@ -41,6 +41,7 @@ 4.2.0 5.1.0 2.0.3 + 0.9.1 @@ -254,6 +255,16 @@ cglib 3.3.0 + + org.msgpack + jackson-dataformat-msgpack + ${msgpack.version} + + + org.xerial.snappy + snappy-java + 1.1.8.4 +