@cryptotaxi247 / kubo / commits / 2a1ee3ae3

use datastore for local data

Jeromy committed Jul 30, 2014 at 20:16 UTC 2a1ee3ae3a8bfa8d5b9aa45eb86bcad9bae33c77
1 file changed +51 -9
routing/dht/dht.go
+51 -9
@@ -3,9 +3,12 @@ package dht
3 import (
4 "sync"
5
6 - peer "github.com/jbenet/go-ipfs/peer"
7 - swarm "github.com/jbenet/go-ipfs/swarm"
8 - u "github.com/jbenet/go-ipfs/util"
6 + peer "github.com/jbenet/go-ipfs/peer"
7 + swarm "github.com/jbenet/go-ipfs/swarm"
8 + u "github.com/jbenet/go-ipfs/util"
9 +
10 + ds "github.com/jbenet/datastore.go"
11 +
12 "code.google.com/p/goprotobuf/proto"
13 )
14
@@ -18,8 +21,11 @@ type IpfsDHT struct {
21
22 network *swarm.Swarm
23
21 - // local data (TEMPORARY: until we formalize data storage with datastore)
22 - data map[string][]byte
24 + // Local peer (yourself)
25 + self *peer.Peer
26 +
27 + // Local data
28 + datastore ds.Datastore
29
30 // map of channels waiting for reply messages
31 listeners map[uint64]chan *swarm.Message
@@ -29,6 +35,15 @@ type IpfsDHT struct {
35 shutdown chan struct{}
36 }
37
38 +func NewDHT(p *peer.Peer) *IpfsDHT {
39 + dht := new(IpfsDHT)
40 + dht.self = p
41 + dht.network = swarm.NewSwarm(p)
42 + dht.listeners = make(map[uint64]chan *swarm.Message)
43 + dht.shutdown = make(chan struct{})
44 + return dht
45 +}
46 +
47 // Read in all messages from swarm and handle them appropriately
48 // NOTE: this function is just a quick sketch
49 func (dht *IpfsDHT) handleMessages() {
@@ -61,10 +76,13 @@ func (dht *IpfsDHT) handleMessages() {
76 case DHTMessage_GET_VALUE:
77 dht.handleGetValue(mes.Peer, pmes)
78 case DHTMessage_PUT_VALUE:
79 + dht.handlePutValue(mes.Peer, pmes)
80 case DHTMessage_FIND_NODE:
81 + dht.handleFindNode(mes.Peer, pmes)
82 case DHTMessage_ADD_PROVIDER:
83 case DHTMessage_GET_PROVIDERS:
84 case DHTMessage_PING:
85 + dht.handleFindNode(mes.Peer, pmes)
86 }
87
88 case <-dht.shutdown:
@@ -74,15 +92,22 @@ func (dht *IpfsDHT) handleMessages() {
92 }
93
94 func (dht *IpfsDHT) handleGetValue(p *peer.Peer, pmes *DHTMessage) {
77 - val, found := dht.data[pmes.GetKey()]
78 - if found {
95 + dskey := ds.NewKey(pmes.GetKey())
96 + i_val, err := dht.datastore.Get(dskey)
97 + if err == nil {
98 isResponse := true
99 resp := new(DHTMessage)
100 resp.Response = &isResponse
101 resp.Id = pmes.Id
102 resp.Key = pmes.Key
103 +
104 + val := i_val.([]byte)
105 resp.Value = val
85 - } else {
106 +
107 + mes := new(swarm.Message)
108 + mes.Peer = p
109 + mes.Data = []byte(resp.String())
110 + } else if err == ds.ErrNotFound {
111 // Find closest node(s) to desired key and reply with that info
112 // TODO: this will need some other metadata in the protobuf message
113 // to signal to the querying node that the data its receiving
@@ -90,8 +115,14 @@ func (dht *IpfsDHT) handleGetValue(p *peer.Peer, pmes *DHTMessage) {
115 }
116 }
117
118 +// Store a value in this nodes local storage
119 func (dht *IpfsDHT) handlePutValue(p *peer.Peer, pmes *DHTMessage) {
94 - panic("Not implemented.")
120 + dskey := ds.NewKey(pmes.GetKey())
121 + err := dht.datastore.Put(dskey, pmes.GetValue())
122 + if err != nil {
123 + // For now, just panic, handle this better later maybe
124 + panic(err)
125 + }
126 }
127
128 func (dht *IpfsDHT) handleFindNode(p *peer.Peer, pmes *DHTMessage) {
@@ -121,6 +152,17 @@ func (dht *IpfsDHT) ListenFor(mesid uint64) <-chan *swarm.Message {
152 return lchan
153 }
154
155 +func (dht *IpfsDHT) Unlisten(mesid uint64) {
156 + dht.listenLock.Lock()
157 + ch, ok := dht.listeners[mesid]
158 + if ok {
159 + delete(dht.listeners, mesid)
160 + }
161 + dht.listenLock.Unlock()
162 + close(ch)
163 +}
164 +
165 +
166 // Stop all communications from this node and shut down
167 func (dht *IpfsDHT) Halt() {
168 dht.shutdown <- struct{}{}