@cryptotaxi247 / kubo / commits / eaa433f27

cmds: implement ipfs dht provide command

License: MIT Signed-off-by: Jeromy <why@ipfs.io>

Jeromy committed Aug 19, 2016 at 19:52 UTC eaa433f27bcafe591f8b2ef02e5f93398ab9e629
2 files changed +116
core/commands/dht.go
+111
@@ -31,6 +31,7 @@ var DhtCmd = &cmds.Command{
31 "findpeer": findPeerDhtCmd,
32 "get": getValueDhtCmd,
33 "put": putValueDhtCmd,
34 + "provide": provideRefDhtCmd,
35 },
36 }
37
@@ -227,6 +228,116 @@ var findProvidersDhtCmd = &cmds.Command{
228 Type: notif.QueryEvent{},
229 }
230
231 +var provideRefDhtCmd = &cmds.Command{
232 + Helptext: cmds.HelpText{
233 + Tagline: "Announce to the network that you are providing given values.",
234 + },
235 +
236 + Arguments: []cmds.Argument{
237 + cmds.StringArg("key", true, true, "The key to find providers for.").EnableStdin(),
238 + },
239 + Options: []cmds.Option{
240 + cmds.BoolOption("verbose", "v", "Print extra information.").Default(false),
241 + },
242 + Run: func(req cmds.Request, res cmds.Response) {
243 + n, err := req.InvocContext().GetNode()
244 + if err != nil {
245 + res.SetError(err, cmds.ErrNormal)
246 + return
247 + }
248 +
249 + dht, ok := n.Routing.(*ipdht.IpfsDHT)
250 + if !ok {
251 + res.SetError(ErrNotDHT, cmds.ErrNormal)
252 + return
253 + }
254 +
255 + var keys []key.Key
256 + for _, arg := range req.Arguments() {
257 + k := key.B58KeyDecode(arg)
258 + if k == "" {
259 + res.SetError(fmt.Errorf("incorrectly formatted key: ", arg), cmds.ErrNormal)
260 + return
261 + }
262 +
263 + has, err := n.Blockstore.Has(k)
264 + if err != nil {
265 + res.SetError(err, cmds.ErrNormal)
266 + return
267 + }
268 +
269 + if !has {
270 + res.SetError(fmt.Errorf("block %s not found locally, cannot provide", k), cmds.ErrNormal)
271 + return
272 + }
273 +
274 + keys = append(keys, k)
275 + }
276 +
277 + outChan := make(chan interface{})
278 + res.SetOutput((<-chan interface{})(outChan))
279 +
280 + events := make(chan *notif.QueryEvent)
281 + ctx := notif.RegisterForQueryEvents(req.Context(), events)
282 +
283 + go func() {
284 + defer close(outChan)
285 + for e := range events {
286 + outChan <- e
287 + }
288 + }()
289 +
290 + go func() {
291 + defer close(outChan)
292 + for _, k := range keys {
293 + err := dht.Provide(ctx, k)
294 + if err != nil {
295 + notif.PublishQueryEvent(ctx, &notif.QueryEvent{
296 + Type: notif.QueryError,
297 + Extra: err.Error(),
298 + })
299 + return
300 + }
301 + }
302 + }()
303 + },
304 + Marshalers: cmds.MarshalerMap{
305 + cmds.Text: func(res cmds.Response) (io.Reader, error) {
306 + outChan, ok := res.Output().(<-chan interface{})
307 + if !ok {
308 + return nil, u.ErrCast()
309 + }
310 +
311 + verbose, _, _ := res.Request().Option("v").Bool()
312 + pfm := pfuncMap{
313 + notif.FinalPeer: func(obj *notif.QueryEvent, out io.Writer, verbose bool) {
314 + if verbose {
315 + fmt.Fprintf(out, "sending provider record to peer %s\n", obj.ID)
316 + }
317 + },
318 + }
319 +
320 + marshal := func(v interface{}) (io.Reader, error) {
321 + obj, ok := v.(*notif.QueryEvent)
322 + if !ok {
323 + return nil, u.ErrCast()
324 + }
325 +
326 + buf := new(bytes.Buffer)
327 + printEvent(obj, buf, verbose, pfm)
328 + return buf, nil
329 + }
330 +
331 + return &cmds.ChannelMarshaler{
332 + Channel: outChan,
333 + Marshaler: marshal,
334 + Res: res,
335 + }, nil
336 + },
337 + },
338 + Type: notif.QueryEvent{},
339 +}
340 +
341 var findPeerDhtCmd = &cmds.Command{
342 Helptext: cmds.HelpText{
343 Tagline: "Query the DHT for all of the multiaddresses associated with a Peer ID.",
routing/dht/routing.go
+5
@@ -263,6 +263,10 @@ func (dht *IpfsDHT) Provide(ctx context.Context, key key.Key) error {
263 go func(p peer.ID) {
264 defer wg.Done()
265 log.Debugf("putProvider(%s, %s)", key, p)
266 + notif.PublishQueryEvent(ctx, &notif.QueryEvent{
267 + Type: notif.FinalPeer,
268 + ID: p,
269 + })
270 err := dht.sendMessage(ctx, p, mes)
271 if err != nil {
272 log.Debug(err)
@@ -272,6 +276,7 @@ func (dht *IpfsDHT) Provide(ctx context.Context, key key.Key) error {
276 wg.Wait()
277 return nil
278 }
279 +
280 func (dht *IpfsDHT) makeProvRecord(skey key.Key) (*pb.Message, error) {
281 pi := pstore.PeerInfo{
282 ID: dht.self,