receiving.go 5.1 KB

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