Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 7 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
11 changes: 11 additions & 0 deletions docs/DEPLOYMENT.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
2 changes: 1 addition & 1 deletion server/cmd/memd/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down
2 changes: 1 addition & 1 deletion server/internal/api/relocate_integration_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
14 changes: 7 additions & 7 deletions server/internal/file/workspace_lock_integration_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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() {
Expand Down Expand Up @@ -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)
Expand Down
75 changes: 57 additions & 18 deletions server/internal/folder/folder.go
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@ import (
"context"
"errors"
"fmt"
"log/slog"
"strings"
"time"

Expand Down Expand Up @@ -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 (
Expand Down Expand Up @@ -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 {
Expand All @@ -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
Expand All @@ -623,31 +642,51 @@ 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 {
return fmt.Errorf("delete folder: %w", err)
}
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) {
Expand Down
154 changes: 154 additions & 0 deletions server/internal/folder/folder_delete_integration_test.go
Original file line number Diff line number Diff line change
@@ -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
}
2 changes: 1 addition & 1 deletion server/internal/folder/folder_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down
Loading