commands/dht: use res.Emit directly
License: MIT Signed-off-by: Overbool <overbool.xu@gmail.com>
Overbool committed
Oct 28, 2018 at 12:41 UTC
4a2c3cb3cf86929f3219c43a5aaf708e4863980d
1 file changed
+43
-85
core/commands/dht.go
+43
-85
@@ -8,18 +8,18 @@ import (
8
"time"
9
10
cmdenv "github.com/ipfs/go-ipfs/core/commands/cmdenv"
11
- dag "gx/ipfs/QmSei8kFMfqdJq7Q68d2LMnHbTWKKg2daA29ezUYFAUNgc/go-merkledag"
12
- path "gx/ipfs/QmT3rzed1ppXefourpmoZ7tyVQfsGPQZ1pHDngLmCvXxd3/go-path"
11
12
cid "gx/ipfs/QmPSQnBKM9g7BaUcZCvswUJVscQ1ipjmwxN5PXCjkp9EQ7/go-cid"
13
ipld "gx/ipfs/QmR7TcHkR9nxkUorfi8XMTAMLUK7GiP64TWWBzY3aacc1o/go-ipld-format"
14
+ cmds "gx/ipfs/QmSXUokcP4TJpFfqozT69AVAYRtzXVMUjzQVkYX41R9Svs/go-ipfs-cmds"
15
+ dag "gx/ipfs/QmSei8kFMfqdJq7Q68d2LMnHbTWKKg2daA29ezUYFAUNgc/go-merkledag"
16
+ path "gx/ipfs/QmT3rzed1ppXefourpmoZ7tyVQfsGPQZ1pHDngLmCvXxd3/go-path"
17
peer "gx/ipfs/QmTRhk7cgjUf2gfQ3p2M9KPECNZEW9XUrmHcFCgog4cPgB/go-libp2p-peer"
18
pstore "gx/ipfs/QmTTJcDL3gsnGDALjh2fDGg1onGRUdVgNL2hU2WEZcVrMX/go-libp2p-peerstore"
19
b58 "gx/ipfs/QmWFAMPqsEyUX7gDUsRVmMWz59FxSpJ1b2v6bJ1yYzo7jY/go-base58-fast/base58"
20
routing "gx/ipfs/QmcQ81jSyWCp1jpkQ8CMbtpXT3jK7Wg6ZtYmoyWFgBoF9c/go-libp2p-routing"
21
notif "gx/ipfs/QmcQ81jSyWCp1jpkQ8CMbtpXT3jK7Wg6ZtYmoyWFgBoF9c/go-libp2p-routing/notifications"
21
- "gx/ipfs/Qmde5VP1qUkyQXKCfmEUA7bP64V2HAptbJ7phuPp7jXWwg/go-ipfs-cmdkit"
22
- cmds "gx/ipfs/QmSXUokcP4TJpFfqozT69AVAYRtzXVMUjzQVkYX41R9Svs/go-ipfs-cmds"
22
+ cmdkit "gx/ipfs/Qmde5VP1qUkyQXKCfmEUA7bP64V2HAptbJ7phuPp7jXWwg/go-ipfs-cmdkit"
23
)
24
25
var ErrNotDHT = errors.New("routing service is not a DHT")
@@ -93,21 +93,13 @@ var queryDhtCmd = &cmds.Command{
93
}
94
}()
95
96
- outChan := make(chan interface{})
97
-
98
- go func() {
99
- defer close(outChan)
100
- for e := range events {
101
- select {
102
- case outChan <- e:
103
- case <-req.Context.Done():
104
- return
105
- }
96
+ for e := range events {
97
+ if err := res.Emit(e); err != nil {
98
+ return err
99
}
107
- }()
108
-
109
- return res.Emit(outChan)
100
+ }
101
102
+ return nil
103
},
104
Encoders: cmds.EncoderMap{
105
cmds.Text: cmds.MakeTypedEncoder(func(req *cmds.Request, w io.Writer, out *notif.QueryEvent) error {
@@ -167,22 +159,10 @@ var findProvidersDhtCmd = &cmds.Command{
159
return err
160
}
161
170
- outChan := make(chan interface{})
171
-
162
ctx, cancel := context.WithCancel(req.Context)
163
ctx, events := notif.RegisterForQueryEvents(ctx)
164
165
pchan := n.Routing.FindProvidersAsync(ctx, c, numProviders)
176
- go func() {
177
- defer close(outChan)
178
- for e := range events {
179
- select {
180
- case outChan <- e:
181
- case <-req.Context.Done():
182
- return
183
- }
184
- }
185
- }()
166
167
go func() {
168
defer cancel()
@@ -194,8 +174,13 @@ var findProvidersDhtCmd = &cmds.Command{
174
})
175
}
176
}()
177
+ for e := range events {
178
+ if err := res.Emit(e); err != nil {
179
+ return err
180
+ }
181
+ }
182
198
- return res.Emit(outChan)
183
+ return nil
184
},
185
Encoders: cmds.EncoderMap{
186
cmds.Text: cmds.MakeTypedEncoder(func(req *cmds.Request, w io.Writer, out *notif.QueryEvent) error {
@@ -279,22 +264,9 @@ var provideRefDhtCmd = &cmds.Command{
264
cids = append(cids, c)
265
}
266
282
- outChan := make(chan interface{})
283
-
267
ctx, cancel := context.WithCancel(req.Context)
268
ctx, events := notif.RegisterForQueryEvents(ctx)
269
287
- go func() {
288
- defer close(outChan)
289
- for e := range events {
290
- select {
291
- case outChan <- e:
292
- case <-req.Context.Done():
293
- return
294
- }
295
- }
296
- }()
297
-
270
go func() {
271
defer cancel()
272
var err error
@@ -311,7 +283,13 @@ var provideRefDhtCmd = &cmds.Command{
283
}
284
}()
285
314
- return res.Emit(outChan)
286
+ for e := range events {
287
+ if err := res.Emit(e); err != nil {
288
+ return err
289
+ }
290
+ }
291
+
292
+ return nil
293
},
294
Encoders: cmds.EncoderMap{
295
cmds.Text: cmds.MakeTypedEncoder(func(req *cmds.Request, w io.Writer, out *notif.QueryEvent) error {
@@ -395,22 +373,9 @@ var findPeerDhtCmd = &cmds.Command{
373
return err
374
}
375
398
- outChan := make(chan interface{})
399
-
376
ctx, cancel := context.WithCancel(req.Context)
377
ctx, events := notif.RegisterForQueryEvents(ctx)
378
403
- go func() {
404
- defer close(outChan)
405
- for v := range events {
406
- select {
407
- case outChan <- v:
408
- case <-req.Context.Done():
409
- }
410
-
411
- }
412
- }()
413
-
379
go func() {
380
defer cancel()
381
pi, err := nd.Routing.FindPeer(ctx, pid)
@@ -428,7 +393,13 @@ var findPeerDhtCmd = &cmds.Command{
393
})
394
}()
395
431
- return res.Emit(outChan)
396
+ for e := range events {
397
+ if err := res.Emit(e); err != nil {
398
+ return err
399
+ }
400
+ }
401
+
402
+ return nil
403
},
404
Encoders: cmds.EncoderMap{
405
cmds.Text: cmds.MakeTypedEncoder(func(req *cmds.Request, w io.Writer, out *notif.QueryEvent) error {
@@ -484,21 +455,9 @@ Different key types can specify other 'best' rules.
455
return err
456
}
457
487
- outChan := make(chan interface{})
488
-
458
ctx, cancel := context.WithCancel(req.Context)
459
ctx, events := notif.RegisterForQueryEvents(ctx)
460
492
- go func() {
493
- defer close(outChan)
494
- for e := range events {
495
- select {
496
- case outChan <- e:
497
- case <-req.Context.Done():
498
- }
499
- }
500
- }()
501
-
461
go func() {
462
defer cancel()
463
val, err := nd.Routing.GetValue(ctx, dhtkey)
@@ -515,7 +474,13 @@ Different key types can specify other 'best' rules.
474
}
475
}()
476
518
- return res.Emit(outChan)
477
+ for e := range events {
478
+ if err := res.Emit(e); err != nil {
479
+ return err
480
+ }
481
+ }
482
+
483
+ return nil
484
},
485
Encoders: cmds.EncoderMap{
486
cmds.Text: cmds.MakeTypedEncoder(func(req *cmds.Request, w io.Writer, out *notif.QueryEvent) error {
@@ -582,24 +547,11 @@ NOTE: A value may not exceed 2048 bytes.
547
return err
548
}
549
585
- outChan := make(chan interface{})
586
-
550
data := req.Arguments[1]
551
552
ctx, cancel := context.WithCancel(req.Context)
553
ctx, events := notif.RegisterForQueryEvents(ctx)
554
592
- go func() {
593
- defer close(outChan)
594
- for e := range events {
595
- select {
596
- case outChan <- e:
597
- case <-req.Context.Done():
598
- return
599
- }
600
- }
601
- }()
602
-
555
go func() {
556
defer cancel()
557
err := nd.Routing.PutValue(ctx, key, []byte(data))
@@ -611,7 +563,13 @@ NOTE: A value may not exceed 2048 bytes.
563
}
564
}()
565
614
- return res.Emit(outChan)
566
+ for e := range events {
567
+ if err := res.Emit(e); err != nil {
568
+ return err
569
+ }
570
+ }
571
+
572
+ return nil
573
},
574
Encoders: cmds.EncoderMap{
575
cmds.Text: cmds.MakeTypedEncoder(func(req *cmds.Request, w io.Writer, out *notif.QueryEvent) error {