@cryptotaxi247 / kubo / commits / 908ff837f

updated peerstream (race)

Juan Batiz-Benet committed Jan 20, 2015 at 11:25 UTC 908ff837fd4b56c99f56f9ed05028775ca585a5e
8 files changed +51 -123
Godeps/Godeps.json
+1 -1
@@ -152,7 +152,7 @@
152 },
153 {
154 "ImportPath": "github.com/jbenet/go-peerstream",
155 - "Rev": "ccc044c2a5999f36743881ff73568660a581f2f2"
155 + "Rev": "530b09b2300da11cc19f479289be5d014c146581"
156 },
157 {
158 "ImportPath": "github.com/jbenet/go-random",
Godeps/_workspace/src/github.com/jbenet/go-peerstream/.travis.yml new
+11
@@ -0,0 +1,11 @@
1 +language: go
2 +
3 +go:
4 + - 1.2
5 + - 1.3
6 + - 1.4
7 + - release
8 + - tip
9 +
10 +script:
11 + - go test -race -cpu=5 -v ./...
Godeps/_workspace/src/github.com/jbenet/go-peerstream/example/closer/closer
Binary files a/Godeps/_workspace/src/github.com/jbenet/go-peerstream/example/closer/closer and /dev/null differ
Godeps/_workspace/src/github.com/jbenet/go-peerstream/example/closer/closer.go deleted
-90
@@ -1,90 +0,0 @@
1 -package main
2 -
3 -import (
4 - "fmt"
5 - "net"
6 - "os"
7 - "time"
8 -
9 - ps "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-peerstream"
10 -)
11 -
12 -func die(err error) {
13 - fmt.Fprintf(os.Stderr, "error: %s\n")
14 - os.Exit(1)
15 -}
16 -
17 -func main() {
18 - // create a new Swarm
19 - swarm := ps.NewSwarm()
20 - defer swarm.Close()
21 -
22 - // tell swarm what to do with a new incoming streams.
23 - // EchoHandler just echos back anything they write.
24 - swarm.SetStreamHandler(ps.EchoHandler)
25 -
26 - l, err := net.Listen("tcp", "localhost:8001")
27 - if err != nil {
28 - die(err)
29 - }
30 -
31 - if _, err := swarm.AddListener(l); err != nil {
32 - die(err)
33 - }
34 -
35 - nc, err := net.Dial("tcp", "localhost:8001")
36 - if err != nil {
37 - die(err)
38 - }
39 -
40 - c, err := swarm.AddConn(nc)
41 - if err != nil {
42 - die(err)
43 - }
44 -
45 - hello := []byte("hello")
46 - goodbye := []byte("goodbye")
47 - swarm.SetStreamHandler(func(s *ps.Stream) {
48 - go func() {
49 - log("handler: got new stream.")
50 - // s.Wait()
51 - // log("handler: done waiting on new stream.")
52 - buf := make([]byte, len(hello))
53 - s.Read(buf)
54 - log("handler: read: %s", buf)
55 - s.Write(goodbye)
56 - log("handler: wrote: %s", goodbye)
57 - s.Close()
58 - log("handler: closed.")
59 - }()
60 - })
61 -
62 - for {
63 - s, err := swarm.NewStreamWithConn(c)
64 - if err != nil {
65 - die(err)
66 - }
67 - // s.Wait()
68 - log("sender: got new stream")
69 - for {
70 - <-time.After(500 * time.Millisecond)
71 - log("sender: writing hello...")
72 - if _, err := s.Write(hello); err != nil {
73 - log("sender: write error: %s", err)
74 - break
75 - }
76 - buf := make([]byte, len(goodbye))
77 - if _, err := s.Read(buf); err != nil {
78 - log("sender: read error: %s", err)
79 - break
80 - }
81 - }
82 - if err := s.Close(); err != nil {
83 - log("sender: close error: %s", err)
84 - }
85 - }
86 -}
87 -
88 -func log(s string, ifs ...interface{}) {
89 - fmt.Fprintf(os.Stderr, s+"\n", ifs...)
90 -}
Godeps/_workspace/src/github.com/jbenet/go-peerstream/swarm.go
+6 -6
@@ -325,20 +325,20 @@ func (s *Swarm) Close() error {
325 var wgl sync.WaitGroup
326 for _, l := range s.Listeners() {
327 wgl.Add(1)
328 - go func() {
329 - l.Close()
328 + go func(list *Listener) {
329 + list.Close()
330 wgl.Done()
331 - }()
331 + }(l)
332 }
333 wgl.Wait()
334
335 var wgc sync.WaitGroup
336 for _, c := range s.Conns() {
337 wgc.Add(1)
338 - go func() {
339 - c.Close()
338 + go func(conn *Conn) {
339 + conn.Close()
340 wgc.Done()
341 - }()
341 + }(c)
342 }
343 wgc.Wait()
344 return nil
Godeps/_workspace/src/github.com/jbenet/go-peerstream/transport/spdystream/spdystream.go
+22 -5
@@ -33,14 +33,31 @@ func (s *stream) Close() error {
33 }
34
35 // Conn is a connection to a remote peer.
36 -type conn ss.Connection
36 +type conn struct {
37 + sc *ss.Connection
38 +
39 + closed chan struct{}
40 +}
41
42 func (c *conn) spdyConn() *ss.Connection {
39 - return (*ss.Connection)(c)
43 + return c.sc
44 }
45
46 func (c *conn) Close() error {
43 - return c.spdyConn().Close()
47 + err := c.spdyConn().CloseWait()
48 + if !c.IsClosed() {
49 + close(c.closed)
50 + }
51 + return err
52 +}
53 +
54 +func (c *conn) IsClosed() bool {
55 + select {
56 + case <-c.closed:
57 + return true
58 + default:
59 + return false
60 + }
61 }
62
63 // OpenStream creates a new stream.
@@ -84,6 +101,6 @@ type transport struct{}
101 var Transport = transport{}
102
103 func (t transport) NewConn(nc net.Conn, isServer bool) (pst.Conn, error) {
87 - c, err := ss.NewConnection(nc, isServer)
88 - return (*conn)(c), err
104 + sc, err := ss.NewConnection(nc, isServer)
105 + return &conn{sc: sc, closed: make(chan struct{})}, err
106 }
Godeps/_workspace/src/github.com/jbenet/go-peerstream/transport/spdystream/spdystream_test.go
+1
@@ -7,5 +7,6 @@ import (
7 )
8
9 func TestSpdyStreamTransport(t *testing.T) {
10 + t.Skip("spdystream is known to be broken")
11 psttest.SubtestAll(t, Transport)
12 }
Godeps/_workspace/src/github.com/jbenet/go-peerstream/transport/test/ttest.go
+10 -21
@@ -44,11 +44,6 @@ func checkErr(t *testing.T, err error) {
44 }
45 }
46
47 -func getNextPort() int {
48 - nextPort++
49 - return nextPort
50 -}
51 -
47 func log(s string, v ...interface{}) {
48 if testing.Verbose() {
49 fmt.Fprintf(os.Stderr, "> "+s+"\n", v...)
@@ -69,17 +64,15 @@ func singleConn(t *testing.T, tr pst.Transport) echoSetup {
64 log("closing stream")
65 })
66
72 - port := getNextPort()
73 - addr := fmt.Sprintf("localhost:%d", port)
74 - log("listening at %s", addr)
75 - l, err := net.Listen("tcp", addr)
67 + log("listening at %s", "localhost:0")
68 + l, err := net.Listen("tcp", "localhost:0")
69 checkErr(t, err)
70
71 _, err = swarm.AddListener(l)
72 checkErr(t, err)
73
81 - log("dialing to %s", addr)
82 - nc1, err := net.Dial("tcp", addr)
74 + log("dialing to %s", l.Addr())
75 + nc1, err := net.Dial("tcp", l.Addr().String())
76 checkErr(t, err)
77
78 c1, err := swarm.AddConn(nc1)
@@ -101,10 +94,8 @@ func makeSwarm(t *testing.T, tr pst.Transport, nListeners int) *ps.Swarm {
94 })
95
96 for i := 0; i < nListeners; i++ {
104 - port := getNextPort()
105 - addr := fmt.Sprintf("localhost:%d", port)
106 - log("%p listening at %s", swarm, addr)
107 - l, err := net.Listen("tcp", addr)
97 + log("%p listening at %s", swarm, "localhost:0")
98 + l, err := net.Listen("tcp", "localhost:0")
99 checkErr(t, err)
100 _, err = swarm.AddListener(l)
101 checkErr(t, err)
@@ -138,17 +129,15 @@ func SubtestSimpleWrite(t *testing.T, tr pst.Transport) {
129 log("closing stream")
130 })
131
141 - port := getNextPort()
142 - addr := fmt.Sprintf("localhost:%d", port)
143 - log("listening at %s", addr)
144 - l, err := net.Listen("tcp", addr)
132 + log("listening at %s", "localhost:0")
133 + l, err := net.Listen("tcp", "localhost:0")
134 checkErr(t, err)
135
136 _, err = swarm.AddListener(l)
137 checkErr(t, err)
138
150 - log("dialing to %s", addr)
151 - nc1, err := net.Dial("tcp", addr)
139 + log("dialing to %s", l.Addr().String())
140 + nc1, err := net.Dial("tcp", l.Addr().String())
141 checkErr(t, err)
142
143 c1, err := swarm.AddConn(nc1)