v2ray-core/transport/internet/kcp/sending.go

383 lines
7.7 KiB
Go
Raw Normal View History

2016-06-26 21:51:17 +00:00
package kcp
2016-07-03 20:14:38 +00:00
import (
"sync"
2017-12-03 21:53:00 +00:00
2017-12-17 00:22:39 +00:00
"v2ray.com/core/common"
2017-12-03 21:53:00 +00:00
"v2ray.com/core/common/buf"
2016-07-03 20:14:38 +00:00
)
2016-07-01 09:57:13 +00:00
type SendingWindow struct {
start uint32
cap uint32
len uint32
last uint32
2016-11-01 11:07:20 +00:00
data []DataSegment
inuse []bool
prev []uint32
next []uint32
2016-07-01 09:57:13 +00:00
2016-07-04 13:34:14 +00:00
totalInFlightSize uint32
writer SegmentWriter
onPacketLoss func(uint32)
2016-07-01 09:57:13 +00:00
}
2016-07-04 13:54:18 +00:00
func NewSendingWindow(size uint32, writer SegmentWriter, onPacketLoss func(uint32)) *SendingWindow {
2016-07-01 09:57:13 +00:00
window := &SendingWindow{
2016-07-03 20:14:38 +00:00
start: 0,
cap: size,
len: 0,
last: 0,
2016-11-01 11:07:20 +00:00
data: make([]DataSegment, size),
2016-07-03 20:14:38 +00:00
prev: make([]uint32, size),
next: make([]uint32, size),
2016-11-01 11:07:20 +00:00
inuse: make([]bool, size),
2016-07-03 20:14:38 +00:00
writer: writer,
onPacketLoss: onPacketLoss,
2016-07-01 09:57:13 +00:00
}
return window
}
2017-12-03 13:56:00 +00:00
func (sw *SendingWindow) Release() {
if sw == nil {
2016-11-21 21:41:12 +00:00
return
}
2017-12-03 13:56:00 +00:00
sw.len = 0
for _, seg := range sw.data {
2016-11-21 21:41:12 +00:00
seg.Release()
}
}
2017-12-03 13:56:00 +00:00
func (sw *SendingWindow) Len() int {
return int(sw.len)
2016-07-01 09:57:13 +00:00
}
2017-12-03 13:56:00 +00:00
func (sw *SendingWindow) IsEmpty() bool {
return sw.len == 0
2016-07-12 21:54:54 +00:00
}
2017-12-03 13:56:00 +00:00
func (sw *SendingWindow) Size() uint32 {
return sw.cap
2016-07-04 13:54:18 +00:00
}
2017-12-03 13:56:00 +00:00
func (sw *SendingWindow) IsFull() bool {
return sw.len == sw.cap
2016-07-04 13:54:18 +00:00
}
2017-12-03 21:53:00 +00:00
func (sw *SendingWindow) Push(number uint32) *buf.Buffer {
2017-12-03 13:56:00 +00:00
pos := (sw.start + sw.len) % sw.cap
sw.data[pos].Number = number
sw.data[pos].timeout = 0
sw.data[pos].transmit = 0
sw.inuse[pos] = true
if sw.len > 0 {
sw.next[sw.last] = pos
sw.prev[pos] = sw.last
2016-07-01 09:57:13 +00:00
}
2017-12-03 13:56:00 +00:00
sw.last = pos
sw.len++
2017-12-03 21:53:00 +00:00
return sw.data[pos].Data()
2016-07-01 09:57:13 +00:00
}
2017-12-03 13:56:00 +00:00
func (sw *SendingWindow) FirstNumber() uint32 {
return sw.data[sw.start].Number
2016-07-01 09:57:13 +00:00
}
2017-12-03 13:56:00 +00:00
func (sw *SendingWindow) Clear(una uint32) {
for !sw.IsEmpty() && sw.data[sw.start].Number < una {
sw.Remove(0)
2016-07-01 09:57:13 +00:00
}
}
2017-12-03 13:56:00 +00:00
func (sw *SendingWindow) Remove(idx uint32) bool {
if sw.IsEmpty() {
2016-11-13 21:27:58 +00:00
return false
2016-07-01 21:27:57 +00:00
}
2017-12-03 13:56:00 +00:00
pos := (sw.start + idx) % sw.cap
if !sw.inuse[pos] {
2016-11-13 21:27:58 +00:00
return false
2016-07-01 10:12:32 +00:00
}
2017-12-03 13:56:00 +00:00
sw.inuse[pos] = false
sw.totalInFlightSize--
if pos == sw.start && pos == sw.last {
sw.len = 0
sw.start = 0
sw.last = 0
} else if pos == sw.start {
delta := sw.next[pos] - sw.start
if sw.next[pos] < sw.start {
delta = sw.next[pos] + sw.cap - sw.start
2016-07-01 09:57:13 +00:00
}
2017-12-03 13:56:00 +00:00
sw.start = sw.next[pos]
sw.len -= delta
} else if pos == sw.last {
sw.last = sw.prev[pos]
2016-07-01 09:57:13 +00:00
} else {
2017-12-03 13:56:00 +00:00
sw.next[sw.prev[pos]] = sw.next[pos]
sw.prev[sw.next[pos]] = sw.prev[pos]
2016-07-01 09:57:13 +00:00
}
2016-11-13 21:27:58 +00:00
return true
2016-07-01 09:57:13 +00:00
}
2017-12-03 13:56:00 +00:00
func (sw *SendingWindow) HandleFastAck(number uint32, rto uint32) {
if sw.IsEmpty() {
2016-07-01 21:27:57 +00:00
return
}
2016-07-01 10:12:32 +00:00
2017-12-03 13:56:00 +00:00
sw.Visit(func(seg *DataSegment) bool {
2016-11-18 15:19:13 +00:00
if number == seg.Number || number-seg.Number > 0x7FFFFFFF {
return false
2016-07-01 09:57:13 +00:00
}
2016-11-18 15:19:13 +00:00
if seg.transmit > 0 && seg.timeout > rto/3 {
seg.timeout -= rto / 3
2016-07-01 09:57:13 +00:00
}
2016-11-18 15:19:13 +00:00
return true
})
}
2017-12-03 13:56:00 +00:00
func (sw *SendingWindow) Visit(visitor func(seg *DataSegment) bool) {
if sw.IsEmpty() {
2016-12-06 23:31:01 +00:00
return
}
2017-12-03 13:56:00 +00:00
for i := sw.start; ; i = sw.next[i] {
if !visitor(&sw.data[i]) || i == sw.last {
2016-07-01 09:57:13 +00:00
break
}
}
}
2017-12-03 13:56:00 +00:00
func (sw *SendingWindow) Flush(current uint32, rto uint32, maxInFlightSize uint32) {
if sw.IsEmpty() {
2016-07-03 20:14:38 +00:00
return
2016-07-01 10:12:32 +00:00
}
2016-07-04 13:34:14 +00:00
var lost uint32
2016-07-04 11:37:42 +00:00
var inFlightSize uint32
2016-07-01 09:57:13 +00:00
2017-12-03 13:56:00 +00:00
sw.Visit(func(segment *DataSegment) bool {
2016-11-18 15:19:13 +00:00
if current-segment.timeout >= 0x7FFFFFFF {
return true
2016-07-01 09:57:13 +00:00
}
2016-11-18 15:19:13 +00:00
if segment.transmit == 0 {
// First time
2017-12-03 13:56:00 +00:00
sw.totalInFlightSize++
2016-11-18 15:19:13 +00:00
} else {
lost++
2016-07-01 09:57:13 +00:00
}
2016-11-18 15:19:13 +00:00
segment.timeout = current + rto
segment.Timestamp = current
segment.transmit++
2017-12-03 13:56:00 +00:00
sw.writer.Write(segment)
2016-11-18 15:19:13 +00:00
inFlightSize++
if inFlightSize >= maxInFlightSize {
return false
2016-07-01 09:57:13 +00:00
}
2016-11-18 15:19:13 +00:00
return true
})
2016-07-01 09:57:13 +00:00
2017-12-03 13:56:00 +00:00
if sw.onPacketLoss != nil && inFlightSize > 0 && sw.totalInFlightSize != 0 {
rate := lost * 100 / sw.totalInFlightSize
sw.onPacketLoss(rate)
2016-07-04 13:34:14 +00:00
}
2016-07-01 09:57:13 +00:00
}
2016-07-03 20:14:38 +00:00
type SendingWorker struct {
2016-07-12 15:56:36 +00:00
sync.RWMutex
2016-10-11 10:24:19 +00:00
conn *Connection
window *SendingWindow
firstUnacknowledged uint32
firstUnacknowledgedUpdated bool
nextNumber uint32
remoteNextNumber uint32
controlWindow uint32
fastResend uint32
2016-07-03 20:14:38 +00:00
}
2016-07-05 21:02:52 +00:00
func NewSendingWorker(kcp *Connection) *SendingWorker {
2016-07-03 20:14:38 +00:00
worker := &SendingWorker{
2016-07-05 21:02:52 +00:00
conn: kcp,
2016-07-03 20:14:38 +00:00
fastResend: 2,
remoteNextNumber: 32,
2016-10-02 21:43:58 +00:00
controlWindow: kcp.Config.GetSendingInFlightSize(),
2016-07-03 20:14:38 +00:00
}
2016-10-02 21:43:58 +00:00
worker.window = NewSendingWindow(kcp.Config.GetSendingBufferSize(), worker, worker.OnPacketLoss)
2016-07-03 20:14:38 +00:00
return worker
}
2016-11-27 20:39:09 +00:00
func (v *SendingWorker) Release() {
2017-02-17 23:04:25 +00:00
v.Lock()
2016-11-27 20:39:09 +00:00
v.window.Release()
2017-02-17 23:04:25 +00:00
v.Unlock()
2016-11-21 21:41:12 +00:00
}
2016-11-27 20:39:09 +00:00
func (v *SendingWorker) ProcessReceivingNext(nextNumber uint32) {
v.Lock()
defer v.Unlock()
2016-07-03 20:14:38 +00:00
2016-11-27 20:39:09 +00:00
v.ProcessReceivingNextWithoutLock(nextNumber)
2016-07-06 14:36:15 +00:00
}
2016-11-27 20:39:09 +00:00
func (v *SendingWorker) ProcessReceivingNextWithoutLock(nextNumber uint32) {
v.window.Clear(nextNumber)
v.FindFirstUnacknowledged()
2016-07-03 20:14:38 +00:00
}
2016-11-27 20:39:09 +00:00
func (v *SendingWorker) FindFirstUnacknowledged() {
first := v.firstUnacknowledged
if !v.window.IsEmpty() {
v.firstUnacknowledged = v.window.FirstNumber()
2016-07-03 20:14:38 +00:00
} else {
2016-11-27 20:39:09 +00:00
v.firstUnacknowledged = v.nextNumber
2016-07-03 20:14:38 +00:00
}
2016-11-27 20:39:09 +00:00
if first != v.firstUnacknowledged {
v.firstUnacknowledgedUpdated = true
2016-10-11 10:24:19 +00:00
}
2016-07-03 20:14:38 +00:00
}
2017-02-26 14:01:50 +00:00
func (v *SendingWorker) processAck(number uint32) bool {
2016-11-27 20:39:09 +00:00
// number < v.firstUnacknowledged || number >= v.nextNumber
if number-v.firstUnacknowledged > 0x7FFFFFFF || number-v.nextNumber < 0x7FFFFFFF {
2016-11-13 21:27:58 +00:00
return false
2016-07-03 20:14:38 +00:00
}
2016-11-27 20:39:09 +00:00
removed := v.window.Remove(number - v.firstUnacknowledged)
2016-11-13 21:27:58 +00:00
if removed {
2016-11-27 20:39:09 +00:00
v.FindFirstUnacknowledged()
2016-11-13 21:27:58 +00:00
}
return removed
2016-07-03 20:14:38 +00:00
}
2016-11-27 20:39:09 +00:00
func (v *SendingWorker) ProcessSegment(current uint32, seg *AckSegment, rto uint32) {
2016-07-15 19:41:15 +00:00
defer seg.Release()
2016-11-27 20:39:09 +00:00
v.Lock()
defer v.Unlock()
2016-07-15 19:41:15 +00:00
2016-11-27 20:39:09 +00:00
if v.remoteNextNumber < seg.ReceivingWindow {
v.remoteNextNumber = seg.ReceivingWindow
2016-07-15 19:41:15 +00:00
}
2016-11-27 20:39:09 +00:00
v.ProcessReceivingNextWithoutLock(seg.ReceivingNext)
2016-07-14 15:38:20 +00:00
2016-12-21 14:37:16 +00:00
if seg.IsEmpty() {
2016-12-02 20:40:58 +00:00
return
}
2016-07-03 20:14:38 +00:00
var maxack uint32
2016-11-13 21:27:58 +00:00
var maxackRemoved bool
2016-12-21 14:37:16 +00:00
for _, number := range seg.NumberList {
2017-02-26 14:01:50 +00:00
removed := v.processAck(number)
2016-07-03 20:14:38 +00:00
if maxack < number {
maxack = number
2016-11-13 21:27:58 +00:00
maxackRemoved = removed
2016-07-03 20:14:38 +00:00
}
}
2016-07-06 14:36:15 +00:00
2016-11-13 21:27:58 +00:00
if maxackRemoved {
2016-11-27 20:39:09 +00:00
v.window.HandleFastAck(maxack, rto)
2016-11-13 21:27:58 +00:00
if current-seg.Timestamp < 10000 {
2016-11-27 20:39:09 +00:00
v.conn.roundTrip.Update(current-seg.Timestamp, current)
2016-11-13 21:27:58 +00:00
}
}
2016-07-03 20:14:38 +00:00
}
2017-12-17 00:22:39 +00:00
func (v *SendingWorker) Push(f buf.Supplier) bool {
2016-11-27 20:39:09 +00:00
v.Lock()
defer v.Unlock()
2016-08-25 09:41:05 +00:00
2017-12-05 17:04:34 +00:00
if v.window.IsFull() {
2017-12-17 00:22:39 +00:00
return false
2016-07-03 20:14:38 +00:00
}
2017-12-05 17:04:34 +00:00
b := v.window.Push(v.nextNumber)
v.nextNumber++
2017-12-17 00:22:39 +00:00
common.Must(b.Reset(f))
return true
2016-07-03 20:14:38 +00:00
}
2016-12-20 21:53:58 +00:00
func (v *SendingWorker) Write(seg Segment) error {
2016-07-03 20:14:38 +00:00
dataSeg := seg.(*DataSegment)
2017-12-14 22:24:40 +00:00
dataSeg.Conv = v.conn.meta.Conversation
2016-11-27 20:39:09 +00:00
dataSeg.SendingNext = v.firstUnacknowledged
2016-07-14 20:52:00 +00:00
dataSeg.Option = 0
2016-11-27 20:39:09 +00:00
if v.conn.State() == StateReadyToClose {
2016-07-14 20:52:00 +00:00
dataSeg.Option = SegmentOptionClose
2016-07-03 20:14:38 +00:00
}
2016-12-20 21:53:58 +00:00
return v.conn.output.Write(dataSeg)
2016-07-12 15:56:36 +00:00
}
2016-11-27 20:39:09 +00:00
func (v *SendingWorker) OnPacketLoss(lossRate uint32) {
if !v.conn.Config.Congestion || v.conn.roundTrip.Timeout() == 0 {
2016-07-03 20:14:38 +00:00
return
}
2016-07-04 13:34:14 +00:00
if lossRate >= 15 {
2016-11-27 20:39:09 +00:00
v.controlWindow = 3 * v.controlWindow / 4
2016-07-04 13:34:14 +00:00
} else if lossRate <= 5 {
2016-11-27 20:39:09 +00:00
v.controlWindow += v.controlWindow / 4
2016-07-03 20:14:38 +00:00
}
2016-11-27 20:39:09 +00:00
if v.controlWindow < 16 {
v.controlWindow = 16
2016-07-03 20:14:38 +00:00
}
2016-11-27 20:39:09 +00:00
if v.controlWindow > 2*v.conn.Config.GetSendingInFlightSize() {
v.controlWindow = 2 * v.conn.Config.GetSendingInFlightSize()
2016-07-03 20:14:38 +00:00
}
}
2016-11-27 20:39:09 +00:00
func (v *SendingWorker) Flush(current uint32) {
v.Lock()
2016-07-03 20:14:38 +00:00
2016-11-27 20:39:09 +00:00
cwnd := v.firstUnacknowledged + v.conn.Config.GetSendingInFlightSize()
if cwnd > v.remoteNextNumber {
cwnd = v.remoteNextNumber
2016-07-03 20:14:38 +00:00
}
2016-11-27 20:39:09 +00:00
if v.conn.Config.Congestion && cwnd > v.firstUnacknowledged+v.controlWindow {
cwnd = v.firstUnacknowledged + v.controlWindow
2016-07-03 20:14:38 +00:00
}
2016-11-27 20:39:09 +00:00
if !v.window.IsEmpty() {
v.window.Flush(current, v.conn.roundTrip.Timeout(), cwnd)
2017-02-17 23:04:25 +00:00
v.firstUnacknowledgedUpdated = false
2016-08-24 13:47:14 +00:00
}
2016-10-11 10:24:19 +00:00
2017-02-17 23:04:25 +00:00
updated := v.firstUnacknowledgedUpdated
2016-11-27 20:39:09 +00:00
v.firstUnacknowledgedUpdated = false
2017-02-17 23:04:25 +00:00
v.Unlock()
if updated {
v.conn.Ping(current, CommandPing)
}
2016-07-03 20:14:38 +00:00
}
2016-11-27 20:39:09 +00:00
func (v *SendingWorker) CloseWrite() {
v.Lock()
defer v.Unlock()
2016-07-03 20:14:38 +00:00
2016-11-27 20:39:09 +00:00
v.window.Clear(0xFFFFFFFF)
2016-07-03 20:14:38 +00:00
}
2016-07-12 21:54:54 +00:00
2016-11-27 20:39:09 +00:00
func (v *SendingWorker) IsEmpty() bool {
v.RLock()
defer v.RUnlock()
2016-07-12 21:54:54 +00:00
2016-11-27 20:39:09 +00:00
return v.window.IsEmpty()
2016-07-12 21:54:54 +00:00
}
2016-11-27 20:39:09 +00:00
func (v *SendingWorker) UpdateNecessary() bool {
return !v.IsEmpty()
}
2017-02-17 23:04:25 +00:00
func (w *SendingWorker) FirstUnacknowledged() uint32 {
w.RLock()
defer w.RUnlock()
return w.firstUnacknowledged
}