diff --git a/im-client/src/main/java/net/sopod/soim/client/config/ClientConfig.java b/im-client/src/main/java/net/sopod/soim/client/config/ClientConfig.java
index efbd92d..7b0a5c9 100644
--- a/im-client/src/main/java/net/sopod/soim/client/config/ClientConfig.java
+++ b/im-client/src/main/java/net/sopod/soim/client/config/ClientConfig.java
@@ -17,6 +17,6 @@ public class ClientConfig {
private String host = "127.0.0.1";
- private Integer port = 8087;
+ private Integer port = 8089;
}
diff --git a/im-common/pom.xml b/im-common/pom.xml
index 27e22d5..2cb898a 100644
--- a/im-common/pom.xml
+++ b/im-common/pom.xml
@@ -26,6 +26,12 @@
com.google.guava
guava
+
+ io.netty
+ netty-common
+ ${netty.version}
+ provided
+
\ No newline at end of file
diff --git a/im-common/src/main/java/net/sopod/soim/common/util/Collects.java b/im-common/src/main/java/net/sopod/soim/common/util/Collects.java
index f5e131e..8b1c9c8 100644
--- a/im-common/src/main/java/net/sopod/soim/common/util/Collects.java
+++ b/im-common/src/main/java/net/sopod/soim/common/util/Collects.java
@@ -21,7 +21,7 @@ public class Collects {
}
public static boolean isEmpty(@Nullable Object[] arr) {
- return arr == null || arr.length > 0;
+ return arr == null || arr.length == 0;
}
public static boolean isNotEmpty(@Nullable Object[] arr) {
diff --git a/im-entry/src/main/java/net/sopod/soim/entry/util/FastThreadLocalThreadFactory.java b/im-common/src/main/java/net/sopod/soim/common/util/netty/FastThreadLocalThreadFactory.java
similarity index 94%
rename from im-entry/src/main/java/net/sopod/soim/entry/util/FastThreadLocalThreadFactory.java
rename to im-common/src/main/java/net/sopod/soim/common/util/netty/FastThreadLocalThreadFactory.java
index bb0990c..c147f75 100644
--- a/im-entry/src/main/java/net/sopod/soim/entry/util/FastThreadLocalThreadFactory.java
+++ b/im-common/src/main/java/net/sopod/soim/common/util/netty/FastThreadLocalThreadFactory.java
@@ -1,4 +1,4 @@
-package net.sopod.soim.entry.util;
+package net.sopod.soim.common.util.netty;
import io.netty.util.concurrent.FastThreadLocalThread;
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
index fe91280..c509183 100644
--- 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
@@ -7,9 +7,8 @@ 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 net.sopod.soim.entry.util.FastThreadLocalThreadFactory;
+import net.sopod.soim.common.util.netty.FastThreadLocalThreadFactory;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
diff --git a/im-entry/src/main/java/net/sopod/soim/entry/server/ImEntryInitializer.java b/im-entry/src/main/java/net/sopod/soim/entry/server/ImEntryInitializer.java
index c3b558a..b3b70f1 100644
--- a/im-entry/src/main/java/net/sopod/soim/entry/server/ImEntryInitializer.java
+++ b/im-entry/src/main/java/net/sopod/soim/entry/server/ImEntryInitializer.java
@@ -20,7 +20,8 @@ public class ImEntryInitializer extends ChannelInitializer {
private static final Logger logger = LoggerFactory.getLogger(ImEntryInitializer.class);
@Override
- protected void initChannel(SocketChannel socketChannel) throws Exception {
+ protected void initChannel(SocketChannel socketChannel) {
+ logger.info("init channel: {}", socketChannel);
LogLevel logLevel = logger.isDebugEnabled() ? LogLevel.DEBUG : LogLevel.INFO;
ChannelPipeline pipeline = socketChannel.pipeline();
pipeline.addLast(new LoggingHandler(logLevel))
diff --git a/im-entry/src/main/java/net/sopod/soim/entry/worker/Worker.java b/im-entry/src/main/java/net/sopod/soim/entry/worker/Worker.java
index a930e39..338177a 100644
--- a/im-entry/src/main/java/net/sopod/soim/entry/worker/Worker.java
+++ b/im-entry/src/main/java/net/sopod/soim/entry/worker/Worker.java
@@ -4,7 +4,7 @@ import com.lmax.disruptor.*;
import com.lmax.disruptor.dsl.Disruptor;
import com.lmax.disruptor.dsl.ProducerType;
import net.sopod.soim.common.util.ImClock;
-import net.sopod.soim.entry.util.FastThreadLocalThreadFactory;
+import net.sopod.soim.common.util.netty.FastThreadLocalThreadFactory;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/cache/RouterUser.java b/im-service/im-router/src/main/java/net/sopod/soim/router/cache/RouterUser.java
index 19d5f9d..04e3895 100644
--- a/im-service/im-router/src/main/java/net/sopod/soim/router/cache/RouterUser.java
+++ b/im-service/im-router/src/main/java/net/sopod/soim/router/cache/RouterUser.java
@@ -1,8 +1,17 @@
package net.sopod.soim.router.cache;
import lombok.Data;
+import lombok.EqualsAndHashCode;
import lombok.experimental.Accessors;
import net.sopod.soim.router.datasync.DataSync;
+import net.sopod.soim.router.datasync.annotation.SyncIgnore;
+
+import java.lang.reflect.Field;
+import java.lang.reflect.Modifier;
+
+class A {
+ private String aName;
+}
/**
* RouterUser
@@ -10,9 +19,12 @@ import net.sopod.soim.router.datasync.DataSync;
* @author tmy
* @date 2022-04-28 11:11
*/
+@EqualsAndHashCode(callSuper = true)
@Data
@Accessors(chain = true)
-public class RouterUser implements DataSync {
+public class RouterUser extends A implements DataSync {
+
+ public static final int a = 1;
private long uid;
@@ -25,4 +37,30 @@ public class RouterUser implements DataSync {
private String imEntryAddr;
+ @SyncIgnore
+ public RouterUser setUid(long uid) {
+ this.uid = uid;
+ return this;
+ }
+
+ public static void main(String[] args) {
+
+ for (Field field : RouterUser.class.getDeclaredFields()) {
+ System.out.println(field.getName());
+ }
+ System.out.println("=============");
+ for (Field field : RouterUser.class.getDeclaredFields()) {
+ System.out.println(field.getName());
+ }
+ System.out.println("=============");
+ for (Field field : RouterUser.class.getSuperclass().getDeclaredFields()) {
+ System.out.println(field.getName());
+ }
+ System.out.println(RouterUser.class.getSuperclass().getSuperclass());
+ for (Field field : RouterUser.class.getSuperclass().getSuperclass().getDeclaredFields()) {
+ System.out.println(field.getName());
+ }
+ // Modifier.isFinal()
+ }
+
}
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 6ac6968..0b13dfb 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,10 @@
package net.sopod.soim.router.cache;
+import net.sf.cglib.proxy.Enhancer;
+import net.sopod.soim.common.util.StringUtil;
+import net.sopod.soim.router.datasync.DataChangeTrigger;
+import net.sopod.soim.router.datasync.DataSyncProxyFactory;
+import net.sopod.soim.router.datasync.SyncTypes;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -26,6 +31,15 @@ public class SoImUserCache {
}
public RouterUser put(Long uid, RouterUser routerUser) {
+ // TODO 这里克隆一个代理对象
+ if (Enhancer.isEnhanced(routerUser.getClass())) {
+ routerUserMap.put(uid, routerUser);
+ return routerUser;
+ }
+ RouterUser proxyRouterUser = DataSyncProxyFactory.newProxyInstance(SyncTypes.ROUTER_USER);
+
+ // 新增数据触发
+ DataChangeTrigger.instance().onAdd(SyncTypes.ROUTER_USER, routerUser);
return routerUserMap.put(uid, routerUser);
}
@@ -34,7 +48,11 @@ public class SoImUserCache {
}
public RouterUser remove(Long uid) {
- return routerUserMap.remove(uid);
+ if (uid != null) {
+ DataChangeTrigger.instance().onRemove(SyncTypes.ROUTER_USER, StringUtil.toString(uid));
+ return routerUserMap.remove(uid);
+ }
+ return null;
}
public Map getRouterUserMap() {
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 1029579..f23bba1 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
@@ -3,6 +3,7 @@ package net.sopod.soim.router.datasync;
import net.sopod.soim.router.datasync.server.SyncLog;
import java.util.Queue;
+import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentLinkedQueue;
import java.util.concurrent.atomic.AtomicInteger;
@@ -15,28 +16,62 @@ import java.util.concurrent.atomic.AtomicInteger;
*/
public class DataChangeTrigger {
- private static final DataChangeTrigger INSTANCE = new DataChangeTrigger();
+ private static DataChangeTrigger INSTANCE;
public static DataChangeTrigger instance() {
+ if (INSTANCE == null) {
+ synchronized (DataChangeTrigger.class) {
+ if (INSTANCE == null) {
+ INSTANCE = new DataChangeTrigger();
+ }
+ }
+ }
return INSTANCE;
}
- private AtomicInteger seqCounter = new AtomicInteger();
+ private final ConcurrentHashMap seqCounterMap = new ConcurrentHashMap<>(128);
- private Queue logQueue = new ConcurrentLinkedQueue<>();
+ private final Queue logQueue = new ConcurrentLinkedQueue<>();
public void onUpdate(SyncTypes.SyncType syncType, String dataKey, String method, Object[] args) {
// 序列化 args,避免后续更改
- // logQueue.add()
+ SyncLog.UpdateLog updateLog = SyncLog.updateLog(getSeq(dataKey), syncType)
+ .setDataKey(dataKey)
+ .setMethod(method)
+ .setArgs(args);
+ publishLog(updateLog);
}
public void onAdd(SyncTypes.SyncType syncType, T data) {
+ String dataKey = syncType.getDataKey(data);
+ AtomicInteger seqCounter = getSeqCounter(dataKey);
// 序列化 data,避免后续更改
-
+ SyncLog.AddLog addLog = SyncLog.addLog(seqCounter.getAndIncrement(), syncType)
+ .addData(data);
+ publishLog(addLog);
}
+ /**
+ * 数据删除日志
+ */
public void onRemove(SyncTypes.SyncType syncType, String dataKey) {
+ SyncLog.RemoveLog removeLog = SyncLog.removeLog(getSeq(dataKey), syncType)
+ .setDataKey(dataKey);
+ publishLog(removeLog);
+ }
+
+ private void publishLog(SyncLog log) {
+ // TODO 判断无订阅者跳过,新增节点同步按 dataKey 单独订阅每一个数据更新
+ logQueue.add(log);
+ System.out.println("publish log: "+log);
+ }
+
+ private int getSeq(String dataKey) {
+ return getSeqCounter(dataKey).getAndIncrement();
+ }
+ private AtomicInteger getSeqCounter(String dataKey) {
+ return seqCounterMap.computeIfAbsent(dataKey, key -> new AtomicInteger());
}
}
diff --git a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/DataSyncProxyFactory.java b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/DataSyncProxyFactory.java
index 05b098f..0a036a5 100644
--- a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/DataSyncProxyFactory.java
+++ b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/DataSyncProxyFactory.java
@@ -3,15 +3,15 @@ package net.sopod.soim.router.datasync;
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.ImClock;
import net.sopod.soim.common.util.Jackson;
import net.sopod.soim.router.cache.RouterUser;
import net.sopod.soim.router.datasync.annotation.SyncIgnore;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
-import java.lang.reflect.Constructor;
-import java.lang.reflect.InvocationTargetException;
-import java.lang.reflect.Method;
+import java.lang.reflect.*;
import java.util.*;
import java.util.concurrent.ConcurrentHashMap;
@@ -31,8 +31,12 @@ public class DataSyncProxyFactory {
*/
private static final Map, Set> typeUpdaterMethodsCache = new ConcurrentHashMap<>();
- @SuppressWarnings("unchecked")
public static T newProxyInstance(SyncTypes.SyncType syncType) {
+ return newProxyInstance(syncType, null);
+ }
+
+ @SuppressWarnings("unchecked")
+ public static T newProxyInstance(SyncTypes.SyncType syncType, T source) {
Class type = syncType.dataType();
T instance;
try {
@@ -78,7 +82,41 @@ public class DataSyncProxyFactory {
Enhancer enhancer = new Enhancer();
enhancer.setSuperclass(type);
enhancer.setCallback(new DataSyncProxyCallback<>(syncType, updaterMethods));
- return (T) enhancer.create();
+ T proxyObj = (T) enhancer.create();
+ if (source != null) {
+ try {
+ cloneFields(source, proxyObj);
+ } catch (Exception e) {
+ // logger.error("克隆对象失败:", e);
+ throw new IllegalStateException("克隆对象失败", e);
+ }
+ }
+ return proxyObj;
+ }
+
+ private static void cloneFields(T source, T proxyObj) throws Exception {
+ cloneFields0(source.getClass(), source, proxyObj);
+ }
+
+ /**
+ * 复制对象属性值
+ * TODO 缓存 避免每次clone从方法区查找
+ */
+ private static void cloneFields0(Class> clazz, T source, T target) throws Exception {
+ Field[] fields = clazz.getDeclaredFields();
+ if (Collects.isNotEmpty(fields)) {
+ for (Field field : fields) {
+ if (!Modifier.isFinal(field.getModifiers())) {
+ field.setAccessible(true);
+ field.set(target, field.get(source));
+ System.out.println(field.getName() + ":" + field.get(source));
+ }
+ }
+ }
+ Class> superClazz = clazz.getSuperclass();
+ if (superClazz != Object.class) {
+ cloneFields0(superClazz, source, target);
+ }
}
public static class DataSyncProxyCallback implements MethodInterceptor {
@@ -103,7 +141,7 @@ public class DataSyncProxyFactory {
return methodProxy.invokeSuper(instance, args);
}
- // TODO 记录更新操作(方法和参数),查询数据订阅者,异步同步数据
+ // 记录数据更新操作(方法和参数)
DataChangeTrigger.instance().onUpdate(syncType, syncType.getDataKey((T) instance), methodName, args);
System.out.println("intercept invoke:" + methodName);
@@ -112,54 +150,16 @@ public class DataSyncProxyFactory {
}
- 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 user1 = new RouterUser();
+ user1.setUid(10086L);
+ user1.setAccount("日月光");
+
+ RouterUser routerUser = newProxyInstance(SyncTypes.ROUTER_USER, user1);
+ System.out.println(routerUser);
+ routerUser.setOnlineTime(ImClock.millis());
+ System.out.println(routerUser);
- 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/datasync/server/SyncDataIncr.java b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncDataIncr.java
deleted file mode 100644
index e92a6b9..0000000
--- a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncDataIncr.java
+++ /dev/null
@@ -1,10 +0,0 @@
-package net.sopod.soim.router.datasync.server;
-
-/**
- * SyncDataIncr
- *
- * @author tmy
- * @date 2022-05-05 22:33
- */
-public class SyncDataIncr {
-}
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
index 5924443..260786f 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/SyncLog.java
@@ -43,8 +43,8 @@ public class SyncLog implements Serializable {
public static final int OPT_REMOVE = 2;
public static final int OPT_UPDATE = 3;
- public static AddLog addLog(SyncTypes.SyncType syncType) {
- return new AddLog<>(syncType);
+ public static AddLog addLog(int logSeq, SyncTypes.SyncType syncType) {
+ return new AddLog<>(logSeq, syncType);
}
public static UpdateLog updateLog(int logSeq, SyncTypes.SyncType syncType) {
@@ -79,7 +79,7 @@ public class SyncLog implements Serializable {
/** ================ 数据id标示:删除,更新用 ===================== */
protected String dataKey;
- /** ================ 更新数据:类,更新方法,更新方法序列化参数(避免修改) ===================== */
+ /** ================ 更新数据:类,更新方法,更新方法序列化后参数(避免修改) ===================== */
protected String clazz;
protected String method;
@@ -280,9 +280,9 @@ public class SyncLog implements Serializable {
.setAccount("画中")
.setImEntryAddr("127.0.0.2")
.setIsOnline(true);
- AddLog addLog = addLog(SyncTypes.ROUTER_USER)
- .addSerializeData(user1)
- .addSerializeData(user2);
+ AddLog addLog = addLog(0, SyncTypes.ROUTER_USER)
+ .addData(user1)
+ .addData(user2);
// String json = Jackson.json().serialize(addLog);
// System.out.println(json);
// System.out.println(json.getBytes(StandardCharsets.UTF_8).length);
@@ -322,9 +322,10 @@ public class SyncLog implements Serializable {
* @param
*/
public static class AddLog extends SyncLog {
- AddLog(SyncTypes.SyncType syncType) {
+ AddLog(int logSeq, SyncTypes.SyncType syncType) {
super(syncType);
this.operateType = OPT_ADD;
+ this.logSeq = logSeq;
}
public AddLog setSerializeDataCollect(List serializeDataCollect) {
@@ -332,7 +333,7 @@ public class SyncLog implements Serializable {
return this;
}
- public AddLog addSerializeData(T data) {
+ public AddLog addData(T data) {
if (this.serializeDataCollect == null) {
this.serializeDataCollect = new ArrayList<>();
}
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
deleted file mode 100644
index 60c6a04..0000000
--- a/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncLogInboundHandler.java
+++ /dev/null
@@ -1,19 +0,0 @@
-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
index c398230..2404790 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
@@ -2,15 +2,19 @@ package net.sopod.soim.router.datasync.server;
import io.netty.bootstrap.ServerBootstrap;
import io.netty.channel.Channel;
+import io.netty.channel.ChannelFutureListener;
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 net.sopod.soim.common.util.netty.FastThreadLocalThreadFactory;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
+import java.util.function.Consumer;
+
/**
* SyncServer
*
@@ -21,15 +25,24 @@ public class SyncServer {
private static final Logger logger = LoggerFactory.getLogger(SyncServer.class);
- public SyncServer() {
- NioEventLoopGroup boss = new NioEventLoopGroup(1);
- NioEventLoopGroup worker = new NioEventLoopGroup(2);
+ private final int port;
+
+ private NioEventLoopGroup boss;
+ private NioEventLoopGroup worker;
+
+ public SyncServer(int port) {
+ this.port = port;
+ }
+
+ public void start(Consumer onFail) throws InterruptedException {
+ 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()
.group(boss, worker)
.channel(NioServerSocketChannel.class)
.childHandler(new ChannelInitializer<>() {
@Override
- protected void initChannel(Channel channel) throws Exception {
+ protected void initChannel(Channel channel) {
LogLevel logLevel = logger.isDebugEnabled() ? LogLevel.DEBUG
: logger.isInfoEnabled() ? LogLevel.INFO
: logger.isWarnEnabled() ? LogLevel.WARN
@@ -37,14 +50,27 @@ public class SyncServer {
ChannelPipeline pipeline = channel.pipeline();
pipeline.addLast(new LoggingHandler(logLevel))
.addLast(new SyncLogDataCodec())
- .addLast(new SyncLogInboundHandler());
+ .addLast(new SyncServerInboundHandler());
}
});
- serverBoot.bind(8080);
+ serverBoot.bind(port).addListener((ChannelFutureListener) future -> {
+ if (!future.isSuccess()) {
+ if (onFail != null) {
+ onFail.accept(future.cause());
+ }
+ return;
+ }
+ logger.info("sync-server listening at {}...", port);
+ });
}
- public void start() {
-
+ public void shutdown() {
+ if (boss != null) {
+ boss.shutdownGracefully();
+ }
+ if (worker != null) {
+ worker.shutdownGracefully();
+ }
}
}
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
new file mode 100644
index 0000000..bfc1f1c
--- /dev/null
+++ b/im-service/im-router/src/main/java/net/sopod/soim/router/datasync/server/SyncServerInboundHandler.java
@@ -0,0 +1,33 @@
+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);
+ }
+
+}