Skip to content
Merged
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
30 changes: 30 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,36 @@ 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.
- **`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

- 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
Expand Down
24 changes: 23 additions & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -385,8 +385,28 @@ 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.

**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
Expand All @@ -397,7 +417,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`) |

---
Expand Down Expand Up @@ -756,6 +777,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
```

Expand Down
50 changes: 40 additions & 10 deletions buckt.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -947,17 +953,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 */
Expand Down Expand Up @@ -1035,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
Expand Down Expand Up @@ -1083,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))
return backend.NewMigrationBackend(log, meter(source), meter(target), bc.MigrationConcurrency, migrationStore)
}

// Non-migration modes
Expand Down
17 changes: 17 additions & 0 deletions buckt_conf.go
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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
}
}

Expand Down
Loading
Loading