diff --git a/net/src/main/java/com/zfoo/net/core/http/HttpServer.java b/net/src/main/java/com/zfoo/net/core/http/HttpServer.java index 36eeb442..2f0e4a7a 100644 --- a/net/src/main/java/com/zfoo/net/core/http/HttpServer.java +++ b/net/src/main/java/com/zfoo/net/core/http/HttpServer.java @@ -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 uriResolver; + private Function uriResolver; - public HttpServer(HostAndPort host, Function uriResolver) { + public HttpServer(HostAndPort host, Function uriResolver) { super(host); this.uriResolver = uriResolver; } diff --git a/net/src/main/java/com/zfoo/net/dispatcher/manager/PacketBus.java b/net/src/main/java/com/zfoo/net/dispatcher/manager/PacketBus.java index d54dc0ae..7930c03f 100644 --- a/net/src/main/java/com/zfoo/net/dispatcher/manager/PacketBus.java +++ b/net/src/main/java/com/zfoo/net/dispatcher/manager/PacketBus.java @@ -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) diff --git a/net/src/main/java/com/zfoo/net/handler/codec/http/HttpCodecHandler.java b/net/src/main/java/com/zfoo/net/handler/codec/http/HttpCodecHandler.java index 7d06ca1f..57cce4aa 100644 --- a/net/src/main/java/com/zfoo/net/handler/codec/http/HttpCodecHandler.java +++ b/net/src/main/java/com/zfoo/net/handler/codec/http/HttpCodecHandler.java @@ -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 uriResolver; + private Function uriResolver; - public HttpCodecHandler(Function uriResolver) { + public HttpCodecHandler(Function uriResolver) { super(); this.uriResolver = uriResolver; } @Override protected void decode(ChannelHandlerContext channelHandlerContext, FullHttpRequest fullHttpRequest, List 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 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; diff --git a/net/src/main/java/com/zfoo/net/packet/common/Message.java b/net/src/main/java/com/zfoo/net/packet/common/Message.java index 49a4960e..a34d8c09 100644 --- a/net/src/main/java/com/zfoo/net/packet/common/Message.java +++ b/net/src/main/java/com/zfoo/net/packet/common/Message.java @@ -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; /** diff --git a/net/src/main/java/com/zfoo/net/packet/model/HttpPacketAttachment.java b/net/src/main/java/com/zfoo/net/packet/model/HttpPacketAttachment.java index 3c6feaa7..ee06b591 100644 --- a/net/src/main/java/com/zfoo/net/packet/model/HttpPacketAttachment.java +++ b/net/src/main/java/com/zfoo/net/packet/model/HttpPacketAttachment.java @@ -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; + } }