implement an HAMT for unixfs directory sharding
License: MIT Signed-off-by: Jeromy <why@ipfs.io>
Jeromy committed
Aug 4, 2016 at 12:45 UTC
bb09ffd756f17876589baf76a008f3c5fa31beca
21 files changed
+1961
-121
assets/assets.go
+12
-2
@@ -59,17 +59,27 @@ func addAssetList(nd *core.IpfsNode, l []string) (*cid.Cid, error) {
59
}
60
61
fname := filepath.Base(p)
62
+
63
c, err := cid.Decode(s)
64
if err != nil {
65
return nil, err
66
}
67
67
- if err := dirb.AddChild(nd.Context(), fname, c); err != nil {
68
+ node, err := nd.DAG.Get(nd.Context(), c)
69
+ if err != nil {
70
+ return nil, err
71
+ }
72
+
73
+ if err := dirb.AddChild(nd.Context(), fname, node); err != nil {
74
return nil, fmt.Errorf("assets: could not add '%s' as a child: %s", fname, err)
75
}
76
}
77
72
- dir := dirb.GetNode()
78
+ dir, err := dirb.GetNode()
79
+ if err != nil {
80
+ return nil, err
81
+ }
82
+
83
dcid, err := nd.DAG.Add(dir)
84
if err != nil {
85
return nil, fmt.Errorf("assets: DAG.Add(dir) failed: %s", err)
commands/request.go
+1
-1
@@ -2,6 +2,7 @@ package commands
2
3
import (
4
"bufio"
5
+ "context"
6
"errors"
7
"fmt"
8
"io"
@@ -10,7 +11,6 @@ import (
11
"strconv"
12
"time"
13
13
- context "context"
14
"github.com/ipfs/go-ipfs/commands/files"
15
"github.com/ipfs/go-ipfs/core"
16
"github.com/ipfs/go-ipfs/repo/config"
core/commands/files/files.go
+8
-11
@@ -265,16 +265,7 @@ func getNodeFromPath(ctx context.Context, node *core.IpfsNode, p string) (node.N
265
ResolveOnce: uio.ResolveUnixfsOnce,
266
}
267
268
- nd, err := core.Resolve(ctx, node.Namesys, resolver, np)
269
- if err != nil {
270
- return nil, err
271
- }
272
- pbnd, ok := nd.(*dag.ProtoNode)
273
- if !ok {
274
- return nil, dag.ErrNotProtobuf
275
- }
276
-
277
- return pbnd, nil
268
+ return core.Resolve(ctx, node.Namesys, resolver, np)
269
default:
270
fsn, err := mfs.Lookup(node.FilesRoot, p)
271
if err != nil {
@@ -357,7 +348,13 @@ Examples:
348
case *mfs.Directory:
349
if !long {
350
var output []mfs.NodeListing
360
- for _, name := range fsn.ListNames() {
351
+ names, err := fsn.ListNames()
352
+ if err != nil {
353
+ res.SetError(err, cmds.ErrNormal)
354
+ return
355
+ }
356
+
357
+ for _, name := range names {
358
output = append(output, mfs.NodeListing{
359
Name: name,
360
})
core/coreunix/add.go
+6
-1
@@ -228,7 +228,12 @@ func (adder *Adder) outputDirs(path string, fsn mfs.FSNode) error {
228
case *mfs.File:
229
return nil
230
case *mfs.Directory:
231
- for _, name := range fsn.ListNames() {
231
+ names, err := fsn.ListNames()
232
+ if err != nil {
233
+ return err
234
+ }
235
+
236
+ for _, name := range names {
237
child, err := fsn.Child(name)
238
if err != nil {
239
return err
fuse/ipns/ipns_unix.go
+1
-2
@@ -16,7 +16,6 @@ import (
16
namesys "github.com/ipfs/go-ipfs/namesys"
17
path "github.com/ipfs/go-ipfs/path"
18
ft "github.com/ipfs/go-ipfs/unixfs"
19
- uio "github.com/ipfs/go-ipfs/unixfs/io"
19
20
ci "gx/ipfs/QmPGxZ1DP2w45WcogpW1h43BvseXbfke9N91qotpoQcUeS/go-libp2p-crypto"
21
logging "gx/ipfs/QmSpJByNKFX1sCsHBEp3R73FL4NF6FnQTEGyNAXHm2GS52/go-log"
@@ -100,7 +99,7 @@ func loadRoot(ctx context.Context, rt *keyRoot, ipfs *core.IpfsNode, name string
99
switch err {
100
case nil:
101
case namesys.ErrResolveFailed:
103
- node = uio.NewEmptyDirectory()
102
+ node = ft.EmptyDirNode()
103
default:
104
log.Errorf("looking up %s: %s", p, err)
105
return nil, err
mfs/dir.go
+69
-47
@@ -12,6 +12,7 @@ import (
12
13
dag "github.com/ipfs/go-ipfs/merkledag"
14
ft "github.com/ipfs/go-ipfs/unixfs"
15
+ uio "github.com/ipfs/go-ipfs/unixfs/io"
16
ufspb "github.com/ipfs/go-ipfs/unixfs/pb"
17
18
node "gx/ipfs/QmYDscK7dmdo2GZ9aumS8s5auUUAH5mR1jvj5pYhWusfK7/go-ipld-node"
@@ -29,25 +30,31 @@ type Directory struct {
30
files map[string]*File
31
32
lock sync.Mutex
32
- node *dag.ProtoNode
33
ctx context.Context
34
35
+ dirbuilder *uio.Directory
36
+
37
modTime time.Time
38
39
name string
40
}
41
40
-func NewDirectory(ctx context.Context, name string, node *dag.ProtoNode, parent childCloser, dserv dag.DAGService) *Directory {
41
- return &Directory{
42
- dserv: dserv,
43
- ctx: ctx,
44
- name: name,
45
- node: node,
46
- parent: parent,
47
- childDirs: make(map[string]*Directory),
48
- files: make(map[string]*File),
49
- modTime: time.Now(),
42
+func NewDirectory(ctx context.Context, name string, node node.Node, parent childCloser, dserv dag.DAGService) (*Directory, error) {
43
+ db, err := uio.NewDirectoryFromNode(dserv, node)
44
+ if err != nil {
45
+ return nil, err
46
}
47
+
48
+ return &Directory{
49
+ dserv: dserv,
50
+ ctx: ctx,
51
+ name: name,
52
+ dirbuilder: db,
53
+ parent: parent,
54
+ childDirs: make(map[string]*Directory),
55
+ files: make(map[string]*File),
56
+ modTime: time.Now(),
57
+ }, nil
58
}
59
60
// closeChild updates the child by the given name to the dag node 'nd'
@@ -81,21 +88,26 @@ func (d *Directory) closeChildUpdate(name string, nd *dag.ProtoNode, sync bool)
88
}
89
90
func (d *Directory) flushCurrentNode() (*dag.ProtoNode, error) {
84
- _, err := d.dserv.Add(d.node)
91
+ nd, err := d.dirbuilder.GetNode()
92
if err != nil {
93
return nil, err
94
}
95
89
- return d.node.Copy().(*dag.ProtoNode), nil
90
-}
96
+ _, err = d.dserv.Add(nd)
97
+ if err != nil {
98
+ return nil, err
99
+ }
100
92
-func (d *Directory) updateChild(name string, nd node.Node) error {
93
- err := d.node.RemoveNodeLink(name)
94
- if err != nil && err != dag.ErrNotFound {
95
- return err
101
+ pbnd, ok := nd.(*dag.ProtoNode)
102
+ if !ok {
103
+ return nil, dag.ErrNotProtobuf
104
}
105
98
- err = d.node.AddNodeLinkClean(name, nd)
106
+ return pbnd.Copy().(*dag.ProtoNode), nil
107
+}
108
+
109
+func (d *Directory) updateChild(name string, nd node.Node) error {
110
+ err := d.dirbuilder.AddChild(d.ctx, name, nd)
111
if err != nil {
112
return err
113
}
@@ -130,8 +142,12 @@ func (d *Directory) cacheNode(name string, nd node.Node) (FSNode, error) {
142
}
143
144
switch i.GetType() {
133
- case ufspb.Data_Directory:
134
- ndir := NewDirectory(d.ctx, name, nd, d, d.dserv)
145
+ case ufspb.Data_Directory, ufspb.Data_HAMTShard:
146
+ ndir, err := NewDirectory(d.ctx, name, nd, d, d.dserv)
147
+ if err != nil {
148
+ return nil, err
149
+ }
150
+
151
d.childDirs[name] = ndir
152
return ndir, nil
153
case ufspb.Data_File, ufspb.Data_Raw, ufspb.Data_Symlink:
@@ -175,15 +191,7 @@ func (d *Directory) Uncache(name string) {
191
// childFromDag searches through this directories dag node for a child link
192
// with the given name
193
func (d *Directory) childFromDag(name string) (node.Node, error) {
178
- pbn, err := d.node.GetLinkedNode(d.ctx, d.dserv, name)
179
- switch err {
180
- case nil:
181
- return pbn, nil
182
- case dag.ErrLinkNotFound:
183
- return nil, os.ErrNotExist
184
- default:
185
- return nil, err
186
- }
194
+ return d.dirbuilder.Find(d.ctx, name)
195
}
196
197
// childUnsync returns the child under this directory by the given name
@@ -209,7 +217,7 @@ type NodeListing struct {
217
Hash string
218
}
219
212
-func (d *Directory) ListNames() []string {
220
+func (d *Directory) ListNames() ([]string, error) {
221
d.lock.Lock()
222
defer d.lock.Unlock()
223
@@ -221,7 +229,12 @@ func (d *Directory) ListNames() []string {
229
names[n] = struct{}{}
230
}
231
224
- for _, l := range d.node.Links() {
232
+ links, err := d.dirbuilder.Links()
233
+ if err != nil {
234
+ return nil, err
235
+ }
236
+
237
+ for _, l := range links {
238
names[l.Name] = struct{}{}
239
}
240
@@ -231,7 +244,7 @@ func (d *Directory) ListNames() []string {
244
}
245
sort.Strings(out)
246
234
- return out
247
+ return out, nil
248
}
249
250
func (d *Directory) List() ([]NodeListing, error) {
@@ -239,7 +252,13 @@ func (d *Directory) List() ([]NodeListing, error) {
252
defer d.lock.Unlock()
253
254
var out []NodeListing
242
- for _, l := range d.node.Links() {
255
+
256
+ links, err := d.dirbuilder.Links()
257
+ if err != nil {
258
+ return nil, err
259
+ }
260
+
261
+ for _, l := range links {
262
child := NodeListing{}
263
child.Name = l.Name
264
@@ -285,20 +304,23 @@ func (d *Directory) Mkdir(name string) (*Directory, error) {
304
}
305
}
306
288
- ndir := new(dag.ProtoNode)
289
- ndir.SetData(ft.FolderPBData())
307
+ ndir := ft.EmptyDirNode()
308
309
_, err = d.dserv.Add(ndir)
310
if err != nil {
311
return nil, err
312
}
313
296
- err = d.node.AddNodeLinkClean(name, ndir)
314
+ err = d.dirbuilder.AddChild(d.ctx, name, ndir)
315
+ if err != nil {
316
+ return nil, err
317
+ }
318
+
319
+ dirobj, err := NewDirectory(d.ctx, name, ndir, d, d.dserv)
320
if err != nil {
321
return nil, err
322
}
323
301
- dirobj := NewDirectory(d.ctx, name, ndir, d, d.dserv)
324
d.childDirs[name] = dirobj
325
return dirobj, nil
326
}
@@ -310,12 +332,7 @@ func (d *Directory) Unlink(name string) error {
332
delete(d.childDirs, name)
333
delete(d.files, name)
334
313
- err := d.node.RemoveNodeLink(name)
314
- if err != nil {
315
- return err
316
- }
317
-
318
- _, err = d.dserv.Add(d.node)
335
+ err := d.dirbuilder.RemoveChild(d.ctx, name)
336
if err != nil {
337
return err
338
}
@@ -350,7 +367,7 @@ func (d *Directory) AddChild(name string, nd node.Node) error {
367
return err
368
}
369
353
- err = d.node.AddNodeLinkClean(name, nd)
370
+ err = d.dirbuilder.AddChild(d.ctx, name, nd)
371
if err != nil {
372
return err
373
}
@@ -406,10 +423,15 @@ func (d *Directory) GetNode() (node.Node, error) {
423
return nil, err
424
}
425
409
- _, err = d.dserv.Add(d.node)
426
+ nd, err := d.dirbuilder.GetNode()
427
+ if err != nil {
428
+ return nil, err
429
+ }
430
+
431
+ _, err = d.dserv.Add(nd)
432
if err != nil {
433
return nil, err
434
}
435
414
- return d.node.Copy().(*dag.ProtoNode), nil
436
+ return nd, err
437
}
mfs/mfs_test.go
+125
@@ -747,6 +747,67 @@ func TestMfsStress(t *testing.T) {
747
}
748
}
749
750
+func TestMfsHugeDir(t *testing.T) {
751
+ ctx, cancel := context.WithCancel(context.Background())
752
+ defer cancel()
753
+ _, rt := setupRoot(ctx, t)
754
+
755
+ for i := 0; i < 100000; i++ {
756
+ err := Mkdir(rt, fmt.Sprintf("/dir%d", i), false, false)
757
+ if err != nil {
758
+ t.Fatal(err)
759
+ }
760
+ }
761
+}
762
+
763
+func TestMkdirP(t *testing.T) {
764
+ ctx, cancel := context.WithCancel(context.Background())
765
+ defer cancel()
766
+ _, rt := setupRoot(ctx, t)
767
+
768
+ err := Mkdir(rt, "/a/b/c/d/e/f", true, true)
769
+ if err != nil {
770
+ t.Fatal(err)
771
+ }
772
+}
773
+
774
+func TestConcurrentWriteAndFlush(t *testing.T) {
775
+ ctx, cancel := context.WithCancel(context.Background())
776
+ defer cancel()
777
+ ds, rt := setupRoot(ctx, t)
778
+
779
+ d := mkdirP(t, rt.GetValue().(*Directory), "foo/bar/baz")
780
+ fn := fileNodeFromReader(t, ds, bytes.NewBuffer(nil))
781
+ err := d.AddChild("file", fn)
782
+ if err != nil {
783
+ t.Fatal(err)
784
+ }
785
+
786
+ nloops := 5000
787
+
788
+ wg := new(sync.WaitGroup)
789
+ wg.Add(1)
790
+ go func() {
791
+ defer wg.Done()
792
+ for i := 0; i < nloops; i++ {
793
+ err := writeFile(rt, "/foo/bar/baz/file", []byte("STUFF"))
794
+ if err != nil {
795
+ t.Error("file write failed: ", err)
796
+ return
797
+ }
798
+ }
799
+ }()
800
+
801
+ for i := 0; i < nloops; i++ {
802
+ _, err := rt.GetValue().GetNode()
803
+ if err != nil {
804
+ t.Fatal(err)
805
+ }
806
+ }
807
+
808
+ wg.Wait()
809
+}
810
+
811
func TestFlushing(t *testing.T) {
812
ctx, cancel := context.WithCancel(context.Background())
813
defer cancel()
@@ -892,6 +953,70 @@ func TestConcurrentReads(t *testing.T) {
953
}
954
wg.Wait()
955
}
956
+func writeFile(rt *Root, path string, data []byte) error {
957
+ n, err := Lookup(rt, path)
958
+ if err != nil {
959
+ return err
960
+ }
961
+
962
+ fi, ok := n.(*File)
963
+ if !ok {
964
+ return fmt.Errorf("expected to receive a file, but didnt get one")
965
+ }
966
+
967
+ fd, err := fi.Open(OpenWriteOnly, true)
968
+ if err != nil {
969
+ return err
970
+ }
971
+ defer fd.Close()
972
+
973
+ nw, err := fd.Write(data)
974
+ if err != nil {
975
+ return err
976
+ }
977
+
978
+ if nw != 10 {
979
+ fmt.Errorf("wrote incorrect amount")
980
+ }
981
+
982
+ return nil
983
+}
984
+
985
+func TestConcurrentWrites(t *testing.T) {
986
+ ctx, cancel := context.WithCancel(context.Background())
987
+ defer cancel()
988
+
989
+ ds, rt := setupRoot(ctx, t)
990
+
991
+ rootdir := rt.GetValue().(*Directory)
992
+
993
+ path := "a/b/c"
994
+ d := mkdirP(t, rootdir, path)
995
+
996
+ fi := fileNodeFromReader(t, ds, bytes.NewReader(make([]byte, 0)))
997
+ err := d.AddChild("afile", fi)
998
+ if err != nil {
999
+ t.Fatal(err)
1000
+ }
1001
+
1002
+ var wg sync.WaitGroup
1003
+ nloops := 100
1004
+ for i := 0; i < 10; i++ {
1005
+ wg.Add(1)
1006
+ go func(me int) {
1007
+ defer wg.Done()
1008
+ mybuf := bytes.Repeat([]byte{byte(me)}, 10)
1009
+ for j := 0; j < nloops; j++ {
1010
+ err := writeFile(rt, "a/b/c/afile", mybuf)
1011
+ if err != nil {
1012
+ t.Error("writefile failed: ", err)
1013
+ return
1014
+ }
1015
+ }
1016
+ }(i)
1017
+ }
1018
+ wg.Wait()
1019
+}
1020
1021
func TestFileDescriptors(t *testing.T) {
1022
ctx, cancel := context.WithCancel(context.Background())
mfs/system.go
+6
-1
@@ -88,7 +88,12 @@ func NewRoot(parent context.Context, ds dag.DAGService, node *dag.ProtoNode, pf
88
89
switch pbn.GetType() {
90
case ft.TDirectory:
91
- root.val = NewDirectory(parent, node.String(), node, root, ds)
91
+ rval, err := NewDirectory(parent, node.String(), node, root, ds)
92
+ if err != nil {
93
+ return nil, err
94
+ }
95
+
96
+ root.val = rval
97
case ft.TFile, ft.TMetadata, ft.TRaw:
98
fi, err := NewFile(node.String(), node, root, ds)
99
if err != nil {
package.json
+6
@@ -300,6 +300,12 @@
300
"hash": "QmaFNtBAXX4nVMQWbUqNysXyhevUj1k4B1y5uS45LC7Vw9",
301
"name": "fuse",
302
"version": "0.1.3"
303
+ },
304
+ {
305
+ "author": "whyrusleeping",
306
+ "hash": "QmfJHywXQu98UeZtGJBQrPAR6AtmDjjbe3qjTo9piXHPnx",
307
+ "name": "murmur3",
308
+ "version": "0.0.0"
309
}
310
],
311
"gxVersion": "0.10.0",
test/sharness/t0250-files-api.sh
+28
@@ -46,6 +46,33 @@ verify_dir_contents() {
46
'
47
}
48
49
+test_sharding() {
50
+ test_expect_success "make a directory" '
51
+ ipfs files mkdir /foo
52
+ '
53
+
54
+ test_expect_success "can make 1100 files in a directory" '
55
+ printf "" > list_exp_raw
56
+ for i in `seq 1100`
57
+ do
58
+ echo $i | ipfs files write --create /foo/file$i
59
+ echo file$i >> list_exp_raw
60
+ done
61
+ '
62
+
63
+ test_expect_success "listing works" '
64
+ ipfs files ls /foo |sort > list_out &&
65
+ sort list_exp_raw > list_exp &&
66
+ test_cmp list_exp list_out
67
+ '
68
+
69
+ test_expect_success "can read a file from sharded directory" '
70
+ ipfs files read /foo/file65 > file_out &&
71
+ echo "65" > file_exp &&
72
+ test_cmp file_out file_exp
73
+ '
74
+}
75
+
76
test_files_api() {
77
test_expect_success "can mkdir in root" '
78
ipfs files mkdir /cats
@@ -491,5 +518,6 @@ test_launch_ipfs_daemon
518
519
ONLINE=1 # set online flag so tests can easily tell
520
test_files_api
521
+test_sharding
522
test_kill_ipfs_daemon
523
test_done
unixfs/format.go
+1
@@ -17,6 +17,7 @@ const (
17
TDirectory = pb.Data_Directory
18
TMetadata = pb.Data_Metadata
19
TSymlink = pb.Data_Symlink
20
+ THAMTShard = pb.Data_HAMTShard
21
)
22
23
var ErrMalformedFileFormat = errors.New("malformed data in file format")
unixfs/hamt/hamt.go
new
+464
@@ -0,0 +1,464 @@
1
+package hamt
2
+
3
+import (
4
+ "context"
5
+ "fmt"
6
+ "math"
7
+ "math/big"
8
+ "os"
9
+
10
+ dag "github.com/ipfs/go-ipfs/merkledag"
11
+ format "github.com/ipfs/go-ipfs/unixfs"
12
+ upb "github.com/ipfs/go-ipfs/unixfs/pb"
13
+
14
+ node "gx/ipfs/QmYDscK7dmdo2GZ9aumS8s5auUUAH5mR1jvj5pYhWusfK7/go-ipld-node"
15
+ proto "gx/ipfs/QmZ4Qi3GaRbjcx28Sme5eMH7RQjGkt8wHxt2a65oLaeFEV/gogo-protobuf/proto"
16
+ "gx/ipfs/QmfJHywXQu98UeZtGJBQrPAR6AtmDjjbe3qjTo9piXHPnx/murmur3"
17
+)
18
+
19
+const (
20
+ HashMurmur3 uint64 = 0x22
21
+)
22
+
23
+type HamtShard struct {
24
+ nd *dag.ProtoNode
25
+
26
+ bitfield *big.Int
27
+
28
+ children []child
29
+
30
+ tableSize int
31
+ tableSizeLg2 int
32
+
33
+ hashFunc uint64
34
+
35
+ prefixPadStr string
36
+ maxpadlen int
37
+
38
+ dserv dag.DAGService
39
+}
40
+
41
+// child can either be another shard, or a leaf node value
42
+type child interface {
43
+ Node() (node.Node, error)
44
+ Label() string
45
+}
46
+
47
+func NewHamtShard(dserv dag.DAGService, size int) *HamtShard {
48
+ ds := makeHamtShard(dserv, size)
49
+ ds.bitfield = big.NewInt(0)
50
+ ds.nd = new(dag.ProtoNode)
51
+ ds.hashFunc = HashMurmur3
52
+ return ds
53
+}
54
+
55
+func makeHamtShard(ds dag.DAGService, size int) *HamtShard {
56
+ maxpadding := fmt.Sprintf("%X", size-1)
57
+ return &HamtShard{
58
+ tableSizeLg2: int(math.Log2(float64(size))),
59
+ prefixPadStr: fmt.Sprintf("%%0%dX", len(maxpadding)),
60
+ maxpadlen: len(maxpadding),
61
+ tableSize: size,
62
+ dserv: ds,
63
+ }
64
+}
65
+
66
+func NewHamtFromDag(dserv dag.DAGService, nd node.Node) (*HamtShard, error) {
67
+ pbnd, ok := nd.(*dag.ProtoNode)
68
+ if !ok {
69
+ return nil, dag.ErrLinkNotFound
70
+ }
71
+
72
+ pbd, err := format.FromBytes(pbnd.Data())
73
+ if err != nil {
74
+ return nil, err
75
+ }
76
+
77
+ if pbd.GetType() != upb.Data_HAMTShard {
78
+ return nil, fmt.Errorf("node was not a dir shard")
79
+ }
80
+
81
+ if pbd.GetHashType() != HashMurmur3 {
82
+ return nil, fmt.Errorf("only murmur3 supported as hash function")
83
+ }
84
+
85
+ ds := makeHamtShard(dserv, int(pbd.GetFanout()))
86
+ ds.nd = pbnd.Copy().(*dag.ProtoNode)
87
+ ds.children = make([]child, len(pbnd.Links()))
88
+ ds.bitfield = new(big.Int).SetBytes(pbd.GetData())
89
+ ds.hashFunc = pbd.GetHashType()
90
+
91
+ return ds, nil
92
+}
93
+
94
+// Node serializes the HAMT structure into a merkledag node with unixfs formatting
95
+func (ds *HamtShard) Node() (node.Node, error) {
96
+ out := new(dag.ProtoNode)
97
+
98
+ // TODO: optimized 'for each set bit'
99
+ for i := 0; i < ds.tableSize; i++ {
100
+ if ds.bitfield.Bit(i) == 0 {
101
+ continue
102
+ }
103
+
104
+ cindex := ds.indexForBitPos(i)
105
+ ch := ds.children[cindex]
106
+ if ch != nil {
107
+ cnd, err := ch.Node()
108
+ if err != nil {
109
+ return nil, err
110
+ }
111
+
112
+ err = out.AddNodeLinkClean(ds.linkNamePrefix(i)+ch.Label(), cnd)
113
+ if err != nil {
114
+ return nil, err
115
+ }
116
+ } else {
117
+ // child unloaded, just copy in link with updated name
118
+ lnk := ds.nd.Links()[cindex]
119
+ label := lnk.Name[ds.maxpadlen:]
120
+
121
+ err := out.AddRawLink(ds.linkNamePrefix(i)+label, lnk)
122
+ if err != nil {
123
+ return nil, err
124
+ }
125
+ }
126
+ }
127
+
128
+ typ := upb.Data_HAMTShard
129
+ data, err := proto.Marshal(&upb.Data{
130
+ Type: &typ,
131
+ Fanout: proto.Uint64(uint64(ds.tableSize)),
132
+ HashType: proto.Uint64(HashMurmur3),
133
+ Data: ds.bitfield.Bytes(),
134
+ })
135
+ if err != nil {
136
+ return nil, err
137
+ }
138
+
139
+ out.SetData(data)
140
+
141
+ _, err = ds.dserv.Add(out)
142
+ if err != nil {
143
+ return nil, err
144
+ }
145
+
146
+ return out, nil
147
+}
148
+
149
+type shardValue struct {
150
+ key string
151
+ val node.Node
152
+}
153
+
154
+func (sv *shardValue) Node() (node.Node, error) {
155
+ return sv.val, nil
156
+}
157
+
158
+func (sv *shardValue) Label() string {
159
+ return sv.key
160
+}
161
+
162
+func hash(val []byte) []byte {
163
+ h := murmur3.New64()
164
+ h.Write(val)
165
+ return h.Sum(nil)
166
+}
167
+
168
+// Label for HamtShards is the empty string, this is used to differentiate them from
169
+// value entries
170
+func (ds *HamtShard) Label() string {
171
+ return ""
172
+}
173
+
174
+// Set sets 'name' = nd in the HAMT
175
+func (ds *HamtShard) Set(ctx context.Context, name string, nd node.Node) error {
176
+ hv := &hashBits{b: hash([]byte(name))}
177
+ return ds.modifyValue(ctx, hv, name, nd)
178
+}
179
+
180
+// Remove deletes the named entry if it exists, this operation is idempotent.
181
+func (ds *HamtShard) Remove(ctx context.Context, name string) error {
182
+ hv := &hashBits{b: hash([]byte(name))}
183
+ return ds.modifyValue(ctx, hv, name, nil)
184
+}
185
+
186
+func (ds *HamtShard) Find(ctx context.Context, name string) (node.Node, error) {
187
+ hv := &hashBits{b: hash([]byte(name))}
188
+
189
+ var out node.Node
190
+ err := ds.getValue(ctx, hv, name, func(sv *shardValue) error {
191
+ out = sv.val
192
+ return nil
193
+ })
194
+
195
+ return out, err
196
+}
197
+
198
+// getChild returns the i'th child of this shard. If it is cached in the
199
+// children array, it will return it from there. Otherwise, it loads the child
200
+// node from disk.
201
+func (ds *HamtShard) getChild(ctx context.Context, i int) (child, error) {
202
+ if i >= len(ds.children) || i < 0 {
203
+ return nil, fmt.Errorf("invalid index passed to getChild (likely corrupt bitfield)")
204
+ }
205
+
206
+ if len(ds.children) != len(ds.nd.Links()) {
207
+ return nil, fmt.Errorf("inconsistent lengths between children array and Links array")
208
+ }
209
+
210
+ c := ds.children[i]
211
+ if c != nil {
212
+ return c, nil
213
+ }
214
+
215
+ return ds.loadChild(ctx, i)
216
+}
217
+
218
+// loadChild reads the i'th child node of this shard from disk and returns it
219
+// as a 'child' interface
220
+func (ds *HamtShard) loadChild(ctx context.Context, i int) (child, error) {
221
+ lnk := ds.nd.Links()[i]
222
+ if len(lnk.Name) < ds.maxpadlen {
223
+ return nil, fmt.Errorf("invalid link name '%s'", lnk.Name)
224
+ }
225
+
226
+ nd, err := lnk.GetNode(ctx, ds.dserv)
227
+ if err != nil {
228
+ return nil, err
229
+ }
230
+
231
+ var c child
232
+ if len(lnk.Name) == ds.maxpadlen {
233
+ pbnd, ok := nd.(*dag.ProtoNode)
234
+ if !ok {
235
+ return nil, dag.ErrNotProtobuf
236
+ }
237
+
238
+ pbd, err := format.FromBytes(pbnd.Data())
239
+ if err != nil {
240
+ return nil, err
241
+ }
242
+
243
+ if pbd.GetType() != format.THAMTShard {
244
+ return nil, fmt.Errorf("HAMT entries must have non-zero length name")
245
+ }
246
+
247
+ cds, err := NewHamtFromDag(ds.dserv, nd)
248
+ if err != nil {
249
+ return nil, err
250
+ }
251
+
252
+ c = cds
253
+ } else {
254
+ c = &shardValue{
255
+ key: lnk.Name[ds.maxpadlen:],
256
+ val: nd,
257
+ }
258
+ }
259
+
260
+ ds.children[i] = c
261
+ return c, nil
262
+}
263
+
264
+func (ds *HamtShard) setChild(i int, c child) {
265
+ ds.children[i] = c
266
+}
267
+
268
+func (ds *HamtShard) insertChild(idx int, key string, val node.Node) error {
269
+ if val == nil {
270
+ return os.ErrNotExist
271
+ }
272
+
273
+ i := ds.indexForBitPos(idx)
274
+ ds.bitfield.SetBit(ds.bitfield, idx, 1)
275
+ sv := &shardValue{
276
+ key: key,
277
+ val: val,
278
+ }
279
+
280
+ ds.children = append(ds.children[:i], append([]child{sv}, ds.children[i:]...)...)
281
+ ds.nd.SetLinks(append(ds.nd.Links()[:i], append([]*node.Link{nil}, ds.nd.Links()[i:]...)...))
282
+ return nil
283
+}
284
+
285
+func (ds *HamtShard) rmChild(i int) error {
286
+ if i < 0 || i >= len(ds.children) || i >= len(ds.nd.Links()) {
287
+ return fmt.Errorf("hamt: attempted to remove child with out of range index")
288
+ }
289
+
290
+ copy(ds.children[i:], ds.children[i+1:])
291
+ ds.children = ds.children[:len(ds.children)-1]
292
+
293
+ copy(ds.nd.Links()[i:], ds.nd.Links()[i+1:])
294
+ ds.nd.SetLinks(ds.nd.Links()[:len(ds.nd.Links())-1])
295
+
296
+ return nil
297
+}
298
+
299
+func (ds *HamtShard) getValue(ctx context.Context, hv *hashBits, key string, cb func(*shardValue) error) error {
300
+ idx := hv.Next(ds.tableSizeLg2)
301
+ if ds.bitfield.Bit(int(idx)) == 1 {
302
+ cindex := ds.indexForBitPos(idx)
303
+
304
+ child, err := ds.getChild(ctx, cindex)
305
+ if err != nil {
306
+ return err
307
+ }
308
+
309
+ switch child := child.(type) {
310
+ case *HamtShard:
311
+ return child.getValue(ctx, hv, key, cb)
312
+ case *shardValue:
313
+ if child.key == key {
314
+ return cb(child)
315
+ }
316
+ }
317
+ }
318
+
319
+ return os.ErrNotExist
320
+}
321
+
322
+func (ds *HamtShard) EnumLinks() ([]*node.Link, error) {
323
+ var links []*node.Link
324
+ err := ds.walkTrie(func(sv *shardValue) error {
325
+ lnk, err := node.MakeLink(sv.val)
326
+ if err != nil {
327
+ return err
328
+ }
329
+
330
+ lnk.Name = sv.key
331
+
332
+ links = append(links, lnk)
333
+ return nil
334
+ })
335
+ if err != nil {
336
+ return nil, err
337
+ }
338
+
339
+ return links, nil
340
+}
341
+
342
+func (ds *HamtShard) walkTrie(cb func(*shardValue) error) error {
343
+ for i := 0; i < ds.tableSize; i++ {
344
+ if ds.bitfield.Bit(i) == 0 {
345
+ continue
346
+ }
347
+
348
+ idx := ds.indexForBitPos(i)
349
+ // NOTE: an optimized version could simply iterate over each
350
+ // element in the 'children' array.
351
+ c, err := ds.getChild(context.TODO(), idx)
352
+ if err != nil {
353
+ return err
354
+ }
355
+
356
+ switch c := c.(type) {
357
+ case *shardValue:
358
+ err := cb(c)
359
+ if err != nil {
360
+ return err
361
+ }
362
+
363
+ case *HamtShard:
364
+ err := c.walkTrie(cb)
365
+ if err != nil {
366
+ return err
367
+ }
368
+ default:
369
+ return fmt.Errorf("unexpected child type: %#v", c)
370
+ }
371
+ }
372
+ return nil
373
+}
374
+
375
+func (ds *HamtShard) modifyValue(ctx context.Context, hv *hashBits, key string, val node.Node) error {
376
+ idx := hv.Next(ds.tableSizeLg2)
377
+
378
+ if ds.bitfield.Bit(idx) != 1 {
379
+ return ds.insertChild(idx, key, val)
380
+ }
381
+
382
+ cindex := ds.indexForBitPos(idx)
383
+
384
+ child, err := ds.getChild(ctx, cindex)
385
+ if err != nil {
386
+ return err
387
+ }
388
+
389
+ switch child := child.(type) {
390
+ case *HamtShard:
391
+ err := child.modifyValue(ctx, hv, key, val)
392
+ if err != nil {
393
+ return err
394
+ }
395
+
396
+ if val == nil {
397
+ switch len(child.children) {
398
+ case 0:
399
+ // empty sub-shard, prune it
400
+ // Note: this shouldnt normally ever happen
401
+ // in the event of another implementation creates flawed
402
+ // structures, this will help to normalize them.
403
+ ds.bitfield.SetBit(ds.bitfield, idx, 0)
404
+ return ds.rmChild(cindex)
405
+ case 1:
406
+ nchild, ok := child.children[0].(*shardValue)
407
+ if ok {
408
+ // sub-shard with a single value element, collapse it
409
+ ds.setChild(cindex, nchild)
410
+ }
411
+ return nil
412
+ }
413
+ }
414
+
415
+ return nil
416
+ case *shardValue:
417
+ switch {
418
+ case val == nil: // passing a nil value signifies a 'delete'
419
+ ds.bitfield.SetBit(ds.bitfield, idx, 0)
420
+ return ds.rmChild(cindex)
421
+
422
+ case child.key == key: // value modification
423
+ child.val = val
424
+ return nil
425
+
426
+ default: // replace value with another shard, one level deeper
427
+ ns := NewHamtShard(ds.dserv, ds.tableSize)
428
+ chhv := &hashBits{
429
+ b: hash([]byte(child.key)),
430
+ consumed: hv.consumed,
431
+ }
432
+
433
+ err := ns.modifyValue(ctx, hv, key, val)
434
+ if err != nil {
435
+ return err
436
+ }
437
+
438
+ err = ns.modifyValue(ctx, chhv, child.key, child.val)
439
+ if err != nil {
440
+ return err
441
+ }
442
+
443
+ ds.setChild(cindex, ns)
444
+ return nil
445
+ }
446
+ default:
447
+ return fmt.Errorf("unexpected type for child: %#v", child)
448
+ }
449
+}
450
+
451
+func (ds *HamtShard) indexForBitPos(bp int) int {
452
+ // TODO: an optimization could reuse the same 'mask' here and change the size
453
+ // as needed. This isnt yet done as the bitset package doesnt make it easy
454
+ // to do.
455
+ mask := new(big.Int).Sub(new(big.Int).Exp(big.NewInt(2), big.NewInt(int64(bp)), nil), big.NewInt(1))
456
+ mask.And(mask, ds.bitfield)
457
+
458
+ return popCount(mask)
459
+}
460
+
461
+// linkNamePrefix takes in the bitfield index of an entry and returns its hex prefix
462
+func (ds *HamtShard) linkNamePrefix(idx int) string {
463
+ return fmt.Sprintf(ds.prefixPadStr, idx)
464
+}
unixfs/hamt/hamt_stress_test.go
new
+280
@@ -0,0 +1,280 @@
1
+package hamt
2
+
3
+import (
4
+ "bufio"
5
+ "context"
6
+ "fmt"
7
+ "math/rand"
8
+ "os"
9
+ "strconv"
10
+ "strings"
11
+ "testing"
12
+ "time"
13
+
14
+ dag "github.com/ipfs/go-ipfs/merkledag"
15
+ mdtest "github.com/ipfs/go-ipfs/merkledag/test"
16
+ ft "github.com/ipfs/go-ipfs/unixfs"
17
+)
18
+
19
+func getNames(prefix string, count int) []string {
20
+ out := make([]string, count)
21
+ for i := 0; i < count; i++ {
22
+ out[i] = fmt.Sprintf("%s%d", prefix, i)
23
+ }
24
+ return out
25
+}
26
+
27
+const (
28
+ opAdd = iota
29
+ opDel
30
+ opFind
31
+)
32
+
33
+type testOp struct {
34
+ Op int
35
+ Val string
36
+}
37
+
38
+func stringArrToSet(arr []string) map[string]bool {
39
+ out := make(map[string]bool)
40
+ for _, s := range arr {
41
+ out[s] = true
42
+ }
43
+ return out
44
+}
45
+
46
+// generate two different random sets of operations to result in the same
47
+// ending directory (same set of entries at the end) and execute each of them
48
+// in turn, then compare to ensure the output is the same on each.
49
+func TestOrderConsistency(t *testing.T) {
50
+ seed := time.Now().UnixNano()
51
+ t.Logf("using seed = %d", seed)
52
+ ds := mdtest.Mock()
53
+
54
+ shardWidth := 1024
55
+
56
+ keep := getNames("good", 4000)
57
+ temp := getNames("tempo", 6000)
58
+
59
+ ops := genOpSet(seed, keep, temp)
60
+ s, err := executeOpSet(t, ds, shardWidth, ops)
61
+ if err != nil {
62
+ t.Fatal(err)
63
+ }
64
+
65
+ err = validateOpSetCompletion(t, s, keep, temp)
66
+ if err != nil {
67
+ t.Fatal(err)
68
+ }
69
+
70
+ ops2 := genOpSet(seed+1000, keep, temp)
71
+ s2, err := executeOpSet(t, ds, shardWidth, ops2)
72
+ if err != nil {
73
+ t.Fatal(err)
74
+ }
75
+
76
+ err = validateOpSetCompletion(t, s2, keep, temp)
77
+ if err != nil {
78
+ t.Fatal(err)
79
+ }
80
+
81
+ nd, err := s.Node()
82
+ if err != nil {
83
+ t.Fatal(err)
84
+ }
85
+
86
+ nd2, err := s2.Node()
87
+ if err != nil {
88
+ t.Fatal(err)
89
+ }
90
+
91
+ k := nd.Cid()
92
+ k2 := nd2.Cid()
93
+
94
+ if !k.Equals(k2) {
95
+ t.Fatal("got different results: ", k, k2)
96
+ }
97
+}
98
+
99
+func validateOpSetCompletion(t *testing.T, s *HamtShard, keep, temp []string) error {
100
+ ctx := context.TODO()
101
+ for _, n := range keep {
102
+ _, err := s.Find(ctx, n)
103
+ if err != nil {
104
+ return fmt.Errorf("couldnt find %s: %s", n, err)
105
+ }
106
+ }
107
+
108
+ for _, n := range temp {
109
+ _, err := s.Find(ctx, n)
110
+ if err != os.ErrNotExist {
111
+ return fmt.Errorf("expected not to find: %s", err)
112
+ }
113
+ }
114
+
115
+ return nil
116
+}
117
+
118
+func executeOpSet(t *testing.T, ds dag.DAGService, width int, ops []testOp) (*HamtShard, error) {
119
+ ctx := context.TODO()
120
+ s := NewHamtShard(ds, width)
121
+ e := ft.EmptyDirNode()
122
+ ds.Add(e)
123
+
124
+ for _, o := range ops {
125
+ switch o.Op {
126
+ case opAdd:
127
+ err := s.Set(ctx, o.Val, e)
128
+ if err != nil {
129
+ return nil, fmt.Errorf("inserting %s: %s", o.Val, err)
130
+ }
131
+ case opDel:
132
+ err := s.Remove(ctx, o.Val)
133
+ if err != nil {
134
+ return nil, fmt.Errorf("deleting %s: %s", o.Val, err)
135
+ }
136
+ case opFind:
137
+ _, err := s.Find(ctx, o.Val)
138
+ if err != nil {
139
+ return nil, fmt.Errorf("finding %s: %s", o.Val, err)
140
+ }
141
+ }
142
+ }
143
+
144
+ return s, nil
145
+}
146
+
147
+func genOpSet(seed int64, keep, temp []string) []testOp {
148
+ tempset := stringArrToSet(temp)
149
+
150
+ allnames := append(keep, temp...)
151
+ shuffle(seed, allnames)
152
+
153
+ var todel []string
154
+
155
+ var ops []testOp
156
+
157
+ for {
158
+ n := len(allnames) + len(todel)
159
+ if n == 0 {
160
+ return ops
161
+ }
162
+
163
+ rn := rand.Intn(n)
164
+
165
+ if rn < len(allnames) {
166
+ next := allnames[0]
167
+ allnames = allnames[1:]
168
+ ops = append(ops, testOp{
169
+ Op: opAdd,
170
+ Val: next,
171
+ })
172
+
173
+ if tempset[next] {
174
+ todel = append(todel, next)
175
+ }
176
+ } else {
177
+ shuffle(seed+100, todel)
178
+ next := todel[0]
179
+ todel = todel[1:]
180
+
181
+ ops = append(ops, testOp{
182
+ Op: opDel,
183
+ Val: next,
184
+ })
185
+ }
186
+ }
187
+}
188
+
189
+// executes the given op set with a repl to allow easier debugging
190
+func debugExecuteOpSet(ds dag.DAGService, width int, ops []testOp) (*HamtShard, error) {
191
+ s := NewHamtShard(ds, width)
192
+ e := ft.EmptyDirNode()
193
+ ds.Add(e)
194
+ ctx := context.TODO()
195
+
196
+ run := 0
197
+
198
+ opnames := map[int]string{
199
+ opAdd: "add",
200
+ opDel: "del",
201
+ }
202
+
203
+mainloop:
204
+ for i := 0; i < len(ops); i++ {
205
+ o := ops[i]
206
+
207
+ fmt.Printf("Op %d: %s %s\n", i, opnames[o.Op], o.Val)
208
+ for run == 0 {
209
+ cmd := readCommand()
210
+ parts := strings.Split(cmd, " ")
211
+ switch parts[0] {
212
+ case "":
213
+ run = 1
214
+ case "find":
215
+ _, err := s.Find(ctx, parts[1])
216
+ if err == nil {
217
+ fmt.Println("success")
218
+ } else {
219
+ fmt.Println(err)
220
+ }
221
+ case "run":
222
+ if len(parts) > 1 {
223
+ n, err := strconv.Atoi(parts[1])
224
+ if err != nil {
225
+ panic(err)
226
+ }
227
+
228
+ run = n
229
+ } else {
230
+ run = -1
231
+ }
232
+ case "lookop":
233
+ for k := 0; k < len(ops); k++ {
234
+ if ops[k].Val == parts[1] {
235
+ fmt.Printf(" Op %d: %s %s\n", k, opnames[ops[k].Op], parts[1])
236
+ }
237
+ }
238
+ case "restart":
239
+ s = NewHamtShard(ds, width)
240
+ i = -1
241
+ continue mainloop
242
+ case "print":
243
+ nd, err := s.Node()
244
+ if err != nil {
245
+ panic(err)
246
+ }
247
+ printDag(ds, nd.(*dag.ProtoNode), 0)
248
+ }
249
+ }
250
+ run--
251
+
252
+ switch o.Op {
253
+ case opAdd:
254
+ err := s.Set(ctx, o.Val, e)
255
+ if err != nil {
256
+ return nil, fmt.Errorf("inserting %s: %s", o.Val, err)
257
+ }
258
+ case opDel:
259
+ fmt.Println("deleting: ", o.Val)
260
+ err := s.Remove(ctx, o.Val)
261
+ if err != nil {
262
+ return nil, fmt.Errorf("deleting %s: %s", o.Val, err)
263
+ }
264
+ case opFind:
265
+ _, err := s.Find(ctx, o.Val)
266
+ if err != nil {
267
+ return nil, fmt.Errorf("finding %s: %s", o.Val, err)
268
+ }
269
+ }
270
+ }
271
+
272
+ return s, nil
273
+}
274
+
275
+func readCommand() string {
276
+ fmt.Print("> ")
277
+ scan := bufio.NewScanner(os.Stdin)
278
+ scan.Scan()
279
+ return scan.Text()
280
+}
unixfs/hamt/hamt_test.go
new
+552
@@ -0,0 +1,552 @@
1
+package hamt
2
+
3
+import (
4
+ "context"
5
+ "fmt"
6
+ "math/rand"
7
+ "os"
8
+ "sort"
9
+ "strings"
10
+ "testing"
11
+ "time"
12
+
13
+ dag "github.com/ipfs/go-ipfs/merkledag"
14
+ mdtest "github.com/ipfs/go-ipfs/merkledag/test"
15
+ dagutils "github.com/ipfs/go-ipfs/merkledag/utils"
16
+ ft "github.com/ipfs/go-ipfs/unixfs"
17
+)
18
+
19
+func shuffle(seed int64, arr []string) {
20
+ r := rand.New(rand.NewSource(seed))
21
+ for i := 0; i < len(arr); i++ {
22
+ a := r.Intn(len(arr))
23
+ b := r.Intn(len(arr))
24
+ arr[a], arr[b] = arr[b], arr[a]
25
+ }
26
+}
27
+
28
+func makeDir(ds dag.DAGService, size int) ([]string, *HamtShard, error) {
29
+ return makeDirWidth(ds, size, 256)
30
+}
31
+
32
+func makeDirWidth(ds dag.DAGService, size, width int) ([]string, *HamtShard, error) {
33
+ s := NewHamtShard(ds, width)
34
+
35
+ var dirs []string
36
+ for i := 0; i < size; i++ {
37
+ dirs = append(dirs, fmt.Sprintf("DIRNAME%d", i))
38
+ }
39
+
40
+ shuffle(time.Now().UnixNano(), dirs)
41
+
42
+ for i := 0; i < len(dirs); i++ {
43
+ nd := ft.EmptyDirNode()
44
+ ds.Add(nd)
45
+ err := s.Set(context.Background(), dirs[i], nd)
46
+ if err != nil {
47
+ return nil, nil, err
48
+ }
49
+ }
50
+
51
+ return dirs, s, nil
52
+}
53
+
54
+func assertLink(s *HamtShard, name string, found bool) error {
55
+ _, err := s.Find(context.Background(), name)
56
+ switch err {
57
+ case os.ErrNotExist:
58
+ if found {
59
+ return err
60
+ }
61
+
62
+ return nil
63
+ case nil:
64
+ if found {
65
+ return nil
66
+ }
67
+
68
+ return fmt.Errorf("expected not to find link named %s", name)
69
+ default:
70
+ return err
71
+ }
72
+}
73
+
74
+func assertSerializationWorks(ds dag.DAGService, s *HamtShard) error {
75
+ nd, err := s.Node()
76
+ if err != nil {
77
+ return err
78
+ }
79
+
80
+ nds, err := NewHamtFromDag(ds, nd)
81
+ if err != nil {
82
+ return err
83
+ }
84
+
85
+ linksA, err := s.EnumLinks()
86
+ if err != nil {
87
+ return err
88
+ }
89
+
90
+ linksB, err := nds.EnumLinks()
91
+ if err != nil {
92
+ return err
93
+ }
94
+
95
+ if len(linksA) != len(linksB) {
96
+ return fmt.Errorf("links arrays are different sizes")
97
+ }
98
+
99
+ for i, a := range linksA {
100
+ b := linksB[i]
101
+ if a.Name != b.Name {
102
+ return fmt.Errorf("links names mismatch")
103
+ }
104
+
105
+ if a.Cid.String() != b.Cid.String() {
106
+ return fmt.Errorf("link hashes dont match")
107
+ }
108
+
109
+ if a.Size != b.Size {
110
+ return fmt.Errorf("link sizes not the same")
111
+ }
112
+ }
113
+
114
+ return nil
115
+}
116
+
117
+func TestBasicSet(t *testing.T) {
118
+ ds := mdtest.Mock()
119
+ for _, w := range []int{128, 256, 512, 1024, 2048, 4096} {
120
+ t.Run(fmt.Sprintf("BasicSet%d", w), func(t *testing.T) {
121
+ names, s, err := makeDirWidth(ds, 1000, w)
122
+ if err != nil {
123
+ t.Fatal(err)
124
+ }
125
+ ctx := context.Background()
126
+
127
+ for _, d := range names {
128
+ _, err := s.Find(ctx, d)
129
+ if err != nil {
130
+ t.Fatal(err)
131
+ }
132
+ }
133
+ })
134
+ }
135
+}
136
+
137
+func TestDirBuilding(t *testing.T) {
138
+ ds := mdtest.Mock()
139
+ s := NewHamtShard(ds, 256)
140
+
141
+ _, s, err := makeDir(ds, 200)
142
+ if err != nil {
143
+ t.Fatal(err)
144
+ }
145
+
146
+ nd, err := s.Node()
147
+ if err != nil {
148
+ t.Fatal(err)
149
+ }
150
+
151
+ //printDag(ds, nd, 0)
152
+
153
+ k := nd.Cid()
154
+
155
+ if k.String() != "QmY89TkSEVHykWMHDmyejSWFj9CYNtvzw4UwnT9xbc4Zjc" {
156
+ t.Fatalf("output didnt match what we expected (got %s)", k.String())
157
+ }
158
+}
159
+
160
+func TestShardReload(t *testing.T) {
161
+ ds := mdtest.Mock()
162
+ s := NewHamtShard(ds, 256)
163
+ ctx := context.Background()
164
+
165
+ _, s, err := makeDir(ds, 200)
166
+ if err != nil {
167
+ t.Fatal(err)
168
+ }
169
+
170
+ nd, err := s.Node()
171
+ if err != nil {
172
+ t.Fatal(err)
173
+ }
174
+
175
+ nds, err := NewHamtFromDag(ds, nd)
176
+ if err != nil {
177
+ t.Fatal(err)
178
+ }
179
+
180
+ lnks, err := nds.EnumLinks()
181
+ if err != nil {
182
+ t.Fatal(err)
183
+ }
184
+
185
+ if len(lnks) != 200 {
186
+ t.Fatal("not enough links back")
187
+ }
188
+
189
+ _, err = nds.Find(ctx, "DIRNAME50")
190
+ if err != nil {
191
+ t.Fatal(err)
192
+ }
193
+
194
+ // Now test roundtrip marshal with no operations
195
+
196
+ nds, err = NewHamtFromDag(ds, nd)
197
+ if err != nil {
198
+ t.Fatal(err)
199
+ }
200
+
201
+ ond, err := nds.Node()
202
+ if err != nil {
203
+ t.Fatal(err)
204
+ }
205
+
206
+ outk := ond.Cid()
207
+ ndk := nd.Cid()
208
+
209
+ if !outk.Equals(ndk) {
210
+ printDiff(ds, nd.(*dag.ProtoNode), ond.(*dag.ProtoNode))
211
+ t.Fatal("roundtrip serialization failed")
212
+ }
213
+}
214
+
215
+func TestRemoveElems(t *testing.T) {
216
+ ds := mdtest.Mock()
217
+ dirs, s, err := makeDir(ds, 500)
218
+ if err != nil {
219
+ t.Fatal(err)
220
+ }
221
+ ctx := context.Background()
222
+
223
+ shuffle(time.Now().UnixNano(), dirs)
224
+
225
+ for _, d := range dirs {
226
+ err := s.Remove(ctx, d)
227
+ if err != nil {
228
+ t.Fatal(err)
229
+ }
230
+ }
231
+
232
+ nd, err := s.Node()
233
+ if err != nil {
234
+ t.Fatal(err)
235
+ }
236
+
237
+ if len(nd.Links()) > 0 {
238
+ t.Fatal("shouldnt have any links here")
239
+ }
240
+
241
+ err = s.Remove(ctx, "doesnt exist")
242
+ if err != os.ErrNotExist {
243
+ t.Fatal("expected error does not exist")
244
+ }
245
+}
246
+
247
+func TestSetAfterMarshal(t *testing.T) {
248
+ ds := mdtest.Mock()
249
+ _, s, err := makeDir(ds, 300)
250
+ if err != nil {
251
+ t.Fatal(err)
252
+ }
253
+ ctx := context.Background()
254
+
255
+ nd, err := s.Node()
256
+ if err != nil {
257
+ t.Fatal(err)
258
+ }
259
+
260
+ nds, err := NewHamtFromDag(ds, nd)
261
+ if err != nil {
262
+ t.Fatal(err)
263
+ }
264
+
265
+ empty := ft.EmptyDirNode()
266
+ for i := 0; i < 100; i++ {
267
+ err := nds.Set(ctx, fmt.Sprintf("moredirs%d", i), empty)
268
+ if err != nil {
269
+ t.Fatal(err)
270
+ }
271
+ }
272
+
273
+ links, err := nds.EnumLinks()
274
+ if err != nil {
275
+ t.Fatal(err)
276
+ }
277
+
278
+ if len(links) != 400 {
279
+ t.Fatal("expected 400 links")
280
+ }
281
+
282
+ err = assertSerializationWorks(ds, nds)
283
+ if err != nil {
284
+ t.Fatal(err)
285
+ }
286
+}
287
+
288
+func TestDuplicateAddShard(t *testing.T) {
289
+ ds := mdtest.Mock()
290
+ dir := NewHamtShard(ds, 256)
291
+ nd := new(dag.ProtoNode)
292
+ ctx := context.Background()
293
+
294
+ err := dir.Set(ctx, "test", nd)
295
+ if err != nil {
296
+ t.Fatal(err)
297
+ }
298
+
299
+ err = dir.Set(ctx, "test", nd)
300
+ if err != nil {
301
+ t.Fatal(err)
302
+ }
303
+
304
+ lnks, err := dir.EnumLinks()
305
+ if err != nil {
306
+ t.Fatal(err)
307
+ }
308
+
309
+ if len(lnks) != 1 {
310
+ t.Fatal("expected only one link")
311
+ }
312
+}
313
+
314
+func TestLoadFailsFromNonShard(t *testing.T) {
315
+ ds := mdtest.Mock()
316
+ nd := ft.EmptyDirNode()
317
+
318
+ _, err := NewHamtFromDag(ds, nd)
319
+ if err == nil {
320
+ t.Fatal("expected dir shard creation to fail when given normal directory")
321
+ }
322
+
323
+ nd = new(dag.ProtoNode)
324
+
325
+ _, err = NewHamtFromDag(ds, nd)
326
+ if err == nil {
327
+ t.Fatal("expected dir shard creation to fail when given normal directory")
328
+ }
329
+}
330
+
331
+func TestFindNonExisting(t *testing.T) {
332
+ ds := mdtest.Mock()
333
+ _, s, err := makeDir(ds, 100)
334
+ if err != nil {
335
+ t.Fatal(err)
336
+ }
337
+ ctx := context.Background()
338
+
339
+ for i := 0; i < 200; i++ {
340
+ _, err := s.Find(ctx, fmt.Sprintf("notfound%d", i))
341
+ if err != os.ErrNotExist {
342
+ t.Fatal("expected ErrNotExist")
343
+ }
344
+ }
345
+}
346
+
347
+func TestRemoveElemsAfterMarshal(t *testing.T) {
348
+ ds := mdtest.Mock()
349
+ dirs, s, err := makeDir(ds, 30)
350
+ if err != nil {
351
+ t.Fatal(err)
352
+ }
353
+ ctx := context.Background()
354
+
355
+ sort.Strings(dirs)
356
+
357
+ err = s.Remove(ctx, dirs[0])
358
+ if err != nil {
359
+ t.Fatal(err)
360
+ }
361
+
362
+ out, err := s.Find(ctx, dirs[0])
363
+ if err == nil {
364
+ t.Fatal("expected error, got: ", out)
365
+ }
366
+
367
+ nd, err := s.Node()
368
+ if err != nil {
369
+ t.Fatal(err)
370
+ }
371
+
372
+ nds, err := NewHamtFromDag(ds, nd)
373
+ if err != nil {
374
+ t.Fatal(err)
375
+ }
376
+
377
+ _, err = nds.Find(ctx, dirs[0])
378
+ if err == nil {
379
+ t.Fatal("expected not to find ", dirs[0])
380
+ }
381
+
382
+ for _, d := range dirs[1:] {
383
+ _, err := nds.Find(ctx, d)
384
+ if err != nil {
385
+ t.Fatal("could not find expected link after unmarshaling")
386
+ }
387
+ }
388
+
389
+ for _, d := range dirs[1:] {
390
+ err := nds.Remove(ctx, d)
391
+ if err != nil {
392
+ t.Fatal(err)
393
+ }
394
+ }
395
+
396
+ links, err := nds.EnumLinks()
397
+ if err != nil {
398
+ t.Fatal(err)
399
+ }
400
+
401
+ if len(links) != 0 {
402
+ t.Fatal("expected all links to be removed")
403
+ }
404
+
405
+ err = assertSerializationWorks(ds, nds)
406
+ if err != nil {
407
+ t.Fatal(err)
408
+ }
409
+}
410
+
411
+func TestBitfieldIndexing(t *testing.T) {
412
+ ds := mdtest.Mock()
413
+ s := NewHamtShard(ds, 256)
414
+
415
+ set := func(i int) {
416
+ s.bitfield.SetBit(s.bitfield, i, 1)
417
+ }
418
+
419
+ assert := func(i int, val int) {
420
+ if s.indexForBitPos(i) != val {
421
+ t.Fatalf("expected index %d to be %d", i, val)
422
+ }
423
+ }
424
+
425
+ assert(50, 0)
426
+ set(4)
427
+ set(5)
428
+ set(60)
429
+
430
+ assert(10, 2)
431
+ set(3)
432
+ assert(10, 3)
433
+ assert(1, 0)
434
+
435
+ assert(100, 4)
436
+ set(50)
437
+ assert(45, 3)
438
+ set(100)
439
+ assert(100, 5)
440
+}
441
+
442
+// test adding a sharded directory node as the child of another directory node.
443
+// if improperly implemented, the parent hamt may assume the child is a part of
444
+// itself.
445
+func TestSetHamtChild(t *testing.T) {
446
+ ds := mdtest.Mock()
447
+ s := NewHamtShard(ds, 256)
448
+ ctx := context.Background()
449
+
450
+ e := ft.EmptyDirNode()
451
+ ds.Add(e)
452
+
453
+ err := s.Set(ctx, "bar", e)
454
+ if err != nil {
455
+ t.Fatal(err)
456
+ }
457
+
458
+ snd, err := s.Node()
459
+ if err != nil {
460
+ t.Fatal(err)
461
+ }
462
+
463
+ _, ns, err := makeDir(ds, 50)
464
+ if err != nil {
465
+ t.Fatal(err)
466
+ }
467
+
468
+ err = ns.Set(ctx, "foo", snd)
469
+ if err != nil {
470
+ t.Fatal(err)
471
+ }
472
+
473
+ nsnd, err := ns.Node()
474
+ if err != nil {
475
+ t.Fatal(err)
476
+ }
477
+
478
+ hs, err := NewHamtFromDag(ds, nsnd)
479
+ if err != nil {
480
+ t.Fatal(err)
481
+ }
482
+
483
+ err = assertLink(hs, "bar", false)
484
+ if err != nil {
485
+ t.Fatal(err)
486
+ }
487
+
488
+ err = assertLink(hs, "foo", true)
489
+ if err != nil {
490
+ t.Fatal(err)
491
+ }
492
+}
493
+
494
+func printDag(ds dag.DAGService, nd *dag.ProtoNode, depth int) {
495
+ padding := strings.Repeat(" ", depth)
496
+ fmt.Println("{")
497
+ for _, l := range nd.Links() {
498
+ fmt.Printf("%s%s: %s", padding, l.Name, l.Cid.String())
499
+ ch, err := ds.Get(context.Background(), l.Cid)
500
+ if err != nil {
501
+ panic(err)
502
+ }
503
+
504
+ printDag(ds, ch.(*dag.ProtoNode), depth+1)
505
+ }
506
+ fmt.Println(padding + "}")
507
+}
508
+
509
+func printDiff(ds dag.DAGService, a, b *dag.ProtoNode) {
510
+ diff, err := dagutils.Diff(context.TODO(), ds, a, b)
511
+ if err != nil {
512
+ panic(err)
513
+ }
514
+
515
+ for _, d := range diff {
516
+ fmt.Println(d)
517
+ }
518
+}
519
+
520
+func BenchmarkHAMTSet(b *testing.B) {
521
+ ds := mdtest.Mock()
522
+ sh := NewHamtShard(ds, 256)
523
+ nd, err := sh.Node()
524
+ if err != nil {
525
+ b.Fatal(err)
526
+ }
527
+
528
+ _, err = ds.Add(nd)
529
+ if err != nil {
530
+ b.Fatal(err)
531
+ }
532
+ ds.Add(ft.EmptyDirNode())
533
+
534
+ for i := 0; i < b.N; i++ {
535
+ s, err := NewHamtFromDag(ds, nd)
536
+ if err != nil {
537
+ b.Fatal(err)
538
+ }
539
+
540
+ err = s.Set(context.TODO(), fmt.Sprint(i), ft.EmptyDirNode())
541
+ if err != nil {
542
+ b.Fatal(err)
543
+ }
544
+
545
+ out, err := s.Node()
546
+ if err != nil {
547
+ b.Fatal(err)
548
+ }
549
+
550
+ nd = out
551
+ }
552
+}
unixfs/hamt/util.go
new
+61
@@ -0,0 +1,61 @@
1
+package hamt
2
+
3
+import (
4
+ "math/big"
5
+)
6
+
7
+type hashBits struct {
8
+ b []byte
9
+ consumed int
10
+}
11
+
12
+func mkmask(n int) byte {
13
+ return (1 << uint(n)) - 1
14
+}
15
+
16
+func (hb *hashBits) Next(i int) int {
17
+ curbi := hb.consumed / 8
18
+ leftb := 8 - (hb.consumed % 8)
19
+
20
+ curb := hb.b[curbi]
21
+ if i == leftb {
22
+ out := int(mkmask(i) & curb)
23
+ hb.consumed += i
24
+ return out
25
+ } else if i < leftb {
26
+ a := curb & mkmask(leftb) // mask out the high bits we don't want
27
+ b := a & ^mkmask(leftb-i) // mask out the low bits we don't want
28
+ c := b >> uint(leftb-i) // shift whats left down
29
+ hb.consumed += i
30
+ return int(c)
31
+ } else {
32
+ out := int(mkmask(leftb) & curb)
33
+ out <<= uint(i - leftb)
34
+ hb.consumed += leftb
35
+ out += hb.Next(i - leftb)
36
+ return out
37
+ }
38
+}
39
+
40
+const (
41
+ m1 = 0x5555555555555555 //binary: 0101...
42
+ m2 = 0x3333333333333333 //binary: 00110011..
43
+ m4 = 0x0f0f0f0f0f0f0f0f //binary: 4 zeros, 4 ones ...
44
+ h01 = 0x0101010101010101 //the sum of 256 to the power of 0,1,2,3...
45
+)
46
+
47
+// from https://en.wikipedia.org/wiki/Hamming_weight
48
+func popCountUint64(x uint64) int {
49
+ x -= (x >> 1) & m1 //put count of each 2 bits into those 2 bits
50
+ x = (x & m2) + ((x >> 2) & m2) //put count of each 4 bits into those 4 bits
51
+ x = (x + (x >> 4)) & m4 //put count of each 8 bits into those 8 bits
52
+ return int((x * h01) >> 56)
53
+}
54
+
55
+func popCount(i *big.Int) int {
56
+ var n int
57
+ for _, v := range i.Bits() {
58
+ n += popCountUint64(uint64(v))
59
+ }
60
+ return n
61
+}
unixfs/hamt/util_test.go
new
+58
@@ -0,0 +1,58 @@
1
+package hamt
2
+
3
+import (
4
+ "math/big"
5
+ "testing"
6
+)
7
+
8
+func TestPopCount(t *testing.T) {
9
+ x := big.NewInt(0)
10
+
11
+ for i := 0; i < 50; i++ {
12
+ x.SetBit(x, i, 1)
13
+ }
14
+
15
+ if popCount(x) != 50 {
16
+ t.Fatal("expected popcount to be 50")
17
+ }
18
+}
19
+
20
+func TestHashBitsEvenSizes(t *testing.T) {
21
+ buf := []byte{255, 127, 79, 45, 116, 99, 35, 17}
22
+ hb := hashBits{b: buf}
23
+
24
+ for _, v := range buf {
25
+ if hb.Next(8) != int(v) {
26
+ t.Fatal("got wrong numbers back")
27
+ }
28
+ }
29
+}
30
+
31
+func TestHashBitsUneven(t *testing.T) {
32
+ buf := []byte{255, 127, 79, 45, 116, 99, 35, 17}
33
+ hb := hashBits{b: buf}
34
+
35
+ v := hb.Next(4)
36
+ if v != 15 {
37
+ t.Fatal("should have gotten 15: ", v)
38
+ }
39
+
40
+ v = hb.Next(4)
41
+ if v != 15 {
42
+ t.Fatal("should have gotten 15: ", v)
43
+ }
44
+
45
+ if v := hb.Next(3); v != 3 {
46
+ t.Fatalf("expected 3, but got %b", v)
47
+ }
48
+ if v := hb.Next(3); v != 7 {
49
+ t.Fatalf("expected 7, but got %b", v)
50
+ }
51
+ if v := hb.Next(3); v != 6 {
52
+ t.Fatalf("expected 6, but got %b", v)
53
+ }
54
+
55
+ if v := hb.Next(15); v != 20269 {
56
+ t.Fatalf("expected 20269, but got %b (%d)", v, v)
57
+ }
58
+}
unixfs/io/dirbuilder.go
+119
-23
@@ -2,48 +2,144 @@ package io
2
3
import (
4
"context"
5
+ "fmt"
6
+ "os"
7
8
mdag "github.com/ipfs/go-ipfs/merkledag"
9
format "github.com/ipfs/go-ipfs/unixfs"
8
- cid "gx/ipfs/QmV5gPoRsjN1Gid3LMdNZTyfCtP2DsvqEbMAmz82RmmiGk/go-cid"
10
+ hamt "github.com/ipfs/go-ipfs/unixfs/hamt"
11
+
12
+ node "gx/ipfs/QmYDscK7dmdo2GZ9aumS8s5auUUAH5mR1jvj5pYhWusfK7/go-ipld-node"
13
)
14
11
-type directoryBuilder struct {
15
+// ShardSplitThreshold specifies how large of an unsharded directory
16
+// the Directory code will generate. Adding entries over this value will
17
+// result in the node being restructured into a sharded object.
18
+var ShardSplitThreshold = 1000
19
+
20
+// DefaultShardWidth is the default value used for hamt sharding width.
21
+var DefaultShardWidth = 256
22
+
23
+type Directory struct {
24
dserv mdag.DAGService
25
dirnode *mdag.ProtoNode
14
-}
26
16
-// NewEmptyDirectory returns an empty merkledag Node with a folder Data chunk
17
-func NewEmptyDirectory() *mdag.ProtoNode {
18
- nd := new(mdag.ProtoNode)
19
- nd.SetData(format.FolderPBData())
20
- return nd
27
+ shard *hamt.HamtShard
28
}
29
23
-// NewDirectory returns a directoryBuilder. It needs a DAGService to add the Children
24
-func NewDirectory(dserv mdag.DAGService) *directoryBuilder {
25
- db := new(directoryBuilder)
30
+// NewDirectory returns a Directory. It needs a DAGService to add the Children
31
+func NewDirectory(dserv mdag.DAGService) *Directory {
32
+ db := new(Directory)
33
db.dserv = dserv
27
- db.dirnode = NewEmptyDirectory()
34
+ db.dirnode = format.EmptyDirNode()
35
return db
36
}
37
31
-// AddChild adds a (name, key)-pair to the root node.
32
-func (d *directoryBuilder) AddChild(ctx context.Context, name string, c *cid.Cid) error {
33
- cnode, err := d.dserv.Get(ctx, c)
38
+func NewDirectoryFromNode(dserv mdag.DAGService, nd node.Node) (*Directory, error) {
39
+ pbnd, ok := nd.(*mdag.ProtoNode)
40
+ if !ok {
41
+ return nil, mdag.ErrNotProtobuf
42
+ }
43
+
44
+ pbd, err := format.FromBytes(pbnd.Data())
45
if err != nil {
35
- return err
46
+ return nil, err
47
}
48
38
- cnpb, ok := cnode.(*mdag.ProtoNode)
39
- if !ok {
40
- return mdag.ErrNotProtobuf
49
+ switch pbd.GetType() {
50
+ case format.TDirectory:
51
+ return &Directory{
52
+ dserv: dserv,
53
+ dirnode: pbnd.Copy().(*mdag.ProtoNode),
54
+ }, nil
55
+ case format.THAMTShard:
56
+ shard, err := hamt.NewHamtFromDag(dserv, nd)
57
+ if err != nil {
58
+ return nil, err
59
+ }
60
+
61
+ return &Directory{
62
+ dserv: dserv,
63
+ shard: shard,
64
+ }, nil
65
+ default:
66
+ return nil, fmt.Errorf("merkledag node was not a directory or shard")
67
}
68
+}
69
43
- return d.dirnode.AddNodeLinkClean(name, cnpb)
70
+// AddChild adds a (name, key)-pair to the root node.
71
+func (d *Directory) AddChild(ctx context.Context, name string, nd node.Node) error {
72
+ if d.shard == nil {
73
+ if len(d.dirnode.Links()) < ShardSplitThreshold {
74
+ _ = d.dirnode.RemoveNodeLink(name)
75
+ return d.dirnode.AddNodeLinkClean(name, nd)
76
+ }
77
+
78
+ err := d.switchToSharding(ctx)
79
+ if err != nil {
80
+ return err
81
+ }
82
+ }
83
+
84
+ return d.shard.Set(ctx, name, nd)
85
}
86
46
-// GetNode returns the root of this directoryBuilder
47
-func (d *directoryBuilder) GetNode() *mdag.ProtoNode {
48
- return d.dirnode
87
+func (d *Directory) switchToSharding(ctx context.Context) error {
88
+ d.shard = hamt.NewHamtShard(d.dserv, DefaultShardWidth)
89
+ for _, lnk := range d.dirnode.Links() {
90
+ cnd, err := d.dserv.Get(ctx, lnk.Cid)
91
+ if err != nil {
92
+ return err
93
+ }
94
+
95
+ err = d.shard.Set(ctx, lnk.Name, cnd)
96
+ if err != nil {
97
+ return err
98
+ }
99
+ }
100
+
101
+ d.dirnode = nil
102
+ return nil
103
+}
104
+
105
+func (d *Directory) Links() ([]*node.Link, error) {
106
+ if d.shard == nil {
107
+ return d.dirnode.Links(), nil
108
+ }
109
+
110
+ return d.shard.EnumLinks()
111
+}
112
+
113
+func (d *Directory) Find(ctx context.Context, name string) (node.Node, error) {
114
+ if d.shard == nil {
115
+ lnk, err := d.dirnode.GetNodeLink(name)
116
+ switch err {
117
+ case mdag.ErrLinkNotFound:
118
+ return nil, os.ErrNotExist
119
+ default:
120
+ return nil, err
121
+ case nil:
122
+ }
123
+
124
+ return d.dserv.Get(ctx, lnk.Cid)
125
+ }
126
+
127
+ return d.shard.Find(ctx, name)
128
+}
129
+
130
+func (d *Directory) RemoveChild(ctx context.Context, name string) error {
131
+ if d.shard == nil {
132
+ return d.dirnode.RemoveNodeLink(name)
133
+ }
134
+
135
+ return d.shard.Remove(ctx, name)
136
+}
137
+
138
+// GetNode returns the root of this Directory
139
+func (d *Directory) GetNode() (node.Node, error) {
140
+ if d.shard == nil {
141
+ return d.dirnode, nil
142
+ }
143
+
144
+ return d.shard.Node()
145
}
unixfs/io/dirbuilder_test.go
+123
-15
@@ -2,49 +2,157 @@ package io
2
3
import (
4
"context"
5
- "io/ioutil"
5
+ "fmt"
6
"testing"
7
8
- testu "github.com/ipfs/go-ipfs/unixfs/test"
8
+ mdtest "github.com/ipfs/go-ipfs/merkledag/test"
9
+ ft "github.com/ipfs/go-ipfs/unixfs"
10
)
11
12
func TestEmptyNode(t *testing.T) {
12
- n := NewEmptyDirectory()
13
+ n := ft.EmptyDirNode()
14
if len(n.Links()) != 0 {
15
t.Fatal("empty node should have 0 links")
16
}
17
}
18
19
+func TestDirectoryGrowth(t *testing.T) {
20
+ ds := mdtest.Mock()
21
+ dir := NewDirectory(ds)
22
+ ctx := context.Background()
23
+
24
+ d := ft.EmptyDirNode()
25
+ ds.Add(d)
26
+
27
+ nelems := 10000
28
+
29
+ for i := 0; i < nelems; i++ {
30
+ err := dir.AddChild(ctx, fmt.Sprintf("dir%d", i), d)
31
+ if err != nil {
32
+ t.Fatal(err)
33
+ }
34
+ }
35
+
36
+ _, err := dir.GetNode()
37
+ if err != nil {
38
+ t.Fatal(err)
39
+ }
40
+
41
+ links, err := dir.Links()
42
+ if err != nil {
43
+ t.Fatal(err)
44
+ }
45
+
46
+ if len(links) != nelems {
47
+ t.Fatal("didnt get right number of elements")
48
+ }
49
+
50
+ dirc := d.Cid()
51
+
52
+ names := make(map[string]bool)
53
+ for _, l := range links {
54
+ names[l.Name] = true
55
+ if !l.Cid.Equals(dirc) {
56
+ t.Fatal("link wasnt correct")
57
+ }
58
+ }
59
+
60
+ for i := 0; i < nelems; i++ {
61
+ dn := fmt.Sprintf("dir%d", i)
62
+ if !names[dn] {
63
+ t.Fatal("didnt find directory: ", dn)
64
+ }
65
+
66
+ _, err := dir.Find(context.Background(), dn)
67
+ if err != nil {
68
+ t.Fatal(err)
69
+ }
70
+ }
71
+}
72
+
73
+func TestDuplicateAddDir(t *testing.T) {
74
+ ds := mdtest.Mock()
75
+ dir := NewDirectory(ds)
76
+ ctx := context.Background()
77
+ nd := ft.EmptyDirNode()
78
+
79
+ err := dir.AddChild(ctx, "test", nd)
80
+ if err != nil {
81
+ t.Fatal(err)
82
+ }
83
+
84
+ err = dir.AddChild(ctx, "test", nd)
85
+ if err != nil {
86
+ t.Fatal(err)
87
+ }
88
+
89
+ lnks, err := dir.Links()
90
+ if err != nil {
91
+ t.Fatal(err)
92
+ }
93
+
94
+ if len(lnks) != 1 {
95
+ t.Fatal("expected only one link")
96
+ }
97
+}
98
+
99
func TestDirBuilder(t *testing.T) {
19
- dserv := testu.GetDAGServ()
20
- ctx, closer := context.WithCancel(context.Background())
21
- defer closer()
22
- inbuf, node := testu.GetRandomNode(t, dserv, 1024)
23
- key := node.Cid()
100
+ ds := mdtest.Mock()
101
+ dir := NewDirectory(ds)
102
+ ctx := context.Background()
103
25
- b := NewDirectory(dserv)
104
+ child := ft.EmptyDirNode()
105
+ _, err := ds.Add(child)
106
+ if err != nil {
107
+ t.Fatal(err)
108
+ }
109
27
- b.AddChild(ctx, "random", key)
110
+ count := 5000
111
29
- dir := b.GetNode()
30
- outn, err := dir.GetLinkedProtoNode(ctx, dserv, "random")
112
+ for i := 0; i < count; i++ {
113
+ err := dir.AddChild(ctx, fmt.Sprintf("entry %d", i), child)
114
+ if err != nil {
115
+ t.Fatal(err)
116
+ }
117
+ }
118
+
119
+ dirnd, err := dir.GetNode()
120
if err != nil {
121
t.Fatal(err)
122
}
123
35
- reader, err := NewDagReader(ctx, outn, dserv)
124
+ links, err := dir.Links()
125
if err != nil {
126
t.Fatal(err)
127
}
128
40
- outbuf, err := ioutil.ReadAll(reader)
129
+ if len(links) != count {
130
+ t.Fatal("not enough links dawg", len(links), count)
131
+ }
132
+
133
+ adir, err := NewDirectoryFromNode(ds, dirnd)
134
if err != nil {
135
t.Fatal(err)
136
}
137
45
- err = testu.ArrComp(inbuf, outbuf)
138
+ links, err = adir.Links()
139
if err != nil {
140
t.Fatal(err)
141
}
142
143
+ names := make(map[string]bool)
144
+ for _, lnk := range links {
145
+ names[lnk.Name] = true
146
+ }
147
+
148
+ for i := 0; i < count; i++ {
149
+ n := fmt.Sprintf("entry %d", i)
150
+ if !names[n] {
151
+ t.Fatal("COULDNT FIND: ", n)
152
+ }
153
+ }
154
+
155
+ if len(links) != count {
156
+ t.Fatal("wrong number of links", len(links), count)
157
+ }
158
}
unixfs/io/resolve.go
+18
-18
@@ -5,26 +5,21 @@ import (
5
6
dag "github.com/ipfs/go-ipfs/merkledag"
7
ft "github.com/ipfs/go-ipfs/unixfs"
8
+ hamt "github.com/ipfs/go-ipfs/unixfs/hamt"
9
10
node "gx/ipfs/QmYDscK7dmdo2GZ9aumS8s5auUUAH5mR1jvj5pYhWusfK7/go-ipld-node"
11
)
12
13
func ResolveUnixfsOnce(ctx context.Context, ds dag.DAGService, nd node.Node, name string) (*node.Link, error) {
13
- pbnd, ok := nd.(*dag.ProtoNode)
14
- if !ok {
15
- lnk, _, err := nd.ResolveLink([]string{name})
16
- return lnk, err
17
- }
18
-
19
- upb, err := ft.FromBytes(pbnd.Data())
20
- if err != nil {
21
- // Not a unixfs node, use standard object traversal code
22
- lnk, _, err := nd.ResolveLink([]string{name})
23
- return lnk, err
24
- }
25
-
26
- switch upb.GetType() {
27
- /*
14
+ switch nd := nd.(type) {
15
+ case *dag.ProtoNode:
16
+ upb, err := ft.FromBytes(nd.Data())
17
+ if err != nil {
18
+ // Not a unixfs node, use standard object traversal code
19
+ return nd.GetNodeLink(name)
20
+ }
21
+
22
+ switch upb.GetType() {
23
case ft.THAMTShard:
24
s, err := hamt.NewHamtFromDag(ds, nd)
25
if err != nil {
@@ -37,10 +32,15 @@ func ResolveUnixfsOnce(ctx context.Context, ds dag.DAGService, nd node.Node, nam
32
return nil, err
33
}
34
40
- return dag.MakeLink(out)
41
- */
35
+ return node.MakeLink(out)
36
+ default:
37
+ return nd.GetNodeLink(name)
38
+ }
39
default:
40
lnk, _, err := nd.ResolveLink([]string{name})
44
- return lnk, err
41
+ if err != nil {
42
+ return nil, err
43
+ }
44
+ return lnk, nil
45
}
46
}
unixfs/pb/unixfs.pb.go
+19
@@ -31,6 +31,7 @@ const (
31
Data_File Data_DataType = 2
32
Data_Metadata Data_DataType = 3
33
Data_Symlink Data_DataType = 4
34
+ Data_HAMTShard Data_DataType = 5
35
)
36
37
var Data_DataType_name = map[int32]string{
@@ -39,6 +40,7 @@ var Data_DataType_name = map[int32]string{
40
2: "File",
41
3: "Metadata",
42
4: "Symlink",
43
+ 5: "HAMTShard",
44
}
45
var Data_DataType_value = map[string]int32{
46
"Raw": 0,
@@ -46,6 +48,7 @@ var Data_DataType_value = map[string]int32{
48
"File": 2,
49
"Metadata": 3,
50
"Symlink": 4,
51
+ "HAMTShard": 5,
52
}
53
54
func (x Data_DataType) Enum() *Data_DataType {
@@ -70,6 +73,8 @@ type Data struct {
73
Data []byte `protobuf:"bytes,2,opt,name=Data" json:"Data,omitempty"`
74
Filesize *uint64 `protobuf:"varint,3,opt,name=filesize" json:"filesize,omitempty"`
75
Blocksizes []uint64 `protobuf:"varint,4,rep,name=blocksizes" json:"blocksizes,omitempty"`
76
+ HashType *uint64 `protobuf:"varint,5,opt,name=hashType" json:"hashType,omitempty"`
77
+ Fanout *uint64 `protobuf:"varint,6,opt,name=fanout" json:"fanout,omitempty"`
78
XXX_unrecognized []byte `json:"-"`
79
}
80
@@ -105,6 +110,20 @@ func (m *Data) GetBlocksizes() []uint64 {
110
return nil
111
}
112
113
+func (m *Data) GetHashType() uint64 {
114
+ if m != nil && m.HashType != nil {
115
+ return *m.HashType
116
+ }
117
+ return 0
118
+}
119
+
120
+func (m *Data) GetFanout() uint64 {
121
+ if m != nil && m.Fanout != nil {
122
+ return *m.Fanout
123
+ }
124
+ return 0
125
+}
126
+
127
type Metadata struct {
128
MimeType *string `protobuf:"bytes,1,opt,name=MimeType" json:"MimeType,omitempty"`
129
XXX_unrecognized []byte `json:"-"`
unixfs/pb/unixfs.proto
+4
@@ -7,12 +7,16 @@ message Data {
7
File = 2;
8
Metadata = 3;
9
Symlink = 4;
10
+ HAMTShard = 5;
11
}
12
13
required DataType Type = 1;
14
optional bytes Data = 2;
15
optional uint64 filesize = 3;
16
repeated uint64 blocksizes = 4;
17
+
18
+ optional uint64 hashType = 5;
19
+ optional uint64 fanout = 6;
20
}
21
22
message Metadata {