mirror of
https://github.com/tiennm99/zfoo.git
synced 2026-08-05 14:24:41 +00:00
ref[net]: refactored the session attribute
This commit is contained in:
@@ -18,7 +18,7 @@ import com.zfoo.net.config.model.NetConfig;
|
||||
import com.zfoo.net.consumer.Consumer;
|
||||
import com.zfoo.net.packet.service.PacketService;
|
||||
import com.zfoo.net.router.Router;
|
||||
import com.zfoo.net.session.manager.SessionManager;
|
||||
import com.zfoo.net.session.SessionManager;
|
||||
import org.springframework.boot.autoconfigure.condition.ConditionalOnBean;
|
||||
import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean;
|
||||
import org.springframework.context.annotation.Bean;
|
||||
|
||||
@@ -19,8 +19,8 @@ import com.zfoo.net.core.AbstractClient;
|
||||
import com.zfoo.net.core.AbstractServer;
|
||||
import com.zfoo.net.packet.service.IPacketService;
|
||||
import com.zfoo.net.router.IRouter;
|
||||
import com.zfoo.net.session.manager.ISessionManager;
|
||||
import com.zfoo.net.session.model.Session;
|
||||
import com.zfoo.net.session.ISessionManager;
|
||||
import com.zfoo.net.session.Session;
|
||||
import com.zfoo.net.task.TaskBus;
|
||||
import com.zfoo.protocol.collection.ArrayUtils;
|
||||
import com.zfoo.protocol.exception.ExceptionUtils;
|
||||
|
||||
+14
-23
@@ -15,8 +15,7 @@ package com.zfoo.net.consumer.balancer;
|
||||
|
||||
import com.zfoo.net.NetContext;
|
||||
import com.zfoo.net.consumer.registry.RegisterVO;
|
||||
import com.zfoo.net.session.model.AttributeType;
|
||||
import com.zfoo.net.session.model.Session;
|
||||
import com.zfoo.net.session.Session;
|
||||
import com.zfoo.protocol.IPacket;
|
||||
import com.zfoo.protocol.ProtocolManager;
|
||||
import com.zfoo.protocol.registration.ProtocolModule;
|
||||
@@ -43,7 +42,7 @@ public abstract class AbstractConsumerLoadBalancer implements IConsumerLoadBalan
|
||||
balancer = ConsistentHashConsumerLoadBalancer.getInstance();
|
||||
break;
|
||||
default:
|
||||
throw new RuntimeException(StringUtils.format("无法识别负载均衡器[{}]", loadBalancer));
|
||||
throw new RuntimeException(StringUtils.format("Load balancer is not recognized[{}]", loadBalancer));
|
||||
}
|
||||
return balancer;
|
||||
}
|
||||
@@ -53,31 +52,24 @@ public abstract class AbstractConsumerLoadBalancer implements IConsumerLoadBalan
|
||||
}
|
||||
|
||||
public List<Session> getSessionsByModule(ProtocolModule module) {
|
||||
var clientSessionMap = NetContext.getSessionManager().getClientSessionMap();
|
||||
var sessions = clientSessionMap.values().stream()
|
||||
.filter(it -> {
|
||||
var attribute = it.getAttribute(AttributeType.CONSUMER);
|
||||
if (Objects.nonNull(attribute)) {
|
||||
var registerVO = (RegisterVO) attribute;
|
||||
return Objects.nonNull(registerVO.getProviderConfig()) && registerVO.getProviderConfig().getProviders().stream().anyMatch(provider -> provider.getProtocolModule().equals(module));
|
||||
} else {
|
||||
return false;
|
||||
}
|
||||
})
|
||||
return NetContext.getSessionManager().getClientSessionMap()
|
||||
.values()
|
||||
.stream()
|
||||
.filter(it -> it.getConsumerAttribute() != null && it.getConsumerAttribute().getProviderConfig() != null)
|
||||
.filter(it -> it.getConsumerAttribute().getProviderConfig().getProviders().stream().anyMatch(provider -> provider.getProtocolModule().equals(module)))
|
||||
.collect(Collectors.toList());
|
||||
return sessions;
|
||||
}
|
||||
|
||||
public List<Session> sessionsByModule(ProtocolModule module) {
|
||||
var clientSessionMap = NetContext.getSessionManager().getClientSessionMap();
|
||||
var sessions = new ArrayList<Session>();
|
||||
for(var clientSession : clientSessionMap.values()) {
|
||||
var attribute = clientSession.getAttribute(AttributeType.CONSUMER);
|
||||
if (attribute == null) {
|
||||
for (var clientSession : clientSessionMap.values()) {
|
||||
var consumerAttribute = clientSession.getConsumerAttribute();
|
||||
if (consumerAttribute == null) {
|
||||
continue;
|
||||
}
|
||||
|
||||
var registerVO = (RegisterVO) attribute;
|
||||
var registerVO = (RegisterVO) consumerAttribute;
|
||||
var providerConfig = registerVO.getProviderConfig();
|
||||
if (providerConfig == null) {
|
||||
continue;
|
||||
@@ -92,13 +84,12 @@ public abstract class AbstractConsumerLoadBalancer implements IConsumerLoadBalan
|
||||
|
||||
|
||||
public boolean sessionHasModule(Session session, IPacket packet) {
|
||||
|
||||
var attribute = session.getAttribute(AttributeType.CONSUMER);
|
||||
if (Objects.isNull(attribute)) {
|
||||
var consumerAttribute = session.getConsumerAttribute();
|
||||
if (Objects.isNull(consumerAttribute)) {
|
||||
return false;
|
||||
}
|
||||
|
||||
var registerVO = (RegisterVO) attribute;
|
||||
var registerVO = (RegisterVO) consumerAttribute;
|
||||
if (Objects.isNull(registerVO.getProviderConfig())) {
|
||||
return false;
|
||||
}
|
||||
|
||||
+3
-4
@@ -14,8 +14,7 @@
|
||||
package com.zfoo.net.consumer.balancer;
|
||||
|
||||
import com.zfoo.net.NetContext;
|
||||
import com.zfoo.net.session.model.AttributeType;
|
||||
import com.zfoo.net.session.model.Session;
|
||||
import com.zfoo.net.session.Session;
|
||||
import com.zfoo.protocol.IPacket;
|
||||
import com.zfoo.protocol.ProtocolManager;
|
||||
import com.zfoo.protocol.collection.CollectionUtils;
|
||||
@@ -84,7 +83,7 @@ public class ConsistentHashConsumerLoadBalancer extends AbstractConsumerLoadBala
|
||||
consistentHash = updateModuleToConsistentHash(module);
|
||||
}
|
||||
if (consistentHash == null) {
|
||||
throw new RunException("一致性hash负载均衡[protocolId:{}]参数[argument:{}],没有服务提供者提供服务[module:{}]", packet.protocolId(), argument, module);
|
||||
throw new RunException("ConsistentHashLoadBalancer [protocolId:{}][argument:{}], no service provides the [module:{}]", packet.protocolId(), argument, module);
|
||||
}
|
||||
var sid = consistentHash.getRealNode(argument).getValue();
|
||||
return NetContext.getSessionManager().getClientSession(sid);
|
||||
@@ -96,7 +95,7 @@ public class ConsistentHashConsumerLoadBalancer extends AbstractConsumerLoadBala
|
||||
private ConsistentHash<String, Long> updateModuleToConsistentHash(ProtocolModule module) {
|
||||
var sessionStringList = getSessionsByModule(module)
|
||||
.stream()
|
||||
.map(session -> new Pair<>(session.getAttribute(AttributeType.CONSUMER).toString(), session.getSid()))
|
||||
.map(session -> new Pair<>(session.getConsumerAttribute().toString(), session.getSid()))
|
||||
.sorted((a, b) -> a.getKey().compareTo(b.getKey()))
|
||||
.collect(Collectors.toList());
|
||||
|
||||
|
||||
@@ -14,7 +14,7 @@
|
||||
package com.zfoo.net.consumer.balancer;
|
||||
|
||||
import com.zfoo.net.router.attachment.SignalAttachment;
|
||||
import com.zfoo.net.session.model.Session;
|
||||
import com.zfoo.net.session.Session;
|
||||
import com.zfoo.protocol.IPacket;
|
||||
import org.springframework.lang.Nullable;
|
||||
|
||||
|
||||
@@ -13,7 +13,7 @@
|
||||
|
||||
package com.zfoo.net.consumer.balancer;
|
||||
|
||||
import com.zfoo.net.session.model.Session;
|
||||
import com.zfoo.net.session.Session;
|
||||
import com.zfoo.protocol.IPacket;
|
||||
import com.zfoo.protocol.ProtocolManager;
|
||||
import com.zfoo.protocol.exception.RunException;
|
||||
@@ -42,7 +42,7 @@ public class RandomConsumerLoadBalancer extends AbstractConsumerLoadBalancer {
|
||||
var sessions = getSessionsByModule(module);
|
||||
|
||||
if (sessions.isEmpty()) {
|
||||
throw new RunException("一致性hash负载均衡[protocolId:{}]参数[argument:{}],没有服务提供者提供服务[module:{}]", packet.protocolId(), argument, module);
|
||||
throw new RunException("RandomConsumerLoadBalancer [protocolId:{}][argument:{}], no service provides the [module:{}]", packet.protocolId(), argument, module);
|
||||
}
|
||||
|
||||
return RandomUtils.randomEle(sessions);
|
||||
|
||||
@@ -15,7 +15,7 @@ package com.zfoo.net.consumer.event;
|
||||
|
||||
import com.zfoo.event.model.event.IEvent;
|
||||
import com.zfoo.net.consumer.registry.RegisterVO;
|
||||
import com.zfoo.net.session.model.Session;
|
||||
import com.zfoo.net.session.Session;
|
||||
|
||||
/**
|
||||
* @author godotg
|
||||
|
||||
@@ -18,7 +18,6 @@ import com.zfoo.net.NetContext;
|
||||
import com.zfoo.net.consumer.event.ConsumerStartEvent;
|
||||
import com.zfoo.net.core.tcp.TcpClient;
|
||||
import com.zfoo.net.core.tcp.TcpServer;
|
||||
import com.zfoo.net.session.model.AttributeType;
|
||||
import com.zfoo.net.util.SessionUtils;
|
||||
import com.zfoo.protocol.collection.ArrayUtils;
|
||||
import com.zfoo.protocol.collection.concurrent.ConcurrentArrayList;
|
||||
@@ -190,13 +189,13 @@ public class ZookeeperRegistry implements IRegistry {
|
||||
switch (state) {
|
||||
// zk客户端与zk服务器失去了连接(忽略此情况,使用本地配置的缓存)
|
||||
case LOST:
|
||||
logger.error("[zookeeper:{}]失去连接,使用缓存", zookeeperConnectStr);
|
||||
logger.error("[zookeeper:{}] lost connection, cache used", zookeeperConnectStr);
|
||||
break;
|
||||
|
||||
// 暂停和只读这2种状态不检测
|
||||
case SUSPENDED:
|
||||
case READ_ONLY:
|
||||
logger.warn("[zookeeper:{}]忽略的[state{}]", zookeeperConnectStr, state);
|
||||
logger.warn("[zookeeper:{}] ignored [state{}]", zookeeperConnectStr, state);
|
||||
break;
|
||||
|
||||
// zk客户端和zk服务器重连了
|
||||
@@ -210,7 +209,7 @@ public class ZookeeperRegistry implements IRegistry {
|
||||
break;
|
||||
|
||||
default:
|
||||
logger.error("[zookeeper:{}]未知状态[state{}]", zookeeperConnectStr, state);
|
||||
logger.error("[zookeeper:{}] unknown [state{}]", zookeeperConnectStr, state);
|
||||
}
|
||||
}
|
||||
}, executor);
|
||||
@@ -219,7 +218,7 @@ public class ZookeeperRegistry implements IRegistry {
|
||||
try {
|
||||
curator.blockUntilConnected();
|
||||
} catch (Throwable t) {
|
||||
throw new RuntimeException("启动zookeeper异常", t);
|
||||
throw new RuntimeException("Start zookeeper exception", t);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -255,7 +254,7 @@ public class ZookeeperRegistry implements IRegistry {
|
||||
|
||||
// 检查zookeeper根节点的内容
|
||||
if (!rootPathData.equals(registryConfig.getCenter())) {
|
||||
throw new RuntimeException(StringUtils.format("zookeeper的rootPath[{}]内容配置错误[{}],期望的内容是[{}],请检查相关节点并重新启动", ROOT_PATH, rootPathData, registryConfig.getCenter()));
|
||||
throw new RuntimeException(StringUtils.format("zookeeper rootPath[{}] misconfigured [{}],expected [{}], check the relevant nodes and restart", ROOT_PATH, rootPathData, registryConfig.getCenter()));
|
||||
}
|
||||
|
||||
// 检查zookeeper根节点的权限
|
||||
@@ -268,7 +267,7 @@ public class ZookeeperRegistry implements IRegistry {
|
||||
var aclList = List.of(new ACL(ZooDefs.Perms.ALL, new Id("digest", DigestAuthenticationProvider.generateDigest(zookeeperAuthorStr))));
|
||||
AssertionUtils.isTrue(providerRootPathAclList.get(0).equals(aclList.get(0)));
|
||||
} catch (Exception e) {
|
||||
throw new RuntimeException(StringUtils.format("zookeeper的rootPath[{}]权限配置错误[{}]", ROOT_PATH, ExceptionUtils.getMessage(e)));
|
||||
throw new RuntimeException(StringUtils.format("zookeeper rootPath[{}] permissions are misconfigured [{}]", ROOT_PATH, ExceptionUtils.getMessage(e)));
|
||||
}
|
||||
}
|
||||
|
||||
@@ -300,7 +299,7 @@ public class ZookeeperRegistry implements IRegistry {
|
||||
// 初始化providerCache
|
||||
providerCuratorCache = CuratorCache.builder(curator, PROVIDER_ROOT_PATH)
|
||||
.withExceptionHandler(e -> {
|
||||
logger.error("providerCuratorCache未知异常", e);
|
||||
logger.error("providerCuratorCache unknown exception", e);
|
||||
initZookeeper();
|
||||
})
|
||||
.build();
|
||||
@@ -310,7 +309,7 @@ public class ZookeeperRegistry implements IRegistry {
|
||||
public void event(Type type, ChildData oldData, ChildData newData) {
|
||||
switch (type) {
|
||||
case NODE_CHANGED:
|
||||
logger.error("不需要处理的[oldData:{}][newData:{}]", childDataToString(oldData), childDataToString(newData));
|
||||
logger.error("No need to deal with [oldData:{}] [newData:{}]", childDataToString(oldData), childDataToString(newData));
|
||||
initZookeeper();
|
||||
break;
|
||||
case NODE_CREATED: // 意味着有可能来了自己作为消费者需要关心的服务提供者
|
||||
@@ -322,7 +321,7 @@ public class ZookeeperRegistry implements IRegistry {
|
||||
if (RegisterVO.providerHasConsumer(provider, localRegisterVO)) {
|
||||
providerHashConsumerSet.add(provider);
|
||||
checkConsumer();
|
||||
logger.info("发现新的订阅服务[{}]", providerStr);
|
||||
logger.info("Discover new subscription service of provider [{}]", providerStr);
|
||||
}
|
||||
break;
|
||||
case NODE_DELETED:
|
||||
@@ -331,7 +330,7 @@ public class ZookeeperRegistry implements IRegistry {
|
||||
if (providerHashConsumerSet.contains(oldProvider)) {
|
||||
providerHashConsumerSet.remove(oldProvider);
|
||||
checkConsumer();
|
||||
logger.info("取消订阅服务[{}]", oldProviderStr);
|
||||
logger.info("Unsubscribe from the service of provider [{}]", oldProviderStr);
|
||||
}
|
||||
break;
|
||||
default:
|
||||
@@ -358,7 +357,7 @@ public class ZookeeperRegistry implements IRegistry {
|
||||
initConsumerCache();
|
||||
} catch (Exception e) {
|
||||
//
|
||||
logger.error("zookeeper初始化失败,等待[{}]秒,重新初始化", RETRY_SECONDS, e);
|
||||
logger.error("Zookeeper failed to initialize, wait [{}] seconds to reinitialize", RETRY_SECONDS, e);
|
||||
SchedulerBus.schedule(() -> initZookeeper(), RETRY_SECONDS, TimeUnit.SECONDS);
|
||||
}
|
||||
});
|
||||
@@ -381,7 +380,7 @@ public class ZookeeperRegistry implements IRegistry {
|
||||
curator.create()
|
||||
.withMode(CreateMode.EPHEMERAL)
|
||||
.forPath(localProviderPath, StringUtils.EMPTY.getBytes());
|
||||
logger.info("注册服务成功[{}]", localProviderVoStr);
|
||||
logger.info("Registration for the provider successful [{}]", localProviderVoStr);
|
||||
} else {
|
||||
// 如果服务提供者已经有节点了,防止这个节点是是上次来不及删除的临时节点
|
||||
var curatorSessionId = curator.getZookeeperClient().getZooKeeper().getSessionId();
|
||||
@@ -392,7 +391,7 @@ public class ZookeeperRegistry implements IRegistry {
|
||||
.deletingChildrenIfNeeded()
|
||||
.withVersion(localProviderStat.getVersion())
|
||||
.forPath(localProviderPath);
|
||||
throw new RuntimeException(StringUtils.format("curator[sessionId:{}]和providerNode[sessionId:{}]的session不一致"
|
||||
throw new RuntimeException(StringUtils.format("session of curator[sessionId:{}] and providerNode[sessionId:{}] can not match"
|
||||
, curatorSessionId, providerNodeSessionId));
|
||||
}
|
||||
}
|
||||
@@ -449,22 +448,20 @@ public class ZookeeperRegistry implements IRegistry {
|
||||
*/
|
||||
private void doCheckConsumer() {
|
||||
if (curator.getState() != CuratorFrameworkState.STARTED) {
|
||||
logger.error("curator还没有启动,忽略本次consumer的检查");
|
||||
logger.error("Curator has not been started yet, ignoring this consumer check");
|
||||
return;
|
||||
}
|
||||
|
||||
logger.info("开始通过providerHashConsumerSet:{}检查[consumer:{}]", providerHashConsumerSet, NetContext.getSessionManager().getClientSessionMap().size());
|
||||
logger.info("start using providerHashConsumerSet:{} to check [consumer:{}]", providerHashConsumerSet, NetContext.getSessionManager().getClientSessionMap().size());
|
||||
|
||||
var recheckFlag = false;
|
||||
|
||||
for (var providerCache : providerHashConsumerSet) {
|
||||
// 先排除已经启动的consumer
|
||||
// getClientSessionMap
|
||||
var consumerClientList = NetContext.getSessionManager().getClientSessionMap().values().stream()
|
||||
.filter(it -> {
|
||||
var attribute = it.getAttribute(AttributeType.CONSUMER);
|
||||
return Objects.nonNull(attribute) && attribute.equals(providerCache);
|
||||
})
|
||||
var consumerClientList = NetContext.getSessionManager().getClientSessionMap().values()
|
||||
.stream()
|
||||
.filter(it -> it.getConsumerAttribute() != null && it.getConsumerAttribute().equals(providerCache))
|
||||
.collect(Collectors.toList());
|
||||
|
||||
if (consumerClientList.size() == 1) {
|
||||
@@ -474,11 +471,11 @@ public class ZookeeperRegistry implements IRegistry {
|
||||
} else {
|
||||
recheckFlag = true;
|
||||
NetContext.getSessionManager().removeClientSession(consumer);
|
||||
logger.error("[consumer:{}]失去连接,从clientSession中移除", consumer);
|
||||
logger.error("[consumer:{}] lost connection, removed from ClientSession", consumer);
|
||||
continue;
|
||||
}
|
||||
} else if (consumerClientList.size() > 1) {
|
||||
logger.error("[consumerClientList:{}]中有多个重复的[RegisterVO:{}]", consumerClientList, providerCache);
|
||||
logger.error("[consumerClientList:{}] are multiple duplicate [RegisterVO:{}]", consumerClientList, providerCache);
|
||||
continue;
|
||||
}
|
||||
|
||||
@@ -488,11 +485,11 @@ public class ZookeeperRegistry implements IRegistry {
|
||||
|
||||
// 自己作为消费者,使用TcpClient连接服务提供者不成功
|
||||
if (Objects.isNull(session)) {
|
||||
logger.error("[consumer:{}]启动失败,等待[{}]秒,重新检查consumer", providerCache, RETRY_SECONDS);
|
||||
logger.error("[consumer:{}] failed to start, wait [{}] seconds to recheck consumer", providerCache, RETRY_SECONDS);
|
||||
recheckFlag = true;
|
||||
} else {
|
||||
// 连接上了服务提供者
|
||||
session.putAttribute(AttributeType.CONSUMER, providerCache);
|
||||
session.setConsumerAttribute(providerCache);
|
||||
EventBus.submit(ConsumerStartEvent.valueOf(providerCache, session));
|
||||
|
||||
try {
|
||||
@@ -509,7 +506,7 @@ public class ZookeeperRegistry implements IRegistry {
|
||||
|
||||
} catch (Exception e) {
|
||||
// 因为并不关心consumer的状态,这种失败只需要记录一个错误日志就可以了
|
||||
logger.error("consumer写入zookeeper失败", e);
|
||||
logger.error("consumer writing to Zookeeper failed", e);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -603,9 +600,9 @@ public class ZookeeperRegistry implements IRegistry {
|
||||
.collect(Collectors.toList());
|
||||
return children;
|
||||
} catch (Exception e) {
|
||||
logger.error("未知异常", e);
|
||||
logger.error("unknown exception", e);
|
||||
} catch (Throwable t) {
|
||||
logger.error("未知错误", t);
|
||||
logger.error("unknown error", t);
|
||||
}
|
||||
return Collections.emptyList();
|
||||
}
|
||||
@@ -625,9 +622,9 @@ public class ZookeeperRegistry implements IRegistry {
|
||||
.collect(Collectors.toSet());
|
||||
return remoteProviderSet;
|
||||
} catch (Exception e) {
|
||||
logger.error("未知异常", e);
|
||||
logger.error("unknown exception", e);
|
||||
} catch (Throwable t) {
|
||||
logger.error("未知错误", t);
|
||||
logger.error("unknown error", t);
|
||||
}
|
||||
return Collections.emptySet();
|
||||
}
|
||||
|
||||
@@ -15,7 +15,7 @@ package com.zfoo.net.core;
|
||||
|
||||
import com.zfoo.net.NetContext;
|
||||
import com.zfoo.net.handler.BaseRouteHandler;
|
||||
import com.zfoo.net.session.model.Session;
|
||||
import com.zfoo.net.session.Session;
|
||||
import com.zfoo.protocol.exception.ExceptionUtils;
|
||||
import com.zfoo.protocol.util.IOUtils;
|
||||
import com.zfoo.util.ThreadUtils;
|
||||
@@ -81,7 +81,7 @@ public abstract class AbstractClient implements IClient {
|
||||
} else if (channelFuture.cause() != null) {
|
||||
logger.error(ExceptionUtils.getMessage(channelFuture.cause()));
|
||||
} else {
|
||||
logger.error("启动客户端[client:{}]未知错误", this);
|
||||
logger.error("[{}] started failed", this.getClass().getSimpleName());
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
@@ -13,7 +13,7 @@
|
||||
|
||||
package com.zfoo.net.core;
|
||||
|
||||
import com.zfoo.net.session.model.Session;
|
||||
import com.zfoo.net.session.Session;
|
||||
|
||||
/**
|
||||
* @author jaysunxiao
|
||||
|
||||
@@ -17,7 +17,7 @@ import com.zfoo.net.core.AbstractServer;
|
||||
import com.zfoo.net.handler.GatewayRouteHandler;
|
||||
import com.zfoo.net.handler.codec.tcp.TcpCodecHandler;
|
||||
import com.zfoo.net.handler.idle.ServerIdleHandler;
|
||||
import com.zfoo.net.session.model.Session;
|
||||
import com.zfoo.net.session.Session;
|
||||
import com.zfoo.protocol.IPacket;
|
||||
import com.zfoo.util.net.HostAndPort;
|
||||
import io.netty.channel.ChannelInitializer;
|
||||
|
||||
@@ -17,7 +17,7 @@ import com.zfoo.net.core.AbstractServer;
|
||||
import com.zfoo.net.handler.GatewayRouteHandler;
|
||||
import com.zfoo.net.handler.codec.websocket.WebSocketCodecHandler;
|
||||
import com.zfoo.net.handler.idle.ServerIdleHandler;
|
||||
import com.zfoo.net.session.model.Session;
|
||||
import com.zfoo.net.session.Session;
|
||||
import com.zfoo.protocol.IPacket;
|
||||
import com.zfoo.protocol.util.IOUtils;
|
||||
import com.zfoo.util.net.HostAndPort;
|
||||
|
||||
@@ -17,7 +17,7 @@ import com.zfoo.net.core.AbstractServer;
|
||||
import com.zfoo.net.handler.GatewayRouteHandler;
|
||||
import com.zfoo.net.handler.codec.websocket.WebSocketCodecHandler;
|
||||
import com.zfoo.net.handler.idle.ServerIdleHandler;
|
||||
import com.zfoo.net.session.model.Session;
|
||||
import com.zfoo.net.session.Session;
|
||||
import com.zfoo.protocol.IPacket;
|
||||
import com.zfoo.protocol.exception.ExceptionUtils;
|
||||
import com.zfoo.protocol.util.IOUtils;
|
||||
|
||||
@@ -17,7 +17,7 @@ import com.zfoo.net.core.AbstractServer;
|
||||
import com.zfoo.net.handler.GatewayRouteHandler;
|
||||
import com.zfoo.net.handler.codec.jprotobuf.JProtobufTcpCodecHandler;
|
||||
import com.zfoo.net.handler.idle.ServerIdleHandler;
|
||||
import com.zfoo.net.session.model.Session;
|
||||
import com.zfoo.net.session.Session;
|
||||
import com.zfoo.protocol.IPacket;
|
||||
import com.zfoo.util.net.HostAndPort;
|
||||
import io.netty.channel.ChannelInitializer;
|
||||
|
||||
@@ -14,7 +14,7 @@
|
||||
package com.zfoo.net.core.tcp.model;
|
||||
|
||||
import com.zfoo.event.model.event.IEvent;
|
||||
import com.zfoo.net.session.model.Session;
|
||||
import com.zfoo.net.session.Session;
|
||||
|
||||
/**
|
||||
* @author jaysunxiao
|
||||
|
||||
@@ -15,7 +15,7 @@ package com.zfoo.net.core.tcp.model;
|
||||
|
||||
import com.zfoo.event.model.event.IEvent;
|
||||
import com.zfoo.net.router.attachment.IAttachment;
|
||||
import com.zfoo.net.session.model.Session;
|
||||
import com.zfoo.net.session.Session;
|
||||
import com.zfoo.protocol.IPacket;
|
||||
|
||||
/**
|
||||
|
||||
@@ -14,7 +14,7 @@
|
||||
package com.zfoo.net.core.tcp.model;
|
||||
|
||||
import com.zfoo.event.model.event.IEvent;
|
||||
import com.zfoo.net.session.model.Session;
|
||||
import com.zfoo.net.session.Session;
|
||||
|
||||
/**
|
||||
* @author jaysunxiao
|
||||
|
||||
@@ -18,7 +18,7 @@ import com.zfoo.net.core.AbstractClient;
|
||||
import com.zfoo.net.handler.BaseRouteHandler;
|
||||
import com.zfoo.net.handler.ClientRouteHandler;
|
||||
import com.zfoo.net.handler.codec.udp.UdpCodecHandler;
|
||||
import com.zfoo.net.session.model.Session;
|
||||
import com.zfoo.net.session.Session;
|
||||
import com.zfoo.protocol.exception.ExceptionUtils;
|
||||
import com.zfoo.util.net.HostAndPort;
|
||||
import io.netty.bootstrap.Bootstrap;
|
||||
@@ -63,7 +63,7 @@ public class UdpClient extends AbstractClient {
|
||||
} else if (channelFuture.cause() != null) {
|
||||
logger.error(ExceptionUtils.getMessage(channelFuture.cause()));
|
||||
} else {
|
||||
logger.error("启动客户端[client:{}]未知错误", this);
|
||||
logger.error("[{}] started failed", this.getClass().getSimpleName());
|
||||
}
|
||||
} catch (Exception e) {
|
||||
logger.error(ExceptionUtils.getMessage(e));
|
||||
|
||||
@@ -15,7 +15,7 @@ package com.zfoo.net.handler;
|
||||
|
||||
import com.zfoo.net.NetContext;
|
||||
import com.zfoo.net.packet.model.DecodedPacketInfo;
|
||||
import com.zfoo.net.session.model.Session;
|
||||
import com.zfoo.net.session.Session;
|
||||
import com.zfoo.net.util.SessionUtils;
|
||||
import com.zfoo.protocol.util.StringUtils;
|
||||
import io.netty.channel.Channel;
|
||||
@@ -43,7 +43,7 @@ public class BaseRouteHandler extends ChannelInboundHandlerAdapter {
|
||||
var setSuccessful = sessionAttr.compareAndSet(null, session);
|
||||
if (!setSuccessful) {
|
||||
channel.close();
|
||||
throw new RuntimeException(StringUtils.format("无法设置[channel:{}]的session", channel));
|
||||
throw new RuntimeException(StringUtils.format("The properties of the session[channel:{}] cannot be set", channel));
|
||||
}
|
||||
return session;
|
||||
}
|
||||
|
||||
@@ -16,7 +16,6 @@ package com.zfoo.net.handler;
|
||||
import com.zfoo.event.manager.EventBus;
|
||||
import com.zfoo.net.NetContext;
|
||||
import com.zfoo.net.core.tcp.model.ClientSessionInactiveEvent;
|
||||
import com.zfoo.net.session.model.AttributeType;
|
||||
import com.zfoo.net.util.SessionUtils;
|
||||
import io.netty.channel.ChannelHandler;
|
||||
import io.netty.channel.ChannelHandlerContext;
|
||||
@@ -48,12 +47,11 @@ public class ClientRouteHandler extends BaseRouteHandler {
|
||||
return;
|
||||
}
|
||||
|
||||
var consumeAttribute = session.getAttribute(AttributeType.CONSUMER);
|
||||
NetContext.getSessionManager().removeClientSession(session);
|
||||
EventBus.submit(ClientSessionInactiveEvent.valueOf(session));
|
||||
|
||||
// 如果是消费者inactive,还需要触发客户端消费者检查事件,以便重新连接
|
||||
if (consumeAttribute != null) {
|
||||
if (session.getConsumerAttribute() != null) {
|
||||
NetContext.getConfigManager().getRegistry().checkConsumer();
|
||||
}
|
||||
|
||||
|
||||
@@ -25,8 +25,7 @@ import com.zfoo.net.packet.model.DecodedPacketInfo;
|
||||
import com.zfoo.net.router.attachment.GatewayAttachment;
|
||||
import com.zfoo.net.router.attachment.IAttachment;
|
||||
import com.zfoo.net.router.attachment.SignalAttachment;
|
||||
import com.zfoo.net.session.model.AttributeType;
|
||||
import com.zfoo.net.session.model.Session;
|
||||
import com.zfoo.net.session.Session;
|
||||
import com.zfoo.net.util.SessionUtils;
|
||||
import com.zfoo.protocol.IPacket;
|
||||
import com.zfoo.protocol.util.JsonUtils;
|
||||
@@ -92,8 +91,8 @@ public class GatewayRouteHandler extends ServerRouteHandler {
|
||||
return;
|
||||
} else {
|
||||
// 使用用户的uid做一致性hash
|
||||
var uid = (Long) session.getAttribute(AttributeType.UID);
|
||||
if (uid != null) {
|
||||
var uid = session.getUid();
|
||||
if (uid < 0) {
|
||||
forwardingPacket(packet, gatewayAttachment, uid);
|
||||
return;
|
||||
}
|
||||
@@ -113,9 +112,9 @@ public class GatewayRouteHandler extends ServerRouteHandler {
|
||||
var consumerSession = ConsistentHashConsumerLoadBalancer.getInstance().loadBalancer(packet, argument);
|
||||
NetContext.getRouter().send(consumerSession, packet, attachment);
|
||||
} catch (Exception e) {
|
||||
logger.error("网关发生异常", e);
|
||||
logger.error("An exception occurred at the gateway", e);
|
||||
} catch (Throwable t) {
|
||||
logger.error("网关发生错误", t);
|
||||
logger.error("An error occurred at the gateway", t);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -127,10 +126,10 @@ public class GatewayRouteHandler extends ServerRouteHandler {
|
||||
}
|
||||
|
||||
var sid = session.getSid();
|
||||
var uid = (Long) session.getAttribute(AttributeType.UID);
|
||||
var uid = session.getUid();
|
||||
|
||||
// 连接到网关的客户端断开了连接
|
||||
EventBus.submit(GatewaySessionInactiveEvent.valueOf(sid, uid == null ? 0 : uid.longValue()));
|
||||
EventBus.submit(GatewaySessionInactiveEvent.valueOf(sid, uid));
|
||||
|
||||
super.channelInactive(ctx);
|
||||
}
|
||||
|
||||
@@ -16,7 +16,7 @@ package com.zfoo.net.router;
|
||||
import com.zfoo.net.router.answer.AsyncAnswer;
|
||||
import com.zfoo.net.router.answer.SyncAnswer;
|
||||
import com.zfoo.net.router.attachment.IAttachment;
|
||||
import com.zfoo.net.session.model.Session;
|
||||
import com.zfoo.net.session.Session;
|
||||
import com.zfoo.protocol.IPacket;
|
||||
import org.springframework.lang.Nullable;
|
||||
|
||||
|
||||
@@ -32,8 +32,7 @@ import com.zfoo.net.router.exception.NetTimeOutException;
|
||||
import com.zfoo.net.router.exception.UnexpectedProtocolException;
|
||||
import com.zfoo.net.router.route.PacketBus;
|
||||
import com.zfoo.net.router.route.SignalBridge;
|
||||
import com.zfoo.net.session.model.AttributeType;
|
||||
import com.zfoo.net.session.model.Session;
|
||||
import com.zfoo.net.session.Session;
|
||||
import com.zfoo.net.task.PacketReceiverTask;
|
||||
import com.zfoo.net.task.TaskBus;
|
||||
import com.zfoo.protocol.IPacket;
|
||||
@@ -126,7 +125,7 @@ public class Router implements IRouter {
|
||||
logger.error("错误的网关授权信息,uid必须大于0");
|
||||
return;
|
||||
}
|
||||
gatewaySession.putAttribute(AttributeType.UID, uid);
|
||||
session.setUid(uid);
|
||||
EventBus.submit(AuthUidToGatewayEvent.valueOf(gatewaySession.getSid(), uid));
|
||||
|
||||
NetContext.getRouter().send(session, AuthUidToGatewayConfirm.valueOf(uid), new GatewayAttachment(gatewaySession, null));
|
||||
@@ -324,9 +323,9 @@ public class Router implements IRouter {
|
||||
PacketBus.route(session, packet, attachment);
|
||||
} catch (Exception e) {
|
||||
EventBus.submit(ServerExceptionEvent.valueOf(session, packet, attachment, e));
|
||||
logger.error(StringUtils.format("e[uid:{}][sid:{}]未知exception异常", session.getAttribute(AttributeType.UID), session.getSid(), e.getMessage()), e);
|
||||
logger.error(StringUtils.format("e[uid:{}][sid:{}] unknown exception", session.getUid(), session.getSid(), e.getMessage()), e);
|
||||
} catch (Throwable t) {
|
||||
logger.error(StringUtils.format("e[uid:{}][sid:{}]未知error错误", session.getAttribute(AttributeType.UID), session.getSid(), t.getMessage()), t);
|
||||
logger.error(StringUtils.format("e[uid:{}][sid:{}] unknown error", session.getUid(), session.getSid(), t.getMessage()), t);
|
||||
} finally {
|
||||
// 如果有服务器在处理同步或者异步消息的时候由于错误没有返回给客户端消息,则可能会残留serverAttachment,所以先移除
|
||||
if (attachment != null) {
|
||||
|
||||
@@ -12,8 +12,7 @@
|
||||
|
||||
package com.zfoo.net.router.attachment;
|
||||
|
||||
import com.zfoo.net.session.model.AttributeType;
|
||||
import com.zfoo.net.session.model.Session;
|
||||
import com.zfoo.net.session.Session;
|
||||
import org.springframework.lang.Nullable;
|
||||
|
||||
/**
|
||||
@@ -59,8 +58,7 @@ public class GatewayAttachment implements IAttachment {
|
||||
public GatewayAttachment(Session session, @Nullable SignalAttachment signalAttachment) {
|
||||
this.client = true;
|
||||
this.sid = session.getSid();
|
||||
var uid = session.getAttribute(AttributeType.UID);
|
||||
this.uid = uid == null ? 0 : (long) uid;
|
||||
this.uid = session.getUid();
|
||||
this.signalAttachment = signalAttachment;
|
||||
}
|
||||
|
||||
|
||||
@@ -14,7 +14,7 @@
|
||||
package com.zfoo.net.router.receiver;
|
||||
|
||||
import com.zfoo.net.router.attachment.IAttachment;
|
||||
import com.zfoo.net.session.model.Session;
|
||||
import com.zfoo.net.session.Session;
|
||||
import com.zfoo.protocol.IPacket;
|
||||
import com.zfoo.protocol.util.StringUtils;
|
||||
import com.zfoo.util.security.IdUtils;
|
||||
|
||||
@@ -14,7 +14,7 @@
|
||||
package com.zfoo.net.router.receiver;
|
||||
|
||||
import com.zfoo.net.router.attachment.IAttachment;
|
||||
import com.zfoo.net.session.model.Session;
|
||||
import com.zfoo.net.session.Session;
|
||||
import com.zfoo.protocol.IPacket;
|
||||
|
||||
/**
|
||||
|
||||
@@ -14,7 +14,7 @@
|
||||
package com.zfoo.net.router.receiver;
|
||||
|
||||
import com.zfoo.net.router.attachment.IAttachment;
|
||||
import com.zfoo.net.session.model.Session;
|
||||
import com.zfoo.net.session.Session;
|
||||
import com.zfoo.protocol.IPacket;
|
||||
import com.zfoo.protocol.util.ReflectionUtils;
|
||||
|
||||
|
||||
@@ -23,7 +23,7 @@ import com.zfoo.net.router.receiver.EnhanceUtils;
|
||||
import com.zfoo.net.router.receiver.IPacketReceiver;
|
||||
import com.zfoo.net.router.receiver.PacketReceiver;
|
||||
import com.zfoo.net.router.receiver.PacketReceiverDefinition;
|
||||
import com.zfoo.net.session.model.Session;
|
||||
import com.zfoo.net.session.Session;
|
||||
import com.zfoo.protocol.IPacket;
|
||||
import com.zfoo.protocol.ProtocolManager;
|
||||
import com.zfoo.protocol.collection.ArrayUtils;
|
||||
|
||||
@@ -19,7 +19,7 @@ import com.zfoo.net.config.model.*;
|
||||
import com.zfoo.net.consumer.Consumer;
|
||||
import com.zfoo.net.packet.service.PacketService;
|
||||
import com.zfoo.net.router.Router;
|
||||
import com.zfoo.net.session.manager.SessionManager;
|
||||
import com.zfoo.net.session.SessionManager;
|
||||
import com.zfoo.protocol.util.DomUtils;
|
||||
import com.zfoo.protocol.util.StringUtils;
|
||||
import org.springframework.beans.factory.config.BeanDefinitionHolder;
|
||||
|
||||
+1
-4
@@ -1,6 +1,5 @@
|
||||
/*
|
||||
* Copyright (C) 2020 The zfoo Authors
|
||||
*
|
||||
* Licensed under the Apache License, Version 2.0 (the "License"); you may not use this file except
|
||||
* in compliance with the License. You may obtain a copy of the License at
|
||||
*
|
||||
@@ -11,9 +10,7 @@
|
||||
* See the License for the specific language governing permissions and limitations under the License.
|
||||
*/
|
||||
|
||||
package com.zfoo.net.session.manager;
|
||||
|
||||
import com.zfoo.net.session.model.Session;
|
||||
package com.zfoo.net.session;
|
||||
|
||||
import java.util.Map;
|
||||
|
||||
+27
-28
@@ -1,6 +1,5 @@
|
||||
/*
|
||||
* Copyright (C) 2020 The zfoo Authors
|
||||
*
|
||||
* Licensed under the Apache License, Version 2.0 (the "License"); you may not use this file except
|
||||
* in compliance with the License. You may obtain a copy of the License at
|
||||
*
|
||||
@@ -11,16 +10,13 @@
|
||||
* See the License for the specific language governing permissions and limitations under the License.
|
||||
*/
|
||||
|
||||
package com.zfoo.net.session.model;
|
||||
package com.zfoo.net.session;
|
||||
|
||||
import com.zfoo.net.consumer.registry.RegisterVO;
|
||||
import com.zfoo.protocol.util.StringUtils;
|
||||
import io.netty.channel.Channel;
|
||||
|
||||
import java.io.Closeable;
|
||||
import java.util.EnumMap;
|
||||
import java.util.Map;
|
||||
import java.util.Objects;
|
||||
import java.util.concurrent.atomic.AtomicLong;
|
||||
|
||||
/**
|
||||
@@ -32,25 +28,26 @@ public class Session implements Closeable {
|
||||
private static final AtomicLong ATOMIC_LONG = new AtomicLong(0);
|
||||
|
||||
/**
|
||||
* session的id
|
||||
* The globally unique ID of the session
|
||||
*/
|
||||
private long sid;
|
||||
|
||||
/**
|
||||
* EN:The default user ID is an ID greater than 0, or less than 0 if there is no login
|
||||
* CN:默认用户的id都是大于0的id,如果没有登录则小于0
|
||||
*/
|
||||
private long uid = -1;
|
||||
|
||||
private Channel channel;
|
||||
|
||||
/**
|
||||
* Session附带的属性参数
|
||||
* Session附带的属性参数,消费者的属性
|
||||
*/
|
||||
private RegisterVO consumerAttribute = null;
|
||||
|
||||
private long uid = Long.MIN_VALUE;
|
||||
|
||||
private Map<AttributeType, Object> attributes = new EnumMap<>(AttributeType.class);
|
||||
|
||||
|
||||
public Session(Channel channel) {
|
||||
if (channel == null) {
|
||||
throw new IllegalArgumentException("channel不能为空");
|
||||
throw new IllegalArgumentException("channel cannot be empty");
|
||||
}
|
||||
this.sid = ATOMIC_LONG.getAndIncrement();
|
||||
this.channel = channel;
|
||||
@@ -59,7 +56,7 @@ public class Session implements Closeable {
|
||||
|
||||
@Override
|
||||
public String toString() {
|
||||
return StringUtils.format("[sid:{}] [channel:{}] [attributes:{}]", sid, channel, attributes);
|
||||
return StringUtils.format("[sid:{}] [uid:{}] [channel:{}]", sid, uid, channel);
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -76,7 +73,7 @@ public class Session implements Closeable {
|
||||
|
||||
@Override
|
||||
public int hashCode() {
|
||||
return Objects.hash(sid);
|
||||
return (int) sid;
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -92,21 +89,23 @@ public class Session implements Closeable {
|
||||
this.sid = sid;
|
||||
}
|
||||
|
||||
public synchronized void putAttribute(AttributeType key, Object value) {
|
||||
attributes.put(key, value);
|
||||
}
|
||||
|
||||
public synchronized void removeAttribute(AttributeType key) {
|
||||
attributes.remove(key);
|
||||
}
|
||||
|
||||
|
||||
public <T> T getAttribute(AttributeType key) {
|
||||
return (T) attributes.get(key);
|
||||
}
|
||||
|
||||
public Channel getChannel() {
|
||||
return channel;
|
||||
}
|
||||
|
||||
public long getUid() {
|
||||
return uid;
|
||||
}
|
||||
|
||||
public void setUid(long uid) {
|
||||
this.uid = uid;
|
||||
}
|
||||
|
||||
public RegisterVO getConsumerAttribute() {
|
||||
return consumerAttribute;
|
||||
}
|
||||
|
||||
public void setConsumerAttribute(RegisterVO consumerAttribute) {
|
||||
this.consumerAttribute = consumerAttribute;
|
||||
}
|
||||
}
|
||||
+5
-7
@@ -1,6 +1,5 @@
|
||||
/*
|
||||
* Copyright (C) 2020 The zfoo Authors
|
||||
*
|
||||
* Licensed under the Apache License, Version 2.0 (the "License"); you may not use this file except
|
||||
* in compliance with the License. You may obtain a copy of the License at
|
||||
*
|
||||
@@ -11,9 +10,8 @@
|
||||
* See the License for the specific language governing permissions and limitations under the License.
|
||||
*/
|
||||
|
||||
package com.zfoo.net.session.manager;
|
||||
package com.zfoo.net.session;
|
||||
|
||||
import com.zfoo.net.session.model.Session;
|
||||
import com.zfoo.net.util.SessionUtils;
|
||||
import com.zfoo.util.security.IdUtils;
|
||||
import org.slf4j.Logger;
|
||||
@@ -51,7 +49,7 @@ public class SessionManager implements ISessionManager {
|
||||
@Override
|
||||
public void addServerSession(Session session) {
|
||||
if (serverSessionMap.containsKey(session.getSid())) {
|
||||
logger.error("server收到重复的[session:{}]", SessionUtils.sessionInfo(session));
|
||||
logger.error("Server received duplicate [session:{}]", SessionUtils.sessionInfo(session));
|
||||
return;
|
||||
}
|
||||
serverSessionMap.put(session.getSid(), session);
|
||||
@@ -60,7 +58,7 @@ public class SessionManager implements ISessionManager {
|
||||
@Override
|
||||
public void removeServerSession(Session session) {
|
||||
if (!serverSessionMap.containsKey(session.getSid())) {
|
||||
logger.error("SessionManager中的serverSession没有包含[session:{}],所以无法移除", SessionUtils.sessionInfo(session));
|
||||
logger.error("[session:{}] does not exist", SessionUtils.sessionInfo(session));
|
||||
return;
|
||||
}
|
||||
serverSessionMap.remove(session.getSid());
|
||||
@@ -80,7 +78,7 @@ public class SessionManager implements ISessionManager {
|
||||
@Override
|
||||
public void addClientSession(Session session) {
|
||||
if (clientSessionMap.containsKey(session.getSid())) {
|
||||
logger.error("client收到重复的[session:{}]", SessionUtils.sessionInfo(session));
|
||||
logger.error("client received duplicate [session:{}]", SessionUtils.sessionInfo(session));
|
||||
return;
|
||||
}
|
||||
clientSessionMap.put(session.getSid(), session);
|
||||
@@ -90,7 +88,7 @@ public class SessionManager implements ISessionManager {
|
||||
@Override
|
||||
public void removeClientSession(Session session) {
|
||||
if (!clientSessionMap.containsKey(session.getSid())) {
|
||||
logger.error("SessionManager中的clientSession没有包含[session:{}],所以无法移除", SessionUtils.sessionInfo(session));
|
||||
logger.error("[session:{}] does not exist", SessionUtils.sessionInfo(session));
|
||||
return;
|
||||
}
|
||||
clientSessionMap.remove(session.getSid());
|
||||
@@ -1,38 +0,0 @@
|
||||
/*
|
||||
* Copyright (C) 2020 The zfoo Authors
|
||||
*
|
||||
* Licensed under the Apache License, Version 2.0 (the "License"); you may not use this file except
|
||||
* in compliance with the License. You may obtain a copy of the License at
|
||||
*
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
*
|
||||
* Unless required by applicable law or agreed to in writing, software distributed under the License is distributed
|
||||
* on an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
* See the License for the specific language governing permissions and limitations under the License.
|
||||
*/
|
||||
|
||||
package com.zfoo.net.session.model;
|
||||
|
||||
/**
|
||||
* @author godotg
|
||||
* @version 3.0
|
||||
*/
|
||||
public enum AttributeType {
|
||||
|
||||
|
||||
/**
|
||||
* 一般是客户端session
|
||||
*/
|
||||
CONSUMER,
|
||||
|
||||
/**
|
||||
* session的uid
|
||||
*/
|
||||
UID,
|
||||
|
||||
/**
|
||||
* 网关ip
|
||||
*/
|
||||
GATEWAY_HOST_AND_PORT,
|
||||
|
||||
}
|
||||
@@ -14,7 +14,7 @@ package com.zfoo.net.task;
|
||||
|
||||
import com.zfoo.net.NetContext;
|
||||
import com.zfoo.net.router.attachment.IAttachment;
|
||||
import com.zfoo.net.session.model.Session;
|
||||
import com.zfoo.net.session.Session;
|
||||
import com.zfoo.protocol.IPacket;
|
||||
|
||||
/**
|
||||
|
||||
@@ -15,7 +15,6 @@ package com.zfoo.net.task;
|
||||
|
||||
import com.zfoo.event.manager.EventBus;
|
||||
import com.zfoo.net.NetContext;
|
||||
import com.zfoo.net.session.model.AttributeType;
|
||||
import com.zfoo.protocol.collection.concurrent.CopyOnWriteHashMapLongObject;
|
||||
import com.zfoo.protocol.util.AssertionUtils;
|
||||
import com.zfoo.protocol.util.StringUtils;
|
||||
@@ -114,11 +113,11 @@ public final class TaskBus {
|
||||
|
||||
if (attachment == null) {
|
||||
var session = task.getSession();
|
||||
var uid = (Long) session.getAttribute(AttributeType.UID);
|
||||
if (uid == null) {
|
||||
var uid = session.getUid();
|
||||
if (uid < 0) {
|
||||
execute((int) session.getSid(), task);
|
||||
} else {
|
||||
execute(uid.intValue(), task);
|
||||
execute(uid, task);
|
||||
}
|
||||
} else {
|
||||
execute(attachment.taskExecutorHash(), task);
|
||||
|
||||
@@ -13,8 +13,7 @@
|
||||
|
||||
package com.zfoo.net.util;
|
||||
|
||||
import com.zfoo.net.session.model.AttributeType;
|
||||
import com.zfoo.net.session.model.Session;
|
||||
import com.zfoo.net.session.Session;
|
||||
import com.zfoo.protocol.util.StringUtils;
|
||||
import io.netty.channel.Channel;
|
||||
import io.netty.channel.ChannelHandlerContext;
|
||||
@@ -65,7 +64,7 @@ public abstract class SessionUtils {
|
||||
// to avoid: io.netty.channel.unix.Errors$NativeIoException: readAddress(..) failed: Connection reset by peer
|
||||
// 有些情况当建立连接过后迅速关闭,这个时候取remoteAddress会有异常
|
||||
}
|
||||
return StringUtils.format(CHANNEL_INFO_TEMPLATE, remoteAddress, session.getSid(), session.getAttribute(AttributeType.UID));
|
||||
return StringUtils.format(CHANNEL_INFO_TEMPLATE, remoteAddress, session.getSid(), session.getUid());
|
||||
}
|
||||
|
||||
public static String sessionSimpleInfo(ChannelHandlerContext ctx) {
|
||||
@@ -80,7 +79,7 @@ public abstract class SessionUtils {
|
||||
if (session == null) {
|
||||
return CHANNEL_SIMPLE_INFO_TEMPLATE;
|
||||
}
|
||||
return StringUtils.format(CHANNEL_SIMPLE_INFO_TEMPLATE, session.getSid(), session.getAttribute(AttributeType.UID));
|
||||
return StringUtils.format(CHANNEL_SIMPLE_INFO_TEMPLATE, session.getSid(), session.getUid());
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -16,7 +16,7 @@ package com.zfoo.net.core.csharp;
|
||||
import com.zfoo.net.NetContext;
|
||||
import com.zfoo.net.packet.csharp.CM_CSharpRequest;
|
||||
import com.zfoo.net.router.receiver.PacketReceiver;
|
||||
import com.zfoo.net.session.model.Session;
|
||||
import com.zfoo.net.session.Session;
|
||||
import com.zfoo.protocol.util.JsonUtils;
|
||||
import org.springframework.stereotype.Component;
|
||||
|
||||
|
||||
@@ -17,7 +17,7 @@ import com.zfoo.net.packet.gateway.GatewayToProviderRequest;
|
||||
import com.zfoo.net.packet.gateway.GatewayToProviderResponse;
|
||||
import com.zfoo.net.router.attachment.GatewayAttachment;
|
||||
import com.zfoo.net.router.receiver.PacketReceiver;
|
||||
import com.zfoo.net.session.model.Session;
|
||||
import com.zfoo.net.session.Session;
|
||||
import com.zfoo.protocol.util.JsonUtils;
|
||||
import com.zfoo.protocol.util.StringUtils;
|
||||
import org.slf4j.Logger;
|
||||
@@ -35,10 +35,6 @@ public class GatewayProviderController {
|
||||
|
||||
/**
|
||||
* 注意:这里第2个请求参数以Request结尾,那么第3个参数必须是 GatewayAttachment类型(参加:PacketBus中扫描时的校验)
|
||||
*
|
||||
* @param session
|
||||
* @param request
|
||||
* @param gatewayAttachment
|
||||
*/
|
||||
@PacketReceiver
|
||||
public void atGatewayToProviderRequest(Session session, GatewayToProviderRequest request, GatewayAttachment gatewayAttachment) {
|
||||
|
||||
@@ -17,9 +17,8 @@ import com.zfoo.net.packet.http.HttpHelloRequest;
|
||||
import com.zfoo.net.packet.http.HttpHelloResponse;
|
||||
import com.zfoo.net.router.attachment.HttpAttachment;
|
||||
import com.zfoo.net.router.receiver.PacketReceiver;
|
||||
import com.zfoo.net.session.model.Session;
|
||||
import com.zfoo.net.session.Session;
|
||||
import com.zfoo.protocol.util.JsonUtils;
|
||||
import io.netty.handler.codec.http.HttpResponseStatus;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
import org.springframework.stereotype.Component;
|
||||
|
||||
+1
-1
@@ -14,7 +14,7 @@ package com.zfoo.net.core.jprotobuf.client;
|
||||
|
||||
import com.zfoo.net.packet.jprotobuf.JProtobufHelloResponse;
|
||||
import com.zfoo.net.router.receiver.PacketReceiver;
|
||||
import com.zfoo.net.session.model.Session;
|
||||
import com.zfoo.net.session.Session;
|
||||
import com.zfoo.protocol.util.JsonUtils;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
@@ -16,7 +16,7 @@ import com.zfoo.net.NetContext;
|
||||
import com.zfoo.net.packet.jprotobuf.JProtobufHelloRequest;
|
||||
import com.zfoo.net.packet.jprotobuf.JProtobufHelloResponse;
|
||||
import com.zfoo.net.router.receiver.PacketReceiver;
|
||||
import com.zfoo.net.session.model.Session;
|
||||
import com.zfoo.net.session.Session;
|
||||
import com.zfoo.protocol.util.JsonUtils;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
@@ -14,9 +14,8 @@
|
||||
package com.zfoo.net.core.json.client;
|
||||
|
||||
import com.zfoo.net.packet.json.JsonHelloResponse;
|
||||
import com.zfoo.net.packet.websocket.WebsocketHelloResponse;
|
||||
import com.zfoo.net.router.receiver.PacketReceiver;
|
||||
import com.zfoo.net.session.model.Session;
|
||||
import com.zfoo.net.session.Session;
|
||||
import com.zfoo.protocol.util.JsonUtils;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
@@ -16,10 +16,8 @@ package com.zfoo.net.core.json.server;
|
||||
import com.zfoo.net.NetContext;
|
||||
import com.zfoo.net.packet.json.JsonHelloRequest;
|
||||
import com.zfoo.net.packet.json.JsonHelloResponse;
|
||||
import com.zfoo.net.packet.websocket.WebsocketHelloRequest;
|
||||
import com.zfoo.net.packet.websocket.WebsocketHelloResponse;
|
||||
import com.zfoo.net.router.receiver.PacketReceiver;
|
||||
import com.zfoo.net.session.model.Session;
|
||||
import com.zfoo.net.session.Session;
|
||||
import com.zfoo.protocol.util.JsonUtils;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
@@ -16,7 +16,7 @@ import com.zfoo.net.NetContext;
|
||||
import com.zfoo.net.packet.provider.ProviderMessAnswer;
|
||||
import com.zfoo.net.packet.provider.ProviderMessAsk;
|
||||
import com.zfoo.net.router.receiver.PacketReceiver;
|
||||
import com.zfoo.net.session.model.Session;
|
||||
import com.zfoo.net.session.Session;
|
||||
import com.zfoo.protocol.util.JsonUtils;
|
||||
import com.zfoo.protocol.util.StringUtils;
|
||||
import org.slf4j.Logger;
|
||||
|
||||
@@ -15,7 +15,7 @@ package com.zfoo.net.core.tcp.client;
|
||||
|
||||
import com.zfoo.net.packet.tcp.TcpHelloResponse;
|
||||
import com.zfoo.net.router.receiver.PacketReceiver;
|
||||
import com.zfoo.net.session.model.Session;
|
||||
import com.zfoo.net.session.Session;
|
||||
import com.zfoo.protocol.util.JsonUtils;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
@@ -16,7 +16,7 @@ import com.zfoo.net.NetContext;
|
||||
import com.zfoo.net.packet.tcp.TcpHelloRequest;
|
||||
import com.zfoo.net.packet.tcp.TcpHelloResponse;
|
||||
import com.zfoo.net.router.receiver.PacketReceiver;
|
||||
import com.zfoo.net.session.model.Session;
|
||||
import com.zfoo.net.session.Session;
|
||||
import com.zfoo.protocol.util.JsonUtils;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
@@ -16,7 +16,7 @@ import com.zfoo.net.NetContext;
|
||||
import com.zfoo.net.packet.tcp.AsyncMessAnswer;
|
||||
import com.zfoo.net.packet.tcp.AsyncMessAsk;
|
||||
import com.zfoo.net.router.receiver.PacketReceiver;
|
||||
import com.zfoo.net.session.model.Session;
|
||||
import com.zfoo.net.session.Session;
|
||||
import com.zfoo.protocol.util.JsonUtils;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
@@ -77,5 +77,5 @@ com.zfoo.net.config.manager.ConfigManager
|
||||
com.zfoo.net.packet.service.PacketService
|
||||
com.zfoo.net.router.Router
|
||||
com.zfoo.net.consumer.Consumer
|
||||
com.zfoo.net.session.manager.SessionManager
|
||||
com.zfoo.net.session.SessionManager
|
||||
*/
|
||||
|
||||
@@ -16,7 +16,7 @@ import com.zfoo.net.NetContext;
|
||||
import com.zfoo.net.packet.tcp.SyncMessAnswer;
|
||||
import com.zfoo.net.packet.tcp.SyncMessAsk;
|
||||
import com.zfoo.net.router.receiver.PacketReceiver;
|
||||
import com.zfoo.net.session.model.Session;
|
||||
import com.zfoo.net.session.Session;
|
||||
import com.zfoo.protocol.util.JsonUtils;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
@@ -15,7 +15,7 @@ package com.zfoo.net.core.udp.client;
|
||||
import com.zfoo.net.packet.udp.UdpHelloResponse;
|
||||
import com.zfoo.net.router.attachment.UdpAttachment;
|
||||
import com.zfoo.net.router.receiver.PacketReceiver;
|
||||
import com.zfoo.net.session.model.Session;
|
||||
import com.zfoo.net.session.Session;
|
||||
import com.zfoo.protocol.util.JsonUtils;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
@@ -17,7 +17,7 @@ import com.zfoo.net.packet.udp.UdpHelloRequest;
|
||||
import com.zfoo.net.packet.udp.UdpHelloResponse;
|
||||
import com.zfoo.net.router.attachment.UdpAttachment;
|
||||
import com.zfoo.net.router.receiver.PacketReceiver;
|
||||
import com.zfoo.net.session.model.Session;
|
||||
import com.zfoo.net.session.Session;
|
||||
import com.zfoo.protocol.util.JsonUtils;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
+1
-1
@@ -15,7 +15,7 @@ package com.zfoo.net.core.websocket.client;
|
||||
|
||||
import com.zfoo.net.packet.websocket.WebsocketHelloResponse;
|
||||
import com.zfoo.net.router.receiver.PacketReceiver;
|
||||
import com.zfoo.net.session.model.Session;
|
||||
import com.zfoo.net.session.Session;
|
||||
import com.zfoo.protocol.util.JsonUtils;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
+1
-1
@@ -16,7 +16,7 @@ import com.zfoo.net.NetContext;
|
||||
import com.zfoo.net.packet.websocket.WebsocketHelloRequest;
|
||||
import com.zfoo.net.packet.websocket.WebsocketHelloResponse;
|
||||
import com.zfoo.net.router.receiver.PacketReceiver;
|
||||
import com.zfoo.net.session.model.Session;
|
||||
import com.zfoo.net.session.Session;
|
||||
import com.zfoo.protocol.util.JsonUtils;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
@@ -14,7 +14,6 @@
|
||||
package com.zfoo.net.session;
|
||||
|
||||
import com.zfoo.net.NetContext;
|
||||
import com.zfoo.net.session.model.Session;
|
||||
import com.zfoo.protocol.util.FileUtils;
|
||||
import com.zfoo.protocol.util.StringUtils;
|
||||
import com.zfoo.util.ThreadUtils;
|
||||
@@ -30,14 +29,14 @@ public abstract class SessionUtils {
|
||||
while (true) {
|
||||
ThreadUtils.sleep(10_000);
|
||||
var builder = new StringBuilder();
|
||||
builder.append(StringUtils.format("clientSession总数:[{}]", NetContext.getSessionManager().getClientSessionMap().size()));
|
||||
builder.append(StringUtils.format("clientSession count:[{}]", NetContext.getSessionManager().getClientSessionMap().size()));
|
||||
builder.append(FileUtils.LS);
|
||||
for (Session session : NetContext.getSessionManager().getClientSessionMap().values()) {
|
||||
builder.append(StringUtils.format("[session:{}]", session.getChannel().remoteAddress()));
|
||||
builder.append(FileUtils.LS);
|
||||
}
|
||||
|
||||
builder.append(StringUtils.format("serverSession总数:[{}]", NetContext.getSessionManager().getServerSessionMap().size()));
|
||||
builder.append(StringUtils.format("serverSession count:[{}]", NetContext.getSessionManager().getServerSessionMap().size()));
|
||||
builder.append(FileUtils.LS);
|
||||
for (Session session : NetContext.getSessionManager().getServerSessionMap().values()) {
|
||||
builder.append(StringUtils.format("[session:{}]", session.getChannel().remoteAddress()));
|
||||
|
||||
Reference in New Issue
Block a user