diff --git a/lib/store/disk/config.go b/lib/store/disk/config.go index 3f3a0e75a..d021c6c74 100644 --- a/lib/store/disk/config.go +++ b/lib/store/disk/config.go @@ -12,5 +12,5 @@ type Config struct { // 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 + ShardLength int // TODO - add a mechanism to migrate in-place from 1 shardLength to another. } diff --git a/lib/store/disk/file.go b/lib/store/disk/file.go new file mode 100644 index 000000000..dc759ad3f --- /dev/null +++ b/lib/store/disk/file.go @@ -0,0 +1,47 @@ +package disk + +import ( + storelib "github.com/uber/kraken/lib/store" + "os" +) + +func newFile(f *os.File) *File { + return &File{ + fd: f, + } +} + +var _ storelib.FileReadWriter = &File{} + +// File represends an open file descriptor to a blob in [Store]. +type File struct { + fd *os.File +} + +func (f *File) Read(p []byte) (n int, err error) { return f.fd.Read(p) } +func (f *File) ReadAt(p []byte, off int64) (n int, err error) { return f.fd.ReadAt(p, off) } +func (f *File) Seek(off int64, whence int) (int64, error) { return f.fd.Seek(off, whence) } +func (f *File) Write(p []byte) (n int, err error) { return f.fd.Write(p) } +func (f *File) WriteAt(p []byte, off int64) (n int, err error) { return f.fd.WriteAt(p, off) } +func (f *File) Close() error { return f.fd.Close() } + +// Size returns the number of bytes the file contains. +func (f *File) Size() int64 { + info, err := f.fd.Stat() + if err != nil { + return 0 + } + return info.Size() +} + +// Cancel is supposed to remove any written content. +// In this implementation file is not actually removed, but rather closed, which is fine, as there won't be key collisions when creating new files. +func (f *File) Cancel() error { + return f.fd.Close() +} + +// Commit is supposed to flush all content for buffered writer. +func (f *File) Commit() error { + // TODO - consider whether we should do f.Sync() before closing the file. + return f.fd.Close() +} diff --git a/lib/store/disk/scoped_store.go b/lib/store/disk/scoped_store.go index 75cbf5389..d137d4004 100644 --- a/lib/store/disk/scoped_store.go +++ b/lib/store/disk/scoped_store.go @@ -1,7 +1,6 @@ package disk import ( - "errors" "os" "github.com/uber-go/tally" @@ -9,10 +8,6 @@ import ( "github.com/uber/kraken/lib/store/metadata" ) -// ErrOutOfScope is returned when the provided key is in the store, but not in the store's [blobScope], -// e.g. if the blob is incomplete, but ScopeComplete was called. -var ErrOutOfScope = errors.New("the blob is in Store but filtered by the selected scope") - // Store is a key-value, persistent, thread-safe, LRU store for blobs and their [metadata.Metadata]. // // - Supports pagination of blobs during reading/writing, such that blobs don't need to be fully loaded into memory. @@ -28,7 +23,7 @@ var ErrOutOfScope = errors.New("the blob is in Store but filtered by the selecte // - Supports directory sharding to speed up disk performance. type Store struct { impl *store - scope blobScope + scope storelib.BlobScope } // NewStore initializes a new [*Store]. If the store has been initialized in the same @@ -45,19 +40,22 @@ func NewStore(config *Config, metrics tally.Scope) (*Store, error) { } return &Store{ impl: s, - scope: blobScopeAny, + scope: storelib.BlobScopeAny, }, nil } // Create adds a new, incomplete blob to the store and reserves space for 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. -func (s *Store) Create(key string, sizeBytes uint64) (storelib.FileReadWriter, error) { +func (s *Store) Create(key string, sizeBytes uint64) (*File, error) { return s.impl.Create(key, sizeBytes) } // Open returns an FD to a file in the store. [os.ErrNotExist] is returned on missing entry. -func (s *Store) Open(key string) (storelib.FileReadWriter, error) { return s.impl.Open(key, s.scope) } +func (s *Store) Open(key string) (*File, error) { return s.impl.Open(key, s.scope) } + +// Has checks if the blob is in the store. +func (s *Store) Has(key string) (inStore bool, inScope bool) { return s.impl.Has(key, s.scope) } // Stat returns [os.FileInfo] about the blob. Returns [os.ErrNotExist] if the blob is not found. func (s *Store) Stat(key string) (os.FileInfo, error) { return s.impl.Stat(key, s.scope) } @@ -66,7 +64,7 @@ func (s *Store) Stat(key string) (os.FileInfo, error) { return s.impl.Stat(key, // Additionally, other store APIs may filter blobs based on completeness. func (s *Store) MarkComplete(key string) error { return s.impl.MarkComplete(key) } -// Delete removes a blob and its [metadata.Metadata] from the store. +// Delete removes a blob and its [metadata.Metadata] from the store. Returns [os.ErrNotExist] on missing blob. 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). @@ -89,7 +87,7 @@ func (s *Store) GetMetadata(key string, md metadata.Metadata) (ok bool, err erro return s.impl.GetMetadata(key, md, s.scope) } -// DeleteMetadata removes any metadata of a blob with `md`'s suffix, if present. +// 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) } @@ -105,9 +103,9 @@ func (s *Store) WriteAtMetadata(key string, md metadata.Metadata, p []byte, off } // ScopeComplete scopes [Store]'s APIs such that they can only operate on complete blobs (except MarkComplete and Create). -// [ErrOutOfScope] is returned if the user tries to operate on an incomplete blob. -func (s *Store) ScopeComplete() *Store { return &Store{s.impl, blobScopeComplete} } +// [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. -// [ErrOutOfScope] is returned if the user tries to operate on a complete blob. -func (s *Store) ScopeIncomplete() *Store { return &Store{s.impl, blobScopeIncomplete} } +// [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/disk/store.go b/lib/store/disk/store.go index 6c3ba7ada..5ecccfa26 100644 --- a/lib/store/disk/store.go +++ b/lib/store/disk/store.go @@ -19,16 +19,6 @@ import ( "go.uber.org/zap" ) -// the set of blobs that the store's APIs can operate on. -type blobScope int - -// flags to scope [store]'s APIs to a subset of blobs. -const ( - blobScopeAny blobScope = iota - blobScopeComplete - blobScopeIncomplete -) - const ( _completeBlob = true _incompleteBlob = false @@ -39,7 +29,7 @@ const ( var _syncEvictionLatencyBuckets = tally.MustMakeExponentialDurationBuckets(100*time.Millisecond, 1.4, 15) -// store implements the APIs of [Store]. [store]'s APIs expose the [blobScope] arg, +// 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. // // Check [Store]'s comments for details on functionality. @@ -73,7 +63,7 @@ func newStore(config *Config, metrics tally.Scope) (*store, error) { } if !ok { log.Info("Initialized a new, empty Store (did not find any previously persisted state to reboot for Store)") - return &store{ + store := &store{ capacity: config.CapacityBytes, size: 0, blobs: make(map[string]*blob), @@ -82,7 +72,10 @@ func newStore(config *Config, metrics tally.Scope) (*store, error) { pather: newPather(config.RootDir, config.ShardLength), log: log, metrics: metrics, - }, nil + } + + store.emitUsageMetrics() + return store, nil } store, err := rebootPersistedStore(config, log, metrics) @@ -92,10 +85,12 @@ func newStore(config *Config, metrics tally.Scope) (*store, error) { return nil, err } log.With("num_blobs", len(store.blobs)).Info("Successfully rebooted Store's previously left state on disk") + + store.emitUsageMetrics() return store, nil } -func (s *store) Open(key string, scope blobScope) (storelib.FileReadWriter, error) { +func (s *store) Open(key string, scope storelib.BlobScope) (*File, error) { s.mu.Lock() defer s.mu.Unlock() @@ -115,10 +110,24 @@ func (s *store) Open(key string, scope blobScope) (storelib.FileReadWriter, erro if err != nil { return nil, fmt.Errorf("open: %w", err) } - return storelib.NewReadWriter(f), nil + return newFile(f), nil } -func (s *store) Stat(key string, scope blobScope) (os.FileInfo, error) { +func (s *store) Has(key string, scope storelib.BlobScope) (inStore bool, inScope bool) { + s.mu.RLock() + defer s.mu.RUnlock() + + b, ok := s.blobs[key] + if !ok { + return false, false + } + if err := isOutOfScope(b, scope); err != nil { + return true, false + } + return true, true +} + +func (s *store) Stat(key string, scope storelib.BlobScope) (os.FileInfo, error) { s.mu.RLock() defer s.mu.RUnlock() @@ -134,7 +143,7 @@ func (s *store) Stat(key string, scope blobScope) (os.FileInfo, error) { return os.Stat(blobPath) } -func (s *store) Create(key string, sizeBytes uint64) (storelib.FileReadWriter, error) { +func (s *store) Create(key string, sizeBytes uint64) (*File, error) { // TODO - we might want some TTI on uploads to the store, after which we cancel the upload, e.g. 1min without the client uploading more data. s.mu.Lock() defer s.mu.Unlock() @@ -180,7 +189,8 @@ func (s *store) Create(key string, sizeBytes uint64) (storelib.FileReadWriter, e evictionBanned: false, } - return storelib.NewReadWriter(f), nil + s.emitUsageMetrics() + return newFile(f), nil } func (s *store) persistBlobSize(key string, sizeBytes uint64) error { @@ -231,7 +241,7 @@ func (s *store) reserveSpace(space uint64) error { func (s *store) releaseSpace(space uint64) { if space > s.size { - s.log.Error("Invariant violation - Store wants to release more disk space than actually reserved. Failing open by releasing all reserved space.") + s.log.Error("Invariant violation - disk.Store wants to release more space than actually reserved. Failing open by setting store.size = 0") s.size = 0 return } @@ -281,7 +291,7 @@ func (s *store) MarkComplete(key string) error { // to avoid an inconsistent state, as failure only costs a negligible amount // of disk until the blob is evicted. func (s *store) tryDeleteImmovableMetadata(key string) { - mdList, err := s.listMetadataNoLock(key, blobScopeAny) + mdList, err := s.listMetadataNoLock(key, storelib.BlobScopeAny) if err != nil { err = fmt.Errorf("list metadata: %w", err) s.log.With("error", err).Error("Failed to delete un-movable metadata upon marking a blob as complete") @@ -310,7 +320,7 @@ func (s *store) checkDiskIfUnevictable(key string, complete bool) (bool, error) return unevictable, nil } -func (s *store) Delete(key string, scope blobScope) error { +func (s *store) Delete(key string, scope storelib.BlobScope) error { s.mu.Lock() defer s.mu.Unlock() @@ -332,10 +342,11 @@ func (s *store) Delete(key string, scope blobScope) error { delete(s.blobs, key) s.releaseSpace(b.size) + s.emitUsageMetrics() return nil } -func (s *store) list(scope blobScope) []string { +func (s *store) list(scope storelib.BlobScope) []string { s.mu.RLock() defer s.mu.RUnlock() @@ -349,7 +360,7 @@ func (s *store) list(scope blobScope) []string { return res } -func (s *store) BanEviction(key string, scope blobScope) error { +func (s *store) BanEviction(key string, scope storelib.BlobScope) error { s.mu.Lock() defer s.mu.Unlock() @@ -381,7 +392,7 @@ func (s *store) BanEviction(key string, scope blobScope) error { return nil } -func (s *store) UnbanEviction(key string, scope blobScope) error { +func (s *store) UnbanEviction(key string, scope storelib.BlobScope) error { s.mu.Lock() defer s.mu.Unlock() @@ -411,7 +422,7 @@ func (s *store) UnbanEviction(key string, scope blobScope) error { return nil } -func (s *store) SetMetadata(key string, md metadata.Metadata, scope blobScope) error { +func (s *store) SetMetadata(key string, md metadata.Metadata, scope storelib.BlobScope) error { s.mu.Lock() defer s.mu.Unlock() @@ -446,7 +457,7 @@ func (s *store) SetMetadata(key string, md metadata.Metadata, scope blobScope) e return nil } -func (s *store) GetMetadata(key string, md metadata.Metadata, scope blobScope) (ok bool, err error) { +func (s *store) GetMetadata(key string, md metadata.Metadata, scope storelib.BlobScope) (ok bool, err error) { s.mu.RLock() defer s.mu.RUnlock() @@ -478,7 +489,8 @@ func (s *store) GetMetadata(key string, md metadata.Metadata, scope blobScope) ( return true, nil } -func (s *store) DeleteMetadata(key string, md metadata.Metadata, scope blobScope) error { +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.` s.mu.Lock() defer s.mu.Unlock() @@ -501,14 +513,14 @@ func (s *store) DeleteMetadata(key string, md metadata.Metadata, scope blobScope return nil } -func (s *store) ListMetadata(key string, scope blobScope) ([]metadata.Metadata, error) { +func (s *store) ListMetadata(key string, scope storelib.BlobScope) ([]metadata.Metadata, error) { s.mu.RLock() defer s.mu.RUnlock() return s.listMetadataNoLock(key, scope) } -func (s *store) listMetadataNoLock(key string, scope blobScope) ([]metadata.Metadata, error) { +func (s *store) listMetadataNoLock(key string, scope storelib.BlobScope) ([]metadata.Metadata, error) { b, ok := s.blobs[key] if !ok { return nil, os.ErrNotExist @@ -538,7 +550,7 @@ func (s *store) listMetadataNoLock(key string, scope blobScope) ([]metadata.Meta return res, nil } -func (s *store) WriteAtMetadata(key string, md metadata.Metadata, p []byte, off int64, scope blobScope) error { +func (s *store) WriteAtMetadata(key string, md metadata.Metadata, p []byte, off int64, scope storelib.BlobScope) error { s.mu.Lock() defer s.mu.Unlock() @@ -592,9 +604,14 @@ func exists(path string) (ok bool, err error) { return false, fmt.Errorf("stat: %w", err) } -func isOutOfScope(b *blob, scope blobScope) error { - if (b.complete && scope == blobScopeIncomplete) || (!b.complete && scope == blobScopeComplete) { - return ErrOutOfScope +func isOutOfScope(b *blob, scope storelib.BlobScope) error { + if (b.complete && scope == storelib.BlobScopeIncomplete) || (!b.complete && scope == storelib.BlobScopeComplete) { + return storelib.ErrOutOfScope } return nil } + +func (s *store) emitUsageMetrics() { + s.metrics.Gauge("num_entries").Update(float64(len(s.blobs))) + s.metrics.Gauge("size_bytes").Update(float64(s.size)) +} diff --git a/lib/store/disk/store_test.go b/lib/store/disk/store_test.go index fcbcf85a6..4e48e00c7 100644 --- a/lib/store/disk/store_test.go +++ b/lib/store/disk/store_test.go @@ -2,6 +2,7 @@ package disk import ( "bytes" + "crypto/rand" "io" "io/fs" "os" @@ -95,7 +96,7 @@ func TestStore(t *testing.T) { require.Nil(f) f, err = store.ScopeComplete().Open(keys[0]) - require.ErrorIs(err, ErrOutOfScope) + require.ErrorIs(err, storelib.ErrOutOfScope) require.Nil(f) require.NoError(store.MarkComplete(keys[0])) @@ -151,8 +152,8 @@ func TestEviction(t *testing.T) { f, fKey := newTestFile(t, store, 4*memsize.KB) require.NoError(f.Close()) require.NoError(store.MarkComplete(fKey)) - keys := store.List() - require.NotContains(keys, cKey) + _, ok := store.Has(cKey) + require.False(ok) // new size == 23KB == 24KB - 5KB (c) + 4KB (f) require.Equal(23*memsize.KB, store.impl.size) require.Equal(5, numBlobsOnDisk(t, store)) @@ -171,9 +172,10 @@ func TestEviction(t *testing.T) { // size == 24KB == 24KB + 15KB (h) - 5KB (b) - 10KB (a) require.Equal(24*memsize.KB, store.impl.size) require.Equal(5, numBlobsOnDisk(t, store)) - keys = store.List() - require.NotContains(keys, bKey) - require.NotContains(keys, aKey) + _, ok = store.Has(bKey) + require.False(ok) + _, ok = store.Has(aKey) + require.False(ok) // allow e to be evicted. require.NoError(store.MarkComplete(eKey)) @@ -187,8 +189,8 @@ func TestEviction(t *testing.T) { i, iKey := newTestFile(t, store, 5*memsize.KB) require.NoError(store.MarkComplete(iKey)) require.NoError(i.Close()) - keys = store.List() - require.NotContains(keys, fKey) + _, ok = store.Has(fKey) + require.False(ok) require.Equal(25*memsize.KB, store.impl.size) require.Equal(5, numBlobsOnDisk(t, store)) // eviction order: h(15KB), e(1KB), g(1KB), i(5KB); d(3KB) is unevictable @@ -196,8 +198,8 @@ func TestEviction(t *testing.T) { j, jKey := newTestFile(t, store, 14*memsize.KB) require.NoError(j.Close()) require.NoError(store.MarkComplete(jKey)) - keys = store.List() - require.NotContains(keys, hKey) + _, ok = store.Has(hKey) + require.False(ok) require.Equal(24*memsize.KB, store.impl.size) require.Equal(5, numBlobsOnDisk(t, store)) // eviction order: e(1KB), g(1KB), i(5KB), j(14KB); d(3KB) is unevictable @@ -205,8 +207,8 @@ func TestEviction(t *testing.T) { k, kKey := newTestFile(t, store, 2*memsize.KB) require.NoError(k.Close()) require.NoError(store.MarkComplete(kKey)) - keys = store.List() - require.NotContains(keys, eKey) + _, ok = store.Has(eKey) + require.False(ok) require.Equal(25*memsize.KB, store.impl.size) require.Equal(5, numBlobsOnDisk(t, store)) // eviction order: g(1KB), i(5KB), j(14KB), k(2KB); d(3KB) is unevictable @@ -214,8 +216,8 @@ func TestEviction(t *testing.T) { l, lKey := newTestFile(t, store, 1*memsize.KB) require.NoError(store.MarkComplete(lKey)) require.NoError(l.Close()) - keys = store.List() - require.NotContains(keys, gKey) + _, ok = store.Has(gKey) + require.False(ok) require.Equal(25*memsize.KB, store.impl.size) require.Equal(5, numBlobsOnDisk(t, store)) // eviction order: i(5KB), j(14KB), k(2KB), l(1KB); d(3KB) is unevictable @@ -335,6 +337,83 @@ func TestOpenedFileAccessibleAfterMarkedComplete(t *testing.T) { require.Equal([]byte("Hello World"), completeFileData) } +func TestMisreportedBlobSize(t *testing.T) { + t.Run("user overreports size", func(t *testing.T) { + require := require.New(t) + store, _ := newTestStore(t, 25*memsize.KB, true) + + // We pad the store's size to 1KB before all other operations to later show that size does not underflow when releasing space. + _, _ = newTestFile(t, store, 1*memsize.KB) + + f, key := newTestFile(t, store, 10*memsize.KB) + data := make([]byte, 5*memsize.KB) + _, err := rand.Read(data) + require.NoError(err) + _, err = io.Copy(f, bytes.NewReader(data)) + require.NoError(err) + require.NoError(f.Close()) + + f, err = store.Open(key) + require.NoError(err) + + // io.ReadAll only returns the 5KB of actual data, despite us reporting the blob as 10KB. There is no 5KB of trailing space. + readData, err := io.ReadAll(f) + require.NoError(err) + require.Equal(data, readData) + + // Stat returns the actual size of the blob. + fi, err := store.Stat(key) + require.NoError(err) + require.Equal(uint64(fi.Size()), 5*memsize.KB) + + // The store always calculates its size based off the client-reported size, not the actual size. + require.Equal(11*memsize.KB, store.impl.size) + + // Both space reservation and space release use the user-reported value, not the actual size of the blob. + // Due to this consistency, the store's size is the same before adding and after deleting/evicting the blob. + require.NoError(store.Delete(key)) + require.Equal(1*memsize.KB, store.impl.size) + }) + + t.Run("user underreports size", func(t *testing.T) { + require := require.New(t) + store, _ := newTestStore(t, 25*memsize.KB, true) + + // We pad the store's size to 1KB before all other operations to later show that size does not underflow when releasing space. + _, _ = newTestFile(t, store, 1*memsize.KB) + + // 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) + require.NoError(err) + _, err = io.Copy(f, bytes.NewReader(data)) + require.NoError(err) + require.NoError(f.Close()) + + f, err = store.Open(key) + require.NoError(err) + + // io.ReadAll only returns the 30KB of actual data, despite us reporting the blob as 24KB. The file is resize-able. + readData, err := io.ReadAll(f) + require.NoError(err) + require.Equal(data, readData) + + // Stat returns the actual size of the blob. + fi, err := store.Stat(key) + require.NoError(err) + require.Equal(uint64(fi.Size()), 30*memsize.KB) + + // The store always calculates its size based off the client-reported size, not the actual size. + require.Equal(25*memsize.KB, store.impl.size) + + // Both space reservation and space release use the user-reported value, not the actual size of the blob. + // Due to this consistency, the store's size is the same before adding and after deleting/evicting the blob. + require.NoError(store.Delete(key)) + require.Equal(1*memsize.KB, store.impl.size) + }) +} + func TestDelete(t *testing.T) { t.Run("incomplete blob", func(t *testing.T) { require := require.New(t) @@ -543,7 +622,7 @@ func TestStat(t *testing.T) { require.NoError(err) _, err = store.ScopeComplete().Stat(key) - require.Equal(ErrOutOfScope, err) + require.Equal(storelib.ErrOutOfScope, err) fInfo, err := store.Stat(key) require.NoError(err) @@ -566,7 +645,7 @@ func TestStat(t *testing.T) { require.NoError(store.BanEviction(key)) _, err = store.ScopeComplete().Stat(key) - require.Equal(ErrOutOfScope, err) + require.Equal(storelib.ErrOutOfScope, err) fInfo, err := store.Stat(key) require.NoError(err) @@ -745,7 +824,7 @@ func TestMetadata(t *testing.T) { var readMd metadata.TorrentMeta ok, err := store.ScopeComplete().GetMetadata(key, &readMd) - require.Equal(ErrOutOfScope, err) + require.Equal(storelib.ErrOutOfScope, err) // incomplete files are ignored require.False(ok) @@ -757,7 +836,7 @@ func TestMetadata(t *testing.T) { // Repeat the tests above for an unevictable file require.NoError(store.BanEviction(key)) ok, err = store.ScopeComplete().GetMetadata(key, &readMd) - require.Equal(ErrOutOfScope, err) + require.Equal(storelib.ErrOutOfScope, err) // incomplete files are ignored require.False(ok) @@ -888,21 +967,23 @@ func TestScopes(t *testing.T) { // While incomplete, ScopeComplete's APIs reject the blob. _, err = store.ScopeComplete().Open(key) - require.ErrorIs(err, ErrOutOfScope) + require.ErrorIs(err, storelib.ErrOutOfScope) _, err = store.ScopeComplete().Stat(key) - require.ErrorIs(err, ErrOutOfScope) - require.ErrorIs(store.ScopeComplete().BanEviction(key), ErrOutOfScope) - require.ErrorIs(store.ScopeComplete().UnbanEviction(key), ErrOutOfScope) - require.ErrorIs(store.ScopeComplete().SetMetadata(key, md), ErrOutOfScope) + require.ErrorIs(err, storelib.ErrOutOfScope) + require.ErrorIs(store.ScopeComplete().BanEviction(key), storelib.ErrOutOfScope) + require.ErrorIs(store.ScopeComplete().UnbanEviction(key), storelib.ErrOutOfScope) + require.ErrorIs(store.ScopeComplete().SetMetadata(key, md), storelib.ErrOutOfScope) var readMd metadata.TorrentMeta ok, err := store.ScopeComplete().GetMetadata(key, &readMd) - require.ErrorIs(err, ErrOutOfScope) + require.ErrorIs(err, storelib.ErrOutOfScope) require.False(ok) _, err = store.ScopeComplete().ListMetadata(key) - require.ErrorIs(err, ErrOutOfScope) - require.ErrorIs(store.ScopeComplete().WriteAtMetadata(key, md, mdData, 0), ErrOutOfScope) - require.ErrorIs(store.ScopeComplete().DeleteMetadata(key, &readMd), ErrOutOfScope) + 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.NotContains(store.ScopeComplete().List(), key) + _, ok = store.ScopeComplete().Has(key) + require.False(ok) // The unscoped store's APIs work regardless of completeness. f, err = store.Open(key) @@ -927,6 +1008,8 @@ func TestScopes(t *testing.T) { require.NoError(store.WriteAtMetadata(key, md, mdData, 0)) require.NoError(store.DeleteMetadata(key, &readMd)) require.Contains(store.List(), key) + _, ok = store.Has(key) + require.True(ok) // ScopeIncomplete's APIs also work, since the blob is (still) incomplete. f, err = store.ScopeIncomplete().Open(key) @@ -951,25 +1034,32 @@ func TestScopes(t *testing.T) { require.NoError(store.ScopeIncomplete().WriteAtMetadata(key, md, mdData, 0)) require.NoError(store.ScopeIncomplete().DeleteMetadata(key, &readMd)) require.Contains(store.ScopeIncomplete().List(), key) + _, ok = store.ScopeIncomplete().Has(key) + require.True(ok) + require.Contains(store.List(), key) + _, ok = store.Has(key) + require.True(ok) require.NoError(store.MarkComplete(key)) // Now that the blob is complete, the roles reverse: ScopeIncomplete rejects it. _, err = store.ScopeIncomplete().Open(key) - require.ErrorIs(err, ErrOutOfScope) + require.ErrorIs(err, storelib.ErrOutOfScope) _, err = store.ScopeIncomplete().Stat(key) - require.ErrorIs(err, ErrOutOfScope) - require.ErrorIs(store.ScopeIncomplete().BanEviction(key), ErrOutOfScope) - require.ErrorIs(store.ScopeIncomplete().UnbanEviction(key), ErrOutOfScope) - require.ErrorIs(store.ScopeIncomplete().SetMetadata(key, md), ErrOutOfScope) + require.ErrorIs(err, storelib.ErrOutOfScope) + require.ErrorIs(store.ScopeIncomplete().BanEviction(key), storelib.ErrOutOfScope) + require.ErrorIs(store.ScopeIncomplete().UnbanEviction(key), storelib.ErrOutOfScope) + require.ErrorIs(store.ScopeIncomplete().SetMetadata(key, md), storelib.ErrOutOfScope) ok, err = store.ScopeIncomplete().GetMetadata(key, &readMd) - require.ErrorIs(err, ErrOutOfScope) + require.ErrorIs(err, storelib.ErrOutOfScope) require.False(ok) _, err = store.ScopeIncomplete().ListMetadata(key) - require.ErrorIs(err, ErrOutOfScope) - require.ErrorIs(store.ScopeIncomplete().WriteAtMetadata(key, md, mdData, 0), ErrOutOfScope) - require.ErrorIs(store.ScopeIncomplete().DeleteMetadata(key, &readMd), ErrOutOfScope) + 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.NotContains(store.ScopeIncomplete().List(), key) + _, ok = store.ScopeIncomplete().Has(key) + require.False(ok) // The unscoped store's APIs still work regardless of completeness. f, err = store.Open(key) @@ -994,6 +1084,8 @@ func TestScopes(t *testing.T) { require.NoError(store.WriteAtMetadata(key, md, mdData, 0)) require.NoError(store.DeleteMetadata(key, &readMd)) require.Contains(store.List(), key) + _, ok = store.Has(key) + require.True(ok) // ScopeComplete's APIs now work, since the blob is complete. f, err = store.ScopeComplete().Open(key) @@ -1018,6 +1110,8 @@ func TestScopes(t *testing.T) { require.NoError(store.ScopeComplete().WriteAtMetadata(key, md, mdData, 0)) require.NoError(store.ScopeComplete().DeleteMetadata(key, &readMd)) require.Contains(store.ScopeComplete().List(), key) + _, ok = store.ScopeComplete().Has(key) + require.True(ok) } func TestScopesDelete(t *testing.T) { @@ -1026,7 +1120,7 @@ func TestScopesDelete(t *testing.T) { f, key := newTestFile(t, store, 1*memsize.KB) require.NoError(f.Close()) - require.ErrorIs(store.ScopeComplete().Delete(key), ErrOutOfScope) + require.ErrorIs(store.ScopeComplete().Delete(key), storelib.ErrOutOfScope) require.NoError(store.ScopeIncomplete().Delete(key)) _, err := store.Stat(key) require.ErrorIs(err, os.ErrNotExist) @@ -1034,7 +1128,7 @@ func TestScopesDelete(t *testing.T) { f, key = newTestFile(t, store, 1*memsize.KB) require.NoError(f.Close()) require.NoError(store.MarkComplete(key)) - require.ErrorIs(store.ScopeIncomplete().Delete(key), ErrOutOfScope) + require.ErrorIs(store.ScopeIncomplete().Delete(key), storelib.ErrOutOfScope) require.NoError(store.ScopeComplete().Delete(key)) _, err = store.Stat(key) require.ErrorIs(err, os.ErrNotExist) diff --git a/lib/store/file.go b/lib/store/file.go index 1e10a941c..596342ecf 100644 --- a/lib/store/file.go +++ b/lib/store/file.go @@ -14,8 +14,6 @@ package store import ( - "os" - "github.com/uber/kraken/lib/store/base" ) @@ -24,37 +22,3 @@ type FileReadWriter = base.FileReadWriter // FileReader is a read-only file. type FileReader = base.FileReader - -func NewReadWriter(f *os.File) FileReadWriter { - return &rwImpl{ - File: f, - } -} - -var _ FileReadWriter = &rwImpl{} - -type rwImpl struct { - *os.File -} - -// Size returns the number of bytes the file contains. -func (i *rwImpl) Size() int64 { - info, err := i.Stat() - if err != nil { - return 0 - } - return info.Size() -} - -// Cancel is supposed to remove any written content. -// In this implementation file is not actually removed, and it's fine since there won't be name -// collision between upload files. -func (i *rwImpl) Cancel() error { - return i.Close() -} - -// Commit is supposed to flush all content for buffered writer. -// In this implementation all writes write to the file directly through syscall. -func (i *rwImpl) Commit() error { - return i.Close() -} diff --git a/lib/store/memory/file.go b/lib/store/memory/file.go new file mode 100644 index 000000000..7fcfb738b --- /dev/null +++ b/lib/store/memory/file.go @@ -0,0 +1,181 @@ +package memory + +import ( + "errors" + "io" + "sync" + "sync/atomic" + + storelib "github.com/uber/kraken/lib/store" +) + +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. +// As soon as the blob is evicted, File's APIs starts returning [ErrEvicted], as File no longer has a reference to its data, ensuring GC can clean it. +type File struct { + data *atomic.Pointer[[]byte] + sliceMu *sync.Mutex // A potential optimization if contention is too high: currently no writes are parallelized. However, writes that only mutate the array but not the slice are parallelizable with each other. Thus, we could transition to a RWMutex to enable that. + off int64 +} + +func newFile(data *atomic.Pointer[[]byte], sliceMu *sync.Mutex) *File { + return &File{ + data: data, + sliceMu: sliceMu, + off: 0, + } +} + +func (f *File) getData() (data []byte, evicted bool) { + buf := f.data.Load() + if buf == nil { + return nil, true + } + return *buf, false +} + +func (f *File) Read(p []byte) (n int, err error) { + if len(p) == 0 { + return 0, nil + } + buf, evicted := f.getData() + if evicted { + return 0, ErrEvicted + } + + if f.off >= int64(len(buf)) { + return 0, io.EOF + } + n = copy(p, buf[f.off:]) + f.off += int64(n) + return n, nil +} + +// ReadAt implements [io.ReaderAt]. Thread-safe. +func (f *File) ReadAt(p []byte, off int64) (n int, err error) { + if len(p) == 0 { + return 0, nil + } + if off < 0 { + return 0, errors.New("negative offset") + } + + buf, evicted := f.getData() + if evicted { + return 0, ErrEvicted + } + if off >= int64(len(buf)) { + return 0, io.EOF + } + + n = copy(p, buf[off:]) + if n < len(p) { + return n, io.EOF + } + return n, nil +} + +// Seek implements [io.Seeker]. Not thread-safe. +func (f *File) Seek(off int64, whence int) (int64, error) { + buf, evicted := f.getData() + if evicted { + return 0, ErrEvicted + } + + var newOff int64 + switch whence { + case io.SeekStart: + newOff = off + case io.SeekCurrent: + newOff = f.off + off + case io.SeekEnd: + newOff = int64(len(buf)) + off + default: + return 0, errors.New("invalid whence") + } + + if newOff < 0 || newOff > int64(len(buf)) { + return 0, errors.New("invalid seek location") + } + f.off = newOff + return newOff, nil +} + +// Stat returns the blob's actual size, even if it differs from the size reported during Create. +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 int64(len(buf)) +} + +// WriteAt implements [io.WriterAt]. It is fully thread-safe. +func (f *File) WriteAt(p []byte, off int64) (n int, err error) { + if off < 0 { + return 0, errors.New("negative offset") + } + + f.sliceMu.Lock() + defer f.sliceMu.Unlock() + + buf, evicted := f.getData() + if evicted { + return 0, ErrEvicted + } + + end := int(off) + len(p) + buf, resized := resizeSliceIfNecessary(buf, end) + n = copy(buf[off:], p) + if resized { + f.data.Store(&buf) + } + return n, nil +} + +func resizeSliceIfNecessary(buf []byte, end int) ([]byte, bool) { + resized := false + if len(buf) < end { + if cap(buf) < end { + newBuf := make([]byte, end) + copy(newBuf, buf) + buf = newBuf + } + buf = buf[:end] + resized = true + } + return buf, resized +} + +// Write implements io.Writer. +func (f *File) Write(p []byte) (n int, err error) { + f.sliceMu.Lock() // We need to lock in case we update the pointer to `data`. + defer f.sliceMu.Unlock() + + buf, evicted := f.getData() + if evicted { + return 0, ErrEvicted + } + + end := int(f.off) + len(p) + buf, resized := resizeSliceIfNecessary(buf, end) + + n = copy(buf[f.off:], p) + if resized { + f.data.Store(&buf) + } + f.off += int64(n) + return n, nil +} + +// Off returns the offset for the next Read or Write operation. +// Thread-safe with other File APIs ONLY after File's blob is evicted. +func (f *File) Off() int64 { + return f.off +} + +func (f *File) Cancel() error { return nil } // no-op +func (f *File) Close() error { return nil } // no-op +func (f *File) Commit() error { return nil } // no-op diff --git a/lib/store/memory/scoped_store.go b/lib/store/memory/scoped_store.go new file mode 100644 index 000000000..68edb11d4 --- /dev/null +++ b/lib/store/memory/scoped_store.go @@ -0,0 +1,105 @@ +package memory + +import ( + "errors" + + "github.com/uber-go/tally" + storelib "github.com/uber/kraken/lib/store" + "github.com/uber/kraken/lib/store/metadata" +) + +// 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") + +// Store is an in-memory, thread-safe, LRU cache for blobs and their [metadata.Metadata]. +// +// - Supports pagination of blobs during reading/writing, such that blobs don't need to be fully loaded into memory. +// +// - New blobs are considered 'incomplete', which unlists them from LRU eviction. The store can be scoped to work on only (in-)complete blobs. +// +// - The store prioritizes writing new blobs over reading existing ones. Therefore, blobs may get evicted while clients hold a [*File] to them. +// In such cases, [ErrEvicted] is returned. +// +// - All APIs are thread-safe. Parallel access to a single blob is allowed but clients must ensure they don't intervene with one another. +// +// - Supports (un-)marking blobs as non-evictable (needed when we want to ensure an entry does not get evicted before the client flushes it to disk). +type Store struct { + impl *store + 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) + if err != nil { + return nil, err + } + return &Store{ + impl: s, + scope: storelib.BlobScopeAny, + }, 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). +func (s *Store) Create(key string, sizeBytes uint64) (*File, error) { + return s.impl.Create(key, sizeBytes) +} + +// Open returns a [*File] pointing to the blob. [*File] APIs returns [ErrEvicted] once the blob gets evicted. +func (s *Store) Open(key string) (*File, error) { return s.impl.Open(key, s.scope) } + +// Has checks if the blob is in the store. +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 during Create. +func (s *Store) Stat(key string) (size int64, err error) { return s.impl.Stat(key, s.scope) } + +// MarkComplete marks the blob as fully written, which enlists it for LRU eviction (unless BanEviction has been called). It is idempotent. +// Additionally, other store APIs may filter blobs based on completeness. +func (s *Store) MarkComplete(key string) error { return s.impl.MarkComplete(key) } + +// Delete removes a blob and its metadata from the store. Returns [os.ErrNotExist] on missing entry. +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) } + +// 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. +func (s *Store) BanEviction(key string) error { return s.impl.BanEviction(key, s.scope) } + +// UnbanEviction removes the effect of BanEviction for a blob. It is idempotent. +func (s *Store) UnbanEviction(key string) error { return s.impl.UnbanEviction(key, 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. Returns [os.ErrNotExist] if key is not in store. +func (s *Store) GetMetadata(key string, md metadata.Metadata) (ok bool, err error) { + return s.impl.GetMetadata(key, md, s.scope) +} + +// ListMetadata returns all [metadata.Metadata] of key. +func (s *Store) ListMetadata(key string) ([]metadata.Metadata, error) { + return s.impl.ListMetadata(key, s.scope) +} + +// DeleteMetadata removes a blob's metadata. No error returned if the metadata 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/memory/store.go b/lib/store/memory/store.go new file mode 100644 index 000000000..cc1935176 --- /dev/null +++ b/lib/store/memory/store.go @@ -0,0 +1,364 @@ +package memory + +import ( + "container/list" + "errors" + "fmt" + "maps" + "os" + "slices" + "sync" + "sync/atomic" + + "github.com/uber-go/tally" + storelib "github.com/uber/kraken/lib/store" + "github.com/uber/kraken/lib/store/metadata" + "github.com/uber/kraken/utils/log" + "go.uber.org/zap" +) + +// 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. +// +// Check [Store]'s comments for details on functionality. +type store struct { + blobs map[string]*blob + evictQueue *list.List // front is next to evict (least recently used entry) + size uint64 + capacity uint64 + mu sync.RWMutex // TODO - benchmark if a [sync.Mutex] has better perf. + log *zap.SugaredLogger + metrics tally.Scope +} + +type blob struct { + data atomic.Pointer[[]byte] // set to nil upon eviction/deletion. + node *list.Element + metadatas map[string]metadata.Metadata + size uint64 + complete bool + evictionBanned bool + 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") + } + + log := log.Default().With("module", "memory_store") + + log.Info("Initialized new, empty *memory.Store") + s := &store{ + blobs: make(map[string]*blob, 0), + evictQueue: list.New(), + capacity: capacityBytes, + size: 0, + log: log, + metrics: metrics, + } + + s.emitUsageMetrics() + return s, nil +} + +func (s *store) Create(key string, sizeBytes uint64) (*File, error) { + s.mu.Lock() + defer s.mu.Unlock() + + _, ok := s.blobs[key] + if ok { + return nil, os.ErrExist + } + if ok := s.reserveSpace(sizeBytes); !ok { + return nil, ErrNoSpace + } + b := &blob{ + size: sizeBytes, + complete: false, + evictionBanned: false, + node: nil, + metadatas: make(map[string]metadata.Metadata, 0), + } + arr := make([]byte, 0, sizeBytes) + b.data.Store(&arr) + s.blobs[key] = b + + s.emitUsageMetrics() + return newFile(&b.data, &b.sliceMu), nil +} + +func (s *store) reserveSpace(space uint64) bool { + // TODO - consider whether it's a worth optimization to check if we can evict enough data BEFORE we start evicting, as to prevent evicting needlessly. + for s.size+space > s.capacity { + if s.evictQueue.Len() == 0 { + return false + } + toEvictNode := s.evictQueue.Front() + toEvictKey := toEvictNode.Value.(string) //nolint:errcheck + b := s.blobs[toEvictKey] + b.sliceMu.Lock() + b.data.Store(nil) // Ensure the byte slice is not referenced by clients outside the store holding [*File], so GC can evict the memory. + b.sliceMu.Unlock() + delete(s.blobs, toEvictKey) + s.evictQueue.Remove(toEvictNode) + s.releaseSpace(b.size) + } + + s.size += space + return true +} + +func (s *store) releaseSpace(space uint64) { + if space > s.size { + s.log.Error("Invariant violation - memory.Store wants to release more space than actually reserved. Failing open by setting store.size = 0") + s.size = 0 + return + } + s.size -= space +} + +func (s *store) Open(key string, scope storelib.BlobScope) (*File, error) { + s.mu.Lock() + defer s.mu.Unlock() + + b, ok := s.blobs[key] + if !ok { + return nil, os.ErrNotExist + } + if err := isOutOfScope(b, scope); err != nil { + return nil, err + } + + if b.node != nil { + s.evictQueue.MoveToBack(b.node) + } + return newFile(&b.data, &b.sliceMu), nil +} + +func (s *store) Has(key string, scope storelib.BlobScope) (inStore bool, inScope bool) { + s.mu.RLock() + defer s.mu.RUnlock() + + b, ok := s.blobs[key] + if !ok { + return false, false + } + if err := isOutOfScope(b, scope); err != nil { + return true, false + } + return true, true +} + +func (s *store) Delete(key string, scope storelib.BlobScope) error { + s.mu.Lock() + defer s.mu.Unlock() + + b, ok := s.blobs[key] + if !ok { + return os.ErrNotExist + } + if err := isOutOfScope(b, scope); err != nil { + return err + } + + b.sliceMu.Lock() + b.data.Store(nil) // Ensure the byte slice is not referenced by clients outside the store holding [*File], so GC can evict the memory. + b.sliceMu.Unlock() + delete(s.blobs, key) + if b.node != nil { + s.evictQueue.Remove(b.node) + } + s.releaseSpace(b.size) + + s.emitUsageMetrics() + return nil +} + +func (s *store) list(scope storelib.BlobScope) []string { + s.mu.RLock() + defer s.mu.RUnlock() + + res := make([]string, 0, len(s.blobs)) + for key, b := range s.blobs { + if err := isOutOfScope(b, scope); err != nil { + continue + } + res = append(res, key) + } + return res +} + +func (s *store) Stat(key string, scope storelib.BlobScope) (size int64, err error) { + s.mu.RLock() + defer s.mu.RUnlock() + + b, ok := s.blobs[key] + if !ok { + return 0, os.ErrNotExist + } + if err := isOutOfScope(b, scope); err != nil { + return 0, err + } + buf := *b.data.Load() + return int64(len(buf)), nil +} + +func (s *store) MarkComplete(key string) error { + s.mu.Lock() + defer s.mu.Unlock() + + b, ok := s.blobs[key] + if !ok { + return os.ErrNotExist + } + if b.complete { + // no-op + return nil + } + + b.complete = true + if !b.evictionBanned { + node := s.evictQueue.PushBack(key) + b.node = node + } + for mdSuffix, md := range b.metadatas { + if !md.Movable() { + delete(b.metadatas, mdSuffix) + } + } + return nil +} + +func (s *store) BanEviction(key string, scope storelib.BlobScope) error { + s.mu.Lock() + defer s.mu.Unlock() + + b, ok := s.blobs[key] + if !ok { + return os.ErrNotExist + } + if err := isOutOfScope(b, scope); err != nil { + return err + } + if b.evictionBanned { + // no-op + return nil + } + + b.evictionBanned = true + if b.complete { + s.evictQueue.Remove(b.node) + b.node = nil + } + return nil +} + +func (s *store) UnbanEviction(key string, scope storelib.BlobScope) error { + s.mu.Lock() + defer s.mu.Unlock() + + b, ok := s.blobs[key] + if !ok { + return os.ErrNotExist + } + if err := isOutOfScope(b, scope); err != nil { + return err + } + if !b.evictionBanned { + // no-op + return nil + } + + b.evictionBanned = false + if b.complete { + node := s.evictQueue.PushBack(key) + b.node = node + } + return nil +} + +func (s *store) SetMetadata(key string, md metadata.Metadata, scope storelib.BlobScope) error { + s.mu.Lock() + defer s.mu.Unlock() + + b, ok := s.blobs[key] + if !ok { + return os.ErrNotExist + } + if err := isOutOfScope(b, scope); err != nil { + return err + } + b.metadatas[md.GetSuffix()] = md + return nil +} + +func (s *store) GetMetadata(key string, md metadata.Metadata, scope storelib.BlobScope) (ok bool, err error) { + s.mu.RLock() + defer s.mu.RUnlock() + + b, ok := s.blobs[key] + if !ok { + return false, os.ErrNotExist + } + if err := isOutOfScope(b, scope); err != nil { + return false, err + } + + res, ok := b.metadatas[md.GetSuffix()] + if !ok { + return false, nil + } + mdData, err := res.Serialize() + if err != nil { + return false, fmt.Errorf("serialize metadata: %w", err) + } + err = md.Deserialize(mdData) + if err != nil { + return false, fmt.Errorf("deserialize metadata: %w", err) + } + return ok, nil +} + +func (s *store) ListMetadata(key string, scope storelib.BlobScope) ([]metadata.Metadata, error) { + s.mu.RLock() + defer s.mu.RUnlock() + + b, ok := s.blobs[key] + if !ok { + return nil, os.ErrNotExist + } + if err := isOutOfScope(b, scope); err != nil { + return nil, err + } + + return slices.Collect(maps.Values(b.metadatas)), nil +} + +func (s *store) DeleteMetadata(key string, mdSuffix string, scope storelib.BlobScope) error { + s.mu.Lock() + defer s.mu.Unlock() + + b, ok := s.blobs[key] + if !ok { + return os.ErrNotExist + } + if err := isOutOfScope(b, scope); err != nil { + return err + } + + delete(b.metadatas, mdSuffix) + return nil +} + +func isOutOfScope(b *blob, scope storelib.BlobScope) error { + if (b.complete && scope == storelib.BlobScopeIncomplete) || (!b.complete && scope == storelib.BlobScopeComplete) { + return storelib.ErrOutOfScope + } + return nil +} + +func (s *store) emitUsageMetrics() { + s.metrics.Gauge("num_entries").Update(float64(len(s.blobs))) + s.metrics.Gauge("size_bytes").Update(float64(s.size)) +} diff --git a/lib/store/memory/store_test.go b/lib/store/memory/store_test.go new file mode 100644 index 000000000..976f2cd64 --- /dev/null +++ b/lib/store/memory/store_test.go @@ -0,0 +1,1034 @@ +package memory + +import ( + "bytes" + "crypto/rand" + "io" + "os" + "regexp" + "sync" + "testing" + + "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/metadata" + "github.com/uber/kraken/utils/memsize" +) + +func newTestFile(t *testing.T, store *Store, size uint64) (f storelib.FileReadWriter, key string) { + require := require.New(t) + key = core.DigestFixture().Hex() + f, err := store.Create(key, size) + require.NoError(err) + return f, key +} + +func TestEviction(t *testing.T) { + require := require.New(t) + store, err := NewStore(25*memsize.KB, tally.NoopScope) + require.NoError(err) + // 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) + c, cKey := newTestFile(t, store, 5*memsize.KB) + _, dKey := newTestFile(t, store, 3*memsize.KB) + e, eKey := newTestFile(t, store, 1*memsize.KB) + + 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) + 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)) + require.NoError(store.MarkComplete(aKey)) + // d is complete but its eviction is banned. + require.NoError(store.MarkComplete(dKey)) + require.NoError(store.BanEviction(dKey)) + // e is banned from eviction before it even becomes complete. + require.NoError(store.BanEviction(eKey)) + require.Equal(24*memsize.KB, store.impl.size) + // Add f (4KB) which should evict c to make space, as d and e are unevictable and c was accessed last (the MarkComplete call). + f, fKey := newTestFile(t, store, 4*memsize.KB) + require.NoError(f.Close()) + require.NoError(store.MarkComplete(fKey)) + _, ok := store.Has(cKey) + require.False(ok) + _, err = c.Read(make([]byte, 5)) + require.ErrorIs(err, ErrEvicted) + // new size == 23KB == 24KB - 5KB (c) + 4KB (f) + require.Equal(23*memsize.KB, store.impl.size) + + // Add g (1KB), which will not evict anything + g, gKey := newTestFile(t, store, 1*memsize.KB) + require.NoError(g.Close()) + require.NoError(store.MarkComplete(gKey)) + require.Equal(24*memsize.KB, store.impl.size) + + // Add h (15KB), which evicts b and a: + h, hKey := newTestFile(t, store, 15*memsize.KB) + require.NoError(h.Close()) + require.NoError(store.MarkComplete(hKey)) + // size == 24KB == 24KB + 15KB (h) - 5KB (b) - 10KB (a) + require.Equal(24*memsize.KB, store.impl.size) + _, ok = store.Has(bKey) + require.False(ok) + _, err = b.Read(make([]byte, 5)) + require.ErrorIs(err, ErrEvicted) + _, ok = store.Has(aKey) + require.False(ok) + _, err = a.Read(make([]byte, 5)) + require.ErrorIs(err, ErrEvicted) + + // allow e to be evicted. + require.NoError(store.MarkComplete(eKey)) + require.NoError(store.UnbanEviction(eKey)) + // eviction order (left-most is next to evict): f(4KB), g(1KB), h(15KB), e(1KB); d(3KB) is unevictable + // we open g to change the order to f, h, e, g + g, err = store.Open(gKey) + require.NoError(err) + require.NoError(g.Close()) + + i, iKey := newTestFile(t, store, 5*memsize.KB) + require.NoError(store.MarkComplete(iKey)) + require.NoError(i.Close()) + _, ok = store.Has(fKey) + require.False(ok) + _, err = f.Read(make([]byte, 5)) + require.ErrorIs(err, ErrEvicted) + require.Equal(25*memsize.KB, store.impl.size) + // eviction order: h(15KB), e(1KB), g(1KB), i(5KB); d(3KB) is unevictable + + j, jKey := newTestFile(t, store, 14*memsize.KB) + require.NoError(j.Close()) + require.NoError(store.MarkComplete(jKey)) + _, ok = store.Has(hKey) + require.False(ok) + _, err = h.Read(make([]byte, 5)) + require.ErrorIs(err, ErrEvicted) + require.Equal(24*memsize.KB, store.impl.size) + // eviction order: e(1KB), g(1KB), i(5KB), j(14KB); d(3KB) is unevictable + + k, kKey := newTestFile(t, store, 2*memsize.KB) + require.NoError(k.Close()) + require.NoError(store.MarkComplete(kKey)) + _, ok = store.Has(eKey) + require.False(ok) + _, err = e.Read(make([]byte, 5)) + require.ErrorIs(err, ErrEvicted) + require.Equal(25*memsize.KB, store.impl.size) + // eviction order: g(1KB), i(5KB), j(14KB), k(2KB); d(3KB) is unevictable + + l, lKey := newTestFile(t, store, 1*memsize.KB) + require.NoError(store.MarkComplete(lKey)) + require.NoError(l.Close()) + _, ok = store.Has(gKey) + require.False(ok) + require.Equal(25*memsize.KB, store.impl.size) + // eviction order: i(5KB), j(14KB), k(2KB), l(1KB); d(3KB) is unevictable + + require.NoError(store.Delete(iKey)) + require.NoError(store.Delete(jKey)) + require.NoError(store.Delete(lKey)) + // evictionOrder: k(2KB); d(3KB) + require.Equal(5*memsize.KB, store.impl.size) +} + +func TestEvictedBlobImmediatelyNotAvailable(t *testing.T) { + require := require.New(t) + store, err := NewStore(1*memsize.KB, tally.NoopScope) + require.NoError(err) + b, key := newTestFile(t, store, 1*memsize.KB) + require.NoError(store.MarkComplete(key)) + + _, newBlobKey := newTestFile(t, store, 1*memsize.KB) + + require.Equal([]string{newBlobKey}, store.List()) + buf := make([]byte, 10) + + n, err := b.Read(buf) + require.Zero(n) + require.Equal(ErrEvicted, err) + + n, err = b.ReadAt(buf, 3) + require.Zero(n) + require.Equal(ErrEvicted, err) + + seekN, err := b.Seek(10, io.SeekStart) + require.True(seekN == 0) + require.Equal(ErrEvicted, err) + + n, err = b.Write(buf) + require.Zero(n) + require.Equal(ErrEvicted, err) + + n, err = b.WriteAt(buf, 3) + require.Zero(n) + require.Equal(ErrEvicted, err) + + require.NoError(b.Cancel()) + require.NoError(b.Close()) + require.NoError(b.Commit()) + + require.Equal(int64(0), b.Size()) +} + +func TestParallelAccessToSingleFile(t *testing.T) { + require := require.New(t) + store, err := NewStore(10*memsize.KB, tally.NoopScope) + require.NoError(err) + + 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. + f, err := store.Create(key, 45*memsize.B) + require.NoError(err) + require.NoError(f.Close()) + // Spawn 5 routines in parallel that write and read to different parts of the file. + var wg sync.WaitGroup + wg.Add(5) + + type res struct { + written, read []byte + err error + } + + results := make([]res, 5) + + for idx := range 5 { + go func(idx int64) { + defer wg.Done() + + f, err := store.Open(key) + if err != nil { + results[idx].err = err + return + } + pos := idx * 10 + writtenData := make([]byte, 10) + for k := range writtenData { + writtenData[k] = byte(idx) + } + _, err = f.WriteAt(writtenData, pos) + if err != nil { + results[idx].err = err + return + } + + readData := make([]byte, 10) + _, err = f.ReadAt(readData, pos) + if err != nil { + results[idx].err = err + return + } + err = f.Close() + if err != nil { + results[idx].err = err + return + } + + results[idx].read = readData + results[idx].written = writtenData + }(int64(idx)) + } + + wg.Wait() + for idx := range 5 { + require.NoError(results[idx].err) + require.Equal(results[idx].written, results[idx].read) + } + + require.NoError(store.MarkComplete(key)) + + f, err = store.ScopeComplete().Open(key) + require.NoError(err) + defer func() { require.NoError(f.Close()) }() + + wantFileData := make([]byte, 50) + for i := range 50 { + wantFileData[i] = byte(i / 10) + } + fData, err := io.ReadAll(f) + require.NoError(err) + require.Equal(wantFileData, fData) +} + +func TestOpenedFileAccessibleAfterMarkedComplete(t *testing.T) { + require := require.New(t) + store, err := NewStore(10*memsize.KB, tally.NoopScope) + require.NoError(err) + + key := core.DigestFixture().Hex() + data := []byte("Hello World") + f, err := store.Create(key, uint64(len(data))) + require.NoError(err) + _, err = io.Copy(f, bytes.NewReader(data)) + require.NoError(err) + require.NoError(f.Close()) + + incompleteFile, err := store.Open(key) + require.NoError(err) + defer func() { require.NoError(incompleteFile.Close()) }() + + require.NoError(store.MarkComplete(key)) + + completeFile, err := store.ScopeComplete().Open(key) + require.NoError(err) + defer func() { require.NoError(completeFile.Close()) }() + + incompleteFileData, err := io.ReadAll(incompleteFile) + require.NoError(err) + completeFileData, err := io.ReadAll(completeFile) + require.NoError(err) + + require.Equal(data, incompleteFileData) + require.Equal(data, completeFileData) +} + +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) + + // 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) + require.NoError(err) + _, err = io.Copy(f, bytes.NewReader(data)) + require.NoError(err) + require.NoError(f.Close()) + + f, err = store.Open(key) + require.NoError(err) + + // io.ReadAll only returns the 5KB of actual data, despite us reporting the blob as 10KB. There is no 5KB of trailing space. + readData, err := io.ReadAll(f) + require.NoError(err) + require.Equal(data, readData) + + // Stat returns the actual size of the blob. + statSize, err := store.Stat(key) + require.NoError(err) + require.Equal(uint64(statSize), 5*memsize.KB) + + // The store always calculates its size based off the client-reported size, not the actual size. + require.Equal(11*memsize.KB, store.impl.size) + + // Both memory reservation and memory release use the user-reported value, not the actual size of the blob. + // Due to this consistency, the store's size is the same before adding and after deleting/evicting the blob. + require.NoError(store.Delete(key)) + require.Equal(1*memsize.KB, store.impl.size) + }) + + t.Run("user underreports size", func(t *testing.T) { + require := require.New(t) + store, err := NewStore(25*memsize.KB, tally.NoopScope) + require.NoError(err) + + // 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) + + // 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) + require.NoError(err) + _, err = io.Copy(f, bytes.NewReader(data)) + require.NoError(err) + require.NoError(f.Close()) + + f, err = store.Open(key) + require.NoError(err) + + // io.ReadAll only returns the 30KB of actual data, despite us reporting the blob as 24KB. + // Despite the store initializing only 24KB of memory in Create, the slice is safely resized when writes exceed its length. + readData, err := io.ReadAll(f) + require.NoError(err) + require.Equal(data, readData) + + // Stat returns the actual size of the blob. + statSize, err := store.Stat(key) + require.NoError(err) + require.Equal(uint64(statSize), 30*memsize.KB) + + // The store always calculates its size based off the client-reported size, not the actual size. + require.Equal(25*memsize.KB, store.impl.size) + + // Both memory reservation and memory release use the user-reported value, not the actual size of the blob. + // Due to this consistency, the store's size is the same before adding and after deleting/evicting the blob. + require.NoError(store.Delete(key)) + require.Equal(1*memsize.KB, store.impl.size) + }) +} + +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) + key := core.DigestFixture().Hex() + f, err := store.Create(key, 100*memsize.B) + require.NoError(err) + _, err = io.Copy(f, bytes.NewReader(make([]byte, 100))) + require.NoError(err) + require.NoError(f.Close()) + require.NoError(store.Delete(key)) + + require.Empty(store.List()) + require.Equal(uint64(0), store.impl.size) + require.Equal(0, len(store.impl.blobs)) + }) + t.Run("incomplete, unevictable blob", func(t *testing.T) { + require := require.New(t) + store, err := NewStore(10*memsize.KB, tally.NoopScope) + require.NoError(err) + key := core.DigestFixture().Hex() + f, err := store.Create(key, 100*memsize.B) + require.NoError(err) + _, err = io.Copy(f, bytes.NewReader(make([]byte, 100))) + require.NoError(err) + require.NoError(f.Close()) + require.NoError(store.BanEviction(key)) + + require.NoError(store.Delete(key)) + + require.Empty(store.List()) + require.Equal(uint64(0), store.impl.size) + require.Equal(0, len(store.impl.blobs)) + }) + t.Run("complete blob", func(t *testing.T) { + require := require.New(t) + store, err := NewStore(10*memsize.KB, tally.NoopScope) + require.NoError(err) + key := core.DigestFixture().Hex() + f, err := store.Create(key, 100*memsize.B) + require.NoError(err) + _, err = io.Copy(f, bytes.NewReader(make([]byte, 100))) + require.NoError(err) + require.NoError(f.Close()) + require.NoError(store.MarkComplete(key)) + + require.NoError(store.Delete(key)) + + require.Empty(store.List()) + require.Equal(uint64(0), store.impl.size) + require.Equal(0, len(store.impl.blobs)) + }) + t.Run("complete, unevictable blob", func(t *testing.T) { + require := require.New(t) + store, err := NewStore(10*memsize.KB, tally.NoopScope) + require.NoError(err) + key := core.DigestFixture().Hex() + f, err := store.Create(key, 100*memsize.B) + require.NoError(err) + _, err = io.Copy(f, bytes.NewReader(make([]byte, 100))) + require.NoError(err) + require.NoError(f.Close()) + require.NoError(store.MarkComplete(key)) + require.NoError(store.BanEviction(key)) + + require.NoError(store.Delete(key)) + + require.Empty(store.List()) + require.Equal(uint64(0), store.impl.size) + require.Equal(0, len(store.impl.blobs)) + }) + t.Run("not found", func(t *testing.T) { + require := require.New(t) + store, err := NewStore(10*memsize.KB, tally.NoopScope) + require.NoError(err) + key := core.DigestFixture().Hex() + + err = store.Delete(key) + require.ErrorIs(err, os.ErrNotExist) + }) +} + +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) + key := core.DigestFixture().Hex() + f, err := store.Create(key, 100*memsize.B) + require.NoError(err) + defer func() { require.NoError(f.Close()) }() + _, err = io.Copy(f, bytes.NewReader(make([]byte, 100))) + require.NoError(err) + + require.NoError(store.MarkComplete(key)) + + require.Equal([]string{key}, store.ScopeComplete().List()) + require.Equal(uint64(100), store.impl.size) + _, err = store.ScopeComplete().Open(key) + require.NoError(err) + }) + 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) + key := core.DigestFixture().Hex() + f, err := store.Create(key, 100*memsize.B) + require.NoError(err) + defer func() { require.NoError(f.Close()) }() + _, err = io.Copy(f, bytes.NewReader(make([]byte, 100))) + require.NoError(err) + require.NoError(store.BanEviction(key)) + + require.NoError(store.MarkComplete(key)) + + require.Equal([]string{key}, store.ScopeComplete().List()) + require.Equal(uint64(100), store.impl.size) + _, err = store.ScopeComplete().Open(key) + require.NoError(err) + }) + t.Run("already complete blob", func(t *testing.T) { + require := require.New(t) + store, err := NewStore(10*memsize.KB, tally.NoopScope) + require.NoError(err) + key := core.DigestFixture().Hex() + f, err := store.Create(key, 100*memsize.B) + require.NoError(err) + defer func() { require.NoError(f.Close()) }() + _, err = io.Copy(f, bytes.NewReader(make([]byte, 100))) + require.NoError(err) + require.NoError(store.MarkComplete(key)) + + require.NoError(store.MarkComplete(key)) + }) + 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) + key := core.DigestFixture().Hex() + f, err := store.Create(key, 100*memsize.B) + require.NoError(err) + defer func() { require.NoError(f.Close()) }() + _, err = io.Copy(f, bytes.NewReader(make([]byte, 100))) + require.NoError(err) + require.NoError(store.MarkComplete(key)) + require.NoError(store.BanEviction(key)) + + require.NoError(store.MarkComplete(key)) + }) + t.Run("not found", func(t *testing.T) { + require := require.New(t) + store, err := NewStore(10*memsize.KB, tally.NoopScope) + require.NoError(err) + key := core.DigestFixture().Hex() + + err = store.MarkComplete(key) + require.ErrorIs(err, os.ErrNotExist) + }) +} + +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) + key := core.DigestFixture().Hex() + f, err := store.Create(key, 10*memsize.B) + require.NoError(err) + defer func() { require.NoError(f.Close()) }() + _, err = io.Copy(f, bytes.NewReader(make([]byte, 10))) + require.NoError(err) + require.NoError(store.MarkComplete(key)) + + size, err := store.Stat(key) + require.NoError(err) + require.Equal(int64(10), size) + }) + t.Run("complete, unevictable blob", func(t *testing.T) { + require := require.New(t) + store, err := NewStore(10*memsize.KB, tally.NoopScope) + require.NoError(err) + key := core.DigestFixture().Hex() + f, err := store.Create(key, 10*memsize.B) + require.NoError(err) + defer func() { require.NoError(f.Close()) }() + _, err = io.Copy(f, bytes.NewReader(make([]byte, 10))) + require.NoError(err) + require.NoError(store.MarkComplete(key)) + require.NoError(store.BanEviction(key)) + + size, err := store.Stat(key) + require.NoError(err) + require.Equal(int64(10), size) + }) + t.Run("incomplete blob", func(t *testing.T) { + require := require.New(t) + store, err := NewStore(10*memsize.KB, tally.NoopScope) + require.NoError(err) + key := core.DigestFixture().Hex() + f, err := store.Create(key, 10*memsize.B) + require.NoError(err) + defer func() { require.NoError(f.Close()) }() + _, err = io.Copy(f, bytes.NewReader(make([]byte, 10))) + require.NoError(err) + + _, err = store.ScopeComplete().Stat(key) + require.Equal(storelib.ErrOutOfScope, err) + + size, err := store.Stat(key) + require.NoError(err) + require.Equal(int64(10), size) + }) + + t.Run("incomplete, unevictable blob", func(t *testing.T) { + require := require.New(t) + store, err := NewStore(10*memsize.KB, tally.NoopScope) + require.NoError(err) + key := core.DigestFixture().Hex() + f, err := store.Create(key, 10*memsize.B) + require.NoError(err) + defer func() { require.NoError(f.Close()) }() + _, err = io.Copy(f, bytes.NewReader(make([]byte, 10))) + require.NoError(err) + require.NoError(store.BanEviction(key)) + + size, err := store.Stat(key) + require.NoError(err) + require.Equal(int64(10), size) + }) + t.Run("non-existent blob", func(t *testing.T) { + require := require.New(t) + store, err := NewStore(10*memsize.KB, tally.NoopScope) + require.NoError(err) + key := core.DigestFixture().Hex() + + size, err := store.Stat(key) + require.ErrorIs(err, os.ErrNotExist) + require.Equal(int64(0), size) + }) +} + +func TestList(t *testing.T) { + require := require.New(t) + store, err := NewStore(10*memsize.KB, tally.NoopScope) + require.NoError(err) + + require.Empty(store.ScopeComplete().List()) + require.Empty(store.ScopeComplete().List()) + + incompleteBlobKey := core.DigestFixture().Hex() + f, err := store.Create(incompleteBlobKey, 10*memsize.B) + require.NoError(err) + defer func(f io.Closer) { require.NoError(f.Close()) }(f) + _, err = io.Copy(f, bytes.NewReader(make([]byte, 10))) + require.NoError(err) + completeBlobKey := core.DigestFixture().Hex() + f, err = store.Create(completeBlobKey, 10*memsize.B) + require.NoError(err) + _, err = io.Copy(f, bytes.NewReader(make([]byte, 10))) + require.NoError(err) + defer func(f io.Closer) { require.NoError(f.Close()) }(f) + require.NoError(store.MarkComplete(completeBlobKey)) + unevictableIncompleteBlobKey := core.DigestFixture().Hex() + f, err = store.Create(unevictableIncompleteBlobKey, 10*memsize.B) + require.NoError(err) + _, err = io.Copy(f, bytes.NewReader(make([]byte, 10))) + require.NoError(err) + defer func(f io.Closer) { require.NoError(f.Close()) }(f) + unevictableCompleteBlobKey := core.DigestFixture().Hex() + f, err = store.Create(unevictableCompleteBlobKey, 10*memsize.B) + require.NoError(err) + _, err = io.Copy(f, bytes.NewReader(make([]byte, 10))) + require.NoError(err) + defer func() { require.NoError(f.Close()) }() + require.NoError(store.MarkComplete(unevictableCompleteBlobKey)) + + wantRes := []string{completeBlobKey, unevictableCompleteBlobKey} + res := store.ScopeComplete().List() + require.ElementsMatch(wantRes, res) + + wantRes = []string{incompleteBlobKey, completeBlobKey, unevictableCompleteBlobKey, unevictableIncompleteBlobKey} + res = store.List() + require.ElementsMatch(wantRes, res) +} + +func init() { + metadata.Register(regexp.MustCompile("immovableMd"), &immovableMdFactory{}) +} + +// used for testing +type immovableMd struct{} +type immovableMdFactory struct{} + +func (f *immovableMdFactory) Create(suffix string) metadata.Metadata { return &immovableMd{} } +func (mda *immovableMd) GetSuffix() string { return "immovableMd" } +func (mda *immovableMd) Movable() bool { return false } +func (mda *immovableMd) Serialize() ([]byte, error) { return nil, nil } +func (mda *immovableMd) Deserialize(b []byte) error { return nil } +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) + key := core.DigestFixture().Hex() + f, err := store.Create(key, 10*memsize.KB) + require.NoError(err) + require.NoError(f.Close()) + + mdStruct := core.MetaInfoFixture() + writtenMd := metadata.NewTorrentMeta(mdStruct) + err = store.SetMetadata(key, writtenMd) + // asserts metadata is not included in LRU eviction calculation. + require.NoError(err) + + var readMd metadata.TorrentMeta + ok, err := store.GetMetadata(key, &readMd) + require.NoError(err) + require.True(ok) + require.Equal(writtenMd.MetaInfo, readMd.MetaInfo) + + mdList, err := store.ListMetadata(key) + require.NoError(err) + require.Len(mdList, 1) + require.Equal(writtenMd.GetSuffix(), mdList[0].GetSuffix()) + + persistMd := metadata.NewPersist(false) + require.NoError(store.SetMetadata(key, persistMd)) + var readPersistMd metadata.Persist + ok, err = store.GetMetadata(key, &readPersistMd) + require.NoError(err) + require.True(ok) + require.False(readPersistMd.Value) + + mdList, err = store.ListMetadata(key) + require.NoError(err) + require.ElementsMatch( + []string{writtenMd.GetSuffix(), persistMd.GetSuffix()}, + []string{mdList[0].GetSuffix(), mdList[1].GetSuffix()}, + ) + + require.NoError(store.DeleteMetadata(key, readMd.GetSuffix())) + ok, err = store.GetMetadata(key, &readMd) + require.NoError(err) + require.False(ok) + // deleting a second time should be a no-op. + 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.GetSuffix())) + mdList, err = store.ListMetadata(key) + require.NoError(err) + require.Empty(mdList) + }) + + t.Run("non-existent blob", func(t *testing.T) { + require := require.New(t) + store, err := NewStore(10*memsize.KB, tally.NoopScope) + require.NoError(err) + nonExistentKey := core.DigestFixture().Hex() + mdStruct := core.MetaInfoFixture() + md := metadata.NewTorrentMeta(mdStruct) + + err = store.SetMetadata(nonExistentKey, md) + require.ErrorIs(err, os.ErrNotExist) + + ok, err := store.GetMetadata(nonExistentKey, md) + require.ErrorIs(err, os.ErrNotExist) + require.False(ok) + + _, err = store.ListMetadata(nonExistentKey) + require.ErrorIs(err, os.ErrNotExist) + + err = store.DeleteMetadata(nonExistentKey, md.GetSuffix()) + require.ErrorIs(err, os.ErrNotExist) + }) + + 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) + key := core.DigestFixture().Hex() + f, err := store.Create(key, 1*memsize.KB) + require.NoError(err) + require.NoError(f.Close()) + + mdStruct := core.MetaInfoFixture() + writtenMd := metadata.NewTorrentMeta(mdStruct) + require.NoError(store.SetMetadata(key, writtenMd)) + + var readMd metadata.TorrentMeta + ok, err := store.ScopeComplete().GetMetadata(key, &readMd) + require.Equal(storelib.ErrOutOfScope, err) + // incomplete files are ignored + require.False(ok) + + ok, err = store.GetMetadata(key, &readMd) + require.NoError(err) + require.True(ok) + require.Equal(writtenMd.MetaInfo, readMd.MetaInfo) + + // Repeat the tests above for an unevictable file + require.NoError(store.BanEviction(key)) + ok, err = store.ScopeComplete().GetMetadata(key, &readMd) + require.Equal(storelib.ErrOutOfScope, err) + // incomplete files are ignored + require.False(ok) + + ok, err = store.GetMetadata(key, &readMd) + require.NoError(err) + require.True(ok) + require.Equal(writtenMd.MetaInfo, readMd.MetaInfo) + + require.NoError(store.MarkComplete(key)) + ok, err = store.ScopeComplete().GetMetadata(key, &readMd) + require.NoError(err) + require.True(ok) + require.Equal(writtenMd.MetaInfo, readMd.MetaInfo) + }) + 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) + keyA := core.DigestFixture().Hex() + fA, err := store.Create(keyA, 10*memsize.KB) + require.NoError(err) + require.NoError(fA.Close()) + require.NoError(store.MarkComplete(keyA)) + + md := metadata.NewTorrentMeta(core.MetaInfoFixture()) + err = store.SetMetadata(keyA, md) + require.NoError(err) + + readMd := metadata.TorrentMeta{} + ok, err := store.GetMetadata(keyA, &readMd) + require.NoError(err) + require.True(ok) + + keyB := core.DigestFixture().Hex() + fB, err := store.Create(keyB, 10*memsize.KB) + require.NoError(err) + defer func() { require.NoError(fB.Close()) }() + + ok, err = store.GetMetadata(keyA, md) + require.ErrorIs(err, os.ErrNotExist) + require.False(ok) + }) + + 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) + key := core.DigestFixture().Hex() + f, err := store.Create(key, 10*memsize.KB) + require.NoError(err) + require.NoError(f.Close()) + + md := &immovableMd{} + require.NoError(store.SetMetadata(key, md)) + movableMd := metadata.NewTorrentMeta(core.MetaInfoFixture()) + require.NoError(store.SetMetadata(key, movableMd)) + + require.NoError(store.MarkComplete(key)) + readMd := immovableMd{} + ok, err := store.GetMetadata(key, &readMd) + require.NoError(err) + require.False(ok) + + readMovableMd := metadata.TorrentMeta{} + ok, err = store.GetMetadata(key, &readMovableMd) + require.NoError(err) + require.True(ok) + }) +} + +func TestScopes(t *testing.T) { + require := require.New(t) + store, err := NewStore(10*memsize.KB, tally.NoopScope) + require.NoError(err) + + // 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) + require.NoError(err) + _, err = io.Copy(f, bytes.NewReader(data)) + require.NoError(err) + require.NoError(f.Close()) + + // While incomplete, ScopeComplete's APIs reject the blob. + _, err = store.ScopeComplete().Open(key) + require.ErrorIs(err, storelib.ErrOutOfScope) + _, err = store.ScopeComplete().Stat(key) + require.ErrorIs(err, storelib.ErrOutOfScope) + require.ErrorIs(store.ScopeComplete().BanEviction(key), storelib.ErrOutOfScope) + require.ErrorIs(store.ScopeComplete().UnbanEviction(key), storelib.ErrOutOfScope) + md := metadata.NewTorrentMeta(core.MetaInfoFixture()) + require.ErrorIs(store.ScopeComplete().SetMetadata(key, md), storelib.ErrOutOfScope) + var readMd metadata.TorrentMeta + ok, err := store.ScopeComplete().GetMetadata(key, &readMd) + require.ErrorIs(err, storelib.ErrOutOfScope) + require.False(ok) + _, err = store.ScopeComplete().ListMetadata(key) + require.ErrorIs(err, 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) + + // The unscoped store's APIs work regardless of completeness. + f, err = store.Open(key) + require.NoError(err) + readData, err := io.ReadAll(f) + require.NoError(err) + require.Equal(data, readData) + require.NoError(f.Close()) + _, err = store.Stat(key) + require.NoError(err) + require.NoError(store.BanEviction(key)) + require.NoError(store.UnbanEviction(key)) + require.NoError(store.SetMetadata(key, md)) + ok, err = store.GetMetadata(key, &readMd) + require.NoError(err) + require.True(ok) + require.Equal(md.MetaInfo, readMd.MetaInfo) + mdList, err := store.ListMetadata(key) + require.NoError(err) + require.Len(mdList, 1) + require.Equal(md.GetSuffix(), mdList[0].GetSuffix()) + require.NoError(store.DeleteMetadata(key, readMd.GetSuffix())) + require.Contains(store.List(), key) + _, ok = store.Has(key) + require.True(ok) + + // ScopeIncomplete's APIs also work, since the blob is (still) incomplete. + f, err = store.ScopeIncomplete().Open(key) + require.NoError(err) + readData, err = io.ReadAll(f) + require.NoError(err) + require.Equal(data, readData) + require.NoError(f.Close()) + _, err = store.ScopeIncomplete().Stat(key) + require.NoError(err) + require.NoError(store.ScopeIncomplete().BanEviction(key)) + require.NoError(store.ScopeIncomplete().UnbanEviction(key)) + require.NoError(store.ScopeIncomplete().SetMetadata(key, md)) + ok, err = store.ScopeIncomplete().GetMetadata(key, &readMd) + require.NoError(err) + require.True(ok) + require.Equal(md.MetaInfo, readMd.MetaInfo) + mdList, err = store.ScopeIncomplete().ListMetadata(key) + require.NoError(err) + require.Len(mdList, 1) + require.Equal(md.GetSuffix(), mdList[0].GetSuffix()) + require.NoError(store.ScopeIncomplete().DeleteMetadata(key, readMd.GetSuffix())) + require.Contains(store.ScopeIncomplete().List(), key) + _, ok = store.ScopeIncomplete().Has(key) + require.True(ok) + require.Contains(store.List(), key) + _, ok = store.Has(key) + require.True(ok) + + require.NoError(store.MarkComplete(key)) + + // Now that the blob is complete, the roles reverse: ScopeIncomplete rejects it. + _, err = store.ScopeIncomplete().Open(key) + require.ErrorIs(err, storelib.ErrOutOfScope) + _, err = store.ScopeIncomplete().Stat(key) + require.ErrorIs(err, storelib.ErrOutOfScope) + require.ErrorIs(store.ScopeIncomplete().BanEviction(key), storelib.ErrOutOfScope) + require.ErrorIs(store.ScopeIncomplete().UnbanEviction(key), storelib.ErrOutOfScope) + require.ErrorIs(store.ScopeIncomplete().SetMetadata(key, md), storelib.ErrOutOfScope) + ok, err = store.ScopeIncomplete().GetMetadata(key, &readMd) + require.ErrorIs(err, storelib.ErrOutOfScope) + require.False(ok) + _, err = store.ScopeIncomplete().ListMetadata(key) + require.ErrorIs(err, 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) + + // The unscoped store's APIs still work regardless of completeness. + f, err = store.Open(key) + require.NoError(err) + readData, err = io.ReadAll(f) + require.NoError(err) + require.Equal(data, readData) + require.NoError(f.Close()) + _, err = store.Stat(key) + require.NoError(err) + require.NoError(store.BanEviction(key)) + require.NoError(store.UnbanEviction(key)) + require.NoError(store.SetMetadata(key, md)) + ok, err = store.GetMetadata(key, &readMd) + require.NoError(err) + require.True(ok) + require.Equal(md.MetaInfo, readMd.MetaInfo) + mdList, err = store.ListMetadata(key) + require.NoError(err) + require.Len(mdList, 1) + require.Equal(md.GetSuffix(), mdList[0].GetSuffix()) + require.NoError(store.DeleteMetadata(key, readMd.GetSuffix())) + require.Contains(store.List(), key) + _, ok = store.Has(key) + require.True(ok) + + // ScopeComplete's APIs now work, since the blob is complete. + f, err = store.ScopeComplete().Open(key) + require.NoError(err) + readData, err = io.ReadAll(f) + require.NoError(err) + require.Equal(data, readData) + require.NoError(f.Close()) + _, err = store.ScopeComplete().Stat(key) + require.NoError(err) + require.NoError(store.ScopeComplete().BanEviction(key)) + require.NoError(store.ScopeComplete().UnbanEviction(key)) + require.NoError(store.ScopeComplete().SetMetadata(key, md)) + ok, err = store.ScopeComplete().GetMetadata(key, &readMd) + require.NoError(err) + require.True(ok) + require.Equal(md.MetaInfo, readMd.MetaInfo) + mdList, err = store.ScopeComplete().ListMetadata(key) + require.NoError(err) + require.Len(mdList, 1) + require.Equal(md.GetSuffix(), mdList[0].GetSuffix()) + require.NoError(store.ScopeComplete().DeleteMetadata(key, readMd.GetSuffix())) + require.Contains(store.ScopeComplete().List(), key) + _, ok = store.ScopeComplete().Has(key) + require.True(ok) +} + +func TestScopesDelete(t *testing.T) { + require := require.New(t) + store, err := NewStore(10*memsize.KB, tally.NoopScope) + require.NoError(err) + + 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) + require.ErrorIs(err, os.ErrNotExist) + + f, key = newTestFile(t, store, 1*memsize.KB) + require.NoError(f.Close()) + require.NoError(store.MarkComplete(key)) + require.ErrorIs(store.ScopeIncomplete().Delete(key), storelib.ErrOutOfScope) + require.NoError(store.ScopeComplete().Delete(key)) + _, err = store.Stat(key) + require.ErrorIs(err, os.ErrNotExist) + + f, key = newTestFile(t, store, 1*memsize.KB) + require.NoError(f.Close()) + require.NoError(store.Delete(key)) + + f, key = newTestFile(t, store, 1*memsize.KB) + require.NoError(f.Close()) + require.NoError(store.MarkComplete(key)) + require.NoError(store.Delete(key)) +} diff --git a/lib/store/scope.go b/lib/store/scope.go new file mode 100644 index 000000000..4622932c2 --- /dev/null +++ b/lib/store/scope.go @@ -0,0 +1,17 @@ +package store + +import "errors" + +// BlobScope is the set of blobs that a store's APIs can operate on. +type BlobScope int + +// Flags to scope a store's APIs to a subset of blobs. +const ( + BlobScopeAny BlobScope = iota + BlobScopeComplete + BlobScopeIncomplete +) + +// ErrOutOfScope is returned when the provided key is in the store, but not in the store's [BlobScope], +// e.g. if the blob is incomplete, but ScopeComplete was called. +var ErrOutOfScope = errors.New("the blob is in store but filtered by the selected scope")