From d1379c71753996d04a3d8568f7534c35e3b5ba72 Mon Sep 17 00:00:00 2001 From: Leigh <351529+leighmcculloch@users.noreply.github.com> Date: Tue, 6 Jan 2026 16:34:41 +1000 Subject: [PATCH 01/11] add filesystem datastore implementation --- support/datastore/datastore.go | 2 + support/datastore/filesystem.go | 307 +++++++++++++++ support/datastore/filesystem_test.go | 534 +++++++++++++++++++++++++++ 3 files changed, 843 insertions(+) create mode 100644 support/datastore/filesystem.go create mode 100644 support/datastore/filesystem_test.go diff --git a/support/datastore/datastore.go b/support/datastore/datastore.go index 8ee7f6f24a..f162d3e0c7 100644 --- a/support/datastore/datastore.go +++ b/support/datastore/datastore.go @@ -57,6 +57,8 @@ func NewDataStore(ctx context.Context, datastoreConfig DataStoreConfig) (DataSto return NewGCSDataStore(ctx, datastoreConfig) case "S3": return NewS3DataStore(ctx, datastoreConfig) + case "Filesystem": + return NewFilesystemDataStore(ctx, datastoreConfig) default: return nil, fmt.Errorf("invalid datastore type %v, not supported", datastoreConfig.Type) diff --git a/support/datastore/filesystem.go b/support/datastore/filesystem.go new file mode 100644 index 0000000000..c02113d0dd --- /dev/null +++ b/support/datastore/filesystem.go @@ -0,0 +1,307 @@ +package datastore + +import ( + "bytes" + "context" + "encoding/json" + "errors" + "fmt" + "io" + "os" + "path/filepath" + "strings" + "time" + + "github.com/stellar/go-stellar-sdk/support/log" +) + +var _ DataStore = &FilesystemDataStore{} + +const metadataSuffix = ".metadata.json" + +// FilesystemDataStore implements DataStore for local filesystem storage. +type FilesystemDataStore struct { + basePath string + writeMetadata bool +} + +// NewFilesystemDataStore creates a new FilesystemDataStore from configuration. +func NewFilesystemDataStore(ctx context.Context, datastoreConfig DataStoreConfig) (DataStore, error) { + destinationPath, ok := datastoreConfig.Params["destination_path"] + if !ok { + return nil, errors.New("invalid Filesystem config, no destination_path") + } + + // write_metadata defaults to true + writeMetadata := true + if val, ok := datastoreConfig.Params["write_metadata"]; ok { + writeMetadata = val != "false" + } + + return NewFilesystemDataStoreWithPath(destinationPath, writeMetadata) +} + +// NewFilesystemDataStoreWithPath creates a FilesystemDataStore with the given base path. +func NewFilesystemDataStoreWithPath(basePath string, writeMetadata bool) (DataStore, error) { + // Ensure the base path exists + if err := os.MkdirAll(basePath, 0755); err != nil { + return nil, fmt.Errorf("failed to create base directory %s: %w", basePath, err) + } + + absPath, err := filepath.Abs(basePath) + if err != nil { + return nil, fmt.Errorf("failed to resolve absolute path: %w", err) + } + + log.Debugf("Creating Filesystem datastore at: %s, writeMetadata: %v", absPath, writeMetadata) + + return &FilesystemDataStore{ + basePath: absPath, + writeMetadata: writeMetadata, + }, nil +} + +// fullPath returns the full filesystem path for a given relative path. +func (f *FilesystemDataStore) fullPath(path string) string { + return filepath.Join(f.basePath, path) +} + +// metadataPath returns the path to the metadata sidecar file. +func (f *FilesystemDataStore) metadataPath(path string) string { + return f.fullPath(path) + metadataSuffix +} + +// Exists checks if a file exists in the filesystem. +func (f *FilesystemDataStore) Exists(ctx context.Context, path string) (bool, error) { + _, err := os.Stat(f.fullPath(path)) + if err == nil { + return true, nil + } + if os.IsNotExist(err) { + return false, nil + } + return false, err +} + +// Size returns the size of a file in bytes. +func (f *FilesystemDataStore) Size(ctx context.Context, path string) (int64, error) { + info, err := os.Stat(f.fullPath(path)) + if err != nil { + if os.IsNotExist(err) { + return 0, os.ErrNotExist + } + return 0, err + } + return info.Size(), nil +} + +// GetFileLastModified returns the last modification time of a file. +func (f *FilesystemDataStore) GetFileLastModified(ctx context.Context, path string) (time.Time, error) { + info, err := os.Stat(f.fullPath(path)) + if err != nil { + if os.IsNotExist(err) { + return time.Time{}, os.ErrNotExist + } + return time.Time{}, err + } + return info.ModTime(), nil +} + +// GetFile returns a reader for the file at the given path. +func (f *FilesystemDataStore) GetFile(ctx context.Context, path string) (io.ReadCloser, error) { + file, err := os.Open(f.fullPath(path)) + if err != nil { + if os.IsNotExist(err) { + return nil, os.ErrNotExist + } + return nil, fmt.Errorf("error opening file %s: %w", path, err) + } + log.Debugf("File retrieved successfully: %s", path) + return file, nil +} + +// GetFileMetadata reads metadata from the sidecar JSON file. +func (f *FilesystemDataStore) GetFileMetadata(ctx context.Context, path string) (map[string]string, error) { + metaPath := f.metadataPath(path) + data, err := os.ReadFile(metaPath) + if err != nil { + if os.IsNotExist(err) { + // Check if the main file exists + if _, mainErr := os.Stat(f.fullPath(path)); os.IsNotExist(mainErr) { + return nil, os.ErrNotExist + } + // Main file exists but no metadata - return empty map + return map[string]string{}, nil + } + return nil, fmt.Errorf("error reading metadata file %s: %w", metaPath, err) + } + + var metadata map[string]string + if err := json.Unmarshal(data, &metadata); err != nil { + return nil, fmt.Errorf("error parsing metadata file %s: %w", metaPath, err) + } + return metadata, nil +} + +// PutFile writes a file to the filesystem with optional metadata sidecar. +func (f *FilesystemDataStore) PutFile(ctx context.Context, path string, in io.WriterTo, metaData map[string]string) error { + fullPath := f.fullPath(path) + + // Create parent directories + dir := filepath.Dir(fullPath) + if err := os.MkdirAll(dir, 0755); err != nil { + return fmt.Errorf("failed to create directory %s: %w", dir, err) + } + + // Write metadata sidecar first if enabled and metadata is provided. + // This ensures that if the data file exists, metadata is assumed to exist too. + if f.writeMetadata && len(metaData) > 0 { + if err := f.writeMetadataFile(path, metaData); err != nil { + return err + } + } + + // Write the data file + file, err := os.Create(fullPath) + if err != nil { + return fmt.Errorf("failed to create file %s: %w", path, err) + } + + if _, err := in.WriteTo(file); err != nil { + file.Close() + return fmt.Errorf("failed to write file %s: %w", path, err) + } + + if err := file.Close(); err != nil { + return fmt.Errorf("failed to close file %s: %w", path, err) + } + + log.Debugf("File uploaded successfully: %s", path) + return nil +} + +// PutFileIfNotExists writes a file only if it doesn't already exist. +func (f *FilesystemDataStore) PutFileIfNotExists(ctx context.Context, path string, in io.WriterTo, metaData map[string]string) (bool, error) { + fullPath := f.fullPath(path) + + // Create parent directories + dir := filepath.Dir(fullPath) + if err := os.MkdirAll(dir, 0755); err != nil { + return false, fmt.Errorf("failed to create directory %s: %w", dir, err) + } + + // Use O_CREATE|O_EXCL for atomic check-and-create + file, err := os.OpenFile(fullPath, os.O_WRONLY|os.O_CREATE|os.O_EXCL, 0644) + if err != nil { + if os.IsExist(err) { + log.Debugf("File already exists: %s", path) + return false, nil + } + return false, fmt.Errorf("failed to create file %s: %w", path, err) + } + + // Write content to the file + buf := &bytes.Buffer{} + if _, err := in.WriteTo(buf); err != nil { + file.Close() + os.Remove(fullPath) // Clean up on error + return false, fmt.Errorf("failed to write file %s: %w", path, err) + } + + if _, err := file.Write(buf.Bytes()); err != nil { + file.Close() + os.Remove(fullPath) // Clean up on error + return false, fmt.Errorf("failed to write file %s: %w", path, err) + } + + if err := file.Close(); err != nil { + return false, fmt.Errorf("failed to close file %s: %w", path, err) + } + + // Write metadata sidecar if enabled and metadata is provided + if f.writeMetadata && len(metaData) > 0 { + if err := f.writeMetadataFile(path, metaData); err != nil { + return true, err + } + } + + log.Debugf("File uploaded successfully: %s", path) + return true, nil +} + +// writeMetadataFile writes metadata to a sidecar JSON file. +func (f *FilesystemDataStore) writeMetadataFile(path string, metaData map[string]string) error { + metaPath := f.metadataPath(path) + + data, err := json.Marshal(metaData) + if err != nil { + return fmt.Errorf("failed to marshal metadata for %s: %w", path, err) + } + + if err := os.WriteFile(metaPath, data, 0644); err != nil { + return fmt.Errorf("failed to write metadata file %s: %w", metaPath, err) + } + return nil +} + +// ListFilePaths lists file paths matching the given options. +// Results are returned in lexicographical order (matching GCS/S3 behavior). +func (f *FilesystemDataStore) ListFilePaths(ctx context.Context, options ListFileOptions) ([]string, error) { + limit := options.Limit + if limit <= 0 || limit > listFilePathsMaxLimit { + limit = listFilePathsMaxLimit + } + + var files []string + err := filepath.WalkDir(f.basePath, func(path string, d os.DirEntry, err error) error { + if err != nil { + return err + } + + // Skip directories + if d.IsDir() { + return nil + } + + // Skip metadata sidecar files + if strings.HasSuffix(d.Name(), metadataSuffix) { + return nil + } + + // Get path relative to basePath and normalize to forward slashes + relPath, err := filepath.Rel(f.basePath, path) + if err != nil { + return err + } + relPath = filepath.ToSlash(relPath) + + // Apply prefix filter + if options.Prefix != "" && !strings.HasPrefix(relPath, options.Prefix) { + return nil + } + + // Apply StartAfter filter (WalkDir walks in lexical order) + if options.StartAfter != "" && relPath <= options.StartAfter { + return nil + } + + files = append(files, relPath) + + // Stop early if we've reached the limit + if uint32(len(files)) >= limit { + return filepath.SkipAll + } + + return nil + }) + if err != nil && err != filepath.SkipAll { + return nil, err + } + + return files, nil +} + +// Close is a no-op for FilesystemDataStore as it doesn't maintain persistent connections. +func (f *FilesystemDataStore) Close() error { + return nil +} diff --git a/support/datastore/filesystem_test.go b/support/datastore/filesystem_test.go new file mode 100644 index 0000000000..65f7edd67d --- /dev/null +++ b/support/datastore/filesystem_test.go @@ -0,0 +1,534 @@ +package datastore + +import ( + "bytes" + "context" + "fmt" + "os" + "path/filepath" + "testing" + + "github.com/stretchr/testify/require" +) + +func TestFilesystemExists(t *testing.T) { + dir := t.TempDir() + store, err := NewFilesystemDataStoreWithPath(dir, true) + require.NoError(t, err) + t.Cleanup(func() { + require.NoError(t, store.Close()) + }) + + // Create a test file + content := []byte("test content") + err = os.WriteFile(filepath.Join(dir, "file.txt"), content, 0644) + require.NoError(t, err) + + exists, err := store.Exists(context.Background(), "file.txt") + require.NoError(t, err) + require.True(t, exists) + + exists, err = store.Exists(context.Background(), "missing-file.txt") + require.NoError(t, err) + require.False(t, exists) +} + +func TestFilesystemSize(t *testing.T) { + dir := t.TempDir() + store, err := NewFilesystemDataStoreWithPath(dir, true) + require.NoError(t, err) + t.Cleanup(func() { + require.NoError(t, store.Close()) + }) + + content := []byte("inside the file") + err = os.WriteFile(filepath.Join(dir, "file.txt"), content, 0644) + require.NoError(t, err) + + size, err := store.Size(context.Background(), "file.txt") + require.NoError(t, err) + require.Equal(t, int64(len(content)), size) + + _, err = store.Size(context.Background(), "missing-file.txt") + require.ErrorIs(t, err, os.ErrNotExist) +} + +func TestFilesystemPutFile(t *testing.T) { + dir := t.TempDir() + store, err := NewFilesystemDataStoreWithPath(dir, true) + require.NoError(t, err) + t.Cleanup(func() { + require.NoError(t, store.Close()) + }) + + content := []byte("inside the file") + writerTo := bytes.NewReader(content) + err = store.PutFile(context.Background(), "file.txt", writerTo, nil) + require.NoError(t, err) + + reader, err := store.GetFile(context.Background(), "file.txt") + require.NoError(t, err) + requireReaderContentEquals(t, reader, content) + + metadata, err := store.GetFileMetadata(context.Background(), "file.txt") + require.NoError(t, err) + require.Equal(t, map[string]string{}, metadata) + + // Test overwriting + otherContent := []byte("other text") + writerTo = bytes.NewReader(otherContent) + err = store.PutFile(context.Background(), "file.txt", writerTo, nil) + require.NoError(t, err) + + reader, err = store.GetFile(context.Background(), "file.txt") + require.NoError(t, err) + requireReaderContentEquals(t, reader, otherContent) +} + +func TestFilesystemPutFileCreatesDirectories(t *testing.T) { + dir := t.TempDir() + store, err := NewFilesystemDataStoreWithPath(dir, true) + require.NoError(t, err) + t.Cleanup(func() { + require.NoError(t, store.Close()) + }) + + content := []byte("nested file content") + writerTo := bytes.NewReader(content) + err = store.PutFile(context.Background(), "a/b/c/file.txt", writerTo, nil) + require.NoError(t, err) + + reader, err := store.GetFile(context.Background(), "a/b/c/file.txt") + require.NoError(t, err) + requireReaderContentEquals(t, reader, content) +} + +func TestFilesystemPutFileIfNotExists(t *testing.T) { + dir := t.TempDir() + store, err := NewFilesystemDataStoreWithPath(dir, true) + require.NoError(t, err) + t.Cleanup(func() { + require.NoError(t, store.Close()) + }) + + existingContent := []byte("existing content") + err = os.WriteFile(filepath.Join(dir, "file.txt"), existingContent, 0644) + require.NoError(t, err) + + // Attempt to overwrite - should fail + newContent := []byte("new content") + writerTo := bytes.NewReader(newContent) + ok, err := store.PutFileIfNotExists(context.Background(), "file.txt", writerTo, nil) + require.NoError(t, err) + require.False(t, ok) + + // Verify content unchanged + reader, err := store.GetFile(context.Background(), "file.txt") + require.NoError(t, err) + requireReaderContentEquals(t, reader, existingContent) + + // Create new file - should succeed + writerTo = bytes.NewReader(newContent) + ok, err = store.PutFileIfNotExists(context.Background(), "other-file.txt", writerTo, nil) + require.NoError(t, err) + require.True(t, ok) + + reader, err = store.GetFile(context.Background(), "other-file.txt") + require.NoError(t, err) + requireReaderContentEquals(t, reader, newContent) +} + +func TestFilesystemGetFileLastModified(t *testing.T) { + dir := t.TempDir() + store, err := NewFilesystemDataStoreWithPath(dir, true) + require.NoError(t, err) + t.Cleanup(func() { + require.NoError(t, store.Close()) + }) + + content := []byte("inside the file") + writerTo := bytes.NewReader(content) + err = store.PutFile(context.Background(), "file.txt", writerTo, nil) + require.NoError(t, err) + + lastModified, err := store.GetFileLastModified(context.Background(), "file.txt") + require.NoError(t, err) + require.NotZero(t, lastModified) +} + +func TestFilesystemPutFileWithMetadata(t *testing.T) { + dir := t.TempDir() + store, err := NewFilesystemDataStoreWithPath(dir, true) + require.NoError(t, err) + t.Cleanup(func() { + require.NoError(t, store.Close()) + }) + + metadataObj := MetaData{ + StartLedger: 1234, + EndLedger: 1234, + StartLedgerCloseTime: 1234, + EndLedgerCloseTime: 1234, + NetworkPassPhrase: "testnet", + CompressionType: "zstd", + ProtocolVersion: 21, + CoreVersion: "v1.2.3", + Version: "1.0.0", + } + + content := []byte("inside the file") + writerTo := bytes.NewReader(content) + err = store.PutFile(context.Background(), "file.txt", writerTo, metadataObj.ToMap()) + require.NoError(t, err) + + metadata, err := store.GetFileMetadata(context.Background(), "file.txt") + require.NoError(t, err) + require.Equal(t, metadataObj.ToMap(), metadata) + + // Verify metadata file was created + _, err = os.Stat(filepath.Join(dir, "file.txt.metadata.json")) + require.NoError(t, err) + + // Update with new metadata + modifiedMetadataObj := MetaData{ + StartLedger: 5678, + EndLedger: 6789, + StartLedgerCloseTime: 1622547800, + EndLedgerCloseTime: 1622548900, + NetworkPassPhrase: "mainnet", + CompressionType: "gzip", + ProtocolVersion: 23, + CoreVersion: "v1.4.0", + Version: "2.0.0", + } + + otherContent := []byte("other text") + writerTo = bytes.NewReader(otherContent) + err = store.PutFile(context.Background(), "file.txt", writerTo, modifiedMetadataObj.ToMap()) + require.NoError(t, err) + + metadata, err = store.GetFileMetadata(context.Background(), "file.txt") + require.NoError(t, err) + require.Equal(t, modifiedMetadataObj.ToMap(), metadata) +} + +func TestFilesystemPutFileWithMetadataDisabled(t *testing.T) { + dir := t.TempDir() + store, err := NewFilesystemDataStoreWithPath(dir, false) // Metadata disabled + require.NoError(t, err) + t.Cleanup(func() { + require.NoError(t, store.Close()) + }) + + metadataObj := MetaData{ + StartLedger: 1234, + EndLedger: 1234, + NetworkPassPhrase: "testnet", + } + + content := []byte("inside the file") + writerTo := bytes.NewReader(content) + err = store.PutFile(context.Background(), "file.txt", writerTo, metadataObj.ToMap()) + require.NoError(t, err) + + // Metadata file should NOT be created + _, err = os.Stat(filepath.Join(dir, "file.txt.metadata.json")) + require.True(t, os.IsNotExist(err)) + + // GetFileMetadata should return empty map + metadata, err := store.GetFileMetadata(context.Background(), "file.txt") + require.NoError(t, err) + require.Equal(t, map[string]string{}, metadata) +} + +func TestFilesystemGetNonExistentFile(t *testing.T) { + dir := t.TempDir() + store, err := NewFilesystemDataStoreWithPath(dir, true) + require.NoError(t, err) + t.Cleanup(func() { + require.NoError(t, store.Close()) + }) + + // Create a different file + content := []byte("inside the file") + err = os.WriteFile(filepath.Join(dir, "file.txt"), content, 0644) + require.NoError(t, err) + + _, err = store.GetFile(context.Background(), "other-file.txt") + require.ErrorIs(t, err, os.ErrNotExist) + + metadata, err := store.GetFileMetadata(context.Background(), "other-file.txt") + require.ErrorIs(t, err, os.ErrNotExist) + require.Nil(t, metadata) +} + +func TestFilesystemListFilePaths(t *testing.T) { + dir := t.TempDir() + store, err := NewFilesystemDataStoreWithPath(dir, true) + require.NoError(t, err) + t.Cleanup(func() { _ = store.Close() }) + + // Create test files + for _, name := range []string{"a", "b", "c"} { + err := os.WriteFile(filepath.Join(dir, name), []byte("1"), 0644) + require.NoError(t, err) + } + + paths, err := store.ListFilePaths(context.Background(), ListFileOptions{Limit: 2}) + require.NoError(t, err) + require.Equal(t, []string{"a", "b"}, paths) +} + +func TestFilesystemListFilePaths_WithPrefix(t *testing.T) { + dir := t.TempDir() + store, err := NewFilesystemDataStoreWithPath(dir, true) + require.NoError(t, err) + t.Cleanup(func() { _ = store.Close() }) + + // Create directory structure + require.NoError(t, os.MkdirAll(filepath.Join(dir, "a"), 0755)) + require.NoError(t, os.MkdirAll(filepath.Join(dir, "b"), 0755)) + require.NoError(t, os.WriteFile(filepath.Join(dir, "a", "x"), []byte("1"), 0644)) + require.NoError(t, os.WriteFile(filepath.Join(dir, "a", "y"), []byte("1"), 0644)) + require.NoError(t, os.WriteFile(filepath.Join(dir, "b", "z"), []byte("1"), 0644)) + + paths, err := store.ListFilePaths(context.Background(), ListFileOptions{Prefix: "a", Limit: 10}) + require.NoError(t, err) + require.Equal(t, []string{"a/x", "a/y"}, paths) +} + +func TestFilesystemListFilePaths_ExcludesMetadataFiles(t *testing.T) { + dir := t.TempDir() + store, err := NewFilesystemDataStoreWithPath(dir, true) + require.NoError(t, err) + t.Cleanup(func() { _ = store.Close() }) + + // Create file with metadata + content := []byte("content") + metadata := map[string]string{"key": "value"} + err = store.PutFile(context.Background(), "file.txt", bytes.NewReader(content), metadata) + require.NoError(t, err) + + // Verify metadata file exists + _, err = os.Stat(filepath.Join(dir, "file.txt.metadata.json")) + require.NoError(t, err) + + // ListFilePaths should only return the main file, not the metadata file + paths, err := store.ListFilePaths(context.Background(), ListFileOptions{}) + require.NoError(t, err) + require.Equal(t, []string{"file.txt"}, paths) +} + +func TestFilesystemListFilePaths_LimitDefaultAndCap(t *testing.T) { + dir := t.TempDir() + store, err := NewFilesystemDataStoreWithPath(dir, true) + require.NoError(t, err) + t.Cleanup(func() { _ = store.Close() }) + + // Create 1200 files + for i := 0; i < 1200; i++ { + err := os.WriteFile(filepath.Join(dir, fmt.Sprintf("%04d", i)), []byte("1"), 0644) + require.NoError(t, err) + } + + // Default limit should cap at 1000 + paths, err := store.ListFilePaths(context.Background(), ListFileOptions{}) + require.NoError(t, err) + require.Equal(t, 1000, len(paths)) + + // Explicit limit over 1000 should also cap at 1000 + paths, err = store.ListFilePaths(context.Background(), ListFileOptions{Limit: 5000}) + require.NoError(t, err) + require.Equal(t, 1000, len(paths)) +} + +func TestFilesystemListFilePaths_StartAfter(t *testing.T) { + t.Run("basic start-after (no Prefix)", func(t *testing.T) { + dir := t.TempDir() + store, err := NewFilesystemDataStoreWithPath(dir, true) + require.NoError(t, err) + t.Cleanup(func() { _ = store.Close() }) + + for i := 0; i < 10; i++ { + err := os.WriteFile(filepath.Join(dir, fmt.Sprintf("%04d", i)), []byte("x"), 0644) + require.NoError(t, err) + } + + paths, err := store.ListFilePaths(context.Background(), ListFileOptions{ + StartAfter: "0005", + }) + require.NoError(t, err) + require.Equal(t, []string{"0006", "0007", "0008", "0009"}, paths) + }) + + t.Run("with Prefix directory and start-after inside it", func(t *testing.T) { + dir := t.TempDir() + store, err := NewFilesystemDataStoreWithPath(dir, true) + require.NoError(t, err) + t.Cleanup(func() { _ = store.Close() }) + + require.NoError(t, os.MkdirAll(filepath.Join(dir, "a"), 0755)) + require.NoError(t, os.MkdirAll(filepath.Join(dir, "b"), 0755)) + require.NoError(t, os.WriteFile(filepath.Join(dir, "a", "0001"), []byte(""), 0644)) + require.NoError(t, os.WriteFile(filepath.Join(dir, "a", "0002"), []byte(""), 0644)) + require.NoError(t, os.WriteFile(filepath.Join(dir, "b", "0002"), []byte(""), 0644)) + + paths, err := store.ListFilePaths(context.Background(), ListFileOptions{ + Prefix: "a/", + StartAfter: "a/0001", + }) + require.NoError(t, err) + require.Equal(t, []string{"a/0002"}, paths) + }) + + t.Run("start-after equals last key -> empty", func(t *testing.T) { + dir := t.TempDir() + store, err := NewFilesystemDataStoreWithPath(dir, true) + require.NoError(t, err) + t.Cleanup(func() { _ = store.Close() }) + + for _, name := range []string{"0000", "0001", "0002"} { + err := os.WriteFile(filepath.Join(dir, name), []byte(""), 0644) + require.NoError(t, err) + } + + paths, err := store.ListFilePaths(context.Background(), ListFileOptions{ + StartAfter: "0002", + }) + require.NoError(t, err) + require.Empty(t, paths) + }) + + t.Run("start-after before first key -> all returned", func(t *testing.T) { + dir := t.TempDir() + store, err := NewFilesystemDataStoreWithPath(dir, true) + require.NoError(t, err) + t.Cleanup(func() { _ = store.Close() }) + + for _, name := range []string{"0001", "0002", "0003"} { + err := os.WriteFile(filepath.Join(dir, name), []byte(""), 0644) + require.NoError(t, err) + } + + paths, err := store.ListFilePaths(context.Background(), ListFileOptions{ + StartAfter: "0000", + }) + require.NoError(t, err) + require.Equal(t, []string{"0001", "0002", "0003"}, paths) + }) + + t.Run("start-after missing-but-between keys -> next greater", func(t *testing.T) { + dir := t.TempDir() + store, err := NewFilesystemDataStoreWithPath(dir, true) + require.NoError(t, err) + t.Cleanup(func() { _ = store.Close() }) + + for _, name := range []string{"0002", "0004", "0006"} { + err := os.WriteFile(filepath.Join(dir, name), []byte(""), 0644) + require.NoError(t, err) + } + + paths, err := store.ListFilePaths(context.Background(), ListFileOptions{ + StartAfter: "0003", + }) + require.NoError(t, err) + require.Equal(t, []string{"0004", "0006"}, paths) + }) + + t.Run("respects limit together with start-after", func(t *testing.T) { + dir := t.TempDir() + store, err := NewFilesystemDataStoreWithPath(dir, true) + require.NoError(t, err) + t.Cleanup(func() { _ = store.Close() }) + + for i := 0; i < 10; i++ { + err := os.WriteFile(filepath.Join(dir, fmt.Sprintf("%04d", i)), []byte("x"), 0644) + require.NoError(t, err) + } + + paths, err := store.ListFilePaths(context.Background(), ListFileOptions{ + StartAfter: "0004", + Limit: 3, + }) + require.NoError(t, err) + require.Equal(t, []string{"0005", "0006", "0007"}, paths) + }) +} + +func TestNewFilesystemDataStore(t *testing.T) { + dir := t.TempDir() + + config := DataStoreConfig{ + Type: "Filesystem", + Params: map[string]string{ + "destination_path": dir, + }, + } + + store, err := NewDataStore(context.Background(), config) + require.NoError(t, err) + require.NotNil(t, store) + require.NoError(t, store.Close()) +} + +func TestNewFilesystemDataStore_MissingDestinationPath(t *testing.T) { + config := DataStoreConfig{ + Type: "Filesystem", + Params: map[string]string{}, + } + + _, err := NewDataStore(context.Background(), config) + require.Error(t, err) + require.Contains(t, err.Error(), "no destination_path") +} + +func TestNewFilesystemDataStore_WriteMetadataConfig(t *testing.T) { + t.Run("default is true", func(t *testing.T) { + dir := t.TempDir() + config := DataStoreConfig{ + Type: "Filesystem", + Params: map[string]string{ + "destination_path": dir, + }, + } + + store, err := NewDataStore(context.Background(), config) + require.NoError(t, err) + defer store.Close() + + // Write file with metadata + content := []byte("content") + metadata := map[string]string{"key": "value"} + err = store.PutFile(context.Background(), "file.txt", bytes.NewReader(content), metadata) + require.NoError(t, err) + + // Metadata file should exist + _, err = os.Stat(filepath.Join(dir, "file.txt.metadata.json")) + require.NoError(t, err) + }) + + t.Run("explicit false disables metadata", func(t *testing.T) { + dir := t.TempDir() + config := DataStoreConfig{ + Type: "Filesystem", + Params: map[string]string{ + "destination_path": dir, + "write_metadata": "false", + }, + } + + store, err := NewDataStore(context.Background(), config) + require.NoError(t, err) + defer store.Close() + + // Write file with metadata + content := []byte("content") + metadata := map[string]string{"key": "value"} + err = store.PutFile(context.Background(), "file.txt", bytes.NewReader(content), metadata) + require.NoError(t, err) + + // Metadata file should NOT exist + _, err = os.Stat(filepath.Join(dir, "file.txt.metadata.json")) + require.True(t, os.IsNotExist(err)) + }) +} From 3b335711cb1ec6d130f6a39ebce88688c17da129 Mon Sep 17 00:00:00 2001 From: Leigh <351529+leighmcculloch@users.noreply.github.com> Date: Wed, 7 Jan 2026 05:42:41 +1000 Subject: [PATCH 02/11] remove metadata support from filesystem datastore --- support/datastore/filesystem.go | 89 ++----------- support/datastore/filesystem_test.go | 190 +++------------------------ 2 files changed, 29 insertions(+), 250 deletions(-) diff --git a/support/datastore/filesystem.go b/support/datastore/filesystem.go index c02113d0dd..3c95be1445 100644 --- a/support/datastore/filesystem.go +++ b/support/datastore/filesystem.go @@ -3,7 +3,6 @@ package datastore import ( "bytes" "context" - "encoding/json" "errors" "fmt" "io" @@ -17,12 +16,12 @@ import ( var _ DataStore = &FilesystemDataStore{} -const metadataSuffix = ".metadata.json" - // FilesystemDataStore implements DataStore for local filesystem storage. +// Note: This implementation does not support storing metadata. The metaData +// parameter in PutFile and PutFileIfNotExists is ignored, and GetFileMetadata +// always returns an empty map. type FilesystemDataStore struct { - basePath string - writeMetadata bool + basePath string } // NewFilesystemDataStore creates a new FilesystemDataStore from configuration. @@ -32,17 +31,11 @@ func NewFilesystemDataStore(ctx context.Context, datastoreConfig DataStoreConfig return nil, errors.New("invalid Filesystem config, no destination_path") } - // write_metadata defaults to true - writeMetadata := true - if val, ok := datastoreConfig.Params["write_metadata"]; ok { - writeMetadata = val != "false" - } - - return NewFilesystemDataStoreWithPath(destinationPath, writeMetadata) + return NewFilesystemDataStoreWithPath(destinationPath) } // NewFilesystemDataStoreWithPath creates a FilesystemDataStore with the given base path. -func NewFilesystemDataStoreWithPath(basePath string, writeMetadata bool) (DataStore, error) { +func NewFilesystemDataStoreWithPath(basePath string) (DataStore, error) { // Ensure the base path exists if err := os.MkdirAll(basePath, 0755); err != nil { return nil, fmt.Errorf("failed to create base directory %s: %w", basePath, err) @@ -53,11 +46,10 @@ func NewFilesystemDataStoreWithPath(basePath string, writeMetadata bool) (DataSt return nil, fmt.Errorf("failed to resolve absolute path: %w", err) } - log.Debugf("Creating Filesystem datastore at: %s, writeMetadata: %v", absPath, writeMetadata) + log.Debugf("Creating Filesystem datastore at: %s", absPath) return &FilesystemDataStore{ - basePath: absPath, - writeMetadata: writeMetadata, + basePath: absPath, }, nil } @@ -66,11 +58,6 @@ func (f *FilesystemDataStore) fullPath(path string) string { return filepath.Join(f.basePath, path) } -// metadataPath returns the path to the metadata sidecar file. -func (f *FilesystemDataStore) metadataPath(path string) string { - return f.fullPath(path) + metadataSuffix -} - // Exists checks if a file exists in the filesystem. func (f *FilesystemDataStore) Exists(ctx context.Context, path string) (bool, error) { _, err := os.Stat(f.fullPath(path)) @@ -120,30 +107,15 @@ func (f *FilesystemDataStore) GetFile(ctx context.Context, path string) (io.Read return file, nil } -// GetFileMetadata reads metadata from the sidecar JSON file. +// GetFileMetadata returns an empty map as filesystem storage does not support metadata. func (f *FilesystemDataStore) GetFileMetadata(ctx context.Context, path string) (map[string]string, error) { - metaPath := f.metadataPath(path) - data, err := os.ReadFile(metaPath) - if err != nil { - if os.IsNotExist(err) { - // Check if the main file exists - if _, mainErr := os.Stat(f.fullPath(path)); os.IsNotExist(mainErr) { - return nil, os.ErrNotExist - } - // Main file exists but no metadata - return empty map - return map[string]string{}, nil - } - return nil, fmt.Errorf("error reading metadata file %s: %w", metaPath, err) - } - - var metadata map[string]string - if err := json.Unmarshal(data, &metadata); err != nil { - return nil, fmt.Errorf("error parsing metadata file %s: %w", metaPath, err) + if _, err := os.Stat(f.fullPath(path)); os.IsNotExist(err) { + return nil, os.ErrNotExist } - return metadata, nil + return map[string]string{}, nil } -// PutFile writes a file to the filesystem with optional metadata sidecar. +// PutFile writes a file to the filesystem. func (f *FilesystemDataStore) PutFile(ctx context.Context, path string, in io.WriterTo, metaData map[string]string) error { fullPath := f.fullPath(path) @@ -153,14 +125,6 @@ func (f *FilesystemDataStore) PutFile(ctx context.Context, path string, in io.Wr return fmt.Errorf("failed to create directory %s: %w", dir, err) } - // Write metadata sidecar first if enabled and metadata is provided. - // This ensures that if the data file exists, metadata is assumed to exist too. - if f.writeMetadata && len(metaData) > 0 { - if err := f.writeMetadataFile(path, metaData); err != nil { - return err - } - } - // Write the data file file, err := os.Create(fullPath) if err != nil { @@ -218,32 +182,10 @@ func (f *FilesystemDataStore) PutFileIfNotExists(ctx context.Context, path strin return false, fmt.Errorf("failed to close file %s: %w", path, err) } - // Write metadata sidecar if enabled and metadata is provided - if f.writeMetadata && len(metaData) > 0 { - if err := f.writeMetadataFile(path, metaData); err != nil { - return true, err - } - } - log.Debugf("File uploaded successfully: %s", path) return true, nil } -// writeMetadataFile writes metadata to a sidecar JSON file. -func (f *FilesystemDataStore) writeMetadataFile(path string, metaData map[string]string) error { - metaPath := f.metadataPath(path) - - data, err := json.Marshal(metaData) - if err != nil { - return fmt.Errorf("failed to marshal metadata for %s: %w", path, err) - } - - if err := os.WriteFile(metaPath, data, 0644); err != nil { - return fmt.Errorf("failed to write metadata file %s: %w", metaPath, err) - } - return nil -} - // ListFilePaths lists file paths matching the given options. // Results are returned in lexicographical order (matching GCS/S3 behavior). func (f *FilesystemDataStore) ListFilePaths(ctx context.Context, options ListFileOptions) ([]string, error) { @@ -263,11 +205,6 @@ func (f *FilesystemDataStore) ListFilePaths(ctx context.Context, options ListFil return nil } - // Skip metadata sidecar files - if strings.HasSuffix(d.Name(), metadataSuffix) { - return nil - } - // Get path relative to basePath and normalize to forward slashes relPath, err := filepath.Rel(f.basePath, path) if err != nil { diff --git a/support/datastore/filesystem_test.go b/support/datastore/filesystem_test.go index 65f7edd67d..783459ab82 100644 --- a/support/datastore/filesystem_test.go +++ b/support/datastore/filesystem_test.go @@ -13,7 +13,7 @@ import ( func TestFilesystemExists(t *testing.T) { dir := t.TempDir() - store, err := NewFilesystemDataStoreWithPath(dir, true) + store, err := NewFilesystemDataStoreWithPath(dir) require.NoError(t, err) t.Cleanup(func() { require.NoError(t, store.Close()) @@ -35,7 +35,7 @@ func TestFilesystemExists(t *testing.T) { func TestFilesystemSize(t *testing.T) { dir := t.TempDir() - store, err := NewFilesystemDataStoreWithPath(dir, true) + store, err := NewFilesystemDataStoreWithPath(dir) require.NoError(t, err) t.Cleanup(func() { require.NoError(t, store.Close()) @@ -55,7 +55,7 @@ func TestFilesystemSize(t *testing.T) { func TestFilesystemPutFile(t *testing.T) { dir := t.TempDir() - store, err := NewFilesystemDataStoreWithPath(dir, true) + store, err := NewFilesystemDataStoreWithPath(dir) require.NoError(t, err) t.Cleanup(func() { require.NoError(t, store.Close()) @@ -87,7 +87,7 @@ func TestFilesystemPutFile(t *testing.T) { func TestFilesystemPutFileCreatesDirectories(t *testing.T) { dir := t.TempDir() - store, err := NewFilesystemDataStoreWithPath(dir, true) + store, err := NewFilesystemDataStoreWithPath(dir) require.NoError(t, err) t.Cleanup(func() { require.NoError(t, store.Close()) @@ -105,7 +105,7 @@ func TestFilesystemPutFileCreatesDirectories(t *testing.T) { func TestFilesystemPutFileIfNotExists(t *testing.T) { dir := t.TempDir() - store, err := NewFilesystemDataStoreWithPath(dir, true) + store, err := NewFilesystemDataStoreWithPath(dir) require.NoError(t, err) t.Cleanup(func() { require.NoError(t, store.Close()) @@ -140,7 +140,7 @@ func TestFilesystemPutFileIfNotExists(t *testing.T) { func TestFilesystemGetFileLastModified(t *testing.T) { dir := t.TempDir() - store, err := NewFilesystemDataStoreWithPath(dir, true) + store, err := NewFilesystemDataStoreWithPath(dir) require.NoError(t, err) t.Cleanup(func() { require.NoError(t, store.Close()) @@ -156,94 +156,9 @@ func TestFilesystemGetFileLastModified(t *testing.T) { require.NotZero(t, lastModified) } -func TestFilesystemPutFileWithMetadata(t *testing.T) { - dir := t.TempDir() - store, err := NewFilesystemDataStoreWithPath(dir, true) - require.NoError(t, err) - t.Cleanup(func() { - require.NoError(t, store.Close()) - }) - - metadataObj := MetaData{ - StartLedger: 1234, - EndLedger: 1234, - StartLedgerCloseTime: 1234, - EndLedgerCloseTime: 1234, - NetworkPassPhrase: "testnet", - CompressionType: "zstd", - ProtocolVersion: 21, - CoreVersion: "v1.2.3", - Version: "1.0.0", - } - - content := []byte("inside the file") - writerTo := bytes.NewReader(content) - err = store.PutFile(context.Background(), "file.txt", writerTo, metadataObj.ToMap()) - require.NoError(t, err) - - metadata, err := store.GetFileMetadata(context.Background(), "file.txt") - require.NoError(t, err) - require.Equal(t, metadataObj.ToMap(), metadata) - - // Verify metadata file was created - _, err = os.Stat(filepath.Join(dir, "file.txt.metadata.json")) - require.NoError(t, err) - - // Update with new metadata - modifiedMetadataObj := MetaData{ - StartLedger: 5678, - EndLedger: 6789, - StartLedgerCloseTime: 1622547800, - EndLedgerCloseTime: 1622548900, - NetworkPassPhrase: "mainnet", - CompressionType: "gzip", - ProtocolVersion: 23, - CoreVersion: "v1.4.0", - Version: "2.0.0", - } - - otherContent := []byte("other text") - writerTo = bytes.NewReader(otherContent) - err = store.PutFile(context.Background(), "file.txt", writerTo, modifiedMetadataObj.ToMap()) - require.NoError(t, err) - - metadata, err = store.GetFileMetadata(context.Background(), "file.txt") - require.NoError(t, err) - require.Equal(t, modifiedMetadataObj.ToMap(), metadata) -} - -func TestFilesystemPutFileWithMetadataDisabled(t *testing.T) { - dir := t.TempDir() - store, err := NewFilesystemDataStoreWithPath(dir, false) // Metadata disabled - require.NoError(t, err) - t.Cleanup(func() { - require.NoError(t, store.Close()) - }) - - metadataObj := MetaData{ - StartLedger: 1234, - EndLedger: 1234, - NetworkPassPhrase: "testnet", - } - - content := []byte("inside the file") - writerTo := bytes.NewReader(content) - err = store.PutFile(context.Background(), "file.txt", writerTo, metadataObj.ToMap()) - require.NoError(t, err) - - // Metadata file should NOT be created - _, err = os.Stat(filepath.Join(dir, "file.txt.metadata.json")) - require.True(t, os.IsNotExist(err)) - - // GetFileMetadata should return empty map - metadata, err := store.GetFileMetadata(context.Background(), "file.txt") - require.NoError(t, err) - require.Equal(t, map[string]string{}, metadata) -} - func TestFilesystemGetNonExistentFile(t *testing.T) { dir := t.TempDir() - store, err := NewFilesystemDataStoreWithPath(dir, true) + store, err := NewFilesystemDataStoreWithPath(dir) require.NoError(t, err) t.Cleanup(func() { require.NoError(t, store.Close()) @@ -264,7 +179,7 @@ func TestFilesystemGetNonExistentFile(t *testing.T) { func TestFilesystemListFilePaths(t *testing.T) { dir := t.TempDir() - store, err := NewFilesystemDataStoreWithPath(dir, true) + store, err := NewFilesystemDataStoreWithPath(dir) require.NoError(t, err) t.Cleanup(func() { _ = store.Close() }) @@ -281,7 +196,7 @@ func TestFilesystemListFilePaths(t *testing.T) { func TestFilesystemListFilePaths_WithPrefix(t *testing.T) { dir := t.TempDir() - store, err := NewFilesystemDataStoreWithPath(dir, true) + store, err := NewFilesystemDataStoreWithPath(dir) require.NoError(t, err) t.Cleanup(func() { _ = store.Close() }) @@ -297,31 +212,9 @@ func TestFilesystemListFilePaths_WithPrefix(t *testing.T) { require.Equal(t, []string{"a/x", "a/y"}, paths) } -func TestFilesystemListFilePaths_ExcludesMetadataFiles(t *testing.T) { - dir := t.TempDir() - store, err := NewFilesystemDataStoreWithPath(dir, true) - require.NoError(t, err) - t.Cleanup(func() { _ = store.Close() }) - - // Create file with metadata - content := []byte("content") - metadata := map[string]string{"key": "value"} - err = store.PutFile(context.Background(), "file.txt", bytes.NewReader(content), metadata) - require.NoError(t, err) - - // Verify metadata file exists - _, err = os.Stat(filepath.Join(dir, "file.txt.metadata.json")) - require.NoError(t, err) - - // ListFilePaths should only return the main file, not the metadata file - paths, err := store.ListFilePaths(context.Background(), ListFileOptions{}) - require.NoError(t, err) - require.Equal(t, []string{"file.txt"}, paths) -} - func TestFilesystemListFilePaths_LimitDefaultAndCap(t *testing.T) { dir := t.TempDir() - store, err := NewFilesystemDataStoreWithPath(dir, true) + store, err := NewFilesystemDataStoreWithPath(dir) require.NoError(t, err) t.Cleanup(func() { _ = store.Close() }) @@ -345,7 +238,7 @@ func TestFilesystemListFilePaths_LimitDefaultAndCap(t *testing.T) { func TestFilesystemListFilePaths_StartAfter(t *testing.T) { t.Run("basic start-after (no Prefix)", func(t *testing.T) { dir := t.TempDir() - store, err := NewFilesystemDataStoreWithPath(dir, true) + store, err := NewFilesystemDataStoreWithPath(dir) require.NoError(t, err) t.Cleanup(func() { _ = store.Close() }) @@ -363,7 +256,7 @@ func TestFilesystemListFilePaths_StartAfter(t *testing.T) { t.Run("with Prefix directory and start-after inside it", func(t *testing.T) { dir := t.TempDir() - store, err := NewFilesystemDataStoreWithPath(dir, true) + store, err := NewFilesystemDataStoreWithPath(dir) require.NoError(t, err) t.Cleanup(func() { _ = store.Close() }) @@ -383,7 +276,7 @@ func TestFilesystemListFilePaths_StartAfter(t *testing.T) { t.Run("start-after equals last key -> empty", func(t *testing.T) { dir := t.TempDir() - store, err := NewFilesystemDataStoreWithPath(dir, true) + store, err := NewFilesystemDataStoreWithPath(dir) require.NoError(t, err) t.Cleanup(func() { _ = store.Close() }) @@ -401,7 +294,7 @@ func TestFilesystemListFilePaths_StartAfter(t *testing.T) { t.Run("start-after before first key -> all returned", func(t *testing.T) { dir := t.TempDir() - store, err := NewFilesystemDataStoreWithPath(dir, true) + store, err := NewFilesystemDataStoreWithPath(dir) require.NoError(t, err) t.Cleanup(func() { _ = store.Close() }) @@ -419,7 +312,7 @@ func TestFilesystemListFilePaths_StartAfter(t *testing.T) { t.Run("start-after missing-but-between keys -> next greater", func(t *testing.T) { dir := t.TempDir() - store, err := NewFilesystemDataStoreWithPath(dir, true) + store, err := NewFilesystemDataStoreWithPath(dir) require.NoError(t, err) t.Cleanup(func() { _ = store.Close() }) @@ -437,7 +330,7 @@ func TestFilesystemListFilePaths_StartAfter(t *testing.T) { t.Run("respects limit together with start-after", func(t *testing.T) { dir := t.TempDir() - store, err := NewFilesystemDataStoreWithPath(dir, true) + store, err := NewFilesystemDataStoreWithPath(dir) require.NoError(t, err) t.Cleanup(func() { _ = store.Close() }) @@ -481,54 +374,3 @@ func TestNewFilesystemDataStore_MissingDestinationPath(t *testing.T) { require.Error(t, err) require.Contains(t, err.Error(), "no destination_path") } - -func TestNewFilesystemDataStore_WriteMetadataConfig(t *testing.T) { - t.Run("default is true", func(t *testing.T) { - dir := t.TempDir() - config := DataStoreConfig{ - Type: "Filesystem", - Params: map[string]string{ - "destination_path": dir, - }, - } - - store, err := NewDataStore(context.Background(), config) - require.NoError(t, err) - defer store.Close() - - // Write file with metadata - content := []byte("content") - metadata := map[string]string{"key": "value"} - err = store.PutFile(context.Background(), "file.txt", bytes.NewReader(content), metadata) - require.NoError(t, err) - - // Metadata file should exist - _, err = os.Stat(filepath.Join(dir, "file.txt.metadata.json")) - require.NoError(t, err) - }) - - t.Run("explicit false disables metadata", func(t *testing.T) { - dir := t.TempDir() - config := DataStoreConfig{ - Type: "Filesystem", - Params: map[string]string{ - "destination_path": dir, - "write_metadata": "false", - }, - } - - store, err := NewDataStore(context.Background(), config) - require.NoError(t, err) - defer store.Close() - - // Write file with metadata - content := []byte("content") - metadata := map[string]string{"key": "value"} - err = store.PutFile(context.Background(), "file.txt", bytes.NewReader(content), metadata) - require.NoError(t, err) - - // Metadata file should NOT exist - _, err = os.Stat(filepath.Join(dir, "file.txt.metadata.json")) - require.True(t, os.IsNotExist(err)) - }) -} From 0218f0a1a277476181e56aa38f6fd193f63f86b5 Mon Sep 17 00:00:00 2001 From: Leigh <351529+leighmcculloch@users.noreply.github.com> Date: Wed, 7 Jan 2026 05:45:13 +1000 Subject: [PATCH 03/11] expand filesystem datastore documentation --- support/datastore/filesystem.go | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) diff --git a/support/datastore/filesystem.go b/support/datastore/filesystem.go index 3c95be1445..9d111a029f 100644 --- a/support/datastore/filesystem.go +++ b/support/datastore/filesystem.go @@ -17,7 +17,11 @@ import ( var _ DataStore = &FilesystemDataStore{} // FilesystemDataStore implements DataStore for local filesystem storage. -// Note: This implementation does not support storing metadata. The metaData +// +// Note: This implementation is not recommended for production use. It is +// intended for development and testing purposes only. +// +// This implementation does not support storing metadata. The metaData // parameter in PutFile and PutFileIfNotExists is ignored, and GetFileMetadata // always returns an empty map. type FilesystemDataStore struct { From a2c97967f601b6f9e93f9770d86dfaa6c632ce13 Mon Sep 17 00:00:00 2001 From: Leigh <351529+leighmcculloch@users.noreply.github.com> Date: Wed, 7 Jan 2026 06:29:54 +1000 Subject: [PATCH 04/11] Fix variable shadowing warnings in filesystem datastore tests --- support/datastore/filesystem_test.go | 14 +++++++------- 1 file changed, 7 insertions(+), 7 deletions(-) diff --git a/support/datastore/filesystem_test.go b/support/datastore/filesystem_test.go index 783459ab82..dc524f4595 100644 --- a/support/datastore/filesystem_test.go +++ b/support/datastore/filesystem_test.go @@ -185,7 +185,7 @@ func TestFilesystemListFilePaths(t *testing.T) { // Create test files for _, name := range []string{"a", "b", "c"} { - err := os.WriteFile(filepath.Join(dir, name), []byte("1"), 0644) + err = os.WriteFile(filepath.Join(dir, name), []byte("1"), 0644) require.NoError(t, err) } @@ -220,7 +220,7 @@ func TestFilesystemListFilePaths_LimitDefaultAndCap(t *testing.T) { // Create 1200 files for i := 0; i < 1200; i++ { - err := os.WriteFile(filepath.Join(dir, fmt.Sprintf("%04d", i)), []byte("1"), 0644) + err = os.WriteFile(filepath.Join(dir, fmt.Sprintf("%04d", i)), []byte("1"), 0644) require.NoError(t, err) } @@ -243,7 +243,7 @@ func TestFilesystemListFilePaths_StartAfter(t *testing.T) { t.Cleanup(func() { _ = store.Close() }) for i := 0; i < 10; i++ { - err := os.WriteFile(filepath.Join(dir, fmt.Sprintf("%04d", i)), []byte("x"), 0644) + err = os.WriteFile(filepath.Join(dir, fmt.Sprintf("%04d", i)), []byte("x"), 0644) require.NoError(t, err) } @@ -281,7 +281,7 @@ func TestFilesystemListFilePaths_StartAfter(t *testing.T) { t.Cleanup(func() { _ = store.Close() }) for _, name := range []string{"0000", "0001", "0002"} { - err := os.WriteFile(filepath.Join(dir, name), []byte(""), 0644) + err = os.WriteFile(filepath.Join(dir, name), []byte(""), 0644) require.NoError(t, err) } @@ -299,7 +299,7 @@ func TestFilesystemListFilePaths_StartAfter(t *testing.T) { t.Cleanup(func() { _ = store.Close() }) for _, name := range []string{"0001", "0002", "0003"} { - err := os.WriteFile(filepath.Join(dir, name), []byte(""), 0644) + err = os.WriteFile(filepath.Join(dir, name), []byte(""), 0644) require.NoError(t, err) } @@ -317,7 +317,7 @@ func TestFilesystemListFilePaths_StartAfter(t *testing.T) { t.Cleanup(func() { _ = store.Close() }) for _, name := range []string{"0002", "0004", "0006"} { - err := os.WriteFile(filepath.Join(dir, name), []byte(""), 0644) + err = os.WriteFile(filepath.Join(dir, name), []byte(""), 0644) require.NoError(t, err) } @@ -335,7 +335,7 @@ func TestFilesystemListFilePaths_StartAfter(t *testing.T) { t.Cleanup(func() { _ = store.Close() }) for i := 0; i < 10; i++ { - err := os.WriteFile(filepath.Join(dir, fmt.Sprintf("%04d", i)), []byte("x"), 0644) + err = os.WriteFile(filepath.Join(dir, fmt.Sprintf("%04d", i)), []byte("x"), 0644) require.NoError(t, err) } From 84611788b2a6ac23ee3a715b3859f85863adb675 Mon Sep 17 00:00:00 2001 From: Leigh <351529+leighmcculloch@users.noreply.github.com> Date: Wed, 7 Jan 2026 06:55:51 +1000 Subject: [PATCH 05/11] Document concurrent write limitations in filesystem datastore --- support/datastore/filesystem.go | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/support/datastore/filesystem.go b/support/datastore/filesystem.go index 9d111a029f..49a7525aeb 100644 --- a/support/datastore/filesystem.go +++ b/support/datastore/filesystem.go @@ -24,6 +24,10 @@ var _ DataStore = &FilesystemDataStore{} // This implementation does not support storing metadata. The metaData // parameter in PutFile and PutFileIfNotExists is ignored, and GetFileMetadata // always returns an empty map. +// +// Concurrent writes to the same file path are not safe and may result in +// data corruption. Callers must ensure proper synchronization when writing +// to the same path from multiple goroutines. type FilesystemDataStore struct { basePath string } From fff78f0b8b3872b0c50888f272795d00c6b299b6 Mon Sep 17 00:00:00 2001 From: Leigh <351529+leighmcculloch@users.noreply.github.com> Date: Wed, 7 Jan 2026 06:59:15 +1000 Subject: [PATCH 06/11] Add context cancellation support in ListFilePaths --- support/datastore/filesystem.go | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/support/datastore/filesystem.go b/support/datastore/filesystem.go index 49a7525aeb..4a76b502bb 100644 --- a/support/datastore/filesystem.go +++ b/support/datastore/filesystem.go @@ -208,6 +208,11 @@ func (f *FilesystemDataStore) ListFilePaths(ctx context.Context, options ListFil return err } + // Check for context cancellation + if ctx.Err() != nil { + return ctx.Err() + } + // Skip directories if d.IsDir() { return nil From 2fc44f35a2c3932f4da4487345ca3c9ec40eaada Mon Sep 17 00:00:00 2001 From: Leigh <351529+leighmcculloch@users.noreply.github.com> Date: Wed, 7 Jan 2026 07:20:49 +1000 Subject: [PATCH 07/11] extract filesystem permission constants --- support/datastore/filesystem.go | 19 ++- support/datastore/filesystem_test.go | 194 +++++++++++++-------------- 2 files changed, 109 insertions(+), 104 deletions(-) diff --git a/support/datastore/filesystem.go b/support/datastore/filesystem.go index 4a76b502bb..e42ac9bc14 100644 --- a/support/datastore/filesystem.go +++ b/support/datastore/filesystem.go @@ -14,6 +14,11 @@ import ( "github.com/stellar/go-stellar-sdk/support/log" ) +const ( + defaultDirPerms os.FileMode = 0755 + defaultFilePerms os.FileMode = 0644 +) + var _ DataStore = &FilesystemDataStore{} // FilesystemDataStore implements DataStore for local filesystem storage. @@ -45,7 +50,7 @@ func NewFilesystemDataStore(ctx context.Context, datastoreConfig DataStoreConfig // NewFilesystemDataStoreWithPath creates a FilesystemDataStore with the given base path. func NewFilesystemDataStoreWithPath(basePath string) (DataStore, error) { // Ensure the base path exists - if err := os.MkdirAll(basePath, 0755); err != nil { + if err := os.MkdirAll(basePath, defaultDirPerms); err != nil { return nil, fmt.Errorf("failed to create base directory %s: %w", basePath, err) } @@ -129,7 +134,7 @@ func (f *FilesystemDataStore) PutFile(ctx context.Context, path string, in io.Wr // Create parent directories dir := filepath.Dir(fullPath) - if err := os.MkdirAll(dir, 0755); err != nil { + if err := os.MkdirAll(dir, defaultDirPerms); err != nil { return fmt.Errorf("failed to create directory %s: %w", dir, err) } @@ -153,17 +158,19 @@ func (f *FilesystemDataStore) PutFile(ctx context.Context, path string, in io.Wr } // PutFileIfNotExists writes a file only if it doesn't already exist. -func (f *FilesystemDataStore) PutFileIfNotExists(ctx context.Context, path string, in io.WriterTo, metaData map[string]string) (bool, error) { +func (f *FilesystemDataStore) PutFileIfNotExists( + ctx context.Context, path string, in io.WriterTo, metaData map[string]string, +) (bool, error) { fullPath := f.fullPath(path) // Create parent directories dir := filepath.Dir(fullPath) - if err := os.MkdirAll(dir, 0755); err != nil { + if err := os.MkdirAll(dir, defaultDirPerms); err != nil { return false, fmt.Errorf("failed to create directory %s: %w", dir, err) } // Use O_CREATE|O_EXCL for atomic check-and-create - file, err := os.OpenFile(fullPath, os.O_WRONLY|os.O_CREATE|os.O_EXCL, 0644) + file, err := os.OpenFile(fullPath, os.O_WRONLY|os.O_CREATE|os.O_EXCL, defaultFilePerms) if err != nil { if os.IsExist(err) { log.Debugf("File already exists: %s", path) @@ -238,7 +245,7 @@ func (f *FilesystemDataStore) ListFilePaths(ctx context.Context, options ListFil files = append(files, relPath) // Stop early if we've reached the limit - if uint32(len(files)) >= limit { + if len(files) >= int(limit) { return filepath.SkipAll } diff --git a/support/datastore/filesystem_test.go b/support/datastore/filesystem_test.go index dc524f4595..08926a2ba2 100644 --- a/support/datastore/filesystem_test.go +++ b/support/datastore/filesystem_test.go @@ -21,7 +21,7 @@ func TestFilesystemExists(t *testing.T) { // Create a test file content := []byte("test content") - err = os.WriteFile(filepath.Join(dir, "file.txt"), content, 0644) + err = os.WriteFile(filepath.Join(dir, "file.txt"), content, 0600) require.NoError(t, err) exists, err := store.Exists(context.Background(), "file.txt") @@ -42,7 +42,7 @@ func TestFilesystemSize(t *testing.T) { }) content := []byte("inside the file") - err = os.WriteFile(filepath.Join(dir, "file.txt"), content, 0644) + err = os.WriteFile(filepath.Join(dir, "file.txt"), content, 0600) require.NoError(t, err) size, err := store.Size(context.Background(), "file.txt") @@ -112,7 +112,7 @@ func TestFilesystemPutFileIfNotExists(t *testing.T) { }) existingContent := []byte("existing content") - err = os.WriteFile(filepath.Join(dir, "file.txt"), existingContent, 0644) + err = os.WriteFile(filepath.Join(dir, "file.txt"), existingContent, 0600) require.NoError(t, err) // Attempt to overwrite - should fail @@ -166,7 +166,7 @@ func TestFilesystemGetNonExistentFile(t *testing.T) { // Create a different file content := []byte("inside the file") - err = os.WriteFile(filepath.Join(dir, "file.txt"), content, 0644) + err = os.WriteFile(filepath.Join(dir, "file.txt"), content, 0600) require.NoError(t, err) _, err = store.GetFile(context.Background(), "other-file.txt") @@ -185,7 +185,7 @@ func TestFilesystemListFilePaths(t *testing.T) { // Create test files for _, name := range []string{"a", "b", "c"} { - err = os.WriteFile(filepath.Join(dir, name), []byte("1"), 0644) + err = os.WriteFile(filepath.Join(dir, name), []byte("1"), 0600) require.NoError(t, err) } @@ -203,9 +203,9 @@ func TestFilesystemListFilePaths_WithPrefix(t *testing.T) { // Create directory structure require.NoError(t, os.MkdirAll(filepath.Join(dir, "a"), 0755)) require.NoError(t, os.MkdirAll(filepath.Join(dir, "b"), 0755)) - require.NoError(t, os.WriteFile(filepath.Join(dir, "a", "x"), []byte("1"), 0644)) - require.NoError(t, os.WriteFile(filepath.Join(dir, "a", "y"), []byte("1"), 0644)) - require.NoError(t, os.WriteFile(filepath.Join(dir, "b", "z"), []byte("1"), 0644)) + require.NoError(t, os.WriteFile(filepath.Join(dir, "a", "x"), []byte("1"), 0600)) + require.NoError(t, os.WriteFile(filepath.Join(dir, "a", "y"), []byte("1"), 0600)) + require.NoError(t, os.WriteFile(filepath.Join(dir, "b", "z"), []byte("1"), 0600)) paths, err := store.ListFilePaths(context.Background(), ListFileOptions{Prefix: "a", Limit: 10}) require.NoError(t, err) @@ -220,7 +220,7 @@ func TestFilesystemListFilePaths_LimitDefaultAndCap(t *testing.T) { // Create 1200 files for i := 0; i < 1200; i++ { - err = os.WriteFile(filepath.Join(dir, fmt.Sprintf("%04d", i)), []byte("1"), 0644) + err = os.WriteFile(filepath.Join(dir, fmt.Sprintf("%04d", i)), []byte("1"), 0600) require.NoError(t, err) } @@ -235,117 +235,115 @@ func TestFilesystemListFilePaths_LimitDefaultAndCap(t *testing.T) { require.Equal(t, 1000, len(paths)) } -func TestFilesystemListFilePaths_StartAfter(t *testing.T) { - t.Run("basic start-after (no Prefix)", func(t *testing.T) { - dir := t.TempDir() - store, err := NewFilesystemDataStoreWithPath(dir) - require.NoError(t, err) - t.Cleanup(func() { _ = store.Close() }) - - for i := 0; i < 10; i++ { - err = os.WriteFile(filepath.Join(dir, fmt.Sprintf("%04d", i)), []byte("x"), 0644) - require.NoError(t, err) - } +func TestFilesystemListFilePaths_StartAfter_Basic(t *testing.T) { + dir := t.TempDir() + store, err := NewFilesystemDataStoreWithPath(dir) + require.NoError(t, err) + t.Cleanup(func() { _ = store.Close() }) - paths, err := store.ListFilePaths(context.Background(), ListFileOptions{ - StartAfter: "0005", - }) + for i := 0; i < 10; i++ { + err = os.WriteFile(filepath.Join(dir, fmt.Sprintf("%04d", i)), []byte("x"), 0600) require.NoError(t, err) - require.Equal(t, []string{"0006", "0007", "0008", "0009"}, paths) - }) + } - t.Run("with Prefix directory and start-after inside it", func(t *testing.T) { - dir := t.TempDir() - store, err := NewFilesystemDataStoreWithPath(dir) - require.NoError(t, err) - t.Cleanup(func() { _ = store.Close() }) - - require.NoError(t, os.MkdirAll(filepath.Join(dir, "a"), 0755)) - require.NoError(t, os.MkdirAll(filepath.Join(dir, "b"), 0755)) - require.NoError(t, os.WriteFile(filepath.Join(dir, "a", "0001"), []byte(""), 0644)) - require.NoError(t, os.WriteFile(filepath.Join(dir, "a", "0002"), []byte(""), 0644)) - require.NoError(t, os.WriteFile(filepath.Join(dir, "b", "0002"), []byte(""), 0644)) - - paths, err := store.ListFilePaths(context.Background(), ListFileOptions{ - Prefix: "a/", - StartAfter: "a/0001", - }) - require.NoError(t, err) - require.Equal(t, []string{"a/0002"}, paths) + paths, err := store.ListFilePaths(context.Background(), ListFileOptions{ + StartAfter: "0005", }) + require.NoError(t, err) + require.Equal(t, []string{"0006", "0007", "0008", "0009"}, paths) +} - t.Run("start-after equals last key -> empty", func(t *testing.T) { - dir := t.TempDir() - store, err := NewFilesystemDataStoreWithPath(dir) - require.NoError(t, err) - t.Cleanup(func() { _ = store.Close() }) +func TestFilesystemListFilePaths_StartAfter_WithPrefix(t *testing.T) { + dir := t.TempDir() + store, err := NewFilesystemDataStoreWithPath(dir) + require.NoError(t, err) + t.Cleanup(func() { _ = store.Close() }) - for _, name := range []string{"0000", "0001", "0002"} { - err = os.WriteFile(filepath.Join(dir, name), []byte(""), 0644) - require.NoError(t, err) - } + require.NoError(t, os.MkdirAll(filepath.Join(dir, "a"), 0755)) + require.NoError(t, os.MkdirAll(filepath.Join(dir, "b"), 0755)) + require.NoError(t, os.WriteFile(filepath.Join(dir, "a", "0001"), []byte(""), 0600)) + require.NoError(t, os.WriteFile(filepath.Join(dir, "a", "0002"), []byte(""), 0600)) + require.NoError(t, os.WriteFile(filepath.Join(dir, "b", "0002"), []byte(""), 0600)) - paths, err := store.ListFilePaths(context.Background(), ListFileOptions{ - StartAfter: "0002", - }) - require.NoError(t, err) - require.Empty(t, paths) + paths, err := store.ListFilePaths(context.Background(), ListFileOptions{ + Prefix: "a/", + StartAfter: "a/0001", }) + require.NoError(t, err) + require.Equal(t, []string{"a/0002"}, paths) +} - t.Run("start-after before first key -> all returned", func(t *testing.T) { - dir := t.TempDir() - store, err := NewFilesystemDataStoreWithPath(dir) - require.NoError(t, err) - t.Cleanup(func() { _ = store.Close() }) - - for _, name := range []string{"0001", "0002", "0003"} { - err = os.WriteFile(filepath.Join(dir, name), []byte(""), 0644) - require.NoError(t, err) - } +func TestFilesystemListFilePaths_StartAfter_EqualsLastKey(t *testing.T) { + dir := t.TempDir() + store, err := NewFilesystemDataStoreWithPath(dir) + require.NoError(t, err) + t.Cleanup(func() { _ = store.Close() }) - paths, err := store.ListFilePaths(context.Background(), ListFileOptions{ - StartAfter: "0000", - }) + for _, name := range []string{"0000", "0001", "0002"} { + err = os.WriteFile(filepath.Join(dir, name), []byte(""), 0600) require.NoError(t, err) - require.Equal(t, []string{"0001", "0002", "0003"}, paths) - }) + } - t.Run("start-after missing-but-between keys -> next greater", func(t *testing.T) { - dir := t.TempDir() - store, err := NewFilesystemDataStoreWithPath(dir) - require.NoError(t, err) - t.Cleanup(func() { _ = store.Close() }) + paths, err := store.ListFilePaths(context.Background(), ListFileOptions{ + StartAfter: "0002", + }) + require.NoError(t, err) + require.Empty(t, paths) +} - for _, name := range []string{"0002", "0004", "0006"} { - err = os.WriteFile(filepath.Join(dir, name), []byte(""), 0644) - require.NoError(t, err) - } +func TestFilesystemListFilePaths_StartAfter_BeforeFirstKey(t *testing.T) { + dir := t.TempDir() + store, err := NewFilesystemDataStoreWithPath(dir) + require.NoError(t, err) + t.Cleanup(func() { _ = store.Close() }) - paths, err := store.ListFilePaths(context.Background(), ListFileOptions{ - StartAfter: "0003", - }) + for _, name := range []string{"0001", "0002", "0003"} { + err = os.WriteFile(filepath.Join(dir, name), []byte(""), 0600) require.NoError(t, err) - require.Equal(t, []string{"0004", "0006"}, paths) + } + + paths, err := store.ListFilePaths(context.Background(), ListFileOptions{ + StartAfter: "0000", }) + require.NoError(t, err) + require.Equal(t, []string{"0001", "0002", "0003"}, paths) +} - t.Run("respects limit together with start-after", func(t *testing.T) { - dir := t.TempDir() - store, err := NewFilesystemDataStoreWithPath(dir) +func TestFilesystemListFilePaths_StartAfter_BetweenKeys(t *testing.T) { + dir := t.TempDir() + store, err := NewFilesystemDataStoreWithPath(dir) + require.NoError(t, err) + t.Cleanup(func() { _ = store.Close() }) + + for _, name := range []string{"0002", "0004", "0006"} { + err = os.WriteFile(filepath.Join(dir, name), []byte(""), 0600) require.NoError(t, err) - t.Cleanup(func() { _ = store.Close() }) + } + + paths, err := store.ListFilePaths(context.Background(), ListFileOptions{ + StartAfter: "0003", + }) + require.NoError(t, err) + require.Equal(t, []string{"0004", "0006"}, paths) +} - for i := 0; i < 10; i++ { - err = os.WriteFile(filepath.Join(dir, fmt.Sprintf("%04d", i)), []byte("x"), 0644) - require.NoError(t, err) - } +func TestFilesystemListFilePaths_StartAfter_WithLimit(t *testing.T) { + dir := t.TempDir() + store, err := NewFilesystemDataStoreWithPath(dir) + require.NoError(t, err) + t.Cleanup(func() { _ = store.Close() }) - paths, err := store.ListFilePaths(context.Background(), ListFileOptions{ - StartAfter: "0004", - Limit: 3, - }) + for i := 0; i < 10; i++ { + err = os.WriteFile(filepath.Join(dir, fmt.Sprintf("%04d", i)), []byte("x"), 0600) require.NoError(t, err) - require.Equal(t, []string{"0005", "0006", "0007"}, paths) + } + + paths, err := store.ListFilePaths(context.Background(), ListFileOptions{ + StartAfter: "0004", + Limit: 3, }) + require.NoError(t, err) + require.Equal(t, []string{"0005", "0006", "0007"}, paths) } func TestNewFilesystemDataStore(t *testing.T) { From ff43078290e53a733ac1b1cc4abce861d64dac43 Mon Sep 17 00:00:00 2001 From: Leigh <351529+leighmcculloch@users.noreply.github.com> Date: Wed, 7 Jan 2026 07:27:00 +1000 Subject: [PATCH 08/11] remove directory creation from NewFilesystemDataStoreWithPath --- support/datastore/filesystem.go | 5 ----- 1 file changed, 5 deletions(-) diff --git a/support/datastore/filesystem.go b/support/datastore/filesystem.go index e42ac9bc14..925a46e32a 100644 --- a/support/datastore/filesystem.go +++ b/support/datastore/filesystem.go @@ -49,11 +49,6 @@ func NewFilesystemDataStore(ctx context.Context, datastoreConfig DataStoreConfig // NewFilesystemDataStoreWithPath creates a FilesystemDataStore with the given base path. func NewFilesystemDataStoreWithPath(basePath string) (DataStore, error) { - // Ensure the base path exists - if err := os.MkdirAll(basePath, defaultDirPerms); err != nil { - return nil, fmt.Errorf("failed to create base directory %s: %w", basePath, err) - } - absPath, err := filepath.Abs(basePath) if err != nil { return nil, fmt.Errorf("failed to resolve absolute path: %w", err) From 328fc599130f923f2280314b5a391169f56f01b1 Mon Sep 17 00:00:00 2001 From: Leigh <351529+leighmcculloch@users.noreply.github.com> Date: Wed, 7 Jan 2026 07:30:54 +1000 Subject: [PATCH 09/11] write file content directly without buffering --- support/datastore/filesystem.go | 11 +---------- 1 file changed, 1 insertion(+), 10 deletions(-) diff --git a/support/datastore/filesystem.go b/support/datastore/filesystem.go index 925a46e32a..f8a314e340 100644 --- a/support/datastore/filesystem.go +++ b/support/datastore/filesystem.go @@ -1,7 +1,6 @@ package datastore import ( - "bytes" "context" "errors" "fmt" @@ -174,15 +173,7 @@ func (f *FilesystemDataStore) PutFileIfNotExists( return false, fmt.Errorf("failed to create file %s: %w", path, err) } - // Write content to the file - buf := &bytes.Buffer{} - if _, err := in.WriteTo(buf); err != nil { - file.Close() - os.Remove(fullPath) // Clean up on error - return false, fmt.Errorf("failed to write file %s: %w", path, err) - } - - if _, err := file.Write(buf.Bytes()); err != nil { + if _, err := in.WriteTo(file); err != nil { file.Close() os.Remove(fullPath) // Clean up on error return false, fmt.Errorf("failed to write file %s: %w", path, err) From 6dbbccc633b87e8d17c6ec9b1ef2b6d067c50839 Mon Sep 17 00:00:00 2001 From: Leigh <351529+leighmcculloch@users.noreply.github.com> Date: Wed, 7 Jan 2026 07:33:47 +1000 Subject: [PATCH 10/11] replace "uploaded" with "written" in log messages --- support/datastore/filesystem.go | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/support/datastore/filesystem.go b/support/datastore/filesystem.go index f8a314e340..bf22a33d6e 100644 --- a/support/datastore/filesystem.go +++ b/support/datastore/filesystem.go @@ -147,7 +147,7 @@ func (f *FilesystemDataStore) PutFile(ctx context.Context, path string, in io.Wr return fmt.Errorf("failed to close file %s: %w", path, err) } - log.Debugf("File uploaded successfully: %s", path) + log.Debugf("File written successfully: %s", path) return nil } @@ -183,7 +183,7 @@ func (f *FilesystemDataStore) PutFileIfNotExists( return false, fmt.Errorf("failed to close file %s: %w", path, err) } - log.Debugf("File uploaded successfully: %s", path) + log.Debugf("File written successfully: %s", path) return true, nil } From 3035b760f86fb95ea442da2bce5bea8540af9812 Mon Sep 17 00:00:00 2001 From: Leigh <351529+leighmcculloch@users.noreply.github.com> Date: Thu, 8 Jan 2026 13:21:54 +1000 Subject: [PATCH 11/11] word Co-authored-by: shawn --- support/datastore/filesystem.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/support/datastore/filesystem.go b/support/datastore/filesystem.go index bf22a33d6e..7c77220d12 100644 --- a/support/datastore/filesystem.go +++ b/support/datastore/filesystem.go @@ -31,7 +31,7 @@ var _ DataStore = &FilesystemDataStore{} // // Concurrent writes to the same file path are not safe and may result in // data corruption. Callers must ensure proper synchronization when writing -// to the same path from multiple goroutines. +// to the same path from multiple processes. type FilesystemDataStore struct { basePath string }