@cryptotaxi247 / kubo / commits / d91955b41

moved handlers to own file

Juan Batiz-Benet committed Sep 16, 2014 at 02:07 UTC d91955b4121e953371ec9da790087f7ed6bfa20c
2 files changed +259 -244
routing/dht/dht.go
-244
@@ -159,250 +159,6 @@ func (dht *IpfsDHT) HandleMessage(ctx context.Context, mes msg.NetMessage) (msg.
159 return rmes, nil
160 }
161
162 -// dhthandler specifies the signature of functions that handle DHT messages.
163 -type dhtHandler func(*peer.Peer, *Message) (*Message, error)
164 -
165 -func (dht *IpfsDHT) handlerForMsgType(t Message_MessageType) dhtHandler {
166 - switch t {
167 - case Message_GET_VALUE:
168 - return dht.handleGetValue
169 - // case Message_PUT_VALUE:
170 - // return dht.handlePutValue
171 - case Message_FIND_NODE:
172 - return dht.handleFindPeer
173 - // case Message_ADD_PROVIDER:
174 - // return dht.handleAddProvider
175 - // case Message_GET_PROVIDERS:
176 - // return dht.handleGetProviders
177 - case Message_PING:
178 - return dht.handlePing
179 - // case Message_DIAGNOSTIC:
180 - // return dht.handleDiagnostic
181 - default:
182 - return nil
183 - }
184 -}
185 -
186 -func (dht *IpfsDHT) putValueToNetwork(p *peer.Peer, key string, value []byte) error {
187 - typ := Message_PUT_VALUE
188 - pmes := &Message{
189 - Type: &typ,
190 - Key: &key,
191 - Value: value,
192 - }
193 -
194 - mes, err := msg.FromObject(p, pmes)
195 - if err != nil {
196 - return err
197 - }
198 - return dht.sender.SendMessage(context.TODO(), mes)
199 -}
200 -
201 -func (dht *IpfsDHT) handleGetValue(p *peer.Peer, pmes *Message) (*Message, error) {
202 - u.DOut("handleGetValue for key: %s\n", pmes.GetKey())
203 -
204 - // setup response
205 - resp := &Message{
206 - Type: pmes.Type,
207 - Key: pmes.Key,
208 - }
209 -
210 - // first, is the key even a key?
211 - key := pmes.GetKey()
212 - if key == "" {
213 - return nil, errors.New("handleGetValue but no key was provided")
214 - }
215 -
216 - // let's first check if we have the value locally.
217 - dskey := ds.NewKey(pmes.GetKey())
218 - iVal, err := dht.datastore.Get(dskey)
219 -
220 - // if we got an unexpected error, bail.
221 - if err != ds.ErrNotFound {
222 - return nil, err
223 - }
224 -
225 - // if we have the value, respond with it!
226 - if err == nil {
227 - u.DOut("handleGetValue success!\n")
228 -
229 - byts, ok := iVal.([]byte)
230 - if !ok {
231 - return nil, fmt.Errorf("datastore had non byte-slice value for %v", dskey)
232 - }
233 -
234 - resp.Value = byts
235 - return resp, nil
236 - }
237 -
238 - // if we know any providers for the requested value, return those.
239 - provs := dht.providers.GetProviders(u.Key(pmes.GetKey()))
240 - if len(provs) > 0 {
241 - u.DOut("handleGetValue returning %d provider[s]\n", len(provs))
242 - resp.ProviderPeers = peersToPBPeers(provs)
243 - return resp, nil
244 - }
245 -
246 - // Find closest peer on given cluster to desired key and reply with that info
247 - closer := dht.betterPeerToQuery(pmes)
248 - if closer == nil {
249 - u.DOut("handleGetValue could not find a closer node than myself.\n")
250 - resp.CloserPeers = nil
251 - return resp, nil
252 - }
253 -
254 - // we got a closer peer, it seems. return it.
255 - u.DOut("handleGetValue returning a closer peer: '%s'\n", closer.ID.Pretty())
256 - resp.CloserPeers = peersToPBPeers([]*peer.Peer{closer})
257 - return resp, nil
258 -}
259 -
260 -// Store a value in this peer local storage
261 -func (dht *IpfsDHT) handlePutValue(p *peer.Peer, pmes *Message) {
262 - dht.dslock.Lock()
263 - defer dht.dslock.Unlock()
264 - dskey := ds.NewKey(pmes.GetKey())
265 - err := dht.datastore.Put(dskey, pmes.GetValue())
266 - if err != nil {
267 - // For now, just panic, handle this better later maybe
268 - panic(err)
269 - }
270 -}
271 -
272 -func (dht *IpfsDHT) handlePing(p *peer.Peer, pmes *Message) (*Message, error) {
273 - u.DOut("[%s] Responding to ping from [%s]!\n", dht.self.ID.Pretty(), p.ID.Pretty())
274 - return &Message{Type: pmes.Type}, nil
275 -}
276 -
277 -func (dht *IpfsDHT) handleFindPeer(p *peer.Peer, pmes *Message) (*Message, error) {
278 - resp := &Message{Type: pmes.Type}
279 - var closest *peer.Peer
280 -
281 - // if looking for self... special case where we send it on CloserPeers.
282 - if peer.ID(pmes.GetKey()).Equal(dht.self.ID) {
283 - closest = dht.self
284 - } else {
285 - closest = dht.betterPeerToQuery(pmes)
286 - }
287 -
288 - if closest == nil {
289 - u.PErr("handleFindPeer: could not find anything.\n")
290 - return resp, nil
291 - }
292 -
293 - if len(closest.Addresses) == 0 {
294 - u.PErr("handleFindPeer: no addresses for connected peer...\n")
295 - return resp, nil
296 - }
297 -
298 - u.DOut("handleFindPeer: sending back '%s'\n", closest.ID.Pretty())
299 - resp.CloserPeers = peersToPBPeers([]*peer.Peer{closest})
300 - return resp, nil
301 -}
302 -
303 -func (dht *IpfsDHT) handleGetProviders(p *peer.Peer, pmes *Message) (*Message, error) {
304 - resp := &Message{
305 - Type: pmes.Type,
306 - Key: pmes.Key,
307 - }
308 -
309 - // check if we have this value, to add ourselves as provider.
310 - has, err := dht.datastore.Has(ds.NewKey(pmes.GetKey()))
311 - if err != nil && err != ds.ErrNotFound {
312 - u.PErr("unexpected datastore error: %v\n", err)
313 - has = false
314 - }
315 -
316 - // setup providers
317 - providers := dht.providers.GetProviders(u.Key(pmes.GetKey()))
318 - if has {
319 - providers = append(providers, dht.self)
320 - }
321 -
322 - // if we've got providers, send thos those.
323 - if providers != nil && len(providers) > 0 {
324 - resp.ProviderPeers = peersToPBPeers(providers)
325 - }
326 -
327 - // Also send closer peers.
328 - closer := dht.betterPeerToQuery(pmes)
329 - if closer != nil {
330 - resp.CloserPeers = peersToPBPeers([]*peer.Peer{closer})
331 - }
332 -
333 - return resp, nil
334 -}
335 -
336 -type providerInfo struct {
337 - Creation time.Time
338 - Value *peer.Peer
339 -}
340 -
341 -func (dht *IpfsDHT) handleAddProvider(p *peer.Peer, pmes *Message) {
342 - key := u.Key(pmes.GetKey())
343 - u.DOut("[%s] Adding [%s] as a provider for '%s'\n",
344 - dht.self.ID.Pretty(), p.ID.Pretty(), peer.ID(key).Pretty())
345 - dht.providers.AddProvider(key, p)
346 -}
347 -
348 -// Halt stops all communications from this peer and shut down
349 -// TODO -- remove this in favor of context
350 -func (dht *IpfsDHT) Halt() {
351 - dht.shutdown <- struct{}{}
352 - dht.network.Close()
353 - dht.providers.Halt()
354 -}
355 -
356 -// NOTE: not yet finished, low priority
357 -func (dht *IpfsDHT) handleDiagnostic(p *peer.Peer, pmes *Message) (*Message, error) {
358 - seq := dht.routingTables[0].NearestPeers(kb.ConvertPeerID(dht.self.ID), 10)
359 -
360 - for _, ps := range seq {
361 - mes, err := msg.FromObject(ps, pmes)
362 - if err != nil {
363 - u.PErr("handleDiagnostics error creating message: %v\n", err)
364 - continue
365 - }
366 - // dht.sender.SendRequest(context.TODO(), mes)
367 - }
368 - return nil, errors.New("not yet ported back")
369 -
370 - // buf := new(bytes.Buffer)
371 - // di := dht.getDiagInfo()
372 - // buf.Write(di.Marshal())
373 - //
374 - // // NOTE: this shouldnt be a hardcoded value
375 - // after := time.After(time.Second * 20)
376 - // count := len(seq)
377 - // for count > 0 {
378 - // select {
379 - // case <-after:
380 - // //Timeout, return what we have
381 - // goto out
382 - // case reqResp := <-listenChan:
383 - // pmesOut := new(Message)
384 - // err := proto.Unmarshal(reqResp.Data, pmesOut)
385 - // if err != nil {
386 - // // It broke? eh, whatever, keep going
387 - // continue
388 - // }
389 - // buf.Write(reqResp.Data)
390 - // count--
391 - // }
392 - // }
393 - //
394 - // out:
395 - // resp := Message{
396 - // Type: Message_DIAGNOSTIC,
397 - // ID: pmes.GetId(),
398 - // Value: buf.Bytes(),
399 - // Response: true,
400 - // }
401 - //
402 - // mes := swarm.NewMessage(p, resp.ToProtobuf())
403 - // dht.netChan.Outgoing <- mes
404 -}
405 -
162 func (dht *IpfsDHT) getValueOrPeers(p *peer.Peer, key u.Key, timeout time.Duration, level int) ([]byte, []*peer.Peer, error) {
163 pmes, err := dht.getValueSingle(p, key, timeout, level)
164 if err != nil {
routing/dht/handlers.go new
+259
@@ -0,0 +1,259 @@
1 +package dht
2 +
3 +import (
4 + "errors"
5 + "fmt"
6 + "time"
7 +
8 + msg "github.com/jbenet/go-ipfs/net/message"
9 + peer "github.com/jbenet/go-ipfs/peer"
10 + kb "github.com/jbenet/go-ipfs/routing/kbucket"
11 + u "github.com/jbenet/go-ipfs/util"
12 +
13 + context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
14 + ds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/datastore.go"
15 +)
16 +
17 +// dhthandler specifies the signature of functions that handle DHT messages.
18 +type dhtHandler func(*peer.Peer, *Message) (*Message, error)
19 +
20 +func (dht *IpfsDHT) handlerForMsgType(t Message_MessageType) dhtHandler {
21 + switch t {
22 + case Message_GET_VALUE:
23 + return dht.handleGetValue
24 + // case Message_PUT_VALUE:
25 + // return dht.handlePutValue
26 + case Message_FIND_NODE:
27 + return dht.handleFindPeer
28 + // case Message_ADD_PROVIDER:
29 + // return dht.handleAddProvider
30 + // case Message_GET_PROVIDERS:
31 + // return dht.handleGetProviders
32 + case Message_PING:
33 + return dht.handlePing
34 + // case Message_DIAGNOSTIC:
35 + // return dht.handleDiagnostic
36 + default:
37 + return nil
38 + }
39 +}
40 +
41 +func (dht *IpfsDHT) putValueToNetwork(p *peer.Peer, key string, value []byte) error {
42 + typ := Message_PUT_VALUE
43 + pmes := &Message{
44 + Type: &typ,
45 + Key: &key,
46 + Value: value,
47 + }
48 +
49 + mes, err := msg.FromObject(p, pmes)
50 + if err != nil {
51 + return err
52 + }
53 + return dht.sender.SendMessage(context.TODO(), mes)
54 +}
55 +
56 +func (dht *IpfsDHT) handleGetValue(p *peer.Peer, pmes *Message) (*Message, error) {
57 + u.DOut("handleGetValue for key: %s\n", pmes.GetKey())
58 +
59 + // setup response
60 + resp := &Message{
61 + Type: pmes.Type,
62 + Key: pmes.Key,
63 + }
64 +
65 + // first, is the key even a key?
66 + key := pmes.GetKey()
67 + if key == "" {
68 + return nil, errors.New("handleGetValue but no key was provided")
69 + }
70 +
71 + // let's first check if we have the value locally.
72 + dskey := ds.NewKey(pmes.GetKey())
73 + iVal, err := dht.datastore.Get(dskey)
74 +
75 + // if we got an unexpected error, bail.
76 + if err != ds.ErrNotFound {
77 + return nil, err
78 + }
79 +
80 + // if we have the value, respond with it!
81 + if err == nil {
82 + u.DOut("handleGetValue success!\n")
83 +
84 + byts, ok := iVal.([]byte)
85 + if !ok {
86 + return nil, fmt.Errorf("datastore had non byte-slice value for %v", dskey)
87 + }
88 +
89 + resp.Value = byts
90 + return resp, nil
91 + }
92 +
93 + // if we know any providers for the requested value, return those.
94 + provs := dht.providers.GetProviders(u.Key(pmes.GetKey()))
95 + if len(provs) > 0 {
96 + u.DOut("handleGetValue returning %d provider[s]\n", len(provs))
97 + resp.ProviderPeers = peersToPBPeers(provs)
98 + return resp, nil
99 + }
100 +
101 + // Find closest peer on given cluster to desired key and reply with that info
102 + closer := dht.betterPeerToQuery(pmes)
103 + if closer == nil {
104 + u.DOut("handleGetValue could not find a closer node than myself.\n")
105 + resp.CloserPeers = nil
106 + return resp, nil
107 + }
108 +
109 + // we got a closer peer, it seems. return it.
110 + u.DOut("handleGetValue returning a closer peer: '%s'\n", closer.ID.Pretty())
111 + resp.CloserPeers = peersToPBPeers([]*peer.Peer{closer})
112 + return resp, nil
113 +}
114 +
115 +// Store a value in this peer local storage
116 +func (dht *IpfsDHT) handlePutValue(p *peer.Peer, pmes *Message) {
117 + dht.dslock.Lock()
118 + defer dht.dslock.Unlock()
119 + dskey := ds.NewKey(pmes.GetKey())
120 + err := dht.datastore.Put(dskey, pmes.GetValue())
121 + if err != nil {
122 + // For now, just panic, handle this better later maybe
123 + panic(err)
124 + }
125 +}
126 +
127 +func (dht *IpfsDHT) handlePing(p *peer.Peer, pmes *Message) (*Message, error) {
128 + u.DOut("[%s] Responding to ping from [%s]!\n", dht.self.ID.Pretty(), p.ID.Pretty())
129 + return &Message{Type: pmes.Type}, nil
130 +}
131 +
132 +func (dht *IpfsDHT) handleFindPeer(p *peer.Peer, pmes *Message) (*Message, error) {
133 + resp := &Message{Type: pmes.Type}
134 + var closest *peer.Peer
135 +
136 + // if looking for self... special case where we send it on CloserPeers.
137 + if peer.ID(pmes.GetKey()).Equal(dht.self.ID) {
138 + closest = dht.self
139 + } else {
140 + closest = dht.betterPeerToQuery(pmes)
141 + }
142 +
143 + if closest == nil {
144 + u.PErr("handleFindPeer: could not find anything.\n")
145 + return resp, nil
146 + }
147 +
148 + if len(closest.Addresses) == 0 {
149 + u.PErr("handleFindPeer: no addresses for connected peer...\n")
150 + return resp, nil
151 + }
152 +
153 + u.DOut("handleFindPeer: sending back '%s'\n", closest.ID.Pretty())
154 + resp.CloserPeers = peersToPBPeers([]*peer.Peer{closest})
155 + return resp, nil
156 +}
157 +
158 +func (dht *IpfsDHT) handleGetProviders(p *peer.Peer, pmes *Message) (*Message, error) {
159 + resp := &Message{
160 + Type: pmes.Type,
161 + Key: pmes.Key,
162 + }
163 +
164 + // check if we have this value, to add ourselves as provider.
165 + has, err := dht.datastore.Has(ds.NewKey(pmes.GetKey()))
166 + if err != nil && err != ds.ErrNotFound {
167 + u.PErr("unexpected datastore error: %v\n", err)
168 + has = false
169 + }
170 +
171 + // setup providers
172 + providers := dht.providers.GetProviders(u.Key(pmes.GetKey()))
173 + if has {
174 + providers = append(providers, dht.self)
175 + }
176 +
177 + // if we've got providers, send thos those.
178 + if providers != nil && len(providers) > 0 {
179 + resp.ProviderPeers = peersToPBPeers(providers)
180 + }
181 +
182 + // Also send closer peers.
183 + closer := dht.betterPeerToQuery(pmes)
184 + if closer != nil {
185 + resp.CloserPeers = peersToPBPeers([]*peer.Peer{closer})
186 + }
187 +
188 + return resp, nil
189 +}
190 +
191 +type providerInfo struct {
192 + Creation time.Time
193 + Value *peer.Peer
194 +}
195 +
196 +func (dht *IpfsDHT) handleAddProvider(p *peer.Peer, pmes *Message) {
197 + key := u.Key(pmes.GetKey())
198 + u.DOut("[%s] Adding [%s] as a provider for '%s'\n",
199 + dht.self.ID.Pretty(), p.ID.Pretty(), peer.ID(key).Pretty())
200 + dht.providers.AddProvider(key, p)
201 +}
202 +
203 +// Halt stops all communications from this peer and shut down
204 +// TODO -- remove this in favor of context
205 +func (dht *IpfsDHT) Halt() {
206 + dht.shutdown <- struct{}{}
207 + dht.network.Close()
208 + dht.providers.Halt()
209 +}
210 +
211 +// NOTE: not yet finished, low priority
212 +func (dht *IpfsDHT) handleDiagnostic(p *peer.Peer, pmes *Message) (*Message, error) {
213 + seq := dht.routingTables[0].NearestPeers(kb.ConvertPeerID(dht.self.ID), 10)
214 +
215 + for _, ps := range seq {
216 + mes, err := msg.FromObject(ps, pmes)
217 + if err != nil {
218 + u.PErr("handleDiagnostics error creating message: %v\n", err)
219 + continue
220 + }
221 + // dht.sender.SendRequest(context.TODO(), mes)
222 + }
223 + return nil, errors.New("not yet ported back")
224 +
225 + // buf := new(bytes.Buffer)
226 + // di := dht.getDiagInfo()
227 + // buf.Write(di.Marshal())
228 + //
229 + // // NOTE: this shouldnt be a hardcoded value
230 + // after := time.After(time.Second * 20)
231 + // count := len(seq)
232 + // for count > 0 {
233 + // select {
234 + // case <-after:
235 + // //Timeout, return what we have
236 + // goto out
237 + // case reqResp := <-listenChan:
238 + // pmesOut := new(Message)
239 + // err := proto.Unmarshal(reqResp.Data, pmesOut)
240 + // if err != nil {
241 + // // It broke? eh, whatever, keep going
242 + // continue
243 + // }
244 + // buf.Write(reqResp.Data)
245 + // count--
246 + // }
247 + // }
248 + //
249 + // out:
250 + // resp := Message{
251 + // Type: Message_DIAGNOSTIC,
252 + // ID: pmes.GetId(),
253 + // Value: buf.Bytes(),
254 + // Response: true,
255 + // }
256 + //
257 + // mes := swarm.NewMessage(p, resp.ToProtobuf())
258 + // dht.netChan.Outgoing <- mes
259 +}