diff --git a/lib/store/disk/config.go b/lib/store/disk/config.go new file mode 100644 index 000000000..3f3a0e75a --- /dev/null +++ b/lib/store/disk/config.go @@ -0,0 +1,16 @@ +package disk + +// Config configures [Store]. +type Config struct { + // The capacity of the store in bytes. When breached, LRU eviction is used. + CapacityBytes uint64 + // The root directory under which [Store]'s blobs and state are stored. + RootDir string + // 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 + // 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 +} diff --git a/lib/store/disk/crash_recovery.go b/lib/store/disk/crash_recovery.go new file mode 100644 index 000000000..04e908a62 --- /dev/null +++ b/lib/store/disk/crash_recovery.go @@ -0,0 +1,194 @@ +package disk + +import ( + "container/list" + "errors" + "fmt" + "io" + "os" + "path/filepath" + "slices" + "strconv" + "time" + + "github.com/uber-go/tally" + "github.com/uber/kraken/utils/closers" + "go.uber.org/zap" +) + +type rebootedBlob struct { + key string + size uint64 + mTime time.Time + evictable bool + complete bool +} + +func rebootPersistedStore(config *Config, log *zap.SugaredLogger, metrics tally.Scope) (*store, error) { + incompleteDirPath := filepath.Join(config.RootDir, _incompleteSubDir) + if !config.RebootIncompleteBlobs { + err := os.RemoveAll(incompleteDirPath) + if err != nil { + return nil, fmt.Errorf("remove incomplete blobs left from a previous service run: %w", err) + } + } + + pather := newPather(config.RootDir, config.ShardLength) + keys, err := pather.rebootKeys(_completeBlob) + if err != nil { + return nil, err + } + numCompleteBlobs := len(keys) + if config.RebootIncompleteBlobs { + incompleteKeys, err := pather.rebootKeys(_incompleteBlob) + if err != nil { + return nil, err + } + keys = append(keys, incompleteKeys...) + } + + completeEvictableBlobs := make([]*rebootedBlob, 0) + otherBlobs := make([]*rebootedBlob, 0) + for i, key := range keys { + complete := i < numCompleteBlobs + b, ok, err := rebootBlob(key, complete, pather) + if err != nil { + return nil, err + } + if !ok { + log.With("key", key).Warn("Could not reboot blob from disk - its parent directory is there but the blob is missing") + continue + } + if b.complete && b.evictable { + completeEvictableBlobs = append(completeEvictableBlobs, b) + } else { + otherBlobs = append(otherBlobs, b) + } + } + + storeSize := uint64(0) + blobs := make(map[string]*blob, 0) + for _, b := range otherBlobs { + blobs[b.key] = &blob{ + size: b.size, + complete: b.complete, + evictionBanned: !b.evictable, + node: nil, + } + storeSize += b.size + } + + slices.SortFunc(completeEvictableBlobs, func(left, right *rebootedBlob) int { + // left-most is oldest, i.e. next-to-evict. + return left.mTime.Compare(right.mTime) + }) + evictQueue := list.New() + for _, b := range completeEvictableBlobs { + node := evictQueue.PushBack(b.key) + blobs[b.key] = &blob{ + size: b.size, + complete: true, + node: node, + evictionBanned: false, + } + storeSize += b.size + } + + store := &store{ + blobs: blobs, + evictQueue: evictQueue, + capacity: config.CapacityBytes, + size: storeSize, + pather: pather, + config: config, + log: log, + metrics: metrics, + } + + if store.size > store.capacity { + prevSize := store.size + // evicts blobs until size <= capacity. + err = store.reserveSpace(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) + } + evictedBytes := prevSize - store.size + log.With("evicted_bytes", evictedBytes).Warn("Store size exceeded its capacity after service reboot. Successfully evicted blobs to reduce size within capacity.") + } + return store, nil +} + +func rebootBlob(key string, complete bool, pather *pather) (res *rebootedBlob, ok bool, err error) { + blobPath := pather.blobPath(key, complete) + fInfo, err := os.Stat(blobPath) + if errors.Is(err, os.ErrNotExist) { + // The directory for the blob exists but not the blob itself. + return nil, false, nil + } + if err != nil { + return nil, false, fmt.Errorf("stat blob file: %w", err) + } + + flagBlobPath := pather.sidecarFilePath(key, complete, _evictionBannedFileName) + isUnevictable, err := exists(flagBlobPath) + if err != nil { + return nil, false, err + } + var size uint64 + if complete { + size = uint64(fInfo.Size()) + } else { + size, ok, err = rebootIncompleteBlobSize(key, pather) + if err != nil { + return nil, false, fmt.Errorf("get incomplete blob size from sidecar file: %w", err) + } + if !ok { + return nil, false, nil + } + } + mTime := fInfo.ModTime() + return &rebootedBlob{ + key: key, + size: size, + mTime: mTime, + evictable: !isUnevictable, + complete: complete, + }, true, nil +} + +func rebootIncompleteBlobSize(key string, pather *pather) (size uint64, ok bool, err error) { + blobSizeFilePath := pather.sidecarFilePath(key, _incompleteBlob, _blobSizeFileName) + blobSizeF, err := os.OpenFile(blobSizeFilePath, os.O_RDONLY, _defaultFilePerm) + if errors.Is(err, os.ErrNotExist) { + // The size metadata file is not present, we fail-open by evicting the blob. + return 0, false, nil + } + if err != nil { + return 0, false, fmt.Errorf("open blob size sidecar file: %w", err) + } + defer closers.Close(blobSizeF) + blobSizeData, err := io.ReadAll(blobSizeF) + if err != nil { + return 0, false, fmt.Errorf("read blob size sidecar file: %w", err) + } + blobSize, err := strconv.Atoi(string(blobSizeData)) + if err != nil { + return 0, false, fmt.Errorf("blob size sidecar file is in unexpected format: %w", err) + } + return uint64(blobSize), true, nil +} + +func existsPersistedStore(rootDir string) (ok bool, err error) { + completeDir, incompleteDir := filepath.Join(rootDir, _completeSubDir), filepath.Join(rootDir, _incompleteSubDir) + completeExists, err := exists(completeDir) + if err != nil { + return false, fmt.Errorf("check if store has persisted state left on disk from previous service runs: %w", err) + } + incompleteExists, err := exists(incompleteDir) + if err != nil { + return false, fmt.Errorf("check if store has persisted state left on disk from previous service runs: %w", err) + } + + return completeExists || incompleteExists, nil +} diff --git a/lib/store/disk/crash_recovery_test.go b/lib/store/disk/crash_recovery_test.go new file mode 100644 index 000000000..4f5fd288b --- /dev/null +++ b/lib/store/disk/crash_recovery_test.go @@ -0,0 +1,361 @@ +package disk + +import ( + "bytes" + "crypto/rand" + "io" + "os" + "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 TestCrashRecovery(t *testing.T) { + t.Run("blobs, evictability, completeness, and size are recovered regardless of sharding", func(t *testing.T) { + for _, shardLength := range []int{0, _defaultShardLength, 4} { + require := require.New(t) + + rootDir, err := os.MkdirTemp("/tmp", "kraken-disk-store") + require.NoError(err) + t.Cleanup(func() { require.NoError(os.RemoveAll(rootDir)) }) + config := &Config{ + CapacityBytes: 10 * memsize.KB, + RootDir: rootDir, + RebootIncompleteBlobs: true, + ShardLength: shardLength, + } + + store, err := NewStore(config, tally.NoopScope) + require.NoError(err) + + completeEvictableF, completeEvictableKey := newTestFile(t, store, 2*memsize.KB) + completeEvictableData := fillWithRandomData(t, completeEvictableF, 2*memsize.KB) + // We don't have to test case where file is not closed, as linux closes all FDs owned by a process upon process death. + require.NoError(completeEvictableF.Close()) + require.NoError(store.MarkComplete(completeEvictableKey)) + + completeUnevictableF, completeUnevictableKey := newTestFile(t, store, 2*memsize.KB) + completeUnevictableData := fillWithRandomData(t, completeUnevictableF, 2*memsize.KB) + require.NoError(completeUnevictableF.Close()) + require.NoError(store.MarkComplete(completeUnevictableKey)) + require.NoError(store.BanEviction(completeUnevictableKey)) + + incompleteEvictableF, incompleteEvictableKey := newTestFile(t, store, 2*memsize.KB) + incompleteEvictableData := fillWithRandomData(t, incompleteEvictableF, 2*memsize.KB) + require.NoError(incompleteEvictableF.Close()) + + incompleteUnevictableF, incompleteUnevictableKey := newTestFile(t, store, 2*memsize.KB) + incompleteUnevictableData := fillWithRandomData(t, incompleteUnevictableF, 2*memsize.KB) + require.NoError(incompleteUnevictableF.Close()) + require.NoError(store.BanEviction(incompleteUnevictableKey)) + + // Assume that the application crashes here. The application would restart and call `NewStore`. + store, err = NewStore(&Config{10 * memsize.KB, rootDir, true, shardLength}, tally.NoopScope) + require.NoError(err) + + require.Equal(8*memsize.KB, store.impl.size) + + // Incomplete files are recovered (since `rebootIncompleteBlobs` is true). + f, err := store.Open(incompleteEvictableKey) + require.NoError(err) + defer func(f io.Closer) { require.NoError(f.Close()) }(f) + data, err := io.ReadAll(f) + require.NoError(err) + require.Equal(incompleteEvictableData, data) + unevictable, err := store.impl.checkDiskIfUnevictable(incompleteEvictableKey, _incompleteBlob) + require.NoError(err) + require.False(unevictable) + + f, err = store.Open(incompleteUnevictableKey) + require.NoError(err) + defer func(f io.Closer) { require.NoError(f.Close()) }(f) + data, err = io.ReadAll(f) + require.NoError(err) + require.Equal(incompleteUnevictableData, data) + unevictable, err = store.impl.checkDiskIfUnevictable(incompleteUnevictableKey, _incompleteBlob) + require.NoError(err) + require.True(unevictable) + + // Complete files are always recovered. + f, err = store.ScopeComplete().Open(completeEvictableKey) + require.NoError(err) + defer func(f io.Closer) { require.NoError(f.Close()) }(f) + data, err = io.ReadAll(f) + require.NoError(err) + require.Equal(completeEvictableData, data) + unevictable, err = store.impl.checkDiskIfUnevictable(completeEvictableKey, _completeBlob) + require.NoError(err) + require.False(unevictable) + + f, err = store.ScopeComplete().Open(completeUnevictableKey) + require.NoError(err) + defer func(f io.Closer) { require.NoError(f.Close()) }(f) + data, err = io.ReadAll(f) + require.NoError(err) + require.Equal(completeUnevictableData, data) + unevictable, err = store.impl.checkDiskIfUnevictable(completeUnevictableKey, _completeBlob) + require.NoError(err) + require.True(unevictable) + + // The sizes of both incomplete and complete files are recovered correctly. + require.Equal(8*memsize.KB, store.impl.size) + + // Run the store with `rebootIncompleteBlobs` as false. + store, err = NewStore(&Config{10 * memsize.KB, rootDir, false, shardLength}, tally.NoopScope) + require.NoError(err) + + // Incomplete files are dropped. + rebootedKeys := store.List() + require.NotContains(rebootedKeys, incompleteEvictableKey) + require.NotContains(rebootedKeys, incompleteUnevictableKey) + + // Complete files are always recovered. + f, err = store.ScopeComplete().Open(completeEvictableKey) + require.NoError(err) + defer func(f io.Closer) { require.NoError(f.Close()) }(f) + data, err = io.ReadAll(f) + require.NoError(err) + require.Equal(completeEvictableData, data) + unevictable, err = store.impl.checkDiskIfUnevictable(completeEvictableKey, _completeBlob) + require.NoError(err) + require.False(unevictable) + + f, err = store.ScopeComplete().Open(completeUnevictableKey) + require.NoError(err) + defer func(f io.Closer) { require.NoError(f.Close()) }(f) + data, err = io.ReadAll(f) + require.NoError(err) + require.Equal(completeUnevictableData, data) + unevictable, err = store.impl.checkDiskIfUnevictable(completeUnevictableKey, _completeBlob) + require.NoError(err) + require.True(unevictable) + + require.Equal(4*memsize.KB, store.impl.size) + } + }) + + t.Run("metadata is recovered for complete blob", func(t *testing.T) { + require := require.New(t) + store, rootDir := newTestStore(t, 10*memsize.KB, true) + + f, key := newTestFile(t, store, 2*memsize.KB) + require.NoError(f.Close()) + require.NoError(store.MarkComplete(key)) + writtenMd := metadata.NewTorrentMeta(core.MetaInfoFixture()) + require.NoError(store.SetMetadata(key, writtenMd)) + + // Assume that the application crashes here. The application would restart and call `NewStore`. + store, err := NewStore(&Config{10 * memsize.KB, rootDir, true, _defaultShardLength}, tally.NoopScope) + require.NoError(err) + + var readMd metadata.TorrentMeta + ok, err := store.ScopeComplete().GetMetadata(key, &readMd) + require.NoError(err) + require.True(ok) + require.Equal(writtenMd.MetaInfo, readMd.MetaInfo) + }) + + t.Run("metadata is recovered for incomplete blob", func(t *testing.T) { + require := require.New(t) + store, rootDir := newTestStore(t, 10*memsize.KB, true) + + f, key := newTestFile(t, store, 2*memsize.KB) + require.NoError(f.Close()) + writtenMd := metadata.NewTorrentMeta(core.MetaInfoFixture()) + require.NoError(store.SetMetadata(key, writtenMd)) + + // Assume that the application crashes here. The application would restart and call `NewStore`. + store, err := NewStore(&Config{10 * memsize.KB, rootDir, true, _defaultShardLength}, tally.NoopScope) + 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) + }) + + t.Run("lru order is approximated and blobs are evicted if store size exceeds capacity", func(t *testing.T) { + require := require.New(t) + store, rootDir := newTestStore(t, 10*memsize.KB, false) + + aF, aKey := newTestFile(t, store, 2*memsize.KB) + _ = fillWithRandomData(t, aF, 2*memsize.KB) + require.NoError(aF.Close()) + require.NoError(store.MarkComplete(aKey)) + require.NoError(store.BanEviction(aKey)) + + bF, _ := newTestFile(t, store, 2*memsize.KB) + _ = fillWithRandomData(t, bF, 1*memsize.KB) + require.NoError(bF.Close()) + + cF, cKey := newTestFile(t, store, 2*memsize.KB) + _ = fillWithRandomData(t, cF, 2*memsize.KB) + require.NoError(cF.Close()) + require.NoError(store.MarkComplete(cKey)) + + dF, dKey := newTestFile(t, store, 2*memsize.KB) + _ = fillWithRandomData(t, dF, 2*memsize.KB) + require.NoError(dF.Close()) + require.NoError(store.MarkComplete(dKey)) + + eF, eKey := newTestFile(t, store, 2*memsize.KB) + _ = fillWithRandomData(t, eF, 2*memsize.KB) + require.NoError(eF.Close()) + require.NoError(store.MarkComplete(eKey)) + require.Equal([]string{cKey, dKey, eKey}, store.impl.evictionOrder()) // a is unevictable and b is incomplete + + // reset the access time for d + dF, err := store.ScopeComplete().Open(dKey) + require.NoError(err) + require.NoError(dF.Close()) + evictionOrderBeforeCrash := store.impl.evictionOrder() + require.Equal([]string{cKey, eKey, dKey}, evictionOrderBeforeCrash) // a is unevictable and b is incomplete + + // Assume that the application restarts here. + store, err = NewStore(&Config{10 * memsize.KB, rootDir, false, _defaultShardLength}, tally.NoopScope) + require.NoError(err) + + // LRU order is approximated, but not exact. + rebootedEvictionOrder := store.impl.evictionOrder() + require.NotEqual(evictionOrderBeforeCrash, rebootedEvictionOrder) + wantEvictionOrder := []string{cKey, dKey, eKey} + require.Equal(wantEvictionOrder, rebootedEvictionOrder) + + // Assume we redeploy the service with a smaller capacity for the disk store: + store, err = NewStore(&Config{6 * memsize.KB, rootDir, false, _defaultShardLength}, tally.NoopScope) + require.NoError(err) + + // since 10KB of blobs are in store, `c` gets evicted to put the store back within its capacity. + require.NotContains(store.List(), cKey) + + require.Equal([]string{dKey, eKey}, store.impl.evictionOrder()) + require.Contains(store.List(), aKey) + }) + + t.Run("reboot fails if unevictable data is more than store capacity", func(t *testing.T) { + require := require.New(t) + store, rootDir := newTestStore(t, 10*memsize.KB, true) + + aF, aKey := newTestFile(t, store, 2*memsize.KB) + _ = fillWithRandomData(t, aF, 2*memsize.KB) + require.NoError(aF.Close()) + require.NoError(store.MarkComplete(aKey)) + require.NoError(store.BanEviction(aKey)) + + bF, _ := newTestFile(t, store, 2*memsize.KB) + _ = fillWithRandomData(t, bF, 1*memsize.KB) + require.NoError(bF.Close()) + + cF, _ := newTestFile(t, store, 2*memsize.KB) + _ = fillWithRandomData(t, cF, 1*memsize.KB) + require.NoError(cF.Close()) + + // 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") + }) +} + +func TestIncompleteBlobDownloadResumedAfterMultipleCrashes(t *testing.T) { + require := require.New(t) + store, rootDir := newTestStore(t, 10*memsize.KB, true) + + f, key := newTestFile(t, store, 4*memsize.KB) + firstData := fillWithRandomData(t, f, 2*memsize.KB) + + // First crash. + store, err := NewStore(&Config{10 * memsize.KB, rootDir, true, _defaultShardLength}, tally.NoopScope) + require.NoError(err) + + f, err = store.Open(key) + require.NoError(err) + + // Write 5KB in total, while reporting only 4KB. + secondData := make([]byte, 3*memsize.KB) + _, err = rand.Read(secondData) + require.NoError(err) + _, err = f.WriteAt(secondData, int64(2*memsize.KB)) + require.NoError(err) + + // Second crash. + store, err = NewStore(&Config{10 * memsize.KB, rootDir, true, _defaultShardLength}, tally.NoopScope) + require.NoError(err) + + wantData := make([]byte, 5*memsize.KB) + copy(wantData, firstData) + copy(wantData[2*memsize.KB:], secondData) + + f, err = store.Open(key) + require.NoError(err) + data, err := io.ReadAll(f) + require.NoError(err) + require.Equal(wantData, data) + // The declared size is recovered correctly both times, without leaking or duplicating reservation. + require.Equal(4*memsize.KB, store.impl.size) + require.NoError(store.MarkComplete(key)) + + // Third crash. + store, err = NewStore(&Config{10 * memsize.KB, rootDir, true, _defaultShardLength}, tally.NoopScope) + require.NoError(err) + // now that the blob is complete, its actual size should be rebooted through stat, instead of trusting the _size sidecar file. + require.Equal(5*memsize.KB, store.impl.size) + f, err = store.ScopeComplete().Open(key) + require.NoError(err) + defer func(f io.Closer) { require.NoError(f.Close()) }(f) + data, err = io.ReadAll(f) + require.NoError(err) + require.Equal(wantData, data) +} + +func TestStoreWorksWhenFileSizeNotCorrect(t *testing.T) { + // Verify that the store works correctly when the reserved size for a file (the one passed by the client in Create) is different + // than its actual size. The store is expected to consistently use EITHER the client-given size OR the actual size of files, but not both. + // If we mix them, this could break the eviction logic - imagine the user uploads a 2GB size but reports it as 1.9GB. Eviction works correctly + // as long as we reserve 2GB upon upload to store and release 2GB upon deletion/eviction from store. BUT if we reserve 2GB and free 1.9GB + // or vice-versa, it could lead to over/under-reservation. + require := require.New(t) + store, _ := newTestStore(t, 10*memsize.KB, false) + + // Declares 8KB but only writes 1KB. + underreportedF, underreportedKey := newTestFile(t, store, 8*memsize.KB) + _ = fillWithRandomData(t, underreportedF, 1*memsize.KB) + require.NoError(underreportedF.Close()) + require.NoError(store.MarkComplete(underreportedKey)) + require.NoError(store.BanEviction(underreportedKey)) + require.Equal(8*memsize.KB, store.impl.size) + + // 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") + + // Declares 2KB (fits exactly within the remaining capacity) but writes 5KB. + overreportedF, overreportedKey := newTestFile(t, store, 2*memsize.KB) + _ = fillWithRandomData(t, overreportedF, 5*memsize.KB) + require.NoError(overreportedF.Close()) + require.NoError(store.MarkComplete(overreportedKey)) + // The store still only accounts for the declared 2KB, not the actual 5KB written. + require.Equal(10*memsize.KB, store.impl.size) + + // Deleting releases exactly the declared size that was reserved, not the actual bytes found on + // disk, keeping reservation accounting self-consistent in both directions. + require.NoError(store.Delete(overreportedKey)) + require.Equal(8*memsize.KB, store.impl.size) + require.NoError(store.Delete(underreportedKey)) + require.Equal(uint64(0), store.impl.size) +} + +func fillWithRandomData(t *testing.T, f storelib.FileReadWriter, sizeBytes uint64) []byte { + data := make([]byte, sizeBytes) + _, err := rand.Read(data) + require.NoError(t, err) + _, err = io.Copy(f, bytes.NewReader(data)) + require.NoError(t, err) + return data +} diff --git a/lib/store/disk/pather.go b/lib/store/disk/pather.go new file mode 100644 index 000000000..86c3b50f3 --- /dev/null +++ b/lib/store/disk/pather.go @@ -0,0 +1,96 @@ +package disk + +import ( + "fmt" + "io/fs" + "path/filepath" + "strings" +) + +const ( + _incompleteSubDir = "incomplete" + _completeSubDir = "complete" + _blobFileName = "data" +) + +type pather struct { + dir string + shardLength int +} + +// shardLength is the number of bytes of blob key to be used for shard ID. +// For every byte (2 HEX char), one more level of directories will be created. +// A shardLength of 0 denotes no sharding. +func newPather(rootDir string, shardLength int) *pather { + return &pather{dir: rootDir, shardLength: shardLength} +} + +func (p *pather) blobPath(key string, complete bool) string { + dirName := p.dirPath(key, complete) + return filepath.Join(dirName, _blobFileName) +} + +func (p *pather) dirPath(key string, complete bool) string { + subDirName := _incompleteSubDir + if complete { + subDirName = _completeSubDir + } + dirPath := filepath.Join(p.dir, subDirName) + for i := 0; i < int(p.shardLength) && i < len(key)/2; i++ { + // (1 byte = 2 char of file name assuming file name is in HEX) + dirName := key[i*2 : i*2+2] + dirPath = filepath.Join(dirPath, dirName) + } + + return filepath.Join(dirPath, key) +} + +func (p *pather) sidecarFilePath(key string, complete bool, sidecarFileName string) string { + dirPath := p.dirPath(key, complete) + return filepath.Join(dirPath, sidecarFileName) +} + +func (p *pather) rebootKeys(complete bool) ([]string, error) { + subDirName := _incompleteSubDir + if complete { + subDirName = _completeSubDir + } + dir := filepath.Join(p.dir, subDirName) + keys := make([]string, 0) + ok, err := exists(dir) + if err != nil { + return nil, fmt.Errorf("exists: %w", err) + } + if !ok { + return []string{}, nil + } + err = filepath.WalkDir(dir, func(path string, entry fs.DirEntry, err error) error { + if err != nil { + return err + } + + if !entry.IsDir() { + return nil + } + if path == dir { + return nil + } + relPath, err := filepath.Rel(dir, path) + if err != nil { + return err + } + nameParts := strings.Split(relPath, string(filepath.Separator)) + numShards := p.shardLength + isBlobDir := len(nameParts) == numShards+1 + if !isBlobDir { + return nil + } + key := nameParts[len(nameParts)-1] + keys = append(keys, key) + return fs.SkipDir + }) + if err != nil { + return nil, fmt.Errorf("walk through dir '%v' to reboot blob keys: %w", dir, err) + } + return keys, nil +} diff --git a/lib/store/disk/pather_test.go b/lib/store/disk/pather_test.go new file mode 100644 index 000000000..6490e4c6b --- /dev/null +++ b/lib/store/disk/pather_test.go @@ -0,0 +1,55 @@ +package disk + +import ( + "os" + "testing" + + "github.com/stretchr/testify/require" + "github.com/uber/kraken/core" + "github.com/uber/kraken/lib/store/metadata" +) + +func TestPather(t *testing.T) { + rootDir, err := os.MkdirTemp("/tmp", "kraken-disk-store") + require.NoError(t, err) + t.Cleanup(func() { require.NoError(t, os.RemoveAll(rootDir)) }) + shardLength := 2 + shardedPather := newPather(rootDir, shardLength) + unshardedPather := newPather(rootDir, 0) + md := metadata.NewTorrentMeta(core.MetaInfoFixture()) + key := "8c6af6ca6458353bfa8cb3d756ca54a4fe7b1de04196bf1b37e0863c3f806a78" + + t.Run("incomplete entry, sharded", func(t *testing.T) { + require := require.New(t) + require.Equal(rootDir+"/incomplete/8c/6a/8c6af6ca6458353bfa8cb3d756ca54a4fe7b1de04196bf1b37e0863c3f806a78", shardedPather.dirPath(key, _incompleteBlob)) + require.Equal(rootDir+"/incomplete/8c/6a/8c6af6ca6458353bfa8cb3d756ca54a4fe7b1de04196bf1b37e0863c3f806a78/data", shardedPather.blobPath(key, _incompleteBlob)) + require.Equal(rootDir+"/incomplete/8c/6a/8c6af6ca6458353bfa8cb3d756ca54a4fe7b1de04196bf1b37e0863c3f806a78/_torrentmeta", shardedPather.sidecarFilePath(key, _incompleteBlob, md.GetSuffix())) + }) + t.Run("incomplete entry, unsharded", func(t *testing.T) { + require := require.New(t) + require.Equal(rootDir+"/incomplete/8c6af6ca6458353bfa8cb3d756ca54a4fe7b1de04196bf1b37e0863c3f806a78", unshardedPather.dirPath(key, _incompleteBlob)) + require.Equal(rootDir+"/incomplete/8c6af6ca6458353bfa8cb3d756ca54a4fe7b1de04196bf1b37e0863c3f806a78/data", unshardedPather.blobPath(key, _incompleteBlob)) + require.Equal(rootDir+"/incomplete/8c6af6ca6458353bfa8cb3d756ca54a4fe7b1de04196bf1b37e0863c3f806a78/_torrentmeta", unshardedPather.sidecarFilePath(key, _incompleteBlob, md.GetSuffix())) + }) + + t.Run("complete entry, sharded", func(t *testing.T) { + require := require.New(t) + require.Equal(rootDir+"/complete/8c/6a/8c6af6ca6458353bfa8cb3d756ca54a4fe7b1de04196bf1b37e0863c3f806a78", shardedPather.dirPath(key, _completeBlob)) + require.Equal(rootDir+"/complete/8c/6a/8c6af6ca6458353bfa8cb3d756ca54a4fe7b1de04196bf1b37e0863c3f806a78/data", shardedPather.blobPath(key, _completeBlob)) + require.Equal(rootDir+"/complete/8c/6a/8c6af6ca6458353bfa8cb3d756ca54a4fe7b1de04196bf1b37e0863c3f806a78/_torrentmeta", shardedPather.sidecarFilePath(key, _completeBlob, md.GetSuffix())) + }) + t.Run("complete entry, unsharded", func(t *testing.T) { + require := require.New(t) + require.Equal(rootDir+"/complete/8c6af6ca6458353bfa8cb3d756ca54a4fe7b1de04196bf1b37e0863c3f806a78", unshardedPather.dirPath(key, _completeBlob)) + require.Equal(rootDir+"/complete/8c6af6ca6458353bfa8cb3d756ca54a4fe7b1de04196bf1b37e0863c3f806a78/data", unshardedPather.blobPath(key, _completeBlob)) + require.Equal(rootDir+"/complete/8c6af6ca6458353bfa8cb3d756ca54a4fe7b1de04196bf1b37e0863c3f806a78/_torrentmeta", unshardedPather.sidecarFilePath(key, _completeBlob, md.GetSuffix())) + }) + + t.Run("unsharded entries with non-digest keys", func(t *testing.T) { + key := "foo" + require := require.New(t) + require.Equal(rootDir+"/complete/foo", unshardedPather.dirPath(key, _completeBlob)) + require.Equal(rootDir+"/complete/foo/data", unshardedPather.blobPath(key, _completeBlob)) + require.Equal(rootDir+"/complete/foo/_torrentmeta", unshardedPather.sidecarFilePath(key, _completeBlob, md.GetSuffix())) + }) +} diff --git a/lib/store/disk/scoped_store.go b/lib/store/disk/scoped_store.go new file mode 100644 index 000000000..75cbf5389 --- /dev/null +++ b/lib/store/disk/scoped_store.go @@ -0,0 +1,113 @@ +package disk + +import ( + "errors" + "os" + + "github.com/uber-go/tally" + storelib "github.com/uber/kraken/lib/store" + "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. +// +// - New blobs are considered 'incomplete', which unlists them from LRU eviction. The store can be scoped to work on only (in-)complete blobs. +// +// - All APIs are thread-safe. Parallel access to a single file is allowed but clients must ensure they don't intervene with one another. +// +// - Supports (un-)marking blobs as non-evictable (may be needed when that data must be written back to remote storage). +// +// - Crash-resistant - all state is restored upon restart (check [newStore] for details). +// +// - Supports directory sharding to speed up disk performance. +type Store struct { + impl *store + scope blobScope +} + +// NewStore initializes a new [*Store]. If the store has been initialized in the same +// directory before, its state is recovered from disk with the following caveats: +// +// - `rebootIncompleteBlobs` configures whether incomplete blobs are evicted or rebooted on restart. +// +// - If the store's size is bigger than its capacity (e.g. configured capacity has been reduced or files have been leaked), +// it evicts blobs until size is within capacity. +func NewStore(config *Config, metrics tally.Scope) (*Store, error) { + s, err := newStore(config, metrics) + if err != nil { + return nil, err + } + return &Store{ + impl: s, + scope: 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) { + 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) } + +// 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) } + +// MarkComplete marks a blob as fully written. It enlists the blob for LRU eviction (unless BanEviction has been called). +// 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. +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. +// Needed when e.g. blobs must be written back to GCS/S3 and eviction before that is unacceptable. +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 atomically sets the respective metadata for a 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) +} + +// DeleteMetadata removes any metadata of a blob with `md`'s suffix, if present. +func (s *Store) DeleteMetadata(key string, md metadata.Metadata) error { + return s.impl.DeleteMetadata(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) +} + +// WriteAtMetadata implements [io.WriterAt] for the metadata file on disk. +func (s *Store) WriteAtMetadata(key string, md metadata.Metadata, p []byte, off int64) error { + return s.impl.WriteAtMetadata(key, md, p, off, s.scope) +} + +// 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} } + +// 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} } diff --git a/lib/store/disk/store.go b/lib/store/disk/store.go new file mode 100644 index 000000000..6c3ba7ada --- /dev/null +++ b/lib/store/disk/store.go @@ -0,0 +1,600 @@ +package disk + +import ( + "container/list" + "errors" + "fmt" + "io" + "os" + "path/filepath" + "strconv" + "sync" + "time" + + "github.com/uber-go/tally" + storelib "github.com/uber/kraken/lib/store" + "github.com/uber/kraken/lib/store/metadata" + "github.com/uber/kraken/utils/closers" + "github.com/uber/kraken/utils/log" + "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 + _defaultFilePerm = 0775 + _evictionBannedFileName = "_eviction_banned" + _blobSizeFileName = "_size" +) + +var _syncEvictionLatencyBuckets = tally.MustMakeExponentialDurationBuckets(100*time.Millisecond, 1.4, 15) + +// store implements the APIs of [Store]. [store]'s APIs expose the [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 { + capacity uint64 + size uint64 // Includes both actively-used and not yet used, but reserved space. + blobs map[string]*blob // TODO - consider whether it's better to use struct instead of pointer to reduce GC stress. + evictQueue *list.List // Front is the next to evict. + // Synchronizes 1) mem state and 2) disk state. Only acquired by public methods. + mu sync.RWMutex // TODO - evaluate whether the read-to-write ratio is more appropriate for a [sync.Mutex] instead. + config *Config + log *zap.SugaredLogger + metrics tally.Scope + *pather +} + +type blob struct { + node *list.Element // Value of [list.Element] is [string]. + size uint64 + complete bool + evictionBanned bool +} + +func newStore(config *Config, metrics tally.Scope) (*store, error) { + log := log.Default().With("module", "disk_store") + ok, err := existsPersistedStore(config.RootDir) + if err != nil { + err = fmt.Errorf("could not check if previously-left persisted state exists on disk: %w", err) + log.With("error", err).Error("Failed to initialize disk store") + return nil, err + } + if !ok { + log.Info("Initialized a new, empty Store (did not find any previously persisted state to reboot for Store)") + return &store{ + capacity: config.CapacityBytes, + size: 0, + blobs: make(map[string]*blob), + evictQueue: list.New(), + config: config, + pather: newPather(config.RootDir, config.ShardLength), + log: log, + metrics: metrics, + }, nil + } + + store, err := rebootPersistedStore(config, log, metrics) + if err != nil { + err = fmt.Errorf("reboot persisted state into memory: %w", err) + log.With("error", err).Error("Failed to initialize disk store") + return nil, err + } + log.With("num_blobs", len(store.blobs)).Info("Successfully rebooted Store's previously left state on disk") + return store, nil +} + +func (s *store) Open(key string, scope blobScope) (storelib.FileReadWriter, 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) + } + path := s.blobPath(key, b.complete) + f, err := os.OpenFile(path, os.O_RDWR, _defaultFilePerm) + if err != nil { + return nil, fmt.Errorf("open: %w", err) + } + return storelib.NewReadWriter(f), nil +} + +func (s *store) Stat(key string, scope blobScope) (os.FileInfo, 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 + } + // TODO - consider whether this should update the access time of the blob. + blobPath := s.blobPath(key, b.complete) + return os.Stat(blobPath) +} + +func (s *store) Create(key string, sizeBytes uint64) (storelib.FileReadWriter, 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() + + 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 err := s.reserveSpace(sizeBytes); err != nil { + return nil, fmt.Errorf("reserve space: %w", err) + } + + dirName := s.dirPath(key, _incompleteBlob) + err := os.MkdirAll(dirName, _defaultFilePerm) + if err != nil { + s.releaseSpace(sizeBytes) + return nil, fmt.Errorf("ensure dir: %w", err) + } + blobPath := s.blobPath(key, _incompleteBlob) + f, err := os.OpenFile(blobPath, os.O_RDWR|os.O_CREATE|os.O_EXCL, _defaultFilePerm) + if err != nil { + s.releaseSpace(sizeBytes) + return nil, fmt.Errorf("open file: %w", err) + } + + if s.config.RebootIncompleteBlobs { + err = s.persistBlobSize(key, sizeBytes) + if err != nil { + // Fail-open: the blob will be discarded upon reboot if incomplete. + s.log.With("error", err).Error("Could not persist client-provided blob size on disk") + } + } + + s.blobs[key] = &blob{ + size: sizeBytes, + node: nil, + complete: false, + evictionBanned: false, + } + + return storelib.NewReadWriter(f), nil +} + +func (s *store) persistBlobSize(key string, sizeBytes uint64) error { + blobSizeFilePath := s.sidecarFilePath(key, _incompleteBlob, _blobSizeFileName) + blobSizeF, err := os.OpenFile(blobSizeFilePath, os.O_WRONLY|os.O_CREATE|os.O_EXCL, _defaultFilePerm) + if err != nil { + return fmt.Errorf("create file: %w", err) + } + defer closers.Close(blobSizeF) + _, err = blobSizeF.Write([]byte(strconv.Itoa(int(sizeBytes)))) + if err != nil { + return fmt.Errorf("write to size file: %w", err) + } + return nil +} + +func (s *store) reserveSpace(space uint64) error { + // TODO - benchmark and consider whether async eviction makes more sense. + startTime := time.Now() + for s.size+space > s.capacity { + if s.evictQueue.Len() == 0 { + s.log.With( + "unevictable_bytes", s.size, + "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") + } + + toEvictNode := s.evictQueue.Front() + toEvictKey := toEvictNode.Value.(string) //nolint:errcheck // We only ever store string in this value. + + err := s.deleteFromDisk(toEvictKey, _completeBlob) + if err != nil { + return fmt.Errorf("delete from disk: %w", err) + } + s.evictQueue.Remove(toEvictNode) + size := s.blobs[toEvictKey].size + s.releaseSpace(size) + delete(s.blobs, toEvictKey) + } + s.size += space + + latency := time.Since(startTime) + s.metrics.Histogram("sync_eviction_latency", _syncEvictionLatencyBuckets).RecordDuration(latency) + return nil +} + +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.size = 0 + return + } + s.size -= space +} + +// Fully deletes the disk state of a blob, including metadata. Works on any blob. +func (s *store) deleteFromDisk(key string, complete bool) error { + dir := s.dirPath(key, complete) + return os.RemoveAll(dir) +} + +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 + } + + oldPathDir := s.dirPath(key, _incompleteBlob) + newPathDir := s.dirPath(key, _completeBlob) + err := os.MkdirAll(filepath.Dir(newPathDir), _defaultFilePerm) + if err != nil { + return fmt.Errorf("mkdirall: %w", err) + } + err = os.Rename(oldPathDir, newPathDir) + if err != nil { + return fmt.Errorf("move dir: %w", err) + } + b.complete = true + if !b.evictionBanned { + node := s.evictQueue.PushBack(key) + b.node = node + } + + s.tryDeleteImmovableMetadata(key) + return nil +} + +// Best-effort attempt to remove immovable metadata. Fail-open on failure +// 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) + 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") + return + } + for _, md := range mdList { + if md.Movable() { + continue + } + mdFilePath := s.sidecarFilePath(key, _completeBlob, md.GetSuffix()) + err = os.Remove(mdFilePath) + if err != nil && !errors.Is(err, os.ErrNotExist) { + err = fmt.Errorf("remove metadata file: %w", err) + s.log.With("error", err).Error("Failed to delete un-movable metadata upon marking a blob as complete") + continue + } + } +} + +func (s *store) checkDiskIfUnevictable(key string, complete bool) (bool, error) { + flagBlobPath := s.sidecarFilePath(key, complete, _evictionBannedFileName) + unevictable, err := exists(flagBlobPath) + if err != nil { + return false, err + } + return unevictable, nil +} + +func (s *store) Delete(key string, scope 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 + } + err := s.deleteFromDisk(key, b.complete) + if err != nil { + return fmt.Errorf("delete from disk: %w", err) + } + if b.node != nil { + s.evictQueue.Remove(b.node) + b.node = nil + } + delete(s.blobs, key) + s.releaseSpace(b.size) + + return nil +} + +func (s *store) list(scope 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) BanEviction(key string, scope 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 + } + + flagBlobPath := s.sidecarFilePath(key, b.complete, _evictionBannedFileName) + // We persist the ban as a flag file on disk for crash-resilience. + f, err := os.OpenFile(flagBlobPath, os.O_RDONLY|os.O_CREATE, _defaultFilePerm) + if err != nil { + return fmt.Errorf("create file that flags eviction as banned: %w", err) + } + closers.Close(f) + + b.evictionBanned = true + if b.complete { + s.evictQueue.Remove(b.node) + b.node = nil + } + return nil +} + +func (s *store) UnbanEviction(key string, scope 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 + } + + flagBlobPath := s.sidecarFilePath(key, b.complete, _evictionBannedFileName) + err := os.Remove(flagBlobPath) + if err != nil { + return fmt.Errorf("remove file that flags eviction as banned: %w", err) + } + + 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 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 + } + + mdData, err := md.Serialize() + if err != nil { + return fmt.Errorf("serialize metadata: %w", err) + } + mdFilePath := s.sidecarFilePath(key, b.complete, md.GetSuffix()) + // Use a tmp file to ensure atomicity. + tmpFilePath := mdFilePath + "-tmp" + tmpFile, err := os.OpenFile(tmpFilePath, os.O_RDWR|os.O_CREATE|os.O_TRUNC, _defaultFilePerm) + if err != nil { + return fmt.Errorf("create tmp file for md: %w", err) + } + defer closers.Close(tmpFile) + _, err = tmpFile.Write(mdData) + if err != nil { + return fmt.Errorf("write to tmp file: %w", err) + } + err = os.Rename(tmpFile.Name(), mdFilePath) + if err != nil { + return fmt.Errorf("rename tmp file: %w", err) + } + return nil +} + +func (s *store) GetMetadata(key string, md metadata.Metadata, scope 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 + } + + mdFilePath := s.sidecarFilePath(key, b.complete, md.GetSuffix()) + mdFile, err := os.OpenFile(mdFilePath, os.O_RDONLY, _defaultFilePerm) + if errors.Is(err, os.ErrNotExist) { + return false, nil + } + if err != nil { + return false, fmt.Errorf("open metadata file: %w", err) + } + defer closers.Close(mdFile) + data, err := io.ReadAll(mdFile) + if err != nil { + return false, fmt.Errorf("read from metadata file: %w", err) + } + err = md.Deserialize(data) + if err != nil { + return false, fmt.Errorf("deserialize into metadata: %w", err) + } + return true, nil +} + +func (s *store) DeleteMetadata(key string, md metadata.Metadata, scope 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 + } + mdFilePath := s.sidecarFilePath(key, b.complete, md.GetSuffix()) + err := os.Remove(mdFilePath) + if errors.Is(err, os.ErrNotExist) { + // no-op + return nil + } + if err != nil { + return fmt.Errorf("remove metadata file: %w", err) + } + return nil +} + +func (s *store) ListMetadata(key string, scope 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) { + b, ok := s.blobs[key] + if !ok { + return nil, os.ErrNotExist + } + if err := isOutOfScope(b, scope); err != nil { + return nil, err + } + + res := make([]metadata.Metadata, 0) + entries, err := os.ReadDir(s.dirPath(key, b.complete)) + if err != nil { + return nil, fmt.Errorf("read metadata files from dir: %w", err) + } + for _, entry := range entries { + name := entry.Name() + if name == _blobFileName || name == _evictionBannedFileName || name == _blobSizeFileName { + continue + } + + md := metadata.CreateFromSuffix(name) + if md == nil { + s.log.With("key", key, "file_name", name).Warn("Found file in blob dir that does not successfully parse as a metadata file") + continue + } + res = append(res, md) + } + return res, nil +} + +func (s *store) WriteAtMetadata(key string, md metadata.Metadata, p []byte, off int64, scope 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 + } + + mdFilePath := s.sidecarFilePath(key, b.complete, md.GetSuffix()) + mdFile, err := os.OpenFile(mdFilePath, os.O_WRONLY, _defaultFilePerm) + if errors.Is(err, os.ErrNotExist) { + return errors.New("metadata does not exist") + } + if err != nil { + return fmt.Errorf("open: %w", err) + } + defer closers.Close(mdFile) + + _, err = mdFile.WriteAt(p, off) + if err != nil { + return fmt.Errorf("write at: %w", err) + } + + return nil +} + +// used during testing +func (s *store) evictionOrder() []string { + s.mu.RLock() + defer s.mu.RUnlock() + + evictionOrder := make([]string, 0) + for curr := s.evictQueue.Front(); curr != nil; curr = curr.Next() { + currKey := curr.Value.(string) //nolint:errcheck // We only ever store string in this value. + evictionOrder = append(evictionOrder, currKey) + } + return evictionOrder +} + +func exists(path string) (ok bool, err error) { + _, err = os.Stat(path) + if err == nil { + return true, nil + } + if errors.Is(err, os.ErrNotExist) { + return false, nil + } + 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 + } + return nil +} diff --git a/lib/store/disk/store_test.go b/lib/store/disk/store_test.go new file mode 100644 index 000000000..fcbcf85a6 --- /dev/null +++ b/lib/store/disk/store_test.go @@ -0,0 +1,1054 @@ +package disk + +import ( + "bytes" + "io" + "io/fs" + "os" + "path/filepath" + "regexp" + "strings" + "sync" + "testing" + "testing/iotest" + "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/metadata" + "github.com/uber/kraken/utils/memsize" +) + +const ( + _defaultShardLength = 2 +) + +func newTestStore(t *testing.T, capacity uint64, rebootIncompleteBlobs bool) (res *Store, rootDir string) { + rootDir, err := os.MkdirTemp("/tmp", "kraken-disk-store") + require.NoError(t, err) + t.Cleanup(func() { require.NoError(t, os.RemoveAll(rootDir)) }) + config := &Config{ + CapacityBytes: capacity, + RootDir: rootDir, + RebootIncompleteBlobs: rebootIncompleteBlobs, + ShardLength: _defaultShardLength, + } + + store, err := NewStore(config, tally.NoopScope) + require.NoError(t, err) + return store, rootDir +} + +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 +} + +// does not count 1) the directories for sharding, 2) metadata files, and 3) the _eviction_banned flag file. +func numBlobsOnDisk(t *testing.T, store *Store) int { + numBlobs := 0 + err := filepath.WalkDir(store.impl.dir, func(path string, _ fs.DirEntry, err error) error { + if err != nil { + return err + } + + if !strings.HasSuffix(path, _blobFileName) { + return nil + } + + numBlobs++ + return nil + }) + require.NoError(t, err) + return numBlobs +} + +func TestStore(t *testing.T) { + require := require.New(t) + store, _ := newTestStore(t, 10*memsize.KB, false) + + keys := []string{} + for i := range 10 { + key := core.DigestFixture().Hex() + keys = append(keys, key) + f, err := store.Create(key, memsize.KB) + defer func(f io.Closer) { require.NoError(f.Close()) }(f) + require.NoError(err) + + data := make([]byte, memsize.KB) + for k := range data { + data[k] = byte(i + 1) + } + n, err := io.Copy(f, bytes.NewReader(data)) + require.Equal(int64(memsize.KB), n) + require.NoError(err) + } + 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.Nil(f) + + f, err = store.ScopeComplete().Open(keys[0]) + require.ErrorIs(err, ErrOutOfScope) + require.Nil(f) + + require.NoError(store.MarkComplete(keys[0])) + f, err = store.Open(keys[0]) + defer func(f io.Closer) { require.NoError(f.Close()) }(f) + require.NoError(err) + wantData := make([]byte, memsize.KB) + for k := range wantData { + wantData[k] = byte(1) + } + require.NoError(iotest.TestReader(f, wantData)) + + // now test that LRU logic works - make sure that the 1 that is complete gets evicted. + f, err = store.Create(core.DigestFixture().Hex(), memsize.KB) + require.NoError(err) + defer func(f io.Closer) { require.NoError(f.Close()) }(f) + f, err = store.Open(keys[0]) + require.ErrorIs(err, os.ErrNotExist) + require.Nil(f) +} + +func TestEviction(t *testing.T) { + require := require.New(t) + store, _ := newTestStore(t, 25*memsize.KB, false) + // create 5 blobs - a, b, c, d, e with different sizes. + a, aKey := newTestFile(t, store, 10*memsize.KB) + require.NoError(a.Close()) + b, bKey := newTestFile(t, store, 5*memsize.KB) + require.NoError(b.Close()) + c, cKey := newTestFile(t, store, 5*memsize.KB) + require.NoError(c.Close()) + d, dKey := newTestFile(t, store, 3*memsize.KB) + require.NoError(d.Close()) + e, eKey := newTestFile(t, store, 1*memsize.KB) + require.NoError(e.Close()) + + require.Equal(24*memsize.KB, store.impl.size) + 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") + // 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)) + keys := store.List() + require.NotContains(keys, cKey) + // new size == 23KB == 24KB - 5KB (c) + 4KB (f) + require.Equal(23*memsize.KB, store.impl.size) + require.Equal(5, numBlobsOnDisk(t, store)) + + // 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) + require.Equal(6, numBlobsOnDisk(t, store)) + + // 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) + require.Equal(5, numBlobsOnDisk(t, store)) + keys = store.List() + require.NotContains(keys, bKey) + require.NotContains(keys, aKey) + + // 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()) + keys = store.List() + require.NotContains(keys, fKey) + 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 + + j, jKey := newTestFile(t, store, 14*memsize.KB) + require.NoError(j.Close()) + require.NoError(store.MarkComplete(jKey)) + keys = store.List() + require.NotContains(keys, hKey) + 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 + + k, kKey := newTestFile(t, store, 2*memsize.KB) + require.NoError(k.Close()) + require.NoError(store.MarkComplete(kKey)) + keys = store.List() + require.NotContains(keys, eKey) + 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 + + l, lKey := newTestFile(t, store, 1*memsize.KB) + require.NoError(store.MarkComplete(lKey)) + require.NoError(l.Close()) + keys = store.List() + require.NotContains(keys, gKey) + 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 + + 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) + require.Equal(2, numBlobsOnDisk(t, store)) +} + +func TestParallelAccessToSingleFile(t *testing.T) { + require := require.New(t) + store, _ := newTestStore(t, 10*memsize.KB, false) + + key := core.DigestFixture().Hex() + f, err := store.Create(key, 1*memsize.KB) + 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, _ := newTestStore(t, 10*memsize.KB, false) + + key := core.DigestFixture().Hex() + f, err := store.Create(key, 1*memsize.KB) + require.NoError(err) + _, err = io.Copy(f, bytes.NewReader([]byte("Hello World"))) + 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([]byte("Hello World"), incompleteFileData) + require.Equal([]byte("Hello World"), completeFileData) +} + +func TestDelete(t *testing.T) { + t.Run("incomplete blob", func(t *testing.T) { + require := require.New(t) + store, _ := newTestStore(t, 10*memsize.KB, false) + 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, numBlobsOnDisk(t, store)) + }) + t.Run("incomplete, unevictable blob", func(t *testing.T) { + require := require.New(t) + store, _ := newTestStore(t, 10*memsize.KB, false) + 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, numBlobsOnDisk(t, store)) + }) + t.Run("complete blob", func(t *testing.T) { + require := require.New(t) + store, _ := newTestStore(t, 10*memsize.KB, false) + 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, numBlobsOnDisk(t, store)) + }) + t.Run("complete, unevictable blob", func(t *testing.T) { + require := require.New(t) + store, _ := newTestStore(t, 10*memsize.KB, false) + 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, numBlobsOnDisk(t, store)) + }) + t.Run("not found", func(t *testing.T) { + require := require.New(t) + store, _ := newTestStore(t, 10*memsize.KB, false) + key := core.DigestFixture().Hex() + + err := store.Delete(key) + require.ErrorIs(os.ErrNotExist, err) + }) +} + +func TestMarkComplete(t *testing.T) { + t.Run("incomplete blob", func(t *testing.T) { + require := require.New(t) + store, _ := newTestStore(t, 10*memsize.KB, false) + 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, _ := newTestStore(t, 10*memsize.KB, false) + 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, _ := newTestStore(t, 10*memsize.KB, false) + 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, _ := newTestStore(t, 10*memsize.KB, false) + 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, _ := newTestStore(t, 10*memsize.KB, false) + key := core.DigestFixture().Hex() + + err := store.MarkComplete(key) + require.ErrorIs(os.ErrNotExist, err) + }) +} + +func TestStat(t *testing.T) { + t.Run("complete blob", func(t *testing.T) { + require := require.New(t) + store, _ := newTestStore(t, 10*memsize.KB, false) + 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)) + + fInfo, err := store.ScopeComplete().Stat(key) + require.NoError(err) + _, err = store.Stat(key) + require.NoError(err) + + require.False(fInfo.IsDir()) + require.WithinDuration(time.Now(), fInfo.ModTime(), 500*time.Millisecond) + require.Equal(_blobFileName, fInfo.Name()) + require.Equal(int64(10), fInfo.Size()) + require.Equal(fs.FileMode(0755), fInfo.Mode()) + }) + t.Run("complete, unevictable blob", func(t *testing.T) { + require := require.New(t) + store, _ := newTestStore(t, 10*memsize.KB, false) + 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)) + + fInfo, err := store.ScopeComplete().Stat(key) + require.NoError(err) + _, err = store.Stat(key) + require.NoError(err) + + require.False(fInfo.IsDir()) + require.WithinDuration(time.Now(), fInfo.ModTime(), 500*time.Millisecond) + require.Equal(_blobFileName, fInfo.Name()) + require.Equal(int64(10), fInfo.Size()) + require.Equal(fs.FileMode(0755), fInfo.Mode()) + }) + t.Run("incomplete blob", func(t *testing.T) { + require := require.New(t) + store, _ := newTestStore(t, 10*memsize.KB, false) + 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(ErrOutOfScope, err) + fInfo, err := store.Stat(key) + require.NoError(err) + + require.False(fInfo.IsDir()) + require.WithinDuration(time.Now(), fInfo.ModTime(), 500*time.Millisecond) + require.Equal(_blobFileName, fInfo.Name()) + require.Equal(int64(10), fInfo.Size()) + require.Equal(fs.FileMode(0755), fInfo.Mode()) + }) + + t.Run("incomplete, unevictable blob", func(t *testing.T) { + require := require.New(t) + store, _ := newTestStore(t, 10*memsize.KB, false) + 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)) + + _, err = store.ScopeComplete().Stat(key) + require.Equal(ErrOutOfScope, err) + fInfo, err := store.Stat(key) + require.NoError(err) + + require.False(fInfo.IsDir()) + require.WithinDuration(time.Now(), fInfo.ModTime(), 500*time.Millisecond) + require.Equal(_blobFileName, fInfo.Name()) + require.Equal(int64(10), fInfo.Size()) + require.Equal(fs.FileMode(0755), fInfo.Mode()) + }) + t.Run("non-existent blob", func(t *testing.T) { + require := require.New(t) + store, _ := newTestStore(t, 10*memsize.KB, false) + key := core.DigestFixture().Hex() + + _, err := store.ScopeComplete().Stat(key) + require.ErrorIs(os.ErrNotExist, err) + _, err = store.Stat(key) + require.ErrorIs(os.ErrNotExist, err) + }) +} + +func TestList(t *testing.T) { + require := require.New(t) + store, _ := newTestStore(t, 10*memsize.KB, false) + + require.Empty(store.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, _ := newTestStore(t, 10*memsize.KB, false) + 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(true) + require.NoError(store.SetMetadata(key, persistMd)) + require.NoError(store.WriteAtMetadata(key, persistMd, []byte("false"), 0)) + 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)) + ok, err = store.GetMetadata(key, &readMd) + require.NoError(err) + require.False(ok) + mdFilePath := store.impl.sidecarFilePath(key, _incompleteBlob, readMd.GetSuffix()) + // ensure the metadata file is deleted from disk + _, err = os.Stat(mdFilePath) + require.ErrorIs(err, os.ErrNotExist) + // deleting a second time should be a no-op. + require.NoError(store.DeleteMetadata(key, &readMd)) + + 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)) + 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, _ := newTestStore(t, 10*memsize.KB, false) + nonExistentKey := core.DigestFixture().Hex() + mdStruct := core.MetaInfoFixture() + md := metadata.NewTorrentMeta(mdStruct) + + err := store.SetMetadata(nonExistentKey, md) + require.ErrorIs(os.ErrNotExist, err) + + ok, err := store.GetMetadata(nonExistentKey, md) + require.ErrorIs(os.ErrNotExist, err) + require.False(ok) + + _, err = store.ListMetadata(nonExistentKey) + require.ErrorIs(os.ErrNotExist, err) + + err = store.WriteAtMetadata(nonExistentKey, md, []byte("data"), 0) + require.ErrorIs(os.ErrNotExist, err) + + err = store.DeleteMetadata(nonExistentKey, md) + require.ErrorIs(os.ErrNotExist, err) + }) + + 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, _ := newTestStore(t, 10*memsize.KB, false) + 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(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(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, _ := newTestStore(t, 10*memsize.KB, false) + 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) + mdFilePath := store.impl.sidecarFilePath(keyA, _completeBlob, md.GetSuffix()) + _, err = os.Stat(mdFilePath) + require.NoError(err) + + 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(os.ErrNotExist, err) + require.False(ok) + _, err = os.Stat(mdFilePath) + require.ErrorIs(err, os.ErrNotExist) + }) + + t.Run("immovable metadata is deleted when calling MarkComplete", func(t *testing.T) { + require := require.New(t) + store, _ := newTestStore(t, 10*memsize.KB, false) + 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) + }) + + t.Run("writeAt on metadata that does not exist", func(t *testing.T) { + require := require.New(t) + store, _ := newTestStore(t, 10*memsize.KB, false) + key := core.DigestFixture().Hex() + f, err := store.Create(key, 10*memsize.KB) + require.NoError(err) + require.NoError(f.Close()) + + md := metadata.NewPersist(true) + err = store.WriteAtMetadata(key, md, []byte("true"), 0) + require.EqualError(err, "metadata does not exist") + }) + + t.Run("list metadata excludes non-metadata sidecar files", func(t *testing.T) { + require := require.New(t) + store, _ := newTestStore(t, 10*memsize.KB, true) + key := core.DigestFixture().Hex() + f, err := store.Create(key, 10*memsize.KB) + require.NoError(err) + require.NoError(f.Close()) + require.NoError(store.BanEviction(key)) + // A sidecar of every blob when [Config.RebootIncompleteBlobs] is on. + sizeSidecarFilePath := store.impl.sidecarFilePath(key, _incompleteBlob, _blobSizeFileName) + _, err = os.Stat(sizeSidecarFilePath) + require.NoError(err) + + mdList, err := store.ListMetadata(key) + require.NoError(err) + require.Empty(mdList) + + torrentMd := metadata.NewTorrentMeta(core.MetaInfoFixture()) + persistMd := metadata.NewPersist(true) + lastAccessMd := metadata.NewLastAccessTime(time.Now()) + require.NoError(store.SetMetadata(key, torrentMd)) + require.NoError(store.SetMetadata(key, persistMd)) + require.NoError(store.SetMetadata(key, lastAccessMd)) + + mdList, err = store.ListMetadata(key) + require.NoError(err) + gotSuffixes := make([]string, len(mdList)) + for i, md := range mdList { + gotSuffixes[i] = md.GetSuffix() + } + require.ElementsMatch( + []string{torrentMd.GetSuffix(), persistMd.GetSuffix(), lastAccessMd.GetSuffix()}, + gotSuffixes, + ) + }) +} + +func TestScopes(t *testing.T) { + require := require.New(t) + store, _ := newTestStore(t, 10*memsize.KB, false) + + f, key := newTestFile(t, store, 2*memsize.KB) + data := fillWithRandomData(t, f, 2*memsize.KB) + require.NoError(f.Close()) + md := metadata.NewTorrentMeta(core.MetaInfoFixture()) + mdData, err := md.Serialize() + require.NoError(err) + + // While incomplete, ScopeComplete's APIs reject the blob. + _, err = store.ScopeComplete().Open(key) + require.ErrorIs(err, 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) + var readMd metadata.TorrentMeta + ok, err := store.ScopeComplete().GetMetadata(key, &readMd) + require.ErrorIs(err, 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.NotContains(store.ScopeComplete().List(), key) + + // 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.WriteAtMetadata(key, md, mdData, 0)) + require.NoError(store.DeleteMetadata(key, &readMd)) + require.Contains(store.List(), key) + + // 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().WriteAtMetadata(key, md, mdData, 0)) + require.NoError(store.ScopeIncomplete().DeleteMetadata(key, &readMd)) + require.Contains(store.ScopeIncomplete().List(), key) + + 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) + _, 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) + ok, err = store.ScopeIncomplete().GetMetadata(key, &readMd) + require.ErrorIs(err, 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.NotContains(store.ScopeIncomplete().List(), key) + + // 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.WriteAtMetadata(key, md, mdData, 0)) + require.NoError(store.DeleteMetadata(key, &readMd)) + require.Contains(store.List(), key) + + // 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().WriteAtMetadata(key, md, mdData, 0)) + require.NoError(store.ScopeComplete().DeleteMetadata(key, &readMd)) + require.Contains(store.ScopeComplete().List(), key) +} + +func TestScopesDelete(t *testing.T) { + require := require.New(t) + store, _ := newTestStore(t, 10*memsize.KB, false) + + f, key := newTestFile(t, store, 1*memsize.KB) + require.NoError(f.Close()) + require.ErrorIs(store.ScopeComplete().Delete(key), 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), 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)) + _, 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.NoError(store.Delete(key)) + _, err = store.Stat(key) + require.ErrorIs(err, os.ErrNotExist) +} diff --git a/lib/store/file.go b/lib/store/file.go index 15bb43733..1e10a941c 100644 --- a/lib/store/file.go +++ b/lib/store/file.go @@ -13,10 +13,48 @@ // limitations under the License. package store -import "github.com/uber/kraken/lib/store/base" +import ( + "os" + + "github.com/uber/kraken/lib/store/base" +) // FileReadWriter is a readable, writable file. 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/metadata/metadata.go b/lib/store/metadata/metadata.go index 1c59ca07f..c8861a125 100644 --- a/lib/store/metadata/metadata.go +++ b/lib/store/metadata/metadata.go @@ -37,7 +37,6 @@ func Register(suffix *regexp.Regexp, factory Factory) { } // CreateFromSuffix creates a Metadata obj based on suffix. -// This is not a very efficient method; It's mostly used during reload. func CreateFromSuffix(suffix string) Metadata { for re, factory := range _factories { if re.MatchString(suffix) { diff --git a/lib/torrent/storage/originstorage/torrent_archive.go b/lib/torrent/storage/originstorage/torrent_archive.go index d190ef3d1..bfcb8d8e3 100644 --- a/lib/torrent/storage/originstorage/torrent_archive.go +++ b/lib/torrent/storage/originstorage/torrent_archive.go @@ -74,7 +74,7 @@ func (a *TorrentArchive) CreateTorrent(namespace string, d core.Digest) (storage } // GetTorrent returns a Torrent for an existing file on disk. If the file does -// not exist, attempts to re-fetch the file from the storae backend configured +// not exist, attempts to re-fetch the file from the storage backend configured // for namespace in a background goroutine, and returns os.ErrNotExist. func (a *TorrentArchive) GetTorrent(namespace string, d core.Digest) (storage.Torrent, error) { mi, err := a.getMetaInfo(namespace, d)