add a copy of daemon to be edited when creating new commands
Brian Tiger Chow committed
Oct 29, 2014 at 00:25 UTC
eff47794aae13f2a9324ebb89aa8c445b9e92403
3 files changed
+359
daemon/daemon.go
new
+166
@@ -0,0 +1,166 @@
1
+package daemon
2
+
3
+import (
4
+ "encoding/json"
5
+ "fmt"
6
+ "io"
7
+ "os"
8
+ "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"
14
+
15
+ 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"
18
+)
19
+
20
+var log = u.Logger("daemon")
21
+
22
+// LockFile is the filename of the daemon lock, relative to config dir
23
+const LockFile = "daemon.lock"
24
+
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
+ }
106
+}
107
+
108
+func (dl *DaemonListener) handleConnection(conn manet.Conn) {
109
+ defer conn.Close()
110
+
111
+ dec := json.NewDecoder(conn)
112
+
113
+ var command Command
114
+ err := dec.Decode(&command)
115
+ if err != nil {
116
+ fmt.Fprintln(conn, err)
117
+ return
118
+ }
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))
166
+}
daemon/daemon_client.go
new
+107
@@ -0,0 +1,107 @@
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
new
+86
@@ -0,0 +1,86 @@
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
+}