From f244448a660a9ce454cce9fa4d2fcb975ca93035 Mon Sep 17 00:00:00 2001 From: IANTHEREAL Date: Sat, 11 Apr 2026 20:45:07 +0000 Subject: [PATCH 1/4] feat(tikv): add keyspace API V2 support via URL parameter When `?keyspace=` is present in the TiKV meta URL, the client uses API V2 with the specified keyspace. This enables per-tenant metadata isolation on shared TiKV clusters. Usage: tikv://pd:2379?keyspace=jfs_t_tenant1 tikv://pd:2379?keyspace=jfs_t_tenant1&ca=/tls/ca.crt&cert=/tls/tls.crt&key=/tls/tls.key Without `?keyspace=`, behavior is unchanged (default keyspace, API V1). Co-Authored-By: Claude Opus 4.6 (1M context) --- pkg/meta/tkv_tikv.go | 12 +++++++++++- 1 file changed, 11 insertions(+), 1 deletion(-) diff --git a/pkg/meta/tkv_tikv.go b/pkg/meta/tkv_tikv.go index cc1f75c87aa1..a6586d7c592d 100644 --- a/pkg/meta/tkv_tikv.go +++ b/pkg/meta/tkv_tikv.go @@ -32,6 +32,7 @@ import ( "github.com/pkg/errors" "github.com/sirupsen/logrus" "github.com/tikv/client-go/v2/config" + "github.com/pingcap/kvproto/pkg/kvrpcpb" tikverr "github.com/tikv/client-go/v2/error" "github.com/tikv/client-go/v2/oracle" "github.com/tikv/client-go/v2/tikv" @@ -86,7 +87,16 @@ func newTikvClient(addr string) (tkvClient, error) { } logger.Infof("TiKV gc interval is set to %s", interval) - client, err := txnkv.NewClient(strings.Split(tUrl.Host, ",")) + var clientOpts []txnkv.ClientOpt + if ks := query.Get("keyspace"); ks != "" { + logger.Infof("Using TiKV API V2 with keyspace: %s", ks) + clientOpts = append(clientOpts, + txnkv.WithKeyspace(ks), + txnkv.WithAPIVersion(kvrpcpb.APIVersion_V2), + ) + } + + client, err := txnkv.NewClient(strings.Split(tUrl.Host, ","), clientOpts...) if err != nil { return nil, err } From 3ff71e097c9efa02a6905c2bb6858c805089fdfc Mon Sep 17 00:00:00 2001 From: IANTHEREAL Date: Wed, 15 Apr 2026 11:34:22 +0000 Subject: [PATCH 2/4] feat: add AdvanceNextChunk to Meta interface Used after TiKV BR clone to prevent slice ID collision between original and cloned volumes. offset=0 reads the current counter, offset>0 atomically increments and returns the new value. Co-Authored-By: Claude Opus 4.6 (1M context) --- pkg/meta/base.go | 7 +++++++ pkg/meta/interface.go | 4 ++++ 2 files changed, 11 insertions(+) diff --git a/pkg/meta/base.go b/pkg/meta/base.go index 808e2442180a..4b4a7d0315c2 100644 --- a/pkg/meta/base.go +++ b/pkg/meta/base.go @@ -2054,6 +2054,13 @@ func (m *baseMeta) Read(ctx Context, inode Ino, indx uint32, slices *[]Slice) (s return 0 } +func (m *baseMeta) AdvanceNextChunk(offset int64) (int64, error) { + if offset == 0 { + return m.en.getCounter("nextChunk") + } + return m.en.incrCounter("nextChunk", offset) +} + func (m *baseMeta) NewSlice(ctx Context, id *uint64) syscall.Errno { m.freeMu.Lock() defer m.freeMu.Unlock() diff --git a/pkg/meta/interface.go b/pkg/meta/interface.go index 1b67dffba039..858b127e3a7b 100644 --- a/pkg/meta/interface.go +++ b/pkg/meta/interface.go @@ -470,6 +470,10 @@ type Meta interface { Read(ctx Context, inode Ino, indx uint32, slices *[]Slice) syscall.Errno // NewSlice returns an id for new slice. NewSlice(ctx Context, id *uint64) syscall.Errno + // AdvanceNextChunk advances the nextChunk counter by offset and returns + // the new value. Used after TiKV BR clone to prevent slice ID collision + // between the original and cloned volumes. + AdvanceNextChunk(offset int64) (int64, error) // Write put a slice of data on top of the given chunk. Write(ctx Context, inode Ino, indx uint32, off uint32, slice Slice, mtime time.Time) syscall.Errno // InvalidateChunkCache invalidate chunk cache From 32df2f5d7b475fd4f5fe30eaf409b18b10e8f271 Mon Sep 17 00:00:00 2001 From: Ian Date: Thu, 21 May 2026 07:45:15 +0200 Subject: [PATCH 3/4] fix: close background workers and cached pages (#2) Co-authored-by: Ubuntu --- pkg/chunk/cached_store.go | 151 +++++++- pkg/chunk/cached_store_close_test.go | 501 +++++++++++++++++++++++++++ pkg/chunk/disk_cache.go | 195 ++++++++++- pkg/chunk/disk_cache_state.go | 9 +- pkg/chunk/mem_cache.go | 38 +- pkg/chunk/prefetch.go | 24 +- pkg/fs/fs.go | 71 +++- pkg/fs/fs_close_test.go | 87 +++++ pkg/meta/base.go | 55 ++- pkg/meta/openfile.go | 31 +- pkg/meta/openfile_close_test.go | 38 ++ pkg/meta/redis.go | 1 + pkg/meta/sql.go | 1 + pkg/meta/tkv.go | 1 + pkg/vfs/reader.go | 84 ++++- pkg/vfs/reader_close_test.go | 117 +++++++ pkg/vfs/writer.go | 214 +++++++++++- pkg/vfs/writer_cancel_test.go | 391 +++++++++++++++++++++ 18 files changed, 1925 insertions(+), 84 deletions(-) create mode 100644 pkg/chunk/cached_store_close_test.go create mode 100644 pkg/fs/fs_close_test.go create mode 100644 pkg/meta/openfile_close_test.go create mode 100644 pkg/vfs/reader_close_test.go create mode 100644 pkg/vfs/writer_cancel_test.go diff --git a/pkg/chunk/cached_store.go b/pkg/chunk/cached_store.go index 1cf5caec73ee..9eb12abbb472 100644 --- a/pkg/chunk/cached_store.go +++ b/pkg/chunk/cached_store.go @@ -397,6 +397,16 @@ func (store *cachedStore) upload(ctx context.Context, key string, block *Page, s return err } +func (store *cachedStore) addWorker() bool { + store.wgMu.Lock() + defer store.wgMu.Unlock() + if store.closed.Load() { + return false + } + store.wg.Add(1) + return true +} + func (s *wSlice) upload(indx int) { blen := s.blockSize(indx) key := s.key(indx) @@ -404,7 +414,15 @@ func (s *wSlice) upload(indx int) { s.pages[indx] = nil s.pendings++ + if !s.store.addWorker() { + for _, p := range pages { + freePage(p) + } + s.errors <- errors.New("chunk store is closed") + return + } go func() { + defer s.store.wg.Done() var block *Page var off int if len(pages) == 1 { @@ -442,9 +460,8 @@ func (s *wSlice) upload(indx int) { } else { s.errors <- nil if s.store.conf.UploadDelay == 0 && s.store.canUpload() { - select { - case s.store.currentUpload <- struct{}{}: - defer func() { <-s.store.currentUpload }() + if s.store.tryAcquireUploadSlot() { + defer s.store.releaseUploadSlot() if err = s.store.upload(ctx, key, block, nil); err == nil { s.store.bcache.uploaded(key, blen) if err := s.store.bcache.removeStage(key); err != nil { @@ -454,7 +471,6 @@ func (s *wSlice) upload(indx int) { s.store.addDelayedStaging(key, stagingPath, time.Now(), false) } return - default: } } block.Release() @@ -462,8 +478,12 @@ func (s *wSlice) upload(indx int) { return } } - s.store.currentUpload <- struct{}{} - defer func() { <-s.store.currentUpload }() + if !s.store.acquireUploadSlot() { + block.Release() + s.errors <- errors.New("chunk store is closed") + return + } + defer s.store.releaseUploadSlot() s.errors <- s.store.upload(ctx, key, block, s) }() } @@ -674,6 +694,11 @@ type cachedStore struct { pendingCh chan *pendingItem pendingKeys map[string]*pendingItem pendingMutex sync.Mutex + done chan struct{} + closeOnce sync.Once + wgMu sync.Mutex + wg sync.WaitGroup + closed atomic.Bool startHour int endHour int compressor compress.Compressor @@ -846,6 +871,7 @@ func NewCachedStore(storage object.ObjectStorage, config Config, reg prometheus. seekable: compressor.CompressBound(0) == 0, pendingCh: make(chan *pendingItem, 100*config.MaxUpload), pendingKeys: make(map[string]*pendingItem), + done: make(chan struct{}), group: NewController(), } if config.UploadLimit > 0 { @@ -871,15 +897,26 @@ func NewCachedStore(storage object.ObjectStorage, config Config, reg prometheus. } }) + store.wg.Add(1) go func() { + defer store.wg.Done() for { + select { + case <-store.done: + return + default: + } if store.bcache.isEmpty() { logger.Warn("cache store is empty, use memory cache") config.CacheSize = 100 << 20 config.CacheDir = "memory" store.bcache = newMemStore(&config, store.bcache.getMetrics()) } - time.Sleep(time.Second) + select { + case <-store.done: + return + case <-time.After(time.Second): + } } }() @@ -906,6 +943,7 @@ func NewCachedStore(storage object.ObjectStorage, config Config, reg prometheus. if store.conf.Writeback { for i := 0; i < store.conf.MaxUpload; i++ { + store.wg.Add(1) go store.uploader() } interval := time.Minute @@ -917,9 +955,15 @@ func NewCachedStore(storage object.ObjectStorage, config Config, reg prometheus. logger.Infof("delay uploading by %s", d) } } + store.wg.Add(1) go func() { + defer store.wg.Done() for { - time.Sleep(interval) + select { + case <-store.done: + return + case <-time.After(interval): + } store.scanDelayedStaging() } }() @@ -1026,10 +1070,10 @@ func parseObjOrigSize(key string) int { } func (store *cachedStore) uploadStagingFile(key string, stagingPath string) { - store.currentUpload <- struct{}{} - defer func() { - <-store.currentUpload - }() + if !store.acquireUploadSlot() { + return + } + defer store.releaseUploadSlot() store.pendingMutex.Lock() item, ok := store.pendingKeys[key] @@ -1090,7 +1134,14 @@ func (store *cachedStore) uploadStagingFile(key string, stagingPath string) { } func (store *cachedStore) addDelayedStaging(key, stagingPath string, added time.Time, force bool) bool { + if store.closed.Load() { + return false + } store.pendingMutex.Lock() + if store.closed.Load() { + store.pendingMutex.Unlock() + return false + } item := store.pendingKeys[key] if item == nil { item = &pendingItem{key, stagingPath, added, atomic.Bool{}} @@ -1102,6 +1153,8 @@ func (store *cachedStore) addDelayedStaging(key, stagingPath string, added time. select { case store.pendingCh <- item: return true + case <-store.done: + item.uploading.Store(false) default: item.uploading.Store(false) } @@ -1126,6 +1179,9 @@ func (store *cachedStore) isPendingValid(key string) bool { } func (store *cachedStore) scanDelayedStaging() { + if store.closed.Load() { + return + } if !store.canUpload() { return } @@ -1135,16 +1191,81 @@ func (store *cachedStore) scanDelayedStaging() { for _, item := range store.pendingKeys { store.pendingMutex.Unlock() if item.ts.Before(cutoff) && item.uploading.CompareAndSwap(false, true) { - store.pendingCh <- item + select { + case store.pendingCh <- item: + case <-store.done: + item.uploading.Store(false) + return + } } store.pendingMutex.Lock() } } func (store *cachedStore) uploader() { - for it := range store.pendingCh { - store.uploadStagingFile(it.key, it.fpath) + defer store.wg.Done() + for { + select { + case it := <-store.pendingCh: + store.uploadStagingFile(it.key, it.fpath) + case <-store.done: + return + } + } +} + +func (store *cachedStore) acquireUploadSlot() bool { + if store.closed.Load() { + return false } + select { + case store.currentUpload <- struct{}{}: + if store.closed.Load() { + store.releaseUploadSlot() + return false + } + return true + case <-store.done: + return false + } +} + +func (store *cachedStore) tryAcquireUploadSlot() bool { + if store.closed.Load() { + return false + } + select { + case store.currentUpload <- struct{}{}: + if store.closed.Load() { + store.releaseUploadSlot() + return false + } + return true + default: + return false + } +} + +func (store *cachedStore) releaseUploadSlot() { + <-store.currentUpload +} + +func (store *cachedStore) Close() error { + store.closeOnce.Do(func() { + store.wgMu.Lock() + store.closed.Store(true) + close(store.done) + store.wgMu.Unlock() + if store.fetcher != nil { + store.fetcher.Close() + } + store.wg.Wait() + store.pendingMutex.Lock() + store.pendingKeys = make(map[string]*pendingItem) + store.pendingMutex.Unlock() + store.bcache.close() + }) + return nil } func (store *cachedStore) canUpload() bool { diff --git a/pkg/chunk/cached_store_close_test.go b/pkg/chunk/cached_store_close_test.go new file mode 100644 index 000000000000..498a5768be88 --- /dev/null +++ b/pkg/chunk/cached_store_close_test.go @@ -0,0 +1,501 @@ +/* + * JuiceFS, Copyright 2020 Juicedata, Inc. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package chunk + +import ( + "bytes" + "context" + "io" + "os" + "path/filepath" + "sync/atomic" + "testing" + "time" + + "github.com/davies/groupcache/consistenthash" + "github.com/juicedata/juicefs/pkg/compress" + "github.com/juicedata/juicefs/pkg/object" + "github.com/twmb/murmur3" +) + +func TestCachedStoreCloseReleasesMemoryCachePages(t *testing.T) { + blob, _ := object.CreateStorage("mem", "", "", "", "") + store := NewCachedStore(blob, Config{ + CacheDir: "memory", + CacheSize: 1 << 20, + BlockSize: 1 << 20, + Compress: "none", + MaxUpload: 1, + MaxDownload: 1, + BufferSize: 32 << 20, + }, nil).(*cachedStore) + + p := NewOffPage(4096) + store.bcache.cache("chunks/0/0/1_0_4096", p, true, false) + p.Release() + if used := store.UsedMemory(); used == 0 { + t.Fatal("memory cache did not retain the test page") + } + + if err := store.Close(); err != nil { + t.Fatalf("close: %s", err) + } + if used := store.UsedMemory(); used != 0 { + t.Fatalf("used memory after close = %d, want 0", used) + } +} + +func TestCachedStoreTryAcquireUploadSlotAfterClose(t *testing.T) { + store := &cachedStore{ + currentUpload: make(chan struct{}, 1), + done: make(chan struct{}), + } + store.closed.Store(true) + close(store.done) + + for i := 0; i < 100; i++ { + if store.tryAcquireUploadSlot() { + t.Fatal("tryAcquireUploadSlot acquired a slot after close") + } + } + if got := len(store.currentUpload); got != 0 { + t.Fatalf("upload slots after close = %d, want 0", got) + } +} + +type blockingPutStorage struct { + object.ObjectStorage + entered chan struct{} + release chan struct{} + deleteCount atomic.Int64 +} + +func (s *blockingPutStorage) Put(ctx context.Context, key string, in io.Reader, getters ...object.AttrGetter) error { + select { + case <-s.entered: + default: + close(s.entered) + } + select { + case <-s.release: + return s.ObjectStorage.Put(ctx, key, in, getters...) + case <-ctx.Done(): + return ctx.Err() + } +} + +func (s *blockingPutStorage) Delete(ctx context.Context, key string, getters ...object.AttrGetter) error { + s.deleteCount.Add(1) + return s.ObjectStorage.Delete(ctx, key, getters...) +} + +type closeTestCacheManager struct { + uploadedCount atomic.Int64 + removeStageCount atomic.Int64 +} + +func (m *closeTestCacheManager) cache(key string, p *Page, force, dropCache bool) {} + +func (m *closeTestCacheManager) remove(key string, staging bool) {} + +func (m *closeTestCacheManager) load(key string) (ReadCloser, error) { + return nil, errNotCached +} + +func (m *closeTestCacheManager) exist(key string) (string, bool) { + return "", false +} + +func (m *closeTestCacheManager) uploaded(key string, size int) { + m.uploadedCount.Add(1) +} + +func (m *closeTestCacheManager) stage(key string, data []byte, tierID uint8) (string, error) { + return "", nil +} + +func (m *closeTestCacheManager) removeStage(key string) error { + m.removeStageCount.Add(1) + return nil +} + +func (m *closeTestCacheManager) stats() (int64, int64) { + return 0, 0 +} + +func (m *closeTestCacheManager) usedMemory() int64 { + return 0 +} + +func (m *closeTestCacheManager) isEmpty() bool { + return false +} + +func (m *closeTestCacheManager) getMetrics() *cacheManagerMetrics { + return nil +} + +func (m *closeTestCacheManager) close() {} + +func newCloseTestCacheStore(id string) *cacheStore { + s := &cacheStore{ + id: id, + dir: id, + pending: make(chan pendingFile, 1), + pages: make(map[string]*Page), + done: make(chan struct{}), + opTs: make(map[time.Duration]func() error), + } + s.state = newDCState(dcUnchanged, s) + return s +} + +func newCloseTestWritableCacheStore(t *testing.T, id string, pending int) *cacheStore { + t.Helper() + keys, err := NewKeyIndex(&Config{CacheEviction: Eviction2Random}) + if err != nil { + t.Fatalf("new key index: %s", err) + } + s := newCloseTestCacheStore(id) + s.capacity = 1 << 20 + s.pending = make(chan pendingFile, pending) + s.keys = keys + s.m = newCacheManagerMetrics(nil) + return s +} + +func newCloseTestDiskCacheManager(stores ...*cacheStore) *cacheManager { + m := &cacheManager{ + consistentMap: consistenthash.New(100, murmur3.Sum32), + storeMap: make(map[string]*cacheStore, len(stores)), + stores: make([]*cacheStore, len(stores)), + allStores: make([]*cacheStore, 0, len(stores)), + done: make(chan struct{}), + } + for i, s := range stores { + m.consistentMap.Add(s.id) + m.storeMap[s.id] = s + m.stores[i] = s + m.allStores = append(m.allStores, s) + } + return m +} + +func assertClosed(t *testing.T, ch <-chan struct{}, name string) { + t.Helper() + select { + case <-ch: + default: + t.Fatalf("%s is still open", name) + } +} + +func waitClosed(t *testing.T, ch <-chan struct{}, name string) { + t.Helper() + select { + case <-ch: + case <-time.After(time.Second): + t.Fatalf("%s is still open", name) + } +} + +func waitForCachePage(t *testing.T, cache *cacheStore, key string) { + t.Helper() + deadline := time.After(time.Second) + for { + cache.Lock() + _, ok := cache.pages[key] + cache.Unlock() + if ok { + return + } + select { + case <-deadline: + t.Fatalf("cache page %s was not inserted", key) + case <-time.After(time.Millisecond): + } + } +} + +func TestCacheStoreForceCacheDoesNotDoubleReleasePageRemovedDuringClose(t *testing.T) { + cache := newCloseTestWritableCacheStore(t, "cache", 0) + key := "chunks/0/0/1_0_5" + page := NewPage([]byte("hello")) + + done := make(chan struct{}) + go func() { + cache.cache(key, page, true, false) + close(done) + }() + waitForCachePage(t, cache, key) + + cache.Lock() + cached := cache.pages[key] + if cached != page { + cache.Unlock() + t.Fatalf("cached page = %p, want %p", cached, page) + } + close(cache.done) + delete(cache.pages, key) + atomic.AddInt64(&cache.totalPages, -int64(cap(page.Data))) + page.Release() + cache.Unlock() + + select { + case <-done: + case <-time.After(time.Second): + t.Fatal("force cache did not return after cache close") + } + if refs := atomic.LoadInt32(&page.refs); refs != 1 { + t.Fatalf("page refs after close race = %d, want 1", refs) + } + if page.Data == nil { + t.Fatal("page data was released by cache after ownership was removed") + } + page.Release() +} + +func TestCacheManagerRemoveStoreClosesRemovedStore(t *testing.T) { + removed := newCloseTestCacheStore("removed") + active := newCloseTestCacheStore("active") + m := newCloseTestDiskCacheManager(removed, active) + + m.removeStore(removed.id) + assertClosed(t, removed.done, "removed store done") + if got := m.length(); got != 1 { + t.Fatalf("manager length after remove = %d, want 1", got) + } + + m.close() + assertClosed(t, active.done, "active store done") + assertClosed(t, removed.done, "removed store done after manager close") +} + +func TestCacheManagerCloseWaitsForRemovedStoreClose(t *testing.T) { + removed := newCloseTestCacheStore("removed") + m := newCloseTestDiskCacheManager(removed) + + releaseStoreClose := make(chan struct{}) + removed.wg.Add(1) + go func() { + <-removed.done + <-releaseStoreClose + removed.wg.Done() + }() + + removeDone := make(chan struct{}) + go func() { + m.removeStore(removed.id) + close(removeDone) + }() + waitClosed(t, removed.done, "removed store done") + + managerCloseDone := make(chan struct{}) + go func() { + m.close() + close(managerCloseDone) + }() + + select { + case <-managerCloseDone: + t.Fatal("manager close returned before removed store close completed") + case <-time.After(100 * time.Millisecond): + } + + close(releaseStoreClose) + select { + case <-managerCloseDone: + case <-time.After(time.Second): + t.Fatal("manager close did not return after removed store close completed") + } + select { + case <-removeDone: + case <-time.After(time.Second): + t.Fatal("removeStore did not return after removed store close completed") + } +} + +func TestCachedStoreCloseWaitsForSliceUpload(t *testing.T) { + mem, _ := object.CreateStorage("mem", "", "", "", "") + blob := &blockingPutStorage{ + ObjectStorage: mem, + entered: make(chan struct{}), + release: make(chan struct{}), + } + store := NewCachedStore(blob, Config{ + CacheDir: "memory", + CacheSize: 1 << 20, + BlockSize: 4, + Compress: "none", + MaxUpload: 1, + MaxDownload: 1, + BufferSize: 32 << 20, + PutTimeout: time.Second, + }, nil).(*cachedStore) + + s := sliceForWrite(1, store, 0) + p := NewOffPage(4) + copy(p.Data, []byte("test")) + s.pages[0] = []*Page{p} + s.length = 4 + s.upload(0) + + select { + case <-blob.entered: + case <-time.After(time.Second): + t.Fatal("upload did not reach object storage") + } + + done := make(chan error, 1) + go func() { + done <- store.Close() + }() + + select { + case err := <-done: + t.Fatalf("Close returned before upload finished: %v", err) + case <-time.After(50 * time.Millisecond): + } + + close(blob.release) + select { + case err := <-done: + if err != nil { + t.Fatalf("close: %s", err) + } + case <-time.After(time.Second): + t.Fatal("Close did not return after upload finished") + } +} + +func TestCachedStoreCloseDoesNotDeleteSuccessfulStagingUpload(t *testing.T) { + mem, _ := object.CreateStorage("mem", "", "", "", "") + blob := &blockingPutStorage{ + ObjectStorage: mem, + entered: make(chan struct{}), + release: make(chan struct{}), + } + bcache := &closeTestCacheManager{} + key := "chunks/0/0/321_0_4" + stagingPath := filepath.Join(t.TempDir(), key) + if err := os.MkdirAll(filepath.Dir(stagingPath), 0700); err != nil { + t.Fatalf("mkdir staging dir: %s", err) + } + if err := os.WriteFile(stagingPath, []byte("good"), 0600); err != nil { + t.Fatalf("write staging file: %s", err) + } + item := &pendingItem{key: key, fpath: stagingPath, ts: time.Now()} + item.uploading.Store(true) + store := &cachedStore{ + storage: blob, + bcache: bcache, + conf: Config{BlockSize: 4, CacheChecksum: CsNone, PutTimeout: time.Second}, + currentUpload: make(chan struct{}, 1), + pendingKeys: map[string]*pendingItem{ + key: item, + }, + done: make(chan struct{}), + compressor: compress.NewCompressor("none"), + } + store.initMetrics() + store.wg.Add(1) + go func() { + defer store.wg.Done() + store.uploadStagingFile(key, stagingPath) + }() + + select { + case <-blob.entered: + case <-time.After(time.Second): + t.Fatal("staging upload did not reach object storage") + } + + done := make(chan error, 1) + go func() { + done <- store.Close() + }() + + select { + case err := <-done: + t.Fatalf("Close returned before staging upload finished: %v", err) + case <-time.After(50 * time.Millisecond): + } + + close(blob.release) + select { + case err := <-done: + if err != nil { + t.Fatalf("close: %s", err) + } + case <-time.After(time.Second): + t.Fatal("Close did not return after staging upload finished") + } + + if got := blob.deleteCount.Load(); got != 0 { + t.Fatalf("object delete count = %d, want 0", got) + } + if got := bcache.uploadedCount.Load(); got != 1 { + t.Fatalf("uploaded count = %d, want 1", got) + } + if got := bcache.removeStageCount.Load(); got != 1 { + t.Fatalf("removeStage count = %d, want 1", got) + } + in, err := mem.Get(context.Background(), key, 0, -1) + if err != nil { + t.Fatalf("uploaded object missing after close: %s", err) + } + defer in.Close() + data, err := io.ReadAll(in) + if err != nil { + t.Fatalf("read uploaded object: %s", err) + } + if string(data) != "good" { + t.Fatalf("uploaded object = %q, want %q", data, "good") + } +} + +func TestCachedStoreUploadAfterCloseReportsError(t *testing.T) { + mem, _ := object.CreateStorage("mem", "", "", "", "") + store := NewCachedStore(mem, Config{ + CacheDir: "memory", + CacheSize: 1 << 20, + BlockSize: 4, + Compress: "none", + MaxUpload: 1, + MaxDownload: 1, + BufferSize: 32 << 20, + }, nil).(*cachedStore) + if err := store.Close(); err != nil { + t.Fatalf("close: %s", err) + } + + s := sliceForWrite(1, store, 0) + p := NewOffPage(4) + copy(p.Data, bytes.Repeat([]byte{'x'}, 4)) + s.pages[0] = []*Page{p} + s.length = 4 + s.upload(0) + + select { + case err := <-s.errors: + if err == nil { + t.Fatal("upload after close returned nil error") + } + case <-time.After(time.Second): + t.Fatal("upload after close did not report an error") + } +} diff --git a/pkg/chunk/disk_cache.go b/pkg/chunk/disk_cache.go index 73aa029bc61d..8fd59d09c56f 100644 --- a/pkg/chunk/disk_cache.go +++ b/pkg/chunk/disk_cache.go @@ -86,6 +86,10 @@ type cacheStore struct { pending chan pendingFile pages map[string]*Page m *cacheManagerMetrics + done chan struct{} + closeOnce sync.Once + wg sync.WaitGroup + closed bool used int64 keys KeyIndex @@ -134,6 +138,7 @@ func newCacheStore(m *cacheManagerMetrics, dir string, cacheSize, maxItems int64 keys: keyIndex, pending: make(chan pendingFile, pendingPages), pages: make(map[string]*Page), + done: make(chan struct{}), uploader: uploader, opTs: make(map[time.Duration]func() error), stagedBlockCooldown: config.CacheExpire / 2, @@ -156,14 +161,21 @@ func newCacheStore(m *cacheManagerMetrics, dir string, cacheSize, maxItems int64 c.setLimitByFreeRatio(usage, c.freeRatio) c.createLockFile() + c.wg.Add(1) go c.checkLockFile() + c.wg.Add(1) go c.flush() + c.wg.Add(1) go c.checkFreeSpace() if c.cacheExpire > 0 { + c.wg.Add(1) go c.cleanupExpire() } + c.wg.Add(1) go c.refreshCacheKeys() + c.wg.Add(1) go c.scanStaging() + c.wg.Add(1) go c.checkTimeout() return c } @@ -224,9 +236,14 @@ func (cache *cacheStore) createLockFile() { } func (cache *cacheStore) checkLockFile() { + defer cache.wg.Done() lockfile := cache.lockFilePath() for cache.available() { - time.Sleep(time.Second * 10) + select { + case <-cache.done: + return + case <-time.After(time.Second * 10): + } if err := cache.statFile(lockfile); err != nil && os.IsNotExist(err) { logger.Infof("lockfile %s is lost, cache device maybe broken", lockfile) if inRootVolume(cache.dir) && cache.freeRatio < 0.2 { @@ -238,6 +255,11 @@ func (cache *cacheStore) checkLockFile() { } func (c *cacheStore) available() bool { + select { + case <-c.done: + return false + default: + } return c.state.state() != dcDown } @@ -284,6 +306,7 @@ func getFunctionName(f interface{}) string { } func (c *cacheStore) checkTimeout() { + defer c.wg.Done() for c.available() { now := utils.Clock() cutOff := now - maxIODur @@ -296,7 +319,11 @@ func (c *cacheStore) checkTimeout() { } } c.opMu.Unlock() - time.Sleep(time.Second) + select { + case <-c.done: + return + case <-time.After(time.Second): + } } } @@ -343,6 +370,7 @@ func (cache *cacheStore) stats() (int64, int64) { } func (cache *cacheStore) checkFreeSpace() { + defer cache.wg.Done() for cache.available() { usage := cache.curFreeRatio() cache.stageFull = usage.br < cache.freeRatio/2 || (usage.inodeCap > 0 && usage.fr < cache.freeRatio/2) @@ -358,12 +386,17 @@ func (cache *cacheStore) checkFreeSpace() { if cache.rawFull { cache.uploadStaging() } - time.Sleep(time.Second) + select { + case <-cache.done: + return + case <-time.After(time.Second): + } } logger.Infof("stop checkFreeSpace at %s", cache.dir) } func (cache *cacheStore) cleanupExpire() { + defer cache.wg.Done() var todel []cacheKey var interval = time.Minute if cache.cacheExpire < time.Minute { @@ -403,18 +436,28 @@ func (cache *cacheStore) cleanupExpire() { _ = cache.removeFile(cache.cachePath(cache.getPathFromKey(k))) } todel = todel[:0] - time.Sleep(interval / 1000 * time.Duration((cnt+1-deleted)*1000/(cnt+1))) + sleep := interval / 1000 * time.Duration((cnt+1-deleted)*1000/(cnt+1)) + select { + case <-cache.done: + return + case <-time.After(sleep): + } } } func (cache *cacheStore) refreshCacheKeys() { + defer cache.wg.Done() if cache.scanInterval < 0 { return } cache.scanCached() if cache.scanInterval > 0 { for { - time.Sleep(cache.scanInterval) + select { + case <-cache.done: + return + case <-time.After(cache.scanInterval): + } cache.scanCached() } } @@ -444,6 +487,9 @@ func (cache *cacheStore) cache(key string, p *Page, force, dropCache bool) { } cache.Lock() defer cache.Unlock() + if cache.closed { + return + } if _, ok := cache.pages[key]; ok { return } @@ -458,9 +504,21 @@ func (cache *cacheStore) cache(key string, p *Page, force, dropCache bool) { case cache.pending <- pendingFile{key, p, dropCache}: default: if force { + sent := false cache.Unlock() - cache.pending <- pendingFile{key, p, dropCache} + select { + case cache.pending <- pendingFile{key, p, dropCache}: + sent = true + case <-cache.done: + } cache.Lock() + if !sent { + if cached, ok := cache.pages[key]; ok && cached == p { + delete(cache.pages, key) + atomic.AddInt64(&cache.totalPages, -int64(cap(p.Data))) + p.Release() + } + } } else { // does not have enough bandwidth to write it into disk, discard it logger.Debugf("Caching queue is full (%s), drop %s (%d bytes)", cache.dir, key, len(p.Data)) @@ -472,6 +530,44 @@ func (cache *cacheStore) cache(key string, p *Page, force, dropCache bool) { } } +func (cache *cacheStore) releasePendingPages() { + for { + select { + case w := <-cache.pending: + cache.Lock() + if _, ok := cache.pages[w.key]; ok { + delete(cache.pages, w.key) + atomic.AddInt64(&cache.totalPages, -int64(cap(w.page.Data))) + } + cache.Unlock() + w.page.Release() + default: + return + } + } +} + +func (cache *cacheStore) close() { + cache.closeOnce.Do(func() { + close(cache.done) + cache.stateLock.Lock() + cache.state.stop() + cache.stateLock.Unlock() + cache.Lock() + cache.closed = true + cache.Unlock() + cache.wg.Wait() + cache.releasePendingPages() + cache.Lock() + for key, page := range cache.pages { + delete(cache.pages, key) + atomic.AddInt64(&cache.totalPages, -int64(cap(page.Data))) + page.Release() + } + cache.Unlock() + }) +} + type DiskFreeRatio struct { br float32 fr float32 @@ -719,8 +815,15 @@ func (cache *cacheStore) stagePath(key string) string { // flush cached block into disk func (cache *cacheStore) flush() { + defer cache.wg.Done() for { - w := <-cache.pending + var w pendingFile + select { + case w = <-cache.pending: + case <-cache.done: + cache.releasePendingPages() + return + } path := cache.cachePath(w.key) if cache.enabled() && cache.flushPage(path, w.page.Data, w.dropCache, 0) == nil { cache.add(w.key, int32(len(w.page.Data)), uint32(time.Now().Unix())) @@ -975,9 +1078,15 @@ func (cache *cacheStore) scanCached() { var pathReg, _ = regexp.Compile(`^chunks/((\d+)|([0-9a-fA-F]{2}))/\d+/\d+_\d+_\d+$`) func (cache *cacheStore) scanStaging() { + defer cache.wg.Done() if cache.uploader == nil { return } + select { + case <-cache.done: + return + default: + } var start = time.Now() var oneMinAgo = start.Add(-time.Minute) @@ -1033,7 +1142,11 @@ type cacheManager struct { consistentMap *consistenthash.Map storeMap map[string]*cacheStore stores []*cacheStore + allStores []*cacheStore metrics *cacheManagerMetrics + done chan struct{} + closeOnce sync.Once + wg sync.WaitGroup } func legacyKeyHash(s string) uint32 { @@ -1093,6 +1206,7 @@ type CacheManager interface { usedMemory() int64 isEmpty() bool getMetrics() *cacheManagerMetrics + close() } func newCacheManager(config *Config, reg prometheus.Registerer, uploader func(key, path string, force bool) bool) CacheManager { @@ -1126,7 +1240,9 @@ func newCacheManager(config *Config, reg prometheus.Registerer, uploader func(ke consistentMap: consistenthash.New(100, murmur3.Sum32), storeMap: make(map[string]*cacheStore, len(dirs)), stores: make([]*cacheStore, len(dirs)), + allStores: make([]*cacheStore, 0, len(dirs)), metrics: metrics, + done: make(chan struct{}), } // 20% of buffer could be used for pending pages @@ -1134,9 +1250,11 @@ func newCacheManager(config *Config, reg prometheus.Registerer, uploader func(ke for i, d := range dirs { store := newCacheStore(metrics, strings.TrimSpace(d)+string(filepath.Separator), dirCacheSize, dirCacheItems, pendingPages, config, uploader) m.stores[i] = store + m.allStores = append(m.allStores, store) m.storeMap[store.id] = store m.consistentMap.Add(store.id) } + m.wg.Add(1) go m.cleanup() return m } @@ -1146,6 +1264,7 @@ func (m *cacheManager) getMetrics() *cacheManagerMetrics { } func (m *cacheManager) cleanup() { + defer m.wg.Done() for !m.isEmpty() { var ids []string m.Lock() @@ -1158,10 +1277,44 @@ func (m *cacheManager) cleanup() { for _, id := range ids { m.removeStore(id) } - time.Sleep(time.Second) + select { + case <-m.done: + return + case <-time.After(time.Second): + } } } +func (m *cacheManager) close() { + m.closeOnce.Do(func() { + close(m.done) + for _, s := range m.storesForClose() { + s.close() + } + m.wg.Wait() + }) +} + +func (m *cacheManager) storesForClose() []*cacheStore { + m.Lock() + defer m.Unlock() + stores := make([]*cacheStore, len(m.allStores)) + copy(stores, m.allStores) + return stores +} + +func (m *cacheManager) activeStores() []*cacheStore { + m.Lock() + defer m.Unlock() + stores := make([]*cacheStore, 0, len(m.stores)) + for _, s := range m.stores { + if s != nil { + stores = append(stores, s) + } + } + return stores +} + func (m *cacheManager) isEmpty() bool { return m.length() == 0 } @@ -1173,11 +1326,13 @@ func (m *cacheManager) length() int { } func (m *cacheManager) removeStore(id string) { + var removed *cacheStore m.Lock() m.consistentMap.Remove(id) var dir string if s := m.storeMap[id]; s != nil { dir = s.dir + removed = s } delete(m.storeMap, id) for i, c := range m.stores { @@ -1186,6 +1341,9 @@ func (m *cacheManager) removeStore(id string) { } } m.Unlock() + if removed != nil { + removed.close() + } logger.Errorf("cache dir `%s`(%s) is unavailable, removed", dir, id) } @@ -1212,27 +1370,28 @@ func (m *cacheManager) removeStage(key string) error { // Deprecated: use getStore instead func (m *cacheManager) getStoreLegacy(key string) *cacheStore { + m.Lock() + defer m.Unlock() + if len(m.stores) == 0 { + return nil + } return m.stores[legacyKeyHash(key)%uint32(len(m.stores))] } func (m *cacheManager) usedMemory() int64 { var used int64 - for _, s := range m.stores { - if s != nil { - used += s.usedMemory() - } + for _, s := range m.activeStores() { + used += s.usedMemory() } return used } func (m *cacheManager) stats() (int64, int64) { var cnt, used int64 - for _, s := range m.stores { - if s != nil { - c, u := s.stats() - cnt += c - used += u - } + for _, s := range m.activeStores() { + c, u := s.stats() + cnt += c + used += u } return cnt, used } diff --git a/pkg/chunk/disk_cache_state.go b/pkg/chunk/disk_cache_state.go index 20f614d835af..e111d629fe5a 100644 --- a/pkg/chunk/disk_cache_state.go +++ b/pkg/chunk/disk_cache_state.go @@ -82,8 +82,9 @@ type dcState interface { } type baseDC struct { - cache *cacheStore - stopCh chan struct{} + cache *cacheStore + stopCh chan struct{} + stopped atomic.Bool } func newDCState(state int, cs *cacheStore) dcState { @@ -109,7 +110,9 @@ func (dc *baseDC) init(cs *cacheStore) { } func (dc *baseDC) stop() { - close(dc.stopCh) + if dc.stopped.CompareAndSwap(false, true) { + close(dc.stopCh) + } } func (dc *baseDC) onIOErr() {} func (dc *baseDC) onIOSucc() {} diff --git a/pkg/chunk/mem_cache.go b/pkg/chunk/mem_cache.go index 267b22db4dbd..7df999e5ec80 100644 --- a/pkg/chunk/mem_cache.go +++ b/pkg/chunk/mem_cache.go @@ -38,6 +38,10 @@ type memcache struct { pages map[string]memItem eviction string cacheExpire time.Duration + done chan struct{} + closeOnce sync.Once + wg sync.WaitGroup + closed bool metrics *cacheManagerMetrics } @@ -49,6 +53,7 @@ func newMemStore(config *Config, metrics *cacheManagerMetrics) *memcache { pages: make(map[string]memItem), eviction: config.CacheEviction, cacheExpire: config.CacheExpire, + done: make(chan struct{}), metrics: metrics, } runtime.SetFinalizer(c, func(c *memcache) { @@ -58,6 +63,7 @@ func newMemStore(config *Config, metrics *cacheManagerMetrics) *memcache { c.pages = nil }) if c.cacheExpire > 0 { + c.wg.Add(1) go c.cleanupExpire() } return c @@ -85,6 +91,9 @@ func (c *memcache) cache(key string, p *Page, force, dropCache bool) { } c.Lock() defer c.Unlock() + if c.closed { + return + } if c.full() && c.eviction == EvictionNone { logger.Debugf("Caching is full, drop %s (%d bytes)", key, len(p.Data)) c.metrics.cacheDrops.Add(1) @@ -177,6 +186,7 @@ func (c *memcache) full() bool { } func (c *memcache) cleanupExpire() { + defer c.wg.Done() var interval = time.Minute if c.cacheExpire < time.Minute { interval = c.cacheExpire @@ -202,7 +212,12 @@ func (c *memcache) cleanupExpire() { if deleted > 0 { logger.Debugf("Expired cache blocks: %d blocks (%s), remaining: %d blocks (%s)", deleted, humanize.IBytes(uint64(freed)), len(c.pages), humanize.IBytes(uint64(c.used))) } - time.Sleep(interval / 1000 * time.Duration((cnt+1-deleted)*1000/(cnt+1))) + sleep := interval / 1000 * time.Duration((cnt+1-deleted)*1000/(cnt+1)) + select { + case <-c.done: + return + case <-time.After(sleep): + } } } @@ -212,3 +227,24 @@ func (c *memcache) stage(key string, data []byte, tierID uint8) (string, error) func (c *memcache) uploaded(key string, size int) {} func (c *memcache) isEmpty() bool { return false } func (c *memcache) getMetrics() *cacheManagerMetrics { return c.metrics } + +func (c *memcache) close() { + c.closeOnce.Do(func() { + close(c.done) + c.wg.Wait() + c.Lock() + c.closed = true + keys := make([]string, 0, len(c.pages)) + for k := range c.pages { + keys = append(keys, k) + } + for _, k := range keys { + item := c.pages[k] + c.delete(k, item.page) + } + c.pages = make(map[string]memItem) + c.used = 0 + c.Unlock() + runtime.SetFinalizer(c, nil) + }) +} diff --git a/pkg/chunk/prefetch.go b/pkg/chunk/prefetch.go index 05ab9ba8b2e5..74287e1fd1c5 100644 --- a/pkg/chunk/prefetch.go +++ b/pkg/chunk/prefetch.go @@ -22,9 +22,12 @@ import ( type prefetcher struct { sync.Mutex - pending chan string - busy map[string]bool - op func(key string) + pending chan string + busy map[string]bool + op func(key string) + closeOnce sync.Once + wg sync.WaitGroup + closed bool } func newPrefetcher(parallel int, fetch func(string)) *prefetcher { @@ -34,12 +37,14 @@ func newPrefetcher(parallel int, fetch func(string)) *prefetcher { op: fetch, } for i := 0; i < parallel; i++ { + p.wg.Add(1) go p.do() } return p } func (p *prefetcher) do() { + defer p.wg.Done() for key := range p.pending { p.op(key) @@ -52,6 +57,9 @@ func (p *prefetcher) do() { func (p *prefetcher) fetch(key string) { p.Lock() defer p.Unlock() + if p.closed { + return + } if _, ok := p.busy[key]; ok { return } @@ -61,3 +69,13 @@ func (p *prefetcher) fetch(key string) { default: } } + +func (p *prefetcher) Close() { + p.closeOnce.Do(func() { + p.Lock() + p.closed = true + close(p.pending) + p.Unlock() + }) + p.wg.Wait() +} diff --git a/pkg/fs/fs.go b/pkg/fs/fs.go index 469ddebb137d..23c48ac04b09 100644 --- a/pkg/fs/fs.go +++ b/pkg/fs/fs.go @@ -147,8 +147,12 @@ type FileSystem struct { cacheM sync.Mutex entries map[Ino]map[string]*entryCache attrs map[Ino]*attrCache + done chan struct{} + closeOnce sync.Once + wg sync.WaitGroup checkAccessFile time.Duration rotateAccessLog int64 + logMu sync.Mutex logBuffer chan string readSizeHistogram prometheus.Histogram @@ -188,6 +192,7 @@ func NewFileSystem(conf *vfs.Config, m meta.Meta, d chunk.ChunkStore, registry * writer: vfs.NewDataWriter(conf, m, d, reader), entries: make(map[meta.Ino]map[string]*entryCache), attrs: make(map[meta.Ino]*attrCache), + done: make(chan struct{}), checkAccessFile: time.Minute, rotateAccessLog: 300 << 20, // 300 MiB @@ -221,6 +226,7 @@ func NewFileSystem(conf *vfs.Config, m meta.Meta, d chunk.ChunkStore, registry * } } + fs.wg.Add(1) go fs.cleanupCache() if conf.AccessLog != "" { f, err := os.OpenFile(conf.AccessLog, os.O_CREATE|os.O_WRONLY|os.O_APPEND, 0666) @@ -229,6 +235,7 @@ func NewFileSystem(conf *vfs.Config, m meta.Meta, d chunk.ChunkStore, registry * } else { _ = os.Chmod(conf.AccessLog, 0666) fs.logBuffer = make(chan string, 1024) + fs.wg.Add(1) go fs.flushLog(f, fs.logBuffer, conf.AccessLog) } } @@ -245,7 +252,13 @@ func (fs *FileSystem) InitMetrics(reg prometheus.Registerer) { } func (fs *FileSystem) cleanupCache() { + defer fs.wg.Done() for { + select { + case <-fs.done: + return + default: + } fs.cacheM.Lock() now := time.Now() var cnt int @@ -275,7 +288,11 @@ func (fs *FileSystem) cleanupCache() { } } fs.cacheM.Unlock() - time.Sleep(time.Second) + select { + case <-fs.done: + return + case <-time.After(time.Second): + } } } @@ -300,14 +317,16 @@ func (fs *FileSystem) InvalidateAttr(ino Ino) { func (fs *FileSystem) log(ctx LogContext, format string, args ...interface{}) { used := ctx.Duration() fs.opsDurationsHistogram.Observe(used.Seconds()) - if fs.logBuffer == nil { - return - } now := utils.Now() cmd := fmt.Sprintf(format, args...) ts := now.Format("2006.01.02 15:04:05.000000") cmd += fmt.Sprintf(" <%.6f>", used.Seconds()) line := fmt.Sprintf("%s [uid:%d,gid:%d,pid:%d] %s\n", ts, ctx.Uid(), ctx.Gid(), ctx.Pid(), cmd) + fs.logMu.Lock() + defer fs.logMu.Unlock() + if fs.logBuffer == nil { + return + } select { case fs.logBuffer <- line: default: @@ -316,15 +335,23 @@ func (fs *FileSystem) log(ctx LogContext, format string, args ...interface{}) { } func (fs *FileSystem) flushLog(f *os.File, logBuffer chan string, path string) { + defer fs.wg.Done() + defer func() { _ = f.Close() }() buf := make([]byte, 0, 128<<10) var lastcheck = time.Now() for { - line := <-logBuffer + line, ok := <-logBuffer + if !ok { + return + } buf = append(buf[:0], []byte(line)...) LOOP: for len(buf) < (128 << 10) { select { - case line = <-logBuffer: + case line, ok = <-logBuffer: + if !ok { + break LOOP + } buf = append(buf, []byte(line)...) default: break LOOP @@ -1076,22 +1103,40 @@ func (fs *FileSystem) Create(ctx meta.Context, p string, mode uint16, umask uint } func (fs *FileSystem) Flush() error { + fs.logMu.Lock() buffer := fs.logBuffer if buffer != nil { buffer <- "" // flush } + fs.logMu.Unlock() fs.Meta().FlushSession() return nil } func (fs *FileSystem) Close() error { - _ = fs.Flush() - buffer := fs.logBuffer - if buffer != nil { - fs.logBuffer = nil - close(buffer) - } - return nil + var err error + fs.closeOnce.Do(func() { + err = errors.Join(err, fs.Flush()) + close(fs.done) + fs.logMu.Lock() + buffer := fs.logBuffer + if buffer != nil { + fs.logBuffer = nil + close(buffer) + } + fs.logMu.Unlock() + fs.wg.Wait() + if c, ok := fs.writer.(io.Closer); ok { + err = errors.Join(err, c.Close()) + } + if c, ok := fs.reader.(io.Closer); ok { + err = errors.Join(err, c.Close()) + } + if c, ok := fs.store.(io.Closer); ok { + err = errors.Join(err, c.Close()) + } + }) + return err } func (fs *FileSystem) Clone(ctx meta.Context, src, dst string, preserve bool) (err syscall.Errno) { diff --git a/pkg/fs/fs_close_test.go b/pkg/fs/fs_close_test.go new file mode 100644 index 000000000000..f782a2fb7aa3 --- /dev/null +++ b/pkg/fs/fs_close_test.go @@ -0,0 +1,87 @@ +/* + * JuiceFS, Copyright 2020 Juicedata, Inc. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package fs + +import ( + "os" + "path/filepath" + "syscall" + "testing" + "time" + + "github.com/juicedata/juicefs/pkg/meta" +) + +func TestFileSystemCloseCascadesToWriter(t *testing.T) { + fs := createTestFS(t) + defer fs.m.Shutdown() //nolint:errcheck + + ctx := meta.NewContext(1, 0, []uint32{0}) + f, errno := fs.Create(ctx, "/x", 0644, 0) + if errno != 0 { + t.Fatalf("create: %s", errno) + } + + if err := fs.Close(); err != nil { + t.Fatalf("close: %s", err) + } + if err := fs.Close(); err != nil { + t.Fatalf("second close: %s", err) + } + if _, errno := f.Pwrite(ctx, []byte("x"), 0); errno != syscall.EBADF { + t.Fatalf("pwrite after close = %s, want %s", errno, syscall.EBADF) + } +} + +func TestFlushLogReturnsWhenBufferClosed(t *testing.T) { + path := filepath.Join(t.TempDir(), "access.log") + f, err := os.OpenFile(path, os.O_CREATE|os.O_WRONLY|os.O_APPEND, 0666) + if err != nil { + t.Fatalf("open access log: %s", err) + } + + fs := &FileSystem{ + checkAccessFile: time.Hour, + rotateAccessLog: 1 << 60, + } + logBuffer := make(chan string, 1) + fs.wg.Add(1) + done := make(chan struct{}) + go func() { + fs.flushLog(f, logBuffer, path) + close(done) + }() + + logBuffer <- "line\n" + close(logBuffer) + select { + case <-done: + case <-time.After(time.Second): + t.Fatal("flushLog did not return after log buffer close") + } + + waited := make(chan struct{}) + go func() { + fs.wg.Wait() + close(waited) + }() + select { + case <-waited: + case <-time.After(time.Second): + t.Fatal("filesystem wait group did not observe flushLog exit") + } +} diff --git a/pkg/meta/base.go b/pkg/meta/base.go index 4b4a7d0315c2..31fac2f2a4bd 100644 --- a/pkg/meta/base.go +++ b/pkg/meta/base.go @@ -534,6 +534,11 @@ func (m *baseMeta) InitSharedMetrics(reg prometheus.Registerer) { if reg == nil { return } + ctx := m.sessCtx + if ctx == nil { + logger.Warnf("InitSharedMetrics called before NewSession; skip shared metrics initialization") + return + } reg.MustRegister(m.usedSpaceG) reg.MustRegister(m.usedInodesG) @@ -562,9 +567,11 @@ func (m *baseMeta) InitSharedMetrics(reg prometheus.Registerer) { } m.subdirInfoG.WithLabelValues(subdir).Set(1) + m.sessWG.Add(1) go func() { + defer m.sessWG.Done() for { - if m.sessCtx != nil && m.sessCtx.Canceled() { + if ctx.Canceled() { return } var totalSpace, availSpace, iused, iavail uint64 @@ -576,21 +583,40 @@ func (m *baseMeta) InitSharedMetrics(reg prometheus.Registerer) { m.totalInodesG.Set(float64(iused + iavail)) } m.updateQuotaMetrics() - utils.SleepWithJitter(time.Second * 10) + if !sleepWithJitterContext(ctx, time.Second*10) { + return + } } }() + m.sessWG.Add(1) go func() { + defer m.sessWG.Done() for { - if m.sessCtx != nil && m.sessCtx.Canceled() { + if ctx.Canceled() { return } m.cleanupQuotaMetrics() - utils.SleepWithJitter(time.Hour) + if !sleepWithJitterContext(ctx, time.Hour) { + return + } } }() } +func sleepWithJitterContext(ctx Context, d time.Duration) bool { + if ctx == nil { + utils.SleepWithJitter(d) + return true + } + select { + case <-ctx.Done(): + return false + case <-time.After(utils.JitterIt(d)): + return true + } +} + func (m *baseMeta) InitMetrics(reg prometheus.Registerer) { if reg == nil { return @@ -739,6 +765,7 @@ func (m *baseMeta) newSessionInfo() []byte { func (m *baseMeta) NewSession(record bool) error { m.sessCtx = Background() ctx := m.sessCtx + m.sessWG.Add(1) go m.refresh(ctx) if err := m.en.cacheACLs(ctx); err != nil { @@ -849,14 +876,19 @@ func (m *baseMeta) OnReload(fn func(f *Format)) { const UmountCode = 11 func (m *baseMeta) refresh(ctx Context) { + defer m.sessWG.Done() for { if ctx.Canceled() { return } if m.conf.Heartbeat > 0 { - utils.SleepWithJitter(m.conf.Heartbeat) + if !sleepWithJitterContext(ctx, m.conf.Heartbeat) { + return + } } else { // use default value - utils.SleepWithJitter(time.Second * 12) + if !sleepWithJitterContext(ctx, time.Second*12) { + return + } } m.sesMu.Lock() if m.umounting { @@ -943,13 +975,22 @@ func (m *baseMeta) CloseSession() error { if m.sid > 0 { err = m.en.doCleanStaleSession(m.sid) } - m.sessCtx.Cancel() + if m.sessCtx != nil { + m.sessCtx.Cancel() + } m.sessWG.Wait() m.stopDeleteSliceTasks() + m.shutdownBase() logger.Infof("close session %d: %v", m.sid, err) return err } +func (m *baseMeta) shutdownBase() { + if m.of != nil { + m.of.Shutdown() + } +} + func (m *baseMeta) FlushSession() { if m.conf.ReadOnly { return diff --git a/pkg/meta/openfile.go b/pkg/meta/openfile.go index e370e6a24640..ca06bc1b1a25 100644 --- a/pkg/meta/openfile.go +++ b/pkg/meta/openfile.go @@ -43,9 +43,12 @@ func (o *openFile) release() { type openfiles struct { sync.Mutex - expire time.Duration - limit uint64 - files map[Ino]*openFile + expire time.Duration + limit uint64 + files map[Ino]*openFile + done chan struct{} + closeOnce sync.Once + wg sync.WaitGroup } func newOpenFiles(expire time.Duration, limit uint64) *openfiles { @@ -53,13 +56,21 @@ func newOpenFiles(expire time.Duration, limit uint64) *openfiles { expire: expire, limit: limit, files: make(map[Ino]*openFile), + done: make(chan struct{}), } + of.wg.Add(1) go of.cleanup() return of } func (o *openfiles) cleanup() { + defer o.wg.Done() for { + select { + case <-o.done: + return + default: + } var ( cnt, deleted, todel int candidateIno Ino @@ -101,10 +112,22 @@ func (o *openfiles) cleanup() { } } o.Unlock() - time.Sleep(time.Millisecond * time.Duration(1000*(cnt+1-deleted*2)/(cnt+1))) + sleep := time.Millisecond * time.Duration(1000*(cnt+1-deleted*2)/(cnt+1)) + select { + case <-o.done: + return + case <-time.After(sleep): + } } } +func (o *openfiles) Shutdown() { + o.closeOnce.Do(func() { + close(o.done) + }) + o.wg.Wait() +} + func (o *openfiles) OpenCheck(ino Ino, attr *Attr) bool { o.Lock() defer o.Unlock() diff --git a/pkg/meta/openfile_close_test.go b/pkg/meta/openfile_close_test.go new file mode 100644 index 000000000000..40ef72b14344 --- /dev/null +++ b/pkg/meta/openfile_close_test.go @@ -0,0 +1,38 @@ +/* + * JuiceFS, Copyright 2020 Juicedata, Inc. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package meta + +import ( + "testing" + "time" +) + +func TestOpenFilesShutdownStopsCleanup(t *testing.T) { + of := newOpenFiles(time.Hour, 0) + done := make(chan struct{}) + go func() { + of.Shutdown() + of.Shutdown() + close(done) + }() + + select { + case <-done: + case <-time.After(time.Second): + t.Fatal("openfiles shutdown did not stop cleanup") + } +} diff --git a/pkg/meta/redis.go b/pkg/meta/redis.go index ccda9fdb9253..3641565f14e5 100644 --- a/pkg/meta/redis.go +++ b/pkg/meta/redis.go @@ -294,6 +294,7 @@ func newRedisMeta(driver, addr string, conf *Config) (Meta, error) { } func (m *redisMeta) Shutdown() error { + m.shutdownBase() if m.cache != nil { m.cache.close() m.cache = nil diff --git a/pkg/meta/sql.go b/pkg/meta/sql.go index d37f611231f2..450a435ab134 100644 --- a/pkg/meta/sql.go +++ b/pkg/meta/sql.go @@ -540,6 +540,7 @@ func newSQLMeta(driver, addr string, conf *Config) (Meta, error) { } func (m *dbMeta) Shutdown() error { + m.shutdownBase() return m.db.Close() } diff --git a/pkg/meta/tkv.go b/pkg/meta/tkv.go index fb4ad4aaac4b..f617d0e222e5 100644 --- a/pkg/meta/tkv.go +++ b/pkg/meta/tkv.go @@ -112,6 +112,7 @@ func newKVMeta(driver, addr string, conf *Config) (Meta, error) { } func (m *kvMeta) Shutdown() error { + m.shutdownBase() return m.client.close() } diff --git a/pkg/vfs/reader.go b/pkg/vfs/reader.go index bfee3ce83e94..9cad4261c460 100644 --- a/pkg/vfs/reader.go +++ b/pkg/vfs/reader.go @@ -263,6 +263,21 @@ func (s *sliceReader) drop() { } } +func (s *sliceReader) close() { + if s.state <= BREAK { + s.state = BREAK + s.cancel() + s.cond.Broadcast() + return + } + if s.refs == 0 { + s.delete() + } else { + s.state = INVALID + s.cond.Broadcast() + } +} + func (s *sliceReader) delete() { *(s.prev) = s.next if s.next != nil { @@ -624,14 +639,26 @@ func (f *fileReader) waitForIO(ctx meta.Context, reqs []*req, buf []byte) (int, } func (f *fileReader) Read(ctx meta.Context, offset uint64, buf []byte) (int, syscall.Errno) { + if f.r.closed.Load() { + return 0, syscall.EBADF + } if f.r.readBufferUsed() > f.r.bufferSize { time.Sleep(time.Millisecond * 10) // slow down for f.r.readBufferUsed() > f.r.bufferSize*2 { // readahead uses 80% of buffer, stop here to avoid OOM + if ctx.Canceled() { + return 0, syscall.EINTR + } + if f.r.closed.Load() { + return 0, syscall.EBADF + } time.Sleep(time.Millisecond * 100) } } f.Lock() defer f.Unlock() + if f.r.closed.Load() { + return 0, syscall.EBADF + } f.acquire() defer f.release() @@ -685,8 +712,11 @@ func (f *fileReader) visit(fn func(s *sliceReader) bool) { func (f *fileReader) Close(ctx meta.Context) { f.Lock() f.closing = true + if f.err == 0 { + f.err = syscall.EBADF + } f.visit(func(s *sliceReader) bool { - s.drop() + s.close() return true }) f.release() @@ -704,6 +734,10 @@ type dataReader struct { readAheadTotal uint64 maxRequests int maxRetries uint32 + done chan struct{} + closeOnce sync.Once + wg sync.WaitGroup + closed atomic.Bool } func NewDataReader(conf *Config, m meta.Meta, store chunk.ChunkStore) DataReader { @@ -722,7 +756,9 @@ func NewDataReader(conf *Config, m meta.Meta, store chunk.ChunkStore) DataReader readAheadMax: uint64(readAheadMax), maxRequests: readAheadMax/conf.Chunk.BlockSize*readSessions + 1, maxRetries: uint32(conf.Meta.Retries), + done: make(chan struct{}), } + r.wg.Add(1) go r.checkReadBuffer() return r } @@ -733,7 +769,13 @@ func (r *dataReader) readBufferUsed() int64 { } func (r *dataReader) checkReadBuffer() { + defer r.wg.Done() for { + select { + case <-r.done: + return + default: + } r.Lock() for _, f := range r.files { for f != nil { @@ -744,7 +786,11 @@ func (r *dataReader) checkReadBuffer() { } } r.Unlock() - time.Sleep(time.Second) + select { + case <-r.done: + return + case <-time.After(time.Second): + } } } @@ -755,8 +801,21 @@ func (r *dataReader) Open(inode Ino, length uint64) FileReader { length: length, } f.last = &(f.slices) + if r.closed.Load() { + f.closing = true + f.err = syscall.EBADF + f.refs = 1 + return f + } r.Lock() + if r.closed.Load() { + r.Unlock() + f.closing = true + f.err = syscall.EBADF + f.refs = 1 + return f + } f.refs = 1 f.next = r.files[inode] r.files[inode] = f @@ -810,6 +869,27 @@ func (r *dataReader) Invalidate(inode Ino, off, length uint64) { }) } +func (r *dataReader) Close() error { + r.closeOnce.Do(func() { + r.closed.Store(true) + close(r.done) + r.Lock() + var files []*fileReader + for _, f := range r.files { + for f != nil { + files = append(files, f) + f = f.next + } + } + r.Unlock() + for _, f := range files { + f.Close(meta.Background()) + } + r.wg.Wait() + }) + return nil +} + func (r *dataReader) readSlice(ctx context.Context, s *meta.Slice, page *chunk.Page, off int) error { buf := page.Data read := 0 diff --git a/pkg/vfs/reader_close_test.go b/pkg/vfs/reader_close_test.go new file mode 100644 index 000000000000..18f3f4ffd0b6 --- /dev/null +++ b/pkg/vfs/reader_close_test.go @@ -0,0 +1,117 @@ +/* + * JuiceFS, Copyright 2020 Juicedata, Inc. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package vfs + +import ( + "context" + "errors" + "sync" + "syscall" + "testing" + "time" + + "github.com/juicedata/juicefs/pkg/chunk" + "github.com/juicedata/juicefs/pkg/meta" +) + +type closeReadMeta struct { + meta.Meta +} + +func (m closeReadMeta) Read(ctx meta.Context, inode Ino, indx uint32, slices *[]meta.Slice) syscall.Errno { + *slices = []meta.Slice{{Id: 1, Size: 4, Len: 4}} + return 0 +} + +type blockingCloseReader struct { + enterOnce sync.Once + cancelOnce sync.Once + entered chan struct{} + canceled chan struct{} +} + +func newBlockingCloseReader() *blockingCloseReader { + return &blockingCloseReader{ + entered: make(chan struct{}), + canceled: make(chan struct{}), + } +} + +func (r *blockingCloseReader) ReadAt(ctx context.Context, p *chunk.Page, off int) (int, error) { + r.enterOnce.Do(func() { close(r.entered) }) + select { + case <-ctx.Done(): + r.cancelOnce.Do(func() { close(r.canceled) }) + return 0, ctx.Err() + case <-time.After(time.Second): + return 0, errors.New("read was not canceled") + } +} + +type closeReadStore struct { + cancelTestStore + reader *blockingCloseReader +} + +func (s *closeReadStore) NewReader(id uint64, length int) chunk.Reader { + return s.reader +} + +func TestDataReaderCloseInterruptsInFlightRead(t *testing.T) { + reader := newBlockingCloseReader() + r := NewDataReader(&Config{ + Meta: &meta.Config{Retries: 1}, + Chunk: &chunk.Config{ + BlockSize: 4, + BufferSize: 1 << 20, + }, + }, closeReadMeta{}, &closeReadStore{reader: reader}).(*dataReader) + + f := r.Open(1, 4) + type result struct { + n int + st syscall.Errno + } + done := make(chan result, 1) + go func() { + n, st := f.Read(meta.Background(), 0, make([]byte, 4)) + done <- result{n: n, st: st} + }() + + select { + case <-reader.entered: + case <-time.After(time.Second): + t.Fatal("read did not start") + } + + if err := r.Close(); err != nil { + t.Fatalf("close: %s", err) + } + select { + case <-reader.canceled: + case <-time.After(time.Second): + t.Fatal("underlying read was not canceled") + } + select { + case got := <-done: + if got.n != 0 || got.st != syscall.EBADF { + t.Fatalf("read after close = (%d, %s), want (0, %s)", got.n, got.st, syscall.EBADF) + } + case <-time.After(time.Second): + t.Fatal("read did not return after close") + } +} diff --git a/pkg/vfs/writer.go b/pkg/vfs/writer.go index 8bdaee774488..6c399aa11bcf 100644 --- a/pkg/vfs/writer.go +++ b/pkg/vfs/writer.go @@ -20,6 +20,7 @@ import ( "math/rand" "runtime" "sync" + "sync/atomic" "syscall" "time" @@ -40,6 +41,28 @@ type FileWriter interface { Truncate(length uint64) } +type closedFileWriter struct { + length uint64 +} + +func (w closedFileWriter) Write(ctx meta.Context, offset uint64, data []byte) syscall.Errno { + return syscall.EBADF +} + +func (w closedFileWriter) Flush(ctx meta.Context) syscall.Errno { + return syscall.EBADF +} + +func (w closedFileWriter) Close(ctx meta.Context) syscall.Errno { + return syscall.EBADF +} + +func (w closedFileWriter) GetLength() uint64 { + return w.length +} + +func (w closedFileWriter) Truncate(length uint64) {} + type DataWriter interface { Open(inode Ino, fleng uint64, tierID uint8) FileWriter Flush(ctx meta.Context, inode Ino) syscall.Errno @@ -69,6 +92,10 @@ func (s *sliceWriter) prepareID(ctx meta.Context, retry bool) { f := s.chunk.file f.Lock() for s.id == 0 { + if ctx.Canceled() { + s.err = syscall.EINTR + break + } var id uint64 f.Unlock() st := f.w.m.NewSlice(ctx, &id) @@ -103,14 +130,18 @@ func (s *sliceWriter) markDone() { } // freezed, no more data -func (s *sliceWriter) flushData() { +func (s *sliceWriter) flushData(ctx meta.Context) { defer s.markDone() if s.slen == 0 { return } - s.prepareID(meta.Background(), true) - if s.err != 0 { - logger.Infof("flush inode: %v chunk: %d err: %s", s.chunk.file.inode, s.id, s.err) + s.prepareID(ctx, true) + f := s.chunk.file + f.Lock() + err := s.err + f.Unlock() + if err != 0 { + logger.Infof("flush inode: %v chunk: %d err: %s", s.chunk.file.inode, s.id, err) s.writer.Abort() return } @@ -119,7 +150,9 @@ func (s *sliceWriter) flushData() { logger.Errorf("upload inode: %v chunk: %v (length: %v) fail: %s", s.chunk.file.inode, s.id, s.length, err) s.writer.Abort() + f.Lock() s.err = syscall.EIO + f.Unlock() } } @@ -137,7 +170,7 @@ func (s *sliceWriter) write(ctx meta.Context, off uint32, data []uint8) syscall. s.lastMod = time.Now() if s.slen == meta.ChunkSize { s.freezed = true - go s.flushData() + go s.flushData(meta.Background()) } else if int(s.slen) >= f.w.blockSize { if s.id > 0 { err := s.writer.FlushTo(int(s.slen)) @@ -167,7 +200,7 @@ func (c *chunkWriter) findWritableSlice(pos uint32, size uint32) *sliceWriter { return s } else if i > 3 { s.freezed = true - go s.flushData() + go s.flushData(meta.Background()) } } if pos < s.off+s.slen && s.off < pos+size { @@ -190,7 +223,7 @@ func (c *chunkWriter) commitThread() { for !s.done { if s.notify.WaitWithTimeout(time.Millisecond*100) && !s.freezed && time.Since(s.started) > flushDuration*2 { s.freezed = true - go s.flushData() + go s.flushData(meta.Background()) } } err := s.err @@ -229,6 +262,7 @@ type fileWriter struct { length uint64 tierID uint8 err syscall.Errno + closing bool flushwaiting uint16 writewaiting uint16 refs uint16 @@ -295,13 +329,28 @@ func (w *dataWriter) usedBufferSize() int64 { } func (f *fileWriter) Write(ctx meta.Context, off uint64, data []byte) syscall.Errno { + if f.w.closed.Load() { + return syscall.EBADF + } for f.totalSlices() >= 1000 { + if ctx.Canceled() { + return syscall.EINTR + } + if f.w.closed.Load() { + return syscall.EBADF + } time.Sleep(time.Millisecond) } if f.w.usedBufferSize() > f.w.bufferSize { // slow down time.Sleep(time.Millisecond * 10) for f.w.usedBufferSize() > f.w.bufferSize*2 { + if ctx.Canceled() { + return syscall.EINTR + } + if f.w.closed.Load() { + return syscall.EBADF + } time.Sleep(time.Millisecond * 100) } } @@ -309,16 +358,24 @@ func (f *fileWriter) Write(ctx meta.Context, off uint64, data []byte) syscall.Er s := time.Now() f.Lock() defer f.Unlock() + if f.closing || f.w.closed.Load() { + return syscall.EBADF + } size := uint64(len(data)) f.writewaiting++ - for f.flushwaiting > 0 { + defer func() { f.writewaiting-- }() + for f.flushwaiting > 0 && !f.closing { if f.writecond.WaitWithTimeout(time.Second) && ctx.Canceled() { - f.writewaiting-- logger.Warnf("write %d interrupted after %d", f.inode, time.Since(s)) return syscall.EINTR } } - f.writewaiting-- + if f.closing { + return syscall.EBADF + } + if f.w.closed.Load() { + return syscall.EBADF + } indx := uint32(off / meta.ChunkSize) pos := uint32(off % meta.ChunkSize) @@ -367,11 +424,15 @@ func (f *fileWriter) flush(ctx meta.Context, writeback bool) syscall.Errno { for _, s := range c.slices { if !s.freezed { s.freezed = true - go s.flushData() + go s.flushData(ctx) } } } - if f.flushcond.WaitWithTimeout(time.Second*3) && ctx.Canceled() && time.Since(s) > f.w.conf.Chunk.PutTimeout*2 { + wait := time.Second * 3 + if f.w.closed.Load() && f.w.conf.Chunk.PutTimeout > 0 && f.w.conf.Chunk.PutTimeout < wait { + wait = f.w.conf.Chunk.PutTimeout + } + if f.flushcond.WaitWithTimeout(wait) && ctx.Canceled() && (f.w.closed.Load() || time.Since(s) > f.w.conf.Chunk.PutTimeout*2) { logger.Warnf("flush %d interrupted after %d", f.inode, time.Since(s)) err = syscall.EINTR break @@ -409,6 +470,15 @@ func (f *fileWriter) Close(ctx meta.Context) syscall.Errno { return f.Flush(ctx) } +func (f *fileWriter) closeForDataWriter() { + f.Lock() + f.closing = true + if f.writewaiting > 0 { + f.writecond.Broadcast() + } + f.Unlock() +} + func (f *fileWriter) GetLength() uint64 { f.Lock() defer f.Unlock() @@ -432,6 +502,10 @@ type dataWriter struct { bufferSize int64 files map[Ino]*fileWriter maxRetries uint32 + done chan struct{} + closeOnce sync.Once + wg sync.WaitGroup + closed atomic.Bool } func NewDataWriter(conf *Config, m meta.Meta, store chunk.ChunkStore, reader DataReader) DataWriter { @@ -444,13 +518,21 @@ func NewDataWriter(conf *Config, m meta.Meta, store chunk.ChunkStore, reader Dat bufferSize: int64(conf.Chunk.BufferSize), files: make(map[Ino]*fileWriter), maxRetries: uint32(conf.Meta.Retries), + done: make(chan struct{}), } + w.wg.Add(1) go w.flushAll() return w } func (w *dataWriter) flushAll() { + defer w.wg.Done() for { + select { + case <-w.done: + return + default: + } w.Lock() now := time.Now() for _, f := range w.files { @@ -466,7 +548,7 @@ func (w *dataWriter) flushAll() { if !s.freezed && (now.Sub(s.started) > flushDuration || now.Sub(s.lastMod) > time.Second && now.Sub(s.started) > time.Second || tooMany && i%2 == lastBit && j <= hs) { s.freezed = true - go s.flushData() + go s.flushData(meta.Background()) } } } @@ -475,13 +557,23 @@ func (w *dataWriter) flushAll() { w.Lock() } w.Unlock() - time.Sleep(time.Millisecond * 100) + select { + case <-w.done: + return + case <-time.After(time.Millisecond * 100): + } } } func (w *dataWriter) Open(inode Ino, len uint64, tierID uint8) FileWriter { + if w.closed.Load() { + return closedFileWriter{length: len} + } w.Lock() defer w.Unlock() + if w.closed.Load() { + return closedFileWriter{length: len} + } f, ok := w.files[inode] if !ok { f = &fileWriter{ @@ -509,7 +601,7 @@ func (w *dataWriter) free(f *fileWriter) { w.Lock() defer w.Unlock() f.refs-- - if f.refs == 0 { + if f.refs == 0 && w.files[f.inode] == f { delete(w.files, f.inode) } } @@ -544,21 +636,107 @@ func (w *dataWriter) UpdateMtime(inode Ino, mtime time.Time) { } } -func (w *dataWriter) FlushAll() error { +func (w *dataWriter) flushAllFiles(ctx meta.Context, abortOnError bool) error { var err error w.Lock() for inode, ind := range w.files { ind.refs++ w.Unlock() - eno := ind.Flush(meta.Background()) - w.free(ind) + eno := ind.Flush(ctx) if eno != 0 { logger.Errorf("flush %s: %s", inode, eno) + if abortOnError { + ind.abortPending() + } + w.free(ind) return eno } + w.free(ind) logger.Debugf("Flush %d", inode) w.Lock() } w.Unlock() return err } + +func (w *dataWriter) FlushAll() error { + return w.flushAllFiles(meta.Background(), false) +} + +func (w *dataWriter) closeTrackedWriters() { + w.Lock() + files := make([]*fileWriter, 0, len(w.files)) + for _, f := range w.files { + files = append(files, f) + } + w.Unlock() + + for _, f := range files { + f.closeForDataWriter() + } +} + +func (w *dataWriter) closeFlushTimeout() time.Duration { + timeout := w.conf.Chunk.PutTimeout * 2 + if timeout <= 0 { + timeout = time.Minute * 2 + } + if timeout < time.Millisecond { + timeout = time.Millisecond + } + return timeout +} + +func (w *dataWriter) Close() error { + var err error + w.closeOnce.Do(func() { + close(w.done) + w.wg.Wait() + w.closed.Store(true) + w.closeTrackedWriters() + ctx := meta.WrapWithTimeout(meta.Background(), w.closeFlushTimeout()) + defer ctx.Cancel() + err = w.flushAllFiles(ctx, true) + if err == nil && ctx.Canceled() { + err = syscall.EINTR + } + if err != nil { + w.abortPending() + } + }) + return err +} + +func (w *dataWriter) abortPending() { + w.Lock() + files := make([]*fileWriter, 0, len(w.files)) + for _, f := range w.files { + f.refs++ + files = append(files, f) + } + w.Unlock() + + for _, f := range files { + f.abortPending() + w.free(f) + } +} + +func (f *fileWriter) abortPending() { + f.Lock() + defer f.Unlock() + for _, c := range f.chunks { + for _, s := range c.slices { + if s.done { + continue + } + s.freezed = true + s.err = syscall.EIO + s.writer.Abort() + s.done = true + s.notify.Signal() + } + } + f.flushcond.Broadcast() + f.writecond.Broadcast() +} diff --git a/pkg/vfs/writer_cancel_test.go b/pkg/vfs/writer_cancel_test.go new file mode 100644 index 000000000000..06f2d979ec9f --- /dev/null +++ b/pkg/vfs/writer_cancel_test.go @@ -0,0 +1,391 @@ +/* + * JuiceFS, Copyright 2020 Juicedata, Inc. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package vfs + +import ( + "errors" + "io" + "sync" + "syscall" + "testing" + "time" + + "github.com/juicedata/juicefs/pkg/chunk" + "github.com/juicedata/juicefs/pkg/meta" + "github.com/juicedata/juicefs/pkg/object" + "github.com/juicedata/juicefs/pkg/utils" +) + +type cancelTestStore struct { + usedMemory int64 +} + +func (s *cancelTestStore) NewReader(id uint64, length int) chunk.Reader { return nil } +func (s *cancelTestStore) NewWriter(id uint64, tierID uint8) chunk.Writer { + return nil +} +func (s *cancelTestStore) Remove(id uint64, length int) error { return nil } +func (s *cancelTestStore) FillCache(id uint64, length uint32) error { + return nil +} +func (s *cancelTestStore) EvictCache(id uint64, length uint32) error { + return nil +} +func (s *cancelTestStore) CheckCache(id uint64, length uint32, handler func(exists bool, loc string, size int)) error { + return nil +} +func (s *cancelTestStore) UsedMemory() int64 { return s.usedMemory } +func (s *cancelTestStore) UpdateLimit(upload, download int64) { +} +func (s *cancelTestStore) BlobStorage() object.ObjectStorage { return nil } + +func TestFileWriterWriteCanceledWhileWaitingForSliceLimit(t *testing.T) { + ctx := meta.NewContext(1, 0, []uint32{0}) + ctx.Cancel() + + f := &fileWriter{ + w: &dataWriter{ + store: &cancelTestStore{}, + bufferSize: 1 << 20, + }, + chunks: map[uint32]*chunkWriter{ + 0: {slices: make([]*sliceWriter, 1000)}, + }, + } + + if st := writeWithTimeout(t, f, ctx); st != syscall.EINTR { + t.Fatalf("Write status = %s, want %s", st, syscall.EINTR) + } +} + +func TestFileWriterWriteCanceledWhileThrottledByBuffer(t *testing.T) { + ctx := meta.NewContext(1, 0, []uint32{0}) + ctx.Cancel() + + buf := utils.Alloc(4) + defer utils.Free(buf) + + f := &fileWriter{ + w: &dataWriter{ + store: &cancelTestStore{}, + bufferSize: 1, + }, + } + + if st := writeWithTimeout(t, f, ctx); st != syscall.EINTR { + t.Fatalf("Write status = %s, want %s", st, syscall.EINTR) + } +} + +func writeWithTimeout(t *testing.T, f *fileWriter, ctx meta.Context) syscall.Errno { + t.Helper() + done := make(chan syscall.Errno, 1) + go func() { + done <- f.Write(ctx, 0, nil) + }() + + select { + case st := <-done: + return st + case <-time.After(200 * time.Millisecond): + t.Fatal("Write did not return after context cancellation") + return 0 + } +} + +type blockingChunkWriter struct { + abortOnce sync.Once + aborted chan struct{} +} + +func newBlockingChunkWriter() *blockingChunkWriter { + return &blockingChunkWriter{aborted: make(chan struct{})} +} + +func (w *blockingChunkWriter) WriteAt(p []byte, off int64) (int, error) { + return len(p), nil +} + +func (w *blockingChunkWriter) ID() uint64 { return 1 } + +func (w *blockingChunkWriter) SetID(id uint64) {} + +func (w *blockingChunkWriter) SetWriteback(enabled bool) {} + +func (w *blockingChunkWriter) FlushTo(offset int) error { return nil } + +func (w *blockingChunkWriter) Finish(length int) error { + <-w.aborted + return errors.New("aborted") +} + +func (w *blockingChunkWriter) Abort() { + w.abortOnce.Do(func() { + close(w.aborted) + }) +} + +var _ chunk.Writer = (*blockingChunkWriter)(nil) +var _ io.WriterAt = (*blockingChunkWriter)(nil) + +type recordingChunkWriter struct { + data []byte +} + +func (w *recordingChunkWriter) WriteAt(p []byte, off int64) (int, error) { + end := int(off) + len(p) + if end > len(w.data) { + buf := make([]byte, end) + copy(buf, w.data) + w.data = buf + } + copy(w.data[off:], p) + return len(p), nil +} + +func (w *recordingChunkWriter) ID() uint64 { return 1 } + +func (w *recordingChunkWriter) SetID(id uint64) {} + +func (w *recordingChunkWriter) SetWriteback(enabled bool) {} + +func (w *recordingChunkWriter) FlushTo(offset int) error { return nil } + +func (w *recordingChunkWriter) Finish(length int) error { return nil } + +func (w *recordingChunkWriter) Abort() {} + +var _ chunk.Writer = (*recordingChunkWriter)(nil) +var _ io.WriterAt = (*recordingChunkWriter)(nil) + +func TestDataWriterCloseAbortsTimedOutFlush(t *testing.T) { + bw := newBlockingChunkWriter() + w := &dataWriter{ + conf: &Config{Chunk: &chunk.Config{ + BlockSize: 1 << 20, + BufferSize: 1 << 20, + PutTimeout: 20 * time.Millisecond, + }}, + store: &cancelTestStore{}, + done: make(chan struct{}), + files: make(map[Ino]*fileWriter), + bufferSize: 1 << 20, + } + f := &fileWriter{ + w: w, + inode: 1, + chunks: make(map[uint32]*chunkWriter), + } + f.flushcond = utils.NewCond(f) + f.writecond = utils.NewCond(f) + c := &chunkWriter{indx: 0, file: f} + s := &sliceWriter{ + id: 1, + chunk: c, + writer: bw, + slen: 1, + notify: utils.NewCond(f), + started: time.Now(), + lastMod: time.Now(), + } + c.slices = []*sliceWriter{s} + f.chunks[c.indx] = c + w.files[f.inode] = f + + done := make(chan error, 1) + go func() { + done <- w.Close() + }() + + select { + case err := <-done: + if !errors.Is(err, syscall.EINTR) { + t.Fatalf("Close error = %v, want %s", err, syscall.EINTR) + } + case <-time.After(500 * time.Millisecond): + t.Fatal("Close did not abort a timed-out flush") + } + + select { + case <-bw.aborted: + default: + t.Fatal("Close did not abort the stuck slice writer") + } +} + +func TestClosedDataWriterOpenCloseDoesNotRemoveTrackedWriter(t *testing.T) { + w := &dataWriter{ + store: &cancelTestStore{}, + done: make(chan struct{}), + files: make(map[Ino]*fileWriter), + bufferSize: 1 << 20, + } + tracked := &fileWriter{ + w: w, + inode: 1, + refs: 1, + chunks: make(map[uint32]*chunkWriter), + } + tracked.flushcond = utils.NewCond(tracked) + tracked.writecond = utils.NewCond(tracked) + w.files[tracked.inode] = tracked + w.closed.Store(true) + + closed := w.Open(tracked.inode, tracked.length, tracked.tierID) + if _, ok := closed.(closedFileWriter); !ok { + t.Fatalf("closed dataWriter Open returned %T, want closedFileWriter", closed) + } + if st := closed.Close(meta.Background()); st != syscall.EBADF { + t.Fatalf("closed writer Close status = %s, want %s", st, syscall.EBADF) + } + if got := w.files[tracked.inode]; got != tracked { + t.Fatalf("closed writer removed tracked writer: got %p, want %p", got, tracked) + } + if tracked.refs != 1 { + t.Fatalf("tracked writer refs = %d, want 1", tracked.refs) + } +} + +func TestFileWriterWriteFailsWhenDataWriterClosesWhileWaitingForFlush(t *testing.T) { + w := &dataWriter{ + store: &cancelTestStore{}, + done: make(chan struct{}), + files: make(map[Ino]*fileWriter), + blockSize: 1 << 20, + bufferSize: 1 << 20, + } + f := &fileWriter{ + w: w, + inode: 1, + refs: 1, + flushwaiting: 1, + chunks: make(map[uint32]*chunkWriter), + } + f.flushcond = utils.NewCond(f) + f.writecond = utils.NewCond(f) + writer := &recordingChunkWriter{} + c := &chunkWriter{indx: 0, file: f} + s := &sliceWriter{ + id: 1, + chunk: c, + writer: writer, + notify: utils.NewCond(f), + started: time.Now(), + lastMod: time.Now(), + } + c.slices = []*sliceWriter{s} + f.chunks[c.indx] = c + w.files[f.inode] = f + + done := make(chan syscall.Errno, 1) + go func() { + done <- f.Write(meta.Background(), 0, []byte("x")) + }() + waitForWriteWaiting(t, f) + + w.closed.Store(true) + w.closeTrackedWriters() + + select { + case st := <-done: + if st != syscall.EBADF { + t.Fatalf("Write status = %s, want %s", st, syscall.EBADF) + } + case <-time.After(500 * time.Millisecond): + t.Fatal("Write did not return after dataWriter close") + } + if s.slen != 0 { + t.Fatalf("slice length = %d, want 0", s.slen) + } + if len(writer.data) != 0 { + t.Fatalf("writer received %d bytes after close, want 0", len(writer.data)) + } +} + +func TestFileWriterWriteFailsWhenDataWriterClosesBeforeFlushWaitEnds(t *testing.T) { + w := &dataWriter{ + store: &cancelTestStore{}, + done: make(chan struct{}), + files: make(map[Ino]*fileWriter), + blockSize: 1 << 20, + bufferSize: 1 << 20, + } + f := &fileWriter{ + w: w, + inode: 1, + refs: 1, + flushwaiting: 1, + chunks: make(map[uint32]*chunkWriter), + } + f.flushcond = utils.NewCond(f) + f.writecond = utils.NewCond(f) + writer := &recordingChunkWriter{} + c := &chunkWriter{indx: 0, file: f} + s := &sliceWriter{ + id: 1, + chunk: c, + writer: writer, + notify: utils.NewCond(f), + started: time.Now(), + lastMod: time.Now(), + } + c.slices = []*sliceWriter{s} + f.chunks[c.indx] = c + w.files[f.inode] = f + + done := make(chan syscall.Errno, 1) + go func() { + done <- f.Write(meta.Background(), 0, []byte("x")) + }() + waitForWriteWaiting(t, f) + + f.Lock() + w.closed.Store(true) + f.flushwaiting = 0 + f.writecond.Broadcast() + f.Unlock() + + select { + case st := <-done: + if st != syscall.EBADF { + t.Fatalf("Write status = %s, want %s", st, syscall.EBADF) + } + case <-time.After(500 * time.Millisecond): + t.Fatal("Write did not return after dataWriter close") + } + if s.slen != 0 { + t.Fatalf("slice length = %d, want 0", s.slen) + } + if len(writer.data) != 0 { + t.Fatalf("writer received %d bytes after close, want 0", len(writer.data)) + } +} + +func waitForWriteWaiting(t *testing.T, f *fileWriter) { + t.Helper() + deadline := time.Now().Add(500 * time.Millisecond) + for time.Now().Before(deadline) { + f.Lock() + waiting := f.writewaiting + f.Unlock() + if waiting > 0 { + return + } + time.Sleep(time.Millisecond) + } + t.Fatal("Write did not start waiting for flush") +} From 16188fa6d58083d4b03c14be8578d17f9918d758 Mon Sep 17 00:00:00 2001 From: Ian Date: Thu, 21 May 2026 08:55:09 +0200 Subject: [PATCH 4/4] fix: release delayed read retries on close --- pkg/vfs/reader.go | 15 ++++++++++++- pkg/vfs/reader_close_test.go | 42 ++++++++++++++++++++++++++++++++++++ 2 files changed, 56 insertions(+), 1 deletion(-) diff --git a/pkg/vfs/reader.go b/pkg/vfs/reader.go index 9cad4261c460..3aa912111b1b 100644 --- a/pkg/vfs/reader.go +++ b/pkg/vfs/reader.go @@ -101,13 +101,14 @@ type sliceReader struct { currentPos uint32 lastAccess time.Time cond *utils.Cond + timer *time.Timer next *sliceReader prev **sliceReader refs uint16 } func (s *sliceReader) delay(delay time.Duration) { - time.AfterFunc(delay, s.run) + s.timer = time.AfterFunc(delay, s.run) } func (s *sliceReader) done(err syscall.Errno, delay time.Duration) { @@ -163,6 +164,7 @@ func (s *sliceReader) run() { f := s.file f.Lock() defer f.Unlock() + s.timer = nil if s.state != NEW || f.shouldStop() { s.done(0, 0) } @@ -265,9 +267,20 @@ func (s *sliceReader) drop() { func (s *sliceReader) close() { if s.state <= BREAK { + stopped := false + if s.timer != nil { + stopped = s.timer.Stop() + s.timer = nil + } s.state = BREAK s.cancel() s.cond.Broadcast() + if stopped { + s.state = INVALID + if s.refs == 0 { + s.delete() + } + } return } if s.refs == 0 { diff --git a/pkg/vfs/reader_close_test.go b/pkg/vfs/reader_close_test.go index 18f3f4ffd0b6..7831f8609e8e 100644 --- a/pkg/vfs/reader_close_test.go +++ b/pkg/vfs/reader_close_test.go @@ -26,6 +26,7 @@ import ( "github.com/juicedata/juicefs/pkg/chunk" "github.com/juicedata/juicefs/pkg/meta" + "github.com/juicedata/juicefs/pkg/utils" ) type closeReadMeta struct { @@ -115,3 +116,44 @@ func TestDataReaderCloseInterruptsInFlightRead(t *testing.T) { t.Fatal("read did not return after close") } } + +func TestDataReaderCloseReleasesDelayedRetryReadBuffer(t *testing.T) { + r := &dataReader{ + files: make(map[Ino]*fileReader), + done: make(chan struct{}), + } + f := &fileReader{ + r: r, + inode: 1, + length: 4, + refs: 1, + } + page := chunk.NewOffPage(4) + s := &sliceReader{ + file: f, + block: &frange{off: 0, len: 4}, + state: NEW, + page: page, + cond: utils.NewCond(&f.Mutex), + } + s.ctx, s.cancel = context.WithCancel(context.Background()) + s.prev = &f.slices + f.slices = s + f.last = &s.next + r.files[f.inode] = f + readBufferUsed.Add(int64(cap(s.page.Data))) + s.delay(time.Hour) + + if err := r.Close(); err != nil { + t.Fatalf("close: %s", err) + } + if page.Data != nil { + t.Fatal("delayed retry page was not released during close") + } + if f.slices != nil { + t.Fatal("delayed retry slice is still linked after close") + } + if got := r.files[f.inode]; got != nil { + t.Fatalf("file reader still tracked after close: %p", got) + } +}