ref[proxy]: refactor read method of TunnelProtocolClient2Server

This commit is contained in:
jaysunxiao committed 2025-06-30 22:47:13 +08:00
1 parent e690d3f224
commit d4f0ceecf0
2 files changed
+42 -23

No files matched your search

@@ -80,23 +80,5 @@ public class TunnelProtocolClient2Server {
}
// -----------------------------------------------------------------------------------------------------------------
public static void read(Channel channel, ByteBuf in) {
var flag = in.readByte();
if (flag == FLAG_PACKET) {
var sid = ByteBufUtils.readLong(in);
var uid = ByteBufUtils.readLong(in);
var session = NetContext.getSessionManager().getServerSession(sid);
if (SessionUtils.isActive(session)) {
session.getChannel().writeAndFlush(TunnelProtocolServer2Client.valueOf(sid, uid, in));
} else {
ReferenceCountUtil.release(in);
}
} else if (flag == FLAG_REGISTER) {
TunnelServer.tunnels.add(channel);
ReferenceCountUtil.release(in);
} else if (flag == FLAG_HEARTBEAT) {
ReferenceCountUtil.release(in);
}
}
}
@@ -13,15 +13,21 @@
package com.zfoo.net.core.proxy.handler;
import com.zfoo.net.NetContext;
import com.zfoo.net.core.proxy.TunnelProtocolClient2Server;
import com.zfoo.net.core.proxy.TunnelProtocolServer2Client;
import com.zfoo.net.core.proxy.TunnelServer;
import com.zfoo.net.packet.PacketService;
import com.zfoo.net.util.SessionUtils;
import com.zfoo.protocol.buffer.ByteBufUtils;
import com.zfoo.protocol.util.IOUtils;
import com.zfoo.protocol.util.StringUtils;
import io.netty.buffer.ByteBuf;
import io.netty.channel.ChannelHandlerContext;
import io.netty.handler.codec.ByteToMessageCodec;
import io.netty.util.ReferenceCountUtil;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.List;
@@ -30,6 +36,8 @@ import java.util.List;
*/
public class TunnelServerCodecHandler extends ByteToMessageCodec<TunnelProtocolServer2Client> {
private static final Logger logger = LoggerFactory.getLogger(TunnelServerCodecHandler.class);
@Override
protected void decode(ChannelHandlerContext ctx, ByteBuf in, List<Object> out) {
// 不够读一个int
@@ -50,12 +58,41 @@ public class TunnelServerCodecHandler extends ByteToMessageCodec<TunnelProtocolS
return;
}
var retainedByteBuf = in.readRetainedSlice(length);
try {
TunnelProtocolClient2Server.read(ctx.channel(), retainedByteBuf);
} catch (Throwable t) {
ReferenceCountUtil.release(retainedByteBuf);
var flag = in.readByte();
length -= 1;
if (flag == TunnelProtocolClient2Server.FLAG_REGISTER) {
TunnelServer.tunnels.add(ctx.channel());
in.readSlice(length);
return;
}
if (flag == TunnelProtocolClient2Server.FLAG_HEARTBEAT) {
// in.readSlice(length);
return;
}
var beforeReaderIndex = in.readerIndex();
var sid = ByteBufUtils.readLong(in);
var uid = ByteBufUtils.readLong(in);
length -= in.readerIndex() - beforeReaderIndex;
var session = NetContext.getSessionManager().getServerSession(sid);
if (!SessionUtils.isActive(session)) {
in.readSlice(length);
logger.warn("session:[{}] is not activate", sid);
return;
}
if (!session.getChannel().isWritable()) {
in.readSlice(length);
logger.warn("session:[{}] is not writable", sid);
return;
}
var retainedByteBuf = in.readRetainedSlice(length);
session.getChannel().writeAndFlush(TunnelProtocolServer2Client.valueOf(sid, uid, retainedByteBuf));
}
@Override