diff --git a/lib/store/disk/config.go b/lib/store/disk/config.go index d021c6c74..238c39fb0 100644 --- a/lib/store/disk/config.go +++ b/lib/store/disk/config.go @@ -1,16 +1,33 @@ package disk +import "errors" + +const _defaultRootDir = "disk_store" + // Config configures [Store]. type Config struct { // The capacity of the store in bytes. When breached, LRU eviction is used. - CapacityBytes uint64 + CapacityBytes uint64 `yaml:"capacity_bytes"` // The root directory under which [Store]'s blobs and state are stored. - RootDir string + RootDir string `yaml:"root_dir"` // Whether after crash/restart, the Store removes incomplete files from disk (usually to prevent leaks) OR // reboots any incomplete files from disk (allowing users to continue the blob download/upload, where it was left off before the crash). - RebootIncompleteBlobs bool + RebootIncompleteBlobs bool `yaml:"reboot_incomplete_blobs"` // If > 0, directory sharding is used to speed up performance, where ShardLength denotes // 1) the length of each directory shard's name and 2) the number of shards. // A value of 0 denotes no sharding. - ShardLength int // TODO - add a mechanism to migrate in-place from 1 shardLength to another. + ShardLength int `yaml:"shard_length"` +} + +func (c *Config) applyDefaults() error { + if c.CapacityBytes == 0 { + return errors.New("capacity_bytes must be explicitly set") + } + if c.ShardLength < 0 { + return errors.New("shard_length must be non-negative") + } + if c.RootDir == "" { + c.RootDir = _defaultRootDir + } + return nil } diff --git a/lib/store/disk/crash_recovery.go b/lib/store/disk/crash_recovery.go index 04e908a62..c104effc2 100644 --- a/lib/store/disk/crash_recovery.go +++ b/lib/store/disk/crash_recovery.go @@ -108,7 +108,7 @@ func rebootPersistedStore(config *Config, log *zap.SugaredLogger, metrics tally. if store.size > store.capacity { prevSize := store.size // evicts blobs until size <= capacity. - err = store.reserveSpace(0) + err = store.ensureFreeSpace(0) if err != nil { log.With("error", err).Error("Store size exceeds its capacity after service reboot. Evicting blobs from disk did not work to reduce size within capacity.") return nil, fmt.Errorf("remove blobs to reduce store size within configured capacity: %w", err) diff --git a/lib/store/disk/crash_recovery_test.go b/lib/store/disk/crash_recovery_test.go index 4f5fd288b..39365d227 100644 --- a/lib/store/disk/crash_recovery_test.go +++ b/lib/store/disk/crash_recovery_test.go @@ -258,7 +258,7 @@ func TestCrashRecovery(t *testing.T) { // Assume that the application restarts here. _, err := NewStore(&Config{5 * memsize.KB, rootDir, true, _defaultShardLength}, tally.NoopScope) - require.ErrorContains(err, "cannot evict enough, the unevictable/incomplete blobs are using up all the space") + require.ErrorIs(err, errNoSpace) }) } @@ -333,7 +333,7 @@ func TestStoreWorksWhenFileSizeNotCorrect(t *testing.T) { // Even though only 1KB is actually used on disk, the store enforces capacity based on the // declared 8KB, so a 3KB blob doesn't fit alongside it (there's nothing evictable to make room). _, err := store.Create(core.DigestFixture().Hex(), 3*memsize.KB) - require.EqualError(err, "reserve space: cannot evict enough, the unevictable/incomplete blobs are using up all the space") + require.ErrorIs(err, errNoSpace) // Declares 2KB (fits exactly within the remaining capacity) but writes 5KB. overreportedF, overreportedKey := newTestFile(t, store, 2*memsize.KB) diff --git a/lib/store/disk/scoped_store.go b/lib/store/disk/scoped_store.go index d137d4004..db17a1d6a 100644 --- a/lib/store/disk/scoped_store.go +++ b/lib/store/disk/scoped_store.go @@ -68,7 +68,7 @@ func (s *Store) MarkComplete(key string) error { return s.impl.MarkComplete(key) func (s *Store) Delete(key string) error { return s.impl.Delete(key, s.scope) } // List returns the keys of all blobs (except those out of scope). -func (s *Store) List() []string { return s.impl.list(s.scope) } +func (s *Store) List() []string { return s.impl.List(s.scope) } // BanEviction marks a blob as unevictable by LRU eviction. It is idempotent. // Needed when e.g. blobs must be written back to GCS/S3 and eviction before that is unacceptable. @@ -88,8 +88,8 @@ func (s *Store) GetMetadata(key string, md metadata.Metadata) (ok bool, err erro } // DeleteMetadata removes a blob's metadata. No error returned if the metadata is not present. -func (s *Store) DeleteMetadata(key string, md metadata.Metadata) error { - return s.impl.DeleteMetadata(key, md, s.scope) +func (s *Store) DeleteMetadata(key string, mdSuffix string) error { + return s.impl.DeleteMetadata(key, mdSuffix, s.scope) } // ListMetadata returns all [metadata.Metadata] of key. @@ -109,3 +109,17 @@ func (s *Store) ScopeComplete() *Store { return &Store{s.impl, storelib.BlobScop // ScopeIncomplete scopes [Store]'s APIs such that they can only operate on incomplete blobs. // [storelib.ErrOutOfScope] is returned if the user tries to operate on a complete blob. func (s *Store) ScopeIncomplete() *Store { return &Store{s.impl, storelib.BlobScopeIncomplete} } + +// Scoped scopes [Store]'s APIs, such that [storelib.ErrOutOfScope] is returned upon attempting to operate on blobs out of scope. +func (s *Store) Scoped(scope storelib.BlobScope) *Store { return &Store{s.impl, scope} } + +// Clean deletes blobs from the store until targetUtilPercent is hit. +// Intended for manual use for incident mitigation/benchmarking. +// Blobs are deleted in the following order: +// 1. complete, non-eviction-banned blobs, in LRU order. +// 2. incomplete, non-eviction-banned blobs, in a random order. +// 3. If respectEvictionBan is true, Clean will not delete blobs banned from eviction. +// Else, it will delete them, in a random order. +func (s *Store) Clean(targetUtilPercent int, respectEvictionBan bool) (newUtil int, err error) { + return s.impl.Clean(targetUtilPercent, respectEvictionBan) +} diff --git a/lib/store/disk/store.go b/lib/store/disk/store.go index 5ecccfa26..72cbdc5fc 100644 --- a/lib/store/disk/store.go +++ b/lib/store/disk/store.go @@ -28,6 +28,7 @@ const ( ) var _syncEvictionLatencyBuckets = tally.MustMakeExponentialDurationBuckets(100*time.Millisecond, 1.4, 15) +var errNoSpace = errors.New("cannot free enough space for new entry, the unevictable/incomplete blobs are using up all the space") // store implements the APIs of [Store]. [store]'s APIs expose the [storelib.BlobScope] arg, // while [Store]'s APIs omit that arg (cleaner interface) and instead expose other APIs to scope the whole store. @@ -54,6 +55,11 @@ type blob struct { } func newStore(config *Config, metrics tally.Scope) (*store, error) { + err := config.applyDefaults() + if err != nil { + return nil, err + } + log := log.Default().With("module", "disk_store") ok, err := existsPersistedStore(config.RootDir) if err != nil { @@ -62,7 +68,7 @@ func newStore(config *Config, metrics tally.Scope) (*store, error) { return nil, err } if !ok { - log.Info("Initialized a new, empty Store (did not find any previously persisted state to reboot for Store)") + log.Info("Initialized a new, empty disk.Store (did not find any previously persisted state to reboot for disk.Store)") store := &store{ capacity: config.CapacityBytes, size: 0, @@ -128,6 +134,8 @@ func (s *store) Has(key string, scope storelib.BlobScope) (inStore bool, inScope } func (s *store) Stat(key string, scope storelib.BlobScope) (os.FileInfo, error) { + // TODO - make stat return `size` and not os.FileInfo to keep its interface consistent with memory.Store and tiered.Store. + // TODO - consider whether Stat (and maybe other APIs) should update the access time of the blob. s.mu.RLock() defer s.mu.RUnlock() @@ -138,7 +146,6 @@ func (s *store) Stat(key string, scope storelib.BlobScope) (os.FileInfo, error) if err := isOutOfScope(b, scope); err != nil { return nil, err } - // TODO - consider whether this should update the access time of the blob. blobPath := s.blobPath(key, b.complete) return os.Stat(blobPath) } @@ -148,18 +155,14 @@ func (s *store) Create(key string, sizeBytes uint64) (*File, error) { s.mu.Lock() defer s.mu.Unlock() - if b, ok := s.blobs[key]; ok { - // TODO - consider whether we need public errors for these cases. - if b.complete { - return nil, errors.New("blob is already in store") - } else { - return nil, errors.New("blob is already in store (it is incomplete)") - } + if _, ok := s.blobs[key]; ok { + return nil, os.ErrExist } - if err := s.reserveSpace(sizeBytes); err != nil { - return nil, fmt.Errorf("reserve space: %w", err) + if err := s.ensureFreeSpace(sizeBytes); err != nil { + return nil, fmt.Errorf("ensure free space: %w", err) } + s.size += sizeBytes dirName := s.dirPath(key, _incompleteBlob) err := os.MkdirAll(dirName, _defaultFilePerm) @@ -207,7 +210,11 @@ func (s *store) persistBlobSize(key string, sizeBytes uint64) error { return nil } -func (s *store) reserveSpace(space uint64) error { +func (s *store) ensureFreeSpace(space uint64) error { + if s.size+space <= s.capacity { + return nil + } + // TODO - benchmark and consider whether async eviction makes more sense. startTime := time.Now() for s.size+space > s.capacity { @@ -217,7 +224,7 @@ func (s *store) reserveSpace(space uint64) error { "required_space", space, "capacity", s.capacity, ).Error("Cannot evict enough data to free space for new entry to Store. The unevictable/incomplete blobs are using up all the space") - return errors.New("cannot evict enough, the unevictable/incomplete blobs are using up all the space") + return errNoSpace } toEvictNode := s.evictQueue.Front() @@ -232,7 +239,6 @@ func (s *store) reserveSpace(space uint64) error { s.releaseSpace(size) delete(s.blobs, toEvictKey) } - s.size += space latency := time.Since(startTime) s.metrics.Histogram("sync_eviction_latency", _syncEvictionLatencyBuckets).RecordDuration(latency) @@ -324,6 +330,10 @@ func (s *store) Delete(key string, scope storelib.BlobScope) error { s.mu.Lock() defer s.mu.Unlock() + return s.deleteNoLock(key, scope) +} + +func (s *store) deleteNoLock(key string, scope storelib.BlobScope) error { b, ok := s.blobs[key] if !ok { return os.ErrNotExist @@ -346,7 +356,7 @@ func (s *store) Delete(key string, scope storelib.BlobScope) error { return nil } -func (s *store) list(scope storelib.BlobScope) []string { +func (s *store) List(scope storelib.BlobScope) []string { s.mu.RLock() defer s.mu.RUnlock() @@ -489,8 +499,7 @@ func (s *store) GetMetadata(key string, md metadata.Metadata, scope storelib.Blo return true, nil } -func (s *store) DeleteMetadata(key string, md metadata.Metadata, scope storelib.BlobScope) error { - // TODO - change interface to take `mdSuffix string` instead of `md metadata.Metadata`, just like memory.Store.` +func (s *store) DeleteMetadata(key, mdSuffix string, scope storelib.BlobScope) error { s.mu.Lock() defer s.mu.Unlock() @@ -501,7 +510,7 @@ func (s *store) DeleteMetadata(key string, md metadata.Metadata, scope storelib. if err := isOutOfScope(b, scope); err != nil { return err } - mdFilePath := s.sidecarFilePath(key, b.complete, md.GetSuffix()) + mdFilePath := s.sidecarFilePath(key, b.complete, mdSuffix) err := os.Remove(mdFilePath) if errors.Is(err, os.ErrNotExist) { // no-op @@ -580,6 +589,66 @@ func (s *store) WriteAtMetadata(key string, md metadata.Metadata, p []byte, off return nil } +func (s *store) Clean(targetUtilPercent int, respectEvictionBan bool) (newUtil int, err error) { + defer func() { + newUtil = int(s.size * 100 / s.capacity) + }() + + if targetUtilPercent < 0 || targetUtilPercent >= 100 { + return 0, errors.New("target_util_percent must be >=0 and <100") + } + s.log.Warn("Starting manual cleanup of disk.Store") + + s.mu.Lock() + defer s.mu.Unlock() + + targetSize := s.capacity * uint64(targetUtilPercent) / 100 + requiredFreeSpace := s.capacity - targetSize + + err = s.ensureFreeSpace(requiredFreeSpace) + if err != nil && !errors.Is(err, errNoSpace) { + return newUtil, fmt.Errorf("ensure free space: %w", err) + } + if err == nil { + return newUtil, err + } + + notEvictionBannedKeys := []string{} + for key, b := range s.blobs { + if !b.evictionBanned { + notEvictionBannedKeys = append(notEvictionBannedKeys, key) + } + } + + for _, key := range notEvictionBannedKeys { + if s.size <= targetSize { + return newUtil, nil + } + + err = s.deleteNoLock(key, storelib.BlobScopeAny) + if err != nil { + return newUtil, fmt.Errorf("delete incomplete, non-eviction-banned blob: %w", err) + } + } + + if respectEvictionBan { + return newUtil, nil + } + + for key := range s.blobs { + if s.size <= targetSize { + return newUtil, nil + } + + err = s.deleteNoLock(key, storelib.BlobScopeAny) + if err != nil { + return newUtil, fmt.Errorf("delete blob banned from eviction: %w", err) + } + } + + return newUtil, nil +} + // used during testing func (s *store) evictionOrder() []string { s.mu.RLock() diff --git a/lib/store/disk/store_test.go b/lib/store/disk/store_test.go index 4e48e00c7..22c4b8627 100644 --- a/lib/store/disk/store_test.go +++ b/lib/store/disk/store_test.go @@ -92,7 +92,7 @@ func TestStore(t *testing.T) { require.Equal(10*memsize.KB, store.impl.size) f, err := store.Create(core.DigestFixture().Hex(), memsize.B) - require.EqualError(err, "reserve space: cannot evict enough, the unevictable/incomplete blobs are using up all the space") + require.ErrorIs(err, errNoSpace) require.Nil(f) f, err = store.ScopeComplete().Open(keys[0]) @@ -137,7 +137,7 @@ func TestEviction(t *testing.T) { require.Equal(5, numBlobsOnDisk(t, store)) // incomplete files cannot be evicted and adding 2KB would result in overreservation. _, err := store.Create(core.DigestFixture().Hex(), 2*memsize.KB) - require.EqualError(err, "reserve space: cannot evict enough, the unevictable/incomplete blobs are using up all the space") + require.ErrorIs(err, errNoSpace) // start marking as complete in this specific order - c, b, a (MarkComplete resets access time). require.NoError(store.MarkComplete(cKey)) require.NoError(store.MarkComplete(bKey)) @@ -764,7 +764,7 @@ func TestMetadata(t *testing.T) { []string{mdList[0].GetSuffix(), mdList[1].GetSuffix()}, ) - require.NoError(store.DeleteMetadata(key, &readMd)) + require.NoError(store.DeleteMetadata(key, readMd.GetSuffix())) ok, err = store.GetMetadata(key, &readMd) require.NoError(err) require.False(ok) @@ -773,14 +773,14 @@ func TestMetadata(t *testing.T) { _, err = os.Stat(mdFilePath) require.ErrorIs(err, os.ErrNotExist) // deleting a second time should be a no-op. - require.NoError(store.DeleteMetadata(key, &readMd)) + require.NoError(store.DeleteMetadata(key, readMd.GetSuffix())) mdList, err = store.ListMetadata(key) require.NoError(err) require.Len(mdList, 1) require.Equal(persistMd.GetSuffix(), mdList[0].GetSuffix()) - require.NoError(store.DeleteMetadata(key, persistMd)) + require.NoError(store.DeleteMetadata(key, persistMd.GetSuffix())) mdList, err = store.ListMetadata(key) require.NoError(err) require.Empty(mdList) @@ -806,7 +806,7 @@ func TestMetadata(t *testing.T) { err = store.WriteAtMetadata(nonExistentKey, md, []byte("data"), 0) require.ErrorIs(os.ErrNotExist, err) - err = store.DeleteMetadata(nonExistentKey, md) + err = store.DeleteMetadata(nonExistentKey, md.GetSuffix()) require.ErrorIs(os.ErrNotExist, err) }) @@ -980,7 +980,7 @@ func TestScopes(t *testing.T) { _, err = store.ScopeComplete().ListMetadata(key) require.ErrorIs(err, storelib.ErrOutOfScope) require.ErrorIs(store.ScopeComplete().WriteAtMetadata(key, md, mdData, 0), storelib.ErrOutOfScope) - require.ErrorIs(store.ScopeComplete().DeleteMetadata(key, &readMd), storelib.ErrOutOfScope) + require.ErrorIs(store.ScopeComplete().DeleteMetadata(key, readMd.GetSuffix()), storelib.ErrOutOfScope) require.NotContains(store.ScopeComplete().List(), key) _, ok = store.ScopeComplete().Has(key) require.False(ok) @@ -1006,7 +1006,7 @@ func TestScopes(t *testing.T) { require.Len(mdList, 1) require.Equal(md.GetSuffix(), mdList[0].GetSuffix()) require.NoError(store.WriteAtMetadata(key, md, mdData, 0)) - require.NoError(store.DeleteMetadata(key, &readMd)) + require.NoError(store.DeleteMetadata(key, readMd.GetSuffix())) require.Contains(store.List(), key) _, ok = store.Has(key) require.True(ok) @@ -1032,7 +1032,7 @@ func TestScopes(t *testing.T) { require.Len(mdList, 1) require.Equal(md.GetSuffix(), mdList[0].GetSuffix()) require.NoError(store.ScopeIncomplete().WriteAtMetadata(key, md, mdData, 0)) - require.NoError(store.ScopeIncomplete().DeleteMetadata(key, &readMd)) + require.NoError(store.ScopeIncomplete().DeleteMetadata(key, readMd.GetSuffix())) require.Contains(store.ScopeIncomplete().List(), key) _, ok = store.ScopeIncomplete().Has(key) require.True(ok) @@ -1056,7 +1056,7 @@ func TestScopes(t *testing.T) { _, err = store.ScopeIncomplete().ListMetadata(key) require.ErrorIs(err, storelib.ErrOutOfScope) require.ErrorIs(store.ScopeIncomplete().WriteAtMetadata(key, md, mdData, 0), storelib.ErrOutOfScope) - require.ErrorIs(store.ScopeIncomplete().DeleteMetadata(key, &readMd), storelib.ErrOutOfScope) + require.ErrorIs(store.ScopeIncomplete().DeleteMetadata(key, readMd.GetSuffix()), storelib.ErrOutOfScope) require.NotContains(store.ScopeIncomplete().List(), key) _, ok = store.ScopeIncomplete().Has(key) require.False(ok) @@ -1082,7 +1082,7 @@ func TestScopes(t *testing.T) { require.Len(mdList, 1) require.Equal(md.GetSuffix(), mdList[0].GetSuffix()) require.NoError(store.WriteAtMetadata(key, md, mdData, 0)) - require.NoError(store.DeleteMetadata(key, &readMd)) + require.NoError(store.DeleteMetadata(key, readMd.GetSuffix())) require.Contains(store.List(), key) _, ok = store.Has(key) require.True(ok) @@ -1108,7 +1108,7 @@ func TestScopes(t *testing.T) { require.Len(mdList, 1) require.Equal(md.GetSuffix(), mdList[0].GetSuffix()) require.NoError(store.ScopeComplete().WriteAtMetadata(key, md, mdData, 0)) - require.NoError(store.ScopeComplete().DeleteMetadata(key, &readMd)) + require.NoError(store.ScopeComplete().DeleteMetadata(key, readMd.GetSuffix())) require.Contains(store.ScopeComplete().List(), key) _, ok = store.ScopeComplete().Has(key) require.True(ok) @@ -1146,3 +1146,138 @@ func TestScopesDelete(t *testing.T) { _, err = store.Stat(key) require.ErrorIs(err, os.ErrNotExist) } + +func TestClean(t *testing.T) { + t.Run("invalid target_util_percent", func(t *testing.T) { + require := require.New(t) + store, _ := newTestStore(t, 100*memsize.KB, false) + + _, err := store.Clean(-1, true) + require.EqualError(err, "target_util_percent must be >=0 and <100") + + _, err = store.Clean(100, true) + require.EqualError(err, "target_util_percent must be >=0 and <100") + }) + + t.Run("no-op when already under target", func(t *testing.T) { + require := require.New(t) + store, _ := newTestStore(t, 100*memsize.KB, false) + f, key := newTestFile(t, store, 20*memsize.KB) + require.NoError(f.Close()) + require.NoError(store.MarkComplete(key)) + + newUtil, err := store.Clean(50, true) + require.NoError(err) + require.Equal(20, newUtil) + require.Equal([]string{key}, store.List()) + }) + + t.Run("LRU eviction alone reaches target", func(t *testing.T) { + require := require.New(t) + store, _ := newTestStore(t, 100*memsize.KB, false) + a, aKey := newTestFile(t, store, 60*memsize.KB) + require.NoError(a.Close()) + require.NoError(store.MarkComplete(aKey)) + b, bKey := newTestFile(t, store, 20*memsize.KB) + require.NoError(b.Close()) + require.NoError(store.MarkComplete(bKey)) + + // LRU-evicting a (60KB) alone is enough to get to <=50% util, so b is spared. + newUtil, err := store.Clean(50, true) + require.NoError(err) + require.Equal(20, newUtil) + require.Equal([]string{bKey}, store.List()) + }) + + t.Run("respectEvictionBan true - deletes incomplete blobs but spares banned ones", func(t *testing.T) { + require := require.New(t) + store, _ := newTestStore(t, 100*memsize.KB, false) + a, aKey := newTestFile(t, store, 30*memsize.KB) + require.NoError(a.Close()) + require.NoError(store.MarkComplete(aKey)) + require.NoError(store.BanEviction(aKey)) + b, _ := newTestFile(t, store, 20*memsize.KB) // incomplete, non-banned. + require.NoError(b.Close()) + c, cKey := newTestFile(t, store, 10*memsize.KB) + require.NoError(c.Close()) + require.NoError(store.MarkComplete(cKey)) + + // Target of 10% cannot be reached without deleting the banned blob a, + // so with respectEvictionBan=true only b (LRU) and c (incomplete) are removed. + newUtil, err := store.Clean(10, true) + require.NoError(err) + require.Equal(30, newUtil) + require.Equal([]string{aKey}, store.List()) + }) + + t.Run("respectEvictionBan true - target not reached, but no error returned", func(t *testing.T) { + require := require.New(t) + store, _ := newTestStore(t, 100*memsize.KB, false) + a, aKey := newTestFile(t, store, 100*memsize.KB) + require.NoError(a.Close()) + require.NoError(store.MarkComplete(aKey)) + require.NoError(store.BanEviction(aKey)) + + // The only blob is banned from eviction, so respectEvictionBan=true leaves + // it untouched even though the target of 0% is nowhere close to being met. + targetUtilPercent := 0 + newUtil, err := store.Clean(targetUtilPercent, true) + require.NoError(err) + require.Greater(newUtil, targetUtilPercent) + require.Equal(100, newUtil) + require.Equal([]string{aKey}, store.List()) + }) + + t.Run("respectEvictionBan false - also deletes banned blobs to reach target", func(t *testing.T) { + require := require.New(t) + store, _ := newTestStore(t, 100*memsize.KB, false) + a, aKey := newTestFile(t, store, 30*memsize.KB) + require.NoError(a.Close()) + require.NoError(store.MarkComplete(aKey)) + require.NoError(store.BanEviction(aKey)) + b, _ := newTestFile(t, store, 20*memsize.KB) // incomplete, non-banned. + require.NoError(b.Close()) + c, cKey := newTestFile(t, store, 10*memsize.KB) + require.NoError(c.Close()) + require.NoError(store.MarkComplete(cKey)) + + newUtil, err := store.Clean(10, false) + require.NoError(err) + require.Equal(0, newUtil) + require.Empty(store.List()) + }) +} + +func TestConfig__applyDefaults(t *testing.T) { + for name, tt := range map[string]struct { + config *Config + wantErr string + }{ + "capacity must be explicitly set": { + config: &Config{ + RootDir: t.TempDir(), + RebootIncompleteBlobs: false, + ShardLength: _defaultShardLength, + }, + wantErr: "capacity_bytes must be explicitly set", + }, + "shard length cannot be negative": { + config: &Config{ + CapacityBytes: 10 * memsize.KB, + RootDir: t.TempDir(), + RebootIncompleteBlobs: false, + ShardLength: -1, + }, + wantErr: "shard_length must be non-negative", + }, + } { + t.Run(name, func(t *testing.T) { + _, err := NewStore(tt.config, tally.NoopScope) + if tt.wantErr != "" { + require.EqualError(t, err, tt.wantErr) + } else { + require.NoError(t, err) + } + }) + } +} diff --git a/lib/store/memory/config.go b/lib/store/memory/config.go new file mode 100644 index 000000000..5358e1d9e --- /dev/null +++ b/lib/store/memory/config.go @@ -0,0 +1,47 @@ +package memory + +import ( + "errors" + "math" + "runtime/debug" +) + +// Config configures [Store]. +type Config struct { + // GOMEMLIMITBytes and GOGC tune the Go GC. In short, it is advised to set GOMEMLIMIT to 90-95% of the container's reserved memory + // and to turn GOGC off, assuming the container has reserved memory. More info on tuning the Go GC: https://go.dev/doc/gc-guide. + // + // If set, they overwrite any previous configurations for GOMEMLIMIT and GOGC. + // GOMEMLIMIT *must* be configured either through [Config] or another way (e.g. an env var). + // GOGC defaults to the Go default (100), if not set and is turned off if a negative value is provided. + GOMEMLIMITBytes int64 `yaml:"gomemlimit_bytes"` + // Check the comment at [Config.GOMEMLIMITBytes]. + GOGC int `yaml:"gogc"` + // CapacityBytes sets a *hard* limit on the memory [Store] can keep reachable (i.e. unreclaimable by the GC) to store blobs. + // Any other memory that [Store] uses (e.g. storing metadata) is NOT capped by CapacityBytes. + CapacityBytes uint64 `yaml:"capacity_bytes"` +} + +func (c *Config) applyDefaults() error { + if c.CapacityBytes <= 0 { + return errors.New("capacity_bytes must be explicitly set") + } + return nil +} + +func (c *Config) configureGC() error { + if c.GOMEMLIMITBytes != 0 { + debug.SetMemoryLimit(c.GOMEMLIMITBytes) + } else { + currentLimit := debug.SetMemoryLimit(-1) + notSet := currentLimit == math.MaxInt64 + if notSet { + return errors.New("GOMEMLIMIT must be configured either through memory.Config or an env var") + } + } + + if c.GOGC != 0 { + debug.SetGCPercent(c.GOGC) + } + return nil +} diff --git a/lib/store/memory/file.go b/lib/store/memory/file.go index 7fcfb738b..560bee9f1 100644 --- a/lib/store/memory/file.go +++ b/lib/store/memory/file.go @@ -103,11 +103,11 @@ func (f *File) Seek(off int64, whence int) (int64, error) { } // Stat returns the blob's actual size, even if it differs from the size reported during Create. +// A return value of -1 represents [ErrEvicted]. func (f *File) Size() int64 { buf, evicted := f.getData() if evicted { - // TODO - consider whether this is ok or whether we need to store the user-provided blob size in [*File]. - return 0 + return -1 } return int64(len(buf)) } diff --git a/lib/store/memory/scoped_store.go b/lib/store/memory/scoped_store.go index 68edb11d4..8a56b7520 100644 --- a/lib/store/memory/scoped_store.go +++ b/lib/store/memory/scoped_store.go @@ -11,8 +11,9 @@ import ( // ErrNoSpace means the store could not free enough space for a new entry. var ErrNoSpace error = errors.New("cannot free enough memory for new entry") -// ErrEvicted is returned when a user tries operating on a blob's [*File] after the blob has been evicted. -var ErrEvicted error = errors.New("the blob has been evicted from the store") +// ErrEvicted is returned when a user tries operating on a blob's [*File] after the blob has been evicted or deleted. +var ErrEvicted error = errors.New("the blob has been evicted or deleted from the store") // TODO - rename to ErrGone +// or something else that shows the entry might be deleted too, not just evicted. // Store is an in-memory, thread-safe, LRU cache for blobs and their [metadata.Metadata]. // @@ -31,9 +32,9 @@ type Store struct { scope storelib.BlobScope } -// NewStore initializes a new, empty [*Store] with the given capacity. -func NewStore(capacityBytes uint64, metrics tally.Scope) (*Store, error) { - s, err := newStore(capacityBytes, metrics) +// NewStore initializes a new, empty [*Store]. +func NewStore(config *Config, metrics tally.Scope) (*Store, error) { + s, err := newStore(config, metrics) if err != nil { return nil, err } @@ -67,7 +68,7 @@ func (s *Store) MarkComplete(key string) error { return s.impl.MarkComplete(key) func (s *Store) Delete(key string) error { return s.impl.Delete(key, s.scope) } // List returns the keys of all blobs (except those out of scope). -func (s *Store) List() []string { return s.impl.list(s.scope) } +func (s *Store) List() []string { return s.impl.List(s.scope) } // BanEviction marks a blob as unevictable by LRU eviction. It is idempotent. // Usually used by clients to ensure a blob is not evicted before being flushed to disk. @@ -103,3 +104,6 @@ func (s *Store) ScopeComplete() *Store { return &Store{s.impl, storelib.BlobScop // ScopeIncomplete scopes [Store]'s APIs such that they can only operate on incomplete blobs. // [storelib.ErrOutOfScope] is returned if the user tries to operate on a complete blob. func (s *Store) ScopeIncomplete() *Store { return &Store{s.impl, storelib.BlobScopeIncomplete} } + +// Scoped scopes [Store]'s APIs, such that [storelib.ErrOutOfScope] is returned upon attempting to operate on blobs out of scope. +func (s *Store) Scoped(scope storelib.BlobScope) *Store { return &Store{s.impl, scope} } diff --git a/lib/store/memory/store.go b/lib/store/memory/store.go index cc1935176..96180981c 100644 --- a/lib/store/memory/store.go +++ b/lib/store/memory/store.go @@ -2,7 +2,6 @@ package memory import ( "container/list" - "errors" "fmt" "maps" "os" @@ -41,18 +40,22 @@ type blob struct { sliceMu sync.Mutex // Synchronizes writes to 1) the atomic pointer and 2) the slice (not the array!). } -func newStore(capacityBytes uint64, metrics tally.Scope) (*store, error) { - if capacityBytes <= 0 { - return nil, errors.New("store capacity must be positive") +func newStore(config *Config, metrics tally.Scope) (*store, error) { + err := config.applyDefaults() + if err != nil { + return nil, err + } + err = config.configureGC() + if err != nil { + return nil, err } log := log.Default().With("module", "memory_store") - - log.Info("Initialized new, empty *memory.Store") + log.Info("Initialized new, empty memory.Store") s := &store{ blobs: make(map[string]*blob, 0), evictQueue: list.New(), - capacity: capacityBytes, + capacity: config.CapacityBytes, size: 0, log: log, metrics: metrics, @@ -175,7 +178,7 @@ func (s *store) Delete(key string, scope storelib.BlobScope) error { return nil } -func (s *store) list(scope storelib.BlobScope) []string { +func (s *store) List(scope storelib.BlobScope) []string { s.mu.RLock() defer s.mu.RUnlock() diff --git a/lib/store/memory/store_test.go b/lib/store/memory/store_test.go index 976f2cd64..78a45f002 100644 --- a/lib/store/memory/store_test.go +++ b/lib/store/memory/store_test.go @@ -4,8 +4,10 @@ import ( "bytes" "crypto/rand" "io" + "math" "os" "regexp" + "runtime/debug" "sync" "testing" @@ -17,6 +19,16 @@ import ( "github.com/uber/kraken/utils/memsize" ) +func newTestStore(t *testing.T, capacityBytes uint64) *Store { + config := &Config{ + GOMEMLIMITBytes: math.MaxInt64, + CapacityBytes: capacityBytes, + } + store, err := NewStore(config, tally.NoopScope) + require.NoError(t, err) + return store +} + func newTestFile(t *testing.T, store *Store, size uint64) (f storelib.FileReadWriter, key string) { require := require.New(t) key = core.DigestFixture().Hex() @@ -27,8 +39,7 @@ func newTestFile(t *testing.T, store *Store, size uint64) (f storelib.FileReadWr func TestEviction(t *testing.T) { require := require.New(t) - store, err := NewStore(25*memsize.KB, tally.NoopScope) - require.NoError(err) + store := newTestStore(t, 25*memsize.KB) // create 5 blobs - a, b, c, d, e with different sizes. a, aKey := newTestFile(t, store, 10*memsize.KB) b, bKey := newTestFile(t, store, 5*memsize.KB) @@ -38,7 +49,7 @@ func TestEviction(t *testing.T) { require.Equal(24*memsize.KB, store.impl.size) // incomplete files cannot be evicted and adding 2KB would result in overreservation. - _, err = store.Create(core.DigestFixture().Hex(), 2*memsize.KB) + _, err := store.Create(core.DigestFixture().Hex(), 2*memsize.KB) require.ErrorIs(err, ErrNoSpace) // start marking as complete in this specific order - c, b, a (MarkComplete resets access time). require.NoError(store.MarkComplete(cKey)) @@ -138,8 +149,7 @@ func TestEviction(t *testing.T) { func TestEvictedBlobImmediatelyNotAvailable(t *testing.T) { require := require.New(t) - store, err := NewStore(1*memsize.KB, tally.NoopScope) - require.NoError(err) + store := newTestStore(t, 1*memsize.KB) b, key := newTestFile(t, store, 1*memsize.KB) require.NoError(store.MarkComplete(key)) @@ -172,13 +182,12 @@ func TestEvictedBlobImmediatelyNotAvailable(t *testing.T) { require.NoError(b.Close()) require.NoError(b.Commit()) - require.Equal(int64(0), b.Size()) + require.Equal(int64(-1), b.Size()) } func TestParallelAccessToSingleFile(t *testing.T) { require := require.New(t) - store, err := NewStore(10*memsize.KB, tally.NoopScope) - require.NoError(err) + store := newTestStore(t, 10*memsize.KB) key := core.DigestFixture().Hex() // we purposefully underreport the size of the blob (real size is 50B) to test if the resizing logic is thread-safe. @@ -256,8 +265,7 @@ func TestParallelAccessToSingleFile(t *testing.T) { func TestOpenedFileAccessibleAfterMarkedComplete(t *testing.T) { require := require.New(t) - store, err := NewStore(10*memsize.KB, tally.NoopScope) - require.NoError(err) + store := newTestStore(t, 10*memsize.KB) key := core.DigestFixture().Hex() data := []byte("Hello World") @@ -289,15 +297,14 @@ func TestOpenedFileAccessibleAfterMarkedComplete(t *testing.T) { func TestMisreportedBlobSize(t *testing.T) { t.Run("user overreports size", func(t *testing.T) { require := require.New(t) - store, err := NewStore(25*memsize.KB, tally.NoopScope) - require.NoError(err) + store := newTestStore(t, 25*memsize.KB) // We pad the store's size to 1KB before all other operations to later show that size does not underflow when releasing memory. _, _ = newTestFile(t, store, 1*memsize.KB) f, key := newTestFile(t, store, 10*memsize.KB) data := make([]byte, 5*memsize.KB) - _, err = rand.Read(data) + _, err := rand.Read(data) require.NoError(err) _, err = io.Copy(f, bytes.NewReader(data)) require.NoError(err) @@ -327,8 +334,7 @@ func TestMisreportedBlobSize(t *testing.T) { t.Run("user underreports size", func(t *testing.T) { require := require.New(t) - store, err := NewStore(25*memsize.KB, tally.NoopScope) - require.NoError(err) + store := newTestStore(t, 25*memsize.KB) // We pad the store's size to 1KB before all other operations to later show that size does not underflow when releasing memory. _, _ = newTestFile(t, store, 1*memsize.KB) @@ -336,7 +342,7 @@ func TestMisreportedBlobSize(t *testing.T) { // We write 30KB to the file even though the store can only store 24KB. The store lets us do this and does not evict the blob. f, key := newTestFile(t, store, 24*memsize.KB) data := make([]byte, 30*memsize.KB) - _, err = rand.Read(data) + _, err := rand.Read(data) require.NoError(err) _, err = io.Copy(f, bytes.NewReader(data)) require.NoError(err) @@ -369,8 +375,7 @@ func TestMisreportedBlobSize(t *testing.T) { func TestDelete(t *testing.T) { t.Run("incomplete blob", func(t *testing.T) { require := require.New(t) - store, err := NewStore(10*memsize.KB, tally.NoopScope) - require.NoError(err) + store := newTestStore(t, 10*memsize.KB) key := core.DigestFixture().Hex() f, err := store.Create(key, 100*memsize.B) require.NoError(err) @@ -385,8 +390,7 @@ func TestDelete(t *testing.T) { }) t.Run("incomplete, unevictable blob", func(t *testing.T) { require := require.New(t) - store, err := NewStore(10*memsize.KB, tally.NoopScope) - require.NoError(err) + store := newTestStore(t, 10*memsize.KB) key := core.DigestFixture().Hex() f, err := store.Create(key, 100*memsize.B) require.NoError(err) @@ -403,8 +407,7 @@ func TestDelete(t *testing.T) { }) t.Run("complete blob", func(t *testing.T) { require := require.New(t) - store, err := NewStore(10*memsize.KB, tally.NoopScope) - require.NoError(err) + store := newTestStore(t, 10*memsize.KB) key := core.DigestFixture().Hex() f, err := store.Create(key, 100*memsize.B) require.NoError(err) @@ -421,8 +424,7 @@ func TestDelete(t *testing.T) { }) t.Run("complete, unevictable blob", func(t *testing.T) { require := require.New(t) - store, err := NewStore(10*memsize.KB, tally.NoopScope) - require.NoError(err) + store := newTestStore(t, 10*memsize.KB) key := core.DigestFixture().Hex() f, err := store.Create(key, 100*memsize.B) require.NoError(err) @@ -440,11 +442,10 @@ func TestDelete(t *testing.T) { }) t.Run("not found", func(t *testing.T) { require := require.New(t) - store, err := NewStore(10*memsize.KB, tally.NoopScope) - require.NoError(err) + store := newTestStore(t, 10*memsize.KB) key := core.DigestFixture().Hex() - err = store.Delete(key) + err := store.Delete(key) require.ErrorIs(err, os.ErrNotExist) }) } @@ -452,8 +453,7 @@ func TestDelete(t *testing.T) { func TestMarkComplete(t *testing.T) { t.Run("incomplete blob", func(t *testing.T) { require := require.New(t) - store, err := NewStore(10*memsize.KB, tally.NoopScope) - require.NoError(err) + store := newTestStore(t, 10*memsize.KB) key := core.DigestFixture().Hex() f, err := store.Create(key, 100*memsize.B) require.NoError(err) @@ -470,8 +470,7 @@ func TestMarkComplete(t *testing.T) { }) t.Run("incomplete blob with forbidden eviction", func(t *testing.T) { require := require.New(t) - store, err := NewStore(10*memsize.KB, tally.NoopScope) - require.NoError(err) + store := newTestStore(t, 10*memsize.KB) key := core.DigestFixture().Hex() f, err := store.Create(key, 100*memsize.B) require.NoError(err) @@ -489,8 +488,7 @@ func TestMarkComplete(t *testing.T) { }) t.Run("already complete blob", func(t *testing.T) { require := require.New(t) - store, err := NewStore(10*memsize.KB, tally.NoopScope) - require.NoError(err) + store := newTestStore(t, 10*memsize.KB) key := core.DigestFixture().Hex() f, err := store.Create(key, 100*memsize.B) require.NoError(err) @@ -503,8 +501,7 @@ func TestMarkComplete(t *testing.T) { }) t.Run("already complete blob with forbidden eviction", func(t *testing.T) { require := require.New(t) - store, err := NewStore(10*memsize.KB, tally.NoopScope) - require.NoError(err) + store := newTestStore(t, 10*memsize.KB) key := core.DigestFixture().Hex() f, err := store.Create(key, 100*memsize.B) require.NoError(err) @@ -518,11 +515,10 @@ func TestMarkComplete(t *testing.T) { }) t.Run("not found", func(t *testing.T) { require := require.New(t) - store, err := NewStore(10*memsize.KB, tally.NoopScope) - require.NoError(err) + store := newTestStore(t, 10*memsize.KB) key := core.DigestFixture().Hex() - err = store.MarkComplete(key) + err := store.MarkComplete(key) require.ErrorIs(err, os.ErrNotExist) }) } @@ -530,8 +526,7 @@ func TestMarkComplete(t *testing.T) { func TestStat(t *testing.T) { t.Run("complete blob", func(t *testing.T) { require := require.New(t) - store, err := NewStore(10*memsize.KB, tally.NoopScope) - require.NoError(err) + store := newTestStore(t, 10*memsize.KB) key := core.DigestFixture().Hex() f, err := store.Create(key, 10*memsize.B) require.NoError(err) @@ -546,8 +541,7 @@ func TestStat(t *testing.T) { }) t.Run("complete, unevictable blob", func(t *testing.T) { require := require.New(t) - store, err := NewStore(10*memsize.KB, tally.NoopScope) - require.NoError(err) + store := newTestStore(t, 10*memsize.KB) key := core.DigestFixture().Hex() f, err := store.Create(key, 10*memsize.B) require.NoError(err) @@ -563,8 +557,7 @@ func TestStat(t *testing.T) { }) t.Run("incomplete blob", func(t *testing.T) { require := require.New(t) - store, err := NewStore(10*memsize.KB, tally.NoopScope) - require.NoError(err) + store := newTestStore(t, 10*memsize.KB) key := core.DigestFixture().Hex() f, err := store.Create(key, 10*memsize.B) require.NoError(err) @@ -582,8 +575,7 @@ func TestStat(t *testing.T) { t.Run("incomplete, unevictable blob", func(t *testing.T) { require := require.New(t) - store, err := NewStore(10*memsize.KB, tally.NoopScope) - require.NoError(err) + store := newTestStore(t, 10*memsize.KB) key := core.DigestFixture().Hex() f, err := store.Create(key, 10*memsize.B) require.NoError(err) @@ -598,8 +590,7 @@ func TestStat(t *testing.T) { }) t.Run("non-existent blob", func(t *testing.T) { require := require.New(t) - store, err := NewStore(10*memsize.KB, tally.NoopScope) - require.NoError(err) + store := newTestStore(t, 10*memsize.KB) key := core.DigestFixture().Hex() size, err := store.Stat(key) @@ -610,8 +601,7 @@ func TestStat(t *testing.T) { func TestList(t *testing.T) { require := require.New(t) - store, err := NewStore(10*memsize.KB, tally.NoopScope) - require.NoError(err) + store := newTestStore(t, 10*memsize.KB) require.Empty(store.ScopeComplete().List()) require.Empty(store.ScopeComplete().List()) @@ -668,8 +658,7 @@ func (mda *immovableMd) Deserialize(b []byte) error { return ni func TestMetadata(t *testing.T) { t.Run("basic functionality", func(t *testing.T) { require := require.New(t) - store, err := NewStore(10*memsize.KB, tally.NoopScope) - require.NoError(err) + store := newTestStore(t, 10*memsize.KB) key := core.DigestFixture().Hex() f, err := store.Create(key, 10*memsize.KB) require.NoError(err) @@ -727,13 +716,12 @@ func TestMetadata(t *testing.T) { t.Run("non-existent blob", func(t *testing.T) { require := require.New(t) - store, err := NewStore(10*memsize.KB, tally.NoopScope) - require.NoError(err) + store := newTestStore(t, 10*memsize.KB) nonExistentKey := core.DigestFixture().Hex() mdStruct := core.MetaInfoFixture() md := metadata.NewTorrentMeta(mdStruct) - err = store.SetMetadata(nonExistentKey, md) + err := store.SetMetadata(nonExistentKey, md) require.ErrorIs(err, os.ErrNotExist) ok, err := store.GetMetadata(nonExistentKey, md) @@ -749,8 +737,7 @@ func TestMetadata(t *testing.T) { t.Run("metadata does not change after marking a file as complete and/or evictable/unevictable", func(t *testing.T) { require := require.New(t) - store, err := NewStore(10*memsize.KB, tally.NoopScope) - require.NoError(err) + store := newTestStore(t, 10*memsize.KB) key := core.DigestFixture().Hex() f, err := store.Create(key, 1*memsize.KB) require.NoError(err) @@ -791,8 +778,7 @@ func TestMetadata(t *testing.T) { }) t.Run("metadata fully gone after blob is evicted", func(t *testing.T) { require := require.New(t) - store, err := NewStore(10*memsize.KB, tally.NoopScope) - require.NoError(err) + store := newTestStore(t, 10*memsize.KB) keyA := core.DigestFixture().Hex() fA, err := store.Create(keyA, 10*memsize.KB) require.NoError(err) @@ -820,8 +806,7 @@ func TestMetadata(t *testing.T) { t.Run("immovable metadata is deleted when calling MarkComplete", func(t *testing.T) { require := require.New(t) - store, err := NewStore(10*memsize.KB, tally.NoopScope) - require.NoError(err) + store := newTestStore(t, 10*memsize.KB) key := core.DigestFixture().Hex() f, err := store.Create(key, 10*memsize.KB) require.NoError(err) @@ -847,13 +832,12 @@ func TestMetadata(t *testing.T) { func TestScopes(t *testing.T) { require := require.New(t) - store, err := NewStore(10*memsize.KB, tally.NoopScope) - require.NoError(err) + store := newTestStore(t, 10*memsize.KB) // Add a file to the store, fill it with data. f, key := newTestFile(t, store, 2*memsize.KB) data := make([]byte, 2*memsize.KB) - _, err = rand.Read(data) + _, err := rand.Read(data) require.NoError(err) _, err = io.Copy(f, bytes.NewReader(data)) require.NoError(err) @@ -1005,14 +989,13 @@ func TestScopes(t *testing.T) { func TestScopesDelete(t *testing.T) { require := require.New(t) - store, err := NewStore(10*memsize.KB, tally.NoopScope) - require.NoError(err) + store := newTestStore(t, 10*memsize.KB) f, key := newTestFile(t, store, 1*memsize.KB) require.NoError(f.Close()) require.ErrorIs(store.ScopeComplete().Delete(key), storelib.ErrOutOfScope) require.NoError(store.ScopeIncomplete().Delete(key)) - _, err = store.Stat(key) + _, err := store.Stat(key) require.ErrorIs(err, os.ErrNotExist) f, key = newTestFile(t, store, 1*memsize.KB) @@ -1032,3 +1015,83 @@ func TestScopesDelete(t *testing.T) { require.NoError(store.MarkComplete(key)) require.NoError(store.Delete(key)) } + +func TestConfig__applyDefaults(t *testing.T) { + for name, tt := range map[string]struct { + config *Config + wantErr string + }{ + "capacity must be explicitly set": { + config: &Config{ + GOMEMLIMITBytes: math.MaxInt64, + GOGC: 100, + }, + wantErr: "capacity_bytes must be explicitly set", + }, + "minimal config": { + config: &Config{ + CapacityBytes: 10 * memsize.KB, + GOMEMLIMITBytes: math.MaxInt64, + }, + wantErr: "", + }, + } { + t.Run(name, func(t *testing.T) { + _, err := NewStore(tt.config, tally.NoopScope) + if tt.wantErr != "" { + require.EqualError(t, err, tt.wantErr) + } else { + require.NoError(t, err) + } + }) + } +} + +func TestConfig__configureGC(t *testing.T) { + t.Run("GOMEMLIMIT must be configured somehow", func(t *testing.T) { + require := require.New(t) + // disable GOMEMLIMIT by setting it to its default value of [math.MaxInt64] + tempChangeGCParams(t, 100, math.MaxInt64) + config := &Config{ + GOGC: 200, + CapacityBytes: 70 * memsize.GB, + } + + _, err := NewStore(config, tally.NoopScope) + require.EqualError(err, "GOMEMLIMIT must be configured either through memory.Config or an env var") + }) + + t.Run("config params are applied to GC even if previously configured", func(t *testing.T) { + require := require.New(t) + + memLimitBeforeStoreCreation, goGCBeforeStoreCreation := int64(101*memsize.GB), 99 + goGCBeforeTest, memLimitBeforeTest := tempChangeGCParams(t, goGCBeforeStoreCreation, memLimitBeforeStoreCreation) + + configMemLimit := int64(100 * memsize.GB) + configGoGC := 200 + config := &Config{ + GOMEMLIMITBytes: configMemLimit, + GOGC: configGoGC, + CapacityBytes: 70 * memsize.GB, + } + + _, err := NewStore(config, tally.NoopScope) + require.NoError(err) + + memLimitAfterStoreCreation := debug.SetMemoryLimit(memLimitBeforeTest) + goGcAfterStoreCreation := debug.SetGCPercent(goGCBeforeTest) + + require.Equal(configMemLimit, memLimitAfterStoreCreation) + require.Equal(configGoGC, goGcAfterStoreCreation) + }) +} + +func tempChangeGCParams(t *testing.T, goGC int, goMemLimit int64) (goGCBeforeTest int, memLimitBeforeTest int64) { + memLimitBeforeTest = debug.SetMemoryLimit(int64(goMemLimit)) + goGCBeforeTest = debug.SetGCPercent(goGC) + t.Cleanup(func() { + debug.SetMemoryLimit(memLimitBeforeTest) + debug.SetGCPercent(goGCBeforeTest) + }) + return goGCBeforeTest, memLimitBeforeTest +} diff --git a/lib/store/tiered/config.go b/lib/store/tiered/config.go new file mode 100644 index 000000000..086bb527f --- /dev/null +++ b/lib/store/tiered/config.go @@ -0,0 +1,30 @@ +package tiered + +import ( + "errors" + + "github.com/uber/kraken/lib/store/disk" + "github.com/uber/kraken/lib/store/memory" +) + +const _defaultNumFlushWorkers = 10 + +type Config struct { + DiskConfig *disk.Config `yaml:"disk_store"` + MemConfig *memory.Config `yaml:"memory_store"` + NumFlushWorkers int `yaml:"num_flush_workers"` + // TODO - support disabling mem layer through the config to support quick incident mitigation. +} + +func (c *Config) applyDefaults() error { + if c.DiskConfig.RebootIncompleteBlobs { + return errors.New("tiered.Store does not support RebootIncompleteBlobs, as it can leak files. Use disk.Store if you need persistence for incomplete blobs") + } + if c.NumFlushWorkers < 0 { + return errors.New("num_flush_workers must be at least 1, otherwise blobs will never get flushed from mem to disk") + } + if c.NumFlushWorkers == 0 { + c.NumFlushWorkers = _defaultNumFlushWorkers + } + return nil +} diff --git a/lib/store/tiered/file.go b/lib/store/tiered/file.go new file mode 100644 index 000000000..8d2a5a1eb --- /dev/null +++ b/lib/store/tiered/file.go @@ -0,0 +1,218 @@ +package tiered + +import ( + "errors" + "fmt" + "io" + "os" + "sync" + + storelib "github.com/uber/kraken/lib/store" + "github.com/uber/kraken/lib/store/disk" + "github.com/uber/kraken/lib/store/memory" + "github.com/uber/kraken/utils/closers" + "go.uber.org/zap" +) + +var _ storelib.FileReadWriter = &File{} + +// File represents an open handle to a blob in [Store], similar to how an [os.File] is an open file descriptor to a file on disk. +// File, just like a Linux file with its page cache, continues to work seemlessly if the blob is initially in memory, but later flushed to disk and evicted from memory. +type File struct { + memF *memory.File + diskF *disk.File + key string + diskStore *disk.Store + openErr error + once sync.Once + log *zap.SugaredLogger +} + +// Either [memory.File] OR [disk.File] must be set, but not both. If disk.File is not set, it will be lazily initialized +// on File's first call after the blob is evicted from memory. +func newFile(key string, memF *memory.File, diskF *disk.File, diskStore *disk.Store, log *zap.SugaredLogger) *File { + f := &File{ + memF: memF, + diskF: nil, + key: key, + diskStore: diskStore, + log: log, + openErr: nil, + } + if diskF != nil { + f.once.Do(func() { f.diskF = diskF }) + } + return f +} + +// A blob might be evicted from the mem store (after being flushed to disk), while a tiered store user is holding a [File]. +// To allow the user to continue operating on the [File] seemlessly, we open the blob in the disk store and +// operate on it instead. This is done once, on the first call to [File] after eviction from memory. +func (f *File) openDiskFileIfNeeded() error { + f.once.Do(func() { + diskF, err := f.diskStore.Open(f.key) + if errors.Is(err, os.ErrNotExist) { + // TODO - Currently, we don't 100% ensure that a blob cannot be evicted from disk before it is evicted from memory. + // We *hope* that it doesn't happen, as there very specific requirements needed to trigger this event + // (a blob needs to 1) be complete, 2) be fully flushed on disk, and 3) not get evicted from memory + // 4) until all other blobs on disk are used more recently than this blob was fully flushed). + // If that's not the case (detected through the error log below), consider adding a `Touch` API in disk.store + // that resets a blob's LAT and call that API each time a blob gets accessed from memory in the tiered store. + f.openErr = errBadSwitch(errors.New("blob not found in neither memory store nor disk store")) + f.log.With("key", f.key). + Error("invariant violation - a blob is neither memory.Store nor disk.Store while a tiered.Store user is trying to operate on the file") + return + } + if err != nil { + f.openErr = errBadSwitch(fmt.Errorf("disk store open: %w", err)) + f.log.With( + "key", f.key, + "error", f.openErr). + Error("tiered.Store user could not operate on a blob, as the tiered.File could not transition from pointing at memory.File to pointing at disk.File after blob eviction from memory store") + return + } + _, err = diskF.Seek(f.memF.Off(), io.SeekStart) + if err != nil { + closers.Close(diskF) + f.openErr = errBadSwitch(fmt.Errorf("disk file seek: %w", err)) + f.log.With( + "key", f.key, + "error", f.openErr). + Error("tiered.Store user could not operate on a blob, as the tiered.File could not transition from pointing at memory.File to pointing at disk.File after blob eviction from memory store") + return + } + + f.diskF = diskF + }) + + return f.openErr +} + +func errBadSwitch(subErr error) error { + return fmt.Errorf("tiered.File could not switch over from pointing to memory.File to pointing to disk.File after blob was evicted from memory: %w", subErr) +} + +func (f *File) Read(p []byte) (n int, err error) { + if f.memF != nil { + n, err = f.memF.Read(p) + if err != memory.ErrEvicted { + return n, err + } + + if err := f.openDiskFileIfNeeded(); err != nil { + return 0, err + } + } + return f.diskF.Read(p) +} + +func (f *File) ReadAt(p []byte, off int64) (n int, err error) { + if f.memF != nil { + n, err = f.memF.ReadAt(p, off) + if err != memory.ErrEvicted { + return n, err + } + if err := f.openDiskFileIfNeeded(); err != nil { + return 0, err + } + } + return f.diskF.ReadAt(p, off) +} + +func (f *File) Seek(off int64, whence int) (int64, error) { + if f.memF != nil { + newOff, err := f.memF.Seek(off, whence) + if err != memory.ErrEvicted { + return newOff, err + } + + if err := f.openDiskFileIfNeeded(); err != nil { + return 0, err + } + } + return f.diskF.Seek(off, whence) +} + +func (f *File) Size() int64 { + if f.memF != nil { + size := f.memF.Size() + if size != -1 { + // Same as [memory.ErrEvicted]. + return size + } + + if err := f.openDiskFileIfNeeded(); err != nil { + return 0 + } + } + return f.diskF.Size() +} + +func (f *File) WriteAt(p []byte, off int64) (n int, err error) { + if f.memF != nil { + n, err = f.memF.WriteAt(p, off) + if err != memory.ErrEvicted { + return n, err + } + + if err := f.openDiskFileIfNeeded(); err != nil { + return 0, err + } + } + return f.diskF.WriteAt(p, off) +} + +func (f *File) Write(p []byte) (n int, err error) { + if f.memF != nil { + n, err = f.memF.Write(p) + if err != memory.ErrEvicted { + return n, err + } + + if err := f.openDiskFileIfNeeded(); err != nil { + return 0, err + } + } + return f.diskF.Write(p) +} + +func (f *File) Cancel() error { + if f.memF != nil { + err := f.memF.Cancel() + if err != memory.ErrEvicted { + return err + } + + if err := f.openDiskFileIfNeeded(); err != nil { + return err + } + } + return f.diskF.Cancel() +} +func (f *File) Close() error { + if f.memF != nil { + err := f.memF.Close() + if err != memory.ErrEvicted { + return err + } + + if err := f.openDiskFileIfNeeded(); err != nil { + return err + } + } + return f.diskF.Close() + +} +func (f *File) Commit() error { + if f.memF != nil { + err := f.memF.Commit() + if err != memory.ErrEvicted { + return err + } + + if err := f.openDiskFileIfNeeded(); err != nil { + return err + } + } + return f.diskF.Commit() +} diff --git a/lib/store/tiered/file_test.go b/lib/store/tiered/file_test.go new file mode 100644 index 000000000..db19e08a5 --- /dev/null +++ b/lib/store/tiered/file_test.go @@ -0,0 +1,138 @@ +package tiered + +import ( + "bytes" + "fmt" + "io" + "sync" + "testing" + "time" + + "github.com/stretchr/testify/require" + "github.com/uber/kraken/core" + "github.com/uber/kraken/utils/memsize" + "github.com/uber/kraken/utils/testutil" +) + +func TestFile_MemToDiskLifecycle(t *testing.T) { + require := require.New(t) + // Mem capacity fits exactly one blob at a time, forcing eviction once a + // second blob is created after the first is flushed and unbanned. + s, _ := newTestStore(t, 10*memsize.KB, 1*memsize.KB, 1) + + key, data := createAndWrite(t, s, 1*memsize.KB) + require.NoError(s.MarkComplete(key)) + waitForFlushed(t, s, key) + + f, err := s.Open(key) + require.NoError(err) + defer func() { require.NoError(f.Close()) }() + + // Creating a same-sized filler blob synchronously evicts key from mem, + // since it's the only evictable, complete blob and capacity is exactly 1KB. + _, _ = createAndWrite(t, s, 1*memsize.KB) + _, inMem := s.impl.mem.Has(key) + require.False(inMem, "key should have been evicted from mem to make room for the filler blob") + + // The already-open File must transparently continue to work, now served from disk. + got, err := io.ReadAll(f) + require.NoError(err) + require.Equal(data, got) +} + +func TestFile_ParallelAccessDuringMemToDiskTransition(t *testing.T) { + require := require.New(t) + s, _ := newTestStore(t, 10*memsize.KB, 1*memsize.KB, 1) + + const chunkSize = 10 + const numChunks = 5 + data := randomData(t, chunkSize*numChunks) + key := core.DigestFixture().Hex() + f, err := s.Create(key, uint64(len(data))) + require.NoError(err) + _, err = f.Write(data) + require.NoError(err) + require.NoError(f.Close()) + require.NoError(s.MarkComplete(key)) + waitForFlushed(t, s, key) + + readF, err := s.Open(key) + require.NoError(err) + defer func() { require.NoError(readF.Close()) }() + + const numReadsPerGoroutine = 50 + var wg sync.WaitGroup + wg.Add(numChunks) + errs := make([]error, numChunks) + for i := range numChunks { + go func(i int) { + defer wg.Done() + want := data[i*chunkSize : (i+1)*chunkSize] + buf := make([]byte, chunkSize) + for range numReadsPerGoroutine { + n, err := readF.ReadAt(buf, int64(i*chunkSize)) + if err != nil { + errs[i] = err + return + } + if n != chunkSize || !bytes.Equal(buf, want) { + errs[i] = fmt.Errorf("unexpected data on chunk %d: got %v, want %v", i, buf, want) + return + } + } + }(i) + } + + // Force a real mem eviction of key while the reads above are in flight. + _, _ = createAndWrite(t, s, 1*memsize.KB) + + wg.Wait() + for i := range numChunks { + require.NoError(errs[i]) + } + + _, inMem := s.impl.mem.Has(key) + require.False(inMem) +} + +// TestFile_BlobEvictedFromDiskBeforeMem exercises the edge case documented in +// file.go's openDiskFileIfNeeded TODO: a blob evicted from disk while still +// resident in mem. This requires very specific LRU timing and is unexpected +// in production - it's exercised here for documentation purposes. +func TestFile_BlobEvictedFromDiskBeforeMem(t *testing.T) { + require := require.New(t) + // Disk capacity fits only one blob; mem capacity fits key + otherKey. + s, _ := newTestStore(t, 1*memsize.KB, 2*memsize.KB, 1) + + key, _ := createAndWrite(t, s, 1*memsize.KB) + require.NoError(s.MarkComplete(key)) + waitForFlushed(t, s, key) + + // Open a handle to key while it's still resident in mem, before disk evicts it. + f, err := s.Open(key) + require.NoError(err) + // f.Close() is a no-op on the mem side regardless of the invariant + // violation triggered below, so it's still safe to assert on here. + defer func() { require.NoError(f.Close()) }() + + // Completing otherKey forces disk (1KB capacity) to evict key to make room. + otherKey, _ := createAndWrite(t, s, 1*memsize.KB) + require.NoError(s.MarkComplete(otherKey)) + err = testutil.PollUntilTrue(5*time.Second, func() bool { + _, inDisk := s.impl.disk.Has(key) + return !inDisk + }) + require.NoError(err) + + // Now evict key from mem too (capacity 2KB fits key + otherKey exactly; + // one more filler forces out key, since it's LRU-oldest). + _, _ = createAndWrite(t, s, 1*memsize.KB) + _, inMem := s.impl.mem.Has(key) + require.False(inMem) + + // The already-open handle can no longer switch over to disk, since disk evicted key first. + buf := make([]byte, 10) + _, err = f.Read(buf) + require.Error(err) + require.ErrorContains(err, "blob not found in neither memory store nor disk store") +} diff --git a/lib/store/tiered/flusher.go b/lib/store/tiered/flusher.go new file mode 100644 index 000000000..e97ef6f94 --- /dev/null +++ b/lib/store/tiered/flusher.go @@ -0,0 +1,325 @@ +package tiered + +import ( + "errors" + "fmt" + "io" + "os" + "sync" + + "github.com/uber/kraken/lib/store/disk" + "github.com/uber/kraken/lib/store/memory" + "github.com/uber/kraken/lib/store/metadata" + "github.com/uber/kraken/utils/closers" + "go.uber.org/zap" +) + +// indirection needed for testing +var ( + memOpen = func(mem *memory.Store, key string) (*memory.File, error) { + return mem.Open(key) + } + ioCopy = io.Copy +) + +type flusher struct { + blobs map[string]*blob + queue []string + mu sync.Mutex + notify chan struct{} + stop chan struct{} + + mem *memory.Store + disk *disk.Store + log *zap.SugaredLogger +} + +type blob struct { + key string + dataDirty bool + dataSize uint64 + dirtyMD map[string]struct{} + mu sync.Mutex +} + +func newFlusher(mem *memory.Store, disk *disk.Store, log *zap.SugaredLogger, numWorkers int) *flusher { + f := &flusher{ + blobs: make(map[string]*blob, 0), + queue: make([]string, 0), + notify: make(chan struct{}, numWorkers), + stop: make(chan struct{}), + mem: mem, + disk: disk, + log: log, + } + + for range numWorkers { + go f.worker() + } + return f +} + +// Enqueues a blob and its metadata for flushing to disk. Should be called once per blob, upon MarkComplete. +// Once a blob has been flushed, UnbanEviction is called on the mem store. +func (f *flusher) markDirty(key string, sizeBytes uint64) { + f.mu.Lock() + defer f.mu.Unlock() + + dirtyMD, err := f.mem.ListMetadata(key) + if err != nil { + f.log.With( + "key", key, + "error", fmt.Errorf("mem store list metadata: %w", err), + ).Error("Could not mark blob as dirty, aborting flushing") + _ = f.mem.UnbanEviction(key) // prevent leak + return + } + dirtyMDMap := make(map[string]struct{}) + for _, dirtyMD := range dirtyMD { + dirtyMDMap[dirtyMD.GetSuffix()] = struct{}{} + } + + f.blobs[key] = &blob{ + key: key, + dataDirty: true, + dataSize: sizeBytes, + dirtyMD: dirtyMDMap, + } + f.queue = append(f.queue, key) + + select { + case f.notify <- struct{}{}: + default: + } +} + +// Must be called on each md mutation. +func (f *flusher) markMetadataDirty(key, mdSuffix string) { + f.mu.Lock() + defer f.mu.Unlock() + + b, dirty := f.blobs[key] + if dirty { + b.mu.Lock() + defer b.mu.Unlock() + + b.dirtyMD[mdSuffix] = struct{}{} + return + } + + _, onDisk := f.disk.Has(key) + if !onDisk { + return // Blob is not yet complete, we can't start flushing. + } + + f.blobs[key] = &blob{ + key: key, + dataDirty: false, + dirtyMD: map[string]struct{}{mdSuffix: {}}, + } + f.queue = append(f.queue, key) + + select { + case f.notify <- struct{}{}: + default: + } +} + +// Aborts the flushing of a blob, possibly midway through. May leave corrupt state on disk, which caller is responsible for cleaning. +func (f *flusher) abort(key string) { + f.mu.Lock() + defer f.mu.Unlock() + + delete(f.blobs, key) +} + +// Stops all flushing gracefully. +func (f *flusher) close() { + // TODO - consider whether there's a race when 1) a blob's data is flushed on disk, but the entry is in mem too, 2) new md is enqueued for flushing, + // 3) close is called, and 4) the workers exit before flushing the metadata, thus the disk blob is corrupt (blob data is persisted, but not metadata). + close(f.stop) +} + +func (f *flusher) worker() { + for { + select { + case <-f.stop: + return + case <-f.notify: + for { + b, ok := f.nextToFlush() + if !ok { + break + } + f.flush(b) + } + } + } +} + +func (f *flusher) nextToFlush() (b *blob, ok bool) { + f.mu.Lock() + defer f.mu.Unlock() + + for len(f.queue) > 0 { + key := f.queue[0] + f.queue = f.queue[1:] + b, ok := f.blobs[key] + if !ok { + continue + } + return b, true + } + return nil, false +} + +func (f *flusher) flush(b *blob) { + key := b.key + defer f.mem.UnbanEviction(key) + + if b.dataDirty { + if err := f.flushData(b); err != nil { + err = fmt.Errorf("flush data: %w", err) + f.log.With( + "key", key, + "error", err, + ).Error("Could not flush data from mem to disk, abandoning flushing operation") + f.handleFlushFailure(key) + return + } + } + + if err := f.flushMetadatasAndUnmarkDirty(key, b); err != nil { + err = fmt.Errorf("flush metadatas: %w", err) + f.log.With( + "key", key, + "error", err, + ).Error("Could not flush metadata from mem to disk, abandoning flushing operation") + f.handleFlushFailure(key) + } +} + +// Tries to prevent corrupt state in disk store upon failed flush. To do so, the flush to disk is aborted, +// simulating similar behavior to the file being flushed and subsequently evicted by the disk store's LRU policy. +// This will break any open [File] handles to the blob after eviction from memory, but it's the best we can do. +func (f *flusher) handleFlushFailure(key string) { + if err := f.disk.Delete(key); err != nil && !errors.Is(err, os.ErrNotExist) { + f.log.With( + "key", key, + "error", err). + Error("Could not clean disk entry after flushing failed, blob is now leaked in disk store") + } + + f.mu.Lock() + defer f.mu.Unlock() + delete(f.blobs, key) +} + +func (f *flusher) flushMetadatasAndUnmarkDirty(key string, b *blob) error { + b.mu.Lock() + for { + dirtyMDSnapshot := b.dirtyMD + b.dirtyMD = make(map[string]struct{}) + b.mu.Unlock() + + for mdSuffix := range dirtyMDSnapshot { + err := f.flushMetadata(key, mdSuffix) + if err != nil { + f.log.With( + "key", key, + "mdSuffix", mdSuffix, + "error", err, + ).Error("Could not flush metadata from mem to disk") + return fmt.Errorf("flush md: %w", err) + } + } + + f.mu.Lock() + b.mu.Lock() + if len(b.dirtyMD) == 0 { + delete(f.blobs, key) + b.mu.Unlock() + f.mu.Unlock() + return nil + } + + f.mu.Unlock() + } +} + +func (f *flusher) flushMetadata(key, mdSuffix string) error { + md := metadata.CreateFromSuffix(mdSuffix) + ok, err := f.mem.GetMetadata(key, md) + if errors.Is(err, os.ErrNotExist) { + return nil + } + if err != nil { + return fmt.Errorf("mem store get md: %w", err) + } + if !ok { + err = f.disk.DeleteMetadata(key, md.GetSuffix()) + if errors.Is(err, os.ErrNotExist) { + return nil + } + if err != nil { + return fmt.Errorf("disk store delete md: %w", err) + } + return nil + } + err = f.disk.SetMetadata(key, md) + if errors.Is(err, os.ErrNotExist) { + return nil + } + if err != nil { + return fmt.Errorf("disk store set md: %w", err) + } + + return nil +} + +func (f *flusher) flushData(b *blob) error { + key := b.key + memF, err := memOpen(f.mem, key) + if errors.Is(err, os.ErrNotExist) { + return nil + } + if err != nil { + return fmt.Errorf("mem store open: %w", err) + } + defer closers.Close(memF) + diskF, err := f.disk.Create(key, b.dataSize) + if err != nil { + return fmt.Errorf("disk store create: %w", err) + } + defer closers.Close(diskF) + f.mu.Lock() + _, ok := f.blobs[b.key] + if !ok { + // abort was called before we created the file, we need to cleanup. + err := f.disk.Delete(key) + if err != nil && !errors.Is(err, os.ErrNotExist) { + f.log.With( + "key", key, + "error", err). + Error("Could not clean disk entry after flushing failed, blob is now leaked in disk store") + } + f.mu.Unlock() + return nil + } + f.mu.Unlock() + _, err = ioCopy(diskF, memF) + if errors.Is(err, memory.ErrEvicted) { + return nil + } + if err != nil { + return fmt.Errorf("io copy from mem file to disk file: %w", err) + } + err = f.disk.MarkComplete(key) + if errors.Is(err, os.ErrNotExist) { + return nil + } + if err != nil { + return fmt.Errorf("disk mark complete: %w", err) + } + return nil +} diff --git a/lib/store/tiered/flusher_test.go b/lib/store/tiered/flusher_test.go new file mode 100644 index 000000000..6eb053f31 --- /dev/null +++ b/lib/store/tiered/flusher_test.go @@ -0,0 +1,113 @@ +package tiered + +import ( + "io" + "testing" + "time" + + "github.com/stretchr/testify/require" + "github.com/uber/kraken/core" + "github.com/uber/kraken/lib/store/metadata" + "github.com/uber/kraken/utils/memsize" + "github.com/uber/kraken/utils/testutil" +) + +func TestFlusher_MetadataOnlyFlush(t *testing.T) { + require := require.New(t) + s, key, _ := newStateBothComplete(t) + + md := metadata.NewTorrentMeta(core.MetaInfoFixture()) + require.NoError(s.SetMetadata(key, md)) + + err := testutil.PollUntilTrue(5*time.Second, func() bool { + var readMd metadata.TorrentMeta + ok, err := s.impl.disk.GetMetadata(key, &readMd) + return err == nil && ok + }) + require.NoError(err) + + var readMd metadata.TorrentMeta + ok, err := s.impl.disk.GetMetadata(key, &readMd) + require.NoError(err) + require.True(ok) + require.Equal(md.MetaInfo, readMd.MetaInfo) +} + +func TestFlusher_HandleFlushFailure(t *testing.T) { + require := require.New(t) + s, _ := newTestStore(t, 1*memsize.B, 10*memsize.KB, 1) + + key, _ := createAndWrite(t, s, 1*memsize.KB) + require.NoError(s.MarkComplete(key)) + + err := testutil.PollUntilTrue(5*time.Second, func() bool { + s.impl.flusher.mu.Lock() + defer s.impl.flusher.mu.Unlock() + _, tracked := s.impl.flusher.blobs[key] + return !tracked + }) + require.NoError(err) + + _, inDisk := s.impl.disk.Has(key) + require.False(inDisk) + + _, ok := s.impl.mem.Has(key) + require.True(ok) + + // The blob should be evictable again after the failed flush left it + // unbanned: filling the rest of mem's capacity evicts it. + _, _ = createAndWrite(t, s, 10*memsize.KB) + _, ok = s.impl.mem.Has(key) + require.False(ok, "blob should be evictable again after a failed flush left it unbanned") +} + +func TestFlusher_DeleteWhileFlushing(t *testing.T) { + require := require.New(t) + s, _ := newTestStore(t, 10*memsize.KB, 10*memsize.KB, 1) + key, _ := createAndWrite(t, s, 1*memsize.KB) + + reached := make(chan struct{}) + proceed := make(chan struct{}) + old := ioCopy + ioCopy = func(dst io.Writer, src io.Reader) (int64, error) { + close(reached) + <-proceed + return old(dst, src) + } + t.Cleanup(func() { + ioCopy = old + select { + case <-proceed: + default: + close(proceed) + } + }) + + require.NoError(s.MarkComplete(key)) + + timedOut := false + select { + case <-reached: + case <-time.After(5 * time.Second): + timedOut = true + } + require.False(timedOut, "timed out waiting for flushData to reach the copy") + + // disk.Create has already succeeded at this point; the entry exists as an + // incomplete disk blob. Deleting now races the in-flight flush. + require.NoError(s.Delete(key)) + close(proceed) + + err := testutil.PollUntilTrue(5*time.Second, func() bool { + s.impl.flusher.mu.Lock() + defer s.impl.flusher.mu.Unlock() + _, tracked := s.impl.flusher.blobs[key] + return !tracked + }) + require.NoError(err) + + _, inDisk := s.impl.disk.Has(key) + require.False(inDisk, "the disk entry created mid-flush must not survive a concurrent Delete") + _, inMem := s.impl.mem.Has(key) + require.False(inMem) +} diff --git a/lib/store/tiered/scoped_store.go b/lib/store/tiered/scoped_store.go new file mode 100644 index 000000000..f74ccce8b --- /dev/null +++ b/lib/store/tiered/scoped_store.go @@ -0,0 +1,92 @@ +package tiered + +import ( + "github.com/uber-go/tally" + storelib "github.com/uber/kraken/lib/store" + "github.com/uber/kraken/lib/store/disk" + "github.com/uber/kraken/lib/store/metadata" +) + +// Store is a tiered (disk + memory), thread-safe, LRU cache for blobs and their [metadata.Metadata]. +// +// - New blobs are initially created in the memory cache to speed up writes/reads and asynchronously flushed to disk +// once MarkComplete is called on them. If the memory cache is full with inevctable blobs, the new blob is created on disk. +// +// - Partially crash-resistant - all complete blobs that were fully flushed to disk are persisted. Use [disk.Store] if you need full crash resistence. +// +// - Supports pagination of blobs during reading/writing, such that blobs don't need to be fully loaded into memory by the client. +// +// - New blobs are considered 'incomplete', which unlists them from automatic LRU eviction. The store can be scoped to work on only (in-)complete blobs. +// +// - All APIs are thread-safe, including operating on a blob in parallel. +type Store struct { + impl *store + scope storelib.BlobScope +} + +// NewStore creates a new [Store] and returns its underlying [disk.Store] in case the +// user wants to directly operate on disk (e.g. if persistence is mandatory). +func NewStore(config *Config, metrics tally.Scope) (*Store, *disk.Store, error) { + impl, diskStore, err := newStore(config, metrics) + if err != nil { + return nil, nil, err + } + return &Store{ + impl: impl, + scope: storelib.BlobScopeAny, + }, diskStore, nil +} + +// Create initializes a new, incomplete blob, reserves space for it, and returns a [*File] pointing to it. +// Incomplete entries cannot be automatically evicted. MarkComplete must be called once the blob is complete. +// The store uses `sizeBytes` for its eviction logic even if the blob's real size differs (which is tolerated if the difference is negligible). +func (s *Store) Create(key string, sizeBytes uint64) (*File, error) { + return s.impl.Create(key, sizeBytes) +} + +// Open returns a [*File] pointing to the blob. +func (s *Store) Open(key string) (*File, error) { return s.impl.Open(key, s.scope) } + +// Has reports if the blob is 1) in the store and 2) in scope. +// If you don't care about the blob's scope and just want to check membership, do: +// +// _, ok := store.Has(key) +func (s *Store) Has(key string) (inStore bool, inScope bool) { return s.impl.Has(key, s.scope) } + +// Stat returns the blob's actual size, even if it differs from the size reported when Create was called. +func (s *Store) Stat(key string) (size int64, err error) { return s.impl.Stat(key, s.scope) } + +// MarkComplete MUST be called once a blob is fully written, as it +// 1) enqueues it to get flushed to disk (after which the blob is eligible for eviction from the memory store) +// 2) enlists it for LRU eviction (unless BanEviction has been called). +// Additionally, other store APIs may filter blobs based on completeness. No-op if the blob is already complete +func (s *Store) MarkComplete(key string) error { return s.impl.MarkComplete(key) } + +// Delete removes a blob and its metadata from the store. +func (s *Store) Delete(key string) error { return s.impl.Delete(key, s.scope) } + +// List returns the keys of all blobs (except those out of scope). +func (s *Store) List() []string { return s.impl.List(s.scope) } + +// SetMetadata sets the respective metadata of the blob. +func (s *Store) SetMetadata(key string, md metadata.Metadata) error { + return s.impl.SetMetadata(key, md, s.scope) +} + +// GetMetadata populates `md` if the metadata is present. +func (s *Store) GetMetadata(key string, md metadata.Metadata) (ok bool, err error) { + return s.impl.GetMetadata(key, md, s.scope) +} + +// DeleteMetadata removes a blob's metadata. No-op if the md is not present. +func (s *Store) DeleteMetadata(key string, mdSuffix string) error { + return s.impl.DeleteMetadata(key, mdSuffix, s.scope) +} + +// ScopeComplete scopes [Store]'s APIs such that they can only operate on complete blobs. +// [storelib.ErrOutOfScope] is returned if the user tries to operate on an incomplete blob. +func (s *Store) ScopeComplete() *Store { return &Store{s.impl, storelib.BlobScopeComplete} } + +// ScopeIncomplete scopes [Store]'s APIs such that they can only operate on incomplete blobs. +// [storelib.ErrOutOfScope] is returned if the user tries to operate on a complete blob. +func (s *Store) ScopeIncomplete() *Store { return &Store{s.impl, storelib.BlobScopeIncomplete} } diff --git a/lib/store/tiered/store.go b/lib/store/tiered/store.go new file mode 100644 index 000000000..b5158dacf --- /dev/null +++ b/lib/store/tiered/store.go @@ -0,0 +1,280 @@ +package tiered + +import ( + "errors" + "fmt" + "maps" + "os" + "slices" + "sync" + + "github.com/uber-go/tally" + storelib "github.com/uber/kraken/lib/store" + "github.com/uber/kraken/lib/store/disk" + "github.com/uber/kraken/lib/store/memory" + "github.com/uber/kraken/lib/store/metadata" + "github.com/uber/kraken/utils/log" + "go.uber.org/zap" +) + +type store struct { + disk *disk.Store + mem *memory.Store + flusher *flusher + log *zap.SugaredLogger + mu sync.RWMutex +} + +func newStore(config *Config, metrics tally.Scope) (*store, *disk.Store, error) { + err := config.applyDefaults() + if err != nil { + return nil, nil, err + } + memStore, err := memory.NewStore(config.MemConfig, metrics) + if err != nil { + return nil, nil, fmt.Errorf("new mem store: %w", err) + } + diskStore, err := disk.NewStore(config.DiskConfig, metrics) + if err != nil { + return nil, nil, fmt.Errorf("new disk store: %w", err) + } + + log := log.Default().With("module", "tiered_store") + + log.Info("Initialized a new tiered.Store") + return &store{ + disk: diskStore, + mem: memStore, + flusher: newFlusher(memStore, diskStore, log, config.NumFlushWorkers), + log: log, + }, diskStore, nil +} + +func (s *store) Create(key string, sizeBytes uint64) (*File, error) { + s.mu.Lock() + defer s.mu.Unlock() + + _, inMem := s.mem.Has(key) + _, inDisk := s.disk.Has(key) + if inMem || inDisk { + return nil, os.ErrExist + } + + memF, err := s.mem.Create(key, sizeBytes) + if err != nil && !errors.Is(err, memory.ErrNoSpace) { + return nil, fmt.Errorf("mem store create: %w", err) + } + if err == nil { + return newFile(key, memF, nil, s.disk, s.log), nil + } + + diskF, err := s.disk.Create(key, sizeBytes) + if err != nil { + return nil, fmt.Errorf("disk store create: %w", err) + } + return newFile(key, nil, diskF, s.disk, s.log), nil +} + +func (s *store) Open(key string, scope storelib.BlobScope) (*File, error) { + memF, err := s.mem.Scoped(scope).Open(key) + if errors.Is(err, storelib.ErrOutOfScope) { + return nil, storelib.ErrOutOfScope + } + if err != nil && !errors.Is(err, os.ErrNotExist) { + return nil, fmt.Errorf("mem store open: %w", err) + } + if err == nil { + return newFile(key, memF, nil, s.disk, s.log), nil + } + + diskF, err := s.disk.Scoped(scope).Open(key) + if errors.Is(err, storelib.ErrOutOfScope) { + return nil, storelib.ErrOutOfScope + } + if errors.Is(err, os.ErrNotExist) { + return nil, os.ErrNotExist + } + if err != nil { + return nil, fmt.Errorf("disk store open: %w", err) + } + return newFile(key, nil, diskF, s.disk, s.log), nil +} + +func (s *store) Has(key string, scope storelib.BlobScope) (inStore bool, inScope bool) { + inMem, inMemScope := s.mem.Scoped(scope).Has(key) + if inMem { + return true, inMemScope + } + + inDisk, inDiskScope := s.disk.Scoped(scope).Has(key) + if inDisk { + return true, inDiskScope + } + return false, false +} + +func (s *store) Delete(key string, scope storelib.BlobScope) error { + s.mu.Lock() + defer s.mu.Unlock() + + err := s.mem.Scoped(scope).Delete(key) + if errors.Is(err, storelib.ErrOutOfScope) { + return storelib.ErrOutOfScope + } + if err != nil && !errors.Is(err, os.ErrNotExist) { + return fmt.Errorf("mem store delete: %w", err) + } + if errors.Is(err, os.ErrNotExist) { + return s.disk.Scoped(scope).Delete(key) + } + + s.flusher.abort(key) + err = s.disk.Delete(key) + if err != nil && !errors.Is(err, os.ErrNotExist) { + err = fmt.Errorf("disk store delete: %w", err) + s.log.With( + "key", key, + "error", err). + Error("Could not delete blob from disk, blob might be leaked and/or in corrupt state") + return err + } + + return nil +} + +func (s *store) List(scope storelib.BlobScope) []string { + diskKeys := s.disk.Scoped(scope).List() + memKeys := s.mem.Scoped(scope).List() + + dedup := make(map[string]struct{}) + for _, key := range diskKeys { + dedup[key] = struct{}{} + } + for _, key := range memKeys { + dedup[key] = struct{}{} + } + if scope == storelib.BlobScopeIncomplete { + // In the code above, we might have added accidentally entries that are complete from the store's POV, + // but still incomplete in disk, as the entry is currently being flushed from mem to disk. We remove these entries. + for _, key := range s.mem.ScopeComplete().List() { + delete(dedup, key) + } + } + return slices.Collect(maps.Keys(dedup)) +} + +func (s *store) Stat(key string, scope storelib.BlobScope) (size int64, err error) { + size, err = s.mem.Scoped(scope).Stat(key) + if errors.Is(err, storelib.ErrOutOfScope) { + return 0, storelib.ErrOutOfScope + } + if err != nil && !errors.Is(err, os.ErrNotExist) { + return 0, fmt.Errorf("mem store stat: %w", err) + } + if err == nil { + return size, nil + } + + fi, err := s.disk.Scoped(scope).Stat(key) + if errors.Is(err, storelib.ErrOutOfScope) { + return 0, storelib.ErrOutOfScope + } + if errors.Is(err, os.ErrNotExist) { + return 0, os.ErrNotExist + } + if err != nil { + return 0, fmt.Errorf("disk store stat: %w", err) + } + return fi.Size(), nil +} + +func (s *store) MarkComplete(key string) error { + s.mu.Lock() + defer s.mu.Unlock() + + if _, ok := s.mem.ScopeComplete().Has(key); ok { + return nil // no-op + } + if _, ok := s.disk.ScopeComplete().Has(key); ok { + return nil // no-op + } + + err := s.mem.BanEviction(key) + if err != nil && !errors.Is(err, os.ErrNotExist) { + return fmt.Errorf("mem store ban eviction: %w", err) + } + if errors.Is(err, os.ErrNotExist) { + return s.disk.MarkComplete(key) + } + + err = s.mem.MarkComplete(key) + if err != nil { + return fmt.Errorf("mem store mark complete: %w", err) + } + size, err := s.mem.Stat(key) + if err != nil { + return fmt.Errorf("mem store stat: %w", err) + } + s.flusher.markDirty(key, uint64(size)) + return nil +} + +func (s *store) SetMetadata(key string, md metadata.Metadata, scope storelib.BlobScope) error { + s.mu.Lock() + defer s.mu.Unlock() + + err := s.mem.Scoped(scope).BanEviction(key) + if errors.Is(err, storelib.ErrOutOfScope) { + return storelib.ErrOutOfScope + } + if err != nil && !errors.Is(err, os.ErrNotExist) { + return fmt.Errorf("mem store ban eviction: %w", err) + } + if errors.Is(err, os.ErrNotExist) { + return s.disk.Scoped(scope).SetMetadata(key, md) + } + + err = s.mem.SetMetadata(key, md) + if err != nil { + return fmt.Errorf("mem store set metadata: %w", err) + } + s.flusher.markMetadataDirty(key, md.GetSuffix()) + return nil +} + +func (s *store) GetMetadata(key string, md metadata.Metadata, scope storelib.BlobScope) (ok bool, err error) { + ok, err = s.mem.Scoped(scope).GetMetadata(key, md) + if errors.Is(err, storelib.ErrOutOfScope) { + return false, storelib.ErrOutOfScope + } + if err != nil && !errors.Is(err, os.ErrNotExist) { + return false, fmt.Errorf("mem store get metadarta: %w", err) + } + if errors.Is(err, os.ErrNotExist) { + return s.disk.Scoped(scope).GetMetadata(key, md) + } + return ok, nil +} + +func (s *store) DeleteMetadata(key string, mdSuffix string, scope storelib.BlobScope) error { + s.mu.Lock() + defer s.mu.Unlock() + + err := s.mem.Scoped(scope).BanEviction(key) + if errors.Is(err, storelib.ErrOutOfScope) { + return storelib.ErrOutOfScope + } + if err != nil && !errors.Is(err, os.ErrNotExist) { + return fmt.Errorf("mem store ban eviction: %w", err) + } + if errors.Is(err, os.ErrNotExist) { + return s.disk.Scoped(scope).DeleteMetadata(key, mdSuffix) + } + + err = s.mem.DeleteMetadata(key, mdSuffix) + if err != nil { + return fmt.Errorf("mem store delete metadata: %w", err) + } + s.flusher.markMetadataDirty(key, mdSuffix) + return nil +} diff --git a/lib/store/tiered/store_test.go b/lib/store/tiered/store_test.go new file mode 100644 index 000000000..6bd5e2530 --- /dev/null +++ b/lib/store/tiered/store_test.go @@ -0,0 +1,487 @@ +package tiered + +import ( + "crypto/rand" + "io" + "math" + "testing" + "time" + + "github.com/stretchr/testify/require" + "github.com/uber-go/tally" + "github.com/uber/kraken/core" + storelib "github.com/uber/kraken/lib/store" + "github.com/uber/kraken/lib/store/disk" + "github.com/uber/kraken/lib/store/memory" + "github.com/uber/kraken/lib/store/metadata" + "github.com/uber/kraken/utils/memsize" + "github.com/uber/kraken/utils/testutil" +) + +const _testBlobSize = 1 * memsize.KB + +func newTieredConfig(diskCapacity, memCapacity uint64, numWorkers int, rootDir string) *Config { + return &Config{ + NumFlushWorkers: numWorkers, + DiskConfig: &disk.Config{ + RootDir: rootDir, + CapacityBytes: diskCapacity, + RebootIncompleteBlobs: false, + }, + MemConfig: &memory.Config{ + CapacityBytes: memCapacity, + GOMEMLIMITBytes: math.MaxInt64, + }, + } +} + +func newTestStore(t *testing.T, diskCapacity, memCapacity uint64, numWorkers int) (s *Store, rootDir string) { + rootDir = t.TempDir() + config := newTieredConfig(diskCapacity, memCapacity, numWorkers, rootDir) + s, _, err := NewStore(config, tally.NoopScope) + require.NoError(t, err) + return s, rootDir +} + +func randomData(t *testing.T, size uint64) []byte { + data := make([]byte, size) + _, err := rand.Read(data) + require.NoError(t, err) + return data +} + +// createAndWrite creates a blob through the public API and writes data to it. +// Whether the blob lands in mem or disk depends on the store's configured capacities. +func createAndWrite(t *testing.T, s *Store, size uint64) (key string, data []byte) { + key = core.DigestFixture().Hex() + data = randomData(t, size) + f, err := s.Create(key, size) + require.NoError(t, err) + _, err = f.Write(data) + require.NoError(t, err) + require.NoError(t, f.Close()) + return key, data +} + +func waitForComplete(t *testing.T, s *Store, key string) { + err := testutil.PollUntilTrue(5*time.Second, func() bool { + _, ok := s.ScopeComplete().Has(key) + return ok + }) + require.NoError(t, err) +} + +// waitForFlushed waits until key has been fully flushed from mem to disk. +func waitForFlushed(t *testing.T, s *Store, key string) { + err := testutil.PollUntilTrue(5*time.Second, func() bool { + _, ok := s.impl.disk.ScopeComplete().Has(key) + return ok + }) + require.NoError(t, err) +} + +// waitForFlusherIdle waits until the flusher is no longer tracking key. +func waitForFlusherIdle(t *testing.T, s *Store, key string) { + err := testutil.PollUntilTrue(5*time.Second, func() bool { + s.impl.flusher.mu.Lock() + defer s.impl.flusher.mu.Unlock() + _, tracked := s.impl.flusher.blobs[key] + return !tracked + }) + require.NoError(t, err) +} + +// The helpers below deterministically put a blob into one of the store's six +// possible lifecycle states, using the real public API for all of them. + +func seedOnlyMemIncomplete(t *testing.T, s *Store) (key string, data []byte) { + return createAndWrite(t, s, _testBlobSize) +} + +func seedMemCompleteEnqueuedNotOnDisk(t *testing.T, s *Store) (key string, data []byte) { + key, data = createAndWrite(t, s, _testBlobSize) + + reached := make(chan struct{}) + blocked := make(chan struct{}) + old := memOpen + memOpen = func(mem *memory.Store, k string) (*memory.File, error) { + if k == key { + close(reached) + <-blocked + } + return old(mem, k) + } + t.Cleanup(func() { + memOpen = old + close(blocked) + waitForFlusherIdle(t, s, key) + }) + + require.NoError(t, s.MarkComplete(key)) + timedOut := false + select { + case <-reached: + case <-time.After(5 * time.Second): + timedOut = true + } + require.False(t, timedOut, "timed out waiting for flusher worker to reach memOpen") + return key, data +} + +func seedMemCompleteDiskIncomplete(t *testing.T, s *Store) (key string, data []byte) { + key, data = createAndWrite(t, s, _testBlobSize) + + reached := make(chan struct{}) + blocked := make(chan struct{}) + old := ioCopy + ioCopy = func(dst io.Writer, src io.Reader) (int64, error) { + close(reached) + <-blocked + return old(dst, src) + } + t.Cleanup(func() { + ioCopy = old + close(blocked) + waitForFlusherIdle(t, s, key) + }) + + require.NoError(t, s.MarkComplete(key)) + timedOut := false + select { + case <-reached: + case <-time.After(5 * time.Second): + timedOut = true + } + require.False(t, timedOut, "timed out waiting for flusher worker to reach ioCopy") + return key, data +} + +func seedBothComplete(t *testing.T, s *Store) (key string, data []byte) { + key, data = createAndWrite(t, s, _testBlobSize) + require.NoError(t, s.MarkComplete(key)) + waitForFlushed(t, s, key) + return key, data +} + +func newStateOnlyMemIncomplete(t *testing.T) (s *Store, key string, data []byte) { + s, _ = newTestStore(t, 10*memsize.KB, 10*memsize.KB, 1) + key, data = seedOnlyMemIncomplete(t, s) + return s, key, data +} + +func newStateMemCompleteEnqueuedNotOnDisk(t *testing.T) (s *Store, key string, data []byte) { + s, _ = newTestStore(t, 10*memsize.KB, 10*memsize.KB, 1) + key, data = seedMemCompleteEnqueuedNotOnDisk(t, s) + return s, key, data +} + +func newStateMemCompleteDiskIncomplete(t *testing.T) (s *Store, key string, data []byte) { + s, _ = newTestStore(t, 10*memsize.KB, 10*memsize.KB, 1) + key, data = seedMemCompleteDiskIncomplete(t, s) + return s, key, data +} + +func newStateBothComplete(t *testing.T) (s *Store, key string, data []byte) { + s, _ = newTestStore(t, 10*memsize.KB, 10*memsize.KB, 1) + key, data = seedBothComplete(t, s) + return s, key, data +} + +// newStateDiskOnlyIncomplete uses a mem capacity of 1 byte so the real Create +// call overflows to disk via the genuine memory.ErrNoSpace fallback path. +func newStateDiskOnlyIncomplete(t *testing.T) (s *Store, key string, data []byte) { + s, _ = newTestStore(t, 10*memsize.KB, 1, 1) + key, data = createAndWrite(t, s, _testBlobSize) + return s, key, data +} + +func newStateDiskOnlyComplete(t *testing.T) (s *Store, key string, data []byte) { + s, key, data = newStateDiskOnlyIncomplete(t) + require.NoError(t, s.MarkComplete(key)) + return s, key, data +} + +type stateBuilder struct { + name string + build func(t *testing.T) (s *Store, key string, data []byte) +} + +func allStates() []stateBuilder { + return []stateBuilder{ + {"only in mem, incomplete", newStateOnlyMemIncomplete}, + {"complete in mem, enqueued, not yet on disk", newStateMemCompleteEnqueuedNotOnDisk}, + {"complete in mem, incomplete on disk", newStateMemCompleteDiskIncomplete}, + {"complete in both mem and disk", newStateBothComplete}, + {"only on disk, incomplete", newStateDiskOnlyIncomplete}, + {"only on disk, complete", newStateDiskOnlyComplete}, + } +} + +func TestConfig__applyDefaults(t *testing.T) { + t.Run("RebootIncompleteBlobs cannot be true", func(t *testing.T) { + require := require.New(t) + config := newTieredConfig(10*memsize.KB, 10*memsize.KB, 10, t.TempDir()) + config.DiskConfig.RebootIncompleteBlobs = true + + _, _, err := NewStore(config, tally.NoopScope) + require.EqualError(err, "tiered.Store does not support RebootIncompleteBlobs, as it can leak files. Use disk.Store if you need persistence for incomplete blobs") + }) + t.Run("NumFlushWorkers cannot be negative", func(t *testing.T) { + require := require.New(t) + config := newTieredConfig(10*memsize.KB, 10*memsize.KB, 10, t.TempDir()) + config.NumFlushWorkers = -1 + + _, _, err := NewStore(config, tally.NoopScope) + require.EqualError(err, "num_flush_workers must be at least 1, otherwise blobs will never get flushed from mem to disk") + }) + t.Run("NumFlushWorkers defaults to non-zero value", func(t *testing.T) { + require := require.New(t) + config := newTieredConfig(10*memsize.KB, 10*memsize.KB, 10, t.TempDir()) + config.NumFlushWorkers = 0 + + _, _, err := NewStore(config, tally.NoopScope) + require.NoError(err) + require.Equal(_defaultNumFlushWorkers, config.NumFlushWorkers) + }) +} + +func TestStore_APICorrectnessAcrossStates(t *testing.T) { + for _, sb := range allStates() { + t.Run(sb.name, func(t *testing.T) { + t.Run("Open", func(t *testing.T) { + require := require.New(t) + s, key, data := sb.build(t) + f, err := s.Open(key) + require.NoError(err) + got, err := io.ReadAll(f) + require.NoError(err) + require.Equal(data, got) + require.NoError(f.Close()) + }) + t.Run("Has", func(t *testing.T) { + require := require.New(t) + s, key, _ := sb.build(t) + inStore, _ := s.Has(key) + require.True(inStore) + }) + t.Run("Stat", func(t *testing.T) { + require := require.New(t) + s, key, data := sb.build(t) + size, err := s.Stat(key) + require.NoError(err) + require.Equal(int64(len(data)), size) + }) + t.Run("MarkComplete", func(t *testing.T) { + require := require.New(t) + s, key, _ := sb.build(t) + _, wasAlreadyComplete := s.ScopeComplete().Has(key) + require.NoError(s.MarkComplete(key)) + waitForComplete(t, s, key) + if !wasAlreadyComplete { + waitForFlushed(t, s, key) + } + require.NoError(s.MarkComplete(key)) + }) + t.Run("Delete", func(t *testing.T) { + require := require.New(t) + s, key, _ := sb.build(t) + require.NoError(s.Delete(key)) + _, ok := s.Has(key) + require.False(ok) + _, inMem := s.impl.mem.Has(key) + require.False(inMem) + _, inDisk := s.impl.disk.Has(key) + require.False(inDisk) + }) + t.Run("List", func(t *testing.T) { + require := require.New(t) + s, key, _ := sb.build(t) + require.Contains(s.List(), key) + }) + }) + } +} + +func TestStore_ScopingAcrossStates(t *testing.T) { + completeStates := []stateBuilder{ + {name: "complete in mem, enqueued, not yet on disk", build: newStateMemCompleteEnqueuedNotOnDisk}, + {name: "complete in mem, incomplete on disk", build: newStateMemCompleteDiskIncomplete}, + {name: "complete in both mem and disk", build: newStateBothComplete}, + {name: "only on disk, complete", build: newStateDiskOnlyComplete}, + } + incompleteStates := []stateBuilder{ + {name: "only in mem, incomplete", build: newStateOnlyMemIncomplete}, + {name: "only on disk, incomplete", build: newStateDiskOnlyIncomplete}, + } + + for _, sb := range completeStates { + t.Run(sb.name, func(t *testing.T) { + require := require.New(t) + s, key, data := sb.build(t) + + _, ok := s.ScopeComplete().Has(key) + require.True(ok) + _, ok = s.ScopeIncomplete().Has(key) + require.False(ok) + + f, err := s.ScopeComplete().Open(key) + require.NoError(err) + got, err := io.ReadAll(f) + require.NoError(err) + require.Equal(data, got) + require.NoError(f.Close()) + + _, err = s.ScopeIncomplete().Open(key) + require.ErrorIs(err, storelib.ErrOutOfScope) + + _, err = s.ScopeComplete().Stat(key) + require.NoError(err) + _, err = s.ScopeIncomplete().Stat(key) + require.ErrorIs(err, storelib.ErrOutOfScope) + + require.ErrorIs(s.ScopeIncomplete().Delete(key), storelib.ErrOutOfScope) + }) + } + + for _, sb := range incompleteStates { + t.Run(sb.name, func(t *testing.T) { + require := require.New(t) + s, key, data := sb.build(t) + + _, ok := s.ScopeIncomplete().Has(key) + require.True(ok) + _, ok = s.ScopeComplete().Has(key) + require.False(ok) + + f, err := s.ScopeIncomplete().Open(key) + require.NoError(err) + got, err := io.ReadAll(f) + require.NoError(err) + require.Equal(data, got) + require.NoError(f.Close()) + + _, err = s.ScopeComplete().Open(key) + require.ErrorIs(err, storelib.ErrOutOfScope) + + _, err = s.ScopeIncomplete().Stat(key) + require.NoError(err) + _, err = s.ScopeComplete().Stat(key) + require.ErrorIs(err, storelib.ErrOutOfScope) + + require.ErrorIs(s.ScopeComplete().Delete(key), storelib.ErrOutOfScope) + }) + } +} + +func TestStore_List_DedupDuringFlush(t *testing.T) { + require := require.New(t) + s, key, _ := newStateMemCompleteDiskIncomplete(t) + + require.NotContains(s.ScopeIncomplete().List(), key) + + all := s.List() + require.Len(all, 1) + require.Contains(all, key) + + require.Equal([]string{key}, s.ScopeComplete().List()) +} + +func TestStore_CrashRecovery(t *testing.T) { + require := require.New(t) + // 2 workers: states 2 and 3 each block a flusher worker via the seams. + diskCap, memCap, numWorkers := 10*memsize.KB, 10*memsize.KB, 2 + s, rootDir := newTestStore(t, diskCap, memCap, numWorkers) + + // Seed state 4 first while both workers are free. + bothCompleteKey, bothCompleteData := seedBothComplete(t, s) + + onlyMemIncompleteKey, _ := seedOnlyMemIncomplete(t, s) + enqueuedKey, _ := seedMemCompleteEnqueuedNotOnDisk(t, s) + memCompleteDiskIncompleteKey, _ := seedMemCompleteDiskIncomplete(t, s) + + // States 5 and 6 write directly to the disk sub-store. The real API's + // overflow-to-disk path (memory.ErrNoSpace) is tested by the matrix via + // newStateDiskOnlyIncomplete/Complete; here we bypass it because the + // shared store has enough mem capacity for the other states. + diskOnlyIncompleteKey := core.DigestFixture().Hex() + diskF, err := s.impl.disk.Create(diskOnlyIncompleteKey, _testBlobSize) + require.NoError(err) + require.NoError(diskF.Close()) + + diskOnlyCompleteKey := core.DigestFixture().Hex() + diskOnlyCompleteData := randomData(t, _testBlobSize) + diskF, err = s.impl.disk.Create(diskOnlyCompleteKey, _testBlobSize) + require.NoError(err) + _, err = diskF.Write(diskOnlyCompleteData) + require.NoError(err) + require.NoError(diskF.Close()) + require.NoError(s.impl.disk.MarkComplete(diskOnlyCompleteKey)) + + config := newTieredConfig(diskCap, memCap, 1, rootDir) + rebooted, _, err := NewStore(config, tally.NoopScope) + require.NoError(err) + + for _, discardedKey := range []string{ + onlyMemIncompleteKey, + enqueuedKey, + memCompleteDiskIncompleteKey, + diskOnlyIncompleteKey, + } { + _, ok := rebooted.Has(discardedKey) + require.False(ok, "key %q should have been discarded on restart", discardedKey) + } + + f, err := rebooted.Open(bothCompleteKey) + require.NoError(err) + got, err := io.ReadAll(f) + require.NoError(err) + require.Equal(bothCompleteData, got) + require.NoError(f.Close()) + + f, err = rebooted.Open(diskOnlyCompleteKey) + require.NoError(err) + got, err = io.ReadAll(f) + require.NoError(err) + require.Equal(diskOnlyCompleteData, got) + require.NoError(f.Close()) +} + +func TestStore_MetadataLifecycleAcrossStates(t *testing.T) { + for _, sb := range allStates() { + t.Run(sb.name, func(t *testing.T) { + require := require.New(t) + s, key, _ := sb.build(t) + + md := metadata.NewTorrentMeta(core.MetaInfoFixture()) + require.NoError(s.SetMetadata(key, md)) + + var readMd metadata.TorrentMeta + ok, err := s.GetMetadata(key, &readMd) + require.NoError(err) + require.True(ok) + require.Equal(md.MetaInfo, readMd.MetaInfo) + + require.NoError(s.DeleteMetadata(key, readMd.GetSuffix())) + ok, err = s.GetMetadata(key, &readMd) + require.NoError(err) + require.False(ok) + }) + } +} + +func TestStore_SetMetadataOnIncompleteBlob_IsNoOpForFlushing(t *testing.T) { + require := require.New(t) + s, _ := newTestStore(t, 10*memsize.KB, 10*memsize.KB, 1) + key, _ := createAndWrite(t, s, _testBlobSize) + + md := metadata.NewTorrentMeta(core.MetaInfoFixture()) + require.NoError(s.SetMetadata(key, md)) + + s.impl.flusher.mu.Lock() + _, tracked := s.impl.flusher.blobs[key] + s.impl.flusher.mu.Unlock() + require.False(tracked) + + _, inDisk := s.impl.disk.Has(key) + require.False(inDisk) +} diff --git a/origin/blobserver/server.go b/origin/blobserver/server.go index 7fd9c9743..3c84105ad 100644 --- a/origin/blobserver/server.go +++ b/origin/blobserver/server.go @@ -39,6 +39,7 @@ import ( "github.com/uber/kraken/lib/persistedretry" "github.com/uber/kraken/lib/persistedretry/writeback" "github.com/uber/kraken/lib/store" + "github.com/uber/kraken/lib/store/disk" "github.com/uber/kraken/lib/store/metadata" "github.com/uber/kraken/origin/blobclient" "github.com/uber/kraken/utils/closers" @@ -64,6 +65,7 @@ type Server struct { addr string hashRing hashring.Ring cas *store.CAStore + diskStore *disk.Store clientProvider blobclient.Provider clusterProvider blobclient.ClusterProvider backends *backend.Manager @@ -154,6 +156,7 @@ func (s *Server) Handler() http.Handler { r.Post("/namespace/{namespace}/blobs/{digest}/remote/{remote}", handler.Wrap(s.replicateToRemoteHandler)) r.Post("/forcecleanup", handler.Wrap(s.forceCleanupHandler)) + r.Post("/forcecleanup/v2", handler.Wrap(s.forceCleanupHandlerV2)) // Internal endpoints: @@ -1054,3 +1057,34 @@ func (s *Server) maybeDelete(name string, ttl time.Duration) (deleted bool, err } return false, nil } + +func (s *Server) forceCleanupHandlerV2(w http.ResponseWriter, r *http.Request) error { + if s.diskStore == nil { + return handler.Errorf("migration to disk.Store not yet done").Status(http.StatusNotImplemented) + } + + rawTargetUtilPercent := r.URL.Query().Get("target_util_percent") + if rawTargetUtilPercent == "" { + return handler.Errorf("query arg target_util_percent required").Status(http.StatusBadRequest) + } + targetUtilPercent, err := strconv.Atoi(rawTargetUtilPercent) + if err != nil { + return handler.Errorf("invalid target_util_percent: %s", err).Status(http.StatusBadRequest) + } + rawRespectEvictionBan := r.URL.Query().Get("respect_eviction_ban") + if rawRespectEvictionBan == "" { + return handler.Errorf("query arg respect_eviction_ban required").Status(http.StatusBadRequest) + } + respectEvictionBan, err := strconv.ParseBool(rawRespectEvictionBan) + if err != nil { + return handler.Errorf("invalid respect_eviction_ban: %s", err).Status(http.StatusBadRequest) + } + + newUtil, err := s.diskStore.Clean(targetUtilPercent, respectEvictionBan) + if err != nil { + return handler.Errorf("Encountered an error while cleaning. New utilization is '%v'. Error: %s", newUtil, err) + } + return json.NewEncoder(w).Encode(map[string]any{ + "new_util": newUtil, + }) +} diff --git a/origin/blobserver/server_test.go b/origin/blobserver/server_test.go index 3237fcf13..1c877a472 100644 --- a/origin/blobserver/server_test.go +++ b/origin/blobserver/server_test.go @@ -16,6 +16,7 @@ package blobserver import ( "bytes" "context" + "encoding/json" "errors" "fmt" "io" @@ -26,12 +27,14 @@ import ( "github.com/golang/mock/gomock" "github.com/stretchr/testify/require" + "github.com/uber-go/tally" "github.com/uber/kraken/core" "github.com/uber/kraken/lib/backend" "github.com/uber/kraken/lib/backend/backenderrors" "github.com/uber/kraken/lib/persistedretry" "github.com/uber/kraken/lib/persistedretry/writeback" + "github.com/uber/kraken/lib/store/disk" "github.com/uber/kraken/lib/store/metadata" "github.com/uber/kraken/origin/blobclient" "github.com/uber/kraken/utils/httputil" @@ -871,3 +874,77 @@ func TestForceCleanupWriteBackFailures(t *testing.T) { ensureHasBlob(t, client, namespace, blob) } + +func newTestDiskStore(t *testing.T) *disk.Store { + d, err := disk.NewStore(&disk.Config{ + CapacityBytes: 100, + RootDir: t.TempDir(), + ShardLength: 2, + }, tally.NoopScope) + require.NoError(t, err) + return d +} + +func TestForceCleanupV2MigrationNotDone(t *testing.T) { + require := require.New(t) + + cp := newTestClientProvider() + s := newTestServer(t, master1, hashRingMaxReplica(), cp) + defer s.cleanup() + + _, err := httputil.Post(fmt.Sprintf( + "http://%s/forcecleanup/v2?target_util_percent=50&respect_eviction_ban=true", s.addr)) + require.True(httputil.IsStatus(err, http.StatusNotImplemented)) +} + +func TestForceCleanupV2InvalidParams(t *testing.T) { + for _, tc := range []struct { + name string + query string + }{ + {"missing target_util_percent", "respect_eviction_ban=true"}, + {"invalid target_util_percent", "target_util_percent=abc&respect_eviction_ban=true"}, + {"missing respect_eviction_ban", "target_util_percent=50"}, + {"invalid respect_eviction_ban", "target_util_percent=50&respect_eviction_ban=abc"}, + } { + t.Run(tc.name, func(t *testing.T) { + require := require.New(t) + + cp := newTestClientProvider() + s := newTestServer(t, master1, hashRingMaxReplica(), cp) + defer s.cleanup() + s.setDiskStore(newTestDiskStore(t)) + + _, err := httputil.Post(fmt.Sprintf("http://%s/forcecleanup/v2?%s", s.addr, tc.query)) + require.True(httputil.IsStatus(err, http.StatusBadRequest)) + }) + } +} + +func TestForceCleanupV2(t *testing.T) { + require := require.New(t) + + cp := newTestClientProvider() + s := newTestServer(t, master1, hashRingMaxReplica(), cp) + defer s.cleanup() + + diskStore := newTestDiskStore(t) + s.setDiskStore(diskStore) + + key := core.DigestFixture().Hex() + f, err := diskStore.Create(key, 60) + require.NoError(err) + require.NoError(f.Close()) + require.NoError(diskStore.MarkComplete(key)) + + // Target of 10% cannot be met while keeping the only blob (60% util), so it gets evicted. + resp, err := httputil.Post(fmt.Sprintf( + "http://%s/forcecleanup/v2?target_util_percent=10&respect_eviction_ban=true", s.addr)) + require.NoError(err) + defer func() { require.NoError(resp.Body.Close()) }() + + var body map[string]int + require.NoError(json.NewDecoder(resp.Body).Decode(&body)) + require.Equal(0, body["new_util"]) + require.Empty(diskStore.List()) +} diff --git a/origin/blobserver/testutils_test.go b/origin/blobserver/testutils_test.go index 32209d0e7..4f78014a4 100644 --- a/origin/blobserver/testutils_test.go +++ b/origin/blobserver/testutils_test.go @@ -34,6 +34,7 @@ import ( "github.com/uber/kraken/lib/hostlist" "github.com/uber/kraken/lib/metainfogen" "github.com/uber/kraken/lib/store" + "github.com/uber/kraken/lib/store/disk" mockbackend "github.com/uber/kraken/mocks/lib/backend" mockpersistedretry "github.com/uber/kraken/mocks/lib/persistedretry" mockblobclient "github.com/uber/kraken/mocks/origin/blobclient" @@ -104,6 +105,11 @@ type testServer struct { writeBackManager *mockpersistedretry.MockManager clk *clock.Mock cleanup func() + + // setDiskStore configures the Server's disk.Store. There's no way to inject + // this via New yet (migration to disk.Store is in progress), so this closure + // captures the Server built below to allow tests to still wire it in. + setDiskStore func(*disk.Store) } func newTestServer( @@ -157,6 +163,7 @@ func newTestServer( writeBackManager: writeBackManager, clk: clk, cleanup: cleanup.Run, + setDiskStore: func(d *disk.Store) { s.diskStore = d }, } }