From 59f4cb8be0af192ee3ebe75c86e5a32dc462880a Mon Sep 17 00:00:00 2001 From: Rhaqim Date: Tue, 11 Aug 2026 22:33:20 +0100 Subject: [PATCH 1/2] migration concurrency worker, retry and failure count --- CHANGELOG.md | 24 ++ README.md | 22 +- buckt.go | 31 ++- buckt_conf.go | 17 ++ buckt_migration_test.go | 130 +++++++++++ client/web/app/api.go | 2 + example/client/web/ui/main.go | 8 +- internal/backend/migrate.go | 199 +++++++++++++--- internal/domain/backend.go | 6 +- internal/migration/interface.go | 30 --- internal/migration/service.go | 397 -------------------------------- internal/mocks/backend.go | 4 +- 12 files changed, 393 insertions(+), 477 deletions(-) delete mode 100644 internal/migration/interface.go delete mode 100644 internal/migration/service.go diff --git a/CHANGELOG.md b/CHANGELOG.md index 8ae3a92..7a3df16 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -25,6 +25,30 @@ Additive; no API removals. (`.pre-commit-config.yaml` + `.gitleaks.toml`) and a `secret-scan` GitHub Actions workflow catch hardcoded credentials before they reach history. +### โšก Changed + +- **`MigrateAll` now copies files concurrently** with a bounded worker pool + instead of one at a time, and checks the target with a single `List` up front + rather than an `Exists` round-trip per file. On a high-latency target (S3/R2) + this is dramatically faster for large buckets. Concurrency is configurable via + `MigrationConfig.Concurrency` (default `DefaultMigrationConcurrency` = 8); + higher values trade memory (each in-flight file is buffered) and provider + rate-limit headroom for throughput. Still idempotent and safe to re-run after + an interruption. Individual file copies now retry transient failures with a + linear backoff, so a flaky target no longer silently drops files. +- **`MigrationStatus` progress now always reaches `total`.** A file that + permanently fails after retries is counted as processed (not dropped), so + `completed` reaches `total` when the run finishes and a progress badge no + longer hangs. The new `Client.MigrationFailures(ctx)` reports how many objects + could not be copied (re-run `MigrateAll` to retry them); the web `/backend` + endpoint gained a `failed` field. + +### ๐Ÿงน Removed + +- Deleted the unused `internal/migration` package (a second, never-wired + migration service). Its worthwhile piece โ€” per-file retry with backoff โ€” was + folded into the live migration path above. + ## [1.8.0] โ€” 2026-08-11 A backward-compatible **minor** release on top of v1.7.0 adding controls to diff --git a/README.md b/README.md index d2697dc..679fbbb 100644 --- a/README.md +++ b/README.md @@ -385,8 +385,26 @@ Migration is **always forward** (primary โ†’ secondary). The primary is treated } time.Sleep(time.Second) } + + // completed counts every processed file, so the loop above always terminates + // even if some files can't be copied. Check how many were left behind: + if failed, _ := client.MigrationFailures(ctx); failed > 0 { + log.Printf("%d file(s) failed after retries โ€” fix the cause and re-run MigrateAll", failed) + } ``` +`MigrateAll` copies files **concurrently** with a bounded worker pool. Tune it for a high-latency target with `Concurrency` (default 8): + +```go + client, _ := buckt.Default(buckt.WithMigration(buckt.MigrationConfig{ + From: buckt.LocalBackend(), + To: s3, + Concurrency: 16, // more parallelism โ†’ faster, but more memory + API pressure + })) +``` + +Each in-flight file is buffered in full, so higher concurrency trades memory and provider rate-limit headroom for throughput. + `client.BackendName()` reports the active backend (`"local"`, `"s3"`, or `"local->s3"` mid-migration) โ€” handy for a status badge. When you're done migrating, swap to the target alone: ```go @@ -397,7 +415,8 @@ Migration is **always forward** (primary โ†’ secondary). The primary is treated | Method | Description | |---|---| | `MigrateAll(ctx)` | Schedule a background copy of all pre-existing objects to the target (idempotent) | -| `MigrationStatus(ctx)` | Report `completed`/`total` copied and whether migration is enabled | +| `MigrationStatus(ctx)` | Report `completed`/`total` processed and whether migration is enabled (`completed == total` when done) | +| `MigrationFailures(ctx)` | How many objects permanently failed after retries (re-run `MigrateAll` to retry them) | | `BackendName()` | Name of the active backend (`local`, `s3`, or `local->s3`) | --- @@ -756,6 +775,7 @@ GetDerivative(fileID, name string) (data []byte, contentType string, err error) // Migration (only when built WithMigration) MigrateAll(ctx context.Context) error MigrationStatus(ctx context.Context) (completed, total int64, ok bool) +MigrationFailures(ctx context.Context) (failed int64, ok bool) BackendName() string ``` diff --git a/buckt.go b/buckt.go index 32ef983..0b3962b 100644 --- a/buckt.go +++ b/buckt.go @@ -947,17 +947,34 @@ func (b *Client) MigrateAll(ctx context.Context) error { return m.MigrateAll(ctx) } -// MigrationStatus reports how many objects have been copied to the target so far -// and the total scheduled by the most recent MigrateAll. ok is false when the -// Client was not created with WithMigration. When a migration finishes, -// completed == total. +// MigrationStatus reports how many objects have been processed by the most +// recent MigrateAll and the total scheduled. ok is false when the Client was not +// created with WithMigration. When a migration finishes, completed == total โ€” +// completed counts every processed object, whether it was copied, already +// present, or permanently failed after retries. Use MigrationFailures to find +// out how many of those permanently failed (successfully copied == completed - +// failed). func (b *Client) MigrationStatus(ctx context.Context) (completed, total int64, ok bool) { m, isMigratable := b.backend.(domain.MigratableBackend) if !isMigratable { return 0, 0, false } - completed, total = m.MigrationStatus(ctx) - return completed, total, true + copied, failed, tot := m.MigrationStatus(ctx) + return copied + failed, tot, true +} + +// MigrationFailures reports how many objects the most recent MigrateAll could +// not copy after retries. ok is false when the Client was not created with +// WithMigration. A non-zero count means the migration finished (completed == +// total) but some objects were left behind โ€” inspect the logs for which, and +// re-run MigrateAll to retry them once the underlying issue is fixed. +func (b *Client) MigrationFailures(ctx context.Context) (failed int64, ok bool) { + m, isMigratable := b.backend.(domain.MigratableBackend) + if !isMigratable { + return 0, false + } + _, failed, _ = m.MigrationStatus(ctx) + return failed, true } /* Helper Methods */ @@ -1083,7 +1100,7 @@ func resolveBackend(mediaDir string, bc BackendConfig, log domain.BucktLogger, l } log.Infof("๐Ÿ”„ Migration mode: %s โ†’ %s", source.Name(), target.Name()) - return backend.NewMigrationBackend(log, meter(source), meter(target)) + return backend.NewMigrationBackend(log, meter(source), meter(target), bc.MigrationConcurrency) } // Non-migration modes diff --git a/buckt_conf.go b/buckt_conf.go index 304d45e..5740cf2 100644 --- a/buckt_conf.go +++ b/buckt_conf.go @@ -130,6 +130,10 @@ type BackendConfig struct { // MigrationEnabled enables dual-write migration mode. MigrationEnabled bool + + // MigrationConcurrency sets how many files MigrateAll copies in parallel. + // 0 uses DefaultMigrationConcurrency. + MigrationConcurrency int } // LocalBackend is a placeholder and is replaced with the actual local backend implementation. @@ -171,6 +175,12 @@ const DefaultMaxTrashBatchSize = 5000 // so this also bounds the worst-case lock-contention window. const DefaultBackendOpTimeout = 5 * time.Minute +// DefaultMigrationConcurrency is how many files MigrateAll copies in parallel +// when MigrationConfig.Concurrency is left at 0. Chosen to overlap network +// latency to a cloud target without hammering it or buffering too many files +// in memory at once. +const DefaultMigrationConcurrency = 8 + type Config struct { MediaDir string FlatNameSpaces bool @@ -333,6 +343,12 @@ type MigrationConfig struct { // To is the destination backend that From is being migrated to. To Backend + + // Concurrency sets how many files MigrateAll copies in parallel. 0 uses a + // sensible default (DefaultMigrationConcurrency). Higher values speed up bulk + // migration to a high-latency target (S3/R2) but use more memory โ€” each + // in-flight file is buffered in full โ€” and can trip provider rate limits. + Concurrency int } // WithBackend configures a single storage backend. @@ -462,6 +478,7 @@ func WithMigration(mc MigrationConfig) ConfigFunc { c.Backend.Source = mc.From c.Backend.Target = mc.To c.Backend.MigrationEnabled = true + c.Backend.MigrationConcurrency = mc.Concurrency } } diff --git a/buckt_migration_test.go b/buckt_migration_test.go index 4a89f04..fee6982 100644 --- a/buckt_migration_test.go +++ b/buckt_migration_test.go @@ -22,6 +22,8 @@ type memBackend struct { mu sync.Mutex objs map[string][]byte puts int // total Put calls, to detect wasteful re-uploads + // failPuts makes the next N Put calls fail, to exercise retry behaviour. + failPuts int } func newMemBackend() *memBackend { return &memBackend{objs: map[string][]byte{}} } @@ -32,6 +34,10 @@ func (m *memBackend) Put(_ context.Context, path string, data []byte) error { m.mu.Lock() defer m.mu.Unlock() m.puts++ + if m.failPuts > 0 { + m.failPuts-- + return fmt.Errorf("mem: transient put failure for %s", path) + } cp := make([]byte, len(data)) copy(cp, data) m.objs[path] = cp @@ -236,6 +242,130 @@ func TestMigration_MigrateAllIsIdempotent(t *testing.T) { assert.Equal(t, objectsAfterFirst, target.count(), "the object set is unchanged on the second pass") } +// TestMigration_ConcurrentCopiesEachFileOnce drives MigrateAll with an explicit +// bounded worker pool over many files and asserts every file is copied exactly +// once (no double-uploads from a race, none skipped), threading the +// MigrationConcurrency config end-to-end. +func TestMigration_ConcurrentCopiesEachFileOnce(t *testing.T) { + dir := t.TempDir() + mediaDir := filepath.Join(dir, "media") + dbPath := filepath.Join(dir, "db.sqlite") + const user = "u1" + const nFiles = 25 + + db1, err := sql.Open("sqlite3", dbPath) + require.NoError(t, err) + c1, err := New(Config{DB: DBConfig{Driver: SQLite, Database: db1}, MediaDir: mediaDir, Log: LogConfig{Silence: true}}) + require.NoError(t, err) + for i := 0; i < nFiles; i++ { + _, err := c1.UploadFile(user, "", fmt.Sprintf("file-%02d.txt", i), "text/plain", []byte(fmt.Sprintf("data-%02d", i))) + require.NoError(t, err) + } + require.NoError(t, c1.Close()) + require.NoError(t, db1.Close()) + + target := newMemBackend() + db2, err := sql.Open("sqlite3", dbPath) + require.NoError(t, err) + c2, err := New(Config{ + DB: DBConfig{Driver: SQLite, Database: db2}, + MediaDir: mediaDir, + Log: LogConfig{Silence: true}, + Backend: BackendConfig{Source: LocalBackend(), Target: target, MigrationEnabled: true, MigrationConcurrency: 4}, + }) + require.NoError(t, err) + t.Cleanup(func() { _ = c2.Close(); _ = db2.Close() }) + + require.NoError(t, c2.MigrateAll(context.Background())) + requireMigrationDone(t, c2) + + assert.Equal(t, nFiles, target.count(), "every file copied to the target") + assert.Equal(t, nFiles, target.putCalls(), "each file copied exactly once (no double-uploads)") + done, total, _ := c2.MigrationStatus(context.Background()) + assert.Equal(t, int64(nFiles), total) + assert.Equal(t, total, done) +} + +// TestMigration_RetriesTransientFailure verifies MigrateFile retries a file +// whose first write fails transiently, so a flaky target doesn't drop files. +func TestMigration_RetriesTransientFailure(t *testing.T) { + dir := t.TempDir() + mediaDir := filepath.Join(dir, "media") + dbPath := filepath.Join(dir, "db.sqlite") + const user = "u1" + + db1, err := sql.Open("sqlite3", dbPath) + require.NoError(t, err) + c1, err := New(Config{DB: DBConfig{Driver: SQLite, Database: db1}, MediaDir: mediaDir, Log: LogConfig{Silence: true}}) + require.NoError(t, err) + _, err = c1.UploadFile(user, "", "only.txt", "text/plain", []byte("payload")) + require.NoError(t, err) + require.NoError(t, c1.Close()) + require.NoError(t, db1.Close()) + + target := newMemBackend() + target.failPuts = 1 // first copy attempt fails; the retry should succeed + + db2, err := sql.Open("sqlite3", dbPath) + require.NoError(t, err) + c2, err := New(Config{ + DB: DBConfig{Driver: SQLite, Database: db2}, + MediaDir: mediaDir, + Log: LogConfig{Silence: true}, + Backend: BackendConfig{Source: LocalBackend(), Target: target, MigrationEnabled: true}, + }) + require.NoError(t, err) + t.Cleanup(func() { _ = c2.Close(); _ = db2.Close() }) + + require.NoError(t, c2.MigrateAll(context.Background())) + requireMigrationDone(t, c2) + + assert.Equal(t, 1, target.count(), "file lands despite the first attempt failing") +} + +// TestMigration_PermanentFailureStillCompletes verifies that a file which fails +// every retry is counted as processed (so status reaches total and a progress +// badge doesn't hang) and surfaced via MigrationFailures. +func TestMigration_PermanentFailureStillCompletes(t *testing.T) { + dir := t.TempDir() + mediaDir := filepath.Join(dir, "media") + dbPath := filepath.Join(dir, "db.sqlite") + const user = "u1" + + db1, err := sql.Open("sqlite3", dbPath) + require.NoError(t, err) + c1, err := New(Config{DB: DBConfig{Driver: SQLite, Database: db1}, MediaDir: mediaDir, Log: LogConfig{Silence: true}}) + require.NoError(t, err) + _, err = c1.UploadFile(user, "", "doomed.txt", "text/plain", []byte("payload")) + require.NoError(t, err) + require.NoError(t, c1.Close()) + require.NoError(t, db1.Close()) + + target := newMemBackend() + target.failPuts = 100 // exceed any retry budget โ†’ permanent failure + + db2, err := sql.Open("sqlite3", dbPath) + require.NoError(t, err) + c2, err := New(Config{ + DB: DBConfig{Driver: SQLite, Database: db2}, + MediaDir: mediaDir, + Log: LogConfig{Silence: true}, + Backend: BackendConfig{Source: LocalBackend(), Target: target, MigrationEnabled: true}, + }) + require.NoError(t, err) + t.Cleanup(func() { _ = c2.Close(); _ = db2.Close() }) + + require.NoError(t, c2.MigrateAll(context.Background())) + // requireMigrationDone asserts completed == total โ€” it would hang (and fail + // the timeout) if a permanently-failed file never counted as processed. + requireMigrationDone(t, c2) + + failed, ok := c2.MigrationFailures(context.Background()) + require.True(t, ok) + assert.Equal(t, int64(1), failed, "the un-copyable file is reported as failed") + assert.Equal(t, 0, target.count(), "nothing landed in the target") +} + func TestMigration_NotEnabledReturnsError(t *testing.T) { c := newTestClient(t) // plain local backend diff --git a/client/web/app/api.go b/client/web/app/api.go index 752a83d..8496a52 100644 --- a/client/web/app/api.go +++ b/client/web/app/api.go @@ -68,12 +68,14 @@ func (svc *APIService) Metrics(c *gin.Context) { // use (and both backends during a migration) plus live status. func (svc *APIService) Backend(c *gin.Context) { completed, total, enabled := svc.client.MigrationStatus(c.Request.Context()) + failed, _ := svc.client.MigrationFailures(c.Request.Context()) c.JSON(200, gin.H{ "name": svc.client.BackendName(), "migration": gin.H{ "enabled": enabled, "running": enabled && total > 0 && completed < total, "completed": completed, + "failed": failed, "total": total, }, }) diff --git a/example/client/web/ui/main.go b/example/client/web/ui/main.go index 2c52da3..9207716 100644 --- a/example/client/web/ui/main.go +++ b/example/client/web/ui/main.go @@ -140,9 +140,13 @@ func bulkMigrate(client *buckt.Client) { log.Println("migration: no pre-existing files to copy") return } - log.Printf("migration: %d/%d objects copied", done, total) + log.Printf("migration: %d/%d objects processed", done, total) if done >= total { - log.Println("โœ… migration complete โ€” you can now restart with -mode=r2") + if failed, _ := client.MigrationFailures(ctx); failed > 0 { + log.Printf("โš ๏ธ migration finished with %d failure(s) โ€” check the logs and re-run to retry them", failed) + } else { + log.Println("โœ… migration complete โ€” you can now restart with -mode=r2") + } return } time.Sleep(time.Second) diff --git a/internal/backend/migrate.go b/internal/backend/migrate.go index cc9210d..c7a6efd 100644 --- a/internal/backend/migrate.go +++ b/internal/backend/migrate.go @@ -4,8 +4,10 @@ import ( "context" "fmt" "io" + "strings" "sync" "sync/atomic" + "time" "github.com/Rhaqim/buckt/internal/domain" ) @@ -16,22 +18,35 @@ type MigrationBackendService struct { primaryBackend domain.FileBackend secondaryBackend domain.FileBackend + // concurrency is how many files MigrateAll copies in parallel (>= 1). + concurrency int + migrating atomic.Bool stats migrationStats } type migrationStats struct { mu sync.Mutex - completed int64 + completed int64 // files copied or already present in the secondary + failed int64 // files that permanently failed after retries total int64 } -func NewMigrationBackend(bucktLogger domain.BucktLogger, primary domain.FileBackend, secondary domain.FileBackend) domain.MigratableBackend { +// defaultMigrationConcurrency mirrors buckt.DefaultMigrationConcurrency. It is +// duplicated here (rather than imported) because the root package imports this +// one, not the other way around. +const defaultMigrationConcurrency = 8 + +func NewMigrationBackend(bucktLogger domain.BucktLogger, primary domain.FileBackend, secondary domain.FileBackend, concurrency int) domain.MigratableBackend { + if concurrency <= 0 { + concurrency = defaultMigrationConcurrency + } bucktLogger.Info("๐Ÿš€ Initialising migration backend: " + primary.Name() + " -> " + secondary.Name()) return &MigrationBackendService{ logger: bucktLogger, primaryBackend: primary, secondaryBackend: secondary, + concurrency: concurrency, } } @@ -177,25 +192,56 @@ func (d *MigrationBackendService) DeleteFolder(ctx context.Context, prefix strin return nil } -// MigrateFile copies a single file from primary to secondary. +const ( + // migrateMaxAttempts bounds how many times MigrateFile tries a single file + // before giving up, so a transient primary/secondary hiccup doesn't drop the + // file from the migration. + migrateMaxAttempts = 3 + // migrateRetryBackoff is the base delay between attempts, scaled by the + // attempt number (linear backoff). + migrateRetryBackoff = 500 * time.Millisecond +) + +// MigrateFile copies a single file from primary to secondary, retrying transient +// read/write failures with a linear backoff. On success it increments the +// completed counter exactly once; on exhaustion it returns the last error. func (d *MigrationBackendService) MigrateFile(ctx context.Context, path string) error { - data, err := d.primaryBackend.Get(ctx, path) - if err != nil { - return fmt.Errorf("failed to read %s from primary: %w", path, err) - } + var lastErr error + for attempt := 1; attempt <= migrateMaxAttempts; attempt++ { + if attempt > 1 { + // Wait before retrying, but bail immediately if cancelled. + select { + case <-ctx.Done(): + return ctx.Err() + case <-time.After(migrateRetryBackoff * time.Duration(attempt-1)): + } + } - if err := d.secondaryBackend.Put(ctx, path, data); err != nil { - return fmt.Errorf("failed to write %s to secondary: %w", path, err) - } + data, err := d.primaryBackend.Get(ctx, path) + if err != nil { + lastErr = fmt.Errorf("failed to read %s from primary: %w", path, err) + continue + } - d.stats.mu.Lock() - d.stats.completed++ - d.stats.mu.Unlock() + if err := d.secondaryBackend.Put(ctx, path, data); err != nil { + lastErr = fmt.Errorf("failed to write %s to secondary: %w", path, err) + continue + } - return nil + d.stats.mu.Lock() + d.stats.completed++ + d.stats.mu.Unlock() + + return nil + } + + return fmt.Errorf("migrate %s failed after %d attempts: %w", path, migrateMaxAttempts, lastErr) } -// MigrateAll copies all files from primary to secondary in the background. +// MigrateAll copies all files from primary to secondary in the background, +// using a bounded pool of d.concurrency workers. It is idempotent: objects +// already present in the secondary are skipped, so it is safe to re-run after +// an interruption. func (d *MigrationBackendService) MigrateAll(ctx context.Context) error { // Atomic check-and-set so two concurrent callers can't both start a // migration. CompareAndSwap returns false if the value was already true. @@ -215,42 +261,123 @@ func (d *MigrationBackendService) MigrateAll(ctx context.Context) error { d.stats.completed = 0 d.stats.mu.Unlock() - // Run migration in a goroutine + // Build a skip-set from a single List of the secondary, instead of one + // Exists round-trip per file (halving the request count on large buckets). + // Keys are normalised with normaliseKey so both the leading-slash and + // slash-stripped forms of a key match โ€” see the cloud backends' altKey. + skip, haveSkipSet := d.secondarySkipSet(ctx) + + // alreadyInSecondary reports whether the secondary already holds path. It + // uses the pre-listed skip-set when available, and otherwise falls back to a + // per-file Exists check (e.g. when the secondary's List failed). + alreadyInSecondary := func(path string) bool { + if haveSkipSet { + _, ok := skip[normaliseKey(path)] + return ok + } + exists, err := d.secondaryBackend.Exists(ctx, path) + return err == nil && exists + } + + workers := d.concurrency + if workers <= 0 { + workers = defaultMigrationConcurrency + } + if workers > len(files) { + workers = len(files) + } + + // Run migration in the background. go func() { defer d.migrating.Store(false) + if len(files) == 0 { + d.logger.Info("โœ… Migration complete (nothing to copy)") + return + } + + paths := make(chan string) + + var wg sync.WaitGroup + wg.Add(workers) + for i := 0; i < workers; i++ { + go func() { + defer wg.Done() + for path := range paths { + if ctx.Err() != nil { + return // cancelled โ€” drain quietly + } + if alreadyInSecondary(path) { + d.stats.mu.Lock() + d.stats.completed++ + d.stats.mu.Unlock() + continue + } + if err := d.MigrateFile(ctx, path); err != nil { + d.logger.Errorf("Failed to migrate %s: %v", path, err) + // Count it as processed-but-failed so progress still + // reaches total (the badge completes) and callers get a + // failure count. Continue with other files. + d.stats.mu.Lock() + d.stats.failed++ + d.stats.mu.Unlock() + } + } + }() + } + + // Feed paths to the workers, stopping early if the context is cancelled. + cancelled := false + feed: for _, path := range files { select { case <-ctx.Done(): - d.logger.Errorf("Migration cancelled: %v", ctx.Err()) - return - default: - } - - // Skip if already exists in secondary - exists, err := d.secondaryBackend.Exists(ctx, path) - if err == nil && exists { - d.stats.mu.Lock() - d.stats.completed++ - d.stats.mu.Unlock() - continue - } - - if err := d.MigrateFile(ctx, path); err != nil { - d.logger.Errorf("Failed to migrate %s: %v", path, err) - // Continue with other files + cancelled = true + break feed + case paths <- path: } } + close(paths) + wg.Wait() + if cancelled || ctx.Err() != nil { + d.logger.Errorf("Migration cancelled: %v", ctx.Err()) + return + } d.logger.Info("โœ… Migration complete") }() return nil } -// MigrationStatus returns the current migration progress. -func (d *MigrationBackendService) MigrationStatus(ctx context.Context) (completed int64, total int64) { +// secondarySkipSet lists the secondary once and returns the set of normalised +// keys it already holds. The bool is false when the listing failed, signalling +// callers to fall back to per-file Exists checks. +func (d *MigrationBackendService) secondarySkipSet(ctx context.Context) (map[string]struct{}, bool) { + keys, err := d.secondaryBackend.List(ctx, "") + if err != nil { + d.logger.Errorf("Could not list secondary for skip-set (falling back to per-file checks): %v", err) + return nil, false + } + set := make(map[string]struct{}, len(keys)) + for _, k := range keys { + set[normaliseKey(k)] = struct{}{} + } + return set, true +} + +// normaliseKey strips a single leading slash so the two key forms buckt can +// produce ("/user/โ€ฆ" from nested-mode writes and "user/โ€ฆ" from a local List) +// compare equal. +func normaliseKey(key string) string { + return strings.TrimPrefix(key, "/") +} + +// MigrationStatus returns the current migration progress: files copied (or +// already present), files that permanently failed after retries, and the total +// scheduled. completed+failed == total once the run finishes. +func (d *MigrationBackendService) MigrationStatus(ctx context.Context) (completed int64, failed int64, total int64) { d.stats.mu.Lock() defer d.stats.mu.Unlock() - return d.stats.completed, d.stats.total + return d.stats.completed, d.stats.failed, d.stats.total } diff --git a/internal/domain/backend.go b/internal/domain/backend.go index 6200692..c4445d6 100644 --- a/internal/domain/backend.go +++ b/internal/domain/backend.go @@ -46,8 +46,10 @@ type MigratableBackend interface { // Migrate a specific file (used for lazy migration on access) MigrateFile(ctx context.Context, path string) error - // Progress info for observability - MigrationStatus(ctx context.Context) (completed int64, total int64) + // Progress info for observability. completed counts files copied (or already + // present); failed counts files that permanently failed after retries; + // total is the number scheduled. completed+failed == total when done. + MigrationStatus(ctx context.Context) (completed int64, failed int64, total int64) } type PlaceholderBackend struct { diff --git a/internal/migration/interface.go b/internal/migration/interface.go deleted file mode 100644 index 3b395df..0000000 --- a/internal/migration/interface.go +++ /dev/null @@ -1,30 +0,0 @@ -package migration - -import ( - "time" -) - -type MigrationMode int - -const ( - MigrateModeNone MigrationMode = iota - MigrateModeToSecondary // primary --> secondary (e.g., local -> s3) - MigrateModeFromSecondary // secondary --> primary (e.g., s3 -> local) -) - -type MigrationConfig struct { - Concurrency int - RetryCount int - RetryBackoff time.Duration - DeleteAfterCopy bool // remove source after successful migration - PersistPath string // where to store checkpoint file if primary is local -} - -type migrationState struct { - Prefix string `json:"prefix"` - Processed map[string]bool `json:"processed"` - Total int64 `json:"total"` - Completed int64 `json:"completed"` - StartedAt time.Time `json:"started_at"` - UpdatedAt time.Time `json:"updated_at"` -} diff --git a/internal/migration/service.go b/internal/migration/service.go deleted file mode 100644 index 5cc97f4..0000000 --- a/internal/migration/service.go +++ /dev/null @@ -1,397 +0,0 @@ -package migration - -import ( - "context" - "encoding/json" - "errors" - "fmt" - "os" - "sync" - "sync/atomic" - "time" - - "github.com/Rhaqim/buckt/internal/domain" -) - -// MigrationBackendService manages dual-backend behaviour and migration. -type MigrationBackendService struct { - logger domain.BucktLogger - - primary domain.FileBackend - secondary domain.FileBackend - - // migration - mu sync.RWMutex - mode MigrationMode - active atomic.Bool - cfg MigrationConfig - checkpointMux sync.Mutex - state *migrationState - statePath string - - cancelMigration context.CancelFunc - // wg sync.WaitGroup -} - -// NewMigrationBackend unchanged except returns *MigrationBackendService -/* Example usage: - -var logger = domain.NewLogger() -var localBackend = domain.NewLocalFileBackend() -var s3Backend = domain.NewS3FileBackend() - -var mgr = NewMigrationBackend(logger, localBackend, s3Backend) - -mgr.EnableMigration(ctx, s3Backend, MigrateModeToSecondary, &MigrationConfig{ - Concurrency: 16, - DeleteAfterCopy: false, - PersistPath: "/var/run/bucket_migration_state.json", -}) -go func() { - err := mgr.MigrateTo(ctx, "images/", func(p string){ fmt.Println("migrated", p) }, func(p string, e error){ fmt.Println("err", p, e) }) - if err != nil { - log.Println("migration finished with error", err) - } -}() -defer func() { - mgr.DisableMigration(ctx) -}() -*/ -func NewMigrationBackend(logger domain.BucktLogger, primary domain.FileBackend, secondary domain.FileBackend) *MigrationBackendService { - return &MigrationBackendService{ - logger: logger, - primary: primary, - secondary: secondary, - mode: MigrateModeNone, - } -} - -func (d *MigrationBackendService) EnableMigration(ctx context.Context, target domain.FileBackend, mode MigrationMode, cfg *MigrationConfig) error { - if target == nil { - return fmt.Errorf("target backend cannot be nil") - } - if cfg == nil { - cfg = &MigrationConfig{} - } - // defaults - if cfg.Concurrency <= 0 { - cfg.Concurrency = 8 - } - if cfg.RetryCount <= 0 { - cfg.RetryCount = 3 - } - if cfg.RetryBackoff == 0 { - cfg.RetryBackoff = 500 * time.Millisecond - } - if cfg.PersistPath == "" { - // simple default: use cwd/.migration_state.json - cfg.PersistPath = ".migration_state.json" - } - - d.mu.Lock() - defer d.mu.Unlock() - - d.secondary = target - d.mode = mode - d.cfg = *cfg - d.statePath = cfg.PersistPath - - return nil -} - -func (d *MigrationBackendService) DisableMigration(ctx context.Context) { - d.mu.Lock() - defer d.mu.Unlock() - if d.cancelMigration != nil { - d.cancelMigration() - } - d.mode = MigrateModeNone - d.secondary = nil - d.active.Store(false) -} - -func (d *MigrationBackendService) Put(ctx context.Context, path string, data []byte) error { - d.mu.RLock() - mode := d.mode - secondary := d.secondary - d.mu.RUnlock() - - switch mode { - case MigrateModeToSecondary: - // write to secondary primarily - if err := secondary.Put(ctx, path, data); err != nil { - // attempt fallback to primary if secondary fails - d.logger.Errorf("secondary put failed: %v", err) - if err2 := d.primary.Put(ctx, path, data); err2 != nil { - d.logger.Errorf("primary fallback put also failed: %v", err2) - return err2 - } - return err - } - // optional: mirror to primary async (disabled by default) - return nil - case MigrateModeFromSecondary: - // primary is main - return d.primary.Put(ctx, path, data) - default: - return d.primary.Put(ctx, path, data) - } -} - -func (d *MigrationBackendService) Get(ctx context.Context, path string) ([]byte, error) { - d.mu.RLock() - mode := d.mode - secondary := d.secondary - d.mu.RUnlock() - - // If migrating to secondary, prefer secondary (new writes go there). - if mode == MigrateModeToSecondary && secondary != nil { - if data, err := secondary.Get(ctx, path); err == nil { - return data, nil - } - // fallback to primary - } - // otherwise primary first - if data, err := d.primary.Get(ctx, path); err == nil { - return data, nil - } - if secondary != nil { - return secondary.Get(ctx, path) - } - return nil, fmt.Errorf("not found") -} - -func (d *MigrationBackendService) Delete(ctx context.Context, path string) error { - // best-effort delete in both - _ = d.primary.Delete(ctx, path) - if d.secondary != nil { - _ = d.secondary.Delete(ctx, path) - } - return nil -} - -// func (d *MigrationBackendService) List(ctx context.Context, prefix string) ([]string, error) { -// d.mu.RLock() -// mode := d.mode -// secondary := d.secondary -// d.mu.RUnlock() - -// // If migrating to secondary, prefer secondary (new writes go there). -// if mode == MigrateModeToSecondary && secondary != nil { -// if paths, err := secondary.List(ctx, prefix); err == nil { -// return paths, nil -// } -// // fallback to primary -// } -// // otherwise primary first -// if paths, err := d.primary.List(ctx, prefix); err == nil { -// return paths, nil -// } -// if secondary != nil { -// return secondary.List(ctx, prefix) -// } -// return nil, fmt.Errorf("not found") -// } - -func (d *MigrationBackendService) loadState(prefix string) (*migrationState, error) { - d.checkpointMux.Lock() - defer d.checkpointMux.Unlock() - // if file exists, read and unmarshal; else create new - if _, err := os.Stat(d.statePath); err == nil { - b, err := os.ReadFile(d.statePath) - if err != nil { - return nil, err - } - var st migrationState - if err := json.Unmarshal(b, &st); err != nil { - return nil, err - } - // if prefix changed, create new state - if st.Prefix != prefix { - st = migrationState{Prefix: prefix, Processed: map[string]bool{}, StartedAt: time.Now()} - } - return &st, nil - } - st := &migrationState{Prefix: prefix, Processed: map[string]bool{}, StartedAt: time.Now()} - return st, nil -} - -func (d *MigrationBackendService) persistState() error { - d.checkpointMux.Lock() - defer d.checkpointMux.Unlock() - if d.state == nil { - return nil - } - d.state.UpdatedAt = time.Now() - b, err := json.MarshalIndent(d.state, "", " ") - if err != nil { - return err - } - return os.WriteFile(d.statePath, b, 0644) -} - -func (d *MigrationBackendService) MigrateTo(ctx context.Context, prefix string, onProgress func(file string), onError func(file string, err error)) error { - d.mu.RLock() - secondary := d.secondary - cfg := d.cfg - d.mu.RUnlock() - if secondary == nil { - return errors.New("no secondary configured") - } - // ensure single migration at a time - if d.active.Load() { - return errors.New("migration already active") - } - - // create cancellable ctx - cctx, cancel := context.WithCancel(ctx) - d.cancelMigration = cancel - d.active.Store(true) - defer func() { - d.active.Store(false) - cancel() - }() - - // load checkpoint - st, err := d.loadState(prefix) - if err != nil { - return err - } - d.state = st - - // list all objects under prefix using primary.List - paths, err := d.primary.List(cctx, prefix) - if err != nil { - return err - } - d.state.Total = int64(len(paths)) - - jobCh := make(chan string, 1024) - errCh := make(chan error, 1) - - // spawn workers - var wg sync.WaitGroup - for i := 0; i < cfg.Concurrency; i++ { - wg.Add(1) - go func() { - defer wg.Done() - for { - select { - case <-cctx.Done(): - return - case p, ok := <-jobCh: - if !ok { - return - } - // skip if already processed - if d.isProcessed(p) { - continue - } - if err := d.migrateOneWithRetries(cctx, p, secondary, cfg.RetryCount, cfg.RetryBackoff, onError); err != nil { - // report but continue - if onError != nil { - onError(p, err) - } - continue - } - // mark processed and persist - d.markProcessed(p) - if onProgress != nil { - onProgress(p) - } - // optional delete source - if cfg.DeleteAfterCopy { - _ = d.primary.Delete(cctx, p) - } - } - } - }() - } - - // feed jobs -FeedLoop: - for _, p := range paths { - select { - case <-cctx.Done(): - break FeedLoop - default: - } - if d.isProcessed(p) { - atomic.AddInt64(&d.state.Completed, 1) - continue - } - jobCh <- p - } - close(jobCh) - close(errCh) - - // wait for workers - wg.Wait() - - // persist final state - _ = d.persistState() - return nil -} - -func (d *MigrationBackendService) isProcessed(path string) bool { - d.checkpointMux.Lock() - defer d.checkpointMux.Unlock() - if d.state == nil || d.state.Processed == nil { - return false - } - return d.state.Processed[path] -} - -func (d *MigrationBackendService) markProcessed(path string) { - d.checkpointMux.Lock() - defer d.checkpointMux.Unlock() - if d.state == nil { - d.state = &migrationState{Processed: map[string]bool{}} - } - if d.state.Processed == nil { - d.state.Processed = map[string]bool{} - } - if !d.state.Processed[path] { - d.state.Processed[path] = true - d.state.Completed++ - } - // persist periodically (you may want to batch) - _ = d.persistState() -} - -func (d *MigrationBackendService) migrateOneWithRetries(ctx context.Context, path string, target domain.FileBackend, retries int, backoff time.Duration, onError func(string, error)) error { - var last error - for i := 0; i <= retries; i++ { - if i > 0 { - time.Sleep(backoff * time.Duration(i)) - } - // read from primary - data, err := d.primary.Get(ctx, path) - if err != nil { - last = err - continue - } - // write to target - if err := target.Put(ctx, path, data); err != nil { - last = err - continue - } - // success - return nil - } - if onError != nil { - onError(path, last) - } - return last -} - -func (d *MigrationBackendService) MigrationStatus(ctx context.Context) (completed int64, total int64) { - if d.state == nil { - return 0, 0 - } - return d.state.Completed, d.state.Total -} - -func (d *MigrationBackendService) IsMigrating() bool { - return d.active.Load() -} diff --git a/internal/mocks/backend.go b/internal/mocks/backend.go index 239df35..b983933 100644 --- a/internal/mocks/backend.go +++ b/internal/mocks/backend.go @@ -75,6 +75,6 @@ func (m *MigrationBackend) MigrateFile(ctx context.Context, path string) error { } // MigrationStatus implements domain.MigratableBackend. -func (m *MigrationBackend) MigrationStatus(ctx context.Context) (completed int64, total int64) { - return 0, 0 +func (m *MigrationBackend) MigrationStatus(ctx context.Context) (completed int64, failed int64, total int64) { + return 0, 0, 0 } From de756f54ea72776914542ad7db3d99f5c7661ad3 Mon Sep 17 00:00:00 2001 From: Rhaqim Date: Tue, 11 Aug 2026 23:19:35 +0100 Subject: [PATCH 2/2] migration state and resumption --- CHANGELOG.md | 6 +++ README.md | 2 + buckt.go | 21 +++++++-- buckt_migration_test.go | 64 +++++++++++++++++++++++++++ buckt_test.go | 10 ++--- internal/backend/migrate.go | 76 ++++++++++++++++++++++++++------ internal/domain/backend.go | 15 +++++++ internal/repository/migration.go | 41 +++++++++++++++++ 8 files changed, 212 insertions(+), 23 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 7a3df16..890e1b1 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -42,6 +42,12 @@ Additive; no API removals. longer hangs. The new `Client.MigrationFailures(ctx)` reports how many objects could not be copied (re-run `MigrateAll` to retry them); the web `/backend` endpoint gained a `failed` field. +- **`MigrateAll` is now resumable.** Each copied object is recorded in the + database (`buckt_migration_models`), so a migration interrupted by a + restart/crash resumes from where it left off โ€” already-copied files are skipped + straight from the persisted state without re-scanning the target, and progress + no longer restarts from zero. Persistence is best-effort (a recording failure + never fails an idempotent copy) and keyed by the target backend's name. ### ๐Ÿงน Removed diff --git a/README.md b/README.md index 679fbbb..319e0ad 100644 --- a/README.md +++ b/README.md @@ -405,6 +405,8 @@ Migration is **always forward** (primary โ†’ secondary). The primary is treated Each in-flight file is buffered in full, so higher concurrency trades memory and provider rate-limit headroom for throughput. +**Resumable.** In migration mode buckt records each copied object in the database (the `buckt_migration_models` table). If the process stops mid-migration, the next `MigrateAll` **resumes from where it left off** โ€” already-copied files are skipped straight from the persisted state, without re-scanning the target. Progress (`MigrationStatus`) reflects the recorded state, so a restarted migration doesn't start its count from zero. Persistence is best-effort: a recording failure never fails a copy (copies are idempotent), and the state is keyed by the target backend's name. + `client.BackendName()` reports the active backend (`"local"`, `"s3"`, or `"local->s3"` mid-migration) โ€” handy for a status badge. When you're done migrating, swap to the target alone: ```go diff --git a/buckt.go b/buckt.go index 0b3962b..6fbcd6b 100644 --- a/buckt.go +++ b/buckt.go @@ -92,8 +92,14 @@ func New(conf Config, opts ...ConfigFunc) (*Client, error) { // Initialize cache cacheManager, lruCache := initializeCache(conf.Cache, bucktLog) - // Initialise Backend - backend := resolveBackend(conf.MediaDir, conf.Backend, bucktLog, lruCache, conf.Metrics) + // Initialise Backend. In migration mode, back the resumable bulk copy with a + // DB-persisted state store (the MigrationModel table, created by db.Migrate + // above) so an interrupted MigrateAll resumes without re-scanning the target. + var migrationStore domain.MigrationStateStore + if conf.Backend.MigrationEnabled { + migrationStore = repository.NewMigrationStateStore(db.DB) + } + backend := resolveBackend(conf.MediaDir, conf.Backend, bucktLog, lruCache, conf.Metrics, migrationStore) // Max file size: 0 means no limit (backward compatible) maxFileSize := conf.MaxFileSize @@ -1052,7 +1058,14 @@ func newAppServices( return folderService, fileService } -func resolveBackend(mediaDir string, bc BackendConfig, log domain.BucktLogger, lru domain.LRUCache, rec metrics.Recorder) Backend { +func resolveBackend( + mediaDir string, + bc BackendConfig, + log domain.BucktLogger, + lru domain.LRUCache, + rec metrics.Recorder, + migrationStore domain.MigrationStateStore, +) Backend { // meter wraps a leaf backend so every operation is recorded. It is a no-op // when rec is nil. In migration mode the leaves (source/target) are metered // individually โ€” never the composite โ€” so per-backend counts stay distinct @@ -1100,7 +1113,7 @@ func resolveBackend(mediaDir string, bc BackendConfig, log domain.BucktLogger, l } log.Infof("๐Ÿ”„ Migration mode: %s โ†’ %s", source.Name(), target.Name()) - return backend.NewMigrationBackend(log, meter(source), meter(target), bc.MigrationConcurrency) + return backend.NewMigrationBackend(log, meter(source), meter(target), bc.MigrationConcurrency, migrationStore) } // Non-migration modes diff --git a/buckt_migration_test.go b/buckt_migration_test.go index fee6982..97410bd 100644 --- a/buckt_migration_test.go +++ b/buckt_migration_test.go @@ -366,6 +366,70 @@ func TestMigration_PermanentFailureStillCompletes(t *testing.T) { assert.Equal(t, 0, target.count(), "nothing landed in the target") } +// TestMigration_ResumesFromPersistedState proves the DB-backed resume: after a +// migration records its progress, a restart (new client on the same DB) skips +// the already-copied files WITHOUT re-scanning the target. The second run uses a +// brand-new empty target that counts Put calls โ€” a zero count proves the skip +// was driven by the persisted state, not by listing the target. +func TestMigration_ResumesFromPersistedState(t *testing.T) { + dir := t.TempDir() + mediaDir := filepath.Join(dir, "media") + dbPath := filepath.Join(dir, "db.sqlite") + const user = "u1" + + // Phase 1 โ€” local-only files. + db1, err := sql.Open("sqlite3", dbPath) + require.NoError(t, err) + c1, err := New(Config{DB: DBConfig{Driver: SQLite, Database: db1}, MediaDir: mediaDir, Log: LogConfig{Silence: true}}) + require.NoError(t, err) + for _, n := range []string{"a.txt", "b.txt", "c.txt"} { + _, err := c1.UploadFile(user, "", n, "text/plain", []byte(n)) + require.NoError(t, err) + } + require.NoError(t, c1.Close()) + require.NoError(t, db1.Close()) + + // Phase 2 โ€” migrate to target1, which records the copies in the DB. + target1 := newMemBackend() + db2, err := sql.Open("sqlite3", dbPath) + require.NoError(t, err) + c2, err := New(Config{ + DB: DBConfig{Driver: SQLite, Database: db2}, + MediaDir: mediaDir, + Log: LogConfig{Silence: true}, + Backend: BackendConfig{Source: LocalBackend(), Target: target1, MigrationEnabled: true}, + }) + require.NoError(t, err) + require.NoError(t, c2.MigrateAll(context.Background())) + requireMigrationDone(t, c2) + require.Equal(t, 3, target1.count(), "first run copies everything") + require.NoError(t, c2.Close()) + require.NoError(t, db2.Close()) + + // Phase 3 โ€” "restart": same DB, a fresh EMPTY target. Because the persisted + // state already records all three keys as committed, MigrateAll must skip + // them without a single Put to the new target. + target2 := newMemBackend() + db3, err := sql.Open("sqlite3", dbPath) + require.NoError(t, err) + c3, err := New(Config{ + DB: DBConfig{Driver: SQLite, Database: db3}, + MediaDir: mediaDir, + Log: LogConfig{Silence: true}, + Backend: BackendConfig{Source: LocalBackend(), Target: target2, MigrationEnabled: true}, + }) + require.NoError(t, err) + t.Cleanup(func() { _ = c3.Close(); _ = db3.Close() }) + + require.NoError(t, c3.MigrateAll(context.Background())) + requireMigrationDone(t, c3) + + assert.Equal(t, 0, target2.putCalls(), "resume skips already-migrated files without re-copying") + done, total, _ := c3.MigrationStatus(context.Background()) + assert.Equal(t, int64(3), total) + assert.Equal(t, total, done, "progress still reaches total on resume") +} + func TestMigration_NotEnabledReturnsError(t *testing.T) { c := newTestClient(t) // plain local backend diff --git a/buckt_test.go b/buckt_test.go index d126284..c5d6915 100644 --- a/buckt_test.go +++ b/buckt_test.go @@ -636,7 +636,7 @@ func TestResolveBackend(t *testing.T) { Source: source, Target: target, } - result := resolveBackend(mediaDir, bc, mockLogger, mockLRU, nil) + result := resolveBackend(mediaDir, bc, mockLogger, mockLRU, nil, nil) _, ok := result.(*backend.MigrationBackendService) assert.True(t, ok) }) @@ -646,7 +646,7 @@ func TestResolveBackend(t *testing.T) { bc := BackendConfig{ Source: source, } - result := resolveBackend(mediaDir, bc, mockLogger, mockLRU, nil) + result := resolveBackend(mediaDir, bc, mockLogger, mockLRU, nil, nil) // Placeholder should be replaced with real local backend _, ok := result.(*backend.LocalFileSystemService) assert.True(t, ok) @@ -657,7 +657,7 @@ func TestResolveBackend(t *testing.T) { bc := BackendConfig{ Source: source, } - result := resolveBackend(mediaDir, bc, mockLogger, mockLRU, nil) + result := resolveBackend(mediaDir, bc, mockLogger, mockLRU, nil, nil) assert.Equal(t, source, result) }) @@ -666,14 +666,14 @@ func TestResolveBackend(t *testing.T) { bc := BackendConfig{ Target: target, } - result := resolveBackend(mediaDir, bc, mockLogger, mockLRU, nil) + result := resolveBackend(mediaDir, bc, mockLogger, mockLRU, nil, nil) _, ok := result.(*backend.LocalFileSystemService) assert.True(t, ok) }) t.Run("No Source or Target", func(t *testing.T) { bc := BackendConfig{} - result := resolveBackend(mediaDir, bc, mockLogger, mockLRU, nil) + result := resolveBackend(mediaDir, bc, mockLogger, mockLRU, nil, nil) _, ok := result.(*backend.LocalFileSystemService) assert.True(t, ok) }) diff --git a/internal/backend/migrate.go b/internal/backend/migrate.go index c7a6efd..7af9c34 100644 --- a/internal/backend/migrate.go +++ b/internal/backend/migrate.go @@ -21,6 +21,11 @@ type MigrationBackendService struct { // concurrency is how many files MigrateAll copies in parallel (>= 1). concurrency int + // store persists per-file migration progress so MigrateAll can resume after + // a restart without re-scanning the target. Nil disables persistence (the + // migration still works, relying on the target's own contents to skip). + store domain.MigrationStateStore + migrating atomic.Bool stats migrationStats } @@ -37,7 +42,7 @@ type migrationStats struct { // one, not the other way around. const defaultMigrationConcurrency = 8 -func NewMigrationBackend(bucktLogger domain.BucktLogger, primary domain.FileBackend, secondary domain.FileBackend, concurrency int) domain.MigratableBackend { +func NewMigrationBackend(bucktLogger domain.BucktLogger, primary domain.FileBackend, secondary domain.FileBackend, concurrency int, store domain.MigrationStateStore) domain.MigratableBackend { if concurrency <= 0 { concurrency = defaultMigrationConcurrency } @@ -47,6 +52,7 @@ func NewMigrationBackend(bucktLogger domain.BucktLogger, primary domain.FileBack primaryBackend: primary, secondaryBackend: secondary, concurrency: concurrency, + store: store, } } @@ -261,18 +267,38 @@ func (d *MigrationBackendService) MigrateAll(ctx context.Context) error { d.stats.completed = 0 d.stats.mu.Unlock() - // Build a skip-set from a single List of the secondary, instead of one - // Exists round-trip per file (halving the request count on large buckets). - // Keys are normalised with normaliseKey so both the leading-slash and - // slash-stripped forms of a key match โ€” see the cloud backends' altKey. - skip, haveSkipSet := d.secondarySkipSet(ctx) - - // alreadyInSecondary reports whether the secondary already holds path. It - // uses the pre-listed skip-set when available, and otherwise falls back to a - // per-file Exists check (e.g. when the secondary's List failed). - alreadyInSecondary := func(path string) bool { - if haveSkipSet { - _, ok := skip[normaliseKey(path)] + // Persisted progress (for resume): the set of keys a previous run already + // recorded as copied. On a restart this lets us skip completed files without + // re-scanning the target. Empty (and best-effort) when no store is wired. + persisted := map[string]struct{}{} + if d.store != nil { + if keys, err := d.store.MigratedKeys(ctx, d.secondaryBackend.Name()); err == nil { + persisted = keys + } else { + d.logger.Errorf("Could not load persisted migration state (resuming from scratch): %v", err) + } + } + + // Live skip-set from a single List of the secondary, instead of one Exists + // round-trip per file (halving the request count on large buckets). Keys are + // normalised with normaliseKey so both the leading-slash and slash-stripped + // forms of a key match โ€” see the cloud backends' altKey. + secondary, haveSecondary := d.secondarySkipSet(ctx) + + inPersisted := func(path string) bool { + _, ok := persisted[normaliseKey(path)] + return ok + } + + // alreadyDone reports whether path is already at the target โ€” recorded in the + // persisted state, present in the pre-listed secondary set, or (as a fallback + // when the secondary List failed) confirmed by a per-file Exists check. + alreadyDone := func(path string) bool { + if inPersisted(path) { + return true + } + if haveSecondary { + _, ok := secondary[normaliseKey(path)] return ok } exists, err := d.secondaryBackend.Exists(ctx, path) @@ -307,10 +333,17 @@ func (d *MigrationBackendService) MigrateAll(ctx context.Context) error { if ctx.Err() != nil { return // cancelled โ€” drain quietly } - if alreadyInSecondary(path) { + if alreadyDone(path) { d.stats.mu.Lock() d.stats.completed++ d.stats.mu.Unlock() + // Backfill the persisted state for files that are already + // in the secondary but not yet recorded (e.g. copied by an + // earlier version, or out of band), so a later resume can + // skip them without listing the target. + if !inPersisted(path) { + d.markMigrated(ctx, path) + } continue } if err := d.MigrateFile(ctx, path); err != nil { @@ -321,7 +354,10 @@ func (d *MigrationBackendService) MigrateAll(ctx context.Context) error { d.stats.mu.Lock() d.stats.failed++ d.stats.mu.Unlock() + continue } + // Record the successful copy so a restart resumes past it. + d.markMigrated(ctx, path) } }() } @@ -373,6 +409,18 @@ func normaliseKey(key string) string { return strings.TrimPrefix(key, "/") } +// markMigrated records, best-effort, that path is now at the secondary so a +// later run can resume past it. A persistence failure never fails the migration +// โ€” the copy already succeeded and is idempotent โ€” it is only logged. +func (d *MigrationBackendService) markMigrated(ctx context.Context, path string) { + if d.store == nil { + return + } + if err := d.store.MarkMigrated(ctx, d.secondaryBackend.Name(), normaliseKey(path), 0); err != nil { + d.logger.Errorf("Failed to persist migration state for %s: %v", path, err) + } +} + // MigrationStatus returns the current migration progress: files copied (or // already present), files that permanently failed after retries, and the total // scheduled. completed+failed == total once the run finishes. diff --git a/internal/domain/backend.go b/internal/domain/backend.go index c4445d6..707ef5e 100644 --- a/internal/domain/backend.go +++ b/internal/domain/backend.go @@ -52,6 +52,21 @@ type MigratableBackend interface { MigrationStatus(ctx context.Context) (completed int64, failed int64, total int64) } +// MigrationStateStore persists which object keys have been copied to a target +// backend, so a bulk migration can resume across a restart without re-scanning +// the target. Keyed by the target backend's Name(). Implementations must be +// safe for concurrent use. Persistence is best-effort from the migration's +// point of view: a failed record never fails a copy (the copy is idempotent, so +// the worst case is re-copying the object on the next run). +type MigrationStateStore interface { + // MigratedKeys returns the set of keys already committed to backend. + MigratedKeys(ctx context.Context, backend string) (map[string]struct{}, error) + + // MarkMigrated records that key (of the given size) was copied to backend. + // It is idempotent โ€” marking an already-recorded key is a no-op. + MarkMigrated(ctx context.Context, backend, key string, size int64) error +} + type PlaceholderBackend struct { Title string } diff --git a/internal/repository/migration.go b/internal/repository/migration.go index 94c4f09..afca438 100644 --- a/internal/repository/migration.go +++ b/internal/repository/migration.go @@ -1,6 +1,9 @@ package repository import ( + "context" + + "github.com/Rhaqim/buckt/internal/domain" "github.com/Rhaqim/buckt/internal/model" "github.com/google/uuid" "gorm.io/gorm" @@ -23,6 +26,44 @@ func NewMigrationRepository(db *gorm.DB) MigrationRepository { } } +// NewMigrationStateStore returns a domain.MigrationStateStore backed by the +// MigrationModel table, used to persist bulk-migration progress for resume. +func NewMigrationStateStore(db *gorm.DB) domain.MigrationStateStore { + return &migrationRepository{db: db} +} + +// MigratedKeys implements domain.MigrationStateStore. +func (repo *migrationRepository) MigratedKeys(ctx context.Context, backend string) (map[string]struct{}, error) { + var keys []string + err := repo.db.WithContext(ctx). + Model(&model.MigrationModel{}). + Where("backend = ? AND status = ?", backend, model.MigrationStatusCommitted). + Pluck("object_key", &keys).Error + if err != nil { + return nil, err + } + set := make(map[string]struct{}, len(keys)) + for _, k := range keys { + set[k] = struct{}{} + } + return set, nil +} + +// MarkMigrated implements domain.MigrationStateStore. It is idempotent: an +// existing (backend, key) row is left untouched rather than duplicated. +func (repo *migrationRepository) MarkMigrated(ctx context.Context, backend, key string, size int64) error { + var m model.MigrationModel + return repo.db.WithContext(ctx). + Where("backend = ? AND object_key = ?", backend, key). + Attrs(model.MigrationModel{ + Backend: model.MigrationBackend(backend), + ObjectKey: key, + Size: size, + Status: model.MigrationStatusCommitted, + }). + FirstOrCreate(&m).Error +} + func (repo *migrationRepository) CreateMigration(migration *model.MigrationModel) error { return repo.db.Create(migration).Error }