package org.kne.cloud.network.klalb; import java.io.DataInputStream; import java.io.DataOutputStream; import java.io.EOFException; import java.io.IOException; import java.io.StreamCorruptedException; import java.nio.ByteBuffer; import java.nio.channels.Channels; import java.nio.channels.ReadableByteChannel; import java.nio.channels.WritableByteChannel; import org.kne.cloud.network.NetworkPacket; import org.kne.cloud.network.ipv6.IPv6Packet; import org.kne.cloud.network.perf.TEST1Packet; import org.kne.cloud.network.perf.TESTNPacket; public abstract class KLALBPacket extends IPv6Packet.IPv6Payload { public static final int KLALB_PROTOCOL_NUMBER=254; public static final int PING=0; public static final int PONG=1; public static final int BWINF=2; //public static final int SYNT=2; //public static final int SACKT=3; public static final int RST=4; public static final int DATAT=5; public static final int ACKT=6; public static final int DATAU=7; public static final int VADDR=8; public static final int ADDLINES=9; public static final int NACKT=10; public static final int TESTN=11; public static final int VADDRACK=12; public static final int VADDRREQ=13; public static final int IPV6OVERKLALB=14; public static final int ADDR=15; public static final int ADDRREQ=16; public static final int TEST1=17; //private static final int HEADER_CAPACITY = 32; private int headerLength=1; //public static final ByteArrayPool dataarraypool=new ByteArrayPool(5000, 8192); protected ByteBuffer klalbHeader; protected KLALBPacket(ByteBuffer klalbHeader,int headerLength) { super(KLALB_PROTOCOL_NUMBER); this.klalbHeader = klalbHeader; this.headerLength=headerLength; } public KLALBPacket(int type,int headerLength) { super(KLALB_PROTOCOL_NUMBER); klalbHeader=NetworkPacket.bufferAllocator.allocate(headerLength); klalbHeader.put((byte) type); this.headerLength=headerLength; } @Override public String toString() { return "KLALBPacket [type=" + getType() + "]"; } public int getType() { return klalbHeader.get(0)&0xff; } private long sndtime,rcvtime; private long joinqueuetime; public void markJoinqueuetime() { joinqueuetime=System.nanoTime(); } public long getJoinqueuetime() { return joinqueuetime; } public void writeToChannel(WritableByteChannel dto) throws IOException { sndtime=System.nanoTime(); dto.write(klalbHeader.slice(0, headerLength)); } @Override protected boolean needEndPosition() { return false; } public void readFromChannel(ReadableByteChannel din,long length) throws IOException { klalbHeader.limit(headerLength); while(klalbHeader.hasRemaining()){ if(din.read(klalbHeader)==-1) { throw new EOFException(); } } rcvtime=System.nanoTime(); } public long getSndtime() { return sndtime; } public long getRcvtime() { return rcvtime; } public long getTotalLength() {//缓冲区limit,实际长度 return headerLength; } protected ByteBuffer getKLALBHeader() { return klalbHeader; } public static KLALBPacket readKLALBPacketFromStream(DataInputStream in) throws IOException { return readKLALBPacketFromChannel(Channels.newChannel(in)); } public static KLALBPacket readKLALBPacketFromChannel(ReadableByteChannel in) throws IOException { while(true) { ByteBuffer bb=NetworkPacket.bufferAllocator.allocate(64); bb.position(0); bb.limit(1); while(bb.hasRemaining()){ if(in.read(bb)==-1) { return null; } } int type=bb.get(0)&0xff; bb.limit(bb.capacity()); KLALBPacket klp; switch(type) { case PING:klp=new PINGPacket(bb); klp.readFromChannel(in); return klp; case PONG: klp=new PONGPacket(bb); klp.readFromChannel(in); return klp; case BWINF: klp=new BWINFPacket(bb); klp.readFromChannel(in); return klp; case RST: klp=new RSTPacket(bb); klp.readFromChannel(in); return klp; case DATAT: klp=new DATATPacket(bb); klp.readFromChannel(in); return klp; case ACKT: klp=new ACKTPacket(bb); klp.readFromChannel(in); return klp; case VADDR: klp=new VADDRPacket(bb); klp.readFromChannel(in); return klp; case ADDLINES: klp=new ADDLINESPacket(bb); klp.readFromChannel(in); return klp; case NACKT: klp=new NACKTPacket(bb); klp.readFromChannel(in); return klp; case TESTN: klp=new TESTNPacket(bb); klp.readFromChannel(in); return klp; case VADDRACK: klp=new VADDRACKPacket(bb); klp.readFromChannel(in); return klp; case VADDRREQ: klp=new VADDRREQPacket(bb); klp.readFromChannel(in); return klp; case IPV6OVERKLALB: klp=new IPv6OverKLALBPacket(bb); klp.readFromChannel(in); return klp; case ADDR: klp=new ADDRPacket(bb); klp.readFromChannel(in); return klp; case ADDRREQ: klp=new ADDRREQPacket(bb); klp.readFromChannel(in); return klp; case TEST1: klp=new TEST1Packet(bb); klp.readFromChannel(in); return klp; } throw new StreamCorruptedException("unknown package type:"+type); //System.err.println("ignore unknown KLALBPacket type:"+type); } } public static void writeKLALBPacketToStream(DataOutputStream out,KLALBPacket klb) throws IOException { writeKLALBPacketToChannel(Channels.newChannel(out),klb); } public static void writeKLALBPacketToChannel(WritableByteChannel writableByteChannel,KLALBPacket klb) throws IOException { klb.writeToChannel(writableByteChannel); } public static KLALBPacket createByBuffer(ByteBuffer data) { return null; } private boolean ce=false; public void setCE(boolean ce) { this.ce=ce; } public boolean isCE() { return ce; } }