@cryptotaxi247 / kubo / commits / 56982b4b3

do not hold locks for multiple filesystem nodes at the same time

License: MIT Signed-off-by: Jeromy <jeromyj@gmail.com>

Jeromy committed Dec 18, 2015 at 21:59 UTC 56982b4b3ada89a1eb638ec28f3dcfb47f581b44
3 files changed +76 -24
mfs/dir.go
+29 -10
@@ -50,19 +50,34 @@ func NewDirectory(ctx context.Context, name string, node *dag.Node, parent child
50 // closeChild updates the child by the given name to the dag node 'nd'
51 // and changes its own dag node, then propogates the changes upward
52 func (d *Directory) closeChild(name string, nd *dag.Node) error {
53 - _, err := d.dserv.Add(nd)
53 + mynd, err := d.closeChildUpdate(name, nd)
54 if err != nil {
55 return err
56 }
57
58 + return d.parent.closeChild(d.name, mynd)
59 +}
60 +
61 +// closeChildUpdate is the portion of closeChild that needs to be locked around
62 +func (d *Directory) closeChildUpdate(name string, nd *dag.Node) (*dag.Node, error) {
63 d.lock.Lock()
64 defer d.lock.Unlock()
60 - err = d.updateChild(name, nd)
65 +
66 + err := d.updateChild(name, nd)
67 if err != nil {
62 - return err
68 + return nil, err
69 }
70
65 - return d.parent.closeChild(d.name, d.node)
71 + return d.flushCurrentNode()
72 +}
73 +
74 +func (d *Directory) flushCurrentNode() (*dag.Node, error) {
75 + _, err := d.dserv.Add(d.node)
76 + if err != nil {
77 + return nil, err
78 + }
79 +
80 + return d.node.Copy(), nil
81 }
82
83 func (d *Directory) updateChild(name string, nd *dag.Node) error {
@@ -263,7 +278,7 @@ func (d *Directory) Mkdir(name string) (*Directory, error) {
278 return nil, err
279 }
280
266 - err = d.parent.closeChild(d.name, d.node)
281 + err = d.flushUp()
282 if err != nil {
283 return nil, err
284 }
@@ -285,13 +300,18 @@ func (d *Directory) Unlink(name string) error {
300 return err
301 }
302
303 + return d.flushUp()
304 +}
305 +
306 +func (d *Directory) flushUp() error {
307 +
308 return d.parent.closeChild(d.name, d.node)
309 }
310
311 // AddChild adds the node 'nd' under this directory giving it the name 'name'
312 func (d *Directory) AddChild(name string, nd *dag.Node) error {
293 - d.Lock()
294 - defer d.Unlock()
313 + d.lock.Lock()
314 + defer d.lock.Unlock()
315
316 _, err := d.childUnsync(name)
317 if err == nil {
@@ -310,7 +330,6 @@ func (d *Directory) AddChild(name string, nd *dag.Node) error {
330
331 d.modTime = time.Now()
332
313 - //return d.parent.closeChild(d.name, d.node)
333 return nil
334 }
335
@@ -353,8 +372,8 @@ func (d *Directory) sync() error {
372 }
373
374 func (d *Directory) GetNode() (*dag.Node, error) {
356 - d.Lock()
357 - defer d.Unlock()
375 + d.lock.Lock()
376 + defer d.lock.Unlock()
377
378 err := d.sync()
379 if err != nil {
mfs/file.go
+30 -14
@@ -16,8 +16,9 @@ type File struct {
16 name string
17 hasChanges bool
18
19 - mod *mod.DagModifier
20 - lock sync.Mutex
19 + dserv dag.DAGService
20 + mod *mod.DagModifier
21 + lock sync.Mutex
22 }
23
24 // NewFile returns a NewFile object with the given parameters
@@ -28,6 +29,7 @@ func NewFile(name string, node *dag.Node, parent childCloser, dserv dag.DAGServi
29 }
30
31 return &File{
32 + dserv: dserv,
33 parent: parent,
34 name: name,
35 mod: dmod,
@@ -60,29 +62,43 @@ func (fi *File) CtxReadFull(ctx context.Context, b []byte) (int, error) {
62 // and signals a republish to occur
63 func (fi *File) Close() error {
64 fi.Lock()
63 - defer fi.Unlock()
65 if fi.hasChanges {
66 err := fi.mod.Sync()
67 if err != nil {
68 return err
69 }
70
70 - nd, err := fi.mod.GetNode()
71 - if err != nil {
72 - return err
73 - }
71 + fi.hasChanges = false
72 +
73 + // explicitly stay locked for flushUp call,
74 + // it will manage the lock for us
75 + return fi.flushUp()
76 + }
77 +
78 + return nil
79 +}
80
81 +// flushUp syncs the file and adds it to the dagservice
82 +// it *must* be called with the File's lock taken
83 +func (fi *File) flushUp() error {
84 + nd, err := fi.mod.GetNode()
85 + if err != nil {
86 fi.Unlock()
76 - err = fi.parent.closeChild(fi.name, nd)
77 - fi.Lock()
78 - if err != nil {
79 - return err
80 - }
87 + return err
88 + }
89
82 - fi.hasChanges = false
90 + _, err = fi.dserv.Add(nd)
91 + if err != nil {
92 + fi.Unlock()
93 + return err
94 }
95
85 - return nil
96 + name := fi.name
97 + parent := fi.parent
98 +
99 + // explicit unlock *only* before closeChild call
100 + fi.Unlock()
101 + return parent.closeChild(name, nd)
102 }
103
104 // Sync flushes the changes in the file to disk
mfs/system.go
+17
@@ -109,6 +109,23 @@ func (kr *Root) GetValue() FSNode {
109 return kr.val
110 }
111
112 +func (kr *Root) Flush() error {
113 + nd, err := kr.GetValue().GetNode()
114 + if err != nil {
115 + return err
116 + }
117 +
118 + k, err := kr.dserv.Add(nd)
119 + if err != nil {
120 + return err
121 + }
122 +
123 + if kr.repub != nil {
124 + kr.repub.Update(k)
125 + }
126 + return nil
127 +}
128 +
129 // closeChild implements the childCloser interface, and signals to the publisher that
130 // there are changes ready to be published
131 func (kr *Root) closeChild(name string, nd *dag.Node) error {