mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-10-06 14:45:51 +00:00
Compare commits
3
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
da39ee65a3 | ||
|
|
e91d43ef08 | ||
|
|
58cf7e9a89 |
@@ -18,6 +18,20 @@ import (
|
||||
// need to have logic similar to FilerStoreWrapper
|
||||
// e.g. fill fileId field for chunks
|
||||
|
||||
// bufferedEvent represents a subscription event captured during a directory refresh.
|
||||
type bufferedEvent struct {
|
||||
oldPath util.FullPath
|
||||
newEntry *filer.Entry
|
||||
}
|
||||
|
||||
// refreshState tracks events that arrive while a directory is being refreshed
|
||||
// from the filer. After the refresh snapshot is applied, buffered events are
|
||||
// replayed so that creates, deletes, and updates that raced with the snapshot
|
||||
// are not lost.
|
||||
type refreshState struct {
|
||||
events []bufferedEvent
|
||||
}
|
||||
|
||||
type MetaCache struct {
|
||||
root util.FullPath
|
||||
localStore filer.VirtualFilerStore
|
||||
@@ -29,6 +43,7 @@ type MetaCache struct {
|
||||
invalidateFunc func(fullpath util.FullPath, entry *filer_pb.Entry)
|
||||
onDirectoryUpdate func(dir util.FullPath)
|
||||
visitGroup singleflight.Group // deduplicates concurrent EnsureVisited calls for the same path
|
||||
refreshing map[util.FullPath]*refreshState
|
||||
}
|
||||
|
||||
func NewMetaCache(dbFolder string, uidGidMapper *UidGidMapper, root util.FullPath,
|
||||
@@ -45,6 +60,7 @@ func NewMetaCache(dbFolder string, uidGidMapper *UidGidMapper, root util.FullPat
|
||||
invalidateFunc: func(fullpath util.FullPath, entry *filer_pb.Entry) {
|
||||
invalidateFunc(fullpath, entry)
|
||||
},
|
||||
refreshing: make(map[util.FullPath]*refreshState),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -69,6 +85,11 @@ func openMetaStore(dbFolder string) (*leveldb.LevelDBStore, filer.VirtualFilerSt
|
||||
func (mc *MetaCache) InsertEntry(ctx context.Context, entry *filer.Entry) error {
|
||||
mc.Lock()
|
||||
defer mc.Unlock()
|
||||
// Buffer the insert if the parent directory is being refreshed
|
||||
dir, _ := entry.DirAndName()
|
||||
if state := mc.isRefreshingDir(util.FullPath(dir)); state != nil {
|
||||
state.events = append(state.events, bufferedEvent{newEntry: entry})
|
||||
}
|
||||
return mc.doInsertEntry(ctx, entry)
|
||||
}
|
||||
|
||||
@@ -76,16 +97,33 @@ func (mc *MetaCache) doInsertEntry(ctx context.Context, entry *filer.Entry) erro
|
||||
return mc.localStore.InsertEntry(ctx, entry)
|
||||
}
|
||||
|
||||
// doBatchInsertEntries inserts multiple entries using LevelDB's batch write.
|
||||
// This is more efficient than inserting entries one by one.
|
||||
func (mc *MetaCache) doBatchInsertEntries(ctx context.Context, entries []*filer.Entry) error {
|
||||
return mc.leveldbStore.BatchInsertEntries(ctx, entries)
|
||||
}
|
||||
|
||||
func (mc *MetaCache) AtomicUpdateEntryFromFiler(ctx context.Context, oldPath util.FullPath, newEntry *filer.Entry) error {
|
||||
mc.Lock()
|
||||
defer mc.Unlock()
|
||||
|
||||
// If the affected directory is being refreshed, buffer the event
|
||||
// instead of applying it. It will be replayed after the snapshot commits.
|
||||
if oldPath != "" {
|
||||
dir, _ := oldPath.DirAndName()
|
||||
if state := mc.isRefreshingDir(util.FullPath(dir)); state != nil {
|
||||
state.events = append(state.events, bufferedEvent{oldPath, newEntry})
|
||||
return nil
|
||||
}
|
||||
}
|
||||
if newEntry != nil {
|
||||
newDir, _ := newEntry.DirAndName()
|
||||
if state := mc.isRefreshingDir(util.FullPath(newDir)); state != nil {
|
||||
state.events = append(state.events, bufferedEvent{oldPath, newEntry})
|
||||
return nil
|
||||
}
|
||||
}
|
||||
|
||||
return mc.doAtomicUpdateEntryFromFiler(ctx, oldPath, newEntry)
|
||||
}
|
||||
|
||||
// doAtomicUpdateEntryFromFiler is the core logic for applying a filer event
|
||||
// to the local cache. Caller must hold mc.Lock().
|
||||
func (mc *MetaCache) doAtomicUpdateEntryFromFiler(ctx context.Context, oldPath util.FullPath, newEntry *filer.Entry) error {
|
||||
entry, err := mc.localStore.FindEntry(ctx, oldPath)
|
||||
if err != nil && err != filer_pb.ErrNotFound {
|
||||
glog.Errorf("Metacache: find entry error: %v", err)
|
||||
@@ -104,8 +142,6 @@ func (mc *MetaCache) AtomicUpdateEntryFromFiler(ctx context.Context, oldPath uti
|
||||
}
|
||||
}
|
||||
}
|
||||
} else {
|
||||
// println("unknown old directory:", oldDir)
|
||||
}
|
||||
|
||||
if newEntry != nil {
|
||||
@@ -143,6 +179,14 @@ func (mc *MetaCache) FindEntry(ctx context.Context, fp util.FullPath) (entry *fi
|
||||
func (mc *MetaCache) DeleteEntry(ctx context.Context, fp util.FullPath) (err error) {
|
||||
mc.Lock()
|
||||
defer mc.Unlock()
|
||||
// Buffer the delete if the parent directory is being refreshed
|
||||
dir, _ := fp.DirAndName()
|
||||
if state := mc.isRefreshingDir(util.FullPath(dir)); state != nil {
|
||||
state.events = append(state.events, bufferedEvent{oldPath: fp})
|
||||
}
|
||||
// Always apply the delete directly as well, so the entry is removed
|
||||
// immediately for the current node's view. CommitRefresh's snapshot
|
||||
// may re-insert it, but the buffered event replay will re-delete it.
|
||||
return mc.localStore.DeleteEntry(ctx, fp)
|
||||
}
|
||||
func (mc *MetaCache) DeleteFolderChildren(ctx context.Context, fp util.FullPath) (err error) {
|
||||
@@ -151,6 +195,58 @@ func (mc *MetaCache) DeleteFolderChildren(ctx context.Context, fp util.FullPath)
|
||||
return mc.localStore.DeleteFolderChildren(ctx, fp)
|
||||
}
|
||||
|
||||
// BeginRefresh starts buffering subscription events for dirPath.
|
||||
// While a refresh is active, AtomicUpdateEntryFromFiler will buffer events
|
||||
// instead of applying them, so they can be replayed after the snapshot is committed.
|
||||
func (mc *MetaCache) BeginRefresh(dirPath util.FullPath) {
|
||||
mc.Lock()
|
||||
defer mc.Unlock()
|
||||
mc.refreshing[dirPath] = &refreshState{}
|
||||
}
|
||||
|
||||
// CommitRefresh atomically replaces a directory's cached entries with the
|
||||
// filer snapshot, then replays any subscription events that were buffered
|
||||
// during the refresh. This ensures no creates, deletes, or updates are lost.
|
||||
func (mc *MetaCache) CommitRefresh(ctx context.Context, dirPath util.FullPath, entries []*filer.Entry) error {
|
||||
mc.Lock()
|
||||
defer mc.Unlock()
|
||||
|
||||
// Clear stale entries and insert the fresh snapshot
|
||||
if err := mc.localStore.DeleteFolderChildren(ctx, dirPath); err != nil {
|
||||
return err
|
||||
}
|
||||
if len(entries) > 0 {
|
||||
if err := mc.leveldbStore.BatchInsertEntries(ctx, entries); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
// Replay buffered events so mutations that raced with the snapshot are applied
|
||||
state := mc.refreshing[dirPath]
|
||||
delete(mc.refreshing, dirPath)
|
||||
if state != nil {
|
||||
for _, ev := range state.events {
|
||||
if err := mc.doAtomicUpdateEntryFromFiler(ctx, ev.oldPath, ev.newEntry); err != nil {
|
||||
glog.Warningf("replay buffered event for %s: %v", dirPath, err)
|
||||
}
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// CancelRefresh discards the refresh state without replaying buffered events.
|
||||
func (mc *MetaCache) CancelRefresh(dirPath util.FullPath) {
|
||||
mc.Lock()
|
||||
defer mc.Unlock()
|
||||
delete(mc.refreshing, dirPath)
|
||||
}
|
||||
|
||||
// isRefreshing returns the refresh state for the directory containing fp,
|
||||
// or nil if no refresh is active. Caller must hold mc.Lock().
|
||||
func (mc *MetaCache) isRefreshingDir(dirPath util.FullPath) *refreshState {
|
||||
return mc.refreshing[dirPath]
|
||||
}
|
||||
|
||||
func (mc *MetaCache) ListDirectoryEntries(ctx context.Context, dirPath util.FullPath, startFileName string, includeStartFile bool, limit int64, eachEntryFunc filer.ListEachEntryFunc) error {
|
||||
mc.RLock()
|
||||
defer mc.RUnlock()
|
||||
|
||||
@@ -49,11 +49,6 @@ func EnsureVisited(mc *MetaCache, client filer_pb.FilerClient, dirPath util.Full
|
||||
return g.Wait()
|
||||
}
|
||||
|
||||
// batchInsertSize is the number of entries to accumulate before flushing to LevelDB.
|
||||
// 100 provides a balance between memory usage (~100 Entry pointers) and write efficiency
|
||||
// (fewer disk syncs). Larger values reduce I/O overhead but increase memory and latency.
|
||||
const batchInsertSize = 100
|
||||
|
||||
func doEnsureVisited(ctx context.Context, mc *MetaCache, client filer_pb.FilerClient, path util.FullPath) error {
|
||||
// Use singleflight to deduplicate concurrent requests for the same path
|
||||
_, err, _ := mc.visitGroup.Do(string(path), func() (interface{}, error) {
|
||||
@@ -69,42 +64,36 @@ func doEnsureVisited(ctx context.Context, mc *MetaCache, client filer_pb.FilerCl
|
||||
|
||||
glog.V(4).Infof("ReadDirAllEntries %s ...", path)
|
||||
|
||||
// Collect entries in batches for efficient LevelDB writes
|
||||
var batch []*filer.Entry
|
||||
// Start buffering subscription events for this directory.
|
||||
// Any events that arrive during the fetch will be replayed
|
||||
// after the snapshot is committed, preventing lost mutations.
|
||||
mc.BeginRefresh(path)
|
||||
|
||||
// Collect all entries from the filer. No lock is held during
|
||||
// network I/O so subscription events can still be buffered.
|
||||
var allEntries []*filer.Entry
|
||||
|
||||
fetchErr := util.Retry("ReadDirAllEntries", func() error {
|
||||
batch = nil // Reset batch on retry, allow GC of previous entries
|
||||
allEntries = nil // Reset on retry
|
||||
return filer_pb.ReadDirAllEntries(ctx, client, path, "", func(pbEntry *filer_pb.Entry, isLast bool) error {
|
||||
entry := filer.FromPbEntry(string(path), pbEntry)
|
||||
if IsHiddenSystemEntry(string(path), entry.Name()) {
|
||||
return nil
|
||||
}
|
||||
|
||||
batch = append(batch, entry)
|
||||
|
||||
// Flush batch when it reaches the threshold
|
||||
// Don't rely on isLast here - hidden entries may cause early return
|
||||
if len(batch) >= batchInsertSize {
|
||||
// No lock needed - LevelDB Write() is thread-safe
|
||||
if err := mc.doBatchInsertEntries(ctx, batch); err != nil {
|
||||
return fmt.Errorf("batch insert for %s: %w", path, err)
|
||||
}
|
||||
// Create new slice to allow GC of flushed entries
|
||||
batch = make([]*filer.Entry, 0, batchInsertSize)
|
||||
}
|
||||
allEntries = append(allEntries, entry)
|
||||
return nil
|
||||
})
|
||||
})
|
||||
|
||||
if fetchErr != nil {
|
||||
mc.CancelRefresh(path)
|
||||
return nil, fmt.Errorf("list %s: %w", path, fetchErr)
|
||||
}
|
||||
|
||||
// Flush any remaining entries in the batch
|
||||
if len(batch) > 0 {
|
||||
if err := mc.doBatchInsertEntries(ctx, batch); err != nil {
|
||||
return nil, fmt.Errorf("batch insert remaining for %s: %w", path, err)
|
||||
}
|
||||
// Atomically replace cached entries with the snapshot and replay
|
||||
// any buffered events that arrived during the fetch.
|
||||
if err := mc.CommitRefresh(ctx, path, allEntries); err != nil {
|
||||
return nil, fmt.Errorf("commit refresh for %s: %w", path, err)
|
||||
}
|
||||
mc.markCachedFn(path)
|
||||
return nil, nil
|
||||
|
||||
@@ -0,0 +1,340 @@
|
||||
package meta_cache
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"path/filepath"
|
||||
"sync"
|
||||
"testing"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/filer"
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/util"
|
||||
)
|
||||
|
||||
func newTestMetaCache(t *testing.T) (*MetaCache, map[util.FullPath]bool) {
|
||||
t.Helper()
|
||||
uidGidMapper, err := NewUidGidMapper("", "")
|
||||
if err != nil {
|
||||
t.Fatalf("create uid/gid mapper: %v", err)
|
||||
}
|
||||
cached := make(map[util.FullPath]bool)
|
||||
mc := NewMetaCache(
|
||||
filepath.Join(t.TempDir(), "meta"),
|
||||
uidGidMapper,
|
||||
util.FullPath("/"),
|
||||
func(path util.FullPath) { cached[path] = true },
|
||||
func(path util.FullPath) bool { return cached[path] },
|
||||
func(util.FullPath, *filer_pb.Entry) {},
|
||||
nil,
|
||||
)
|
||||
t.Cleanup(func() { mc.Shutdown() })
|
||||
return mc, cached
|
||||
}
|
||||
|
||||
func makeEntry(dir, name string) *filer.Entry {
|
||||
return &filer.Entry{
|
||||
FullPath: util.NewFullPath(dir, name),
|
||||
Attr: filer.Attr{Mode: 0644},
|
||||
}
|
||||
}
|
||||
|
||||
func listEntries(t *testing.T, mc *MetaCache, dir string) []string {
|
||||
t.Helper()
|
||||
var names []string
|
||||
err := mc.ListDirectoryEntries(context.Background(), util.FullPath(dir), "", false, 10000, func(entry *filer.Entry) (bool, error) {
|
||||
names = append(names, entry.Name())
|
||||
return true, nil
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("list %s: %v", dir, err)
|
||||
}
|
||||
return names
|
||||
}
|
||||
|
||||
// TestCommitRefreshDeleteDuringRefresh verifies that a delete event arriving
|
||||
// during a directory refresh is not overwritten by the stale snapshot.
|
||||
func TestCommitRefreshDeleteDuringRefresh(t *testing.T) {
|
||||
mc, cached := newTestMetaCache(t)
|
||||
ctx := context.Background()
|
||||
dir := util.FullPath("/testdir")
|
||||
cached[dir] = true
|
||||
|
||||
// Pre-populate the cache with 3 entries
|
||||
for _, name := range []string{"a", "b", "c"} {
|
||||
if err := mc.InsertEntry(ctx, makeEntry("/testdir", name)); err != nil {
|
||||
t.Fatalf("insert %s: %v", name, err)
|
||||
}
|
||||
}
|
||||
|
||||
// Start a refresh — simulate fetching a snapshot that includes all 3
|
||||
mc.BeginRefresh(dir)
|
||||
|
||||
// While refresh is in progress, a subscription event deletes "b"
|
||||
if err := mc.AtomicUpdateEntryFromFiler(ctx, util.NewFullPath("/testdir", "b"), nil); err != nil {
|
||||
t.Fatalf("atomic delete b: %v", err)
|
||||
}
|
||||
|
||||
// Commit the snapshot (which still includes "b") — the buffered delete
|
||||
// should be replayed, removing "b" from the final result
|
||||
snapshot := []*filer.Entry{
|
||||
makeEntry("/testdir", "a"),
|
||||
makeEntry("/testdir", "b"),
|
||||
makeEntry("/testdir", "c"),
|
||||
}
|
||||
if err := mc.CommitRefresh(ctx, dir, snapshot); err != nil {
|
||||
t.Fatalf("commit refresh: %v", err)
|
||||
}
|
||||
|
||||
names := listEntries(t, mc, "/testdir")
|
||||
expected := map[string]bool{"a": true, "c": true}
|
||||
if len(names) != len(expected) {
|
||||
t.Fatalf("expected %v, got %v", expected, names)
|
||||
}
|
||||
for _, n := range names {
|
||||
if !expected[n] {
|
||||
t.Errorf("unexpected entry %q in listing", n)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// TestCommitRefreshCreateDuringRefresh verifies that a create event arriving
|
||||
// during a directory refresh is preserved after the snapshot is applied.
|
||||
func TestCommitRefreshCreateDuringRefresh(t *testing.T) {
|
||||
mc, cached := newTestMetaCache(t)
|
||||
ctx := context.Background()
|
||||
dir := util.FullPath("/testdir")
|
||||
cached[dir] = true
|
||||
|
||||
// Pre-populate with "a"
|
||||
if err := mc.InsertEntry(ctx, makeEntry("/testdir", "a")); err != nil {
|
||||
t.Fatalf("insert a: %v", err)
|
||||
}
|
||||
|
||||
// Start refresh — snapshot will have "a" only
|
||||
mc.BeginRefresh(dir)
|
||||
|
||||
// During refresh, a subscription event creates "d"
|
||||
newEntry := makeEntry("/testdir", "d")
|
||||
if err := mc.AtomicUpdateEntryFromFiler(ctx, "", newEntry); err != nil {
|
||||
t.Fatalf("atomic create d: %v", err)
|
||||
}
|
||||
|
||||
// Commit snapshot (only "a") — buffered create of "d" should be replayed
|
||||
snapshot := []*filer.Entry{makeEntry("/testdir", "a")}
|
||||
if err := mc.CommitRefresh(ctx, dir, snapshot); err != nil {
|
||||
t.Fatalf("commit refresh: %v", err)
|
||||
}
|
||||
|
||||
names := listEntries(t, mc, "/testdir")
|
||||
expected := map[string]bool{"a": true, "d": true}
|
||||
if len(names) != len(expected) {
|
||||
t.Fatalf("expected %v, got %v", expected, names)
|
||||
}
|
||||
for _, n := range names {
|
||||
if !expected[n] {
|
||||
t.Errorf("unexpected entry %q", n)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// TestCommitRefreshStaleEntriesRemoved verifies that entries deleted on the
|
||||
// filer (absent from the snapshot) are removed from the local cache.
|
||||
func TestCommitRefreshStaleEntriesRemoved(t *testing.T) {
|
||||
mc, cached := newTestMetaCache(t)
|
||||
ctx := context.Background()
|
||||
dir := util.FullPath("/testdir")
|
||||
cached[dir] = true
|
||||
|
||||
// Pre-populate with a, b, c
|
||||
for _, name := range []string{"a", "b", "c"} {
|
||||
if err := mc.InsertEntry(ctx, makeEntry("/testdir", name)); err != nil {
|
||||
t.Fatalf("insert %s: %v", name, err)
|
||||
}
|
||||
}
|
||||
|
||||
// Refresh with snapshot that only has "a" and "c" — "b" was deleted on filer
|
||||
mc.BeginRefresh(dir)
|
||||
snapshot := []*filer.Entry{
|
||||
makeEntry("/testdir", "a"),
|
||||
makeEntry("/testdir", "c"),
|
||||
}
|
||||
if err := mc.CommitRefresh(ctx, dir, snapshot); err != nil {
|
||||
t.Fatalf("commit refresh: %v", err)
|
||||
}
|
||||
|
||||
names := listEntries(t, mc, "/testdir")
|
||||
for _, n := range names {
|
||||
if n == "b" {
|
||||
t.Errorf("stale entry 'b' should have been removed by refresh")
|
||||
}
|
||||
}
|
||||
if len(names) != 2 {
|
||||
t.Fatalf("expected [a, c], got %v", names)
|
||||
}
|
||||
}
|
||||
|
||||
// TestCommitRefreshLocalDeleteBuffered verifies that a local DeleteEntry
|
||||
// during refresh is both applied immediately and replayed after commit.
|
||||
func TestCommitRefreshLocalDeleteBuffered(t *testing.T) {
|
||||
mc, cached := newTestMetaCache(t)
|
||||
ctx := context.Background()
|
||||
dir := util.FullPath("/testdir")
|
||||
cached[dir] = true
|
||||
|
||||
// Pre-populate
|
||||
for _, name := range []string{"a", "b"} {
|
||||
if err := mc.InsertEntry(ctx, makeEntry("/testdir", name)); err != nil {
|
||||
t.Fatalf("insert %s: %v", name, err)
|
||||
}
|
||||
}
|
||||
|
||||
mc.BeginRefresh(dir)
|
||||
|
||||
// Local delete of "a" — simulates Unlink calling DeleteEntry directly
|
||||
if err := mc.DeleteEntry(ctx, util.NewFullPath("/testdir", "a")); err != nil {
|
||||
t.Fatalf("delete a: %v", err)
|
||||
}
|
||||
|
||||
// Verify "a" is gone immediately (before commit)
|
||||
_, err := mc.FindEntry(ctx, util.NewFullPath("/testdir", "a"))
|
||||
if err != filer_pb.ErrNotFound {
|
||||
t.Fatalf("expected ErrNotFound for deleted entry, got: %v", err)
|
||||
}
|
||||
|
||||
// Commit with snapshot that includes "a" — the buffered delete replays
|
||||
snapshot := []*filer.Entry{
|
||||
makeEntry("/testdir", "a"),
|
||||
makeEntry("/testdir", "b"),
|
||||
}
|
||||
if err := mc.CommitRefresh(ctx, dir, snapshot); err != nil {
|
||||
t.Fatalf("commit refresh: %v", err)
|
||||
}
|
||||
|
||||
// "a" should still be gone after commit
|
||||
names := listEntries(t, mc, "/testdir")
|
||||
for _, n := range names {
|
||||
if n == "a" {
|
||||
t.Errorf("locally deleted entry 'a' reappeared after commit")
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// TestCommitRefreshLocalCreateBuffered verifies that a local InsertEntry
|
||||
// during refresh is preserved after the snapshot is applied.
|
||||
func TestCommitRefreshLocalCreateBuffered(t *testing.T) {
|
||||
mc, cached := newTestMetaCache(t)
|
||||
ctx := context.Background()
|
||||
dir := util.FullPath("/testdir")
|
||||
cached[dir] = true
|
||||
|
||||
mc.BeginRefresh(dir)
|
||||
|
||||
// Local create during refresh
|
||||
if err := mc.InsertEntry(ctx, makeEntry("/testdir", "new")); err != nil {
|
||||
t.Fatalf("insert new: %v", err)
|
||||
}
|
||||
|
||||
// Commit with empty snapshot — the buffered create should still appear
|
||||
if err := mc.CommitRefresh(ctx, dir, nil); err != nil {
|
||||
t.Fatalf("commit refresh: %v", err)
|
||||
}
|
||||
|
||||
names := listEntries(t, mc, "/testdir")
|
||||
if len(names) != 1 || names[0] != "new" {
|
||||
t.Fatalf("expected [new], got %v", names)
|
||||
}
|
||||
}
|
||||
|
||||
// TestCancelRefreshDiscardsBuffer verifies that CancelRefresh discards
|
||||
// buffered events and resumes normal event processing.
|
||||
func TestCancelRefreshDiscardsBuffer(t *testing.T) {
|
||||
mc, cached := newTestMetaCache(t)
|
||||
ctx := context.Background()
|
||||
dir := util.FullPath("/testdir")
|
||||
cached[dir] = true
|
||||
|
||||
if err := mc.InsertEntry(ctx, makeEntry("/testdir", "a")); err != nil {
|
||||
t.Fatalf("insert a: %v", err)
|
||||
}
|
||||
|
||||
mc.BeginRefresh(dir)
|
||||
|
||||
// Buffer a delete event
|
||||
if err := mc.AtomicUpdateEntryFromFiler(ctx, util.NewFullPath("/testdir", "a"), nil); err != nil {
|
||||
t.Fatalf("atomic delete: %v", err)
|
||||
}
|
||||
|
||||
// Cancel — buffered events are discarded
|
||||
mc.CancelRefresh(dir)
|
||||
|
||||
// "a" should still be in cache since the buffered delete was discarded
|
||||
entry, err := mc.FindEntry(ctx, util.NewFullPath("/testdir", "a"))
|
||||
if err != nil {
|
||||
t.Fatalf("find a: %v", err)
|
||||
}
|
||||
if entry == nil {
|
||||
t.Fatal("entry 'a' should still exist after cancel")
|
||||
}
|
||||
|
||||
// After cancel, events should apply normally (not buffered)
|
||||
if err := mc.AtomicUpdateEntryFromFiler(ctx, util.NewFullPath("/testdir", "a"), nil); err != nil {
|
||||
t.Fatalf("atomic delete after cancel: %v", err)
|
||||
}
|
||||
_, err = mc.FindEntry(ctx, util.NewFullPath("/testdir", "a"))
|
||||
if err != filer_pb.ErrNotFound {
|
||||
t.Fatalf("expected ErrNotFound after direct delete, got: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
// TestConcurrentDeletesDuringRefresh simulates the scenario from issue #8442:
|
||||
// multiple concurrent deletes racing with a directory refresh.
|
||||
func TestConcurrentDeletesDuringRefresh(t *testing.T) {
|
||||
mc, cached := newTestMetaCache(t)
|
||||
ctx := context.Background()
|
||||
dir := util.FullPath("/testdir")
|
||||
cached[dir] = true
|
||||
|
||||
const numFiles = 1000
|
||||
|
||||
// Pre-populate
|
||||
for i := 0; i < numFiles; i++ {
|
||||
name := fmt.Sprintf("file_%04d", i)
|
||||
if err := mc.InsertEntry(ctx, makeEntry("/testdir", name)); err != nil {
|
||||
t.Fatalf("insert %s: %v", name, err)
|
||||
}
|
||||
}
|
||||
|
||||
// Build snapshot (taken before deletes happen)
|
||||
var snapshot []*filer.Entry
|
||||
for i := 0; i < numFiles; i++ {
|
||||
snapshot = append(snapshot, makeEntry("/testdir", fmt.Sprintf("file_%04d", i)))
|
||||
}
|
||||
|
||||
mc.BeginRefresh(dir)
|
||||
|
||||
// Concurrently delete all files via subscription events (simulates
|
||||
// deletes from two mount nodes)
|
||||
var wg sync.WaitGroup
|
||||
for i := 0; i < numFiles; i++ {
|
||||
wg.Add(1)
|
||||
go func(idx int) {
|
||||
defer wg.Done()
|
||||
name := fmt.Sprintf("file_%04d", idx)
|
||||
fp := util.NewFullPath("/testdir", name)
|
||||
mc.AtomicUpdateEntryFromFiler(ctx, fp, nil)
|
||||
}(i)
|
||||
}
|
||||
wg.Wait()
|
||||
|
||||
// Commit with the stale snapshot — all deletes should be replayed
|
||||
if err := mc.CommitRefresh(ctx, dir, snapshot); err != nil {
|
||||
t.Fatalf("commit refresh: %v", err)
|
||||
}
|
||||
|
||||
names := listEntries(t, mc, "/testdir")
|
||||
if len(names) != 0 {
|
||||
t.Errorf("expected 0 entries after all deletes, got %d: first few: %v",
|
||||
len(names), names[:min(5, len(names))])
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user