kcp.go 11 KB


  1. // Package kcp - A Fast and Reliable ARQ Protocol
  2. //
  3. // Acknowledgement:
  4. // skywind3000@github for inventing the KCP protocol
  5. // xtaci@github for translating to Golang
  6. package kcp
  7. import (
  8. "github.com/v2ray/v2ray-core/common/alloc"
  9. v2io "github.com/v2ray/v2ray-core/common/io"
  10. "github.com/v2ray/v2ray-core/common/log"
  11. )
  12. const (
  13. IKCP_RTO_NDL = 30 // no delay min rto
  14. IKCP_RTO_MIN = 100 // normal min rto
  15. IKCP_RTO_DEF = 200
  16. IKCP_RTO_MAX = 60000
  17. IKCP_WND_SND = 32
  18. IKCP_WND_RCV = 32
  19. IKCP_INTERVAL = 100
  20. )
  21. func _itimediff(later, earlier uint32) int32 {
  22. return (int32)(later - earlier)
  23. }
  24. type State int
  25. const (
  26. StateActive State = 0
  27. StateReadyToClose State = 1
  28. StatePeerClosed State = 2
  29. StateTerminating State = 3
  30. StateTerminated State = 4
  31. )
  32. // KCP defines a single KCP connection
  33. type KCP struct {
  34. conv uint16
  35. state State
  36. stateBeginTime uint32
  37. lastIncomingTime uint32
  38. lastPayloadTime uint32
  39. sendingUpdated bool
  40. receivingUpdated bool
  41. lastPingTime uint32
  42. mtu, mss uint32
  43. snd_una, snd_nxt, rcv_nxt uint32
  44. rx_rttvar, rx_srtt, rx_rto uint32
  45. snd_wnd, rcv_wnd, rmt_wnd, cwnd uint32
  46. current, interval uint32
  47. snd_queue *SendingQueue
  48. rcv_queue *ReceivingQueue
  49. snd_buf []*DataSegment
  50. rcv_buf *ReceivingWindow
  51. acklist *ACKList
  52. fastresend int32
  53. congestionControl bool
  54. output *SegmentWriter
  55. }
  56. // NewKCP create a new kcp control object, 'conv' must equal in two endpoint
  57. // from the same connection.
  58. func NewKCP(conv uint16, mtu uint32, sendingWindowSize uint32, receivingWindowSize uint32, sendingQueueSize uint32, output v2io.Writer) *KCP {
  59. log.Debug("KCP|Core: creating KCP ", conv)
  60. kcp := new(KCP)
  61. kcp.conv = conv
  62. kcp.snd_wnd = sendingWindowSize
  63. kcp.rcv_wnd = receivingWindowSize
  64. kcp.rmt_wnd = IKCP_WND_RCV
  65. kcp.mtu = mtu
  66. kcp.mss = kcp.mtu - DataSegmentOverhead
  67. kcp.rx_rto = IKCP_RTO_DEF
  68. kcp.interval = IKCP_INTERVAL
  69. kcp.output = NewSegmentWriter(mtu, output)
  70. kcp.rcv_buf = NewReceivingWindow(receivingWindowSize)
  71. kcp.snd_queue = NewSendingQueue(sendingQueueSize)
  72. kcp.rcv_queue = NewReceivingQueue()
  73. kcp.acklist = new(ACKList)
  74. kcp.cwnd = kcp.snd_wnd
  75. return kcp
  76. }
  77. func (kcp *KCP) SetState(state State) {
  78. kcp.state = state
  79. kcp.stateBeginTime = kcp.current
  80. switch state {
  81. case StateReadyToClose:
  82. kcp.rcv_queue.Close()
  83. case StatePeerClosed:
  84. kcp.ClearSendQueue()
  85. case StateTerminating:
  86. kcp.rcv_queue.Close()
  87. case StateTerminated:
  88. kcp.rcv_queue.Close()
  89. }
  90. }
  91. func (kcp *KCP) HandleOption(opt SegmentOption) {
  92. if (opt & SegmentOptionClose) == SegmentOptionClose {
  93. kcp.OnPeerClosed()
  94. }
  95. }
  96. func (kcp *KCP) OnPeerClosed() {
  97. if kcp.state == StateReadyToClose {
  98. kcp.SetState(StateTerminating)
  99. }
  100. if kcp.state == StateActive {
  101. kcp.SetState(StatePeerClosed)
  102. }
  103. }
  104. func (kcp *KCP) OnClose() {
  105. if kcp.state == StateActive {
  106. kcp.SetState(StateReadyToClose)
  107. }
  108. if kcp.state == StatePeerClosed {
  109. kcp.SetState(StateTerminating)
  110. }
  111. }
  112. // DumpReceivingBuf moves available data from rcv_buf -> rcv_queue
  113. // @Private
  114. func (kcp *KCP) DumpReceivingBuf() {
  115. for {
  116. seg := kcp.rcv_buf.RemoveFirst()
  117. if seg == nil {
  118. break
  119. }
  120. kcp.rcv_queue.Put(seg.Data)
  121. seg.Data = nil
  122. kcp.rcv_buf.Advance()
  123. kcp.rcv_nxt++
  124. kcp.receivingUpdated = true
  125. }
  126. }
  127. // Send is user/upper level send, returns below zero for error
  128. func (kcp *KCP) Send(buffer []byte) int {
  129. nBytes := 0
  130. for len(buffer) > 0 && !kcp.snd_queue.IsFull() {
  131. var size int
  132. if len(buffer) > int(kcp.mss) {
  133. size = int(kcp.mss)
  134. } else {
  135. size = len(buffer)
  136. }
  137. seg := &DataSegment{
  138. Data: alloc.NewSmallBuffer().Clear().Append(buffer[:size]),
  139. }
  140. kcp.snd_queue.Push(seg)
  141. buffer = buffer[size:]
  142. nBytes += size
  143. }
  144. return nBytes
  145. }
  146. // https://tools.ietf.org/html/rfc6298
  147. func (kcp *KCP) update_ack(rtt int32) {
  148. if kcp.rx_srtt == 0 {
  149. kcp.rx_srtt = uint32(rtt)
  150. kcp.rx_rttvar = uint32(rtt) / 2
  151. } else {
  152. delta := rtt - int32(kcp.rx_srtt)
  153. if delta < 0 {
  154. delta = -delta
  155. }
  156. kcp.rx_rttvar = (3*kcp.rx_rttvar + uint32(delta)) / 4
  157. kcp.rx_srtt = (7*kcp.rx_srtt + uint32(rtt)) / 8
  158. if kcp.rx_srtt < kcp.interval {
  159. kcp.rx_srtt = kcp.interval
  160. }
  161. }
  162. var rto uint32
  163. if kcp.interval < 4*kcp.rx_rttvar {
  164. rto = kcp.rx_srtt + 4*kcp.rx_rttvar
  165. } else {
  166. rto = kcp.rx_srtt + kcp.interval
  167. }
  168. if rto > IKCP_RTO_MAX {
  169. rto = IKCP_RTO_MAX
  170. }
  171. kcp.rx_rto = rto * 3 / 2
  172. }
  173. func (kcp *KCP) shrink_buf() {
  174. prevUna := kcp.snd_una
  175. if len(kcp.snd_buf) > 0 {
  176. seg := kcp.snd_buf[0]
  177. kcp.snd_una = seg.Number
  178. } else {
  179. kcp.snd_una = kcp.snd_nxt
  180. }
  181. if kcp.snd_una != prevUna {
  182. kcp.sendingUpdated = true
  183. }
  184. }
  185. func (kcp *KCP) parse_ack(sn uint32) {
  186. if _itimediff(sn, kcp.snd_una) < 0 || _itimediff(sn, kcp.snd_nxt) >= 0 {
  187. return
  188. }
  189. for k, seg := range kcp.snd_buf {
  190. if sn == seg.Number {
  191. kcp.snd_buf = append(kcp.snd_buf[:k], kcp.snd_buf[k+1:]...)
  192. seg.Release()
  193. break
  194. }
  195. if _itimediff(sn, seg.Number) < 0 {
  196. break
  197. }
  198. }
  199. }
  200. func (kcp *KCP) parse_fastack(sn uint32) {
  201. if _itimediff(sn, kcp.snd_una) < 0 || _itimediff(sn, kcp.snd_nxt) >= 0 {
  202. return
  203. }
  204. for _, seg := range kcp.snd_buf {
  205. if _itimediff(sn, seg.Number) < 0 {
  206. break
  207. } else if sn != seg.Number {
  208. seg.ackSkipped++
  209. }
  210. }
  211. }
  212. func (kcp *KCP) HandleReceivingNext(receivingNext uint32) {
  213. count := 0
  214. for _, seg := range kcp.snd_buf {
  215. if _itimediff(receivingNext, seg.Number) > 0 {
  216. seg.Release()
  217. count++
  218. } else {
  219. break
  220. }
  221. }
  222. kcp.snd_buf = kcp.snd_buf[count:]
  223. }
  224. func (kcp *KCP) HandleSendingNext(sendingNext uint32) {
  225. if kcp.acklist.Clear(sendingNext) {
  226. kcp.receivingUpdated = true
  227. }
  228. }
  229. func (kcp *KCP) parse_data(newseg *DataSegment) {
  230. sn := newseg.Number
  231. if _itimediff(sn, kcp.rcv_nxt+kcp.rcv_wnd) >= 0 ||
  232. _itimediff(sn, kcp.rcv_nxt) < 0 {
  233. return
  234. }
  235. idx := sn - kcp.rcv_nxt
  236. if !kcp.rcv_buf.Set(idx, newseg) {
  237. newseg.Release()
  238. }
  239. kcp.DumpReceivingBuf()
  240. }
  241. // Input when you received a low level packet (eg. UDP packet), call it
  242. func (kcp *KCP) Input(data []byte) int {
  243. kcp.lastIncomingTime = kcp.current
  244. var seg ISegment
  245. var maxack uint32
  246. var flag int
  247. for {
  248. seg, data = ReadSegment(data)
  249. if seg == nil {
  250. break
  251. }
  252. switch seg := seg.(type) {
  253. case *DataSegment:
  254. kcp.HandleOption(seg.Opt)
  255. kcp.HandleSendingNext(seg.SendingNext)
  256. kcp.acklist.Add(seg.Number, seg.Timestamp)
  257. kcp.receivingUpdated = true
  258. kcp.parse_data(seg)
  259. kcp.lastPayloadTime = kcp.current
  260. case *ACKSegment:
  261. kcp.HandleOption(seg.Opt)
  262. if kcp.rmt_wnd < seg.ReceivingWindow {
  263. kcp.rmt_wnd = seg.ReceivingWindow
  264. }
  265. kcp.HandleReceivingNext(seg.ReceivingNext)
  266. for i := 0; i < int(seg.Count); i++ {
  267. ts := seg.TimestampList[i]
  268. sn := seg.NumberList[i]
  269. if _itimediff(kcp.current, ts) >= 0 {
  270. kcp.update_ack(_itimediff(kcp.current, ts))
  271. }
  272. kcp.parse_ack(sn)
  273. if flag == 0 {
  274. flag = 1
  275. maxack = sn
  276. } else if _itimediff(sn, maxack) > 0 {
  277. maxack = sn
  278. }
  279. }
  280. kcp.lastPayloadTime = kcp.current
  281. case *CmdOnlySegment:
  282. kcp.HandleOption(seg.Opt)
  283. if seg.Cmd == SegmentCommandTerminated {
  284. if kcp.state == StateActive ||
  285. kcp.state == StateReadyToClose ||
  286. kcp.state == StatePeerClosed {
  287. kcp.SetState(StateTerminating)
  288. } else if kcp.state == StateTerminating {
  289. kcp.SetState(StateTerminated)
  290. }
  291. }
  292. kcp.HandleReceivingNext(seg.ReceivinNext)
  293. kcp.HandleSendingNext(seg.SendingNext)
  294. default:
  295. }
  296. kcp.shrink_buf()
  297. }
  298. if flag != 0 {
  299. kcp.parse_fastack(maxack)
  300. }
  301. return 0
  302. }
  303. // flush pending data
  304. func (kcp *KCP) flush() {
  305. if kcp.state == StateTerminated {
  306. return
  307. }
  308. if kcp.state == StateActive && _itimediff(kcp.current, kcp.lastPayloadTime) >= 30000 {
  309. kcp.OnClose()
  310. }
  311. if kcp.state == StateTerminating {
  312. kcp.output.Write(&CmdOnlySegment{
  313. Conv: kcp.conv,
  314. Cmd: SegmentCommandTerminated,
  315. })
  316. kcp.output.Flush()
  317. if _itimediff(kcp.current, kcp.stateBeginTime) > 8000 {
  318. kcp.SetState(StateTerminated)
  319. }
  320. return
  321. }
  322. if kcp.state == StateReadyToClose && _itimediff(kcp.current, kcp.stateBeginTime) > 15000 {
  323. kcp.SetState(StateTerminating)
  324. }
  325. current := kcp.current
  326. lost := false
  327. // flush acknowledges
  328. //if kcp.receivingUpdated {
  329. ackSeg := kcp.acklist.AsSegment()
  330. if ackSeg != nil {
  331. ackSeg.Conv = kcp.conv
  332. ackSeg.ReceivingWindow = uint32(kcp.rcv_nxt + kcp.rcv_wnd)
  333. ackSeg.ReceivingNext = kcp.rcv_nxt
  334. kcp.output.Write(ackSeg)
  335. kcp.receivingUpdated = false
  336. }
  337. //}
  338. // calculate window size
  339. cwnd := kcp.snd_una + kcp.snd_wnd
  340. if cwnd < kcp.rmt_wnd {
  341. cwnd = kcp.rmt_wnd
  342. }
  343. if kcp.congestionControl && cwnd < kcp.snd_una+kcp.cwnd {
  344. cwnd = kcp.snd_una + kcp.cwnd
  345. }
  346. for !kcp.snd_queue.IsEmpty() && _itimediff(kcp.snd_nxt, cwnd) < 0 {
  347. seg := kcp.snd_queue.Pop()
  348. seg.Conv = kcp.conv
  349. seg.Number = kcp.snd_nxt
  350. seg.timeout = current
  351. seg.ackSkipped = 0
  352. seg.transmit = 0
  353. kcp.snd_buf = append(kcp.snd_buf, seg)
  354. kcp.snd_nxt++
  355. }
  356. // calculate resent
  357. resent := uint32(kcp.fastresend)
  358. if kcp.fastresend <= 0 {
  359. resent = 0xffffffff
  360. }
  361. // flush data segments
  362. for _, segment := range kcp.snd_buf {
  363. needsend := false
  364. if segment.transmit == 0 {
  365. needsend = true
  366. segment.transmit++
  367. segment.timeout = current + kcp.rx_rto
  368. } else if _itimediff(current, segment.timeout) >= 0 {
  369. needsend = true
  370. segment.transmit++
  371. segment.timeout = current + kcp.rx_rto
  372. lost = true
  373. } else if segment.ackSkipped >= resent {
  374. needsend = true
  375. segment.transmit++
  376. segment.ackSkipped = 0
  377. segment.timeout = current + kcp.rx_rto
  378. lost = true
  379. }
  380. if needsend {
  381. segment.Timestamp = current
  382. segment.SendingNext = kcp.snd_una
  383. segment.Opt = 0
  384. if kcp.state == StateReadyToClose {
  385. segment.Opt = SegmentOptionClose
  386. }
  387. kcp.output.Write(segment)
  388. kcp.sendingUpdated = false
  389. }
  390. }
  391. if kcp.sendingUpdated || kcp.receivingUpdated || _itimediff(kcp.current, kcp.lastPingTime) >= 5000 {
  392. seg := &CmdOnlySegment{
  393. Conv: kcp.conv,
  394. Cmd: SegmentCommandPing,
  395. ReceivinNext: kcp.rcv_nxt,
  396. SendingNext: kcp.snd_una,
  397. }
  398. if kcp.state == StateReadyToClose {
  399. seg.Opt = SegmentOptionClose
  400. }
  401. kcp.output.Write(seg)
  402. kcp.lastPingTime = kcp.current
  403. kcp.sendingUpdated = false
  404. kcp.receivingUpdated = false
  405. }
  406. // flash remain segments
  407. kcp.output.Flush()
  408. if kcp.congestionControl {
  409. if lost {
  410. kcp.cwnd = 3 * kcp.cwnd / 4
  411. } else {
  412. kcp.cwnd += kcp.cwnd / 4
  413. }
  414. if kcp.cwnd < 4 {
  415. kcp.cwnd = 4
  416. }
  417. if kcp.cwnd > kcp.snd_wnd {
  418. kcp.cwnd = kcp.snd_wnd
  419. }
  420. }
  421. }
  422. // Update updates state (call it repeatedly, every 10ms-100ms), or you can ask
  423. // ikcp_check when to call it again (without ikcp_input/_send calling).
  424. // 'current' - current timestamp in millisec.
  425. func (kcp *KCP) Update(current uint32) {
  426. kcp.current = current
  427. kcp.flush()
  428. }
  429. // NoDelay options
  430. // fastest: ikcp_nodelay(kcp, 1, 20, 2, 1)
  431. // nodelay: 0:disable(default), 1:enable
  432. // interval: internal update timer interval in millisec, default is 100ms
  433. // resend: 0:disable fast resend(default), 1:enable fast resend
  434. // nc: 0:normal congestion control(default), 1:disable congestion control
  435. func (kcp *KCP) NoDelay(interval uint32, resend int, congestionControl bool) int {
  436. kcp.interval = interval
  437. if resend >= 0 {
  438. kcp.fastresend = int32(resend)
  439. }
  440. kcp.congestionControl = congestionControl
  441. return 0
  442. }
  443. // WaitSnd gets how many packet is waiting to be sent
  444. func (kcp *KCP) WaitSnd() uint32 {
  445. return uint32(len(kcp.snd_buf)) + kcp.snd_queue.Len()
  446. }
  447. func (this *KCP) ClearSendQueue() {
  448. this.snd_queue.Clear()
  449. for _, seg := range this.snd_buf {
  450. seg.Release()
  451. }
  452. this.snd_buf = nil
  453. }