From f0e631799ebed5cd4229305cfccbab873d704f3c Mon Sep 17 00:00:00 2001 From: baihailong Date: Tue, 3 Sep 2024 09:48:59 +0800 Subject: [PATCH] fix(sdk): client update extents from meta and drop extents in cache leadto ek conflict. Signed-off-by: baihailong --- sdk/data/stream/extent_cache.go | 2 +- sdk/data/stream/stream_reader.go | 5 +++++ sdk/data/stream/stream_writer.go | 2 ++ 3 files changed, 8 insertions(+), 1 deletion(-) diff --git a/sdk/data/stream/extent_cache.go b/sdk/data/stream/extent_cache.go index f656577b9..2135256bc 100644 --- a/sdk/data/stream/extent_cache.go +++ b/sdk/data/stream/extent_cache.go @@ -121,7 +121,7 @@ func (cache *ExtentCache) update(gen, size uint64, force bool, eks []proto.Exten cache.root.Clear(false) for _, ek := range eks { extent := ek - log.LogDebugf("action[update] update cache replace or insert ek [%v]", ek.String()) + log.LogDebugf("action[update] update cache ino(%v) replace or insert ek [%v]", cache.inode, ek.String()) cache.root.ReplaceOrInsert(&extent) } } diff --git a/sdk/data/stream/stream_reader.go b/sdk/data/stream/stream_reader.go index 92ee0e67c..b3883d2fa 100644 --- a/sdk/data/stream/stream_reader.go +++ b/sdk/data/stream/stream_reader.go @@ -54,6 +54,7 @@ type Streamer struct { pendingCache chan bcacheKey verSeq uint64 needUpdateVer int32 + extentsLock sync.Mutex } type bcacheKey struct { @@ -91,6 +92,8 @@ func (s *Streamer) String() string { // TODO should we call it RefreshExtents instead? func (s *Streamer) GetExtents() error { + s.extentsLock.Lock() + defer s.extentsLock.Unlock() if s.client.disableMetaCache || !s.needBCache { return s.extents.RefreshForce(s.inode, s.client.getExtents) } @@ -99,6 +102,8 @@ func (s *Streamer) GetExtents() error { } func (s *Streamer) GetExtentsForce() error { + s.extentsLock.Lock() + defer s.extentsLock.Unlock() return s.extents.RefreshForce(s.inode, s.client.getExtents) } diff --git a/sdk/data/stream/stream_writer.go b/sdk/data/stream/stream_writer.go index 8c104e200..4d259df08 100644 --- a/sdk/data/stream/stream_writer.go +++ b/sdk/data/stream/stream_writer.go @@ -299,6 +299,8 @@ func (s *Streamer) handleRequest(request interface{}) { } func (s *Streamer) write(data []byte, offset, size, flags int, checkFunc func() error) (total int, err error) { + s.extentsLock.Lock() + defer s.extentsLock.Unlock() var ( direct bool retryTimes int8