@cryptotaxi247 / kubo / commits / d67e1e506

daemon: Replaced daemon package with lock functions

Matt Bell committed Oct 27, 2014 at 17:49 UTC d67e1e50683dafd3fcb8786a122dd101908278b4
3 files changed +8 -342
daemon/daemon.go
+8 -149
@@ -1,166 +1,25 @@
1 package daemon
2
3 import (
4 - "encoding/json"
5 - "fmt"
4 "io"
7 - "os"
5 "path"
9 - "sync"
10 -
11 - core "github.com/jbenet/go-ipfs/core"
12 - "github.com/jbenet/go-ipfs/core/commands"
13 - u "github.com/jbenet/go-ipfs/util"
6
7 lock "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/camlistore/lock"
16 - ma "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr"
17 - manet "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr/net"
8 )
9
20 -var log = u.Logger("daemon")
21 -
10 // LockFile is the filename of the daemon lock, relative to config dir
11 const LockFile = "daemon.lock"
12
25 -// DaemonListener listens to an initialized IPFS node and can send it commands instead of
26 -// starting up a new set of connections
27 -type DaemonListener struct {
28 - node *core.IpfsNode
29 - list manet.Listener
30 - closed bool
31 - wg sync.WaitGroup
32 - lk io.Closer
33 -}
34 -
35 -//Command accepts user input and can be sent to the running IPFS node
36 -type Command struct {
37 - Command string
38 - Args []string
39 - Opts map[string]interface{}
40 -}
41 -
42 -func NewDaemonListener(ipfsnode *core.IpfsNode, addr ma.Multiaddr, confdir string) (*DaemonListener, error) {
43 - var err error
44 - confdir, err = u.TildeExpansion(confdir)
45 - if err != nil {
46 - return nil, err
47 - }
48 -
49 - lk, err := daemonLock(confdir)
50 - if err != nil {
51 - return nil, err
52 - }
53 -
54 - ofi, err := os.Create(confdir + "/rpcaddress")
55 - if err != nil {
56 - log.Warningf("Could not create rpcaddress file: %s", err)
57 - return nil, err
58 - }
59 -
60 - _, err = ofi.Write([]byte(addr.String()))
61 - if err != nil {
62 - log.Warningf("Could not write to rpcaddress file: %s", err)
63 - return nil, err
64 - }
65 - ofi.Close()
66 -
67 - list, err := manet.Listen(addr)
68 - if err != nil {
69 - return nil, err
70 - }
71 - log.Info("New daemon listener initialized.")
72 -
73 - return &DaemonListener{
74 - node: ipfsnode,
75 - list: list,
76 - lk: lk,
77 - }, nil
78 -}
79 -
80 -func NewCommand() *Command {
81 - return &Command{
82 - Opts: make(map[string]interface{}),
83 - }
84 -}
85 -
86 -func (dl *DaemonListener) Listen() {
87 - if dl.closed {
88 - panic("attempting to listen on a closed daemon Listener")
89 - }
90 -
91 - // add ourselves to workgroup. and remove ourselves when done.
92 - dl.wg.Add(1)
93 - defer dl.wg.Done()
94 -
95 - log.Info("daemon listening")
96 - for {
97 - conn, err := dl.list.Accept()
98 - if err != nil {
99 - if !dl.closed {
100 - log.Warning("DaemonListener Accept: %v", err)
101 - }
102 - return
103 - }
104 - go dl.handleConnection(conn)
105 - }
13 +func Lock(confdir string) (io.Closer, error) {
14 + return lock.Lock(path.Join(confdir, LockFile))
15 }
16
108 -func (dl *DaemonListener) handleConnection(conn manet.Conn) {
109 - defer conn.Close()
110 -
111 - dec := json.NewDecoder(conn)
17 +func Locked(confdir string) bool {
18 + if lk, err := Lock(confdir); err != nil {
19 + return true
20
113 - var command Command
114 - err := dec.Decode(&command)
115 - if err != nil {
116 - fmt.Fprintln(conn, err)
117 - return
21 + } else {
22 + lk.Close()
23 + return false
24 }
119 -
120 - log.Debug("Got command: %v", command)
121 - switch command.Command {
122 - case "add":
123 - err = commands.Add(dl.node, command.Args, command.Opts, conn)
124 - case "cat":
125 - err = commands.Cat(dl.node, command.Args, command.Opts, conn)
126 - case "ls":
127 - err = commands.Ls(dl.node, command.Args, command.Opts, conn)
128 - case "pin":
129 - err = commands.Pin(dl.node, command.Args, command.Opts, conn)
130 - case "publish":
131 - err = commands.Publish(dl.node, command.Args, command.Opts, conn)
132 - case "resolve":
133 - err = commands.Resolve(dl.node, command.Args, command.Opts, conn)
134 - case "diag":
135 - err = commands.Diag(dl.node, command.Args, command.Opts, conn)
136 - case "blockGet":
137 - err = commands.BlockGet(dl.node, command.Args, command.Opts, conn)
138 - case "blockPut":
139 - err = commands.BlockPut(dl.node, command.Args, command.Opts, conn)
140 - case "log":
141 - err = commands.Log(dl.node, command.Args, command.Opts, conn)
142 - case "unpin":
143 - err = commands.Unpin(dl.node, command.Args, command.Opts, conn)
144 - case "updateApply":
145 - command.Opts["onDaemon"] = true
146 - err = commands.UpdateApply(dl.node, command.Args, command.Opts, conn)
147 - default:
148 - err = fmt.Errorf("Invalid Command: '%s'", command.Command)
149 - }
150 - if err != nil {
151 - log.Errorf("%s: %s", command.Command, err)
152 - fmt.Fprintln(conn, err)
153 - }
154 -}
155 -
156 -func (dl *DaemonListener) Close() error {
157 - dl.closed = true
158 - err := dl.list.Close()
159 - dl.wg.Wait() // wait till done before releasing lock.
160 - dl.lk.Close()
161 - return err
162 -}
163 -
164 -func daemonLock(confdir string) (io.Closer, error) {
165 - return lock.Lock(path.Join(confdir, LockFile))
25 }
daemon/daemon_client.go deleted
-107
@@ -1,107 +0,0 @@
1 -package daemon
2 -
3 -import (
4 - "bufio"
5 - "encoding/json"
6 - "errors"
7 - "io"
8 - "os"
9 -
10 - ma "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr"
11 - manet "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr/net"
12 -
13 - u "github.com/jbenet/go-ipfs/util"
14 -)
15 -
16 -// ErrDaemonNotRunning is returned when attempting to retrieve the daemon's
17 -// address and the daemon is not actually running.
18 -var ErrDaemonNotRunning = errors.New("daemon not running")
19 -
20 -func getDaemonAddr(confdir string) (string, error) {
21 - var err error
22 - confdir, err = u.TildeExpansion(confdir)
23 - if err != nil {
24 - return "", err
25 - }
26 - fi, err := os.Open(confdir + "/rpcaddress")
27 - if err != nil {
28 - log.Debug("getDaemonAddr failed: %s", err)
29 - if err == os.ErrNotExist {
30 - return "", ErrDaemonNotRunning
31 - }
32 - return "", err
33 - }
34 -
35 - read := bufio.NewReader(fi)
36 -
37 - // TODO: operating system agostic line delim
38 - line, err := read.ReadBytes('\n')
39 - if err != nil && err != io.EOF {
40 - return "", err
41 - }
42 - return string(line), nil
43 -}
44 -
45 -// SendCommand attempts to run the command over a currently-running daemon.
46 -// If there is no running daemon, returns ErrDaemonNotRunning. This is done
47 -// over network RPC API. The address of the daemon is retrieved from the config
48 -// directory, where live daemons write their addresses to special files.
49 -func SendCommand(command *Command, confdir string) error {
50 - server := os.Getenv("IPFS_ADDRESS_RPC")
51 -
52 - if server == "" {
53 - //check if daemon is running
54 - log.Info("Checking if daemon is running...")
55 - if !serverIsRunning(confdir) {
56 - return ErrDaemonNotRunning
57 - }
58 -
59 - log.Info("Daemon is running!")
60 -
61 - var err error
62 - server, err = getDaemonAddr(confdir)
63 - if err != nil {
64 - return err
65 - }
66 - }
67 -
68 - return serverComm(server, command)
69 -}
70 -
71 -func serverIsRunning(confdir string) bool {
72 - var err error
73 - confdir, err = u.TildeExpansion(confdir)
74 - if err != nil {
75 - log.Errorf("Tilde Expansion Failed: %s", err)
76 - return false
77 - }
78 - lk, err := daemonLock(confdir)
79 - if err == nil {
80 - lk.Close()
81 - return false
82 - }
83 - return true
84 -}
85 -
86 -func serverComm(server string, command *Command) error {
87 - log.Info("Daemon address: %s", server)
88 - maddr, err := ma.NewMultiaddr(server)
89 - if err != nil {
90 - return err
91 - }
92 -
93 - conn, err := manet.Dial(maddr)
94 - if err != nil {
95 - return err
96 - }
97 -
98 - enc := json.NewEncoder(conn)
99 - err = enc.Encode(command)
100 - if err != nil {
101 - return err
102 - }
103 -
104 - io.Copy(os.Stdout, conn)
105 -
106 - return nil
107 -}
daemon/daemon_test.go deleted
-86
@@ -1,86 +0,0 @@
1 -package daemon
2 -
3 -import (
4 - "encoding/base64"
5 - "os"
6 - "testing"
7 -
8 - ma "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr"
9 - config "github.com/jbenet/go-ipfs/config"
10 - core "github.com/jbenet/go-ipfs/core"
11 - ci "github.com/jbenet/go-ipfs/crypto"
12 - peer "github.com/jbenet/go-ipfs/peer"
13 -)
14 -
15 -func TestInitializeDaemonListener(t *testing.T) {
16 -
17 - priv, pub, err := ci.GenerateKeyPair(ci.RSA, 512)
18 - if err != nil {
19 - t.Fatal(err)
20 - }
21 - prbytes, err := priv.Bytes()
22 - if err != nil {
23 - t.Fatal(err)
24 - }
25 -
26 - ident, _ := peer.IDFromPubKey(pub)
27 - privKey := base64.StdEncoding.EncodeToString(prbytes)
28 - pID := ident.Pretty()
29 -
30 - id := config.Identity{
31 - PeerID: pID,
32 - PrivKey: privKey,
33 - }
34 -
35 - nodeConfigs := []*config.Config{
36 - &config.Config{
37 - Identity: id,
38 - Datastore: config.Datastore{
39 - Type: "memory",
40 - },
41 - Addresses: config.Addresses{
42 - Swarm: "/ip4/0.0.0.0/tcp/4001",
43 - API: "/ip4/127.0.0.1/tcp/8000",
44 - },
45 - },
46 -
47 - &config.Config{
48 - Identity: id,
49 - Datastore: config.Datastore{
50 - Type: "leveldb",
51 - Path: ".test/datastore",
52 - },
53 - Addresses: config.Addresses{
54 - Swarm: "/ip4/0.0.0.0/tcp/4001",
55 - API: "/ip4/127.0.0.1/tcp/8000",
56 - },
57 - },
58 - }
59 -
60 - var tempConfigDir = ".test"
61 - err = os.MkdirAll(tempConfigDir, os.ModeDir|0777)
62 - if err != nil {
63 - t.Fatalf("error making temp config dir: %v", err)
64 - }
65 -
66 - for _, c := range nodeConfigs {
67 -
68 - node, _ := core.NewIpfsNode(c, false)
69 - addr, err := ma.NewMultiaddr("/ip4/127.0.0.1/tcp/1327")
70 - if err != nil {
71 - t.Fatal(err)
72 - }
73 -
74 - dl, initErr := NewDaemonListener(node, addr, tempConfigDir)
75 - if initErr != nil {
76 - t.Fatal(initErr)
77 - }
78 -
79 - closeErr := dl.Close()
80 - if closeErr != nil {
81 - t.Fatal(closeErr)
82 - }
83 -
84 - }
85 -
86 -}