perf[net]: zfoo基于netty的http服务器实现

This commit is contained in:
jaysunxiao
2021-08-10 16:24:22 +08:00
parent 9a8c93abe5
commit b403f6ff26
5 changed files with 52 additions and 50 deletions
@@ -20,6 +20,7 @@ import com.zfoo.protocol.IPacket;
import com.zfoo.util.net.HostAndPort;
import io.netty.channel.ChannelInitializer;
import io.netty.channel.socket.SocketChannel;
import io.netty.handler.codec.http.FullHttpRequest;
import io.netty.handler.codec.http.HttpObjectAggregator;
import io.netty.handler.codec.http.HttpServerCodec;
import io.netty.handler.stream.ChunkedWriteHandler;
@@ -35,9 +36,9 @@ public class HttpServer extends AbstractServer {
/**
* http的地址解析器
*/
private Function<String, IPacket> uriResolver;
private Function<FullHttpRequest, IPacket> uriResolver;
public HttpServer(HostAndPort host, Function<String, IPacket> uriResolver) {
public HttpServer(HostAndPort host, Function<FullHttpRequest, IPacket> uriResolver) {
super(host);
this.uriResolver = uriResolver;
}
@@ -19,6 +19,7 @@ import com.zfoo.net.dispatcher.model.vo.EnhanceUtils;
import com.zfoo.net.dispatcher.model.vo.IPacketReceiver;
import com.zfoo.net.dispatcher.model.vo.PacketReceiverDefinition;
import com.zfoo.net.packet.model.GatewayPacketAttachment;
import com.zfoo.net.packet.model.HttpPacketAttachment;
import com.zfoo.net.packet.model.IPacketAttachment;
import com.zfoo.net.packet.model.UdpPacketAttachment;
import com.zfoo.net.packet.service.PacketService;
@@ -112,7 +113,7 @@ public abstract class PacketBus {
// 如果以Ask结尾的请求,那么attachment不能为GatewayAttachment
if (attachmentClazz != null) {
if (packetName.endsWith(PacketService.NET_REQUEST_SUFFIX)) {
AssertionUtils.isTrue(attachmentClazz.equals(GatewayPacketAttachment.class) || attachmentClazz.equals(UdpPacketAttachment.class)
AssertionUtils.isTrue(attachmentClazz.equals(GatewayPacketAttachment.class) || attachmentClazz.equals(UdpPacketAttachment.class) || attachmentClazz.equals(HttpPacketAttachment.class)
, "[class:{}] [method:{}] [packet:{}] must use [attachment:{}]!", bean.getClass().getName(), methodName, packetName, GatewayPacketAttachment.class.getCanonicalName());
} else if (packetName.endsWith(PacketService.NET_ASK_SUFFIX)) {
AssertionUtils.isTrue(!attachmentClazz.equals(GatewayPacketAttachment.class)
@@ -12,24 +12,18 @@
package com.zfoo.net.handler.codec.http;
import com.zfoo.net.NetContext;
import com.zfoo.net.packet.common.Message;
import com.zfoo.net.packet.model.DecodedPacketInfo;
import com.zfoo.net.packet.model.EncodedPacketInfo;
import com.zfoo.net.packet.model.HttpPacketAttachment;
import com.zfoo.net.packet.service.PacketService;
import com.zfoo.net.util.SessionUtils;
import com.zfoo.protocol.IPacket;
import com.zfoo.protocol.util.JsonUtils;
import com.zfoo.protocol.util.StringUtils;
import io.netty.buffer.ByteBuf;
import io.netty.buffer.Unpooled;
import io.netty.channel.ChannelHandlerContext;
import io.netty.handler.codec.MessageToMessageCodec;
import io.netty.handler.codec.http.DefaultFullHttpResponse;
import io.netty.handler.codec.http.FullHttpRequest;
import io.netty.handler.codec.http.HttpResponseStatus;
import io.netty.handler.codec.http.HttpVersion;
import io.netty.util.ReferenceCountUtil;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -44,66 +38,48 @@ public class HttpCodecHandler extends MessageToMessageCodec<FullHttpRequest, Enc
private static final Logger logger = LoggerFactory.getLogger(HttpCodecHandler.class);
private Function<String, IPacket> uriResolver;
private Function<FullHttpRequest, IPacket> uriResolver;
public HttpCodecHandler(Function<String, IPacket> uriResolver) {
public HttpCodecHandler(Function<FullHttpRequest, IPacket> uriResolver) {
super();
this.uriResolver = uriResolver;
}
@Override
protected void decode(ChannelHandlerContext channelHandlerContext, FullHttpRequest fullHttpRequest, List<Object> list) {
var uri = fullHttpRequest.uri();
ByteBuf in = fullHttpRequest.content();
// 不够读一个int
if (in.readableBytes() <= PacketService.PACKET_HEAD_LENGTH) {
return;
}
in.markReaderIndex();
var length = in.readInt();
// 如果长度非法,则抛出异常断开连接
if (length < 0) {
throw new IllegalArgumentException(StringUtils.format("[session:{}]的包头长度[length:{}]非法"
, SessionUtils.sessionInfo(channelHandlerContext), length));
}
// ByteBuf里的数据太小
if (in.readableBytes() < length) {
in.resetReaderIndex();
return;
}
ByteBuf tmpByteBuf = null;
try {
tmpByteBuf = in.readRetainedSlice(length);
DecodedPacketInfo packetInfo = NetContext.getPacketService().read(tmpByteBuf);
packetInfo.setPacketAttachment(HttpPacketAttachment.valueOf());
list.add(packetInfo);
var packet = uriResolver.apply(fullHttpRequest);
var attachment = HttpPacketAttachment.valueOf(fullHttpRequest, HttpResponseStatus.OK);
var decodedPacketInfo = DecodedPacketInfo.valueOf(packet, attachment);
list.add(decodedPacketInfo);
} catch (Exception e) {
logger.error("exception异常", e);
throw e;
} catch (Throwable t) {
logger.error("throwable错误", t);
throw t;
} finally {
ReferenceCountUtil.release(tmpByteBuf);
}
}
@Override
protected void encode(ChannelHandlerContext channelHandlerContext, EncodedPacketInfo out, List<Object> list) {
try {
var byteBuf = channelHandlerContext.alloc().ioBuffer();
var httpPacketAttachment = (HttpPacketAttachment) out.getPacketAttachment();
var packet = (IPacket) out.getPacket();
var attachment = (HttpPacketAttachment) out.getPacketAttachment();
var fullHttpResponse = new DefaultFullHttpResponse(HttpVersion.HTTP_1_1, HttpResponseStatus.OK, Unpooled.wrappedBuffer("I am ok".getBytes()));
list.add(fullHttpResponse);
var protocolVersion = attachment.getFullHttpRequest().protocolVersion();
var httpResponseStatus = attachment.getHttpResponseStatus();
if (packet.protocolId() == Message.PROTOCOL_ID) {
var fullHttpResponse = new DefaultFullHttpResponse(protocolVersion, httpResponseStatus);
list.add(fullHttpResponse);
} else {
var byteBuf = channelHandlerContext.alloc().ioBuffer();
var jsonStr = JsonUtils.object2StringTurbo(packet);
byteBuf.writeBytes(StringUtils.bytes(jsonStr));
var fullHttpResponse = new DefaultFullHttpResponse(protocolVersion, httpResponseStatus, byteBuf);
list.add(fullHttpResponse);
}
} catch (Exception e) {
logger.error("[{}]编码exception异常", JsonUtils.object2String(out), e);
throw e;
@@ -26,6 +26,8 @@ public class Message implements IPacket {
public static final transient short PROTOCOL_ID = 100;
public static final Message DEFAULT = new Message();
private byte module;
/**
@@ -14,6 +14,8 @@
package com.zfoo.net.packet.model;
import com.zfoo.util.math.RandomUtils;
import io.netty.handler.codec.http.FullHttpRequest;
import io.netty.handler.codec.http.HttpResponseStatus;
/**
* @author jaysunxiao
@@ -23,9 +25,14 @@ public class HttpPacketAttachment implements IPacketAttachment {
public static final transient short PROTOCOL_ID = 3;
private transient FullHttpRequest fullHttpRequest;
public static HttpPacketAttachment valueOf() {
private transient HttpResponseStatus httpResponseStatus;
public static HttpPacketAttachment valueOf(FullHttpRequest fullHttpRequest, HttpResponseStatus httpResponseStatus) {
var attachment = new HttpPacketAttachment();
attachment.fullHttpRequest = fullHttpRequest;
attachment.httpResponseStatus = httpResponseStatus;
return attachment;
}
@@ -44,4 +51,19 @@ public class HttpPacketAttachment implements IPacketAttachment {
return PROTOCOL_ID;
}
public FullHttpRequest getFullHttpRequest() {
return fullHttpRequest;
}
public void setFullHttpRequest(FullHttpRequest fullHttpRequest) {
this.fullHttpRequest = fullHttpRequest;
}
public HttpResponseStatus getHttpResponseStatus() {
return httpResponseStatus;
}
public void setHttpResponseStatus(HttpResponseStatus httpResponseStatus) {
this.httpResponseStatus = httpResponseStatus;
}
}