fixed preload deadlock
This commit is contained in:
parent
b010b3a93e
commit
70414c4941
1 changed files with 59 additions and 33 deletions
|
|
@ -304,11 +304,12 @@ func (mh *MetadataHandler) flush() {
|
||||||
func (mh *MetadataHandler) close() {
|
func (mh *MetadataHandler) close() {
|
||||||
mh.mu.Lock()
|
mh.mu.Lock()
|
||||||
for _, md := range mh.md {
|
for _, md := range mh.md {
|
||||||
if md.f != nil {
|
md.Close()
|
||||||
close(md.wc)
|
//if md.f != nil {
|
||||||
md.f.Sync()
|
// close(md.wc)
|
||||||
md.f.Close()
|
// md.f.Sync()
|
||||||
}
|
// md.f.Close()
|
||||||
|
//}
|
||||||
}
|
}
|
||||||
mh.mu.Unlock()
|
mh.mu.Unlock()
|
||||||
mh.flush()
|
mh.flush()
|
||||||
|
|
@ -345,30 +346,43 @@ func IsIOErr(err error) bool {
|
||||||
}
|
}
|
||||||
|
|
||||||
type Metadata struct {
|
type Metadata struct {
|
||||||
mu sync.RWMutex `json:"-"`
|
mu sync.Mutex `json:"-"`
|
||||||
fs *FS `json:"-"`
|
cmu sync.RWMutex `json:"-"`
|
||||||
f provider.File `json:"-"`
|
fs *FS `json:"-"`
|
||||||
wc chan writeAt `json:"-"`
|
f provider.File `json:"-"`
|
||||||
err error `json:"-"`
|
wc chan writeAt `json:"-"`
|
||||||
name string `json:"-"`
|
done chan struct{} `json:"-"`
|
||||||
Size int64 `json:"s"`
|
err atomic.Pointer[mdErr] `json:"-"`
|
||||||
Atime int64 `json:"a"`
|
name string `json:"-"`
|
||||||
Chunks chunk.Chunks `json:"c"`
|
Size int64 `json:"s"`
|
||||||
|
Atime int64 `json:"a"`
|
||||||
|
Chunks chunk.Chunks `json:"c"`
|
||||||
}
|
}
|
||||||
|
|
||||||
func (md *Metadata) Close() error {
|
func (md *Metadata) Close() (err error) {
|
||||||
md.mu.Lock()
|
md.mu.Lock()
|
||||||
defer md.mu.Unlock()
|
defer md.mu.Unlock()
|
||||||
close(md.wc)
|
if md.wc != nil {
|
||||||
err := md.f.Close()
|
close(md.wc)
|
||||||
md.f = nil
|
if md.done != nil {
|
||||||
return err
|
<-md.done
|
||||||
|
}
|
||||||
|
md.wc = nil
|
||||||
|
}
|
||||||
|
if md.f != nil {
|
||||||
|
md.f.Sync()
|
||||||
|
err = md.f.Close()
|
||||||
|
md.f = nil
|
||||||
|
}
|
||||||
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
func (md *Metadata) Delete() error {
|
func (md *Metadata) Delete() error {
|
||||||
md.mu.Lock()
|
md.mu.Lock()
|
||||||
defer md.mu.Unlock()
|
defer md.mu.Unlock()
|
||||||
|
md.cmu.Lock()
|
||||||
md.Chunks = chunk.Chunks{}
|
md.Chunks = chunk.Chunks{}
|
||||||
|
md.cmu.Unlock()
|
||||||
close(md.wc)
|
close(md.wc)
|
||||||
md.f.Close()
|
md.f.Close()
|
||||||
md.f = nil
|
md.f = nil
|
||||||
|
|
@ -377,26 +391,26 @@ func (md *Metadata) Delete() error {
|
||||||
}
|
}
|
||||||
|
|
||||||
func (md *Metadata) FullyCached() bool {
|
func (md *Metadata) FullyCached() bool {
|
||||||
md.mu.RLock()
|
md.cmu.RLock()
|
||||||
defer md.mu.RUnlock()
|
defer md.cmu.RUnlock()
|
||||||
return len(md.Chunks) == 1 && md.Size == md.Chunks[0][1]
|
return len(md.Chunks) == 1 && md.Size == md.Chunks[0][1]
|
||||||
}
|
}
|
||||||
|
|
||||||
func (md *Metadata) HasChunk(off int64, n int) bool {
|
func (md *Metadata) HasChunk(off int64, n int) bool {
|
||||||
md.mu.RLock()
|
md.cmu.RLock()
|
||||||
defer md.mu.RUnlock()
|
defer md.cmu.RUnlock()
|
||||||
return md.Chunks.Exists(off, n, md.Size)
|
return md.Chunks.Exists(off, n, md.Size)
|
||||||
}
|
}
|
||||||
|
|
||||||
func (md *Metadata) ChunkSize() int64 {
|
func (md *Metadata) ChunkSize() int64 {
|
||||||
md.mu.RLock()
|
md.cmu.RLock()
|
||||||
defer md.mu.RUnlock()
|
defer md.cmu.RUnlock()
|
||||||
return md.Chunks.Size()
|
return md.Chunks.Size()
|
||||||
}
|
}
|
||||||
|
|
||||||
func (md *Metadata) addChunk(off int64, n int) {
|
func (md *Metadata) addChunk(off int64, n int) {
|
||||||
md.mu.Lock()
|
md.cmu.Lock()
|
||||||
defer md.mu.Unlock()
|
defer md.cmu.Unlock()
|
||||||
md.Chunks.Add(off, n)
|
md.Chunks.Add(off, n)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -411,14 +425,24 @@ func (md *Metadata) addChunk(off int64, n int) {
|
||||||
// return md.f.ReadAt(p, pos)
|
// return md.f.ReadAt(p, pos)
|
||||||
//}
|
//}
|
||||||
|
|
||||||
|
type mdErr struct {
|
||||||
|
err error
|
||||||
|
}
|
||||||
|
|
||||||
type writeAt struct {
|
type writeAt struct {
|
||||||
pos int64
|
pos int64
|
||||||
data []byte
|
data []byte
|
||||||
ret *[]byte
|
ret *[]byte
|
||||||
}
|
}
|
||||||
|
|
||||||
func (md *Metadata) writer() {
|
func (md *Metadata) resetChans() {
|
||||||
md.wc = make(chan writeAt, 10)
|
md.wc = make(chan writeAt, 10)
|
||||||
|
md.done = make(chan struct{})
|
||||||
|
}
|
||||||
|
|
||||||
|
func (md *Metadata) writer() {
|
||||||
|
md.resetChans()
|
||||||
|
md.err.Store(nil)
|
||||||
go func() {
|
go func() {
|
||||||
for wa := range md.wc {
|
for wa := range md.wc {
|
||||||
atomic.StoreInt64(&md.Atime, now())
|
atomic.StoreInt64(&md.Atime, now())
|
||||||
|
|
@ -427,18 +451,20 @@ func (md *Metadata) writer() {
|
||||||
md.addChunk(wa.pos, n)
|
md.addChunk(wa.pos, n)
|
||||||
md.fs.q.Add(n)
|
md.fs.q.Add(n)
|
||||||
} else {
|
} else {
|
||||||
md.err = err
|
md.err.Store(&mdErr{err})
|
||||||
}
|
}
|
||||||
if wa.ret != nil {
|
if wa.ret != nil {
|
||||||
streamingPool.Put(wa.ret)
|
streamingPool.Put(wa.ret)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
close(md.done)
|
||||||
}()
|
}()
|
||||||
}
|
}
|
||||||
|
|
||||||
func (md *Metadata) WriteAt(data []byte, pos int64) (int, error) {
|
func (md *Metadata) WriteAt(data []byte, pos int64) (int, error) {
|
||||||
if md.err != nil {
|
err := md.err.Load()
|
||||||
return 0, md.err
|
if err != nil && err.err != nil {
|
||||||
|
return 0, err.err
|
||||||
}
|
}
|
||||||
if md.f == nil {
|
if md.f == nil {
|
||||||
err := md.openCacheFile()
|
err := md.openCacheFile()
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue