fixed panic and quota

This commit is contained in:
ston1th 2022-10-08 10:34:16 +02:00
commit 37a6b0ae5f
3 changed files with 41 additions and 25 deletions

View file

@ -28,8 +28,6 @@ var (
log logr.Logger log logr.Logger
) )
const gib = 1024 * 1024 * 1024
func main() { func main() {
klog.InitFlags(nil) klog.InitFlags(nil)
flag.StringVar(&cfg, "config", "", "path to config file") flag.StringVar(&cfg, "config", "", "path to config file")

View file

@ -74,7 +74,7 @@ func Validate(cfg *Config) error {
if cfg.Server.HTTP.Addr == "" { if cfg.Server.HTTP.Addr == "" {
cfg.Server.HTTP.Addr = defListenHTTP cfg.Server.HTTP.Addr = defListenHTTP
} }
c := cfg.Cache c := &cfg.Cache
if c.Metadata == "" { if c.Metadata == "" {
return errors.New("cache.metadata is empty") return errors.New("cache.metadata is empty")
} }

View file

@ -199,7 +199,14 @@ func (mh *MetadataHandler) Metadata(name string, size int64) (md *Metadata) {
} }
return md 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 mh.md[name] = md
return return
} }
@ -221,6 +228,7 @@ func (mh *MetadataHandler) init() {
} }
return true return true
} }
v.ctx, v.cancel = context.WithCancel(context.Background())
v.fs = mh.fs v.fs = mh.fs
v.name = k v.name = k
return false return false
@ -352,6 +360,8 @@ type Metadata struct {
f provider.File `json:"-"` f provider.File `json:"-"`
wc chan writeAt `json:"-"` wc chan writeAt `json:"-"`
done chan struct{} `json:"-"` done chan struct{} `json:"-"`
ctx context.Context `json:"-"`
cancel func() `json:"-"`
err atomic.Pointer[mdErr] `json:"-"` err atomic.Pointer[mdErr] `json:"-"`
name string `json:"-"` name string `json:"-"`
Size int64 `json:"s"` Size int64 `json:"s"`
@ -366,13 +376,8 @@ func (md *Metadata) Close() (err error) {
} }
func (md *Metadata) close() (err error) { func (md *Metadata) close() (err error) {
if md.wc != nil { md.cancel()
close(md.wc) md.wc = nil
if md.done != nil {
<-md.done
}
md.wc = nil
}
if md.f != nil { if md.f != nil {
md.f.Sync() md.f.Sync()
err = md.f.Close() err = md.f.Close()
@ -439,31 +444,37 @@ type writeAt struct {
func (md *Metadata) resetChans() { 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() { func (md *Metadata) writer() {
md.resetChans() md.resetChans()
md.err.Store(nil) md.err.Store(nil)
go func() { go func() {
for wa := range md.wc { for {
atomic.StoreInt64(&md.Atime, now()) select {
n, err := md.f.WriteAt(wa.data, wa.pos) case <-md.ctx.Done():
if err == nil { return
md.addChunk(wa.pos, n) case wa := <-md.wc:
md.fs.q.Add(n) atomic.StoreInt64(&md.Atime, now())
} else { n, err := md.f.WriteAt(wa.data, wa.pos)
md.err.Store(&mdErr{err}) if err == nil {
} md.addChunk(wa.pos, n)
if wa.ret != nil { md.fs.q.Add(n)
streamingPool.Put(wa.ret) } 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) { func (md *Metadata) WriteAt(data []byte, pos int64) (int, error) {
if md.ctx.Err() == context.Canceled {
return 0, context.Canceled
}
err := md.err.Load() err := md.err.Load()
if err != nil && err.err != nil { if err != nil && err.err != nil {
return 0, err.err return 0, err.err
@ -486,7 +497,10 @@ func (md *Metadata) WriteAt(data []byte, pos int64) (int, error) {
} }
copy(buf, data) copy(buf, data)
wa.data = buf[:n] wa.data = buf[:n]
md.wc <- wa select {
case <-md.ctx.Done():
case md.wc <- wa:
}
return n, nil return n, nil
} }
@ -501,6 +515,10 @@ func (md *Metadata) openCacheFile() error {
return err return err
} }
md.f = f md.f = f
if md.ctx.Err() != nil {
md.cancel()
md.ctx, md.cancel = context.WithCancel(context.Background())
}
md.writer() md.writer()
return nil return nil
} }