From 6e6e064ad42a6a22b8d7cf51ea2961b92dac3f16 Mon Sep 17 00:00:00 2001 From: jaysunxiao Date: Thu, 12 Jun 2025 18:19:07 +0800 Subject: [PATCH] ref[net]: refactor session of client --- .../com/zfoo/net/core/AbstractClient.java | 45 ++++++++++++------- .../zfoo/net/handler/ClientRouteHandler.java | 7 +-- .../java/com/zfoo/net/util/SessionUtils.java | 7 ++- 3 files changed, 38 insertions(+), 21 deletions(-) diff --git a/net/src/main/java/com/zfoo/net/core/AbstractClient.java b/net/src/main/java/com/zfoo/net/core/AbstractClient.java index 09dc93b7..e5484860 100644 --- a/net/src/main/java/com/zfoo/net/core/AbstractClient.java +++ b/net/src/main/java/com/zfoo/net/core/AbstractClient.java @@ -13,12 +13,12 @@ package com.zfoo.net.core; -import com.zfoo.net.NetContext; -import com.zfoo.net.handler.BaseRouteHandler; import com.zfoo.net.session.Session; -import com.zfoo.protocol.exception.ExceptionUtils; +import com.zfoo.net.util.SessionUtils; +import com.zfoo.protocol.exception.RunException; import com.zfoo.protocol.util.IOUtils; import com.zfoo.protocol.util.ThreadUtils; +import com.zfoo.scheduler.util.TimeUtils; import io.netty.bootstrap.Bootstrap; import io.netty.channel.*; import io.netty.channel.epoll.Epoll; @@ -30,6 +30,8 @@ import io.netty.util.concurrent.DefaultThreadFactory; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import java.sql.SQLOutput; + /** * @author godotg */ @@ -66,20 +68,31 @@ public abstract class AbstractClient extends ChannelInitializ var channelFuture = bootstrap.connect(hostAddress, port); channelFuture.syncUninterruptibly(); - if (channelFuture.isSuccess()) { - if (channelFuture.channel().isActive()) { - var channel = channelFuture.channel(); - var session = BaseRouteHandler.initChannel(channel); - NetContext.getSessionManager().addClientSession(session); - logger.info("{} started at [{}]", this.getClass().getSimpleName(), channel.localAddress()); - return session; - } - } else if (channelFuture.cause() != null) { - logger.error(ExceptionUtils.getMessage(channelFuture.cause())); - } else { - logger.error("[{}] started failed", this.getClass().getSimpleName()); + if (channelFuture.cause() != null) { + throw new RuntimeException(channelFuture.cause()); } - return null; + + if (!channelFuture.isSuccess() || !channelFuture.channel().isActive()) { + throw new RunException("[{}] started failed", this.getClass().getSimpleName()); + } + + var channel = channelFuture.channel(); + var session = SessionUtils.getSession(channel); + var loop = 128; + var sleepMillisSeconds = 100; + for (int i = 0; i < loop; i++) { + ThreadUtils.sleep(sleepMillisSeconds); + session = SessionUtils.getSession(channel); + if (session != null) { + break; + } + } + if (session == null) { + channel.close(); + throw new RunException("[{}] start client timeout after [{}] seconds", this.getClass().getSimpleName(), loop * sleepMillisSeconds / TimeUtils.MILLIS_PER_SECOND); + } + logger.info("{} started at [{}]", this.getClass().getSimpleName(), channel.localAddress()); + return session; } diff --git a/net/src/main/java/com/zfoo/net/handler/ClientRouteHandler.java b/net/src/main/java/com/zfoo/net/handler/ClientRouteHandler.java index 74a8cce5..634d57fe 100644 --- a/net/src/main/java/com/zfoo/net/handler/ClientRouteHandler.java +++ b/net/src/main/java/com/zfoo/net/handler/ClientRouteHandler.java @@ -17,6 +17,7 @@ import com.zfoo.event.manager.EventBus; import com.zfoo.net.NetContext; import com.zfoo.net.core.event.ClientSessionActiveEvent; import com.zfoo.net.core.event.ClientSessionInactiveEvent; +import com.zfoo.net.core.event.ServerSessionActiveEvent; import com.zfoo.net.util.SessionUtils; import io.netty.channel.ChannelHandler; import io.netty.channel.ChannelHandlerContext; @@ -34,10 +35,10 @@ public class ClientRouteHandler extends BaseRouteHandler { @Override public void channelActive(ChannelHandlerContext ctx) throws Exception { super.channelActive(ctx); - // 客户端的session初始化在启动的时候已经做了,这边直接获取session - var session = SessionUtils.getSession(ctx); - EventBus.post(ClientSessionActiveEvent.valueOf(session)); + var session = initChannel(ctx.channel()); + NetContext.getSessionManager().addClientSession(session); logger.info("client channel is active {}", SessionUtils.sessionInfo(ctx)); + EventBus.post(ClientSessionActiveEvent.valueOf(session)); } @Override diff --git a/net/src/main/java/com/zfoo/net/util/SessionUtils.java b/net/src/main/java/com/zfoo/net/util/SessionUtils.java index 385b38e5..5cf0875a 100644 --- a/net/src/main/java/com/zfoo/net/util/SessionUtils.java +++ b/net/src/main/java/com/zfoo/net/util/SessionUtils.java @@ -40,8 +40,11 @@ public abstract class SessionUtils { } public static Session getSession(ChannelHandlerContext ctx) { - var sessionAttr = ctx.channel().attr(SESSION_KEY); - return sessionAttr.get(); + return getSession(ctx.channel()); + } + + public static Session getSession(Channel channel) { + return channel.attr(SESSION_KEY).get(); } public static String toIp(Session session) {