From 36f4929e5224ec13b2fed509c97670e3209812b6 Mon Sep 17 00:00:00 2001 From: Pavel Karpy Date: Thu, 9 Jun 2022 19:45:51 +0300 Subject: [PATCH] [#1507] node: Do not handle object concurrently by the policer Cache object that are being processed. That prevents concurrent object handling when there is a few number of objects and object handling takes more time that the policer needs for starting that object handling one more time. Signed-off-by: Pavel Karpy --- CHANGELOG.md | 1 + pkg/services/policer/policer.go | 31 +++++++++++++++++++++++++++++++ pkg/services/policer/process.go | 15 ++++++++++++--- 3 files changed, 44 insertions(+), 3 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index f294cb7b..5637cccf 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -10,6 +10,7 @@ Changelog for NeoFS Node ### Fixed - Do not replicate object twice to the same node (#1410) +- Concurrent object handling by the Policer (#1411) ### Removed diff --git a/pkg/services/policer/policer.go b/pkg/services/policer/policer.go index ac747f58..21807c32 100644 --- a/pkg/services/policer/policer.go +++ b/pkg/services/policer/policer.go @@ -1,6 +1,7 @@ package policer import ( + "sync" "time" lru "github.com/hashicorp/golang-lru" @@ -22,12 +23,39 @@ type nodeLoader interface { ObjectServiceLoad() float64 } +type objectsInWork struct { + m sync.RWMutex + objs map[oid.Address]struct{} +} + +func (oiw *objectsInWork) inWork(addr oid.Address) bool { + oiw.m.RLock() + _, ok := oiw.objs[addr] + oiw.m.RUnlock() + + return ok +} + +func (oiw *objectsInWork) remove(addr oid.Address) { + oiw.m.Lock() + delete(oiw.objs, addr) + oiw.m.Unlock() +} + +func (oiw *objectsInWork) add(addr oid.Address) { + oiw.m.Lock() + oiw.objs[addr] = struct{}{} + oiw.m.Unlock() +} + // Policer represents the utility that verifies // compliance with the object storage policy. type Policer struct { *cfg cache *lru.Cache + + objsInWork *objectsInWork } // Option is an option for Policer constructor. @@ -95,6 +123,9 @@ func New(opts ...Option) *Policer { return &Policer{ cfg: c, cache: cache, + objsInWork: &objectsInWork{ + objs: make(map[oid.Address]struct{}, c.maxCapacity), + }, } } diff --git a/pkg/services/policer/process.go b/pkg/services/policer/process.go index 6834d735..ca360968 100644 --- a/pkg/services/policer/process.go +++ b/pkg/services/policer/process.go @@ -48,15 +48,24 @@ func (p *Policer) shardPolicyWorker(ctx context.Context) { return default: addr := addrs[i] - addrStr := addr.EncodeToString() + if p.objsInWork.inWork(addr) { + // do not process an object + // that is in work + continue + } + err = p.taskPool.Submit(func() { - v, ok := p.cache.Get(addrStr) + v, ok := p.cache.Get(addr) if ok && time.Since(v.(time.Time)) < p.evictDuration { return } + p.objsInWork.add(addr) + p.processObject(ctx, addr) - p.cache.Add(addrStr, time.Now()) + + p.cache.Add(addr, time.Now()) + p.objsInWork.remove(addr) }) if err != nil { p.log.Warn("pool submission", zap.Error(err))