mirror of
https://github.com/tiennm99/zfoo.git
synced 2026-10-03 09:14:01 +00:00
ref: refactor packet service
This commit is contained in:
1 parent
75d186923c
commit
be11fdc9be
22 files changed
+264
-123
No files matched your search
@@ -44,19 +44,4 @@ public class TunnelClient extends AbstractClient<SocketChannel> {
|
||||
channel.pipeline().addLast(new TunnelClientRouteHandler());
|
||||
}
|
||||
|
||||
public static class DecodedPacketInfo {
|
||||
|
||||
public long sid;
|
||||
|
||||
public Object packet;
|
||||
|
||||
public Object attachment;
|
||||
|
||||
public DecodedPacketInfo(long sid, Object packet, Object attachment) {
|
||||
this.sid = sid;
|
||||
this.packet = packet;
|
||||
this.attachment = attachment;
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
@@ -0,0 +1,113 @@
|
||||
/*
|
||||
* 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.core.proxy;
|
||||
|
||||
import com.zfoo.net.NetContext;
|
||||
import com.zfoo.net.packet.EncodedPacketInfo;
|
||||
import com.zfoo.net.packet.PacketService;
|
||||
import com.zfoo.protocol.buffer.ByteBufUtils;
|
||||
import io.netty.buffer.ByteBuf;
|
||||
import io.netty.channel.Channel;
|
||||
|
||||
|
||||
/**
|
||||
* protocol schema:
|
||||
* <p>
|
||||
* 4byte header length
|
||||
* 1byte flag, 0 is packet, 10 is register client, 30 is register uid
|
||||
* body
|
||||
*
|
||||
* @author jaysunxiao
|
||||
*/
|
||||
public class TunnelProtocolClient2Server {
|
||||
|
||||
public static final byte FLAG_PACKET = 0;
|
||||
public static final byte FLAG_REGISTER = 10;
|
||||
|
||||
private long sid;
|
||||
|
||||
private long uid;
|
||||
|
||||
private ByteBuf byteBuf;
|
||||
|
||||
public static TunnelProtocolClient2Server valueOf(long sid, long uid, ByteBuf byteBuf) {
|
||||
var tunnelProtocol = new TunnelProtocolClient2Server();
|
||||
tunnelProtocol.sid = sid;
|
||||
tunnelProtocol.uid = uid;
|
||||
tunnelProtocol.byteBuf = byteBuf;
|
||||
return tunnelProtocol;
|
||||
}
|
||||
|
||||
public long getSid() {
|
||||
return sid;
|
||||
}
|
||||
|
||||
public long getUid() {
|
||||
return uid;
|
||||
}
|
||||
|
||||
public ByteBuf getByteBuf() {
|
||||
return byteBuf;
|
||||
}
|
||||
|
||||
|
||||
// -----------------------------------------------------------------------------------------------------------------
|
||||
public static void writePacket(ByteBuf out, EncodedPacketInfo encodedPacketInfo) {
|
||||
out.ensureWritable(4);
|
||||
out.writerIndex(PacketService.PACKET_HEAD_LENGTH);
|
||||
out.writeByte(FLAG_PACKET);
|
||||
ByteBufUtils.writeLong(out, encodedPacketInfo.getSid());
|
||||
ByteBufUtils.writeLong(out, encodedPacketInfo.getUid());
|
||||
|
||||
var packet = encodedPacketInfo.getPacket();
|
||||
var attachment = encodedPacketInfo.getAttachment();
|
||||
NetContext.getPacketService().write(out, packet, attachment);
|
||||
NetContext.getPacketService().writeHeaderBefore(out);
|
||||
}
|
||||
|
||||
// -----------------------------------------------------------------------------------------------------------------
|
||||
public static class TunnelRegister {
|
||||
public long sid;
|
||||
|
||||
public TunnelRegister(long sid) {
|
||||
this.sid = sid;
|
||||
}
|
||||
}
|
||||
|
||||
public static void writeRegister(ByteBuf out, TunnelRegister register) {
|
||||
// out.ensureWritable(4);
|
||||
// out.writerIndex(PacketService.PACKET_HEAD_LENGTH);
|
||||
// out.writeByte(FLAG_PACKET);
|
||||
// ByteBufUtils.writeLong(out, encodedPacketInfo.getSid());
|
||||
// ByteBufUtils.writeLong(out, encodedPacketInfo.getUid());
|
||||
//
|
||||
// var packet = encodedPacketInfo.getPacket();
|
||||
// var attachment = encodedPacketInfo.getAttachment();
|
||||
// NetContext.getPacketService().write(out, packet, attachment);
|
||||
// NetContext.getPacketService().writeHeaderBefore(out);
|
||||
}
|
||||
|
||||
// -----------------------------------------------------------------------------------------------------------------
|
||||
public static void read(Channel channel, ByteBuf in) {
|
||||
var flag = in.readByte();
|
||||
if (flag == FLAG_REGISTER) {
|
||||
TunnelServer.tunnels.add(channel);
|
||||
} else if (flag == FLAG_PACKET) {
|
||||
var sid = ByteBufUtils.readLong(in);
|
||||
var uid = ByteBufUtils.readLong(in);
|
||||
var session = NetContext.getSessionManager().getServerSession(sid);
|
||||
session.getChannel().writeAndFlush(TunnelProtocolServer2Client.valueOf(sid, uid, in));
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
@@ -12,6 +12,9 @@
|
||||
|
||||
package com.zfoo.net.core.proxy;
|
||||
|
||||
import com.zfoo.net.NetContext;
|
||||
import com.zfoo.net.packet.PacketService;
|
||||
import com.zfoo.protocol.buffer.ByteBufUtils;
|
||||
import io.netty.buffer.ByteBuf;
|
||||
|
||||
|
||||
@@ -22,28 +25,66 @@ public class TunnelProtocolServer2Client {
|
||||
|
||||
private long sid;
|
||||
|
||||
private long uid;
|
||||
|
||||
private ByteBuf byteBuf;
|
||||
|
||||
public static TunnelProtocolServer2Client valueOf(long sid, ByteBuf byteBuf) {
|
||||
public static TunnelProtocolServer2Client valueOf(long sid, long uid, ByteBuf byteBuf) {
|
||||
var tunnelProtocol = new TunnelProtocolServer2Client();
|
||||
tunnelProtocol.sid = sid;
|
||||
tunnelProtocol.uid = uid;
|
||||
tunnelProtocol.byteBuf = byteBuf;
|
||||
return tunnelProtocol;
|
||||
}
|
||||
|
||||
// -----------------------------------------------------------------------------------------------------------------
|
||||
public static class TunnelPacketInfo {
|
||||
|
||||
public long sid;
|
||||
|
||||
public long uid;
|
||||
|
||||
public Object packet;
|
||||
|
||||
public Object attachment;
|
||||
|
||||
public TunnelPacketInfo(long sid, long uid, Object packet, Object attachment) {
|
||||
this.sid = sid;
|
||||
this.uid = uid;
|
||||
this.packet = packet;
|
||||
this.attachment = attachment;
|
||||
}
|
||||
}
|
||||
|
||||
public static TunnelPacketInfo read(ByteBuf in) {
|
||||
var sid = ByteBufUtils.readLong(in);
|
||||
var uid = ByteBufUtils.readLong(in);
|
||||
var packetInfo = NetContext.getPacketService().read(in);
|
||||
return new TunnelPacketInfo(sid, uid, packetInfo.getPacket(), packetInfo.getAttachment());
|
||||
}
|
||||
// -----------------------------------------------------------------------------------------------------------------
|
||||
|
||||
|
||||
public void write(ByteBuf out) {
|
||||
out.ensureWritable(4);
|
||||
out.writerIndex(PacketService.PACKET_HEAD_LENGTH);
|
||||
ByteBufUtils.writeLong(out, sid);
|
||||
ByteBufUtils.writeLong(out, uid);
|
||||
out.writeBytes(byteBuf);
|
||||
NetContext.getPacketService().writeHeaderBefore(out);
|
||||
}
|
||||
// -----------------------------------------------------------------------------------------------------------------
|
||||
|
||||
|
||||
public long getSid() {
|
||||
return sid;
|
||||
}
|
||||
|
||||
public void setSid(long sid) {
|
||||
this.sid = sid;
|
||||
public long getUid() {
|
||||
return uid;
|
||||
}
|
||||
|
||||
public ByteBuf getByteBuf() {
|
||||
return byteBuf;
|
||||
}
|
||||
|
||||
public void setByteBuf(ByteBuf byteBuf) {
|
||||
this.byteBuf = byteBuf;
|
||||
}
|
||||
}
|
||||
@@ -13,6 +13,7 @@
|
||||
|
||||
package com.zfoo.net.core.proxy.handler;
|
||||
|
||||
import com.zfoo.net.NetContext;
|
||||
import com.zfoo.net.core.proxy.TunnelProtocolServer2Client;
|
||||
import com.zfoo.net.core.proxy.TunnelServer;
|
||||
import com.zfoo.net.packet.PacketService;
|
||||
@@ -56,11 +57,15 @@ public class ProxyCodecHandler extends ByteToMessageCodec<TunnelProtocolServer2C
|
||||
|
||||
var session = SessionUtils.getSession(ctx);
|
||||
var tunnel = RandomUtils.randomEle(TunnelServer.tunnels);
|
||||
tunnel.writeAndFlush(TunnelProtocolServer2Client.valueOf(session.getSid(), sliceByteBuf));
|
||||
tunnel.writeAndFlush(TunnelProtocolServer2Client.valueOf(session.getSid(), session.getUid(), sliceByteBuf));
|
||||
}
|
||||
|
||||
@Override
|
||||
protected void encode(ChannelHandlerContext ctx, TunnelProtocolServer2Client broker, ByteBuf out) {
|
||||
protected void encode(ChannelHandlerContext ctx, TunnelProtocolServer2Client tunnelProtocol, ByteBuf out) {
|
||||
out.ensureWritable(7);
|
||||
out.writerIndex(PacketService.PACKET_HEAD_LENGTH);
|
||||
out.writeBytes(tunnelProtocol.getByteBuf());
|
||||
NetContext.getPacketService().writeHeaderBefore(out);
|
||||
}
|
||||
|
||||
}
|
||||
@@ -13,12 +13,10 @@
|
||||
|
||||
package com.zfoo.net.core.proxy.handler;
|
||||
|
||||
import com.zfoo.net.NetContext;
|
||||
import com.zfoo.net.core.proxy.TunnelClient;
|
||||
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.EncodedPacketInfo;
|
||||
import com.zfoo.net.packet.PacketService;
|
||||
import com.zfoo.protocol.buffer.ByteBufUtils;
|
||||
import com.zfoo.protocol.util.IOUtils;
|
||||
import com.zfoo.protocol.util.StringUtils;
|
||||
import io.netty.buffer.ByteBuf;
|
||||
@@ -32,9 +30,8 @@ import java.util.List;
|
||||
/**
|
||||
* @author jaysunxiao
|
||||
*/
|
||||
public class TunnelClientCodecHandler extends ByteToMessageCodec<TunnelProtocolServer2Client> {
|
||||
public class TunnelClientCodecHandler extends ByteToMessageCodec<Object> {
|
||||
|
||||
private static final Logger logger = LoggerFactory.getLogger(TunnelClientCodecHandler.class);
|
||||
|
||||
@Override
|
||||
protected void decode(ChannelHandlerContext ctx, ByteBuf in, List<Object> out) {
|
||||
@@ -58,13 +55,15 @@ public class TunnelClientCodecHandler extends ByteToMessageCodec<TunnelProtocolS
|
||||
|
||||
var sliceByteBuf = in.readSlice(length);
|
||||
|
||||
var sid = ByteBufUtils.readLong(sliceByteBuf);
|
||||
var packetInfo = NetContext.getPacketService().read(sliceByteBuf);
|
||||
out.add(new TunnelClient.DecodedPacketInfo(sid, packetInfo.getPacket(), packetInfo.getAttachment()));
|
||||
var tunnelPacketInfo = TunnelProtocolServer2Client.read(sliceByteBuf);
|
||||
out.add(tunnelPacketInfo);
|
||||
}
|
||||
|
||||
@Override
|
||||
protected void encode(ChannelHandlerContext ctx, TunnelProtocolServer2Client tunnelProtocol, ByteBuf out) {
|
||||
protected void encode(ChannelHandlerContext ctx, Object msg, ByteBuf out) {
|
||||
if (msg instanceof EncodedPacketInfo) {
|
||||
TunnelProtocolClient2Server.writePacket(out, (EncodedPacketInfo) msg);
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
@@ -15,22 +15,25 @@ package com.zfoo.net.core.proxy.handler;
|
||||
|
||||
import com.zfoo.net.NetContext;
|
||||
import com.zfoo.net.core.proxy.TunnelClient;
|
||||
import com.zfoo.net.handler.BaseRouteHandler;
|
||||
import com.zfoo.net.packet.DecodedPacketInfo;
|
||||
import com.zfoo.net.core.proxy.TunnelProtocolClient2Server;
|
||||
import com.zfoo.net.core.proxy.TunnelProtocolServer2Client;
|
||||
import com.zfoo.net.handler.ClientRouteHandler;
|
||||
import com.zfoo.net.session.Session;
|
||||
import io.netty.channel.ChannelHandler;
|
||||
import io.netty.channel.ChannelHandlerContext;
|
||||
|
||||
/**
|
||||
* @author godotg
|
||||
* @author jaysunxiao
|
||||
*/
|
||||
@ChannelHandler.Sharable
|
||||
public class TunnelClientRouteHandler extends BaseRouteHandler {
|
||||
public class TunnelClientRouteHandler extends ClientRouteHandler {
|
||||
|
||||
@Override
|
||||
public void channelActive(ChannelHandlerContext ctx) throws Exception {
|
||||
super.channelActive(ctx);
|
||||
TunnelClient.tunnels.add(ctx.channel());
|
||||
|
||||
ctx.channel().writeAndFlush(new TunnelProtocolClient2Server.TunnelRegister(1));
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -41,9 +44,9 @@ public class TunnelClientRouteHandler extends BaseRouteHandler {
|
||||
|
||||
@Override
|
||||
public void channelRead(ChannelHandlerContext ctx, Object msg) {
|
||||
var decodedPacketInfo = (TunnelClient.DecodedPacketInfo) msg;
|
||||
var session = new Session(decodedPacketInfo.sid, ctx.channel(), 0);
|
||||
NetContext.getRouter().receive(session, decodedPacketInfo.packet, decodedPacketInfo.attachment);
|
||||
|
||||
var tunnelPacketInfo = (TunnelProtocolServer2Client.TunnelPacketInfo) msg;
|
||||
var session = new Session(tunnelPacketInfo.sid, tunnelPacketInfo.uid, ctx.channel());
|
||||
NetContext.getRouter().receive(session, tunnelPacketInfo.packet, tunnelPacketInfo.attachment);
|
||||
}
|
||||
|
||||
}
|
||||
@@ -13,11 +13,9 @@
|
||||
|
||||
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.protocol.buffer.ByteBufUtils;
|
||||
import com.zfoo.protocol.util.IOUtils;
|
||||
import com.zfoo.protocol.util.StringUtils;
|
||||
import io.netty.buffer.ByteBuf;
|
||||
@@ -56,29 +54,12 @@ public class TunnelServerCodecHandler extends ByteToMessageCodec<TunnelProtocolS
|
||||
}
|
||||
|
||||
var sliceByteBuf = in.readSlice(length);
|
||||
var messageType = sliceByteBuf.readByte();
|
||||
if (messageType == -1) {
|
||||
TunnelServer.tunnels.add(ctx.channel());
|
||||
return;
|
||||
}
|
||||
|
||||
|
||||
var packetInfo = NetContext.getPacketService().read(sliceByteBuf);
|
||||
out.add(packetInfo);
|
||||
TunnelProtocolClient2Server.read(ctx.channel(), sliceByteBuf);
|
||||
}
|
||||
|
||||
@Override
|
||||
protected void encode(ChannelHandlerContext ctx, TunnelProtocolServer2Client tunnelProtocol, ByteBuf out) {
|
||||
out.ensureWritable(4);
|
||||
out.writerIndex(PacketService.PACKET_HEAD_LENGTH);
|
||||
ByteBufUtils.writeLong(out, tunnelProtocol.getSid());
|
||||
out.writeBytes(tunnelProtocol.getByteBuf());
|
||||
|
||||
int length = out.writerIndex();
|
||||
int packetLength = length - PacketService.PACKET_HEAD_LENGTH;
|
||||
out.writerIndex(0);
|
||||
out.writeInt(packetLength);
|
||||
out.writerIndex(length);
|
||||
tunnelProtocol.write(out);
|
||||
}
|
||||
|
||||
}
|
||||
@@ -19,7 +19,7 @@ import io.netty.channel.ChannelHandler;
|
||||
import io.netty.channel.ChannelHandlerContext;
|
||||
|
||||
/**
|
||||
* @author godotg
|
||||
* @author jaysunxiao
|
||||
*/
|
||||
@ChannelHandler.Sharable
|
||||
public class TunnelServerRouteHandler extends ServerRouteHandler {
|
||||
|
||||
@@ -60,7 +60,7 @@ public class TcpCodecHandler extends ByteToMessageCodec<EncodedPacketInfo> {
|
||||
|
||||
@Override
|
||||
protected void encode(ChannelHandlerContext ctx, EncodedPacketInfo packetInfo, ByteBuf out) {
|
||||
NetContext.getPacketService().write(out, packetInfo.getPacket(), packetInfo.getAttachment());
|
||||
NetContext.getPacketService().writeHeaderAndBody(out, packetInfo.getPacket(), packetInfo.getAttachment());
|
||||
}
|
||||
|
||||
}
|
||||
@@ -66,7 +66,7 @@ public class UdpCodecHandler extends MessageToMessageCodec<DatagramPacket, Encod
|
||||
var byteBuf = channelHandlerContext.alloc().ioBuffer();
|
||||
var udpAttachment = (UdpAttachment) out.getAttachment();
|
||||
|
||||
NetContext.getPacketService().write(byteBuf, out.getPacket(), out.getAttachment());
|
||||
NetContext.getPacketService().writeHeaderAndBody(byteBuf, out.getPacket(), out.getAttachment());
|
||||
list.add(new DatagramPacket(byteBuf, new InetSocketAddress(udpAttachment.getHost(), udpAttachment.getPort())));
|
||||
}
|
||||
}
|
||||
@@ -49,7 +49,7 @@ public class WebSocketCodecHandler extends MessageToMessageCodec<WebSocketFrame,
|
||||
@Override
|
||||
protected void encode(ChannelHandlerContext channelHandlerContext, EncodedPacketInfo out, List<Object> list) {
|
||||
var byteBuf = channelHandlerContext.alloc().ioBuffer();
|
||||
NetContext.getPacketService().write(byteBuf, out.getPacket(), out.getAttachment());
|
||||
NetContext.getPacketService().writeHeaderAndBody(byteBuf, out.getPacket(), out.getAttachment());
|
||||
list.add(new BinaryWebSocketFrame(byteBuf));
|
||||
}
|
||||
|
||||
|
||||
@@ -30,7 +30,7 @@ public class ClientIdleHandler extends ChannelDuplexHandler {
|
||||
|
||||
private static final Logger logger = LoggerFactory.getLogger(ClientIdleHandler.class);
|
||||
|
||||
private static final EncodedPacketInfo heartbeatPacket = EncodedPacketInfo.valueOf(new Heartbeat(), null);
|
||||
private static final EncodedPacketInfo heartbeatPacket = EncodedPacketInfo.valueOf(0, 0, new Heartbeat(), null);
|
||||
|
||||
@Override
|
||||
public void userEventTriggered(ChannelHandlerContext ctx, Object evt) {
|
||||
|
||||
@@ -21,6 +21,9 @@ import org.springframework.lang.Nullable;
|
||||
*/
|
||||
public class EncodedPacketInfo {
|
||||
|
||||
private long sid;
|
||||
private long uid;
|
||||
|
||||
/**
|
||||
* 解码后的包
|
||||
*/
|
||||
@@ -32,26 +35,28 @@ public class EncodedPacketInfo {
|
||||
private Object attachment;
|
||||
|
||||
|
||||
public static EncodedPacketInfo valueOf(Object packet, @Nullable Object attachment) {
|
||||
public static EncodedPacketInfo valueOf(long sid, long uid, Object packet, @Nullable Object attachment) {
|
||||
EncodedPacketInfo packetInfo = new EncodedPacketInfo();
|
||||
packetInfo.sid = sid;
|
||||
packetInfo.uid = uid;
|
||||
packetInfo.packet = packet;
|
||||
packetInfo.attachment = attachment;
|
||||
return packetInfo;
|
||||
}
|
||||
|
||||
public long getSid() {
|
||||
return sid;
|
||||
}
|
||||
|
||||
public long getUid() {
|
||||
return uid;
|
||||
}
|
||||
|
||||
public Object getPacket() {
|
||||
return packet;
|
||||
}
|
||||
|
||||
public void setPacket(Object packet) {
|
||||
this.packet = packet;
|
||||
}
|
||||
|
||||
public Object getAttachment() {
|
||||
return attachment;
|
||||
}
|
||||
|
||||
public void setAttachment(Object attachment) {
|
||||
this.attachment = attachment;
|
||||
}
|
||||
}
|
||||
@@ -26,4 +26,7 @@ public interface IPacketService {
|
||||
|
||||
void write(ByteBuf buffer, Object packet, @Nullable Object attachment);
|
||||
|
||||
void writeHeaderAndBody(ByteBuf buffer, Object packet, @Nullable Object attachment);
|
||||
|
||||
void writeHeaderBefore(ByteBuf buffer);
|
||||
}
|
||||
@@ -159,36 +159,46 @@ public class PacketService implements IPacketService {
|
||||
|
||||
@Override
|
||||
public void write(ByteBuf buffer, Object packet, Object attachment) {
|
||||
// 写入包packet
|
||||
ProtocolManager.write(buffer, packet);
|
||||
|
||||
// 写入包的附加包attachment
|
||||
if (attachment == null) {
|
||||
ByteBufUtils.writeBool(buffer, false);
|
||||
} else {
|
||||
ByteBufUtils.writeBool(buffer, true);
|
||||
// 写入包的附加包attachment
|
||||
ProtocolManager.write(buffer, attachment);
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public void writeHeaderAndBody(ByteBuf buffer, Object packet, Object attachment) {
|
||||
try {
|
||||
// 预留写入包的长度,一个int字节大小
|
||||
buffer.ensureWritable(7);
|
||||
buffer.writerIndex(PACKET_HEAD_LENGTH);
|
||||
|
||||
// 写入包packet
|
||||
ProtocolManager.write(buffer, packet);
|
||||
write(buffer, packet, attachment);
|
||||
|
||||
// 写入包的附加包attachment
|
||||
if (attachment == null) {
|
||||
ByteBufUtils.writeBool(buffer, false);
|
||||
} else {
|
||||
ByteBufUtils.writeBool(buffer, true);
|
||||
// 写入包的附加包attachment
|
||||
ProtocolManager.write(buffer, attachment);
|
||||
}
|
||||
|
||||
int length = buffer.writerIndex();
|
||||
|
||||
int packetLength = length - PACKET_HEAD_LENGTH;
|
||||
|
||||
buffer.writerIndex(0);
|
||||
|
||||
buffer.writeInt(packetLength);
|
||||
|
||||
buffer.writerIndex(length);
|
||||
writeHeaderBefore(buffer);
|
||||
} catch (Exception e) {
|
||||
logger.error("write packet exception", e);
|
||||
} catch (Throwable t) {
|
||||
logger.error("write packet error", t);
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public void writeHeaderBefore(ByteBuf buffer) {
|
||||
int length = buffer.writerIndex();
|
||||
|
||||
int packetLength = length - PACKET_HEAD_LENGTH;
|
||||
|
||||
buffer.writerIndex(0);
|
||||
|
||||
buffer.writeInt(packetLength);
|
||||
|
||||
buffer.writerIndex(length);
|
||||
}
|
||||
}
|
||||
@@ -201,7 +201,7 @@ public class Router implements IRouter {
|
||||
if (!channel.isActive()) {
|
||||
return;
|
||||
}
|
||||
var packetInfo = EncodedPacketInfo.valueOf(packet, attachment);
|
||||
var packetInfo = EncodedPacketInfo.valueOf(session.getSid(), session.getUid(), packet, attachment);
|
||||
if (!channel.isWritable()) {
|
||||
logger.warn("send msg error, protocol [{}] sid=[{}] uid=[{}] isActive=[{}] isWritable=[{}]"
|
||||
, packet.getClass().getSimpleName(), session.getSid(), session.getUid(), channel.isActive(), channel.isWritable());
|
||||
|
||||
@@ -55,10 +55,10 @@ public class Session implements Closeable {
|
||||
this.sid = ATOMIC_LONG.incrementAndGet();
|
||||
}
|
||||
|
||||
public Session(long sid, Channel channel, long uid) {
|
||||
public Session(long sid, long uid, Channel channel) {
|
||||
this.sid = sid;
|
||||
this.channel = channel;
|
||||
this.uid = uid;
|
||||
this.channel = channel;
|
||||
}
|
||||
|
||||
@Override
|
||||
|
||||
+2
-3
@@ -15,7 +15,6 @@ package com.zfoo.net.core.proxy.client;
|
||||
|
||||
import com.zfoo.net.anno.PacketReceiver;
|
||||
import com.zfoo.net.packet.proxy.ProxyHelloResponse;
|
||||
import com.zfoo.net.packet.tcp.TcpHelloResponse;
|
||||
import com.zfoo.net.session.Session;
|
||||
import com.zfoo.protocol.util.JsonUtils;
|
||||
import org.slf4j.Logger;
|
||||
@@ -26,10 +25,10 @@ import org.springframework.stereotype.Component;
|
||||
* @author jaysunxiao
|
||||
*/
|
||||
@Component
|
||||
public class ReverseProxyClientController {
|
||||
public class ProxyClientController {
|
||||
|
||||
|
||||
private static final Logger logger = LoggerFactory.getLogger(ReverseProxyClientController.class);
|
||||
private static final Logger logger = LoggerFactory.getLogger(ProxyClientController.class);
|
||||
|
||||
@PacketReceiver
|
||||
public void atProxyHelloResponse(Session session, ProxyHelloResponse response) {
|
||||
+2
-4
@@ -26,15 +26,13 @@ import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
import org.springframework.context.support.ClassPathXmlApplicationContext;
|
||||
|
||||
import java.util.function.Consumer;
|
||||
|
||||
/**
|
||||
* @author jaysunxiao
|
||||
*/
|
||||
@Ignore
|
||||
public class ReverseProxyClientTest {
|
||||
public class ProxyClientTest {
|
||||
|
||||
private static final Logger logger = LoggerFactory.getLogger(ReverseProxyClientTest.class);
|
||||
private static final Logger logger = LoggerFactory.getLogger(ProxyClientTest.class);
|
||||
|
||||
@Test
|
||||
public void startClient() throws Exception {
|
||||
+3
-3
@@ -25,7 +25,7 @@ import org.springframework.context.support.ClassPathXmlApplicationContext;
|
||||
* @author jaysunxiao
|
||||
*/
|
||||
@Ignore
|
||||
public class ReverseProxyServerTest {
|
||||
public class ProxyServerTest {
|
||||
|
||||
/**
|
||||
* ReverseProxyServerTest reverse proxy TargetServerTest
|
||||
@@ -37,8 +37,8 @@ public class ReverseProxyServerTest {
|
||||
var server = new TcpServer(HostAndPort.valueOf("0.0.0.0:9000"));
|
||||
server.start();
|
||||
|
||||
var reverseProxyServer = new TunnelServer(HostAndPort.valueOf("0.0.0.0:9001"));
|
||||
reverseProxyServer.start();
|
||||
var tunnelServer = new TunnelServer(HostAndPort.valueOf("0.0.0.0:9001"));
|
||||
tunnelServer.start();
|
||||
|
||||
ThreadUtils.sleep(Long.MAX_VALUE);
|
||||
}
|
||||
@@ -14,8 +14,7 @@
|
||||
package com.zfoo.net.core.proxy.server;
|
||||
|
||||
import com.zfoo.net.core.HostAndPort;
|
||||
import com.zfoo.net.core.tcp.TcpClient;
|
||||
import com.zfoo.net.core.tcp.TcpServer;
|
||||
import com.zfoo.net.core.proxy.TunnelClient;
|
||||
import com.zfoo.protocol.util.ThreadUtils;
|
||||
import org.junit.Ignore;
|
||||
import org.junit.Test;
|
||||
@@ -36,8 +35,8 @@ public class TargetServerTest {
|
||||
public void startServer() {
|
||||
var context = new ClassPathXmlApplicationContext("config.xml");
|
||||
|
||||
var reverseProxyClient = new TcpClient(HostAndPort.valueOf("127.0.0.1:9001"));
|
||||
var session = reverseProxyClient.start();
|
||||
var tunnelClient = new TunnelClient(HostAndPort.valueOf("127.0.0.1:9001"));
|
||||
var session = tunnelClient.start();
|
||||
|
||||
ThreadUtils.sleep(Long.MAX_VALUE);
|
||||
}
|
||||
|
||||
@@ -66,7 +66,7 @@ public class ProtocolTest {
|
||||
cm.setF("Hello Jaysunxiao,this is the World!!!!!!!!!!!!!!!!!!!!!!!!!!!!");
|
||||
|
||||
ByteBuf writeBuff = Unpooled.directBuffer();
|
||||
packetService.write(writeBuff, cm, attachment);
|
||||
packetService.writeHeaderAndBody(writeBuff, cm, attachment);
|
||||
|
||||
writeBuff.readerIndex(PacketService.PACKET_HEAD_LENGTH);// 信息头的长度
|
||||
|
||||
@@ -83,7 +83,7 @@ public class ProtocolTest {
|
||||
cm.setB(objectA0);
|
||||
|
||||
ByteBuf writeBuff = Unpooled.buffer();
|
||||
packetService.write(writeBuff, cm, null);
|
||||
packetService.writeHeaderAndBody(writeBuff, cm, null);
|
||||
|
||||
writeBuff.readerIndex(PacketService.PACKET_HEAD_LENGTH);// 信息头的长度
|
||||
|
||||
@@ -102,7 +102,7 @@ public class ProtocolTest {
|
||||
cm.setD(Double.MIN_VALUE);
|
||||
|
||||
ByteBuf writeBuff = Unpooled.buffer();
|
||||
packetService.write(writeBuff, cm, null);
|
||||
packetService.writeHeaderAndBody(writeBuff, cm, null);
|
||||
|
||||
writeBuff.readerIndex(PacketService.PACKET_HEAD_LENGTH);// 信息头的长度
|
||||
|
||||
@@ -121,7 +121,7 @@ public class ProtocolTest {
|
||||
cm.setD(100.1);
|
||||
|
||||
ByteBuf writeBuff = Unpooled.buffer();
|
||||
packetService.write(writeBuff, cm, null);
|
||||
packetService.writeHeaderAndBody(writeBuff, cm, null);
|
||||
|
||||
writeBuff.readerIndex(PacketService.PACKET_HEAD_LENGTH);// 信息头的长度
|
||||
|
||||
@@ -156,7 +156,7 @@ public class ProtocolTest {
|
||||
cm.setListListWithMap(listListWithMap);
|
||||
|
||||
ByteBuf writeBuff = Unpooled.buffer();
|
||||
packetService.write(writeBuff, cm, null);
|
||||
packetService.writeHeaderAndBody(writeBuff, cm, null);
|
||||
|
||||
writeBuff.readerIndex(PacketService.PACKET_HEAD_LENGTH);// 信息头的长度
|
||||
|
||||
@@ -186,7 +186,7 @@ public class ProtocolTest {
|
||||
|
||||
|
||||
ByteBuf writeBuff = Unpooled.buffer();
|
||||
packetService.write(writeBuff, cm, null);
|
||||
packetService.writeHeaderAndBody(writeBuff, cm, null);
|
||||
|
||||
writeBuff.readerIndex(PacketService.PACKET_HEAD_LENGTH);// 信息头的长度
|
||||
|
||||
@@ -205,7 +205,7 @@ public class ProtocolTest {
|
||||
cm.setB(array);
|
||||
|
||||
ByteBuf writeBuff = Unpooled.buffer();
|
||||
packetService.write(writeBuff, cm, null);
|
||||
packetService.writeHeaderAndBody(writeBuff, cm, null);
|
||||
|
||||
writeBuff.readerIndex(PacketService.PACKET_HEAD_LENGTH);// 信息头的长度
|
||||
|
||||
@@ -233,7 +233,7 @@ public class ProtocolTest {
|
||||
cm.setMapWithListAndMap(mapWithListAndMap);
|
||||
|
||||
ByteBuf writeBuff = Unpooled.buffer();
|
||||
packetService.write(writeBuff, cm, null);
|
||||
packetService.writeHeaderAndBody(writeBuff, cm, null);
|
||||
|
||||
writeBuff.readerIndex(PacketService.PACKET_HEAD_LENGTH);// 信息头的长度
|
||||
|
||||
|
||||
Reference in new issue
Block a user