Synchronize virtfs rdwr state
This commit is contained in:
parent
c0ce185051
commit
562e1ef780
|
|
@ -5,6 +5,7 @@ import (
|
|||
"fmt"
|
||||
"os"
|
||||
"strings"
|
||||
"sync"
|
||||
)
|
||||
|
||||
// BuildTree walks an FsNodeDecl spec and produces a *Tree.
|
||||
|
|
@ -257,31 +258,48 @@ func openLeaf(d *FsNodeDecl, uid, gid string, mode os.FileMode, name string) (Fi
|
|||
isStreaming = true
|
||||
case d.Rdwr != nil:
|
||||
fn := d.Rdwr
|
||||
var stateMu sync.Mutex
|
||||
var writeMu sync.Mutex
|
||||
var result []byte
|
||||
var reqCtx context.Context
|
||||
var reqCancel context.CancelFunc
|
||||
blockFn = notBlocking
|
||||
readFn = func() ([]byte, error) {
|
||||
stateMu.Lock()
|
||||
if result != nil {
|
||||
return result, nil
|
||||
out := make([]byte, len(result))
|
||||
copy(out, result)
|
||||
stateMu.Unlock()
|
||||
return out, nil
|
||||
}
|
||||
stateMu.Unlock()
|
||||
if d.Read != nil {
|
||||
return d.Read()
|
||||
}
|
||||
return nil, nil
|
||||
}
|
||||
writeFn = func(data []byte) error {
|
||||
reqCtx, reqCancel = context.WithCancel(context.Background())
|
||||
r, err := fn(reqCtx, data)
|
||||
if err != nil {
|
||||
return err
|
||||
writeMu.Lock()
|
||||
defer writeMu.Unlock()
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
stateMu.Lock()
|
||||
reqCancel = cancel
|
||||
stateMu.Unlock()
|
||||
r, err := fn(ctx, data)
|
||||
stateMu.Lock()
|
||||
reqCancel = nil
|
||||
if err == nil {
|
||||
result = append(result[:0], r...)
|
||||
}
|
||||
result = r
|
||||
return nil
|
||||
stateMu.Unlock()
|
||||
cancel()
|
||||
return err
|
||||
}
|
||||
closeFn := func() error {
|
||||
if reqCancel != nil {
|
||||
reqCancel()
|
||||
stateMu.Lock()
|
||||
cancel := reqCancel
|
||||
stateMu.Unlock()
|
||||
if cancel != nil {
|
||||
cancel()
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
|
|
|||
Loading…
Reference in New Issue