receiving.go 4.9 KB

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