receiving.go 5.4 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287
  1. package kcp
  2. import (
  3. "sync"
  4. "v2ray.com/core/common/buf"
  5. )
  6. type ReceivingWindow struct {
  7. start uint32
  8. size uint32
  9. list []*DataSegment
  10. }
  11. func NewReceivingWindow(size uint32) *ReceivingWindow {
  12. return &ReceivingWindow{
  13. start: 0,
  14. size: size,
  15. list: make([]*DataSegment, size),
  16. }
  17. }
  18. func (w *ReceivingWindow) Size() uint32 {
  19. return w.size
  20. }
  21. func (w *ReceivingWindow) Position(idx uint32) (uint32, bool) {
  22. if idx >= w.size {
  23. return 0, false
  24. }
  25. return (w.start + idx) % w.size, true
  26. }
  27. func (w *ReceivingWindow) Set(idx uint32, value *DataSegment) bool {
  28. pos, ok := w.Position(idx)
  29. if !ok {
  30. return false
  31. }
  32. if w.list[pos] != nil {
  33. return false
  34. }
  35. w.list[pos] = value
  36. return true
  37. }
  38. func (w *ReceivingWindow) Remove(idx uint32) *DataSegment {
  39. pos, ok := w.Position(idx)
  40. if !ok {
  41. return nil
  42. }
  43. e := w.list[pos]
  44. w.list[pos] = nil
  45. return e
  46. }
  47. func (w *ReceivingWindow) RemoveFirst() *DataSegment {
  48. return w.Remove(0)
  49. }
  50. func (w *ReceivingWindow) HasFirst() bool {
  51. pos, _ := w.Position(0)
  52. return w.list[pos] != nil
  53. }
  54. func (w *ReceivingWindow) Advance() {
  55. w.start++
  56. if w.start == w.size {
  57. w.start = 0
  58. }
  59. }
  60. type AckList struct {
  61. writer SegmentWriter
  62. timestamps []uint32
  63. numbers []uint32
  64. nextFlush []uint32
  65. flushCandidates []uint32
  66. dirty bool
  67. }
  68. func NewAckList(writer SegmentWriter) *AckList {
  69. return &AckList{
  70. writer: writer,
  71. timestamps: make([]uint32, 0, 128),
  72. numbers: make([]uint32, 0, 128),
  73. nextFlush: make([]uint32, 0, 128),
  74. flushCandidates: make([]uint32, 0, 128),
  75. }
  76. }
  77. func (l *AckList) Add(number uint32, timestamp uint32) {
  78. l.timestamps = append(l.timestamps, timestamp)
  79. l.numbers = append(l.numbers, number)
  80. l.nextFlush = append(l.nextFlush, 0)
  81. l.dirty = true
  82. }
  83. func (l *AckList) Clear(una uint32) {
  84. count := 0
  85. for i := 0; i < len(l.numbers); i++ {
  86. if l.numbers[i] < una {
  87. continue
  88. }
  89. if i != count {
  90. l.numbers[count] = l.numbers[i]
  91. l.timestamps[count] = l.timestamps[i]
  92. l.nextFlush[count] = l.nextFlush[i]
  93. }
  94. count++
  95. }
  96. if count < len(l.numbers) {
  97. l.numbers = l.numbers[:count]
  98. l.timestamps = l.timestamps[:count]
  99. l.nextFlush = l.nextFlush[:count]
  100. l.dirty = true
  101. }
  102. }
  103. func (l *AckList) Flush(current uint32, rto uint32) {
  104. l.flushCandidates = l.flushCandidates[:0]
  105. seg := NewAckSegment()
  106. for i := 0; i < len(l.numbers); i++ {
  107. if l.nextFlush[i] > current {
  108. if len(l.flushCandidates) < cap(l.flushCandidates) {
  109. l.flushCandidates = append(l.flushCandidates, l.numbers[i])
  110. }
  111. continue
  112. }
  113. seg.PutNumber(l.numbers[i])
  114. seg.PutTimestamp(l.timestamps[i])
  115. timeout := rto / 2
  116. if timeout < 20 {
  117. timeout = 20
  118. }
  119. l.nextFlush[i] = current + timeout
  120. if seg.IsFull() {
  121. l.writer.Write(seg)
  122. seg.Release()
  123. seg = NewAckSegment()
  124. l.dirty = false
  125. }
  126. }
  127. if l.dirty || !seg.IsEmpty() {
  128. for _, number := range l.flushCandidates {
  129. if seg.IsFull() {
  130. break
  131. }
  132. seg.PutNumber(number)
  133. }
  134. l.writer.Write(seg)
  135. seg.Release()
  136. l.dirty = false
  137. }
  138. }
  139. type ReceivingWorker struct {
  140. sync.RWMutex
  141. conn *Connection
  142. leftOver buf.MultiBuffer
  143. window *ReceivingWindow
  144. acklist *AckList
  145. nextNumber uint32
  146. windowSize uint32
  147. }
  148. func NewReceivingWorker(kcp *Connection) *ReceivingWorker {
  149. windowsSize := kcp.Config.GetReceivingInFlightSize()
  150. if windowsSize > kcp.Config.GetReceivingBufferSize() {
  151. windowsSize = kcp.Config.GetReceivingBufferSize()
  152. }
  153. worker := &ReceivingWorker{
  154. conn: kcp,
  155. window: NewReceivingWindow(kcp.Config.GetReceivingBufferSize()),
  156. windowSize: windowsSize,
  157. }
  158. worker.acklist = NewAckList(worker)
  159. return worker
  160. }
  161. func (w *ReceivingWorker) Release() {
  162. w.Lock()
  163. w.leftOver.Release()
  164. w.Unlock()
  165. }
  166. func (w *ReceivingWorker) ProcessSendingNext(number uint32) {
  167. w.Lock()
  168. defer w.Unlock()
  169. w.acklist.Clear(number)
  170. }
  171. func (w *ReceivingWorker) ProcessSegment(seg *DataSegment) {
  172. w.Lock()
  173. defer w.Unlock()
  174. number := seg.Number
  175. idx := number - w.nextNumber
  176. if idx >= w.windowSize {
  177. return
  178. }
  179. w.acklist.Clear(seg.SendingNext)
  180. w.acklist.Add(number, seg.Timestamp)
  181. if !w.window.Set(idx, seg) {
  182. seg.Release()
  183. }
  184. }
  185. func (w *ReceivingWorker) ReadMultiBuffer() buf.MultiBuffer {
  186. if w.leftOver != nil {
  187. mb := w.leftOver
  188. w.leftOver = nil
  189. return mb
  190. }
  191. mb := buf.NewMultiBufferCap(32)
  192. w.Lock()
  193. defer w.Unlock()
  194. for {
  195. seg := w.window.RemoveFirst()
  196. if seg == nil {
  197. break
  198. }
  199. w.window.Advance()
  200. w.nextNumber++
  201. mb.Append(seg.Detach())
  202. seg.Release()
  203. }
  204. return mb
  205. }
  206. func (w *ReceivingWorker) Read(b []byte) int {
  207. mb := w.ReadMultiBuffer()
  208. nBytes, _ := mb.Read(b)
  209. if !mb.IsEmpty() {
  210. w.leftOver = mb
  211. }
  212. return nBytes
  213. }
  214. func (w *ReceivingWorker) IsDataAvailable() bool {
  215. w.RLock()
  216. defer w.RUnlock()
  217. return w.window.HasFirst()
  218. }
  219. func (w *ReceivingWorker) NextNumber() uint32 {
  220. w.RLock()
  221. defer w.RUnlock()
  222. return w.nextNumber
  223. }
  224. func (w *ReceivingWorker) Flush(current uint32) {
  225. w.Lock()
  226. defer w.Unlock()
  227. w.acklist.Flush(current, w.conn.roundTrip.Timeout())
  228. }
  229. func (w *ReceivingWorker) Write(seg Segment) error {
  230. ackSeg := seg.(*AckSegment)
  231. ackSeg.Conv = w.conn.meta.Conversation
  232. ackSeg.ReceivingNext = w.nextNumber
  233. ackSeg.ReceivingWindow = w.nextNumber + w.windowSize
  234. if w.conn.State() == StateReadyToClose {
  235. ackSeg.Option = SegmentOptionClose
  236. }
  237. return w.conn.output.Write(ackSeg)
  238. }
  239. func (*ReceivingWorker) CloseRead() {
  240. }
  241. func (w *ReceivingWorker) UpdateNecessary() bool {
  242. w.RLock()
  243. defer w.RUnlock()
  244. return len(w.acklist.numbers) > 0
  245. }