diff --git a/CHANGELOG.md b/CHANGELOG.md index 104b02f..12723a4 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -13,6 +13,13 @@ The project publishes 0.x prerelease versions; a stable release line is not yet to the canonical `bytefolk` organization while retaining the published npm scope, MCP identity, and existing cache paths. +### Fixed + +- Recursive folder delete now removes associated objects from bucket storage + after the database transaction commits. Previously only DB rows were deleted, + leaving orphan blobs in the bucket. Blob deletion is best-effort and logged + on failure, matching the existing single-file delete behavior. + ### Security - Normalize the client-declared MIME type of a stored file before deciding how diff --git a/docs/DEPLOYMENT.md b/docs/DEPLOYMENT.md index 801e41d..ddfbf7f 100644 --- a/docs/DEPLOYMENT.md +++ b/docs/DEPLOYMENT.md @@ -225,6 +225,17 @@ Redis AOF protects normal restarts but is not in the portable backup. A restore therefore starts with an empty queue/replay window. Requeue or reindex any file whose processing did not reach a terminal state before the backup. +### Object retention + +Deleting a file or recursively deleting a folder removes the database rows and +then deletes the associated objects from the bucket. Object deletion is +best-effort and happens after the database transaction commits, so a process +crash between commit and blob removal may leave orphan objects in the bucket. +These orphans are unreferenced — no live database row points at them — and are +safe to leave in place. They are not automatically reclaimed; operators may +remove them manually if bucket accounting matters. This crash window also +applies to single-file deletion. + ### Restore drill Restore only into an empty installation. The script verifies every checksum diff --git a/server/cmd/memd/main.go b/server/cmd/memd/main.go index 37cb47d..3681230 100644 --- a/server/cmd/memd/main.go +++ b/server/cmd/memd/main.go @@ -117,7 +117,7 @@ func run() error { "providers", cfg.ManagedEmbeddingProviders, ) } - folderSvc := folder.New(database.Pool) + folderSvc := folder.New(database.Pool, store, logger) fileSvc := file.New(database.Pool, store, folderSvc) memorySvc := memory.New(database.Pool) durableContextSvc := durablecontext.New(database.Pool, memorySvc) diff --git a/server/internal/api/content_disposition_integration_test.go b/server/internal/api/content_disposition_integration_test.go index 29baedf..cb2e0a7 100644 --- a/server/internal/api/content_disposition_integration_test.go +++ b/server/internal/api/content_disposition_integration_test.go @@ -105,7 +105,7 @@ func TestContentDispositionHTTPPostgres(t *testing.T) { authService := auth.New(database.Pool) workspaceService := workspace.New(database.Pool) - folderService := folder.New(database.Pool) + folderService := folder.New(database.Pool, nil, nil) user, err := authService.CreateUser( ctx, "content-disposition-"+uuid.NewString()+"@example.test", diff --git a/server/internal/api/relocate_integration_test.go b/server/internal/api/relocate_integration_test.go index 855a449..b18c5e7 100644 --- a/server/internal/api/relocate_integration_test.go +++ b/server/internal/api/relocate_integration_test.go @@ -103,7 +103,7 @@ func TestRelocateHTTPPostgres(t *testing.T) { t.Fatalf("create restricted relocation token: %v", err) } - folderService := folder.New(database.Pool) + folderService := folder.New(database.Pool, nil, nil) fileService := file.New(database.Pool, nil, folderService) server := httptest.NewServer((&Server{ Auth: authService, diff --git a/server/internal/file/workspace_lock_integration_test.go b/server/internal/file/workspace_lock_integration_test.go index 49d0b9c..f47a02b 100644 --- a/server/internal/file/workspace_lock_integration_test.go +++ b/server/internal/file/workspace_lock_integration_test.go @@ -50,15 +50,15 @@ func TestFilePathLockingIntegration(t *testing.T) { t.Run("Put then Rename has one path", func(t *testing.T) { userID, _ := createFileLockTenant(t, ctx, database.Pool, "put-rename") - folders := folder.New(database.Pool) + folders := folder.New(database.Pool, nil, nil) if _, err := folders.Create(ctx, userID, "/A"); err != nil { t.Fatalf("create destination: %v", err) } putBackend := pglockwait.NewBackend(t, ctx, dsn, "file-put") renameBackend := pglockwait.NewBackend(t, ctx, dsn, "file-put-rename") store := newBlockingObjectStore() - service := New(putBackend.Pool, store, folder.New(putBackend.Pool)) - renameFolders := folder.New(renameBackend.Pool) + service := New(putBackend.Pool, store, folder.New(putBackend.Pool, nil, nil)) + renameFolders := folder.New(renameBackend.Pool, nil, nil) type putOutcome struct { result *PutResult @@ -130,7 +130,7 @@ func TestFilePathLockingIntegration(t *testing.T) { t.Run("Move then Rename has one path", func(t *testing.T) { userID, _ := createFileLockTenant(t, ctx, database.Pool, "move-rename") - folders := folder.New(database.Pool) + folders := folder.New(database.Pool, nil, nil) source, err := folders.Create(ctx, userID, "/Source") if err != nil { t.Fatalf("create source: %v", err) @@ -161,9 +161,9 @@ func TestFilePathLockingIntegration(t *testing.T) { service := New( moveBackend.Pool, &recordingObjectStore{}, - folder.New(moveBackend.Pool), + folder.New(moveBackend.Pool, nil, nil), ) - renameFolders := folder.New(renameBackend.Pool) + renameFolders := folder.New(renameBackend.Pool, nil, nil) fileGate, fileGatePID := lockFilesTable(t, ctx, database.Pool) moveDone := make(chan error, 1) go func() { @@ -206,7 +206,7 @@ func TestFilePathLockingIntegration(t *testing.T) { t.Run("Relocate validates and commits path with name", func(t *testing.T) { userID, _ := createFileLockTenant(t, ctx, database.Pool, "atomic-relocate") - folders := folder.New(database.Pool) + folders := folder.New(database.Pool, nil, nil) source, err := folders.Create(ctx, userID, "/RelocateSource") if err != nil { t.Fatalf("create relocate source: %v", err) diff --git a/server/internal/folder/folder.go b/server/internal/folder/folder.go index 0dd0469..88e805c 100644 --- a/server/internal/folder/folder.go +++ b/server/internal/folder/folder.go @@ -13,6 +13,7 @@ import ( "context" "errors" "fmt" + "log/slog" "strings" "time" @@ -58,13 +59,29 @@ type Node struct { Children []*Node `json:"children,omitempty"` } +// blobDeleter is the narrow contract for best-effort object removal. The folder +// package depends on this interface rather than the concrete storage.Store so +// that tests and callers that do not need blob cleanup can pass nil. +type blobDeleter interface { + Delete(ctx context.Context, key string) error +} + // Service is the folder service. type Service struct { - pool *pgxpool.Pool + pool *pgxpool.Pool + store blobDeleter + log *slog.Logger } -// New constructs a folder Service. -func New(pool *pgxpool.Pool) *Service { return &Service{pool: pool} } +// New constructs a folder Service. A nil store skips object-storage cleanup on +// recursive delete (useful for tests that have no bucket). A nil logger falls +// back to slog.Default. +func New(pool *pgxpool.Pool, store blobDeleter, log *slog.Logger) *Service { + if log == nil { + log = slog.Default() + } + return &Service{pool: pool, store: store, log: log} +} // Sentinel errors. var ( @@ -586,9 +603,10 @@ func rewritePrefixTx(ctx context.Context, tx pgx.Tx, userID, srcID uuid.UUID, ol // // - recursive=false (default): folder must be empty (no subfolders, no files) // or ErrNotEmpty is returned. -// - recursive=true: subfolders + files are deleted from the DB. S3 cleanup -// is TODO — for now we only purge the DB rows; orphan blobs will be -// reaped by a future garbage-collection pass. +// - recursive=true: subfolders + files are deleted from the DB and their +// objects are removed from bucket storage. Blob deletion is best-effort +// and happens after the database transaction commits, so a crash between +// commit and blob removal may leave orphan objects (see docs/DEPLOYMENT.md). func (s *Service) Delete(ctx context.Context, userID uuid.UUID, path string, recursive bool) error { norm, err := pathx.Normalize(path) if err != nil { @@ -597,7 +615,8 @@ func (s *Service) Delete(ctx context.Context, userID uuid.UUID, path string, rec if norm == pathx.Root { return ErrRootOp } - return s.withPathMutationTx(ctx, userID, func(tx pgx.Tx) error { + var orphanKeys []string + txErr := s.withPathMutationTx(ctx, userID, func(tx pgx.Tx) error { src, err := selectFolderByPathTx(ctx, tx, userID, norm) if err != nil { return err @@ -623,24 +642,33 @@ func (s *Service) Delete(ctx context.Context, userID uuid.UUID, path string, rec return fmt.Errorf("check recursive delete memories: %w", err) } if containsMemories { - // A folder operation must never become an implicit memory - // deletion. The caller has to use the memory lifecycle's - // explicit forget operation first. return ErrContainsMemories } - // Hard delete: remove all descendant files first (FKs cascade - // from folders → files would only NULL out folder_id, so we have - // to delete files explicitly). - if _, err := tx.Exec(ctx, + rows, err := tx.Query(ctx, `DELETE FROM files WHERE user_id = $1 AND (folder_id = $2 OR path = $3 - OR left(path, length($3) + 1) = $3 || '/')`, - userID, src.ID, src.Path); err != nil { + OR left(path, length($3) + 1) = $3 || '/') + RETURNING storage_key`, + userID, src.ID, src.Path) + if err != nil { return fmt.Errorf("recursive delete files: %w", err) } - // Subfolder rows cascade via the FK ON DELETE CASCADE when we - // drop the parent below. + for rows.Next() { + var key string + if err := rows.Scan(&key); err != nil { + rows.Close() + return fmt.Errorf("scan storage_key: %w", err) + } + if key != "" { + orphanKeys = append(orphanKeys, key) + } + } + if err := rows.Err(); err != nil { + rows.Close() + return fmt.Errorf("iterate deleted storage_keys: %w", err) + } + rows.Close() } if _, err := tx.Exec(ctx, `DELETE FROM folders WHERE id = $1 AND user_id = $2`, src.ID, userID); err != nil { @@ -648,6 +676,17 @@ func (s *Service) Delete(ctx context.Context, userID uuid.UUID, path string, rec } return nil }) + if txErr != nil { + return txErr + } + if s.store != nil { + for _, key := range orphanKeys { + if derr := s.store.Delete(ctx, key); derr != nil { + s.log.Warn("folder.blob_delete_failed", "storage_key", key, "err", derr) + } + } + } + return nil } func isEmptyTx(ctx context.Context, tx pgx.Tx, userID, folderID uuid.UUID, path string) (bool, error) { diff --git a/server/internal/folder/folder_delete_integration_test.go b/server/internal/folder/folder_delete_integration_test.go new file mode 100644 index 0000000..9ef2b7f --- /dev/null +++ b/server/internal/folder/folder_delete_integration_test.go @@ -0,0 +1,154 @@ +package folder + +import ( + "bytes" + "context" + "fmt" + "os" + "strings" + "testing" + + "github.com/google/uuid" + "github.com/jackc/pgx/v5/pgxpool" + + "github.com/PeterGuy326/mem/server/internal/storage" +) + +// TestRecursiveDeleteCleansBlobs verifies that a recursive folder delete +// removes objects from bucket storage, not just DB rows. It requires a live +// PostgreSQL instance (MEM_TEST_DB) and a MinIO endpoint (MEM_TEST_S3_ENDPOINT). +func TestRecursiveDeleteCleansBlobs(t *testing.T) { + dsn := os.Getenv("MEM_TEST_DB") + if dsn == "" { + t.Skip("MEM_TEST_DB not set; skipping") + } + s3Endpoint := os.Getenv("MEM_TEST_S3_ENDPOINT") + if s3Endpoint == "" { + t.Skip("MEM_TEST_S3_ENDPOINT not set; skipping") + } + + s3AccessKey := os.Getenv("MEM_TEST_S3_ACCESS_KEY") + if s3AccessKey == "" { + s3AccessKey = "mem" + } + s3SecretKey := os.Getenv("MEM_TEST_S3_SECRET_KEY") + if s3SecretKey == "" { + s3SecretKey = "mem-minio-password" + } + s3Bucket := os.Getenv("MEM_TEST_S3_BUCKET") + if s3Bucket == "" { + s3Bucket = "mem" + } + + ctx := context.Background() + + cfg, err := pgxpool.ParseConfig(dsn) + if err != nil { + t.Fatalf("parse MEM_TEST_DB: %v", err) + } + if !strings.HasSuffix(cfg.ConnConfig.Database, "_test") { + t.Fatalf("refusing non-test database %q", cfg.ConnConfig.Database) + } + pool, err := pgxpool.NewWithConfig(ctx, cfg) + if err != nil { + t.Fatalf("pgxpool: %v", err) + } + defer pool.Close() + + store, err := storage.New(ctx, s3Endpoint, s3AccessKey, s3SecretKey, s3Bucket, "", false) + if err != nil { + t.Fatalf("storage.New: %v", err) + } + + userID := folderTestUUID(t, ctx, pool, + `INSERT INTO users (email, password_hash) VALUES ($1, 'x') RETURNING id`, + "folder-blob-cleanup-"+uuid.NewString()+"@example.com") + defer func() { + _, _ = pool.Exec(ctx, `DELETE FROM users WHERE id = $1`, userID) + }() + + svc := New(pool, store, nil) + + if _, err := svc.Create(ctx, userID, "/CleanupParent/Child"); err != nil { + t.Fatalf("create folders: %v", err) + } + + childFolder, err := svc.Get(ctx, userID, "/CleanupParent/Child") + if err != nil { + t.Fatalf("get child folder: %v", err) + } + + type testFile struct { + id uuid.UUID + name string + path string + folderID uuid.UUID + storageKey string + } + files := []testFile{ + { + id: uuid.New(), + name: "a.txt", + path: "/CleanupParent", + folderID: mustFolderID(t, ctx, pool, userID, "/CleanupParent"), + storageKey: fmt.Sprintf("users/%s/%s/a.txt", userID, uuid.New()), + }, + { + id: uuid.New(), + name: "b.txt", + path: "/CleanupParent/Child", + folderID: childFolder.ID, + storageKey: fmt.Sprintf("users/%s/%s/b.txt", userID, uuid.New()), + }, + } + + for _, f := range files { + if err := store.Put(ctx, f.storageKey, bytes.NewReader([]byte("content-"+f.name)), int64(len("content-"+f.name)), "text/plain"); err != nil { + t.Fatalf("put object %s: %v", f.storageKey, err) + } + if _, err := pool.Exec(ctx, + `INSERT INTO files (id, user_id, name, path, folder_id, size, sha256, mime, storage_key, tags, index_status) + VALUES ($1, $2, $3, $4, $5, $6, $7, 'text/plain', $8, '{}', 'pending')`, + f.id, userID, f.name, f.path, f.folderID, len("content-"+f.name), + "sha256-"+f.id.String(), f.storageKey, + ); err != nil { + t.Fatalf("insert file row %s: %v", f.name, err) + } + } + + for _, f := range files { + if _, err := store.Get(ctx, f.storageKey); err != nil { + t.Fatalf("precondition: object %s should exist before delete: %v", f.storageKey, err) + } + } + + if err := svc.Delete(ctx, userID, "/CleanupParent", true); err != nil { + t.Fatalf("recursive delete: %v", err) + } + + var remaining int + if err := pool.QueryRow(ctx, + `SELECT COUNT(*) FROM files WHERE user_id = $1 AND path LIKE '/CleanupParent%'`, + userID).Scan(&remaining); err != nil { + t.Fatalf("count remaining files: %v", err) + } + if remaining != 0 { + t.Fatalf("expected 0 file rows after recursive delete, got %d", remaining) + } + + for _, f := range files { + if _, err := store.Get(ctx, f.storageKey); err == nil { + t.Errorf("object %s still exists in bucket after recursive delete", f.storageKey) + } + } +} + +func mustFolderID(t *testing.T, ctx context.Context, pool *pgxpool.Pool, userID uuid.UUID, path string) uuid.UUID { + t.Helper() + var id uuid.UUID + if err := pool.QueryRow(ctx, + `SELECT id FROM folders WHERE user_id = $1 AND path = $2`, userID, path).Scan(&id); err != nil { + t.Fatalf("lookup folder %q: %v", path, err) + } + return id +} diff --git a/server/internal/folder/folder_test.go b/server/internal/folder/folder_test.go index dcb9fc9..21fd464 100644 --- a/server/internal/folder/folder_test.go +++ b/server/internal/folder/folder_test.go @@ -228,7 +228,7 @@ func TestMemoryPathLifecycleIntegration(t *testing.T) { VALUES ('folder memory test', $1) RETURNING id`, userID) - svc := New(pool) + svc := New(pool, nil, nil) for _, p := range []string{ "/A%_/Child", "/AtomicSource/original/Child", diff --git a/server/internal/folder/workspace_lock_integration_test.go b/server/internal/folder/workspace_lock_integration_test.go index 53d3eee..caf433d 100644 --- a/server/internal/folder/workspace_lock_integration_test.go +++ b/server/internal/folder/workspace_lock_integration_test.go @@ -53,12 +53,12 @@ func TestWorkspacePathLockingIntegration(t *testing.T) { t.Run("content lock blocks folder rewrite", func(t *testing.T) { userID, _ := createWorkspaceLockTenant(t, ctx, database.Pool, "content-blocks-folder") - service := New(database.Pool) + service := New(database.Pool, nil, nil) if _, err := service.Create(ctx, userID, "/Shared/Child"); err != nil { t.Fatalf("create source: %v", err) } renameBackend := pglockwait.NewBackend(t, ctx, dsn, "folder-rewrite") - renameService := New(renameBackend.Pool) + renameService := New(renameBackend.Pool, nil, nil) writerTx, err := database.Pool.BeginTx(ctx, pgx.TxOptions{}) if err != nil { @@ -105,7 +105,7 @@ func TestWorkspacePathLockingIntegration(t *testing.T) { createBackend := pglockwait.NewBackend(t, ctx, dsn, "folder-create") memoryBackend := pglockwait.NewBackend(t, ctx, dsn, "memory-remember") checkpointBackend := pglockwait.NewBackend(t, ctx, dsn, "handoff-checkpoint") - folderService := New(createBackend.Pool) + folderService := New(createBackend.Pool, nil, nil) memoryService := memory.New(memoryBackend.Pool) handoffService := handoff.New(checkpointBackend.Pool) @@ -177,7 +177,7 @@ func TestWorkspacePathLockingIntegration(t *testing.T) { t.Run("concurrent renames cannot split a subtree", func(t *testing.T) { userID, _ := createWorkspaceLockTenant(t, ctx, database.Pool, "rename-race") - service := New(database.Pool) + service := New(database.Pool, nil, nil) child, err := service.Create(ctx, userID, "/Race/Child") if err != nil { t.Fatalf("create source subtree: %v", err) @@ -207,8 +207,8 @@ func TestWorkspacePathLockingIntegration(t *testing.T) { service *Service } attempts := []renameAttempt{ - {name: "RaceB", backend: renameBBackend, service: New(renameBBackend.Pool)}, - {name: "RaceC", backend: renameCBackend, service: New(renameCBackend.Pool)}, + {name: "RaceB", backend: renameBBackend, service: New(renameBBackend.Pool, nil, nil)}, + {name: "RaceC", backend: renameCBackend, service: New(renameCBackend.Pool, nil, nil)}, } gateTx, err := database.Pool.BeginTx(ctx, pgx.TxOptions{})