reprovider: Implement 'bitswap reprovide' command
License: MIT Signed-off-by: Łukasz Magiera <magik6k@gmail.com>
Łukasz Magiera committed
Aug 1, 2017 at 17:45 UTC
17ae331be2be53358445d91f6356185f28465887
4 files changed
+70
-15
core/commands/bitswap.go
+32
-4
@@ -21,10 +21,11 @@ var BitswapCmd = &cmds.Command{
21
ShortDescription: ``,
22
},
23
Subcommands: map[string]*cmds.Command{
24
- "wantlist": showWantlistCmd,
25
- "stat": bitswapStatCmd,
26
- "unwant": unwantCmd,
27
- "ledger": ledgerCmd,
24
+ "wantlist": showWantlistCmd,
25
+ "stat": bitswapStatCmd,
26
+ "unwant": unwantCmd,
27
+ "ledger": ledgerCmd,
28
+ "reprovide": reprovideCmd,
29
},
30
}
31
@@ -242,3 +243,30 @@ prints the ledger associated with a given peer.
243
},
244
},
245
}
246
+
247
+var reprovideCmd = &cmds.Command{
248
+ Helptext: cmds.HelpText{
249
+ Tagline: "Trigger reprovider.",
250
+ ShortDescription: `
251
+Trigger reprovider to announce our data to network.
252
+`,
253
+ },
254
+ Run: func(req cmds.Request, res cmds.Response) {
255
+ nd, err := req.InvocContext().GetNode()
256
+ if err != nil {
257
+ res.SetError(err, cmds.ErrNormal)
258
+ return
259
+ }
260
+
261
+ if !nd.OnlineMode() {
262
+ res.SetError(errNotOnline, cmds.ErrClient)
263
+ return
264
+ }
265
+
266
+ err = nd.Reprovider.Trigger(req.Context())
267
+ if err != nil {
268
+ res.SetError(err, cmds.ErrNormal)
269
+ return
270
+ }
271
+ },
272
+}
core/core.go
+2
-2
@@ -277,7 +277,7 @@ func (n *IpfsNode) startLateOnlineServices(ctx context.Context) error {
277
default:
278
return fmt.Errorf("unknown reprovider strtaegy '%s'", cfg.Reprovider.Strategy)
279
}
280
- n.Reprovider = rp.NewReprovider(n.Routing, keyProvider)
280
+ n.Reprovider = rp.NewReprovider(ctx, n.Routing, keyProvider)
281
282
if cfg.Reprovider.Interval != "0" {
283
interval := kReprovideFrequency
@@ -290,7 +290,7 @@ func (n *IpfsNode) startLateOnlineServices(ctx context.Context) error {
290
interval = dur
291
}
292
293
- go n.Reprovider.ProvideEvery(ctx, interval)
293
+ go n.Reprovider.ProvideEvery(interval)
294
}
295
296
return nil
exchange/reprovide/providers.go
+1
@@ -26,6 +26,7 @@ func NewPinnedProvider(pinning pin.Pinner, dag merkledag.DAGService, onlyRoots b
26
27
outCh := make(chan *cid.Cid)
28
go func() {
29
+ defer close(outCh)
30
set.ForEach(func(c *cid.Cid) error {
31
select {
32
case <-ctx.Done():
exchange/reprovide/reprovide.go
+35
-9
@@ -16,35 +16,48 @@ var log = logging.Logger("reprovider")
16
type KeyChanFunc func(context.Context) (<-chan *cid.Cid, error)
17
18
type Reprovider struct {
19
+ ctx context.Context
20
+ trigger chan context.CancelFunc
21
+
22
// The routing system to provide values through
23
rsys routing.ContentRouting
24
25
keyProvider KeyChanFunc
26
}
27
25
-func NewReprovider(rsys routing.ContentRouting, keyProvider KeyChanFunc) *Reprovider {
28
+func NewReprovider(ctx context.Context, rsys routing.ContentRouting, keyProvider KeyChanFunc) *Reprovider {
29
return &Reprovider{
27
- rsys: rsys,
30
+ ctx: ctx,
31
+ trigger: make(chan context.CancelFunc),
32
+
33
+ rsys: rsys,
34
keyProvider: keyProvider,
35
}
36
}
37
32
-func (rp *Reprovider) ProvideEvery(ctx context.Context, tick time.Duration) {
38
+func (rp *Reprovider) ProvideEvery(tick time.Duration) {
39
// dont reprovide immediately.
40
// may have just started the daemon and shutting it down immediately.
41
// probability( up another minute | uptime ) increases with uptime.
42
after := time.After(time.Minute)
43
+ var done context.CancelFunc
44
for {
45
select {
39
- case <-ctx.Done():
46
+ case <-rp.ctx.Done():
47
return
48
+ case done = <-rp.trigger:
49
case <-after:
42
- err := rp.Reprovide(ctx)
43
- if err != nil {
44
- log.Debug(err)
45
- }
46
- after = time.After(tick)
50
}
51
+
52
+ err := rp.Reprovide(rp.ctx)
53
+ if err != nil {
54
+ log.Debug(err)
55
+ }
56
+
57
+ if done != nil {
58
+ done()
59
+ }
60
+ after = time.After(tick)
61
}
62
}
63
@@ -72,3 +85,16 @@ func (rp *Reprovider) Reprovide(ctx context.Context) error {
85
}
86
return nil
87
}
88
+
89
+func (rp *Reprovider) Trigger(ctx context.Context) error {
90
+ progressCtx, done := context.WithCancel(ctx)
91
+ select {
92
+ case <-rp.ctx.Done():
93
+ return context.Canceled
94
+ case <-ctx.Done():
95
+ return context.Canceled
96
+ case rp.trigger <- done:
97
+ <-progressCtx.Done()
98
+ return nil
99
+ }
100
+}