@cryptotaxi247 / kubo / commits / 3a76ef047

a little error handling and some work on providers

Jeromy committed Aug 5, 2014 at 09:38 UTC 3a76ef047836b4585e90a0d7893174751fa51aaa
3 files changed +44 -5
routing/dht/dht.go
+12 -4
@@ -35,6 +35,7 @@ type IpfsDHT struct {
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 + providerLock sync.RWMutex
39
40 // map of channels waiting for reply messages
41 listeners map[uint64]chan *swarm.Message
@@ -46,6 +47,9 @@ type IpfsDHT struct {
47
48 // Create a new DHT object with the given peer as the 'local' host
49 func NewDHT(p *peer.Peer) (*IpfsDHT, error) {
50 + if p == nil {
51 + panic("Tried to create new dht with nil peer")
52 + }
53 network := swarm.NewSwarm(p)
54 err := network.Listen()
55 if err != nil {
@@ -68,24 +72,27 @@ func (dht *IpfsDHT) Start() {
72 }
73
74 // Connect to a new peer at the given address
71 -func (dht *IpfsDHT) Connect(addr *ma.Multiaddr) error {
75 +func (dht *IpfsDHT) Connect(addr *ma.Multiaddr) (*peer.Peer, error) {
76 + if addr == nil {
77 + panic("addr was nil!")
78 + }
79 peer := new(peer.Peer)
80 peer.AddAddress(addr)
81
82 conn,err := swarm.Dial("tcp", peer)
83 if err != nil {
77 - return err
84 + return nil, err
85 }
86
87 err = identify.Handshake(dht.self, peer, conn.Incoming.MsgChan, conn.Outgoing.MsgChan)
88 if err != nil {
82 - return err
89 + return nil, err
90 }
91
92 dht.network.StartConn(conn)
93
94 dht.routes.Update(peer)
88 - return nil
95 + return peer, nil
96 }
97
98 // Read in all messages from swarm and handle them appropriately
@@ -195,6 +202,7 @@ func (dht *IpfsDHT) handleGetProviders(p *peer.Peer, pmes *DHTMessage) {
202 // ?????
203 }
204
205 + // This is just a quick hack, formalize method of sending addrs later
206 var addrs []string
207 for _,prov := range providers {
208 ma := prov.NetAddress("tcp")
routing/dht/routing.go
+29 -1
@@ -3,9 +3,12 @@ package dht
3 import (
4 "math/rand"
5 "time"
6 + "encoding/json"
7
8 proto "code.google.com/p/goprotobuf/proto"
9
10 + ma "github.com/jbenet/go-multiaddr"
11 +
12 peer "github.com/jbenet/go-ipfs/peer"
13 swarm "github.com/jbenet/go-ipfs/swarm"
14 u "github.com/jbenet/go-ipfs/util"
@@ -125,7 +128,32 @@ func (s *IpfsDHT) FindProviders(key u.Key, timeout time.Duration) ([]*peer.Peer,
128 if err != nil {
129 return nil, err
130 }
128 - panic("Not yet implemented.")
131 + var addrs map[string]string
132 + err := json.Unmarshal(pmes_out.GetValue(), &addrs)
133 + if err != nil {
134 + return nil, err
135 + }
136 +
137 + for key,addr := range addrs {
138 + p := s.network.Find(u.Key(key))
139 + if p == nil {
140 + maddr,err := ma.NewMultiaddr(addr)
141 + if err != nil {
142 + u.PErr("error connecting to new peer: %s", err)
143 + continue
144 + }
145 + p, err := s.Connect(maddr)
146 + if err != nil {
147 + u.PErr("error connecting to new peer: %s", err)
148 + continue
149 + }
150 + }
151 + s.providerLock.Lock()
152 + prov_arr := s.providers[key]
153 + s.providers[key] = append(prov_arr, p)
154 + s.providerLock.Unlock()
155 + }
156 +
157 }
158 }
159
swarm/swarm.go
+3
@@ -290,3 +290,6 @@ Loop:
290 delete(s.conns, conn.Peer.Key())
291 s.connsLock.Unlock()
292 }
293 +
294 +func (s *Swarm) Find(addr *ma.Multiaddr) {
295 +}