@cryptotaxi247 / kubo / commits / ec5276c29

really ugly impl of 'ipfs dht query' command

Jeromy committed Jan 22, 2015 at 07:59 UTC ec5276c29c55156960eb926be1604f9e9d7fd276
5 files changed +297 -75
core/commands/dht.go new
+108
@@ -0,0 +1,108 @@
1 +package commands
2 +
3 +import (
4 + "bytes"
5 + "errors"
6 + "fmt"
7 + "io"
8 + "time"
9 +
10 + cmds "github.com/jbenet/go-ipfs/commands"
11 + ipdht "github.com/jbenet/go-ipfs/routing/dht"
12 + u "github.com/jbenet/go-ipfs/util"
13 +)
14 +
15 +var DhtCmd = &cmds.Command{
16 + Helptext: cmds.HelpText{
17 + Tagline: "Issue commands directly through the DHT",
18 + ShortDescription: ``,
19 + },
20 +
21 + Subcommands: map[string]*cmds.Command{
22 + "query": queryDhtCmd,
23 + },
24 +}
25 +
26 +var queryDhtCmd = &cmds.Command{
27 + Helptext: cmds.HelpText{
28 + Tagline: "Run a 'findClosestPeers' query through the DHT",
29 + ShortDescription: ``,
30 + },
31 +
32 + Arguments: []cmds.Argument{
33 + cmds.StringArg("peerID", true, true, "The peerID to run the query against"),
34 + },
35 + Options: []cmds.Option{
36 + cmds.BoolOption("verbose", "v", "Write extra information"),
37 + },
38 + Run: func(req cmds.Request) (interface{}, error) {
39 + n, err := req.Context().GetNode()
40 + if err != nil {
41 + return nil, err
42 + }
43 +
44 + dht, ok := n.Routing.(*ipdht.IpfsDHT)
45 + if !ok {
46 + return nil, errors.New("Routing service was not a dht")
47 + }
48 +
49 + events := make(chan *ipdht.QueryEvent)
50 + closestPeers, err := dht.GetClosestPeers(req.Context().Context, u.Key(req.Arguments()[0]), events)
51 +
52 + go func() {
53 + defer close(events)
54 + for p := range closestPeers {
55 + events <- &ipdht.QueryEvent{
56 + ID: p,
57 + Type: ipdht.FinalPeer,
58 + }
59 + }
60 + }()
61 +
62 + outChan := make(chan interface{})
63 + go func() {
64 + defer close(outChan)
65 + for e := range events {
66 + outChan <- e
67 + }
68 + }()
69 + return outChan, nil
70 + },
71 + Marshalers: cmds.MarshalerMap{
72 + cmds.Text: func(res cmds.Response) (io.Reader, error) {
73 + outChan, ok := res.Output().(<-chan interface{})
74 + if !ok {
75 + return nil, u.ErrCast()
76 + }
77 +
78 + marshal := func(v interface{}) (io.Reader, error) {
79 + obj, ok := v.(*ipdht.QueryEvent)
80 + if !ok {
81 + return nil, u.ErrCast()
82 + }
83 +
84 + buf := new(bytes.Buffer)
85 + fmt.Fprintf(buf, "%s: ", time.Now().Format("15:04:05.000"))
86 + switch obj.Type {
87 + case ipdht.FinalPeer:
88 + fmt.Fprintf(buf, "%s\n", obj.ID)
89 + case ipdht.PeerResponse:
90 + fmt.Fprintf(buf, "* %s says use ", obj.ID)
91 + for _, p := range obj.Responses {
92 + fmt.Fprintf(buf, "%s ", p.ID)
93 + }
94 + fmt.Fprintln(buf)
95 + case ipdht.SendingQuery:
96 + fmt.Fprintf(buf, "* querying %s\n", obj.ID)
97 + }
98 + return buf, nil
99 + }
100 +
101 + return &cmds.ChannelMarshaler{
102 + Channel: outChan,
103 + Marshaler: marshal,
104 + }, nil
105 + },
106 + },
107 + Type: ipdht.QueryEvent{},
108 +}
core/commands/root.go
+1
@@ -71,6 +71,7 @@ var rootSubcommands = map[string]*cmds.Command{
71 "cat": CatCmd,
72 "commands": CommandsDaemonCmd,
73 "config": ConfigCmd,
74 + "dht": DhtCmd,
75 "diag": DiagCmd,
76 "id": IDCmd,
77 "log": LogCmd,
p2p/peer/peer.go
+30
@@ -3,6 +3,7 @@ package peer
3
4 import (
5 "encoding/hex"
6 + "encoding/json"
7 "fmt"
8
9 b58 "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-base58"
@@ -131,6 +132,35 @@ type PeerInfo struct {
132 Addrs []ma.Multiaddr
133 }
134
135 +func (pi *PeerInfo) MarshalJSON() ([]byte, error) {
136 + out := make(map[string]interface{})
137 + out["ID"] = IDB58Encode(pi.ID)
138 + var addrs []string
139 + for _, a := range pi.Addrs {
140 + addrs = append(addrs, a.String())
141 + }
142 + out["Addrs"] = addrs
143 + return json.Marshal(out)
144 +}
145 +
146 +func (pi *PeerInfo) UnmarshalJSON(b []byte) error {
147 + var data map[string]interface{}
148 + err := json.Unmarshal(b, &data)
149 + if err != nil {
150 + return err
151 + }
152 + pid, err := IDB58Decode(data["ID"].(string))
153 + if err != nil {
154 + return err
155 + }
156 + pi.ID = pid
157 + addrs := data["Addrs"].([]interface{})
158 + for _, a := range addrs {
159 + pi.Addrs = append(pi.Addrs, ma.StringCast(a.(string)))
160 + }
161 + return nil
162 +}
163 +
164 // IDSlice for sorting peers
165 type IDSlice []ID
166
routing/dht/lookup.go new
+156
@@ -0,0 +1,156 @@
1 +package dht
2 +
3 +import (
4 + "encoding/json"
5 +
6 + context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
7 +
8 + peer "github.com/jbenet/go-ipfs/p2p/peer"
9 + kb "github.com/jbenet/go-ipfs/routing/kbucket"
10 + u "github.com/jbenet/go-ipfs/util"
11 + errors "github.com/jbenet/go-ipfs/util/debugerror"
12 + pset "github.com/jbenet/go-ipfs/util/peerset"
13 +)
14 +
15 +type QueryEventType int
16 +
17 +const (
18 + SendingQuery QueryEventType = iota
19 + PeerResponse
20 + FinalPeer
21 +)
22 +
23 +type QueryEvent struct {
24 + ID peer.ID
25 + Type QueryEventType
26 + Responses []*peer.PeerInfo
27 +}
28 +
29 +func pointerizePeerInfos(pis []peer.PeerInfo) []*peer.PeerInfo {
30 + out := make([]*peer.PeerInfo, len(pis))
31 + for i, p := range pis {
32 + np := p
33 + out[i] = &np
34 + }
35 + return out
36 +}
37 +
38 +// Kademlia 'node lookup' operation. Returns a channel of the K closest peers
39 +// to the given key
40 +func (dht *IpfsDHT) GetClosestPeers(ctx context.Context, key u.Key, events chan<- *QueryEvent) (<-chan peer.ID, error) {
41 + e := log.EventBegin(ctx, "getClosestPeers", &key)
42 + tablepeers := dht.routingTable.NearestPeers(kb.ConvertKey(key), AlphaValue)
43 + if len(tablepeers) == 0 {
44 + return nil, errors.Wrap(kb.ErrLookupFailure)
45 + }
46 +
47 + out := make(chan peer.ID, KValue)
48 + peerset := pset.NewLimited(KValue)
49 +
50 + for _, p := range tablepeers {
51 + select {
52 + case out <- p:
53 + case <-ctx.Done():
54 + return nil, ctx.Err()
55 + }
56 + peerset.Add(p)
57 + }
58 +
59 + query := dht.newQuery(key, func(ctx context.Context, p peer.ID) (*dhtQueryResult, error) {
60 + // For DHT query command
61 + select {
62 + case events <- &QueryEvent{
63 + Type: SendingQuery,
64 + ID: p,
65 + }:
66 + }
67 +
68 + closer, err := dht.closerPeersSingle(ctx, key, p)
69 + if err != nil {
70 + log.Errorf("error getting closer peers: %s", err)
71 + return nil, err
72 + }
73 +
74 + var filtered []peer.PeerInfo
75 + for _, clp := range closer {
76 + if kb.Closer(clp, dht.self, key) && peerset.TryAdd(clp) {
77 + select {
78 + case out <- clp:
79 + log.Error("Sending out peer: %s", clp.Pretty())
80 + case <-ctx.Done():
81 + return nil, ctx.Err()
82 + }
83 + filtered = append(filtered, dht.peerstore.PeerInfo(clp))
84 + }
85 + }
86 + log.Errorf("filtered: %v", filtered)
87 +
88 + // For DHT query command
89 + select {
90 + case events <- &QueryEvent{
91 + Type: PeerResponse,
92 + ID: p,
93 + Responses: pointerizePeerInfos(filtered),
94 + }:
95 + }
96 +
97 + return &dhtQueryResult{closerPeers: filtered}, nil
98 + })
99 +
100 + go func() {
101 + defer close(out)
102 + defer e.Done()
103 + // run it!
104 + _, err := query.Run(ctx, tablepeers)
105 + if err != nil {
106 + log.Debugf("closestPeers query run error: %s", err)
107 + }
108 + }()
109 +
110 + return out, nil
111 +}
112 +
113 +func (dht *IpfsDHT) closerPeersSingle(ctx context.Context, key u.Key, p peer.ID) ([]peer.ID, error) {
114 + pmes, err := dht.findPeerSingle(ctx, p, peer.ID(key))
115 + if err != nil {
116 + return nil, err
117 + }
118 +
119 + var out []peer.ID
120 + for _, pbp := range pmes.GetCloserPeers() {
121 + pid := peer.ID(pbp.GetId())
122 + if pid != dht.self { // dont add self
123 + dht.peerstore.AddAddresses(pid, pbp.Addresses())
124 + out = append(out, pid)
125 + }
126 + }
127 + return out, nil
128 +}
129 +
130 +func (qe *QueryEvent) MarshalJSON() ([]byte, error) {
131 + out := make(map[string]interface{})
132 + out["ID"] = peer.IDB58Encode(qe.ID)
133 + out["Type"] = int(qe.Type)
134 + out["Responses"] = qe.Responses
135 + return json.Marshal(out)
136 +}
137 +
138 +func (qe *QueryEvent) UnmarshalJSON(b []byte) error {
139 + temp := struct {
140 + ID string
141 + Type int
142 + Responses []*peer.PeerInfo
143 + }{}
144 + err := json.Unmarshal(b, &temp)
145 + if err != nil {
146 + return err
147 + }
148 + pid, err := peer.IDB58Decode(temp.ID)
149 + if err != nil {
150 + return err
151 + }
152 + qe.ID = pid
153 + qe.Type = QueryEventType(temp.Type)
154 + qe.Responses = temp.Responses
155 + return nil
156 +}
routing/dht/routing.go
+2 -75
@@ -48,7 +48,7 @@ func (dht *IpfsDHT) PutValue(ctx context.Context, key u.Key, value []byte) error
48 return err
49 }
50
51 - pchan, err := dht.getClosestPeers(ctx, key)
51 + pchan, err := dht.GetClosestPeers(ctx, key, nil)
52 if err != nil {
53 return err
54 }
@@ -134,7 +134,7 @@ func (dht *IpfsDHT) Provide(ctx context.Context, key u.Key) error {
134 // add self locally
135 dht.providers.AddProvider(key, dht.self)
136
137 - peers, err := dht.getClosestPeers(ctx, key)
137 + peers, err := dht.GetClosestPeers(ctx, key, nil)
138 if err != nil {
139 return err
140 }
@@ -164,79 +164,6 @@ func (dht *IpfsDHT) FindProviders(ctx context.Context, key u.Key) ([]peer.PeerIn
164 return providers, nil
165 }
166
167 -// Kademlia 'node lookup' operation. Returns a channel of the K closest peers
168 -// to the given key
169 -func (dht *IpfsDHT) getClosestPeers(ctx context.Context, key u.Key) (<-chan peer.ID, error) {
170 - e := log.EventBegin(ctx, "getClosestPeers", &key)
171 - tablepeers := dht.routingTable.ListPeers()
172 - if len(tablepeers) == 0 {
173 - return nil, errors.Wrap(kb.ErrLookupFailure)
174 - }
175 -
176 - out := make(chan peer.ID, KValue)
177 - peerset := pset.NewLimited(KValue)
178 -
179 - for _, p := range tablepeers {
180 - select {
181 - case out <- p:
182 - case <-ctx.Done():
183 - return nil, ctx.Err()
184 - }
185 - peerset.Add(p)
186 - }
187 -
188 - query := dht.newQuery(key, func(ctx context.Context, p peer.ID) (*dhtQueryResult, error) {
189 - closer, err := dht.closerPeersSingle(ctx, key, p)
190 - if err != nil {
191 - log.Errorf("error getting closer peers: %s", err)
192 - return nil, err
193 - }
194 -
195 - var filtered []peer.PeerInfo
196 - for _, p := range closer {
197 - if kb.Closer(p, dht.self, key) && peerset.TryAdd(p) {
198 - select {
199 - case out <- p:
200 - case <-ctx.Done():
201 - return nil, ctx.Err()
202 - }
203 - filtered = append(filtered, dht.peerstore.PeerInfo(p))
204 - }
205 - }
206 -
207 - return &dhtQueryResult{closerPeers: filtered}, nil
208 - })
209 -
210 - go func() {
211 - defer close(out)
212 - defer e.Done()
213 - // run it!
214 - _, err := query.Run(ctx, tablepeers)
215 - if err != nil {
216 - log.Debugf("closestPeers query run error: %s", err)
217 - }
218 - }()
219 -
220 - return out, nil
221 -}
222 -
223 -func (dht *IpfsDHT) closerPeersSingle(ctx context.Context, key u.Key, p peer.ID) ([]peer.ID, error) {
224 - pmes, err := dht.findPeerSingle(ctx, p, peer.ID(key))
225 - if err != nil {
226 - return nil, err
227 - }
228 -
229 - var out []peer.ID
230 - for _, pbp := range pmes.GetCloserPeers() {
231 - pid := peer.ID(pbp.GetId())
232 - if pid != dht.self { // dont add self
233 - dht.peerstore.AddAddresses(pid, pbp.Addresses())
234 - out = append(out, pid)
235 - }
236 - }
237 - return out, nil
238 -}
239 -
167 // FindProvidersAsync is the same thing as FindProviders, but returns a channel.
168 // Peers will be returned on the channel as soon as they are found, even before
169 // the search query completes.