public RpcMessageDecoder() { // lengthFieldOffset: magic code is 4B, and version is 1B, and then full length. so value is 5 // lengthFieldLength: full length is 4B. so value is 4 // lengthAdjustment: full length include all data and read 9 bytes before, so the left length is (fullLength-9). so values is -9 // initialBytesToStrip: we will check magic code and version manually, so do not strip any bytes. so values is 0 this(RpcConstants.MAX_FRAME_LENGTH, 5, 4, -9, 0); }
/** * @param maxFrameLength Maximum frame length. It decide the maximum length of data that can be received. * If it exceeds, the data will be discarded. * @param lengthFieldOffset Length field offset. The length field is the one that skips the specified length of byte. * @param lengthFieldLength The number of bytes in the length field. * @param lengthAdjustment The compensation value to add to the value of the length field * @param initialBytesToStrip Number of bytes skipped. * If you need to receive all of the header+body data, this value is 0 * if you only want to receive the body data, then you need to skip the number of bytes consumed by the header. */ public RpcMessageDecoder(int maxFrameLength, int lengthFieldOffset, int lengthFieldLength, int lengthAdjustment, int initialBytesToStrip) { super(maxFrameLength, lengthFieldOffset, lengthFieldLength, lengthAdjustment, initialBytesToStrip); }
private void checkVersion(ByteBuf in) { // read the version and compare byte version = in.readByte(); if (version != RpcConstants.VERSION) { throw new RuntimeException("version isn't compatible" + version); } }
就是读取version,并且检查
checkMagicNumber
private void checkMagicNumber(ByteBuf in) { // read the first 4 bit, which is the magic number, and compare int len = RpcConstants.MAGIC_NUMBER.length; byte[] tmp = new byte[len]; in.readBytes(tmp); for (int i = 0; i < len; i++) { if (tmp[i] != RpcConstants.MAGIC_NUMBER[i]) { throw new IllegalArgumentException("Unknown magic code: " + Arrays.toString(tmp)); } } }
就是一位一位去比较魔术对不对
RpcMessageEncoder
@Slf4j public class RpcMessageEncoder extends MessageToByteEncoder<RpcMessage> { private static final AtomicInteger ATOMIC_INTEGER = new AtomicInteger(0);
@Override protected void encode(ChannelHandlerContext ctx, RpcMessage rpcMessage, ByteBuf out) { try { out.writeBytes(RpcConstants.MAGIC_NUMBER); out.writeByte(RpcConstants.VERSION); // leave a place to write the value of full length out.writerIndex(out.writerIndex() + 4); byte messageType = rpcMessage.getMessageType(); out.writeByte(messageType); out.writeByte(rpcMessage.getCodec()); out.writeByte(CompressTypeEnum.GZIP.getCode()); out.writeInt(ATOMIC_INTEGER.getAndIncrement()); // build full length byte[] bodyBytes = null; int fullLength = RpcConstants.HEAD_LENGTH; // if messageType is not heartbeat message,fullLength = head length + body length if (messageType != RpcConstants.HEARTBEAT_REQUEST_TYPE && messageType != RpcConstants.HEARTBEAT_RESPONSE_TYPE) { // serialize the object String codecName = SerializationTypeEnum.getName(rpcMessage.getCodec()); log.info("codec name: [{}] ", codecName); Serializer serializer = ExtensionLoader.getExtensionLoader(Serializer.class) .getExtension(codecName); bodyBytes = serializer.serialize(rpcMessage.getData()); // compress the bytes String compressName = CompressTypeEnum.getName(rpcMessage.getCompress()); Compress compress = ExtensionLoader.getExtensionLoader(Compress.class) .getExtension(compressName); bodyBytes = compress.compress(bodyBytes); fullLength += bodyBytes.length; }