peerstream updated (data race)
Juan Batiz-Benet committed
Dec 18, 2014 at 14:04 UTC
8b4c7cf34bce8fc91662ce96443bfac3659ce9fa
4 files changed
+84
-3
Godeps/Godeps.json
+1
-1
@@ -141,7 +141,7 @@
141
},
142
{
143
"ImportPath": "github.com/jbenet/go-peerstream",
144
- "Rev": "5023d0d6b3efeb50c2c30535d011bdcb2351e212"
144
+ "Rev": "1c71a3e04eeef9297a12ecdff75a0b28ffa8bf90"
145
},
146
{
147
"ImportPath": "github.com/jbenet/go-random",
Godeps/_workspace/src/github.com/jbenet/go-peerstream/conn.go
-2
@@ -159,7 +159,6 @@ func (s *Swarm) addConn(netConn net.Conn, server bool) (*Conn, error) {
159
// first, check if we already have it...
160
for c := range s.conns {
161
if c.netConn == netConn {
162
- s.connLock.Unlock()
162
return c, nil
163
}
164
}
@@ -167,7 +166,6 @@ func (s *Swarm) addConn(netConn net.Conn, server bool) (*Conn, error) {
166
// create a new spdystream connection
167
ssConn, err := ss.NewConnection(netConn, server)
168
if err != nil {
170
- s.connLock.Unlock()
169
return nil, err
170
}
171
Godeps/_workspace/src/github.com/jbenet/go-peerstream/example/blockhandler/blockhandler.go
new
+77
@@ -0,0 +1,77 @@
1
+package main
2
+
3
+import (
4
+ "bufio"
5
+ "fmt"
6
+ "net"
7
+ "os"
8
+ "time"
9
+
10
+ ps "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-peerstream"
11
+)
12
+
13
+func die(err error) {
14
+ fmt.Fprintf(os.Stderr, "error: %s\n")
15
+ os.Exit(1)
16
+}
17
+
18
+func main() {
19
+ // create a new Swarm
20
+ swarm := ps.NewSwarm()
21
+ defer swarm.Close()
22
+
23
+ // tell swarm what to do with a new incoming streams.
24
+ // EchoHandler just echos back anything they write.
25
+ swarm.SetStreamHandler(ps.EchoHandler)
26
+
27
+ l, err := net.Listen("tcp", "localhost:8001")
28
+ if err != nil {
29
+ die(err)
30
+ }
31
+
32
+ if _, err := swarm.AddListener(l); err != nil {
33
+ die(err)
34
+ }
35
+
36
+ nc, err := net.Dial("tcp", "localhost:8001")
37
+ if err != nil {
38
+ die(err)
39
+ }
40
+
41
+ c, err := swarm.AddConn(nc)
42
+ if err != nil {
43
+ die(err)
44
+ }
45
+
46
+ nRcvStream := 0
47
+ bio := bufio.NewReader(os.Stdin)
48
+ swarm.SetStreamHandler(func(s *ps.Stream) {
49
+ log("handling new stream %d", nRcvStream)
50
+ nRcvStream++
51
+
52
+ line, err := bio.ReadString('\n')
53
+ if err != nil {
54
+ die(err)
55
+ }
56
+ _ = line
57
+ // line = "read: " + line
58
+ // s.Write([]byte(line))
59
+ s.Close()
60
+ })
61
+
62
+ nSndStream := 0
63
+ for {
64
+ <-time.After(200 * time.Millisecond)
65
+ s, err := swarm.NewStreamWithConn(c)
66
+ if err != nil {
67
+ die(err)
68
+ }
69
+ log("sender got new stream %d", nSndStream)
70
+ nSndStream++
71
+ s.Wait()
72
+ }
73
+}
74
+
75
+func log(s string, ifs ...interface{}) {
76
+ fmt.Fprintf(os.Stderr, s+"\n", ifs...)
77
+}
Godeps/_workspace/src/github.com/jbenet/go-peerstream/swarm.go
+6
@@ -110,28 +110,34 @@ func (s *Swarm) SelectConn() SelectConn {
110
111
// Conns returns all the connections associated with this Swarm.
112
func (s *Swarm) Conns() []*Conn {
113
+ s.connLock.RLock()
114
conns := make([]*Conn, 0, len(s.conns))
115
for c := range s.conns {
116
conns = append(conns, c)
117
}
118
+ s.connLock.RUnlock()
119
return conns
120
}
121
122
// Listeners returns all the listeners associated with this Swarm.
123
func (s *Swarm) Listeners() []*Listener {
124
+ s.listenerLock.RLock()
125
out := make([]*Listener, 0, len(s.listeners))
126
for c := range s.listeners {
127
out = append(out, c)
128
}
129
+ s.listenerLock.RUnlock()
130
return out
131
}
132
133
// Streams returns all the streams associated with this Swarm.
134
func (s *Swarm) Streams() []*Stream {
135
+ s.streamLock.RLock()
136
out := make([]*Stream, 0, len(s.streams))
137
for c := range s.streams {
138
out = append(out, c)
139
}
140
+ s.streamLock.RUnlock()
141
return out
142
}
143