@cryptotaxi247 / kubo / commits / 9c221719f

switch over to using a lock file for daemon connections

Jeromy committed Sep 23, 2014 at 00:17 UTC 9c221719f909301e86ac6e1d252b7caeebd0da1b
6 files changed +99 -12
cmd/ipfs/cat.go
+6 -5
@@ -27,16 +27,17 @@ func catCmd(c *commander.Command, inp []string) error {
27 return nil
28 }
29
30 + conf, err := getConfigDir(c.Parent)
31 + if err != nil {
32 + return err
33 + }
34 +
35 com := daemon.NewCommand()
36 com.Command = "cat"
37 com.Args = inp
38
34 - err := daemon.SendCommand(com, "localhost:12345")
39 + err = daemon.SendCommand(com, conf)
40 if err != nil {
36 - conf, err := getConfigDir(c.Parent)
37 - if err != nil {
38 - return err
39 - }
41 n, err := localNode(conf, false)
42 if err != nil {
43 return err
cmd/ipfs/gen.go
+2
@@ -12,6 +12,8 @@ import (
12 u "github.com/jbenet/go-ipfs/util"
13 )
14
15 +// CommanderFunc is a function that can be passed into the Commander library as
16 +// a command handler. Defined here because commander lacks this definition.
17 type CommanderFunc func(*commander.Command, []string) error
18
19 // MakeCommand Wraps a commands.CmdFunc so that it may be safely run by the
cmd/ipfs/mount_unix.go
+1 -1
@@ -56,7 +56,7 @@ func mountCmd(c *commander.Command, inp []string) error {
56 return err
57 }
58
59 - dl, err := daemon.NewDaemonListener(n, maddr)
59 + dl, err := daemon.NewDaemonListener(n, maddr, conf)
60 if err != nil {
61 fmt.Println("Failed to create daemon listener.")
62 return err
core/commands/add.go
+15
@@ -17,12 +17,25 @@ import (
17 // Error indicating the max depth has been exceded.
18 var ErrDepthLimitExceeded = fmt.Errorf("depth limit exceeded")
19
20 +// Add is a command that imports files and directories -- given as arguments -- into ipfs.
21 func Add(n *core.IpfsNode, args []string, opts map[string]interface{}, out io.Writer) error {
22 depth := 1
23 +
24 + // if recursive, set depth to reflect so
25 if r, ok := opts["r"].(bool); r && ok {
26 depth = -1
27 }
28 +
29 + // add every path in args
30 for _, path := range args {
31 +
32 + // get absolute path, as incoming arg may be relative
33 + path, err := filepath.Abs(path)
34 + if err != nil {
35 + return fmt.Errorf("addFile error: %v", err)
36 + }
37 +
38 + // Add the file
39 nd, err := AddPath(n, path, depth)
40 if err != nil {
41 if err == ErrDepthLimitExceeded && depth == 1 {
@@ -31,6 +44,7 @@ func Add(n *core.IpfsNode, args []string, opts map[string]interface{}, out io.Wr
44 return fmt.Errorf("addFile error: %v", err)
45 }
46
47 + // get the key to print it
48 k, err := nd.Key()
49 if err != nil {
50 return fmt.Errorf("addFile error: %v", err)
@@ -41,6 +55,7 @@ func Add(n *core.IpfsNode, args []string, opts map[string]interface{}, out io.Wr
55 return nil
56 }
57
58 +// AddPath adds a particular path to ipfs.
59 func AddPath(n *core.IpfsNode, fpath string, depth int) (*dag.Node, error) {
60 if depth == 0 {
61 return nil, ErrDepthLimitExceeded
daemon/daemon.go
+35 -2
@@ -3,21 +3,28 @@ package daemon
3 import (
4 "encoding/json"
5 "fmt"
6 + "io"
7 "net"
8 + "os"
9
10 core "github.com/jbenet/go-ipfs/core"
11 "github.com/jbenet/go-ipfs/core/commands"
12 u "github.com/jbenet/go-ipfs/util"
13 + "github.com/op/go-logging"
14
15 + "github.com/camlistore/lock"
16 ma "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr"
17 )
18
19 +var log = logging.MustGetLogger("daemon")
20 +
21 // DaemonListener listens to an initialized IPFS node and can send it commands instead of
22 // starting up a new set of connections
23 type DaemonListener struct {
24 node *core.IpfsNode
25 list net.Listener
26 closed bool
27 + lk io.Closer
28 }
29
30 //Command accepts user input and can be sent to the running IPFS node
@@ -27,21 +34,46 @@ type Command struct {
34 Opts map[string]interface{}
35 }
36
30 -func NewDaemonListener(ipfsnode *core.IpfsNode, addr *ma.Multiaddr) (*DaemonListener, error) {
37 +func NewDaemonListener(ipfsnode *core.IpfsNode, addr *ma.Multiaddr, confdir string) (*DaemonListener, error) {
38 + var err error
39 + confdir, err = u.TildeExpansion(confdir)
40 + if err != nil {
41 + return nil, err
42 + }
43 +
44 + lk, err := lock.Lock(confdir + "/daemon.lock")
45 + if err != nil {
46 + return nil, err
47 + }
48 +
49 network, host, err := addr.DialArgs()
50 if err != nil {
51 return nil, err
52 }
53
54 + ofi, err := os.Create(confdir + "/rpcaddress")
55 + if err != nil {
56 + log.Warning("Could not create rpcaddress file: %s", err)
57 + return nil, err
58 + }
59 +
60 + _, err = ofi.Write([]byte(host))
61 + if err != nil {
62 + log.Warning("Could not write to rpcaddress file: %s", err)
63 + return nil, err
64 + }
65 + ofi.Close()
66 +
67 list, err := net.Listen(network, host)
68 if err != nil {
69 return nil, err
70 }
40 - fmt.Println("New daemon listener initialized.")
71 + log.Info("New daemon listener initialized.")
72
73 return &DaemonListener{
74 node: ipfsnode,
75 list: list,
76 + lk: lk,
77 }, nil
78 }
79
@@ -60,6 +92,7 @@ func (dl *DaemonListener) Listen() {
92 if !dl.closed {
93 u.PErr("DaemonListener Accept: %v\n", err)
94 }
95 + dl.lk.Close()
96 return
97 }
98 go dl.handleConnection(conn)
daemon/daemon_client.go
+40 -4
@@ -1,27 +1,63 @@
1 package daemon
2
3 import (
4 + "bufio"
5 "encoding/json"
6 + "errors"
7 "io"
8 "net"
9 "os"
10
11 ma "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr"
12 + u "github.com/jbenet/go-ipfs/util"
13 )
14
12 -//SendCommand connects to the address on the network with a timeout and encodes the connection into JSON
13 -func SendCommand(command *Command, server string) error {
15 +// ErrDaemonNotRunning is returned when attempting to retrieve the daemon's
16 +// address and the daemon is not actually running.
17 +var ErrDaemonNotRunning = errors.New("daemon not running")
18
15 - maddr, err := ma.NewMultiaddr(server)
19 +func getDaemonAddr(confdir string) (string, error) {
20 + var err error
21 + confdir, err = u.TildeExpansion(confdir)
22 + if err != nil {
23 + return "", err
24 + }
25 + fi, err := os.Open(confdir + "/rpcaddress")
26 + if err != nil {
27 + log.Debug("getDaemonAddr failed: %s", err)
28 + if err == os.ErrNotExist {
29 + return "", ErrDaemonNotRunning
30 + }
31 + return "", err
32 + }
33 +
34 + read := bufio.NewReader(fi)
35 +
36 + // TODO: operating system agostic line delim
37 + line, err := read.ReadBytes('\n')
38 + if err != nil && err != io.EOF {
39 + return "", err
40 + }
41 + return string(line), nil
42 +}
43 +
44 +// SendCommand issues a command (of type daemon.Command) to the daemon, if it
45 +// is running (if not, errors out). This is done over network RPC API. The
46 +// address of the daemon is retrieved from the configuration directory, where
47 +// live daemons write their addresses to special files.
48 +func SendCommand(command *Command, confdir string) error {
49 + server, err := getDaemonAddr(confdir)
50 if err != nil {
51 return err
52 }
53
20 - network, host, err := maddr.DialArgs()
54 + maddr, err := ma.NewMultiaddr(server)
55 if err != nil {
56 return err
57 }
58
59 + network, host, err := maddr.DialArgs()
60 +
61 conn, err := net.Dial(network, host)
62 if err != nil {
63 return err