@cryptotaxi247 / kubo / commits / 375a38c5f

add basic publish command, needs polish

Jeromy committed Sep 25, 2014 at 15:10 UTC 375a38c5f77aea82ef5677a8bae5e7d1d7e4118b
7 files changed +30 -6
cmd/ipfs/ipfs.go
+2
@@ -4,6 +4,7 @@ import (
4 "errors"
5 "fmt"
6 "os"
7 + "runtime/pprof"
8
9 "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/gonuts/flag"
10 "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/commander"
@@ -51,6 +52,7 @@ Use "ipfs help <command>" for more information about a command.
52 cmdIpfsInit,
53 cmdIpfsServe,
54 cmdIpfsRun,
55 + cmdIpfsPub,
56 },
57 Flag: *flag.NewFlagSet("ipfs", flag.ExitOnError),
58 }
cmd/ipfs/run.go
+1 -1
@@ -41,7 +41,7 @@ func runCmd(c *commander.Command, inp []string) error {
41 return err
42 }
43
44 - dl, err := daemon.NewDaemonListener(n, maddr)
44 + dl, err := daemon.NewDaemonListener(n, maddr, conf)
45 if err != nil {
46 return err
47 }
daemon/daemon.go
+2
@@ -131,6 +131,8 @@ func (dl *DaemonListener) handleConnection(conn net.Conn) {
131 err = commands.Ls(dl.node, command.Args, command.Opts, conn)
132 case "pin":
133 err = commands.Pin(dl.node, command.Args, command.Opts, conn)
134 + case "publish":
135 + err = commands.Publish(dl.node, command.Args, command.Opts, conn)
136 default:
137 err = fmt.Errorf("Invalid Command: '%s'", command.Command)
138 }
namesys/publisher.go
+14 -2
@@ -1,6 +1,8 @@
1 package namesys
2
3 import (
4 + "time"
5 +
6 "code.google.com/p/go.net/context"
7 "code.google.com/p/goprotobuf/proto"
8
@@ -15,8 +17,16 @@ type IpnsPublisher struct {
17 routing routing.IpfsRouting
18 }
19
20 +func NewPublisher(dag *mdag.DAGService, route routing.IpfsRouting) *IpnsPublisher {
21 + return &IpnsPublisher{
22 + dag: dag,
23 + routing: route,
24 + }
25 +}
26 +
27 // Publish accepts a keypair and a value,
28 func (p *IpnsPublisher) Publish(k ci.PrivKey, value u.Key) error {
29 + log.Debug("namesys: Publish %s", value.Pretty())
30 ctx := context.TODO()
31 data, err := CreateEntryData(k, value)
32 if err != nil {
@@ -40,13 +50,15 @@ func (p *IpnsPublisher) Publish(k ci.PrivKey, value u.Key) error {
50 }
51
52 // Store associated public key
43 - err = p.routing.PutValue(ctx, u.Key(nameb), pkbytes)
53 + timectx, _ := context.WithDeadline(ctx, time.Now().Add(time.Second*4))
54 + err = p.routing.PutValue(timectx, u.Key(nameb), pkbytes)
55 if err != nil {
56 return err
57 }
58
59 // Store ipns entry at h("/ipns/"+b58(h(pubkey)))
49 - err = p.routing.PutValue(ctx, u.Key(ipnskey), data)
60 + timectx, _ = context.WithDeadline(ctx, time.Now().Add(time.Second*4))
61 + err = p.routing.PutValue(timectx, u.Key(ipnskey), data)
62 if err != nil {
63 return err
64 }
routing/dht/dht.go
+3
@@ -13,6 +13,7 @@ import (
13 peer "github.com/jbenet/go-ipfs/peer"
14 kb "github.com/jbenet/go-ipfs/routing/kbucket"
15 u "github.com/jbenet/go-ipfs/util"
16 + "github.com/op/go-logging"
17
18 context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
19 ds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/datastore.go"
@@ -21,6 +22,8 @@ import (
22 "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/goprotobuf/proto"
23 )
24
25 +var log = logging.MustGetLogger("dht")
26 +
27 // TODO. SEE https://github.com/jbenet/node-ipfs/blob/master/submodules/ipfs-dht/index.js
28
29 // IpfsDHT is an implementation of Kademlia with Coral and S/Kademlia modifications.
routing/dht/query.go
+6 -1
@@ -101,6 +101,11 @@ func newQueryRunner(ctx context.Context, q *dhtQuery) *dhtQueryRunner {
101 }
102
103 func (r *dhtQueryRunner) Run(peers []*peer.Peer) (*dhtQueryResult, error) {
104 + log.Debug("Run query with %d peers.", len(peers))
105 + if len(peers) == 0 {
106 + log.Warning("Running query with no peers!")
107 + return nil, nil
108 + }
109 // setup concurrency rate limiting
110 for i := 0; i < r.query.concurrency; i++ {
111 r.rateLimit <- struct{}{}
@@ -164,7 +169,7 @@ func (r *dhtQueryRunner) addPeerToQuery(next *peer.Peer, benchmark *peer.Peer) {
169 r.peersSeen[next.Key()] = next
170 r.Unlock()
171
167 - u.DOut("adding peer to query: %v\n", next.ID.Pretty())
172 + log.Debug("adding peer to query: %v\n", next.ID.Pretty())
173
174 // do this after unlocking to prevent possible deadlocks.
175 r.peersRemaining.Increment(1)
routing/dht/routing.go
+2 -2
@@ -18,6 +18,7 @@ import (
18 // PutValue adds value corresponding to given Key.
19 // This is the top level "Store" operation of the DHT
20 func (dht *IpfsDHT) PutValue(ctx context.Context, key u.Key, value []byte) error {
21 + log.Debug("[%s] PutValue %v %v", dht.self.ID.Pretty(), key, value)
22 err := dht.putLocal(key, value)
23 if err != nil {
24 return err
@@ -30,7 +31,7 @@ func (dht *IpfsDHT) PutValue(ctx context.Context, key u.Key, value []byte) error
31 }
32
33 query := newQuery(key, func(ctx context.Context, p *peer.Peer) (*dhtQueryResult, error) {
33 - u.DOut("[%s] PutValue qry part %v\n", dht.self.ID.Pretty(), p.ID.Pretty())
34 + log.Debug("[%s] PutValue qry part %v", dht.self.ID.Pretty(), p.ID.Pretty())
35 err := dht.putValueToNetwork(ctx, p, string(key), value)
36 if err != nil {
37 return nil, err
@@ -39,7 +40,6 @@ func (dht *IpfsDHT) PutValue(ctx context.Context, key u.Key, value []byte) error
40 })
41
42 _, err = query.Run(ctx, peers)
42 - u.DOut("[%s] PutValue %v %v\n", dht.self.ID.Pretty(), key, value)
43 return err
44 }
45