Fixed connections all over.
Juan Batiz-Benet committed
Sep 19, 2014 at 07:51 UTC
9dd39de491551d5c9628b4ea0f767721257ba11f
12 files changed
+237
-194
crypto/spipe/handshake.go
+2
@@ -221,6 +221,7 @@ func (s *SecurePipe) handleSecureIn(hashType string, tIV, tCKey, tMKey []byte) {
221
return
222
}
223
224
+ // u.DOut("[peer %s] secure in [from = %s] %d\n", s.local.ID.Pretty(), s.remote.ID.Pretty(), len(data))
225
if len(data) <= macSize {
226
continue
227
}
@@ -268,6 +269,7 @@ func (s *SecurePipe) handleSecureOut(hashType string, mIV, mCKey, mMKey []byte)
269
copy(buff[len(data):], myMac.Sum(nil))
270
myMac.Reset()
271
272
+ // u.DOut("[peer %s] secure out [to = %s] %d\n", s.local.ID.Pretty(), s.remote.ID.Pretty(), len(buff))
273
s.insecure.Out <- buff
274
}
275
}
net/message/message.go
+21
@@ -6,11 +6,13 @@ import (
6
proto "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/goprotobuf/proto"
7
)
8
9
+// NetMessage is the interface for the message
10
type NetMessage interface {
11
Peer() *peer.Peer
12
Data() []byte
13
}
14
15
+// New is the interface for constructing a new message.
16
func New(p *peer.Peer, data []byte) NetMessage {
17
return &message{peer: p, data: data}
18
}
@@ -55,3 +57,22 @@ func NewPipe(bufsize int) *Pipe {
57
Outgoing: make(chan NetMessage, bufsize),
58
}
59
}
60
+
61
+// ConnectTo connects this pipe to another, using a context for termination.
62
+func (p *Pipe) ConnectTo(p2 *Pipe) {
63
+ connectChans(p.Outgoing, p2.Outgoing)
64
+ connectChans(p2.Incoming, p.Incoming)
65
+}
66
+
67
+func connectChans(a, b chan NetMessage) {
68
+ go func() {
69
+ for {
70
+ m, more := <-a
71
+ if !more {
72
+ close(b)
73
+ return
74
+ }
75
+ b <- m
76
+ }
77
+ }()
78
+}
net/mux/mux.go
+8
-2
@@ -94,7 +94,10 @@ func (m *Muxer) handleIncomingMessages(ctx context.Context) {
94
}
95
96
select {
97
- case msg := <-m.Incoming:
97
+ case msg, more := <-m.Incoming:
98
+ if !more {
99
+ return
100
+ }
101
go m.handleIncomingMessage(ctx, msg)
102
103
case <-ctx.Done():
@@ -132,7 +135,10 @@ func (m *Muxer) handleIncomingMessage(ctx context.Context, m1 msg.NetMessage) {
135
func (m *Muxer) handleOutgoingMessages(ctx context.Context, pid ProtocolID, proto Protocol) {
136
for {
137
select {
135
- case msg := <-proto.GetPipe().Outgoing:
138
+ case msg, more := <-proto.GetPipe().Outgoing:
139
+ if !more {
140
+ return
141
+ }
142
go m.handleOutgoingMessage(ctx, pid, msg)
143
144
case <-ctx.Done():
net/net.go
+3
@@ -53,6 +53,9 @@ func NewIpfsNetwork(ctx context.Context, local *peer.Peer,
53
return nil, err
54
}
55
56
+ // remember to wire components together.
57
+ in.muxer.Pipe.ConnectTo(in.swarm.Pipe)
58
+
59
return in, nil
60
}
61
net/service/service.go
+6
-1
@@ -82,6 +82,8 @@ func (s *Service) sendMessage(ctx context.Context, m msg.NetMessage, rid Request
82
return err
83
}
84
85
+ // u.DOut("Service send message [to = %s]\n", m.Peer().ID.Pretty())
86
+
87
// send message
88
m2 := msg.New(m.Peer(), data)
89
select {
@@ -150,7 +152,10 @@ func (s *Service) SendRequest(ctx context.Context, m msg.NetMessage) (msg.NetMes
152
func (s *Service) handleIncomingMessages(ctx context.Context) {
153
for {
154
select {
153
- case m := <-s.Incoming:
155
+ case m, more := <-s.Incoming:
156
+ if !more {
157
+ return
158
+ }
159
go s.handleIncomingMessage(ctx, m)
160
161
case <-ctx.Done():
net/swarm/conn.go
+11
-4
@@ -101,13 +101,13 @@ func (s *Swarm) connSetup(c *conn.Conn) error {
101
return errors.New("Tried to start nil connection.")
102
}
103
104
- // u.DOut("Starting connection: %s\n", c.Peer.Key().Pretty())
104
+ u.DOut("Starting connection: %s\n", c.Peer.Key().Pretty())
105
106
if err := s.connSecure(c); err != nil {
107
return fmt.Errorf("Conn securing error: %v", err)
108
}
109
110
- // u.DOut("Secured connection: %s\n", c.Peer.Key().Pretty())
110
+ u.DOut("Secured connection: %s\n", c.Peer.Key().Pretty())
111
112
// add to conns
113
s.connsLock.Lock()
@@ -139,6 +139,7 @@ func (s *Swarm) connSecure(c *conn.Conn) error {
139
return err
140
}
141
142
+ c.Secure = sp
143
return nil
144
}
145
@@ -165,8 +166,11 @@ func (s *Swarm) fanOut() {
166
continue
167
}
168
169
+ // u.DOut("[peer: %s] Sent message [to = %s]\n",
170
+ // s.local.ID.Pretty(), msg.Peer().ID.Pretty())
171
+
172
// queue it in the connection's buffer
169
- conn.Outgoing.MsgChan <- msg.Data()
173
+ conn.Secure.Out <- msg.Data()
174
}
175
}
176
}
@@ -184,13 +188,16 @@ func (s *Swarm) fanIn(c *conn.Conn) {
188
case <-c.Closed:
189
goto out
190
187
- case data, ok := <-c.Incoming.MsgChan:
191
+ case data, ok := <-c.Secure.In:
192
if !ok {
193
e := fmt.Errorf("Error retrieving from conn: %v", c.Peer.Key().Pretty())
194
s.errChan <- e
195
goto out
196
}
197
198
+ // u.DOut("[peer: %s] Received message [from = %s]\n",
199
+ // s.local.ID.Pretty(), c.Peer.ID.Pretty())
200
+
201
msg := msg.New(c.Peer, data)
202
s.Incoming <- msg
203
}
routing/dht/Message.go
+1
-1
@@ -46,7 +46,7 @@ func peersToPBPeers(peers []*peer.Peer) []*Message_Peer {
46
func (m *Message) GetClusterLevel() int {
47
level := m.GetClusterLevelRaw() - 1
48
if level < 0 {
49
- u.PErr("handleGetValue: no routing level specified, assuming 0\n")
49
+ u.PErr("GetClusterLevel: no routing level specified, assuming 0\n")
50
level = 0
51
}
52
return int(level)
routing/dht/dht.go
+16
-6
@@ -125,7 +125,7 @@ func (dht *IpfsDHT) HandleMessage(ctx context.Context, mes msg.NetMessage) (msg.
125
dht.Update(mPeer)
126
127
// Print out diagnostic
128
- u.DOut("[peer: %s]\nGot message type: '%s' [from = %s]\n",
128
+ u.DOut("[peer: %s] Got message type: '%s' [from = %s]\n",
129
dht.self.ID.Pretty(),
130
Message_MessageType_name[int32(pmes.GetType())], mPeer.ID.Pretty())
131
@@ -141,6 +141,11 @@ func (dht *IpfsDHT) HandleMessage(ctx context.Context, mes msg.NetMessage) (msg.
141
return nil, err
142
}
143
144
+ // if nil response, return it before serializing
145
+ if rpmes == nil {
146
+ return nil, nil
147
+ }
148
+
149
// serialize response msg
150
rmes, err := msg.FromObject(mPeer, rpmes)
151
if err != nil {
@@ -161,6 +166,11 @@ func (dht *IpfsDHT) sendRequest(ctx context.Context, p *peer.Peer, pmes *Message
166
167
start := time.Now()
168
169
+ // Print out diagnostic
170
+ u.DOut("[peer: %s] Sent message type: '%s' [to = %s]\n",
171
+ dht.self.ID.Pretty(),
172
+ Message_MessageType_name[int32(pmes.GetType())], p.ID.Pretty())
173
+
174
rmes, err := dht.sender.SendRequest(ctx, mes)
175
if err != nil {
176
return nil, err
@@ -209,10 +219,10 @@ func (dht *IpfsDHT) getValueOrPeers(ctx context.Context, p *peer.Peer,
219
return nil, nil, err
220
}
221
212
- u.POut("pmes.GetValue() %v\n", pmes.GetValue())
222
+ u.DOut("pmes.GetValue() %v\n", pmes.GetValue())
223
if value := pmes.GetValue(); value != nil {
224
// Success! We were given the value
215
- u.POut("getValueOrPeers: got value\n")
225
+ u.DOut("getValueOrPeers: got value\n")
226
return value, nil, nil
227
}
228
@@ -222,7 +232,7 @@ func (dht *IpfsDHT) getValueOrPeers(ctx context.Context, p *peer.Peer,
232
if err != nil {
233
return nil, nil, err
234
}
225
- u.POut("getValueOrPeers: get from providers\n")
235
+ u.DOut("getValueOrPeers: get from providers\n")
236
return val, nil, nil
237
}
238
@@ -250,11 +260,11 @@ func (dht *IpfsDHT) getValueOrPeers(ctx context.Context, p *peer.Peer,
260
}
261
262
if len(peers) > 0 {
253
- u.POut("getValueOrPeers: peers\n")
263
+ u.DOut("getValueOrPeers: peers\n")
264
return nil, peers, nil
265
}
266
257
- u.POut("getValueOrPeers: u.ErrNotFound\n")
267
+ u.DOut("getValueOrPeers: u.ErrNotFound\n")
268
return nil, nil, u.ErrNotFound
269
}
270
routing/dht/dht_test.go
+158
-172
@@ -1,194 +1,180 @@
1
package dht
2
3
-// import (
4
-// "testing"
5
-//
6
-// context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
7
-//
8
-// ds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/datastore.go"
9
-// ma "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr"
10
-//
11
-// ci "github.com/jbenet/go-ipfs/crypto"
12
-// spipe "github.com/jbenet/go-ipfs/crypto/spipe"
13
-// inet "github.com/jbenet/go-ipfs/net"
14
-// mux "github.com/jbenet/go-ipfs/net/mux"
15
-// netservice "github.com/jbenet/go-ipfs/net/service"
16
-// peer "github.com/jbenet/go-ipfs/peer"
17
-// u "github.com/jbenet/go-ipfs/util"
18
-//
19
-// "bytes"
20
-// "fmt"
21
-// "time"
22
-// )
23
-//
24
-// func setupDHT(t *testing.T, p *peer.Peer) *IpfsDHT {
25
-// ctx := context.TODO()
26
-//
27
-// peerstore := peer.NewPeerstore()
28
-//
29
-// ctx, _ = context.WithCancel(ctx)
30
-// dhts := netservice.NewService(nil) // nil handler for now, need to patch it
31
-// if err := dhts.Start(ctx); err != nil {
32
-// t.Fatal(err)
33
-// }
34
-//
35
-// net, err := inet.NewIpfsNetwork(context.TODO(), p, &mux.ProtocolMap{
36
-// mux.ProtocolID_Routing: dhts,
37
-// })
38
-// if err != nil {
39
-// t.Fatal(err)
40
-// }
41
-//
42
-// d := NewDHT(p, peerstore, net, dhts, ds.NewMapDatastore())
43
-// dhts.Handler = d
44
-// return d
45
-// }
46
-//
47
-// func setupDHTS(n int, t *testing.T) ([]*ma.Multiaddr, []*peer.Peer, []*IpfsDHT) {
48
-// var addrs []*ma.Multiaddr
49
-// for i := 0; i < 4; i++ {
50
-// a, err := ma.NewMultiaddr(fmt.Sprintf("/ip4/127.0.0.1/tcp/%d", 5000+i))
51
-// if err != nil {
52
-// t.Fatal(err)
53
-// }
54
-// addrs = append(addrs, a)
55
-// }
56
-//
57
-// var peers []*peer.Peer
58
-// for i := 0; i < 4; i++ {
59
-// p := new(peer.Peer)
60
-// p.AddAddress(addrs[i])
61
-// sk, pk, err := ci.GenerateKeyPair(ci.RSA, 512)
62
-// if err != nil {
63
-// panic(err)
64
-// }
65
-// p.PubKey = pk
66
-// p.PrivKey = sk
67
-// id, err := spipe.IDFromPubKey(pk)
68
-// if err != nil {
69
-// panic(err)
70
-// }
71
-// p.ID = id
72
-// peers = append(peers, p)
73
-// }
74
-//
75
-// var dhts []*IpfsDHT
76
-// for i := 0; i < 4; i++ {
77
-// dhts[i] = setupDHT(t, peers[i])
78
-// }
79
-//
80
-// return addrs, peers, dhts
81
-// }
82
-//
83
-// func makePeer(addr *ma.Multiaddr) *peer.Peer {
84
-// p := new(peer.Peer)
85
-// p.AddAddress(addr)
86
-// sk, pk, err := ci.GenerateKeyPair(ci.RSA, 512)
87
-// if err != nil {
88
-// panic(err)
89
-// }
90
-// p.PrivKey = sk
91
-// p.PubKey = pk
92
-// id, err := spipe.IDFromPubKey(pk)
93
-// if err != nil {
94
-// panic(err)
95
-// }
96
-//
97
-// p.ID = id
98
-// return p
99
-// }
100
-//
101
-// func TestPing(t *testing.T) {
102
-// u.Debug = true
103
-// addrA, err := ma.NewMultiaddr("/ip4/127.0.0.1/tcp/2222")
104
-// if err != nil {
105
-// t.Fatal(err)
106
-// }
107
-// addrB, err := ma.NewMultiaddr("/ip4/127.0.0.1/tcp/5678")
108
-// if err != nil {
109
-// t.Fatal(err)
110
-// }
111
-//
112
-// peerA := makePeer(addrA)
113
-// peerB := makePeer(addrB)
114
-//
115
-// dhtA := setupDHT(t, peerA)
116
-// dhtB := setupDHT(t, peerB)
117
-//
118
-// defer dhtA.Halt()
119
-// defer dhtB.Halt()
120
-//
121
-// _, err = dhtA.Connect(addrB)
122
-// if err != nil {
123
-// t.Fatal(err)
124
-// }
125
-//
126
-// //Test that we can ping the node
127
-// err = dhtA.Ping(peerB, time.Second*2)
128
-// if err != nil {
129
-// t.Fatal(err)
130
-// }
131
-// }
132
-//
133
-// func TestValueGetSet(t *testing.T) {
134
-// u.Debug = false
135
-// addrA, err := ma.NewMultiaddr("/ip4/127.0.0.1/tcp/1235")
136
-// if err != nil {
137
-// t.Fatal(err)
138
-// }
139
-// addrB, err := ma.NewMultiaddr("/ip4/127.0.0.1/tcp/5679")
140
-// if err != nil {
141
-// t.Fatal(err)
142
-// }
143
-//
144
-// peerA := makePeer(addrA)
145
-// peerB := makePeer(addrB)
146
-//
147
-// dhtA := setupDHT(t, peerA)
148
-// dhtB := setupDHT(t, peerB)
149
-//
150
-// defer dhtA.Halt()
151
-// defer dhtB.Halt()
152
-//
153
-// _, err = dhtA.Connect(addrB)
154
-// if err != nil {
155
-// t.Fatal(err)
156
-// }
157
-//
158
-// dhtA.PutValue("hello", []byte("world"))
159
-//
160
-// val, err := dhtA.GetValue("hello", time.Second*2)
161
-// if err != nil {
162
-// t.Fatal(err)
163
-// }
164
-//
165
-// if string(val) != "world" {
166
-// t.Fatalf("Expected 'world' got '%s'", string(val))
167
-// }
168
-//
169
-// }
170
-//
3
+import (
4
+ "testing"
5
+
6
+ context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
7
+
8
+ ds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/datastore.go"
9
+ ma "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr"
10
+
11
+ ci "github.com/jbenet/go-ipfs/crypto"
12
+ spipe "github.com/jbenet/go-ipfs/crypto/spipe"
13
+ inet "github.com/jbenet/go-ipfs/net"
14
+ mux "github.com/jbenet/go-ipfs/net/mux"
15
+ netservice "github.com/jbenet/go-ipfs/net/service"
16
+ peer "github.com/jbenet/go-ipfs/peer"
17
+ u "github.com/jbenet/go-ipfs/util"
18
+
19
+ "fmt"
20
+ "time"
21
+)
22
+
23
+func setupDHT(t *testing.T, p *peer.Peer) *IpfsDHT {
24
+ ctx, _ := context.WithCancel(context.TODO())
25
+
26
+ peerstore := peer.NewPeerstore()
27
+
28
+ dhts := netservice.NewService(nil) // nil handler for now, need to patch it
29
+ if err := dhts.Start(ctx); err != nil {
30
+ t.Fatal(err)
31
+ }
32
+
33
+ net, err := inet.NewIpfsNetwork(ctx, p, &mux.ProtocolMap{
34
+ mux.ProtocolID_Routing: dhts,
35
+ })
36
+ if err != nil {
37
+ t.Fatal(err)
38
+ }
39
+
40
+ d := NewDHT(p, peerstore, net, dhts, ds.NewMapDatastore())
41
+ dhts.Handler = d
42
+ return d
43
+}
44
+
45
+func setupDHTS(n int, t *testing.T) ([]*ma.Multiaddr, []*peer.Peer, []*IpfsDHT) {
46
+ var addrs []*ma.Multiaddr
47
+ for i := 0; i < 4; i++ {
48
+ a, err := ma.NewMultiaddr(fmt.Sprintf("/ip4/127.0.0.1/tcp/%d", 5000+i))
49
+ if err != nil {
50
+ t.Fatal(err)
51
+ }
52
+ addrs = append(addrs, a)
53
+ }
54
+
55
+ var peers []*peer.Peer
56
+ for i := 0; i < 4; i++ {
57
+ p := makePeer(addrs[i])
58
+ peers = append(peers, p)
59
+ }
60
+
61
+ var dhts []*IpfsDHT
62
+ for i := 0; i < 4; i++ {
63
+ dhts[i] = setupDHT(t, peers[i])
64
+ }
65
+
66
+ return addrs, peers, dhts
67
+}
68
+
69
+func makePeer(addr *ma.Multiaddr) *peer.Peer {
70
+ p := new(peer.Peer)
71
+ p.AddAddress(addr)
72
+ sk, pk, err := ci.GenerateKeyPair(ci.RSA, 512)
73
+ if err != nil {
74
+ panic(err)
75
+ }
76
+ p.PrivKey = sk
77
+ p.PubKey = pk
78
+ id, err := spipe.IDFromPubKey(pk)
79
+ if err != nil {
80
+ panic(err)
81
+ }
82
+
83
+ p.ID = id
84
+ return p
85
+}
86
+
87
+func TestPing(t *testing.T) {
88
+ u.Debug = true
89
+ addrA, err := ma.NewMultiaddr("/ip4/127.0.0.1/tcp/2222")
90
+ if err != nil {
91
+ t.Fatal(err)
92
+ }
93
+ addrB, err := ma.NewMultiaddr("/ip4/127.0.0.1/tcp/5678")
94
+ if err != nil {
95
+ t.Fatal(err)
96
+ }
97
+
98
+ peerA := makePeer(addrA)
99
+ peerB := makePeer(addrB)
100
+
101
+ dhtA := setupDHT(t, peerA)
102
+ dhtB := setupDHT(t, peerB)
103
+
104
+ defer dhtA.Halt()
105
+ defer dhtB.Halt()
106
+
107
+ _, err = dhtA.Connect(peerB)
108
+ if err != nil {
109
+ t.Fatal(err)
110
+ }
111
+
112
+ //Test that we can ping the node
113
+ err = dhtA.Ping(peerB, time.Second*2)
114
+ if err != nil {
115
+ t.Fatal(err)
116
+ }
117
+}
118
+
119
+func TestValueGetSet(t *testing.T) {
120
+ u.Debug = false
121
+ addrA, err := ma.NewMultiaddr("/ip4/127.0.0.1/tcp/1235")
122
+ if err != nil {
123
+ t.Fatal(err)
124
+ }
125
+ addrB, err := ma.NewMultiaddr("/ip4/127.0.0.1/tcp/5679")
126
+ if err != nil {
127
+ t.Fatal(err)
128
+ }
129
+
130
+ peerA := makePeer(addrA)
131
+ peerB := makePeer(addrB)
132
+
133
+ dhtA := setupDHT(t, peerA)
134
+ dhtB := setupDHT(t, peerB)
135
+
136
+ defer dhtA.Halt()
137
+ defer dhtB.Halt()
138
+
139
+ _, err = dhtA.Connect(peerB)
140
+ if err != nil {
141
+ t.Fatal(err)
142
+ }
143
+
144
+ dhtA.PutValue("hello", []byte("world"))
145
+
146
+ val, err := dhtA.GetValue("hello", time.Second*2)
147
+ if err != nil {
148
+ t.Fatal(err)
149
+ }
150
+
151
+ if string(val) != "world" {
152
+ t.Fatalf("Expected 'world' got '%s'", string(val))
153
+ }
154
+
155
+}
156
+
157
// func TestProvides(t *testing.T) {
158
// u.Debug = false
159
//
174
-// addrs, _, dhts := setupDHTS(4, t)
160
+// _, peers, dhts := setupDHTS(4, t)
161
// defer func() {
162
// for i := 0; i < 4; i++ {
163
// dhts[i].Halt()
164
// }
165
// }()
166
//
181
-// _, err := dhts[0].Connect(addrs[1])
167
+// _, err := dhts[0].Connect(peers[1])
168
// if err != nil {
169
// t.Fatal(err)
170
// }
171
//
186
-// _, err = dhts[1].Connect(addrs[2])
172
+// _, err = dhts[1].Connect(peers[2])
173
// if err != nil {
174
// t.Fatal(err)
175
// }
176
//
191
-// _, err = dhts[1].Connect(addrs[3])
177
+// _, err = dhts[1].Connect(peers[3])
178
// if err != nil {
179
// t.Fatal(err)
180
// }
routing/dht/handlers.go
+1
@@ -97,6 +97,7 @@ func (dht *IpfsDHT) handlePutValue(p *peer.Peer, pmes *Message) (*Message, error
97
defer dht.dslock.Unlock()
98
dskey := ds.NewKey(pmes.GetKey())
99
err := dht.datastore.Put(dskey, pmes.GetValue())
100
+ u.DOut("[%s] handlePutValue %v %v", dht.self.ID.Pretty(), dskey, pmes.GetValue())
101
return nil, err
102
}
103
routing/dht/query.go
+8
-8
@@ -163,7 +163,7 @@ func (r *dhtQueryRunner) addPeerToQuery(next *peer.Peer, benchmark *peer.Peer) {
163
r.peersSeen[next.Key()] = next
164
r.Unlock()
165
166
- u.POut("adding peer to query: %v\n", next.ID.Pretty())
166
+ u.DOut("adding peer to query: %v\n", next.ID.Pretty())
167
168
// do this after unlocking to prevent possible deadlocks.
169
r.peersRemaining.Increment(1)
@@ -187,14 +187,14 @@ func (r *dhtQueryRunner) spawnWorkers() {
187
if !more {
188
return // channel closed.
189
}
190
- u.POut("spawning worker for: %v\n", p.ID.Pretty())
190
+ u.DOut("spawning worker for: %v\n", p.ID.Pretty())
191
go r.queryPeer(p)
192
}
193
}
194
}
195
196
func (r *dhtQueryRunner) queryPeer(p *peer.Peer) {
197
- u.POut("spawned worker for: %v\n", p.ID.Pretty())
197
+ u.DOut("spawned worker for: %v\n", p.ID.Pretty())
198
199
// make sure we rate limit concurrency.
200
select {
@@ -204,33 +204,33 @@ func (r *dhtQueryRunner) queryPeer(p *peer.Peer) {
204
return
205
}
206
207
- u.POut("running worker for: %v\n", p.ID.Pretty())
207
+ u.DOut("running worker for: %v\n", p.ID.Pretty())
208
209
// finally, run the query against this peer
210
res, err := r.query.qfunc(r.ctx, p)
211
212
if err != nil {
213
- u.POut("ERROR worker for: %v %v\n", p.ID.Pretty(), err)
213
+ u.DOut("ERROR worker for: %v %v\n", p.ID.Pretty(), err)
214
r.Lock()
215
r.errs = append(r.errs, err)
216
r.Unlock()
217
218
} else if res.success {
219
- u.POut("SUCCESS worker for: %v\n", p.ID.Pretty(), res)
219
+ u.DOut("SUCCESS worker for: %v\n", p.ID.Pretty(), res)
220
r.Lock()
221
r.result = res
222
r.Unlock()
223
r.cancel() // signal to everyone that we're done.
224
225
} else if res.closerPeers != nil {
226
- u.POut("PEERS CLOSER -- worker for: %v\n", p.ID.Pretty())
226
+ u.DOut("PEERS CLOSER -- worker for: %v\n", p.ID.Pretty())
227
for _, next := range res.closerPeers {
228
r.addPeerToQuery(next, p)
229
}
230
}
231
232
// signal we're done proccessing peer p
233
- u.POut("completing worker for: %v\n", p.ID.Pretty())
233
+ u.DOut("completing worker for: %v\n", p.ID.Pretty())
234
r.peersRemaining.Decrement(1)
235
r.rateLimit <- struct{}{}
236
}
routing/dht/routing.go
+2
@@ -30,6 +30,7 @@ func (dht *IpfsDHT) PutValue(key u.Key, value []byte) error {
30
}
31
32
query := newQuery(key, func(ctx context.Context, p *peer.Peer) (*dhtQueryResult, error) {
33
+ u.DOut("[%s] PutValue qry part %v\n", dht.self.ID.Pretty(), p.ID.Pretty())
34
err := dht.putValueToNetwork(ctx, p, string(key), value)
35
if err != nil {
36
return nil, err
@@ -38,6 +39,7 @@ func (dht *IpfsDHT) PutValue(key u.Key, value []byte) error {
39
})
40
41
_, err := query.Run(ctx, peers)
42
+ u.DOut("[%s] PutValue %v %v\n", dht.self.ID.Pretty(), key, value)
43
return err
44
}
45