@cryptotaxi247 / kubo / commits / b2a218816

initial hack at turning ipfs into a daemon, just implemented simple rpc at this point

Jeromy committed Sep 9, 2014 at 05:11 UTC b2a218816dfa619400e15fb9958edb67db8b610d
9 files changed +224 -261
cmd/ipfs/add.go
+12 -91
@@ -2,21 +2,14 @@ package main
2
3 import (
4 "fmt"
5 - "io/ioutil"
5 "os"
7 - "path/filepath"
6
7 "github.com/gonuts/flag"
8 "github.com/jbenet/commander"
11 - core "github.com/jbenet/go-ipfs/core"
12 - importer "github.com/jbenet/go-ipfs/importer"
13 - dag "github.com/jbenet/go-ipfs/merkledag"
9 + daemon "github.com/jbenet/go-ipfs/daemon"
10 u "github.com/jbenet/go-ipfs/util"
11 )
12
17 -// Error indicating the max depth has been exceded.
18 -var ErrDepthLimitExceeded = fmt.Errorf("depth limit exceeded")
19 -
13 var cmdIpfsAdd = &commander.Command{
14 UsageLine: "add",
15 Short: "Add an object to ipfs.",
@@ -41,92 +34,20 @@ func addCmd(c *commander.Command, inp []string) error {
34 return nil
35 }
36
44 - n, err := localNode(false)
37 + cmd := daemon.NewCommand()
38 + cmd.Command = "add"
39 + fmt.Println(inp)
40 + cmd.Args = inp
41 + cmd.Opts["r"] = c.Flag.Lookup("r").Value.Get()
42 + err := daemon.SendCommand(cmd, "localhost:12345")
43 if err != nil {
46 - return err
47 - }
48 -
49 - recursive := c.Flag.Lookup("r").Value.Get().(bool)
50 - var depth int
51 - if recursive {
52 - depth = -1
53 - } else {
54 - depth = 1
55 - }
56 -
57 - for _, fpath := range inp {
58 - _, err := addPath(n, fpath, depth)
44 + // Do locally
45 + n, err := localNode(false)
46 if err != nil {
60 - if !recursive {
61 - return fmt.Errorf("%s is a directory. Use -r to add recursively", fpath)
62 - }
63 -
64 - u.PErr("error adding %s: %v\n", fpath, err)
47 + return err
48 }
66 - }
67 - return err
68 -}
69 -
70 -func addPath(n *core.IpfsNode, fpath string, depth int) (*dag.Node, error) {
71 - if depth == 0 {
72 - return nil, ErrDepthLimitExceeded
73 - }
49
75 - fi, err := os.Stat(fpath)
76 - if err != nil {
77 - return nil, err
78 - }
79 -
80 - if fi.IsDir() {
81 - return addDir(n, fpath, depth)
82 - }
83 -
84 - return addFile(n, fpath, depth)
85 -}
86 -
87 -func addDir(n *core.IpfsNode, fpath string, depth int) (*dag.Node, error) {
88 - tree := &dag.Node{Data: dag.FolderPBData()}
89 -
90 - files, err := ioutil.ReadDir(fpath)
91 - if err != nil {
92 - return nil, err
50 + daemon.ExecuteCommand(cmd, n, os.Stdout)
51 }
94 -
95 - // construct nodes for containing files.
96 - for _, f := range files {
97 - fp := filepath.Join(fpath, f.Name())
98 - nd, err := addPath(n, fp, depth-1)
99 - if err != nil {
100 - return nil, err
101 - }
102 -
103 - if err = tree.AddNodeLink(f.Name(), nd); err != nil {
104 - return nil, err
105 - }
106 - }
107 -
108 - return tree, addNode(n, tree, fpath)
109 -}
110 -
111 -func addFile(n *core.IpfsNode, fpath string, depth int) (*dag.Node, error) {
112 - root, err := importer.NewDagFromFile(fpath)
113 - if err != nil {
114 - return nil, err
115 - }
116 -
117 - return root, addNode(n, root, fpath)
118 -}
119 -
120 -// addNode adds the node to the graph + local storage
121 -func addNode(n *core.IpfsNode, nd *dag.Node, fpath string) error {
122 - // add the file to the graph + local storage
123 - err := n.DAG.AddRecursive(nd)
124 - if err != nil {
125 - return err
126 - }
127 -
128 - u.POut("added %s\n", fpath)
129 -
130 - // ensure we keep it. atm no-op
131 - return n.PinDagNode(nd)
52 + return nil
53 }
cmd/ipfs/cat.go
+8 -20
@@ -1,13 +1,11 @@
1 package main
2
3 import (
4 - "fmt"
5 - "io"
4 "os"
5
6 "github.com/gonuts/flag"
7 "github.com/jbenet/commander"
10 - dag "github.com/jbenet/go-ipfs/merkledag"
8 + "github.com/jbenet/go-ipfs/daemon"
9 u "github.com/jbenet/go-ipfs/util"
10 )
11
@@ -29,28 +27,18 @@ func catCmd(c *commander.Command, inp []string) error {
27 return nil
28 }
29
32 - n, err := localNode(false)
33 - if err != nil {
34 - return err
35 - }
30 + com := daemon.NewCommand()
31 + com.Command = "cat"
32 + com.Args = inp
33
37 - for _, fn := range inp {
38 - nd, err := n.Resolver.ResolvePath(fn)
34 + err := daemon.SendCommand(com, "localhost:12345")
35 + if err != nil {
36 + n, err := localNode(false)
37 if err != nil {
38 return err
39 }
40
43 - read, err := dag.NewDagReader(nd, n.DAG)
44 - if err != nil {
45 - fmt.Println(err)
46 - continue
47 - }
48 -
49 - _, err = io.Copy(os.Stdout, read)
50 - if err != nil {
51 - fmt.Println(err)
52 - continue
53 - }
41 + daemon.ExecuteCommand(com, n, os.Stdout)
42 }
43 return nil
44 }
cmd/ipfs/ls.go
+13 -9
@@ -1,8 +1,12 @@
1 package main
2
3 import (
4 + "fmt"
5 + "os"
6 +
7 "github.com/gonuts/flag"
8 "github.com/jbenet/commander"
9 + "github.com/jbenet/go-ipfs/daemon"
10 u "github.com/jbenet/go-ipfs/util"
11 )
12
@@ -27,20 +31,20 @@ func lsCmd(c *commander.Command, inp []string) error {
31 return nil
32 }
33
30 - n, err := localNode(false)
34 + fmt.Println("hello")
35 + com := daemon.NewCommand()
36 + com.Command = "ls"
37 + com.Args = inp
38 + err := daemon.SendCommand(com, "localhost:12345")
39 if err != nil {
32 - return err
33 - }
34 -
35 - for _, fn := range inp {
36 - nd, err := n.Resolver.ResolvePath(fn)
40 + fmt.Println(err)
41 + n, err := localNode(false)
42 if err != nil {
43 return err
44 }
45
41 - for _, link := range nd.Links {
42 - u.POut("%s %d %s\n", link.Hash.B58String(), link.Size, link.Name)
43 - }
46 + daemon.ExecuteCommand(com, n, os.Stdout)
47 }
48 +
49 return nil
50 }
cmd/ipfs/mount_unix.go
+11
@@ -7,6 +7,7 @@ import (
7
8 "github.com/gonuts/flag"
9 "github.com/jbenet/commander"
10 + "github.com/jbenet/go-ipfs/daemon"
11 rofs "github.com/jbenet/go-ipfs/fuse/readonly"
12 u "github.com/jbenet/go-ipfs/util"
13 )
@@ -26,16 +27,26 @@ var cmdIpfsMount = &commander.Command{
27 }
28
29 func mountCmd(c *commander.Command, inp []string) error {
30 + u.Debug = true
31 if len(inp) < 1 || len(inp[0]) == 0 {
32 u.POut(c.Long)
33 return nil
34 }
35 + fmt.Println("wtf.")
36
37 n, err := localNode(true)
38 if err != nil {
39 return err
40 }
41
42 + fmt.Println("starting new daemon listener...")
43 + dl, err := daemon.NewDaemonListener(n, "localhost:12345")
44 + if err != nil {
45 + return err
46 + }
47 + go dl.Listen()
48 + defer dl.Close()
49 +
50 mp := inp[0]
51 fmt.Printf("Mounting at %s\n", mp)
52
core/commands/add.go new
+80
@@ -0,0 +1,80 @@
1 +package commands
2 +
3 +import (
4 + "fmt"
5 + "io/ioutil"
6 + "os"
7 + "path/filepath"
8 +
9 + "github.com/jbenet/go-ipfs/core"
10 + "github.com/jbenet/go-ipfs/importer"
11 + dag "github.com/jbenet/go-ipfs/merkledag"
12 + u "github.com/jbenet/go-ipfs/util"
13 +)
14 +
15 +// Error indicating the max depth has been exceded.
16 +var ErrDepthLimitExceeded = fmt.Errorf("depth limit exceeded")
17 +
18 +func AddPath(n *core.IpfsNode, fpath string, depth int) (*dag.Node, error) {
19 + if depth == 0 {
20 + return nil, ErrDepthLimitExceeded
21 + }
22 +
23 + fi, err := os.Stat(fpath)
24 + if err != nil {
25 + return nil, err
26 + }
27 +
28 + if fi.IsDir() {
29 + return addDir(n, fpath, depth)
30 + }
31 +
32 + return addFile(n, fpath, depth)
33 +}
34 +
35 +func addDir(n *core.IpfsNode, fpath string, depth int) (*dag.Node, error) {
36 + tree := &dag.Node{Data: dag.FolderPBData()}
37 +
38 + files, err := ioutil.ReadDir(fpath)
39 + if err != nil {
40 + return nil, err
41 + }
42 +
43 + // construct nodes for containing files.
44 + for _, f := range files {
45 + fp := filepath.Join(fpath, f.Name())
46 + nd, err := AddPath(n, fp, depth-1)
47 + if err != nil {
48 + return nil, err
49 + }
50 +
51 + if err = tree.AddNodeLink(f.Name(), nd); err != nil {
52 + return nil, err
53 + }
54 + }
55 +
56 + return tree, addNode(n, tree, fpath)
57 +}
58 +
59 +func addFile(n *core.IpfsNode, fpath string, depth int) (*dag.Node, error) {
60 + root, err := importer.NewDagFromFile(fpath)
61 + if err != nil {
62 + return nil, err
63 + }
64 +
65 + return root, addNode(n, root, fpath)
66 +}
67 +
68 +// addNode adds the node to the graph + local storage
69 +func addNode(n *core.IpfsNode, nd *dag.Node, fpath string) error {
70 + // add the file to the graph + local storage
71 + err := n.DAG.AddRecursive(nd)
72 + if err != nil {
73 + return err
74 + }
75 +
76 + u.POut("added %s\n", fpath)
77 +
78 + // ensure we keep it. atm no-op
79 + return n.PinDagNode(nd)
80 +}
core/core.go
+16 -11
@@ -133,19 +133,24 @@ func loadBitswap(cfg *config.Config, d ds.Datastore) (*bitswap.BitSwap, error) {
133 route := dht.NewDHT(local, net, d)
134 route.Start()
135
136 - for _, p := range cfg.Peers {
137 - maddr, err := ma.NewMultiaddr(p.Address)
138 - if err != nil {
139 - u.PErr("error: %v\n", err)
140 - continue
136 + go func() {
137 + u.DOut("setup: connecting to peers.\n")
138 + for _, p := range cfg.Peers {
139 + maddr, err := ma.NewMultiaddr(p.Address)
140 + if err != nil {
141 + u.PErr("error: %v\n", err)
142 + continue
143 + }
144 +
145 + u.DOut("setup: connect.\n")
146 + _, err = route.Connect(maddr)
147 + if err != nil {
148 + u.PErr("Bootstrapping error: %v\n", err)
149 + }
150 }
151 + }()
152
143 - _, err = route.Connect(maddr)
144 - if err != nil {
145 - u.PErr("Bootstrapping error: %v\n", err)
146 - }
147 - }
148 -
153 + u.DOut("setup: return new bitswap\n")
154 return bitswap.NewBitSwap(local, net, d, route), nil
155 }
156
daemon/daemon.go
+76 -48
@@ -3,41 +3,54 @@ package daemon
3 import (
4 "encoding/json"
5 "errors"
6 + "fmt"
7 + "io"
8 "net"
7 - "strings"
9
10 + core "github.com/jbenet/go-ipfs/core"
11 + commands "github.com/jbenet/go-ipfs/core/commands"
12 + dag "github.com/jbenet/go-ipfs/merkledag"
13 u "github.com/jbenet/go-ipfs/util"
14 )
15
16 var ErrInvalidCommand = errors.New("invalid command")
17
18 type DaemonListener struct {
15 - list net.Listener
16 - CommChan chan *Command
17 - closed bool
19 + node *core.IpfsNode
20 + list net.Listener
21 + closed bool
22 }
23
20 -func NewDaemonListener(addr string) (*DaemonListener, error) {
24 +func NewDaemonListener(node *core.IpfsNode, addr string) (*DaemonListener, error) {
25 list, err := net.Listen("tcp", addr)
26 if err != nil {
27 return nil, err
28 }
29 + fmt.Println("new daemon listener.")
30
31 return &DaemonListener{
27 - list: list,
28 - CommChan: make(chan *Command),
32 + node: node,
33 + list: list,
34 }, nil
35 }
36
37 type Command struct {
33 - command string
34 - args []string
35 - resp chan string
38 + Command string
39 + Args []string
40 + Opts map[string]interface{}
41 +}
42 +
43 +func NewCommand() *Command {
44 + return &Command{
45 + Opts: make(map[string]interface{}),
46 + }
47 }
48
49 func (dl *DaemonListener) Listen() {
50 + fmt.Println("listen.")
51 for {
52 c, err := dl.list.Accept()
53 + fmt.Println("Loop!")
54 if err != nil {
55 if !dl.closed {
56 u.PErr("DaemonListener Accept: %v\n", err)
@@ -49,59 +62,74 @@ func (dl *DaemonListener) Listen() {
62 }
63
64 func (dl *DaemonListener) handleConnection(c net.Conn) {
65 + defer c.Close()
66 +
67 dec := json.NewDecoder(c)
53 - enc := json.NewEncoder(c)
54 - var com string
68 +
69 + var com Command
70 err := dec.Decode(&com)
71 if err != nil {
57 - err := enc.Encode(err.Error())
58 - if err != nil {
59 - u.PErr("DaemonListener decode: %v\n", err)
60 - }
72 + fmt.Fprintln(c, err)
73 return
74 }
75 +
76 u.DOut("Got command: %v\n", com)
77 + ExecuteCommand(&com, dl.node, c)
78 +}
79
65 - cmd, err := parseCommand(com)
66 - if err != nil {
67 - err := enc.Encode(err.Error())
68 - if err != nil {
69 - u.PErr("DaemonListener parse: %v\n", err)
80 +func ExecuteCommand(com *Command, n *core.IpfsNode, out io.Writer) {
81 + u.DOut("executing command: %s\n", com.Command)
82 + switch com.Command {
83 + case "add":
84 + depth := 1
85 + if r, ok := com.Opts["r"].(bool); r && ok {
86 + depth = -1
87 }
71 - return
72 - }
88 + for _, path := range com.Args {
89 + _, err := commands.AddPath(n, path, depth)
90 + if err != nil {
91 + fmt.Fprintf(out, "addFile error: %v\n", err)
92 + continue
93 + }
94 + }
95 + case "cat":
96 + for _, fn := range com.Args {
97 + nd, err := n.Resolver.ResolvePath(fn)
98 + if err != nil {
99 + fmt.Fprintf(out, "catFile error: %v\n", err)
100 + return
101 + }
102
74 - select {
75 - case dl.CommChan <- cmd:
76 - default:
77 - u.PErr("Recieved command after closing...")
78 - return
79 - }
103 + read, err := dag.NewDagReader(nd, n.DAG)
104 + if err != nil {
105 + fmt.Fprintln(out, err)
106 + continue
107 + }
108
81 - resp := <-cmd.resp
82 - err = enc.Encode(resp)
83 - if err != nil {
84 - u.PErr("handleConnection: %v\n", err)
85 - }
86 -}
109 + _, err = io.Copy(out, read)
110 + if err != nil {
111 + fmt.Fprintln(out, err)
112 + continue
113 + }
114 + }
115 + case "ls":
116 + for _, fn := range com.Args {
117 + nd, err := n.Resolver.ResolvePath(fn)
118 + if err != nil {
119 + fmt.Fprintf(out, "ls: %v\n", err)
120 + return
121 + }
122
88 -func parseCommand(cmdi string) (*Command, error) {
89 - params := strings.Split(cmdi, " ")
90 - if len(params) == 0 {
91 - return nil, ErrInvalidCommand
123 + for _, link := range nd.Links {
124 + fmt.Fprintf(out, "%s %d %s\n", link.Hash.B58String(), link.Size, link.Name)
125 + }
126 + }
127 + default:
128 + fmt.Fprintf(out, "Invalid Command: '%s'\n", com.Command)
129 }
93 -
94 - //TODO: some sort of validation here
95 -
96 - return &Command{
97 - command: params[0],
98 - args: params[1:],
99 - resp: make(chan string),
100 - }, nil
130 }
131
132 func (dl *DaemonListener) Close() error {
133 dl.closed = true
105 - close(dl.CommChan)
134 return dl.list.Close()
135 }
daemon/daemon_client.go
+8 -12
@@ -2,28 +2,24 @@ package daemon
2
3 import (
4 "encoding/json"
5 + "io"
6 "net"
7 + "os"
8 )
9
8 -func SendCommand(command, server string) (string, error) {
10 +func SendCommand(com *Command, server string) error {
11 con, err := net.Dial("tcp", server)
12 if err != nil {
11 - return "", err
13 + return err
14 }
15
16 enc := json.NewEncoder(con)
15 - err = enc.Encode(command)
17 + err = enc.Encode(com)
18 if err != nil {
17 - return "", err
19 + return err
20 }
21
20 - dec := json.NewDecoder(con)
22 + io.Copy(os.Stdout, con)
23
22 - var resp string
23 - err = dec.Decode(&resp)
24 - if err != nil {
25 - return "", err
26 - }
27 -
28 - return resp, nil
24 + return nil
25 }
daemon/daemon_test.go deleted
-70
@@ -1,70 +0,0 @@
1 -package daemon
2 -
3 -import (
4 - "fmt"
5 - "testing"
6 -)
7 -
8 -func TestCommandCall(t *testing.T) {
9 - dl, err := NewDaemonListener("localhost:12345")
10 - if err != nil {
11 - t.Fatal(err)
12 - }
13 -
14 - go dl.Listen()
15 - defer dl.Close()
16 -
17 - go func() {
18 - _, err := SendCommand("test command for fun", "localhost:12345")
19 - if err != nil {
20 - t.Fatal(err)
21 - }
22 - }()
23 -
24 - cmd := <-dl.CommChan
25 - if cmd.command != "test" {
26 - t.Fatal("command parsing failed.")
27 - }
28 -
29 - if cmd.args[0] != "command" ||
30 - cmd.args[1] != "for" ||
31 - cmd.args[2] != "fun" {
32 - t.Fatal("Args parsed incorrectly.")
33 - }
34 -}
35 -
36 -func TestFailures(t *testing.T) {
37 - dl, err := NewDaemonListener("localhost:12345")
38 - if err != nil {
39 - t.Fatal(err)
40 - }
41 -
42 - go dl.Listen()
43 - defer dl.Close()
44 -
45 - go func() {
46 - _, err := SendCommand("test", "localhost:12345")
47 - if err != nil {
48 - t.Fatal(err)
49 - }
50 - }()
51 -
52 - cmd := <-dl.CommChan
53 - if cmd.command != "test" || len(cmd.args) > 0 {
54 - t.Fatal("Parsing Failed.")
55 - }
56 -
57 - go func() {
58 - _, err := SendCommand("", "localhost:12345")
59 - if err != nil {
60 - t.Fatal(err)
61 - }
62 - }()
63 -
64 - cmd = <-dl.CommChan
65 - if cmd.command != "" || len(cmd.args) > 0 {
66 - fmt.Println(cmd)
67 - t.Fatal("Parsing Failed.")
68 - }
69 -
70 -}