From 37a6b0ae5ff48137086532a28ac55ab195d5698f Mon Sep 17 00:00:00 2001 From: ston1th Date: Sat, 8 Oct 2022 10:34:16 +0200 Subject: [PATCH] fixed panic and quota --- cmd/cachefs/main.go | 2 -- pkg/config/config.go | 2 +- pkg/fs/metadata.go | 62 ++++++++++++++++++++++++++++---------------- 3 files changed, 41 insertions(+), 25 deletions(-) diff --git a/cmd/cachefs/main.go b/cmd/cachefs/main.go index b39a7bf..c2b6799 100644 --- a/cmd/cachefs/main.go +++ b/cmd/cachefs/main.go @@ -28,8 +28,6 @@ var ( log logr.Logger ) -const gib = 1024 * 1024 * 1024 - func main() { klog.InitFlags(nil) flag.StringVar(&cfg, "config", "", "path to config file") diff --git a/pkg/config/config.go b/pkg/config/config.go index c13b948..7f95139 100644 --- a/pkg/config/config.go +++ b/pkg/config/config.go @@ -74,7 +74,7 @@ func Validate(cfg *Config) error { if cfg.Server.HTTP.Addr == "" { cfg.Server.HTTP.Addr = defListenHTTP } - c := cfg.Cache + c := &cfg.Cache if c.Metadata == "" { return errors.New("cache.metadata is empty") } diff --git a/pkg/fs/metadata.go b/pkg/fs/metadata.go index b736110..bc2befa 100644 --- a/pkg/fs/metadata.go +++ b/pkg/fs/metadata.go @@ -199,7 +199,14 @@ func (mh *MetadataHandler) Metadata(name string, size int64) (md *Metadata) { } return md } - md = &Metadata{Size: size, fs: mh.fs, name: name} + ctx, cancel := context.WithCancel(context.Background()) + md = &Metadata{ + Size: size, + fs: mh.fs, + name: name, + ctx: ctx, + cancel: cancel, + } mh.md[name] = md return } @@ -221,6 +228,7 @@ func (mh *MetadataHandler) init() { } return true } + v.ctx, v.cancel = context.WithCancel(context.Background()) v.fs = mh.fs v.name = k return false @@ -352,6 +360,8 @@ type Metadata struct { f provider.File `json:"-"` wc chan writeAt `json:"-"` done chan struct{} `json:"-"` + ctx context.Context `json:"-"` + cancel func() `json:"-"` err atomic.Pointer[mdErr] `json:"-"` name string `json:"-"` Size int64 `json:"s"` @@ -366,13 +376,8 @@ func (md *Metadata) Close() (err error) { } func (md *Metadata) close() (err error) { - if md.wc != nil { - close(md.wc) - if md.done != nil { - <-md.done - } - md.wc = nil - } + md.cancel() + md.wc = nil if md.f != nil { md.f.Sync() err = md.f.Close() @@ -439,31 +444,37 @@ type writeAt struct { func (md *Metadata) resetChans() { md.wc = make(chan writeAt, 10) - md.done = make(chan struct{}) } func (md *Metadata) writer() { md.resetChans() md.err.Store(nil) go func() { - for wa := range md.wc { - atomic.StoreInt64(&md.Atime, now()) - n, err := md.f.WriteAt(wa.data, wa.pos) - if err == nil { - md.addChunk(wa.pos, n) - md.fs.q.Add(n) - } else { - md.err.Store(&mdErr{err}) - } - if wa.ret != nil { - streamingPool.Put(wa.ret) + for { + select { + case <-md.ctx.Done(): + return + case wa := <-md.wc: + atomic.StoreInt64(&md.Atime, now()) + n, err := md.f.WriteAt(wa.data, wa.pos) + if err == nil { + md.addChunk(wa.pos, n) + md.fs.q.Add(n) + } else { + md.err.Store(&mdErr{err}) + } + if wa.ret != nil { + streamingPool.Put(wa.ret) + } } } - close(md.done) }() } func (md *Metadata) WriteAt(data []byte, pos int64) (int, error) { + if md.ctx.Err() == context.Canceled { + return 0, context.Canceled + } err := md.err.Load() if err != nil && err.err != nil { return 0, err.err @@ -486,7 +497,10 @@ func (md *Metadata) WriteAt(data []byte, pos int64) (int, error) { } copy(buf, data) wa.data = buf[:n] - md.wc <- wa + select { + case <-md.ctx.Done(): + case md.wc <- wa: + } return n, nil } @@ -501,6 +515,10 @@ func (md *Metadata) openCacheFile() error { return err } md.f = f + if md.ctx.Err() != nil { + md.cancel() + md.ctx, md.cancel = context.WithCancel(context.Background()) + } md.writer() return nil }