@cryptotaxi247 / kubo / commits / a865fde21

reprovider: make reprovide cmd error if reprovider is active

License: MIT Signed-off-by: Łukasz Magiera <magik6k@gmail.com>

Łukasz Magiera committed Aug 3, 2017 at 15:13 UTC a865fde21d2fd3c5a763c0d70bac9bb190296f8c
1 file changed +35 -6
exchange/reprovide/reprovide.go
+35 -6
@@ -14,10 +14,11 @@ import (
14 var log = logging.Logger("reprovider")
15
16 type keyChanFunc func(context.Context) (<-chan *cid.Cid, error)
17 +type doneFunc func(error)
18
19 type Reprovider struct {
20 ctx context.Context
20 - trigger chan context.CancelFunc
21 + trigger chan doneFunc
22
23 // The routing system to provide values through
24 rsys routing.ContentRouting
@@ -29,7 +30,7 @@ type Reprovider struct {
30 func NewReprovider(ctx context.Context, rsys routing.ContentRouting, keyProvider keyChanFunc) *Reprovider {
31 return &Reprovider{
32 ctx: ctx,
32 - trigger: make(chan context.CancelFunc),
33 + trigger: make(chan doneFunc),
34
35 rsys: rsys,
36 keyProvider: keyProvider,
@@ -42,7 +43,7 @@ func (rp *Reprovider) ProvideEvery(tick time.Duration) {
43 // may have just started the daemon and shutting it down immediately.
44 // probability( up another minute | uptime ) increases with uptime.
45 after := time.After(time.Minute)
45 - var done context.CancelFunc
46 + var done doneFunc
47 for {
48 select {
49 case <-rp.ctx.Done():
@@ -51,14 +52,19 @@ func (rp *Reprovider) ProvideEvery(tick time.Duration) {
52 case <-after:
53 }
54
55 + unmute := rp.muteTrigger()
56 +
57 err := rp.Reprovide()
58 if err != nil {
59 log.Debug(err)
60 }
61
62 if done != nil {
60 - done()
63 + done(err)
64 }
65 +
66 + unmute()
67 +
68 after = time.After(tick)
69 }
70 }
@@ -93,13 +99,36 @@ func (rp *Reprovider) Reprovide() error {
99 func (rp *Reprovider) Trigger(ctx context.Context) error {
100 progressCtx, done := context.WithCancel(ctx)
101
102 + var err error
103 + df := func(e error) {
104 + err = e
105 + done()
106 + }
107 +
108 select {
109 case <-rp.ctx.Done():
110 return context.Canceled
111 case <-ctx.Done():
112 return context.Canceled
101 - case rp.trigger <- done:
113 + case rp.trigger <- df:
114 <-progressCtx.Done()
103 - return nil
115 + return err
116 }
117 }
118 +
119 +func (rp *Reprovider) muteTrigger() context.CancelFunc {
120 + ctx, cf := context.WithCancel(rp.ctx)
121 + go func() {
122 + defer cf()
123 + for {
124 + select {
125 + case <-ctx.Done():
126 + return
127 + case done := <-rp.trigger:
128 + done(fmt.Errorf("reprovider is already running"))
129 + }
130 + }
131 + }()
132 +
133 + return cf
134 +}