1
0
mirror of https://github.com/v2fly/v2ray-core.git synced 2025-01-21 16:56:27 -05:00
v2fly/transport/internet/kcp/sending.go

383 lines
7.7 KiB
Go
Raw Normal View History

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