KcpPipe.go 33 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273127412751276127712781279128012811282128312841285
  1. package Common
  2. import (
  3. "encoding/binary"
  4. "errors"
  5. "github.com/astaxie/beego/logs"
  6. "github.com/vzex/zappy"
  7. "math/rand"
  8. "net"
  9. "sync"
  10. "forward-core/ikcp"
  11. "github.com/klauspost/reedsolomon"
  12. "time"
  13. )
  14. const WriteBufferSize = 5000 //udp writer will add some data for checksum or encrypt
  15. const ReadBufferSize = WriteBufferSize * 2 //so reader must be larger
  16. const dataLimit = WriteBufferSize
  17. const mainV = 0
  18. const subV = 1
  19. func init() {
  20. }
  21. func NewKCP(addr string, setting *KcpSetting) (*UDPListener, error) {
  22. udpAddr, err := net.ResolveUDPAddr("udp", addr)
  23. if err != nil {
  24. return nil, err
  25. }
  26. udpSocket, _err := net.ListenUDP("udp", udpAddr)
  27. if _err != nil {
  28. return nil, _err
  29. }
  30. listener := &UDPListener{connChan: make(chan *UDPMakeSession), quitChan: make(chan struct{}), sock: udpSocket, readBuffer: make([]byte, ReadBufferSize*2), sessions: make(map[string]*UDPMakeSession), setting: setting}
  31. go listener.loop()
  32. return listener, nil
  33. }
  34. func DefaultKcpSetting() *KcpSetting {
  35. return &KcpSetting{Nodelay: 1, Interval: 10, Resend: 2, Nc: 1, Sndwnd: 1024, Rcvwnd: 1024, Mtu: 1400}
  36. }
  37. func DialKcpTimeout(addr string, timeout int) (*UDPMakeSession, error) {
  38. return DialKcpTimeoutWithSetting(addr, timeout, DefaultKcpSetting(), 0, 0, false, false)
  39. }
  40. func DialKcpTimeoutWithSetting(addr string, timeout int, setting *KcpSetting, ds, ps int, comp, confuse bool) (*UDPMakeSession, error) {
  41. bReset := false
  42. if timeout < 5 {
  43. bReset = true
  44. timeout = 5
  45. } else if timeout > 255 {
  46. bReset = true
  47. timeout = 255
  48. }
  49. if bReset {
  50. logs.Debug("timeout should in [5, 255], force reset timeout to", timeout)
  51. }
  52. udpAddr, err := net.ResolveUDPAddr("udp", addr)
  53. if err != nil {
  54. return nil, err
  55. }
  56. sock, _err := net.ListenUDP("udp", &net.UDPAddr{})
  57. if _err != nil {
  58. logs.Debug("dial addr fail", _err.Error())
  59. return nil, _err
  60. }
  61. session := &UDPMakeSession{readBuffer: make([]byte, ReadBufferSize*2), do: make(chan Action), do2: make(chan Action), quitChan: make(chan struct{}), recvChan: make(chan KcpCache), processBuffer: make([]byte, ReadBufferSize), closeChan: make(chan struct{}), encodeBuffer: make([]byte, 7), checkCanWrite: make(chan (chan struct{})), xor: setting.Xor}
  62. session.remote = udpAddr
  63. session.sock = sock
  64. session.status = "firstsyn"
  65. session.timeout = int64(timeout)
  66. if confuse {
  67. session.confuseSeed = rand.Intn(int(WriteBufferSize/2)) + 10
  68. }
  69. if ds != 0 && ps != 0 {
  70. er := session.SetFec(ds, ps)
  71. if er != nil {
  72. logs.Debug("set fec error:", er.Error())
  73. } else {
  74. logs.Debug("set fec ok:", ds, ps)
  75. }
  76. }
  77. if comp {
  78. session.compressCache = make([]byte, zappy.MaxEncodedLen(ReadBufferSize))
  79. }
  80. _timeout := int(timeout / 2)
  81. if _timeout < 5 {
  82. timeout = 5
  83. }
  84. arg := int(int32(timeout) + int32(mainV<<24) + int32(subV<<16))
  85. arg2 := int16((ds << 8) | ps)
  86. if comp {
  87. arg2 = arg2 | (1 << 7)
  88. } else {
  89. arg2 = arg2 & ^(1 << 7)
  90. }
  91. info := makeEncode(session.encodeBuffer, FirstSYN, arg, arg2, session.xor)
  92. code := session.doAndWait(func() {
  93. sock.WriteToUDP(info, udpAddr)
  94. }, _timeout, func(status byte, arg int32, arg2 int16) int {
  95. if status == ResetAck {
  96. _mainV, _subV := int(byte(arg>>24)), int(byte(arg>>16))
  97. logs.Debug("pipe version not eq,%d.%d=>%d.%d", mainV, subV, _mainV, _subV)
  98. return 1
  99. }
  100. if status != FirstACK {
  101. logs.Debug("recv not FirstACK", status)
  102. return -1
  103. } else {
  104. session.status = "firstack"
  105. session.id = int(arg)
  106. return 0
  107. }
  108. })
  109. if code != 0 {
  110. logs.Debug("handshakefail with code", code)
  111. return nil, errors.New("handshake fail,1")
  112. }
  113. code = session.doAndWait(func() {
  114. if session.confuseSeed > 0 {
  115. sock.WriteToUDP(makeEncode(session.encodeBuffer, SndSYN, session.id, 1, session.xor), udpAddr)
  116. } else {
  117. sock.WriteToUDP(makeEncode(session.encodeBuffer, SndSYN, session.id, 0, session.xor), udpAddr)
  118. }
  119. }, _timeout, func(status byte, arg int32, arg2 int16) int {
  120. if status != SndACK {
  121. return -1
  122. } else if session.id != int(arg) {
  123. return 2
  124. } else {
  125. session.status = "ok"
  126. return 0
  127. }
  128. })
  129. if code != 0 {
  130. return nil, errors.New("handshake fail,2")
  131. }
  132. session.kcp = ikcp.Ikcp_create(uint32(session.id), session)
  133. session.kcp.Output = udp_output
  134. ikcp.Ikcp_wndsize(session.kcp, setting.Sndwnd, setting.Rcvwnd)
  135. ikcp.Ikcp_nodelay(session.kcp, setting.Nodelay, setting.Interval, setting.Resend, setting.Nc)
  136. ikcp.Ikcp_setmtu(session.kcp, setting.Mtu)
  137. go session.loop()
  138. return session, nil
  139. }
  140. type UDPListener struct {
  141. connChan chan *UDPMakeSession
  142. quitChan chan struct{}
  143. sock *net.UDPConn
  144. readBuffer []byte
  145. sessions map[string]*UDPMakeSession
  146. sessionsLock sync.RWMutex
  147. setting *KcpSetting
  148. }
  149. type UDPMakeSession struct {
  150. id int
  151. idstr string
  152. status string
  153. overTime int64
  154. quitChan chan struct{}
  155. closeChan chan struct{}
  156. recvChan chan KcpCache
  157. handShakeChan chan string
  158. handShakeChanQuit chan struct{}
  159. sock *net.UDPConn
  160. remote *net.UDPAddr
  161. kcp *ikcp.Ikcpcb
  162. do chan Action
  163. do2 chan Action
  164. checkCanWrite chan (chan struct{})
  165. listener *UDPListener
  166. closed bool
  167. readBuffer []byte
  168. processBuffer []byte
  169. encodeBuffer []byte
  170. timeout int64
  171. xor string
  172. compressCache []byte
  173. fecDataShards int
  174. fecParityShards int
  175. fecW *reedsolomon.Encoder
  176. fecR *reedsolomon.Encoder
  177. fecRCacheTbl map[uint]*fecInfo
  178. fecWCacheTbl *fecInfo
  179. fecWriteId uint //uint16
  180. fecSendC uint
  181. fecSendL int
  182. fecRecvId uint
  183. confuseSeed int
  184. }
  185. type KcpSetting struct {
  186. Nodelay int32
  187. Interval int32 //not for set
  188. Resend int32
  189. Nc int32
  190. Sndwnd int32
  191. Rcvwnd int32
  192. Mtu int32
  193. Xor string
  194. }
  195. type KcpCache struct {
  196. b []byte
  197. l int
  198. c chan int
  199. }
  200. type Action struct {
  201. t string
  202. args []interface{}
  203. }
  204. type fecInfo struct {
  205. bytes [][]byte
  206. overTime int64
  207. }
  208. func iclock() int32 {
  209. return int32((time.Now().UnixNano() / 1000000) & 0xffffffff)
  210. }
  211. func udp_output(buf []byte, _len int32, kcp *ikcp.Ikcpcb, user interface{}) int32 {
  212. c := user.(*UDPMakeSession)
  213. //logs.Debug("send udp", _len, c.remote.String())
  214. c.output(buf[:_len])
  215. return 0
  216. }
  217. func (l *UDPListener) Accept() (net.Conn, error) {
  218. var c *UDPMakeSession
  219. var err error
  220. select {
  221. case c = <-l.connChan:
  222. case <-l.quitChan:
  223. }
  224. if c == nil {
  225. err = errors.New("listener quit")
  226. }
  227. return net.Conn(c), err
  228. }
  229. func (l *UDPListener) Dump() {
  230. l.sessionsLock.RLock()
  231. defer l.sessionsLock.RUnlock()
  232. for addr, session := range l.sessions {
  233. logs.Debug("listener", addr, session.status)
  234. }
  235. }
  236. func (l *UDPListener) inner_loop() {
  237. sock := l.sock
  238. for {
  239. n, from, err := sock.ReadFromUDP(l.readBuffer)
  240. if err == nil {
  241. //logs.Debug("recv", n, from)
  242. addr := from.String()
  243. l.sessionsLock.RLock()
  244. session, bHave := l.sessions[addr]
  245. l.sessionsLock.RUnlock()
  246. if bHave {
  247. if session.status == "ok" {
  248. if session.remote.String() == from.String() && (n >= int(ikcp.IKCP_OVERHEAD) || session.compressCache != nil) {
  249. __xor(l.readBuffer, session.xor)
  250. var buf []byte
  251. if n <= 7 || session.compressCache == nil {
  252. buf = make([]byte, n)
  253. copy(buf, l.readBuffer[:n])
  254. } else {
  255. _b, _er := zappy.Decode(nil, l.readBuffer[:n])
  256. if _er != nil {
  257. logs.Debug("decompress fail", _er.Error())
  258. //go session.Close()
  259. //don't close pipe, just drop this data
  260. continue
  261. }
  262. buf = _b
  263. //logs.Debug("decompress", n, len(_b))
  264. }
  265. session.DoAction2("input", buf, len(buf))
  266. }
  267. continue
  268. } else {
  269. session.serverDo(string(l.readBuffer[:n]))
  270. }
  271. } else {
  272. status, _, fec := makeDecode(l.readBuffer[:n], l.setting.Xor)
  273. if status != FirstSYN {
  274. go sock.WriteToUDP([]byte("0"), from)
  275. logs.Debug("invalid package,reset", from, status)
  276. continue
  277. }
  278. sessionId := GetId("udp")
  279. session = &UDPMakeSession{status: "init", overTime: time.Now().Unix() + 10, remote: from, sock: sock, recvChan: make(chan KcpCache), quitChan: make(chan struct{}), readBuffer: make([]byte, ReadBufferSize*2), processBuffer: make([]byte, ReadBufferSize), timeout: 30, do: make(chan Action), do2: make(chan Action), id: sessionId, handShakeChan: make(chan string), handShakeChanQuit: make(chan struct{}), listener: l, closeChan: make(chan struct{}), encodeBuffer: make([]byte, 7), checkCanWrite: make(chan (chan struct{})), xor: l.setting.Xor}
  280. if fec & ^(1<<7) != 0 {
  281. er := session.SetFec((int(fec)>>8)&0xff, int(fec&^(1<<7)&0xff))
  282. if er != nil {
  283. logs.Debug("set fec error:", fec, er.Error())
  284. } else {
  285. logs.Debug("set fec ok:", (int(fec)>>8)&0xff, int(fec&^(1<<7)&0xff))
  286. }
  287. }
  288. if int(fec)&(1<<7) != 0 {
  289. session.compressCache = make([]byte, zappy.MaxEncodedLen(ReadBufferSize))
  290. }
  291. l.sessionsLock.Lock()
  292. l.sessions[addr] = session
  293. l.sessionsLock.Unlock()
  294. session.serverInit(l, l.setting)
  295. session.serverDo(string(l.readBuffer[:n]))
  296. }
  297. //logs.Debug("debug out.........")
  298. } else {
  299. e, ok := err.(net.Error)
  300. if !ok || !e.Timeout() {
  301. logs.Debug("recv error", err.Error(), from)
  302. l.remove(from.String())
  303. //time.Sleep(time.Second)
  304. break
  305. }
  306. }
  307. }
  308. }
  309. func (l *UDPListener) remove(addr string) {
  310. logs.Debug("listener remove", addr)
  311. l.sessionsLock.Lock()
  312. defer l.sessionsLock.Unlock()
  313. session, bHave := l.sessions[addr]
  314. if bHave {
  315. RmId("udp", session.id)
  316. }
  317. delete(l.sessions, addr)
  318. }
  319. func (l *UDPListener) Close() error {
  320. if l.sock != nil {
  321. l.sock.Close()
  322. l.sock = nil
  323. } else {
  324. return nil
  325. }
  326. close(l.quitChan)
  327. return nil
  328. }
  329. func (l *UDPListener) Addr() net.Addr {
  330. return nil
  331. }
  332. func (l *UDPListener) loop() {
  333. l.inner_loop()
  334. }
  335. const (
  336. Reset byte = 0
  337. FirstSYN byte = 6
  338. FirstACK byte = 1
  339. SndSYN byte = 2
  340. SndACK byte = 2
  341. Data byte = 4
  342. Ping byte = 5
  343. Close byte = 7
  344. CloseBack byte = 8
  345. ResetAck byte = 9
  346. Fake byte = 50
  347. )
  348. func __xor(s []byte, xor string) []byte {
  349. if len(xor) == 0 {
  350. return s
  351. }
  352. encodingData := []byte(xor)
  353. encodingLen := len(encodingData)
  354. n := len(s)
  355. if n == 0 {
  356. return s
  357. }
  358. for i := 0; i < n; i++ {
  359. s[i] = s[i] ^ encodingData[i%encodingLen]
  360. }
  361. return s
  362. }
  363. func _xor(s []byte, xor string) []byte {
  364. if len(xor) == 0 {
  365. return s
  366. }
  367. encodingData := []byte(xor)
  368. encodingLen := len(encodingData)
  369. n := len(s)
  370. if n == 0 {
  371. return s
  372. }
  373. r := make([]byte, n)
  374. for i := 0; i < n; i++ {
  375. r[i] = s[i] ^ encodingData[i%encodingLen]
  376. }
  377. return r
  378. }
  379. func makeEncode(buf []byte, status byte, arg int, arg2 int16, xor string) []byte {
  380. buf[0] = status
  381. binary.LittleEndian.PutUint32(buf[1:], uint32(arg))
  382. binary.LittleEndian.PutUint16(buf[5:], uint16(arg2))
  383. return _xor(buf, xor)
  384. }
  385. func makeDecode(data []byte, xor string) (status byte, arg int32, arg2 int16) {
  386. if len(data) < 7 {
  387. return Reset, 0, 0
  388. }
  389. data = _xor(data, xor)
  390. status = data[0]
  391. arg = int32(binary.LittleEndian.Uint32(data[1:5]))
  392. arg2 = int16(binary.LittleEndian.Uint16(data[5:7]))
  393. return
  394. }
  395. func (session *UDPMakeSession) SetFec(DataShards, ParityShards int) (er error) {
  396. session.fecDataShards = DataShards
  397. session.fecParityShards = ParityShards
  398. var fec reedsolomon.Encoder
  399. fec, er = reedsolomon.New(DataShards, ParityShards)
  400. if er != nil {
  401. return
  402. }
  403. session.fecR = &fec
  404. fec, er = reedsolomon.New(DataShards, ParityShards)
  405. if er == nil {
  406. session.fecRCacheTbl = make(map[uint]*fecInfo)
  407. session.fecWCacheTbl = nil
  408. session.fecW = &fec
  409. } else {
  410. session.fecR = nil
  411. }
  412. return
  413. }
  414. func (session *UDPMakeSession) doAndWait(f func(), sec int, readf func(status byte, arg int32, arg2 int16) int) (code int) {
  415. t := time.NewTicker(50 * time.Millisecond)
  416. currT := time.Now().Unix()
  417. f()
  418. out:
  419. for {
  420. select {
  421. case <-t.C:
  422. if time.Now().Unix()-currT >= int64(sec) {
  423. logs.Debug("session timeout")
  424. code = -1
  425. break out
  426. }
  427. session.sock.SetReadDeadline(time.Now().Add(2 * time.Second))
  428. n, from, err := session.sock.ReadFromUDP(session.readBuffer)
  429. if err != nil {
  430. e, ok := err.(net.Error)
  431. if !ok || !e.Timeout() {
  432. logs.Debug("recv error", err.Error(), from)
  433. code = -2
  434. break out
  435. }
  436. } else {
  437. code = readf(makeDecode(session.readBuffer[:n], session.xor))
  438. if code >= 0 {
  439. break out
  440. }
  441. }
  442. go f()
  443. }
  444. }
  445. t.Stop()
  446. if code > 0 {
  447. logs.Debug("handshake fail,got code", code)
  448. }
  449. return
  450. }
  451. func (session *UDPMakeSession) writeTo(b []byte) {
  452. if session.compressCache != nil && len(b) > 7 {
  453. enc, er := zappy.Encode(session.compressCache, b)
  454. if er != nil {
  455. logs.Debug("compress error", er.Error())
  456. go session.Close()
  457. return
  458. }
  459. //logs.Debug("compress", len(b), len(enc))
  460. session.sock.WriteTo(__xor(enc, session.xor), session.remote)
  461. } else {
  462. session.sock.WriteTo(__xor(b, session.xor), session.remote)
  463. }
  464. }
  465. func (session *UDPMakeSession) output(b []byte) {
  466. if session.fecW == nil {
  467. session.writeTo(b)
  468. } else {
  469. id := session.fecWriteId
  470. session.fecSendC++
  471. info := session.fecWCacheTbl
  472. if info == nil {
  473. info = &fecInfo{make([][]byte, session.fecDataShards+session.fecParityShards), time.Now().Unix() + 15}
  474. session.fecWCacheTbl = info
  475. }
  476. _b := make([]byte, len(b)+7)
  477. _len := len(b)
  478. _b[0] = byte(_len & 0xff)
  479. _b[1] = byte((_len >> 8) & 0xff)
  480. _b[2] = byte(id & 0xff)
  481. _b[3] = byte((id >> 8) & 0xff)
  482. _b[4] = byte((id >> 16) & 0xff)
  483. _b[5] = byte((id >> 32) & 0xff)
  484. _b[6] = byte(session.fecSendC - 1)
  485. copy(_b[7:], b)
  486. info.bytes[session.fecSendC-1] = _b
  487. if session.fecSendL < len(_b) {
  488. session.fecSendL = len(_b)
  489. }
  490. session.writeTo(_b)
  491. if session.fecSendC >= uint(session.fecDataShards) {
  492. for i := 0; i < session.fecDataShards; i++ {
  493. if session.fecSendL > len(info.bytes[i]) {
  494. __b := make([]byte, session.fecSendL)
  495. copy(__b, info.bytes[i])
  496. info.bytes[i] = __b
  497. }
  498. }
  499. for i := 0; i < session.fecParityShards; i++ {
  500. info.bytes[i+session.fecDataShards] = make([]byte, session.fecSendL)
  501. }
  502. er := (*session.fecW).Encode(info.bytes)
  503. if er != nil {
  504. //logs.Debug("wocao encode err", er.Error())
  505. go session.Close()
  506. return
  507. }
  508. for i := session.fecDataShards; i < session.fecDataShards+session.fecParityShards; i++ {
  509. //if rand.Intn(100) >= 15 {
  510. _info := info.bytes[i]
  511. session.writeTo(_info)
  512. //_len := int(_info[0]) | (int(_info[1]) << 8)
  513. //logs.Debug("output udp id", id, i, _len, len(_info))
  514. //} else {
  515. // logs.Debug("drop output udp id", id, i, _len, len(_info))
  516. //}
  517. }
  518. session.fecWCacheTbl = nil
  519. session.fecSendC = 0
  520. session.fecSendL = 0
  521. session.fecWriteId++
  522. //logs.Debug("flush id", id)
  523. }
  524. //logs.Debug("output sn", c.fecWriteId, c.fecSendC, _len)
  525. }
  526. }
  527. func (session *UDPMakeSession) serverDo(s string) {
  528. go func() {
  529. //logs.Debug("prepare handshake", session.remote)
  530. select {
  531. case <-session.handShakeChanQuit:
  532. case session.handShakeChan <- s:
  533. }
  534. }()
  535. }
  536. func (session *UDPMakeSession) serverInit(l *UDPListener, setting *KcpSetting) {
  537. go func() {
  538. c := time.NewTicker(50 * time.Millisecond)
  539. defer func() {
  540. close(session.handShakeChanQuit)
  541. c.Stop()
  542. if session.status != "ok" {
  543. session.Close()
  544. }
  545. }()
  546. overTime := time.Now().Unix() + session.timeout
  547. for {
  548. select {
  549. case s := <-session.handShakeChan:
  550. //logs.Debug("process handshake", session.remote)
  551. status, arg, arg2 := makeDecode([]byte(s), session.xor)
  552. switch session.status {
  553. case "init":
  554. if status != FirstSYN {
  555. logs.Debug("status != FirstSYN, reset", session.remote, status)
  556. session.sock.WriteToUDP(makeEncode(session.encodeBuffer, Reset, 0, 0, session.xor), session.remote)
  557. return
  558. }
  559. _mainV, _subV := int(byte(arg>>24)), int(byte(arg>>16))
  560. if _mainV != mainV || _subV != subV {
  561. session.sock.WriteToUDP(makeEncode(session.encodeBuffer, ResetAck, (mainV<<24)+(subV<<16), 0, session.xor), session.remote)
  562. logs.Debug("pipe version not eq,kickout,%d.%d=>%d.%d", mainV, subV, _mainV, _subV)
  563. return
  564. }
  565. session.status = "firstack"
  566. session.timeout = int64(arg & 0xff)
  567. session.sock.WriteToUDP(makeEncode(session.encodeBuffer, FirstACK, session.id, 0, session.xor), session.remote)
  568. overTime = time.Now().Unix() + session.timeout
  569. case "firstack":
  570. if status != SndSYN {
  571. logs.Debug("status != SndSYN, nothing", session.remote, status)
  572. /*
  573. if status == FirstSYN {
  574. session.sock.WriteToUDP(makeEncode(session.encodeBuffer, FirstACK, session.id), session.remote, session.xor)
  575. }*/
  576. return
  577. }
  578. if arg2 > 0 {
  579. session.confuseSeed = rand.Intn(int(WriteBufferSize/2)) + 10
  580. logs.Debug("confuse!")
  581. }
  582. session.status = "ok"
  583. session.kcp = ikcp.Ikcp_create(uint32(session.id), session)
  584. session.kcp.Output = udp_output
  585. ikcp.Ikcp_wndsize(session.kcp, setting.Sndwnd, setting.Rcvwnd)
  586. ikcp.Ikcp_nodelay(session.kcp, setting.Nodelay, setting.Interval, setting.Resend, setting.Nc)
  587. ikcp.Ikcp_setmtu(session.kcp, setting.Mtu)
  588. go session.loop()
  589. go func() {
  590. select {
  591. case l.connChan <- session:
  592. case <-l.quitChan:
  593. }
  594. }()
  595. session.sock.WriteToUDP(makeEncode(session.encodeBuffer, SndACK, session.id, 0, session.xor), session.remote)
  596. overTime = time.Now().Unix() + session.timeout
  597. }
  598. case <-c.C:
  599. if time.Now().Unix() > overTime {
  600. return
  601. }
  602. switch session.status {
  603. case "firstack":
  604. session.sock.WriteToUDP(makeEncode(session.encodeBuffer, FirstACK, session.id, 0, session.xor), session.remote)
  605. case "ok":
  606. buf := make([]byte, 7)
  607. session.sock.WriteToUDP(makeEncode(buf, SndACK, session.id, 0, session.xor), session.remote)
  608. }
  609. }
  610. }
  611. }()
  612. }
  613. func (session *UDPMakeSession) loop() {
  614. curr := time.Now().Unix()
  615. session.overTime = curr + session.timeout
  616. ping := make(chan struct{})
  617. pingC := 0
  618. callUpdate := false
  619. updateC := make(chan struct{})
  620. if session.listener == nil {
  621. go func() {
  622. tmp := session.readBuffer
  623. t := time.Time{}
  624. session.sock.SetReadDeadline(t)
  625. for {
  626. n, from, err := session.sock.ReadFromUDP(tmp)
  627. if err != nil {
  628. e, ok := err.(net.Error)
  629. if !ok || !e.Timeout() {
  630. break
  631. }
  632. }
  633. if session.remote.String() == from.String() {
  634. //logs.Debug("===", n, len(session.compressCache))
  635. if n >= int(ikcp.IKCP_OVERHEAD) || n <= 7 || session.compressCache != nil {
  636. __xor(session.readBuffer[:n], session.xor)
  637. var buf []byte
  638. if n <= 7 || session.compressCache == nil {
  639. buf = make([]byte, n)
  640. copy(buf, session.readBuffer[:n])
  641. } else {
  642. _b, _er := zappy.Decode(nil, session.readBuffer[:n])
  643. if _er != nil {
  644. logs.Debug("decompress fail", _er.Error())
  645. //go session.Close()
  646. //don't close pipe, just drop data
  647. continue
  648. }
  649. buf = _b
  650. //logs.Debug("decompress", n, len(_b))
  651. }
  652. session.DoAction2("input", buf, len(buf))
  653. }
  654. }
  655. }
  656. }()
  657. }
  658. updateF := func(n time.Duration) {
  659. if !callUpdate {
  660. callUpdate = true
  661. time.AfterFunc(n*time.Millisecond, func() {
  662. select {
  663. case updateC <- struct{}{}:
  664. case <-session.quitChan:
  665. }
  666. })
  667. }
  668. }
  669. fastCheck := false
  670. waitList := [](chan struct{}){}
  671. recoverChan := make(chan struct{})
  672. var waitRecvCache *KcpCache
  673. var forceWait int64 = 0
  674. go func() {
  675. out:
  676. for {
  677. select {
  678. //session.wait.Done()
  679. case <-ping:
  680. updateF(50)
  681. pingC++
  682. if pingC >= 4 {
  683. pingC = 0
  684. if ikcp.Ikcp_waitsnd(session.kcp) <= dataLimit/2 {
  685. go session.DoWrite(makeEncode(session.encodeBuffer, Ping, 0, 0, session.xor))
  686. }
  687. if session.fecR != nil {
  688. curr := time.Now().Unix()
  689. for id, info := range session.fecRCacheTbl {
  690. if curr >= info.overTime {
  691. delete(session.fecRCacheTbl, id)
  692. if session.fecRecvId <= id {
  693. session.fecRecvId = id + 1
  694. }
  695. //logs.Debug("timeout after del", id, len(c.fecRCacheTbl))
  696. }
  697. }
  698. }
  699. if forceWait > 0 {
  700. if time.Now().Unix() > forceWait && ikcp.Ikcp_waitsnd(session.kcp) <= dataLimit/2 {
  701. forceWait = 0
  702. go func() {
  703. select {
  704. case <-session.quitChan:
  705. case recoverChan <- struct{}{}:
  706. }
  707. }()
  708. }
  709. }
  710. }
  711. if time.Now().Unix() > session.overTime {
  712. logs.Debug("overtime close", session.LocalAddr().String(), session.RemoteAddr().String())
  713. go session.Close()
  714. } else {
  715. time.AfterFunc(300*time.Millisecond, func() {
  716. select {
  717. case ping <- struct{}{}:
  718. case <-session.quitChan:
  719. }
  720. })
  721. }
  722. case <-recoverChan:
  723. fastCheck = false
  724. for _, r := range waitList {
  725. //logs.Debug("recover writing data")
  726. select {
  727. case r <- struct{}{}:
  728. case <-session.quitChan:
  729. }
  730. }
  731. waitList = [](chan struct{}){}
  732. case c := <-session.checkCanWrite:
  733. if session.confuseSeed > 0 {
  734. if WriteBufferSize/2-rand.Intn(session.confuseSeed) < 500 || rand.Intn(100) < 5 {
  735. //make a sleep
  736. forceWait = time.Now().Add(time.Millisecond * time.Duration(rand.Intn(1000))).Unix()
  737. }
  738. }
  739. if ikcp.Ikcp_waitsnd(session.kcp) > dataLimit {
  740. //logs.Debug("wait for data limit")
  741. waitList = append(waitList, c)
  742. if !fastCheck {
  743. fastCheck = true
  744. var f func()
  745. f = func() {
  746. n := ikcp.Ikcp_waitsnd(session.kcp)
  747. //logs.Debug("fast check!", n, len(waitList))
  748. if n <= dataLimit/2 {
  749. select {
  750. case <-session.quitChan:
  751. logs.Debug("recover writing data quit")
  752. case recoverChan <- struct{}{}:
  753. }
  754. } else {
  755. updateF(20)
  756. time.AfterFunc(40*time.Millisecond, f)
  757. }
  758. }
  759. time.AfterFunc(20*time.Millisecond, f)
  760. }
  761. } else if forceWait > 0 && time.Now().Unix() < forceWait {
  762. waitList = append(waitList, c)
  763. } else {
  764. select {
  765. case c <- struct{}{}:
  766. case <-session.quitChan:
  767. }
  768. }
  769. case ca := <-session.recvChan:
  770. tmp := session.processBuffer
  771. for {
  772. hr := ikcp.Ikcp_recv(session.kcp, tmp, ReadBufferSize)
  773. if hr > 0 {
  774. status := tmp[0]
  775. if status == Data {
  776. n := 1
  777. l := hr
  778. if session.confuseSeed > 0 {
  779. n = 3
  780. l = (int32(tmp[1]) | (int32(tmp[2]) << 8)) + int32(n)
  781. }
  782. //logs.Debug("try recv", hr, n, l)
  783. copy(ca.b, tmp[n:int(l)])
  784. select {
  785. case ca.c <- int(int(l) - n):
  786. case <-session.quitChan:
  787. }
  788. break
  789. } else if status == Fake {
  790. } else {
  791. session.DoAction("recv", status)
  792. }
  793. } else {
  794. waitRecvCache = &ca
  795. break
  796. }
  797. }
  798. case action := <-session.do2:
  799. switch action.t {
  800. case "input":
  801. session.overTime = time.Now().Unix() + session.timeout
  802. args := action.args
  803. s := args[0].([]byte)
  804. n := args[1].(int)
  805. if session.fecR != nil {
  806. if len(s) <= 7 {
  807. if n < 7 {
  808. logs.Debug("recv reset")
  809. go session._Close(false)
  810. break
  811. } else if n == 7 {
  812. status, _, _ := makeDecode(s, session.xor)
  813. if status == Reset || status == ResetAck {
  814. logs.Debug("recv reset2", status)
  815. go session._Close(false)
  816. }
  817. break
  818. }
  819. break
  820. }
  821. id := uint(int(s[2]) | (int(s[3]) << 8) | (int(s[4]) << 16) | (int(s[5]) << 24))
  822. var seq uint = uint(s[6])
  823. _len := int(s[0]) | (int(s[1]) << 8)
  824. //binary.Read(head[:4], binary.LittleEndian, &id)
  825. if id < session.fecRecvId {
  826. //logs.Debug("drop id for noneed", id, seq)
  827. break
  828. }
  829. if seq < uint(session.fecDataShards) {
  830. ikcp.Ikcp_input(session.kcp, s[7:], _len)
  831. //logs.Debug("direct input udp", id, seq, _len)
  832. }
  833. if seq >= uint(session.fecDataShards+session.fecParityShards) {
  834. logs.Debug("-ds and -ps must be equal on both sides")
  835. go session.Close()
  836. break
  837. }
  838. tbl, have := session.fecRCacheTbl[id]
  839. if !have {
  840. tbl = &fecInfo{make([][]byte, session.fecDataShards+session.fecParityShards), time.Now().Unix() + 3}
  841. session.fecRCacheTbl[id] = tbl
  842. }
  843. //logs.Debug("got", id, seq, n, _len)
  844. if tbl.bytes[seq] != nil {
  845. //dup, drop
  846. break
  847. } else {
  848. tbl.bytes[seq] = s
  849. }
  850. count := 0
  851. reaL := 0
  852. for _, v := range tbl.bytes {
  853. if v != nil {
  854. count++
  855. if reaL < len(v) {
  856. reaL = len(v)
  857. }
  858. }
  859. }
  860. if count >= session.fecDataShards {
  861. markTbl := make(map[int]bool, len(tbl.bytes))
  862. for _seq, _b := range tbl.bytes {
  863. if _b != nil {
  864. markTbl[_seq] = true
  865. }
  866. }
  867. bNeedRebuild := false
  868. for i, v := range tbl.bytes {
  869. if v != nil {
  870. if i >= session.fecDataShards {
  871. bNeedRebuild = true
  872. }
  873. if len(v) < reaL {
  874. _b := make([]byte, reaL)
  875. copy(_b, v)
  876. tbl.bytes[i] = _b
  877. }
  878. }
  879. }
  880. if bNeedRebuild {
  881. er := (*session.fecR).Reconstruct(tbl.bytes)
  882. if er != nil {
  883. //logs.Debug("2Reconstruct fail, close pipe", count, session.fecDataShards, session.fecParityShards, er.Error())
  884. //break //broken data, may be should be closed, now just keep input
  885. } else {
  886. //logs.Debug("Reconstruct ok, input", id)
  887. for i := 0; i < session.fecDataShards; i++ {
  888. if _, have := markTbl[i]; !have {
  889. _len := int(tbl.bytes[i][0]) | (int(tbl.bytes[i][1]) << 8)
  890. ikcp.Ikcp_input(session.kcp, tbl.bytes[i][7:], int(_len))
  891. //logs.Debug("fec input for mark ok", i, id, _len)
  892. }
  893. }
  894. }
  895. }
  896. delete(session.fecRCacheTbl, id)
  897. //logs.Debug("after del", id, len(c.fecRCacheTbl))
  898. if session.fecRecvId <= id {
  899. session.fecRecvId = id + 1
  900. }
  901. }
  902. } else {
  903. if n < 7 {
  904. logs.Debug("recv reset")
  905. go session._Close(false)
  906. break
  907. } else if n == 7 {
  908. status, _, _ := makeDecode(s, session.xor)
  909. if status == Reset || status == ResetAck {
  910. logs.Debug("recv reset2", status)
  911. go session._Close(false)
  912. }
  913. break
  914. }
  915. ikcp.Ikcp_input(session.kcp, s, n)
  916. }
  917. if waitRecvCache != nil {
  918. ca := *waitRecvCache
  919. tmp := session.processBuffer
  920. for {
  921. hr := ikcp.Ikcp_recv(session.kcp, tmp, ReadBufferSize)
  922. if hr > 0 {
  923. status := tmp[0]
  924. if status == Data {
  925. n := 1
  926. l := hr
  927. if session.confuseSeed > 0 {
  928. n = 3
  929. l = (int32(tmp[1]) | (int32(tmp[2]) << 8)) + int32(n)
  930. }
  931. //logs.Debug("try recv", n, l)
  932. waitRecvCache = nil
  933. copy(ca.b, tmp[n:int(l)])
  934. select {
  935. case ca.c <- (int(l) - n):
  936. case <-session.quitChan:
  937. }
  938. break
  939. } else if status == Fake {
  940. } else {
  941. session.DoAction("recv", status)
  942. }
  943. } else {
  944. break
  945. }
  946. }
  947. }
  948. updateF(10)
  949. case "write":
  950. b := action.args[0].([]byte)
  951. ikcp.Ikcp_send(session.kcp, b, len(b))
  952. updateF(10)
  953. }
  954. case <-session.quitChan:
  955. break out
  956. case <-updateC:
  957. now := uint32(iclock())
  958. ikcp.Ikcp_update(session.kcp, now)
  959. callUpdate = false
  960. }
  961. }
  962. }()
  963. select {
  964. case ping <- struct{}{}:
  965. case <-session.quitChan:
  966. }
  967. out:
  968. for {
  969. select {
  970. case action := <-session.do:
  971. switch action.t {
  972. case "quit":
  973. if session.closed {
  974. break
  975. }
  976. session._Close(true)
  977. case "closebegin":
  978. //A tell B to close and wait
  979. time.AfterFunc(time.Millisecond*500, func() {
  980. //logs.Debug("close over, step3", session.LocalAddr().String(), session.RemoteAddr().String())
  981. go session.DoAction("closeover")
  982. })
  983. buf := make([]byte, 7)
  984. go session.DoWrite(makeEncode(buf, Close, 0, 0, session.xor))
  985. case "closeover":
  986. //A call timeover
  987. close(session.closeChan)
  988. break out
  989. case "closeend":
  990. //B call close
  991. session._Close(false)
  992. break out
  993. case "recv":
  994. status := (action.args[0]).(byte)
  995. switch status {
  996. case CloseBack:
  997. //A call from B
  998. //logs.Debug("recv back close, step2", session.LocalAddr().String(), session.RemoteAddr().String())
  999. go session.DoAction("closeover")
  1000. case Close:
  1001. if session.closed {
  1002. break
  1003. }
  1004. if session.status != "ok" {
  1005. session._Close(false)
  1006. } else {
  1007. //logs.Debug("recv remote close, step1", session.LocalAddr().String(), session.RemoteAddr().String())
  1008. buf := make([]byte, 7)
  1009. go session.DoWrite(makeEncode(buf, CloseBack, 0, 0, session.xor))
  1010. time.AfterFunc(time.Millisecond*500, func() {
  1011. //logs.Debug("close remote over, step4", session.LocalAddr().String(), session.RemoteAddr().String())
  1012. if session.closed {
  1013. return
  1014. }
  1015. go session.DoAction("closeend")
  1016. })
  1017. }
  1018. case Reset:
  1019. logs.Debug("recv reset")
  1020. go session._Close(false)
  1021. case Ping:
  1022. default:
  1023. if session.status != "ok" {
  1024. session.Close()
  1025. }
  1026. }
  1027. }
  1028. }
  1029. }
  1030. }
  1031. func (session *UDPMakeSession) _Close(bFirstCall bool) {
  1032. if session.closed {
  1033. return
  1034. }
  1035. session.closed = true
  1036. //session.wait.Wait()
  1037. go func() {
  1038. //logs.Debug("pipe begin close", session.LocalAddr().String(), session.RemoteAddr().String())
  1039. if bFirstCall {
  1040. go session.DoAction("closebegin", session.LocalAddr().String())
  1041. <-session.closeChan
  1042. } else {
  1043. close(session.closeChan)
  1044. }
  1045. //logs.Debug("pipe end close", session.id)
  1046. close(session.quitChan)
  1047. if session.listener != nil {
  1048. session.listener.remove(session.remote.String())
  1049. } else {
  1050. if session.sock != nil {
  1051. session.sock.Close()
  1052. }
  1053. }
  1054. }()
  1055. }
  1056. func (session *UDPMakeSession) Close() error {
  1057. if session.closed {
  1058. return nil
  1059. }
  1060. if session.status != "ok" {
  1061. session._Close(false)
  1062. return nil
  1063. }
  1064. go session.DoAction("quit")
  1065. <-session.closeChan
  1066. return nil
  1067. }
  1068. func (session *UDPMakeSession) DoAction2(action string, args ...interface{}) {
  1069. //session.wait.Add(1)
  1070. //logs.Debug(action, len(args))
  1071. select {
  1072. case session.do2 <- Action{t: action, args: args}:
  1073. case <-session.quitChan:
  1074. //session.wait.Done()
  1075. }
  1076. }
  1077. func (session *UDPMakeSession) DoAction(action string, args ...interface{}) {
  1078. //session.wait.Add(1)
  1079. //logs.Debug(action, len(args))
  1080. select {
  1081. case session.do <- Action{t: action, args: args}:
  1082. case <-session.quitChan:
  1083. //session.wait.Done()
  1084. }
  1085. }
  1086. func (session *UDPMakeSession) processInput(s string, n int) {
  1087. }
  1088. func (session *UDPMakeSession) LocalAddr() net.Addr {
  1089. return session.sock.LocalAddr()
  1090. }
  1091. func (session *UDPMakeSession) RemoteAddr() net.Addr {
  1092. return net.Addr(session.remote)
  1093. }
  1094. func (session *UDPMakeSession) SetDeadline(t time.Time) error {
  1095. return session.sock.SetDeadline(t)
  1096. }
  1097. func (session *UDPMakeSession) SetReadDeadline(t time.Time) error {
  1098. return session.sock.SetReadDeadline(t)
  1099. }
  1100. func (session *UDPMakeSession) SetWriteDeadline(t time.Time) error {
  1101. return session.sock.SetWriteDeadline(t)
  1102. }
  1103. func (session *UDPMakeSession) DoWrite(s []byte) bool {
  1104. wc := make(chan struct{})
  1105. select {
  1106. case session.checkCanWrite <- wc:
  1107. select {
  1108. case <-wc:
  1109. case <-session.quitChan:
  1110. return false
  1111. }
  1112. session.DoAction2("write", s)
  1113. return true
  1114. case <-session.quitChan:
  1115. return false
  1116. }
  1117. }
  1118. func (session *UDPMakeSession) Write(b []byte) (n int, err error) {
  1119. sendL := len(b)
  1120. if sendL == 0 || session.status != "ok" {
  1121. return 0, nil
  1122. }
  1123. var data []byte
  1124. if session.confuseSeed <= 0 {
  1125. data = make([]byte, sendL+1)
  1126. data[0] = Data
  1127. copy(data[1:], b)
  1128. } else {
  1129. remain := WriteBufferSize - sendL - 3
  1130. if remain > 0 && (WriteBufferSize/2-rand.Intn(session.confuseSeed) < 500 || rand.Intn(100) < 5) {
  1131. if remain > 1000 {
  1132. remain = 1000
  1133. }
  1134. remain = rand.Intn(remain / 2)
  1135. } else {
  1136. remain = 0
  1137. }
  1138. data = make([]byte, sendL+remain+3)
  1139. data[0] = Data
  1140. data[1] = byte(sendL & 0xff)
  1141. data[2] = byte((sendL >> 8) & 0xff)
  1142. copy(data[3:], b)
  1143. //logs.Debug("try send", len(data), sendL)
  1144. //copy(data[3+sendL:], b)
  1145. }
  1146. ok := session.DoWrite(data)
  1147. if !ok {
  1148. return 0, errors.New("closed")
  1149. }
  1150. if session.confuseSeed > 0 {
  1151. if rand.Intn(WriteBufferSize/2) > session.confuseSeed/2 {
  1152. n := rand.Intn(session.confuseSeed)
  1153. if n > 10 {
  1154. d := make([]byte, n)
  1155. d[0] = Fake
  1156. session.DoWrite(d)
  1157. }
  1158. }
  1159. }
  1160. return sendL, err
  1161. }
  1162. //udp read does not relay on the len(p), please make a big enough array to cache data
  1163. func (session *UDPMakeSession) Read(p []byte) (n int, err error) {
  1164. wc := KcpCache{p, 0, make(chan int)}
  1165. select {
  1166. case session.recvChan <- wc:
  1167. select {
  1168. case n = <-wc.c:
  1169. case <-session.quitChan:
  1170. n = -1
  1171. }
  1172. case <-session.quitChan:
  1173. n = -1
  1174. }
  1175. //logs.Debug("real recv", l, string(b[:l]))
  1176. if n == -1 {
  1177. return 0, errors.New("force quit for read error")
  1178. } else {
  1179. return n, nil
  1180. }
  1181. }
  1182. type _reuseTbl struct {
  1183. tbl map[int]bool
  1184. }
  1185. var currIdMap map[string]int
  1186. var reuseTbl map[string]*_reuseTbl
  1187. func GetId(name string) int {
  1188. if currIdMap == nil {
  1189. currIdMap = make(map[string]int)
  1190. currIdMap[name] = 0
  1191. }
  1192. i, _ := currIdMap[name]
  1193. i++
  1194. if i >= 2147483647 {
  1195. i = 0
  1196. }
  1197. currIdMap[name] = i
  1198. // println("gen new id", currIdMap[name])
  1199. return i
  1200. }
  1201. func RmId(name string, id int) {
  1202. return
  1203. }
粤ICP备19079148号