From 31e8785c630608a90dd7bfe3baa4d807fa26c745 Mon Sep 17 00:00:00 2001 From: xtaci Date: Sun, 27 Oct 2019 13:18:53 +0000 Subject: [PATCH] add a timeout policy for fec buffer --- fec.go | 34 ++++++++++++++++++++++++++++------ 1 file changed, 28 insertions(+), 6 deletions(-) diff --git a/fec.go b/fec.go index 17ff29f..86ce4e0 100644 --- a/fec.go +++ b/fec.go @@ -12,6 +12,7 @@ const ( fecHeaderSizePlus2 = fecHeaderSize + 2 // plus 2B data size typeData = 0xf1 typeParity = 0xf2 + fecExpire = 60000 ) // fecPacket is a decoded FEC packet @@ -21,13 +22,19 @@ func (bts fecPacket) seqid() uint32 { return binary.LittleEndian.Uint32(bts) } func (bts fecPacket) flag() uint16 { return binary.LittleEndian.Uint16(bts[4:]) } func (bts fecPacket) data() []byte { return bts[6:] } +// fecElement has auxcilliary time field +type fecElement struct { + fecPacket + ts uint32 +} + // fecDecoder for decoding incoming packets type fecDecoder struct { rxlimit int // queue size limit dataShards int parityShards int shardSize int - rx []fecPacket // ordered receive queue + rx []fecElement // ordered receive queue // caches decodeCache [][]byte @@ -81,14 +88,15 @@ func (dec *fecDecoder) decode(in fecPacket) (recovered [][]byte) { // make a copy pkt := fecPacket(xmitBuf.Get().([]byte)[:len(in)]) copy(pkt, in) + elem := fecElement{pkt, currentMs()} // insert into ordered rx queue if insertIdx == n+1 { - dec.rx = append(dec.rx, pkt) + dec.rx = append(dec.rx, elem) } else { - dec.rx = append(dec.rx, fecPacket{}) + dec.rx = append(dec.rx, fecElement{}) copy(dec.rx[insertIdx+1:], dec.rx[insertIdx:]) // shift right - dec.rx[insertIdx] = pkt + dec.rx[insertIdx] = elem } // shard range for current packet @@ -171,13 +179,27 @@ func (dec *fecDecoder) decode(in fecPacket) (recovered [][]byte) { } dec.rx = dec.freeRange(0, 1, dec.rx) } + + // timeout policy + current := currentMs() + expiredIdx := -1 + for k := range dec.rx { + if _itimediff(current, dec.rx[k].ts) > 10 { //fecExpire { + expiredIdx = k + continue + } + break + } + if expiredIdx >= 0 { + dec.rx = dec.freeRange(0, expiredIdx, dec.rx) + } return } // free a range of fecPacket -func (dec *fecDecoder) freeRange(first, n int, q []fecPacket) []fecPacket { +func (dec *fecDecoder) freeRange(first, n int, q []fecElement) []fecElement { for i := first; i < first+n; i++ { // recycle buffer - xmitBuf.Put([]byte(q[i])) + xmitBuf.Put([]byte(q[i].fecPacket)) } if first == 0 && n < cap(q)/2 {