@cryptotaxi247 / kubo / commits / 075009642

swarm connection using msgio

http://github.com/jbenet/msgio

Juan Batiz-Benet committed Jul 8, 2014 at 15:55 UTC 0750096421b1273f7fa4658755c1a49b749428a6
3 files changed +174
swarm/conn.go new
+72
@@ -0,0 +1,72 @@
1 +package swarm
2 +
3 +import (
4 + "fmt"
5 + "net"
6 + ma "github.com/jbenet/go-multiaddr"
7 + peer "github.com/jbenet/go-ipfs/peer"
8 + msgio "github.com/jbenet/go-msgio"
9 +)
10 +
11 +const ChanBuffer = 10
12 +
13 +type Conn struct {
14 + Peer *peer.Peer
15 + Addr *ma.Multiaddr
16 + Conn net.Conn
17 +
18 + Closed chan bool
19 + Outgoing *msgio.Chan
20 + Incoming *msgio.Chan
21 +}
22 +
23 +
24 +func Dial(network string, peer *peer.Peer) (*Conn, error) {
25 + addr := peer.NetAddress(network)
26 + if addr == nil {
27 + return nil, fmt.Errorf("No address for network %s", network)
28 + }
29 +
30 + network, host, err := addr.DialArgs()
31 + if err != nil {
32 + return nil, err
33 + }
34 +
35 + nconn, err := net.Dial(network, host)
36 + if err != nil {
37 + return nil, err
38 + }
39 +
40 + out := msgio.NewChan(10)
41 + inc := msgio.NewChan(10)
42 +
43 + conn := &Conn{
44 + Peer: peer,
45 + Addr: addr,
46 + Conn: nconn,
47 +
48 + Outgoing: out,
49 + Incoming: inc,
50 + Closed: make(chan bool, 1),
51 + }
52 +
53 + go out.WriteTo(nconn)
54 + go inc.ReadFrom(nconn, 1 << 12)
55 +
56 + return conn, nil
57 +}
58 +
59 +func (s *Conn) Close() error {
60 + if s.Conn == nil {
61 + return fmt.Errorf("Already closed.") // already closed
62 + }
63 +
64 + // closing net connection
65 + err := s.Conn.Close()
66 + s.Conn = nil
67 + // closing channels
68 + s.Incoming.Close()
69 + s.Outgoing.Close()
70 + s.Closed<- true
71 + return err
72 +}
swarm/conn_test.go new
+93
@@ -0,0 +1,93 @@
1 +package swarm
2 +
3 +import (
4 + "fmt"
5 + "net"
6 + "testing"
7 + ma "github.com/jbenet/go-multiaddr"
8 + mh "github.com/jbenet/go-multihash"
9 + peer "github.com/jbenet/go-ipfs/peer"
10 +)
11 +
12 +func setupPeer(id string, addr string) (*peer.Peer, error) {
13 + tcp, err := ma.NewMultiaddr(addr)
14 + if err != nil {
15 + return nil, err
16 + }
17 +
18 + mh, err := mh.FromHexString(id)
19 + if err != nil {
20 + return nil, err
21 + }
22 +
23 + p := &peer.Peer{Id: peer.PeerId(mh)}
24 + p.AddAddress(tcp)
25 + return p, nil
26 +}
27 +
28 +func echoListen(listener *net.TCPListener) {
29 + for {
30 + c, err := listener.Accept()
31 + if err == nil {
32 + fmt.Println("accepeted")
33 + go echo(c)
34 + }
35 + }
36 +}
37 +
38 +func echo(c net.Conn) {
39 + for {
40 + data := make([]byte, 1024)
41 + i, err := c.Read(data)
42 + if err != nil {
43 + fmt.Printf("error %v\n", err)
44 + return
45 + }
46 + _, err = c.Write(data[:i])
47 + if err != nil {
48 + fmt.Printf("error %v\n", err)
49 + return
50 + }
51 + fmt.Println("echoing", data[:i])
52 + }
53 +}
54 +
55 +func TestDial(t *testing.T) {
56 +
57 + listener, err := net.Listen("tcp", "127.0.0.1:1234")
58 + if err != nil {
59 + t.Fatal("error setting up listener", err)
60 + }
61 + go echoListen(listener.(*net.TCPListener))
62 +
63 + p, err := setupPeer("11140beec7b5ea3f0fdbc95d0dd47f3c5bc275da8a33", "/ip4/127.0.0.1/tcp/1234")
64 + if err != nil {
65 + t.Fatal("error setting up peer", err)
66 + }
67 +
68 + c, err := Dial("tcp", p)
69 + if err != nil {
70 + t.Fatal("error dialing peer", err)
71 + }
72 +
73 +
74 + fmt.Println("sending")
75 + c.Outgoing.MsgChan<- []byte("beep")
76 + c.Outgoing.MsgChan<- []byte("boop")
77 + out := <-c.Incoming.MsgChan
78 + fmt.Println("recving", string(out))
79 + if string(out) != "beep" {
80 + t.Error("unexpected conn output")
81 + }
82 +
83 + out = <-c.Incoming.MsgChan
84 + if string(out) != "boop" {
85 + t.Error("unexpected conn output")
86 + }
87 +
88 + fmt.Println("closing")
89 + c.Close()
90 + listener.Close()
91 +}
92 +
93 +
swarm/swarm.go new
+9
@@ -0,0 +1,9 @@
1 +package swarm
2 +
3 +import (
4 +)
5 +
6 +type Swarm struct {
7 + Conns map[string]*Conn
8 +}
9 +