@cryptotaxi247 / kubo / commits / 7fcf56e8e

namesys: select on output

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

Łukasz Magiera committed Oct 16, 2018 at 17:45 UTC 7fcf56e8e532997f7b93f51587b74ebf7e37c5b4
4 files changed +32 -54
namesys/base.go
+17 -20
@@ -49,6 +49,11 @@ func resolveAsync(ctx context.Context, r resolver, name string, options opts.Res
49 defer close(outCh)
50 var subCh <-chan Result
51 var cancelSub context.CancelFunc
52 + defer func() {
53 + if cancelSub != nil {
54 + cancelSub()
55 + }
56 + }()
57
58 for {
59 select {
@@ -59,20 +64,17 @@ func resolveAsync(ctx context.Context, r resolver, name string, options opts.Res
64 }
65
66 if res.err != nil {
62 - outCh <- Result{Err: res.err}
63 - if cancelSub != nil {
64 - cancelSub()
65 - }
67 + emitResult(ctx, outCh, Result{Err: res.err})
68 return
69 }
70 log.Debugf("resolved %s to %s", name, res.value.String())
71 if !strings.HasPrefix(res.value.String(), ipnsPrefix) {
70 - outCh <- Result{Path: res.value}
72 + emitResult(ctx, outCh, Result{Path: res.value})
73 break
74 }
75
76 if depth == 1 {
75 - outCh <- Result{Path: res.value, Err: ErrResolveRecursion}
77 + emitResult(ctx, outCh, Result{Path: res.value, Err: ErrResolveRecursion})
78 break
79 }
80
@@ -87,6 +89,7 @@ func resolveAsync(ctx context.Context, r resolver, name string, options opts.Res
89 cancelSub()
90 }
91 subCtx, cancelSub = context.WithCancel(ctx)
92 + _ = cancelSub
93
94 p := strings.TrimPrefix(res.value.String(), ipnsPrefix)
95 subCh = resolveAsync(subCtx, r, p, subopts)
@@ -96,27 +99,21 @@ func resolveAsync(ctx context.Context, r resolver, name string, options opts.Res
99 break
100 }
101
99 - select {
100 - case outCh <- res:
101 - case <-ctx.Done():
102 - if cancelSub != nil {
103 - cancelSub()
104 - }
105 - return
106 - }
102 + emitResult(ctx, outCh, res)
103 case <-ctx.Done():
108 - if cancelSub != nil {
109 - cancelSub()
110 - }
104 return
105 }
106 if resCh == nil && subCh == nil {
114 - if cancelSub != nil {
115 - cancelSub()
116 - }
107 return
108 }
109 }
110 }()
111 return outCh
112 }
113 +
114 +func emitResult(ctx context.Context, outCh chan<- Result, r Result) {
115 + select {
116 + case outCh <- r:
117 + case <-ctx.Done():
118 + }
119 +}
namesys/dns.go
+2 -8
@@ -80,10 +80,7 @@ func (r *DNSResolver) resolveOnceAsync(ctx context.Context, name string, options
80 }
81 if subRes.error == nil {
82 p, err := appendPath(subRes.path)
83 - select {
84 - case out <- onceResult{value: p, err: err}:
85 - case <-ctx.Done():
86 - }
83 + emitOnceResult(ctx, out, onceResult{value: p, err: err})
84 return
85 }
86 case rootRes, ok := <-rootChan:
@@ -93,10 +90,7 @@ func (r *DNSResolver) resolveOnceAsync(ctx context.Context, name string, options
90 }
91 if rootRes.error == nil {
92 p, err := appendPath(rootRes.path)
96 - select {
97 - case out <- onceResult{value: p, err: err}:
98 - case <-ctx.Done():
99 - }
93 + emitOnceResult(ctx, out, onceResult{value: p, err: err})
94 }
95 case <-ctx.Done():
96 return
namesys/namesys.go
+9 -10
@@ -141,19 +141,11 @@ func (ns *mpns) resolveOnceAsync(ctx context.Context, name string, options opts.
141 if len(segments) > 3 {
142 p, err := path.FromSegments("", strings.TrimRight(p.String(), "/"), segments[3])
143 if err != nil {
144 - select {
145 - case out <- onceResult{value: p, err: err}:
146 - case <-ctx.Done():
147 - }
148 - return
144 + emitOnceResult(ctx, out, onceResult{value: p, ttl: res.ttl, err: err})
145 }
146 }
147
152 - select {
153 - case out <- onceResult{value: p, ttl: res.ttl, err: res.err}:
154 - case <-ctx.Done():
155 - return
156 - }
148 + emitOnceResult(ctx, out, onceResult{value: p, ttl: res.ttl, err: res.err})
149 case <-ctx.Done():
150 return
151 }
@@ -163,6 +155,13 @@ func (ns *mpns) resolveOnceAsync(ctx context.Context, name string, options opts.
155 return out
156 }
157
158 +func emitOnceResult(ctx context.Context, outCh chan<- onceResult, r onceResult) {
159 + select {
160 + case outCh <- r:
161 + case <-ctx.Done():
162 + }
163 +}
164 +
165 // Publish implements Publisher
166 func (ns *mpns) Publish(ctx context.Context, name ci.PrivKey, value path.Path) error {
167 return ns.PublishWithEOL(ctx, name, value, time.Now().Add(DefaultRecordTTL))
namesys/routing.go
+4 -16
@@ -112,10 +112,7 @@ func (r *IpnsResolver) resolveOnceAsync(ctx context.Context, name string, option
112 err = proto.Unmarshal(val, entry)
113 if err != nil {
114 log.Debugf("RoutingResolver: could not unmarshal value for name %s: %s", name, err)
115 - select {
116 - case out <- onceResult{err: err}:
117 - case <-ctx.Done():
118 - }
115 + emitOnceResult(ctx, out, onceResult{err: err})
116 return
117 }
118
@@ -129,10 +126,7 @@ func (r *IpnsResolver) resolveOnceAsync(ctx context.Context, name string, option
126 // Not a multihash, probably a new style record
127 p, err = path.ParsePath(string(entry.GetValue()))
128 if err != nil {
132 - select {
133 - case out <- onceResult{err: err}:
134 - case <-ctx.Done():
135 - }
129 + emitOnceResult(ctx, out, onceResult{err: err})
130 return
131 }
132 }
@@ -154,17 +148,11 @@ func (r *IpnsResolver) resolveOnceAsync(ctx context.Context, name string, option
148 }
149 default:
150 log.Errorf("encountered error when parsing EOL: %s", err)
157 - select {
158 - case out <- onceResult{err: err}:
159 - case <-ctx.Done():
160 - }
151 + emitOnceResult(ctx, out, onceResult{err: err})
152 return
153 }
154
164 - select {
165 - case out <- onceResult{value: p, ttl: ttl}:
166 - case <-ctx.Done():
167 - }
155 + emitOnceResult(ctx, out, onceResult{value: p, ttl: ttl})
156 case <-ctx.Done():
157 return
158 }