net.go 17 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825
  1. package tcptransport
  2. import (
  3. "encoding/gob"
  4. "fmt"
  5. "log"
  6. "net"
  7. "sync"
  8. "sync/atomic"
  9. "time"
  10. "trial/go-chord"
  11. )
  12. /*
  13. TCPTransport provides a TCP based Chord transport layer. This allows Chord
  14. to be implemented over a network, instead of only using the LocalTransport. It is
  15. meant to be a simple implementation, optimizing for simplicity instead of performance.
  16. Messages are sent with a header frame, followed by a body frame. All data is encoded
  17. using the GOB format for simplicity.
  18. Internally, there is 1 Goroutine listening for inbound connections, 1 Goroutine PER
  19. inbound connection.
  20. */
  21. type TCPTransport struct {
  22. sock *net.TCPListener
  23. timeout time.Duration
  24. maxIdle time.Duration
  25. lock sync.RWMutex
  26. local map[string]*chord.LocalRPC
  27. inbound map[*net.TCPConn]struct{}
  28. poolLock sync.Mutex
  29. pool map[string][]*tcpOutConn
  30. shutdown int32
  31. }
  32. type tcpOutConn struct {
  33. host string
  34. sock *net.TCPConn
  35. header tcpHeader
  36. enc *gob.Encoder
  37. dec *gob.Decoder
  38. used time.Time
  39. }
  40. const (
  41. tcpPing = iota
  42. tcpListReq
  43. tcpGetPredReq
  44. tcpNotifyReq
  45. tcpFindSucReq
  46. tcpClearPredReq
  47. tcpSkipSucReq
  48. )
  49. type tcpHeader struct {
  50. ReqType int
  51. }
  52. // Potential body types
  53. type tcpBodyError struct {
  54. Err error
  55. }
  56. type tcpBodyString struct {
  57. S string
  58. }
  59. type tcpBodyVnode struct {
  60. Vn *chord.Vnode
  61. }
  62. type tcpBodyTwoVnode struct {
  63. Target *chord.Vnode
  64. Vn *chord.Vnode
  65. }
  66. type tcpBodyFindSuc struct {
  67. Target *chord.Vnode
  68. Num int
  69. Key []byte
  70. }
  71. type tcpBodyVnodeError struct {
  72. Vnode *chord.Vnode
  73. Err error
  74. }
  75. type tcpBodyVnodeListError struct {
  76. Vnodes []*chord.Vnode
  77. Err error
  78. }
  79. type tcpBodyBoolError struct {
  80. B bool
  81. Err error
  82. }
  83. // Creates a new TCP transport on the given listen address with the
  84. // configured timeout duration.
  85. func InitTCPTransport(listen string, timeout time.Duration) (*TCPTransport, error) {
  86. // Try to start the listener
  87. sock, err := net.Listen("tcp", listen)
  88. if err != nil {
  89. return nil, err
  90. }
  91. // allocate maps
  92. local := make(map[string]*chord.LocalRPC)
  93. inbound := make(map[*net.TCPConn]struct{})
  94. pool := make(map[string][]*tcpOutConn)
  95. // Maximum age of a connection
  96. maxIdle := time.Duration(300 * time.Second)
  97. // Setup the transport
  98. tcp := &TCPTransport{sock: sock.(*net.TCPListener),
  99. timeout: timeout,
  100. maxIdle: maxIdle,
  101. local: local,
  102. inbound: inbound,
  103. pool: pool}
  104. // Listen for connections
  105. go tcp.listen()
  106. // Reap old connections
  107. go tcp.reapOld()
  108. // Done
  109. return tcp, nil
  110. }
  111. // Checks for a local vnode
  112. func (t *TCPTransport) get(vn *chord.Vnode) (chord.VnodeRPC, bool) {
  113. key := vn.String()
  114. t.lock.RLock()
  115. defer t.lock.RUnlock()
  116. w, ok := t.local[key]
  117. if ok {
  118. return w.Obj, ok
  119. } else {
  120. return nil, ok
  121. }
  122. }
  123. // Gets an outbound connection to a host
  124. func (t *TCPTransport) getConn(host string) (*tcpOutConn, error) {
  125. // Check if we have a conn cached
  126. var out *tcpOutConn
  127. t.poolLock.Lock()
  128. if atomic.LoadInt32(&t.shutdown) == 1 {
  129. t.poolLock.Unlock()
  130. return nil, fmt.Errorf("TCP transport is shutdown")
  131. }
  132. list, ok := t.pool[host]
  133. if ok && len(list) > 0 {
  134. out = list[len(list)-1]
  135. list = list[:len(list)-1]
  136. t.pool[host] = list
  137. }
  138. t.poolLock.Unlock()
  139. if out != nil {
  140. // Verify that the socket is valid. Might be closed.
  141. if _, err := out.sock.Read(nil); err == nil {
  142. return out, nil
  143. }
  144. out.sock.Close()
  145. }
  146. // Try to establish a connection
  147. conn, err := net.DialTimeout("tcp", host, t.timeout)
  148. if err != nil {
  149. return nil, err
  150. }
  151. // Setup the socket
  152. sock := conn.(*net.TCPConn)
  153. t.setupConn(sock)
  154. enc := gob.NewEncoder(sock)
  155. dec := gob.NewDecoder(sock)
  156. now := time.Now()
  157. // Wrap the sock
  158. out = &tcpOutConn{host: host, sock: sock, enc: enc, dec: dec, used: now}
  159. return out, nil
  160. }
  161. // Returns an outbound TCP connection to the pool
  162. func (t *TCPTransport) returnConn(o *tcpOutConn) {
  163. // Update the last used time
  164. o.used = time.Now()
  165. // Push back into the pool
  166. t.poolLock.Lock()
  167. defer t.poolLock.Unlock()
  168. if atomic.LoadInt32(&t.shutdown) == 1 {
  169. o.sock.Close()
  170. return
  171. }
  172. list, _ := t.pool[o.host]
  173. t.pool[o.host] = append(list, o)
  174. }
  175. // Setup a connection
  176. func (t *TCPTransport) setupConn(c *net.TCPConn) {
  177. c.SetNoDelay(true)
  178. c.SetKeepAlive(true)
  179. }
  180. // Gets a list of the Vnodes on the box
  181. func (t *TCPTransport) ListVnodes(host string) ([]*chord.Vnode, error) {
  182. // Get a conn
  183. out, err := t.getConn(host)
  184. if err != nil {
  185. return nil, err
  186. }
  187. // Response channels
  188. respChan := make(chan []*chord.Vnode, 1)
  189. errChan := make(chan error, 1)
  190. go func() {
  191. // Send a list command
  192. out.header.ReqType = tcpListReq
  193. body := tcpBodyString{S: host}
  194. if err := out.enc.Encode(&out.header); err != nil {
  195. errChan <- err
  196. return
  197. }
  198. if err := out.enc.Encode(&body); err != nil {
  199. errChan <- err
  200. return
  201. }
  202. // Read in the response
  203. resp := tcpBodyVnodeListError{}
  204. if err := out.dec.Decode(&resp); err != nil {
  205. errChan <- err
  206. }
  207. // Return the connection
  208. t.returnConn(out)
  209. if resp.Err == nil {
  210. respChan <- resp.Vnodes
  211. } else {
  212. errChan <- resp.Err
  213. }
  214. }()
  215. select {
  216. case <-time.After(t.timeout):
  217. return nil, fmt.Errorf("Command timed out!")
  218. case err := <-errChan:
  219. return nil, err
  220. case res := <-respChan:
  221. return res, nil
  222. }
  223. }
  224. // Ping a Vnode, check for liveness
  225. func (t *TCPTransport) Ping(vn *chord.Vnode) (bool, error) {
  226. // Get a conn
  227. out, err := t.getConn(vn.Host)
  228. if err != nil {
  229. return false, err
  230. }
  231. // Response channels
  232. respChan := make(chan bool, 1)
  233. errChan := make(chan error, 1)
  234. go func() {
  235. // Send a list command
  236. out.header.ReqType = tcpPing
  237. body := tcpBodyVnode{Vn: vn}
  238. if err := out.enc.Encode(&out.header); err != nil {
  239. errChan <- err
  240. return
  241. }
  242. if err := out.enc.Encode(&body); err != nil {
  243. errChan <- err
  244. return
  245. }
  246. // Read in the response
  247. resp := tcpBodyBoolError{}
  248. if err := out.dec.Decode(&resp); err != nil {
  249. errChan <- err
  250. return
  251. }
  252. // Return the connection
  253. t.returnConn(out)
  254. if resp.Err == nil {
  255. respChan <- resp.B
  256. } else {
  257. errChan <- resp.Err
  258. }
  259. }()
  260. select {
  261. case <-time.After(t.timeout):
  262. return false, fmt.Errorf("Command timed out!")
  263. case err := <-errChan:
  264. return false, err
  265. case res := <-respChan:
  266. return res, nil
  267. }
  268. }
  269. // Request a nodes predecessor
  270. func (t *TCPTransport) GetPredecessor(vn *chord.Vnode) (*chord.Vnode, error) {
  271. // Get a conn
  272. out, err := t.getConn(vn.Host)
  273. if err != nil {
  274. return nil, err
  275. }
  276. respChan := make(chan *chord.Vnode, 1)
  277. errChan := make(chan error, 1)
  278. go func() {
  279. // Send a list command
  280. out.header.ReqType = tcpGetPredReq
  281. body := tcpBodyVnode{Vn: vn}
  282. if err := out.enc.Encode(&out.header); err != nil {
  283. errChan <- err
  284. return
  285. }
  286. if err := out.enc.Encode(&body); err != nil {
  287. errChan <- err
  288. return
  289. }
  290. // Read in the response
  291. resp := tcpBodyVnodeError{}
  292. if err := out.dec.Decode(&resp); err != nil {
  293. errChan <- err
  294. return
  295. }
  296. // Return the connection
  297. t.returnConn(out)
  298. if resp.Err == nil {
  299. respChan <- resp.Vnode
  300. } else {
  301. errChan <- resp.Err
  302. }
  303. }()
  304. select {
  305. case <-time.After(t.timeout):
  306. return nil, fmt.Errorf("Command timed out!")
  307. case err := <-errChan:
  308. return nil, err
  309. case res := <-respChan:
  310. return res, nil
  311. }
  312. }
  313. // Notify our successor of ourselves
  314. func (t *TCPTransport) Notify(target, self *chord.Vnode) ([]*chord.Vnode, error) {
  315. // Get a conn
  316. out, err := t.getConn(target.Host)
  317. if err != nil {
  318. return nil, err
  319. }
  320. respChan := make(chan []*chord.Vnode, 1)
  321. errChan := make(chan error, 1)
  322. go func() {
  323. // Send a list command
  324. out.header.ReqType = tcpNotifyReq
  325. body := tcpBodyTwoVnode{Target: target, Vn: self}
  326. if err := out.enc.Encode(&out.header); err != nil {
  327. errChan <- err
  328. return
  329. }
  330. if err := out.enc.Encode(&body); err != nil {
  331. errChan <- err
  332. return
  333. }
  334. // Read in the response
  335. resp := tcpBodyVnodeListError{}
  336. if err := out.dec.Decode(&resp); err != nil {
  337. errChan <- err
  338. return
  339. }
  340. // Return the connection
  341. t.returnConn(out)
  342. if resp.Err == nil {
  343. respChan <- resp.Vnodes
  344. } else {
  345. errChan <- resp.Err
  346. }
  347. }()
  348. select {
  349. case <-time.After(t.timeout):
  350. return nil, fmt.Errorf("Command timed out!")
  351. case err := <-errChan:
  352. return nil, err
  353. case res := <-respChan:
  354. return res, nil
  355. }
  356. }
  357. // Find a successor
  358. func (t *TCPTransport) FindSuccessors(vn *chord.Vnode, n int, k []byte) ([]*chord.Vnode, error) {
  359. // Get a conn
  360. out, err := t.getConn(vn.Host)
  361. if err != nil {
  362. return nil, err
  363. }
  364. respChan := make(chan []*chord.Vnode, 1)
  365. errChan := make(chan error, 1)
  366. go func() {
  367. // Send a list command
  368. out.header.ReqType = tcpFindSucReq
  369. body := tcpBodyFindSuc{Target: vn, Num: n, Key: k}
  370. if err := out.enc.Encode(&out.header); err != nil {
  371. errChan <- err
  372. return
  373. }
  374. if err := out.enc.Encode(&body); err != nil {
  375. errChan <- err
  376. return
  377. }
  378. // Read in the response
  379. resp := tcpBodyVnodeListError{}
  380. if err := out.dec.Decode(&resp); err != nil {
  381. errChan <- err
  382. return
  383. }
  384. // Return the connection
  385. t.returnConn(out)
  386. if resp.Err == nil {
  387. respChan <- resp.Vnodes
  388. } else {
  389. errChan <- resp.Err
  390. }
  391. }()
  392. select {
  393. case <-time.After(t.timeout):
  394. return nil, fmt.Errorf("Command timed out!")
  395. case err := <-errChan:
  396. return nil, err
  397. case res := <-respChan:
  398. return res, nil
  399. }
  400. }
  401. // Clears a predecessor if it matches a given vnode. Used to leave.
  402. func (t *TCPTransport) ClearPredecessor(target, self *chord.Vnode) error {
  403. // Get a conn
  404. out, err := t.getConn(target.Host)
  405. if err != nil {
  406. return err
  407. }
  408. respChan := make(chan bool, 1)
  409. errChan := make(chan error, 1)
  410. go func() {
  411. // Send a list command
  412. out.header.ReqType = tcpClearPredReq
  413. body := tcpBodyTwoVnode{Target: target, Vn: self}
  414. if err := out.enc.Encode(&out.header); err != nil {
  415. errChan <- err
  416. return
  417. }
  418. if err := out.enc.Encode(&body); err != nil {
  419. errChan <- err
  420. return
  421. }
  422. // Read in the response
  423. resp := tcpBodyError{}
  424. if err := out.dec.Decode(&resp); err != nil {
  425. errChan <- err
  426. return
  427. }
  428. // Return the connection
  429. t.returnConn(out)
  430. if resp.Err == nil {
  431. respChan <- true
  432. } else {
  433. errChan <- resp.Err
  434. }
  435. }()
  436. select {
  437. case <-time.After(t.timeout):
  438. return fmt.Errorf("Command timed out!")
  439. case err := <-errChan:
  440. return err
  441. case <-respChan:
  442. return nil
  443. }
  444. }
  445. // Instructs a node to skip a given successor. Used to leave.
  446. func (t *TCPTransport) SkipSuccessor(target, self *chord.Vnode) error {
  447. // Get a conn
  448. out, err := t.getConn(target.Host)
  449. if err != nil {
  450. return err
  451. }
  452. respChan := make(chan bool, 1)
  453. errChan := make(chan error, 1)
  454. go func() {
  455. // Send a list command
  456. out.header.ReqType = tcpSkipSucReq
  457. body := tcpBodyTwoVnode{Target: target, Vn: self}
  458. if err := out.enc.Encode(&out.header); err != nil {
  459. errChan <- err
  460. return
  461. }
  462. if err := out.enc.Encode(&body); err != nil {
  463. errChan <- err
  464. return
  465. }
  466. // Read in the response
  467. resp := tcpBodyError{}
  468. if err := out.dec.Decode(&resp); err != nil {
  469. errChan <- err
  470. return
  471. }
  472. // Return the connection
  473. t.returnConn(out)
  474. if resp.Err == nil {
  475. respChan <- true
  476. } else {
  477. errChan <- resp.Err
  478. }
  479. }()
  480. select {
  481. case <-time.After(t.timeout):
  482. return fmt.Errorf("Command timed out!")
  483. case err := <-errChan:
  484. return err
  485. case <-respChan:
  486. return nil
  487. }
  488. }
  489. // Register for an RPC callbacks
  490. func (t *TCPTransport) Register(v *chord.Vnode, o chord.VnodeRPC) {
  491. key := v.String()
  492. t.lock.Lock()
  493. t.local[key] = &chord.LocalRPC{v, o}
  494. t.lock.Unlock()
  495. }
  496. // Shutdown the TCP transport
  497. func (t *TCPTransport) Shutdown() {
  498. atomic.StoreInt32(&t.shutdown, 1)
  499. t.sock.Close()
  500. // Close all the inbound connections
  501. t.lock.RLock()
  502. for conn := range t.inbound {
  503. conn.Close()
  504. }
  505. t.lock.RUnlock()
  506. // Close all the outbound
  507. t.poolLock.Lock()
  508. for _, conns := range t.pool {
  509. for _, out := range conns {
  510. out.sock.Close()
  511. }
  512. }
  513. t.pool = nil
  514. t.poolLock.Unlock()
  515. }
  516. // Closes old outbound connections
  517. func (t *TCPTransport) reapOld() {
  518. for {
  519. if atomic.LoadInt32(&t.shutdown) == 1 {
  520. return
  521. }
  522. time.Sleep(30 * time.Second)
  523. t.reapOnce()
  524. }
  525. }
  526. func (t *TCPTransport) reapOnce() {
  527. t.poolLock.Lock()
  528. defer t.poolLock.Unlock()
  529. for host, conns := range t.pool {
  530. max := len(conns)
  531. for i := 0; i < max; i++ {
  532. if time.Since(conns[i].used) > t.maxIdle {
  533. conns[i].sock.Close()
  534. conns[i], conns[max-1] = conns[max-1], nil
  535. max--
  536. i--
  537. }
  538. }
  539. // Trim any idle conns
  540. t.pool[host] = conns[:max]
  541. }
  542. }
  543. // Listens for inbound connections
  544. func (t *TCPTransport) listen() {
  545. for {
  546. conn, err := t.sock.AcceptTCP()
  547. if err != nil {
  548. if atomic.LoadInt32(&t.shutdown) == 0 {
  549. fmt.Printf("[ERR] Error accepting TCP connection! %s", err)
  550. continue
  551. } else {
  552. return
  553. }
  554. }
  555. // Setup the conn
  556. t.setupConn(conn)
  557. // Register the inbound conn
  558. t.lock.Lock()
  559. t.inbound[conn] = struct{}{}
  560. t.lock.Unlock()
  561. // Start handler
  562. go t.handleConn(conn)
  563. }
  564. }
  565. // Handles inbound TCP connections
  566. func (t *TCPTransport) handleConn(conn *net.TCPConn) {
  567. // Defer the cleanup
  568. defer func() {
  569. t.lock.Lock()
  570. delete(t.inbound, conn)
  571. t.lock.Unlock()
  572. conn.Close()
  573. }()
  574. dec := gob.NewDecoder(conn)
  575. enc := gob.NewEncoder(conn)
  576. header := tcpHeader{}
  577. var sendResp interface{}
  578. for {
  579. // Get the header
  580. if err := dec.Decode(&header); err != nil {
  581. if atomic.LoadInt32(&t.shutdown) == 0 && err.Error() != "EOF" {
  582. log.Printf("[ERR] Failed to decode TCP header! Got %s", err)
  583. }
  584. return
  585. }
  586. // Read in the body and process request
  587. switch header.ReqType {
  588. case tcpPing:
  589. body := tcpBodyVnode{}
  590. if err := dec.Decode(&body); err != nil {
  591. log.Printf("[ERR] Failed to decode TCP body! Got %s", err)
  592. return
  593. }
  594. // Generate a response
  595. _, ok := t.get(body.Vn)
  596. if ok {
  597. sendResp = tcpBodyBoolError{B: ok, Err: nil}
  598. } else {
  599. sendResp = tcpBodyBoolError{B: ok, Err: fmt.Errorf("Target VN not found! Target %s:%s",
  600. body.Vn.Host, body.Vn.String())}
  601. }
  602. case tcpListReq:
  603. body := tcpBodyString{}
  604. if err := dec.Decode(&body); err != nil {
  605. log.Printf("[ERR] Failed to decode TCP body! Got %s", err)
  606. return
  607. }
  608. // Generate all the local clients
  609. res := make([]*chord.Vnode, 0, len(t.local))
  610. // Build list
  611. t.lock.RLock()
  612. for _, v := range t.local {
  613. res = append(res, v.Vnode)
  614. }
  615. t.lock.RUnlock()
  616. // Make response
  617. sendResp = tcpBodyVnodeListError{Vnodes: trimSlice(res)}
  618. case tcpGetPredReq:
  619. body := tcpBodyVnode{}
  620. if err := dec.Decode(&body); err != nil {
  621. log.Printf("[ERR] Failed to decode TCP body! Got %s", err)
  622. return
  623. }
  624. // Generate a response
  625. obj, ok := t.get(body.Vn)
  626. resp := tcpBodyVnodeError{}
  627. sendResp = &resp
  628. if ok {
  629. node, err := obj.GetPredecessor()
  630. resp.Vnode = node
  631. resp.Err = err
  632. } else {
  633. resp.Err = fmt.Errorf("Target VN not found! Target %s:%s",
  634. body.Vn.Host, body.Vn.String())
  635. }
  636. case tcpNotifyReq:
  637. body := tcpBodyTwoVnode{}
  638. if err := dec.Decode(&body); err != nil {
  639. log.Printf("[ERR] Failed to decode TCP body! Got %s", err)
  640. return
  641. }
  642. if body.Target == nil {
  643. return
  644. }
  645. // Generate a response
  646. obj, ok := t.get(body.Target)
  647. resp := tcpBodyVnodeListError{}
  648. sendResp = &resp
  649. if ok {
  650. nodes, err := obj.Notify(body.Vn)
  651. resp.Vnodes = trimSlice(nodes)
  652. resp.Err = err
  653. } else {
  654. resp.Err = fmt.Errorf("Target VN not found! Target %s:%s",
  655. body.Target.Host, body.Target.String())
  656. }
  657. case tcpFindSucReq:
  658. body := tcpBodyFindSuc{}
  659. if err := dec.Decode(&body); err != nil {
  660. log.Printf("[ERR] Failed to decode TCP body! Got %s", err)
  661. return
  662. }
  663. // Generate a response
  664. obj, ok := t.get(body.Target)
  665. resp := tcpBodyVnodeListError{}
  666. sendResp = &resp
  667. if ok {
  668. nodes, err := obj.FindSuccessors(body.Num, body.Key)
  669. resp.Vnodes = trimSlice(nodes)
  670. resp.Err = err
  671. } else {
  672. resp.Err = fmt.Errorf("Target VN not found! Target %s:%s",
  673. body.Target.Host, body.Target.String())
  674. }
  675. case tcpClearPredReq:
  676. body := tcpBodyTwoVnode{}
  677. if err := dec.Decode(&body); err != nil {
  678. log.Printf("[ERR] Failed to decode TCP body! Got %s", err)
  679. return
  680. }
  681. // Generate a response
  682. obj, ok := t.get(body.Target)
  683. resp := tcpBodyError{}
  684. sendResp = &resp
  685. if ok {
  686. resp.Err = obj.ClearPredecessor(body.Vn)
  687. } else {
  688. resp.Err = fmt.Errorf("Target VN not found! Target %s:%s",
  689. body.Target.Host, body.Target.String())
  690. }
  691. case tcpSkipSucReq:
  692. body := tcpBodyTwoVnode{}
  693. if err := dec.Decode(&body); err != nil {
  694. log.Printf("[ERR] Failed to decode TCP body! Got %s", err)
  695. return
  696. }
  697. // Generate a response
  698. obj, ok := t.get(body.Target)
  699. resp := tcpBodyError{}
  700. sendResp = &resp
  701. if ok {
  702. resp.Err = obj.SkipSuccessor(body.Vn)
  703. } else {
  704. resp.Err = fmt.Errorf("Target VN not found! Target %s:%s",
  705. body.Target.Host, body.Target.String())
  706. }
  707. default:
  708. log.Printf("[ERR] Unknown request type! Got %d", header.ReqType)
  709. return
  710. }
  711. // Send the response
  712. if err := enc.Encode(sendResp); err != nil {
  713. log.Printf("[ERR] Failed to send TCP body! Got %s", err)
  714. return
  715. }
  716. }
  717. }
  718. // Trims the slice to remove nil elements
  719. func trimSlice(vn []*chord.Vnode) []*chord.Vnode {
  720. if vn == nil {
  721. return vn
  722. }
  723. // Find a non-nil index
  724. idx := len(vn) - 1
  725. for vn[idx] == nil {
  726. idx--
  727. }
  728. return vn[:idx+1]
  729. }