@cryptotaxi247 / kubo / commits / 248e06f75

working towards Providers implementation

Jeromy committed Aug 3, 2014 at 21:46 UTC 248e06f759288d07100412c5b1838ce6600cb3ac
3 files changed +112 -10
routing/dht/dht.go
+40 -2
@@ -3,6 +3,7 @@ package dht
3 import (
4 "sync"
5 "time"
6 + "encoding/json"
7
8 peer "github.com/jbenet/go-ipfs/peer"
9 swarm "github.com/jbenet/go-ipfs/swarm"
@@ -31,6 +32,10 @@ type IpfsDHT struct {
32 // Local data
33 datastore ds.Datastore
34
35 + // Map keys to peers that can provide their value
36 + // TODO: implement a TTL on each of these keys
37 + providers map[u.Key][]*peer.Peer
38 +
39 // map of channels waiting for reply messages
40 listeners map[uint64]chan *swarm.Message
41 listenLock sync.RWMutex
@@ -185,11 +190,44 @@ func (dht *IpfsDHT) handleFindNode(p *peer.Peer, pmes *DHTMessage) {
190 }
191
192 func (dht *IpfsDHT) handleGetProviders(p *peer.Peer, pmes *DHTMessage) {
188 - panic("Not implemented.")
193 + providers := dht.providers[u.Key(pmes.GetKey())]
194 + if providers == nil || len(providers) == 0 {
195 + // ?????
196 + }
197 +
198 + var addrs []string
199 + for _,prov := range providers {
200 + ma := prov.NetAddress("tcp")
201 + str,err := ma.String()
202 + if err != nil {
203 + u.PErr("Error: %s", err)
204 + continue
205 + }
206 +
207 + addrs = append(addrs, str)
208 + }
209 +
210 + data,err := json.Marshal(addrs)
211 + if err != nil {
212 + panic(err)
213 + }
214 +
215 + resp := pDHTMessage{
216 + Type: DHTMessage_GET_PROVIDERS,
217 + Key: pmes.GetKey(),
218 + Value: data,
219 + Id: pmes.GetId(),
220 + }
221 +
222 + mes := swarm.NewMessage(p, resp.ToProtobuf())
223 + dht.network.Chan.Outgoing <-mes
224 }
225
226 func (dht *IpfsDHT) handleAddProvider(p *peer.Peer, pmes *DHTMessage) {
192 - panic("Not implemented.")
227 + //TODO: need to implement TTLs on providers
228 + key := u.Key(pmes.GetKey())
229 + parr := dht.providers[key]
230 + dht.providers[key] = append(parr, p)
231 }
232
233
routing/dht/routing.go
+71 -7
@@ -11,6 +11,9 @@ import (
11 u "github.com/jbenet/go-ipfs/util"
12 )
13
14 +// Pool size is the number of nodes used for group find/set RPC calls
15 +var PoolSize = 6
16 +
17 // TODO: determine a way of creating and managing message IDs
18 func GenerateMessageID() uint64 {
19 return uint64(rand.Uint32()) << 32 & uint64(rand.Uint32())
@@ -25,8 +28,6 @@ func (s *IpfsDHT) PutValue(key u.Key, value []byte) error {
28 var p *peer.Peer
29 p = s.routes.NearestPeer(convertKey(key))
30 if p == nil {
28 - u.POut("nbuckets: %d", len(s.routes.Buckets))
29 - u.POut("%d", s.routes.Buckets[0].Len())
31 panic("Table returned nil peer!")
32 }
33
@@ -64,7 +65,7 @@ func (s *IpfsDHT) GetValue(key u.Key, timeout time.Duration) ([]byte, error) {
65 timeup := time.After(timeout)
66 select {
67 case <-timeup:
67 - // TODO: unregister listener
68 + s.Unlisten(pmes.Id)
69 return nil, u.ErrTimeout
70 case resp := <-response_chan:
71 pmes_out := new(DHTMessage)
@@ -81,17 +82,80 @@ func (s *IpfsDHT) GetValue(key u.Key, timeout time.Duration) ([]byte, error) {
82
83 // Announce that this node can provide value for given key
84 func (s *IpfsDHT) Provide(key u.Key) error {
84 - return u.ErrNotImplemented
85 + peers := s.routes.NearestPeers(convertKey(key), PoolSize)
86 + if len(peers) == 0 {
87 + //return an error
88 + }
89 +
90 + pmes := pDHTMessage{
91 + Type: DHTMessage_ADD_PROVIDER,
92 + Key: string(key),
93 + }
94 + pbmes := pmes.ToProtobuf()
95 +
96 + for _,p := range peers {
97 + mes := swarm.NewMessage(p, pbmes)
98 + s.network.Chan.Outgoing <-mes
99 + }
100 + return nil
101 }
102
103 // FindProviders searches for peers who can provide the value for given key.
88 -func (s *IpfsDHT) FindProviders(key u.Key, timeout time.Duration) (*peer.Peer, error) {
89 - return nil, u.ErrNotImplemented
104 +func (s *IpfsDHT) FindProviders(key u.Key, timeout time.Duration) ([]*peer.Peer, error) {
105 + p := s.routes.NearestPeer(convertKey(key))
106 +
107 + pmes := pDHTMessage{
108 + Type: DHTMessage_GET_PROVIDERS,
109 + Key: string(key),
110 + Id: GenerateMessageID(),
111 + }
112 +
113 + mes := swarm.NewMessage(p, pmes.ToProtobuf())
114 +
115 + listen_chan := s.ListenFor(pmes.Id)
116 + s.network.Chan.Outgoing <-mes
117 + after := time.After(timeout)
118 + select {
119 + case <-after:
120 + s.Unlisten(pmes.Id)
121 + return nil, u.ErrTimeout
122 + case resp := <-listen_chan:
123 + pmes_out := new(DHTMessage)
124 + err := proto.Unmarshal(resp.Data, pmes_out)
125 + if err != nil {
126 + return nil, err
127 + }
128 + panic("Not yet implemented.")
129 + }
130 }
131
132 // Find specific Peer
133
134 // FindPeer searches for a peer with given ID.
135 func (s *IpfsDHT) FindPeer(id peer.ID, timeout time.Duration) (*peer.Peer, error) {
96 - return nil, u.ErrNotImplemented
136 + p := s.routes.NearestPeer(convertPeerID(id))
137 +
138 + pmes := pDHTMessage{
139 + Type: DHTMessage_FIND_NODE,
140 + Key: string(id),
141 + Id: GenerateMessageID(),
142 + }
143 +
144 + mes := swarm.NewMessage(p, pmes.ToProtobuf())
145 +
146 + listen_chan := s.ListenFor(pmes.Id)
147 + s.network.Chan.Outgoing <-mes
148 + after := time.After(timeout)
149 + select {
150 + case <-after:
151 + s.Unlisten(pmes.Id)
152 + return nil, u.ErrTimeout
153 + case resp := <-listen_chan:
154 + pmes_out := new(DHTMessage)
155 + err := proto.Unmarshal(resp.Data, pmes_out)
156 + if err != nil {
157 + return nil, err
158 + }
159 + panic("Not yet implemented.")
160 + }
161 }
routing/routing.go
+1 -1
@@ -25,7 +25,7 @@ type IpfsRouting interface {
25 Provide(key u.Key) error
26
27 // FindProviders searches for peers who can provide the value for given key.
28 - FindProviders(key u.Key, timeout time.Duration) (*peer.Peer, error)
28 + FindProviders(key u.Key, timeout time.Duration) ([]*peer.Peer, error)
29
30 // Find specific Peer
31