ref[net]: refactor udp server and client

This commit is contained in:
jaysunxiao committed 2025-06-13 10:42:19 +08:00
1 parent 6e6e064ad4
commit 7a8395b951
3 files changed
+58 -30

No files matched your search

@@ -17,10 +17,9 @@ import com.zfoo.net.NetContext;
import com.zfoo.net.core.AbstractClient;
import com.zfoo.net.core.HostAndPort;
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.Session;
import com.zfoo.protocol.exception.ExceptionUtils;
import com.zfoo.protocol.exception.RunException;
import io.netty.bootstrap.Bootstrap;
import io.netty.channel.Channel;
import io.netty.channel.ChannelOption;
@@ -39,39 +38,31 @@ public class UdpClient extends AbstractClient<Channel> {
@Override
public synchronized Session start() {
try {
this.bootstrap = new Bootstrap();
this.bootstrap.group(nioEventLoopGroup)
.channel(Epoll.isAvailable() ? EpollDatagramChannel.class : NioDatagramChannel.class)
.option(ChannelOption.SO_BROADCAST, true)
.handler(this);
this.bootstrap = new Bootstrap();
this.bootstrap.group(nioEventLoopGroup)
.channel(Epoll.isAvailable() ? EpollDatagramChannel.class : NioDatagramChannel.class)
.option(ChannelOption.SO_BROADCAST, true)
.handler(this);
// bind(0)随机选择一个端口
var channelFuture = bootstrap.bind(0).sync();
channelFuture.syncUninterruptibly();
// bind(0)随机选择一个端口
var channelFuture = bootstrap.bind(0);
channelFuture.syncUninterruptibly();
if (channelFuture.isSuccess()) {
if (channelFuture.channel().isActive()) {
var channel = channelFuture.channel();
var session = BaseRouteHandler.initChannel(channel);
NetContext.getSessionManager().addClientSession(session);
logger.info("UdpClient started at [{}]", channel.localAddress());
return session;
}
} else if (channelFuture.cause() != null) {
logger.error(ExceptionUtils.getMessage(channelFuture.cause()));
} else {
logger.error("[{}] started failed", this.getClass().getSimpleName());
}
} catch (Exception e) {
logger.error(ExceptionUtils.getMessage(e));
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 session = BaseRouteHandler.initChannel(channelFuture.channel());
return session;
}
@Override
protected void initChannel(Channel channel) {
channel.pipeline().addLast(new UdpCodecHandler());
channel.pipeline().addLast(new ClientRouteHandler());
channel.pipeline().addLast(new UdpRouteHandler());
}
}
@@ -0,0 +1,38 @@
/*
* 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.udp;
import com.zfoo.net.NetContext;
import com.zfoo.net.handler.BaseRouteHandler;
import com.zfoo.net.packet.DecodedPacketInfo;
import com.zfoo.net.util.SessionUtils;
import io.netty.channel.ChannelHandler;
import io.netty.channel.ChannelHandlerContext;
/**
* @author godotg
*/
@ChannelHandler.Sharable
public class UdpRouteHandler extends BaseRouteHandler {
@Override
public void channelRead(ChannelHandlerContext ctx, Object msg) {
var session = SessionUtils.getSession(ctx);
if (session == null) {
session = initChannel(ctx.channel());
}
DecodedPacketInfo decodedPacketInfo = (DecodedPacketInfo) msg;
NetContext.getRouter().receive(session, decodedPacketInfo.getPacket(), decodedPacketInfo.getAttachment());
}
}
@@ -14,7 +14,6 @@ package com.zfoo.net.core.udp;
import com.zfoo.net.core.AbstractServer;
import com.zfoo.net.core.HostAndPort;
import com.zfoo.net.handler.ServerRouteHandler;
import com.zfoo.net.handler.codec.udp.UdpCodecHandler;
import io.netty.bootstrap.Bootstrap;
import io.netty.channel.Channel;
@@ -65,6 +64,6 @@ public class UdpServer extends AbstractServer<Channel> {
@Override
protected void initChannel(Channel channel) {
channel.pipeline().addLast(new UdpCodecHandler());
channel.pipeline().addLast(new ServerRouteHandler());
channel.pipeline().addLast(new UdpRouteHandler());
}
}