From 9daf94e2a22b564f377944f85bc3ba69acadb59b Mon Sep 17 00:00:00 2001 From: tangmingyou <234767776@qq.com> Date: Thu, 5 May 2022 17:57:19 +0800 Subject: [PATCH] =?UTF-8?q?im-router=20=E6=95=B0=E6=8D=AE=E5=90=8C?= =?UTF-8?q?=E6=AD=A5?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../net/sopod/soim/common/util/Jackson.java | 4 + .../net/sopod/soim/router/cache/DataSync.java | 10 +- .../router/cache/DataSyncProxyFactory.java | 141 ++++++++++-- .../soim/router/cache/SoImUserCache.java | 13 +- .../soim/router/cache/annotation/Sync.java | 19 ++ .../{DataSyncIgnore.java => SyncIgnore.java} | 7 +- .../router/config/ImRouterAppOnReady.java | 1 + .../datasync/server/DataChangeTrigger.java | 42 ++++ .../router/datasync/server/SyncClient.java | 30 +++ .../soim/router/datasync/server/SyncLog.java | 200 ++++++++++++++++++ .../datasync/server/SyncLogDataCodec.java | 50 +++++ .../server/SyncLogInboundHandler.java | 19 ++ .../router/datasync/server/SyncServer.java | 50 +++++ .../router/datasync/server/SyncTypes.java | 126 +++++++++++ 14 files changed, 689 insertions(+), 23 deletions(-) create mode 100644 im-service/im-router/src/main/java/net/sopod/soim/router/cache/annotation/Sync.java rename im-service/im-router/src/main/java/net/sopod/soim/router/cache/annotation/{DataSyncIgnore.java => SyncIgnore.java} (50%) create mode 100644 im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/DataChangeTrigger.java create mode 100644 im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncClient.java create mode 100644 im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncLog.java create mode 100644 im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncLogDataCodec.java create mode 100644 im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncLogInboundHandler.java create mode 100644 im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncServer.java create mode 100644 im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncTypes.java 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 5c9a753..9845667 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 @@ -176,6 +176,10 @@ public class Jackson { return XML_INSTANCE; } + public boolean canSerialize(Class clazz) { + return objectMapper.canSerialize(clazz); + } + public T deserialize(String content, Class valueType) { try { return objectMapper.readValue(content, valueType); diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/cache/DataSync.java b/im-service/im-router/src/main/java/net/sopod/soim/router/cache/DataSync.java index b566fdf..0e5a335 100644 --- a/im-service/im-router/src/main/java/net/sopod/soim/router/cache/DataSync.java +++ b/im-service/im-router/src/main/java/net/sopod/soim/router/cache/DataSync.java @@ -9,7 +9,10 @@ package net.sopod.soim.router.cache; public interface DataSync { /** 不是更新数据的方法开头 */ - String[] nonUpdateMethodStart = new String[]{"get", "select", "list"}; + String[] nonUpdateMethodStart = new String[]{"get", "find", "list", "is", "select"}; + + /** 不是更新数据的方法 */ + String[] nonUpdateMethod = new String[]{"toString", "equals", "hashCode", "wait", "getClass", "notify", "notifyAll"}; /** * 如果是更新方法会同步数据,到新增节点或备份节点 @@ -20,6 +23,11 @@ public interface DataSync { return false; } } + for (String method : nonUpdateMethod) { + if (method.equals(methodName)) { + return false; + } + } return true; } diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/cache/DataSyncProxyFactory.java b/im-service/im-router/src/main/java/net/sopod/soim/router/cache/DataSyncProxyFactory.java index 271b483..dff73fc 100644 --- a/im-service/im-router/src/main/java/net/sopod/soim/router/cache/DataSyncProxyFactory.java +++ b/im-service/im-router/src/main/java/net/sopod/soim/router/cache/DataSyncProxyFactory.java @@ -4,11 +4,21 @@ import net.sf.cglib.proxy.Enhancer; import net.sf.cglib.proxy.MethodInterceptor; import net.sf.cglib.proxy.MethodProxy; import net.sopod.soim.common.util.Collects; +import net.sopod.soim.common.util.Jackson; +import net.sopod.soim.router.cache.annotation.Sync; +import net.sopod.soim.router.cache.annotation.SyncIgnore; +import net.sopod.soim.router.datasync.server.DataChangeTrigger; +import net.sopod.soim.router.datasync.server.SyncTypes; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import java.io.Serializable; +import java.lang.reflect.Constructor; +import java.lang.reflect.InvocationTargetException; import java.lang.reflect.Method; +import java.util.*; +import java.util.concurrent.ConcurrentHashMap; +import java.util.stream.Collectors; /** * BiSyncProxyManager @@ -21,46 +31,139 @@ public class DataSyncProxyFactory { private static final Logger logger = LoggerFactory.getLogger(DataSyncProxyFactory.class); + /** + * 缓存数据类型更新方法列表 + */ + private static final Map, Set> cacheTypeUpdaterMethods = new ConcurrentHashMap<>(); + @SuppressWarnings("unchecked") - public static T newProxyInstance(Class biSyncClazz) { + public static T newProxyInstance(SyncTypes.SyncType syncType) { + Class type = syncType.dataType(); + T instance; + try { + Constructor constructor = type.getDeclaredConstructor(); + constructor.setAccessible(true); + instance = constructor.newInstance(); + } catch (InstantiationException | InvocationTargetException | NoSuchMethodException | IllegalAccessException e) { + throw new IllegalStateException(type + "实例创建失败,没有无参构造函数", e); + } + + // 获取查询更新方法列表 + Set updaterMethods = cacheTypeUpdaterMethods.computeIfAbsent(type, dataType -> { + Method[] methods = dataType.getMethods(); + Set updaterMethodNames = new HashSet<>(); + for (Method m : methods) { + SyncIgnore syncIgnore = m.getDeclaredAnnotation(SyncIgnore.class); + String methodName = m.getName(); + if (syncIgnore != null) { + continue; + } + boolean isUpdater = m.getDeclaringClass() != Object.class + && instance.isUpdateMethod(methodName); + if (isUpdater) { + if (updaterMethodNames.contains(methodName)) { + throw new IllegalStateException(dataType + "重复的数据更新方法名" + methodName); + } + // 检查参数可序列化 + Class[] paramTypes = m.getParameterTypes(); + for (int i = 0; i < paramTypes.length; i++) { + if (!Jackson.json().canSerialize(paramTypes[i])) { + throw new IllegalStateException(String.format("类:%s 更新方法:%s 第%di个参数不可序列化", + instance.getClass().getName(), methodName, i + 1)); + } + } + updaterMethodNames.add(methodName); + } + } + logger.info("{} 更新方法列表: {}", dataType, updaterMethodNames); + return updaterMethodNames; + }); + + // 创建代理对象 Enhancer enhancer = new Enhancer(); - enhancer.setSuperclass(biSyncClazz); - enhancer.setCallback(new DataSyncProxyCallback()); + enhancer.setSuperclass(type); + enhancer.setCallback(new DataSyncProxyCallback<>(syncType, updaterMethods)); return (T) enhancer.create(); } - public static class DataSyncProxyCallback implements MethodInterceptor { + public static class DataSyncProxyCallback implements MethodInterceptor { + + private final SyncTypes.SyncType syncType; + + private final Set updaterMethods; + + public DataSyncProxyCallback(SyncTypes.SyncType syncType, Set updaterMethods) { + this.syncType = syncType; + this.updaterMethods = updaterMethods; + } @Override public Object intercept(Object instance, Method method, Object[] args, MethodProxy methodProxy) throws Throwable { String methodName = method.getName(); - // 是判断是否更新的方法跳过 - if ("isUpdateMethod".equals(methodName)) { + // 是判断不是更新的方法跳过 + if (null != method.getDeclaredAnnotation(SyncIgnore.class) + || !updaterMethods.contains(methodName) + || "isUpdateMethod".equals(methodName)) { return methodProxy.invokeSuper(instance, args); } - boolean isUpdateMethod = ((DataSync) instance).isUpdateMethod(methodName); - if (!isUpdateMethod) { - return methodProxy.invokeSuper(instance, args); - } - if (Collects.isNotEmpty(args)) { - // TODO 参数值可能为 null,检查参数可序列化放在生成代理对象时 - for (int i = 0; i < args.length; i++) { - if (!(args[i] instanceof Serializable)) { - logger.error("类:{} 更新方法:{} 第{}i个参数不可序列化", instance.getClass(), methodName, i+1); - } - } - } + // TODO 记录更新操作(方法和参数),查询数据订阅者,异步同步数据 + DataChangeTrigger.instance().onUpdate(syncType, syncType.getDataKey((T) instance), methodName, args); + System.out.println("intercept invoke:" + methodName); return methodProxy.invokeSuper(instance, args); } } + public static interface A { + default void setName(String name) { + System.out.println("setName: " + name); + } + + } + + public static class B implements A { + private int age; + + public void setAge(int age) { + this.age = age; + } + + @Override + public void setName(String name) { + System.out.println("over setName: " + name); + } + } + public static void main(String[] args) { - RouterUser routerUser = newProxyInstance(RouterUser.class); + + RouterUser routerUser = newProxyInstance(SyncTypes.ROUTER_USER); routerUser.setAccount("日月光"); System.out.println(routerUser.getAccount()); + +// Enhancer enhancer = new Enhancer(); +// enhancer.setSuperclass(B.class); +// enhancer.setCallback(new MethodInterceptor() { +// @Override +// public Object intercept(Object o, Method method, Object[] args, MethodProxy methodProxy) throws Throwable { +// System.out.println("proxy method: " + method.getName()); +// return methodProxy.invokeSuper(o, args); +// } +// }); +// +// B b = (B) enhancer.create(); +// b.setAge(12); +// b.setName("沧海"); +// +// System.out.println(b.age); +// +// for (Method method : B.class.getMethods()) { +// System.out.println(method.getName() + "-" + method.hashCode() + ": " + method.getDeclaringClass()); +// } +// System.out.println(Serializable.class.isAssignableFrom(String.class)); +// System.out.println(Serializable.class.isAssignableFrom(Integer.class)); +// System.out.println(Integer.class.isAssignableFrom(Serializable.class)); } } diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/cache/SoImUserCache.java b/im-service/im-router/src/main/java/net/sopod/soim/router/cache/SoImUserCache.java index 468e2a5..6ac6968 100644 --- a/im-service/im-router/src/main/java/net/sopod/soim/router/cache/SoImUserCache.java +++ b/im-service/im-router/src/main/java/net/sopod/soim/router/cache/SoImUserCache.java @@ -1,5 +1,8 @@ package net.sopod.soim.router.cache; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + import java.util.Map; import java.util.concurrent.ConcurrentHashMap; @@ -12,6 +15,8 @@ import java.util.concurrent.ConcurrentHashMap; */ public class SoImUserCache { + private static final Logger logger = LoggerFactory.getLogger(SoImUserCache.class); + private static final SoImUserCache INSTANCE = new SoImUserCache(); private final ConcurrentHashMap routerUserMap = new ConcurrentHashMap<>(128); @@ -20,14 +25,18 @@ public class SoImUserCache { return INSTANCE; } - public void put(Long uid, RouterUser routerUser) { - routerUserMap.put(uid, routerUser); + public RouterUser put(Long uid, RouterUser routerUser) { + return routerUserMap.put(uid, routerUser); } public RouterUser get(Long uid) { return routerUserMap.get(uid); } + public RouterUser remove(Long uid) { + return routerUserMap.remove(uid); + } + public Map getRouterUserMap() { return routerUserMap; } diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/cache/annotation/Sync.java b/im-service/im-router/src/main/java/net/sopod/soim/router/cache/annotation/Sync.java new file mode 100644 index 0000000..65c0f5e --- /dev/null +++ b/im-service/im-router/src/main/java/net/sopod/soim/router/cache/annotation/Sync.java @@ -0,0 +1,19 @@ +package net.sopod.soim.router.cache.annotation; + +import java.lang.annotation.*; + +/** + * Sync + * 修改数据的方法操作 + * + * @author tmy + * @date 2022-05-05 17:18 + */ +@Documented +@Retention(RetentionPolicy.RUNTIME) +@Target({ElementType.METHOD}) +public @interface Sync { + + + +} diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/cache/annotation/DataSyncIgnore.java b/im-service/im-router/src/main/java/net/sopod/soim/router/cache/annotation/SyncIgnore.java similarity index 50% rename from im-service/im-router/src/main/java/net/sopod/soim/router/cache/annotation/DataSyncIgnore.java rename to im-service/im-router/src/main/java/net/sopod/soim/router/cache/annotation/SyncIgnore.java index 9dfa142..26333cd 100644 --- a/im-service/im-router/src/main/java/net/sopod/soim/router/cache/annotation/DataSyncIgnore.java +++ b/im-service/im-router/src/main/java/net/sopod/soim/router/cache/annotation/SyncIgnore.java @@ -1,5 +1,7 @@ package net.sopod.soim.router.cache.annotation; +import java.lang.annotation.*; + /** * BiSyncIgnore * 忽略方法同步 @@ -7,6 +9,9 @@ package net.sopod.soim.router.cache.annotation; * @author tmy * @date 2022-05-04 17:39 */ -public @interface DataSyncIgnore { +@Documented +@Retention(RetentionPolicy.RUNTIME) +@Target({ElementType.METHOD}) +public @interface SyncIgnore { } diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/config/ImRouterAppOnReady.java b/im-service/im-router/src/main/java/net/sopod/soim/router/config/ImRouterAppOnReady.java index ead6b06..9acf3ae 100644 --- a/im-service/im-router/src/main/java/net/sopod/soim/router/config/ImRouterAppOnReady.java +++ b/im-service/im-router/src/main/java/net/sopod/soim/router/config/ImRouterAppOnReady.java @@ -81,6 +81,7 @@ public class ImRouterAppOnReady implements ApplicationListener logQueue = new ConcurrentLinkedQueue<>(); + + public void onUpdate(SyncTypes.SyncType syncType, String dataKey, String method, Object[] args) { + // 序列化 args,避免后续更改 + // logQueue.add() + } + + public void onAdd(SyncTypes.SyncType syncType, T data) { + // 序列化 data,避免后续更改 + + } + + public void onRemove(SyncTypes.SyncType syncType, String dataKey) { + + } + +} 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 new file mode 100644 index 0000000..cbb6acc --- /dev/null +++ b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncClient.java @@ -0,0 +1,30 @@ +package net.sopod.soim.router.datasync.server; + +import io.netty.bootstrap.Bootstrap; +import io.netty.channel.ChannelInitializer; +import io.netty.channel.nio.NioEventLoopGroup; +import io.netty.channel.socket.SocketChannel; + +/** + * SyncClient + * + * @author tmy + * @date 2022-05-05 15:03 + */ +public class SyncClient { + + public SyncClient() { + NioEventLoopGroup group = new NioEventLoopGroup(2); + new Bootstrap() + .group(group) + .handler(new ChannelInitializer() { + @Override + protected void initChannel(SocketChannel channel) throws Exception { + channel.pipeline() + .addLast(new SyncLogDataCodec()); + } + }); + } + + +} 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/SyncLog.java new file mode 100644 index 0000000..a411674 --- /dev/null +++ b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncLog.java @@ -0,0 +1,200 @@ +package net.sopod.soim.router.datasync.server; + +import io.netty.buffer.ByteBuf; +import lombok.Data; +import lombok.experimental.Accessors; +import net.sopod.soim.common.util.Jackson; +import org.apache.dubbo.common.io.Bytes; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import javax.annotation.Nullable; +import java.io.Serializable; +import java.lang.reflect.Method; +import java.nio.charset.StandardCharsets; + +/** + * SyncLog + * syncLog serialize + * + * @author tmy + * @date 2022-05-05 11:37 + */ +@Data +@Accessors(chain = true) +public class SyncLog implements Serializable { + + private static final Logger logger = LoggerFactory.getLogger(SyncLog.class); + + private static final long serialVersionUID = -8925525709590423526L; + + private static final short MAGIC = 0x7a21; + + public static final int OPT_ADD = 1; + public static final int OPT_REMOVE = 2; + public static final int OPT_UPDATE = 3; + + /** 日志序列号保证顺序 */ + private int logSeq; + + /** + * 操作类型: + * 1.新增 + * 2.删除 + * 3.更新 + */ + private int operateType; + + /** + * {@link SyncTypes} ordinal + */ + private int syncDataType; + + /** 数据id标示 */ + private String dataKey; + + /** ================ 新增参数:序列化后的数据 ===================== */ + private String addSerializeData; + + /** ================ 更新参数 ===================== */ + private String clazz; + + private String method; + + private Object[] args; + + public byte[] toBytes() { + byte[] clazzBytes = clazz.getBytes(); + byte[] methodBytes = method.getBytes(); + int argSize = args == null ? 0 : args.length; + + byte[][] byteArgs = new byte[argSize][]; + if (args != null) { + for (int i = 0; i < args.length; i++) { + // 反序列化时根据方法参数类型json反序列化 + // TODO null + String argJson = Jackson.json().serialize(args[i]); + byteArgs[i] = argJson.getBytes(StandardCharsets.UTF_8); + } + } + int argsByteLen = 0; + for (byte[] byteArg : byteArgs) { + argsByteLen += 4; + argsByteLen += byteArg.length; + } + // 总长 + 同步数据类型(byte) + + byte[] bytes = new byte[ + 2 // 魔术 + + 4 // 请求体总长度 + + 1 // 同步数据类型 + + 4 // clazz 字节长度 + + clazzBytes.length // clazz字节 + + 4 // method 字节长度 + + methodBytes.length // method 字节 + + 4 // args参数个数 + + argsByteLen // args参数字节 + ]; + int offset = 0; + Bytes.short2bytes(MAGIC, bytes, offset); + offset += 2; + + Bytes.int2bytes(bytes.length - 6, bytes, offset); + offset += 4; + + bytes[offset] = (byte)syncDataType; + offset += 1; + + Bytes.int2bytes(clazzBytes.length, bytes, offset); + offset += 4; + System.arraycopy(clazzBytes, 0, bytes, offset, clazzBytes.length); + offset += clazzBytes.length; + + Bytes.int2bytes(methodBytes.length, bytes, offset); + offset += 4; + System.arraycopy(methodBytes, 0, bytes, offset, methodBytes.length); + offset += methodBytes.length; + + Bytes.int2bytes(argSize, bytes, offset); + offset += 4; + if (argsByteLen > 0) { + for (byte[] byteArg : byteArgs) { + Bytes.int2bytes(byteArg.length, bytes, offset); + offset += 4; + System.arraycopy(byteArg, 0, bytes, offset, byteArg.length); + offset += byteArg.length; + } + } + return bytes; + } + + public static SyncLog read(ByteBuf buf) { + short magic = buf.readShort(); + if (magic != MAGIC) { + throw new IllegalStateException("unknown bytes magic error"); + } + SyncLog log = new SyncLog(); + // 后续bytes长度 + int dataLen = buf.readInt(); + log.syncDataType = buf.readByte(); + int clazzLen = buf.readInt(); + byte[] clazzBytes = new byte[clazzLen]; + buf.readBytes(clazzBytes); + log.clazz = new String(clazzBytes, StandardCharsets.UTF_8); + int methodLen = buf.readInt(); + byte[] methodBytes = clazzLen >= methodLen ? clazzBytes : new byte[methodLen]; + buf.readBytes(methodBytes, 0, methodLen); + log.method = new String(methodBytes, 0, methodLen, StandardCharsets.UTF_8); + int argSize = buf.readInt(); + log.args = new Object[argSize]; + if (argSize > 0) { + Method method = getClassMethod(log.clazz, log.method); + if (method == null) { + throw new IllegalStateException("类" + log.clazz + "方法" + log.method + "未找到"); + } + Class[] paramTypes = method.getParameterTypes(); + if (paramTypes.length != argSize) { + throw new IllegalStateException("类" + log.clazz + "方法" + log.method + "指定参数" + argSize + "个,查到参数" + paramTypes.length + "个"); + } + // 复用 bytes + byte[] argBytes = new byte[0]; + for (int i = 0; i < argSize; i++) { + int argLen = buf.readInt(); + argBytes = argBytes.length >= argLen ? argBytes : new byte[argLen]; + buf.readBytes(argBytes, 0, argLen); + String argJson = new String(argBytes, 0, argLen, StandardCharsets.UTF_8); + Object arg = Jackson.json().deserialize(argJson, paramTypes[i]); + log.args[i] = arg; + } + } + return log; + } + + /** + * TODO 缓存反射结果 + */ + @Nullable + private static Method getClassMethod(String clazzName, String methodName) { + try { + Class clazz = Class.forName(clazzName); + Method[] methods = clazz.getDeclaredMethods(); + for (Method method : methods) { + if (method.getName().equals(methodName)) { + return method; + } + } + } catch (ClassNotFoundException e) { + logger.error("获取类型方法失败: {}.{}", clazzName, methodName, e); + } + return null; + } + + public static void main(String[] args) { + new SyncLog() + .setSyncDataType(SyncTypes.ROUTER_USER.ordinal()) + .setClazz("") + .setMethod("") + .setArgs(new Object[]{}); + + } + +} 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 new file mode 100644 index 0000000..3d86700 --- /dev/null +++ b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncLogDataCodec.java @@ -0,0 +1,50 @@ +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/SyncLogInboundHandler.java b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncLogInboundHandler.java new file mode 100644 index 0000000..60c6a04 --- /dev/null +++ b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncLogInboundHandler.java @@ -0,0 +1,19 @@ +package net.sopod.soim.router.datasync.server; + +import io.netty.channel.ChannelHandlerContext; +import io.netty.channel.SimpleChannelInboundHandler; + +/** + * SyncLogInboundHandler + * + * @author tmy + * @date 2022-05-05 14:57 + */ +public class SyncLogInboundHandler 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/SyncServer.java b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncServer.java new file mode 100644 index 0000000..c398230 --- /dev/null +++ b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncServer.java @@ -0,0 +1,50 @@ +package net.sopod.soim.router.datasync.server; + +import io.netty.bootstrap.ServerBootstrap; +import io.netty.channel.Channel; +import io.netty.channel.ChannelInitializer; +import io.netty.channel.ChannelPipeline; +import io.netty.channel.nio.NioEventLoopGroup; +import io.netty.channel.socket.nio.NioServerSocketChannel; +import io.netty.handler.logging.LogLevel; +import io.netty.handler.logging.LoggingHandler; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * SyncServer + * + * @author tmy + * @date 2022-05-05 10:02 + */ +public class SyncServer { + + private static final Logger logger = LoggerFactory.getLogger(SyncServer.class); + + public SyncServer() { + NioEventLoopGroup boss = new NioEventLoopGroup(1); + NioEventLoopGroup worker = new NioEventLoopGroup(2); + ServerBootstrap serverBoot = new ServerBootstrap() + .group(boss, worker) + .channel(NioServerSocketChannel.class) + .childHandler(new ChannelInitializer<>() { + @Override + protected void initChannel(Channel channel) throws Exception { + LogLevel logLevel = logger.isDebugEnabled() ? LogLevel.DEBUG + : logger.isInfoEnabled() ? LogLevel.INFO + : logger.isWarnEnabled() ? LogLevel.WARN + : logger.isErrorEnabled() ? LogLevel.ERROR : LogLevel.INFO; + ChannelPipeline pipeline = channel.pipeline(); + pipeline.addLast(new LoggingHandler(logLevel)) + .addLast(new SyncLogDataCodec()) + .addLast(new SyncLogInboundHandler()); + } + }); + serverBoot.bind(8080); + } + + public void start() { + + } + +} diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncTypes.java b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncTypes.java new file mode 100644 index 0000000..5629e37 --- /dev/null +++ b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncTypes.java @@ -0,0 +1,126 @@ +package net.sopod.soim.router.datasync.server; + +import net.sopod.soim.common.util.StringUtil; +import net.sopod.soim.router.cache.DataSync; +import net.sopod.soim.router.cache.RouterUser; +import net.sopod.soim.router.cache.SoImUserCache; + +import javax.annotation.Nullable; +import java.util.ArrayList; +import java.util.List; +import java.util.concurrent.ConcurrentHashMap; + +/** + * SyncTypeEnum + * 需要同步的数据类型 + * + * @author tmy + * @date 2022-05-05 10:40 + */ +public class SyncTypes { + + /** + * 模拟枚举根据顺序值获取数据类型对象 + */ + public static SyncType getSyncType(int ordinal) { + return SyncType.getSyncType(ordinal); + } + + /** + * 用户数据同步处理: + * 改: 代理对象方法监控(getData(dataKey), updateMethod(args)) + * 增: 增数据(data) + * 删: 删数据(dataKey) + */ + public static final SyncType ROUTER_USER = new SyncType<>(RouterUser.class, "ROUTER_USER") { + @Override + public String getDataKey(RouterUser data) { + return StringUtil.toString(data.getUid()); + } + + @Override + public RouterUser getData(String uid) { + return SoImUserCache.getInstance().get(Long.valueOf(uid)); + } + + @Override + public boolean addData(RouterUser data) { + SoImUserCache.getInstance().put(data.getUid(), data); + return true; + } + + @Override + public boolean removeData(String uid) { + return null != SoImUserCache.getInstance().remove(Long.valueOf(uid)); + } + }; + + public static abstract class SyncType { + private static final List> TYPES = new ArrayList<>(); + private static int ordinalSeq = Byte.MIN_VALUE; + + private final int ordinal; + private final Class dataType; + private final String name; + + SyncType(Class dataType, String name) { + this.ordinal = ordinalSeq++; + if (ordinal > Byte.MAX_VALUE) { + throw new IllegalStateException("类型超过255个,需修改协议"); + } + this.dataType = dataType; + this.name = name; + TYPES.add(this); + } + + public Class dataType() { + return dataType; + } + + public String name() { + return name; + } + + public int ordinal() { + return ordinal; + } + + public abstract String getDataKey(T data); + + public abstract T getData(String key); + + public abstract boolean addData(T data); + + public abstract boolean removeData(String key); + + @Nullable + @SuppressWarnings("unchecked") + static SyncType getSyncType(int ordinal) { + return (SyncType) (ordinal < TYPES.size() ? TYPES.get(ordinal) : null); + } + + @Override + public String toString() { + return "SyncType{" + + "ordinal=" + ordinal + + ", dataType=" + dataType + + ", name='" + name + '\'' + + '}' + "@" + Integer.toHexString(this.hashCode()); + } + + } + + public static void main(String[] args) { + System.out.println(ROUTER_USER.ordinal()); + System.out.println(SyncTypes.getSyncType(0)); + System.out.println(SyncTypes.getSyncType(1)); + System.out.println(SyncTypes.getSyncType(2)); + Class routerUserClass = SyncTypes.ROUTER_USER.dataType(); + + ConcurrentHashMap map = new ConcurrentHashMap<>(); + map.put("1", "2"); + System.out.println(map.remove("2")); + System.out.println(map.remove("1")); + } + +}