@cryptotaxi247 / kubo / commits / e61c59758

implement initial ipns filesystem interface as well as plumbing command for publishing

Jeromy committed Sep 25, 2014 at 18:29 UTC e61c59758bba88d13d0f0f9e81f6686b513313b9
5 files changed +417 -3
cmd/ipfs/publish.go new
+59
@@ -0,0 +1,59 @@
1 +package main
2 +
3 +import (
4 + "os"
5 +
6 + "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/gonuts/flag"
7 + "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/commander"
8 + "github.com/jbenet/go-ipfs/core/commands"
9 + "github.com/jbenet/go-ipfs/daemon"
10 + u "github.com/jbenet/go-ipfs/util"
11 +)
12 +
13 +var cmdIpfsPub = &commander.Command{
14 + UsageLine: "publish",
15 + Short: "Publish an object to ipns under your key.",
16 + Long: `ipfs publish <path> - Publish object to ipns.
17 +
18 +`,
19 + Run: pubCmd,
20 + Flag: *flag.NewFlagSet("ipfs-publish", flag.ExitOnError),
21 +}
22 +
23 +func init() {
24 + cmdIpfsPub.Flag.String("k", "", "Specify key to use for publishing.")
25 +}
26 +
27 +func pubCmd(c *commander.Command, inp []string) error {
28 + u.Debug = true
29 + if len(inp) < 1 {
30 + u.POut(c.Long)
31 + return nil
32 + }
33 +
34 + conf, err := getConfigDir(c.Parent)
35 + if err != nil {
36 + return err
37 + }
38 +
39 + cmd := daemon.NewCommand()
40 + cmd.Command = "publish"
41 + cmd.Args = inp
42 + cmd.Opts["k"] = c.Flag.Lookup("k").Value.Get()
43 + err = daemon.SendCommand(cmd, conf)
44 + if err != nil {
45 + u.DOut("Executing command locally.\n")
46 + // Do locally
47 + conf, err := getConfigDir(c.Parent)
48 + if err != nil {
49 + return err
50 + }
51 + n, err := localNode(conf, true)
52 + if err != nil {
53 + return err
54 + }
55 +
56 + return commands.Publish(n, cmd.Args, cmd.Opts, os.Stdout)
57 + }
58 + return nil
59 +}
core/commands/publish.go new
+39
@@ -0,0 +1,39 @@
1 +package commands
2 +
3 +import (
4 + "errors"
5 + "fmt"
6 + "io"
7 +
8 + "github.com/jbenet/go-ipfs/core"
9 + u "github.com/jbenet/go-ipfs/util"
10 + "github.com/op/go-logging"
11 +
12 + nsys "github.com/jbenet/go-ipfs/namesys"
13 +)
14 +
15 +var log = logging.MustGetLogger("commands")
16 +
17 +func Publish(n *core.IpfsNode, args []string, opts map[string]interface{}, out io.Writer) error {
18 + log.Debug("Begin Publish")
19 + if n.Identity == nil {
20 + return errors.New("Identity not loaded!")
21 + }
22 +
23 + k := n.Identity.PrivKey
24 + val := u.Key(args[0])
25 +
26 + pub := nsys.NewPublisher(n.DAG, n.Routing)
27 + err := pub.Publish(k, val)
28 + if err != nil {
29 + return err
30 + }
31 +
32 + hash, err := k.GetPublic().Hash()
33 + if err != nil {
34 + return err
35 + }
36 + fmt.Fprintf(out, "Published %s to %s\n", val, u.Key(hash).Pretty())
37 +
38 + return nil
39 +}
fuse/ipns/ipns_unix.go new
+316
@@ -0,0 +1,316 @@
1 +package ipns
2 +
3 +import (
4 + "fmt"
5 + "io/ioutil"
6 + "os"
7 + "os/exec"
8 + "os/signal"
9 + "path/filepath"
10 + "runtime"
11 + "syscall"
12 + "time"
13 +
14 + "bazil.org/fuse"
15 + "bazil.org/fuse/fs"
16 + "code.google.com/p/goprotobuf/proto"
17 + "github.com/jbenet/go-ipfs/core"
18 + ci "github.com/jbenet/go-ipfs/crypto"
19 + mdag "github.com/jbenet/go-ipfs/merkledag"
20 + u "github.com/jbenet/go-ipfs/util"
21 + "github.com/op/go-logging"
22 +)
23 +
24 +var log = logging.MustGetLogger("ipns")
25 +
26 +// FileSystem is the readonly Ipfs Fuse Filesystem.
27 +type FileSystem struct {
28 + Ipfs *core.IpfsNode
29 + RootNode *Root
30 +}
31 +
32 +// NewFileSystem constructs new fs using given core.IpfsNode instance.
33 +func NewIpns(ipfs *core.IpfsNode, ipfspath string) (*FileSystem, error) {
34 + root, err := CreateRoot(ipfs, []ci.PrivKey{ipfs.Identity.PrivKey}, ipfspath)
35 + if err != nil {
36 + return nil, err
37 + }
38 + return &FileSystem{Ipfs: ipfs, RootNode: root}, nil
39 +}
40 +
41 +func CreateRoot(n *core.IpfsNode, keys []ci.PrivKey, ipfsroot string) (*Root, error) {
42 + root := new(Root)
43 + root.LocalDirs = make(map[string]*Node)
44 + root.Ipfs = n
45 + abspath, err := filepath.Abs(ipfsroot)
46 + if err != nil {
47 + return nil, err
48 + }
49 + root.IpfsRoot = abspath
50 +
51 + root.Keys = keys
52 +
53 + if len(keys) == 0 {
54 + log.Warning("No keys given for ipns root creation")
55 + } else {
56 + k := keys[0]
57 + pub := k.GetPublic()
58 + hash, err := pub.Hash()
59 + if err != nil {
60 + log.Error("Read Root Error: %s", err)
61 + return nil, err
62 + }
63 + root.LocalLink = &Link{u.Key(hash).Pretty()}
64 + }
65 +
66 + return root, nil
67 +}
68 +
69 +// Root constructs the Root of the filesystem, a Root object.
70 +func (f FileSystem) Root() (fs.Node, fuse.Error) {
71 + return f.RootNode, nil
72 +}
73 +
74 +// Root is the root object of the filesystem tree.
75 +type Root struct {
76 + Ipfs *core.IpfsNode
77 + Keys []ci.PrivKey
78 +
79 + // Used for symlinking into ipfs
80 + IpfsRoot string
81 + LocalDirs map[string]*Node
82 +
83 + LocalLink *Link
84 +}
85 +
86 +// Attr returns file attributes.
87 +func (*Root) Attr() fuse.Attr {
88 + return fuse.Attr{Mode: os.ModeDir | 0111} // -rw+x
89 +}
90 +
91 +// Lookup performs a lookup under this node.
92 +func (s *Root) Lookup(name string, intr fs.Intr) (fs.Node, fuse.Error) {
93 + log.Debug("ipns: Root Lookup: '%s'", name)
94 + switch name {
95 + case "mach_kernel", ".hidden", "._.":
96 + // Just quiet some log noise on OS X.
97 + return nil, fuse.ENOENT
98 + }
99 +
100 + if name == "local" {
101 + if s.LocalLink == nil {
102 + return nil, fuse.ENOENT
103 + }
104 + return s.LocalLink, nil
105 + }
106 +
107 + nd, ok := s.LocalDirs[name]
108 + if ok {
109 + return nd, nil
110 + }
111 +
112 + log.Debug("ipns: Falling back to resolution.")
113 + resolved, err := s.Ipfs.Namesys.Resolve(name)
114 + if err != nil {
115 + log.Error("ipns: namesys resolve error: %s", err)
116 + return nil, fuse.ENOENT
117 + }
118 +
119 + return &Link{s.IpfsRoot + "/" + resolved}, nil
120 +}
121 +
122 +// ReadDir reads a particular directory. Disallowed for root.
123 +func (r *Root) ReadDir(intr fs.Intr) ([]fuse.Dirent, fuse.Error) {
124 + u.DOut("Read Root.\n")
125 + listing := []fuse.Dirent{
126 + fuse.Dirent{
127 + Name: "local",
128 + Type: fuse.DT_Link,
129 + },
130 + }
131 + for _, k := range r.Keys {
132 + pub := k.GetPublic()
133 + hash, err := pub.Hash()
134 + if err != nil {
135 + log.Error("Read Root Error: %s", err)
136 + continue
137 + }
138 + ent := fuse.Dirent{
139 + Name: u.Key(hash).Pretty(),
140 + Type: fuse.DT_Dir,
141 + }
142 + listing = append(listing, ent)
143 + }
144 + return listing, nil
145 +}
146 +
147 +// Node is the core object representing a filesystem tree node.
148 +type Node struct {
149 + nsRoot *Node
150 + Ipfs *core.IpfsNode
151 + Nd *mdag.Node
152 + fd *mdag.DagReader
153 + cached *mdag.PBData
154 +}
155 +
156 +func (s *Node) loadData() error {
157 + s.cached = new(mdag.PBData)
158 + return proto.Unmarshal(s.Nd.Data, s.cached)
159 +}
160 +
161 +// Attr returns the attributes of a given node.
162 +func (s *Node) Attr() fuse.Attr {
163 + u.DOut("Node attr.\n")
164 + if s.cached == nil {
165 + s.loadData()
166 + }
167 + switch s.cached.GetType() {
168 + case mdag.PBData_Directory:
169 + u.DOut("this is a directory.\n")
170 + return fuse.Attr{Mode: os.ModeDir | 0555}
171 + case mdag.PBData_File, mdag.PBData_Raw:
172 + u.DOut("this is a file.\n")
173 + size, _ := s.Nd.Size()
174 + return fuse.Attr{
175 + Mode: 0444,
176 + Size: uint64(size),
177 + Blocks: uint64(len(s.Nd.Links)),
178 + }
179 + default:
180 + u.PErr("Invalid data type.")
181 + return fuse.Attr{}
182 + }
183 +}
184 +
185 +// Lookup performs a lookup under this node.
186 +func (s *Node) Lookup(name string, intr fs.Intr) (fs.Node, fuse.Error) {
187 + u.DOut("Lookup '%s'\n", name)
188 + nd, err := s.Ipfs.Resolver.ResolveLinks(s.Nd, []string{name})
189 + if err != nil {
190 + // todo: make this error more versatile.
191 + return nil, fuse.ENOENT
192 + }
193 +
194 + return &Node{Ipfs: s.Ipfs, Nd: nd}, nil
195 +}
196 +
197 +// ReadDir reads the link structure as directory entries
198 +func (s *Node) ReadDir(intr fs.Intr) ([]fuse.Dirent, fuse.Error) {
199 + u.DOut("Node ReadDir\n")
200 + entries := make([]fuse.Dirent, len(s.Nd.Links))
201 + for i, link := range s.Nd.Links {
202 + n := link.Name
203 + if len(n) == 0 {
204 + n = link.Hash.B58String()
205 + }
206 + entries[i] = fuse.Dirent{Name: n, Type: fuse.DT_File}
207 + }
208 +
209 + if len(entries) > 0 {
210 + return entries, nil
211 + }
212 + return nil, fuse.ENOENT
213 +}
214 +
215 +// ReadAll reads the object data as file data
216 +func (s *Node) ReadAll(intr fs.Intr) ([]byte, fuse.Error) {
217 + log.Debug("ipns: ReadAll Node")
218 + r, err := mdag.NewDagReader(s.Nd, s.Ipfs.DAG)
219 + if err != nil {
220 + return nil, err
221 + }
222 + // this is a terrible function... 'ReadAll'?
223 + // what if i have a 6TB file? GG RAM.
224 + return ioutil.ReadAll(r)
225 +}
226 +
227 +// Mount mounts an IpfsNode instance at a particular path. It
228 +// serves until the process receives exit signals (to Unmount).
229 +func Mount(ipfs *core.IpfsNode, fpath string, ipfspath string) error {
230 +
231 + sigc := make(chan os.Signal, 1)
232 + signal.Notify(sigc, syscall.SIGHUP, syscall.SIGINT,
233 + syscall.SIGTERM, syscall.SIGQUIT)
234 +
235 + go func() {
236 + <-sigc
237 + for {
238 + err := Unmount(fpath)
239 + if err == nil {
240 + return
241 + }
242 + time.Sleep(time.Millisecond * 10)
243 + }
244 + ipfs.Network.Close()
245 + }()
246 +
247 + c, err := fuse.Mount(fpath)
248 + if err != nil {
249 + return err
250 + }
251 + defer c.Close()
252 +
253 + fsys, err := NewIpns(ipfs, ipfspath)
254 + if err != nil {
255 + return err
256 + }
257 +
258 + err = fs.Serve(c, fsys)
259 + if err != nil {
260 + return err
261 + }
262 +
263 + // check if the mount process has an error to report
264 + <-c.Ready
265 + if err := c.MountError; err != nil {
266 + return err
267 + }
268 + return nil
269 +}
270 +
271 +// Unmount attempts to unmount the provided FUSE mount point, forcibly
272 +// if necessary.
273 +func Unmount(point string) error {
274 + fmt.Printf("Unmounting %s...\n", point)
275 +
276 + var cmd *exec.Cmd
277 + switch runtime.GOOS {
278 + case "darwin":
279 + cmd = exec.Command("diskutil", "umount", "force", point)
280 + case "linux":
281 + cmd = exec.Command("fusermount", "-u", point)
282 + default:
283 + return fmt.Errorf("unmount: unimplemented")
284 + }
285 +
286 + errc := make(chan error, 1)
287 + go func() {
288 + if err := exec.Command("umount", point).Run(); err == nil {
289 + errc <- err
290 + }
291 + // retry to unmount with the fallback cmd
292 + errc <- cmd.Run()
293 + }()
294 +
295 + select {
296 + case <-time.After(1 * time.Second):
297 + return fmt.Errorf("umount timeout")
298 + case err := <-errc:
299 + return err
300 + }
301 +}
302 +
303 +type Link struct {
304 + Target string
305 +}
306 +
307 +func (l *Link) Attr() fuse.Attr {
308 + log.Debug("Link attr.")
309 + return fuse.Attr{
310 + Mode: os.ModeSymlink | 0555,
311 + }
312 +}
313 +
314 +func (l *Link) Readlink(req *fuse.ReadlinkRequest, intr fs.Intr) (string, fuse.Error) {
315 + return l.Target, nil
316 +}
routing/dht/query.go
+1
@@ -106,6 +106,7 @@ func (r *dhtQueryRunner) Run(peers []*peer.Peer) (*dhtQueryResult, error) {
106 log.Warning("Running query with no peers!")
107 return nil, nil
108 }
109 +
110 // setup concurrency rate limiting
111 for i := 0; i < r.query.concurrency; i++ {
112 r.rateLimit <- struct{}{}
routing/dht/routing.go
+2 -3
@@ -47,14 +47,13 @@ func (dht *IpfsDHT) PutValue(ctx context.Context, key u.Key, value []byte) error
47 // If the search does not succeed, a multiaddr string of a closer peer is
48 // returned along with util.ErrSearchIncomplete
49 func (dht *IpfsDHT) GetValue(ctx context.Context, key u.Key) ([]byte, error) {
50 - ll := startNewRPC("GET")
51 - defer ll.EndAndPrint()
50 + log.Debug("Get Value [%s]", key.Pretty())
51
52 // If we have it local, dont bother doing an RPC!
53 // NOTE: this might not be what we want to do...
54 val, err := dht.getLocal(key)
55 if err == nil {
57 - ll.Success = true
56 + log.Debug("Got value locally!")
57 return val, nil
58 }
59