Skip to content

Commit f66c74a

Browse files
authored
Merge pull request #1344 from jcechace/PBM-1776-gcs-parallel-upload
PBM-1776 GCS parallel upload (with storage close)
2 parents 3b00391 + fb06a23 commit f66c74a

26 files changed

Lines changed: 399 additions & 131 deletions

File tree

cmd/pbm-agent/agent.go

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -460,6 +460,7 @@ func (a *Agent) storStatus(
460460
if err != nil {
461461
return topo.SubsysStatus{Err: fmt.Sprintf("unable to get storage: %v", err)}
462462
}
463+
defer storage.Close(stg, log)
463464

464465
ok, err := storage.IsInitialized(ctx, stg)
465466
if err != nil {

cmd/pbm-agent/delete.go

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -98,6 +98,7 @@ func (a *Agent) Delete(ctx context.Context, d *ctrl.DeleteBackupCmd, opid ctrl.O
9898
l.Error("get storage: %v", err)
9999
return
100100
}
101+
defer storage.Close(stg, l)
101102
l.Info("deleting backups older than %v %s", t, util.LogProfileArg(d.Profile))
102103
err = backup.DeleteBackupBefore(ctx, a.leadConn, stg, d.Profile, bcpType, t)
103104
if err != nil {
@@ -271,6 +272,7 @@ func (a *Agent) Cleanup(ctx context.Context, d *ctrl.CleanupCmd, opid ctrl.OPID,
271272
l.Error("get storage: " + err.Error())
272273
return
273274
}
275+
defer storage.Close(stg, l)
274276

275277
cr, err := backup.MakeCleanupInfo(ctx, a.leadConn, d.OlderThan, d.Profile)
276278
if err != nil {
@@ -308,6 +310,7 @@ func (a *Agent) deletePITRImpl(ctx context.Context, ts bson.Timestamp) error {
308310
if err != nil {
309311
return errors.Wrap(err, "get storage")
310312
}
313+
defer storage.Close(stg, l)
311314

312315
eg := &errgroup.Group{}
313316
eg.SetLimit(runtime.NumCPU())

cmd/pbm-agent/pitr.go

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -20,6 +20,7 @@ import (
2020
"github.com/percona/percona-backup-mongodb/pbm/oplog"
2121
"github.com/percona/percona-backup-mongodb/pbm/prio"
2222
"github.com/percona/percona-backup-mongodb/pbm/slicer"
23+
"github.com/percona/percona-backup-mongodb/pbm/storage"
2324
"github.com/percona/percona-backup-mongodb/pbm/topo"
2425
"github.com/percona/percona-backup-mongodb/pbm/util"
2526
)
@@ -366,6 +367,7 @@ func (a *Agent) pitr(ctx context.Context) error {
366367
err = s.Catchup(ctx)
367368
}
368369
if err != nil {
370+
storage.Close(stg, l)
369371
if err := lck.Release(); err != nil {
370372
l.Error("release lock: %v", err)
371373
}
@@ -375,6 +377,7 @@ func (a *Agent) pitr(ctx context.Context) error {
375377

376378
go func() {
377379
stopSlicingCtx, stopSlicing := context.WithCancel(ctx)
380+
defer storage.Close(stg, l)
378381
defer stopSlicing()
379382
stopC := make(chan struct{})
380383

cmd/pbm-agent/profile.go

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -91,6 +91,7 @@ func (a *Agent) handleAddConfigProfile(
9191
err = errors.Wrap(err, "storage from config")
9292
return
9393
}
94+
defer storage.Close(stg, l)
9495

9596
err = storage.HasReadAccess(ctx, stg)
9697
if err != nil {

cmd/pbm-speed-test/main.go

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,7 @@ import (
1515
"github.com/percona/percona-backup-mongodb/pbm/connect"
1616
"github.com/percona/percona-backup-mongodb/pbm/errors"
1717
"github.com/percona/percona-backup-mongodb/pbm/log"
18+
"github.com/percona/percona-backup-mongodb/pbm/storage"
1819
"github.com/percona/percona-backup-mongodb/pbm/storage/blackhole"
1920
"github.com/percona/percona-backup-mongodb/pbm/util"
2021
"github.com/percona/percona-backup-mongodb/pbm/version"
@@ -209,6 +210,7 @@ func testStorage(mURL string, compression compress.CompressionType, level *int,
209210
if err != nil {
210211
stdlog.Fatalln("Error: get storage:", err)
211212
}
213+
defer storage.Close(stg, nil)
212214
done := make(chan struct{})
213215
go printw(done)
214216
r, err := doTest(sess, stg, compression, level, sizeGb, collection)

cmd/pbm/backup.go

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -464,12 +464,15 @@ func describeBackup(
464464

465465
var stg storage.Storage
466466
if b.coll || bcp.Size == 0 {
467+
l := log.LogEventFromContext(ctx)
468+
467469
// to read backed up collection names
468470
// or calculate size of files for legacy backups
469-
stg, err = util.StorageFromConfig(&bcp.Store.StorageConf, node, log.LogEventFromContext(ctx))
471+
stg, err = util.StorageFromConfig(&bcp.Store.StorageConf, node, l)
470472
if err != nil {
471473
return nil, errors.Wrap(err, "get storage")
472474
}
475+
defer storage.Close(stg, l)
473476

474477
err = storage.HasReadAccess(ctx, stg)
475478
if err != nil && !errors.Is(err, storage.ErrUninitialized) {

cmd/pbm/oplog.go

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -11,6 +11,7 @@ import (
1111
"github.com/percona/percona-backup-mongodb/pbm/defs"
1212
"github.com/percona/percona-backup-mongodb/pbm/errors"
1313
"github.com/percona/percona-backup-mongodb/pbm/log"
14+
"github.com/percona/percona-backup-mongodb/pbm/storage"
1415
"github.com/percona/percona-backup-mongodb/pbm/util"
1516
"github.com/percona/percona-backup-mongodb/sdk"
1617
)
@@ -79,6 +80,7 @@ func replayOplog(
7980
if err != nil {
8081
return nil, errors.Wrap(err, "get storage")
8182
}
83+
defer storage.Close(stg, l)
8284

8385
name := time.Now().UTC().Format(time.RFC3339Nano)
8486
cmd := ctrl.Cmd{

cmd/pbm/restore.go

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -188,6 +188,7 @@ func runRestore(
188188
if err != nil {
189189
return nil, errors.Wrap(err, "get storage")
190190
}
191+
defer storage.Close(stg, l)
191192

192193
m, err := doRestore(
193194
ctx,
@@ -578,6 +579,7 @@ func runFinishRestore(o descrRestoreOpts, node string) (fmt.Stringer, error) {
578579
if err != nil {
579580
return nil, errors.Wrap(err, "get storage")
580581
}
582+
defer storage.Close(stg, nil)
581583

582584
path := fmt.Sprintf("%s/%s/cluster", defs.PhysRestoresDir, o.restore)
583585
msg := outMsg{"Command sent. Check `pbm describe-restore ...` for the result."}
@@ -779,8 +781,9 @@ func describeRestore(
779781
if err != nil {
780782
return nil, errors.Wrap(err, "get storage")
781783
}
782-
meta, err = restore.GetPhysRestoreMeta(o.restore, stg, log.New(nil, "cli", "").
783-
NewEvent("", "", "", bson.Timestamp{}))
784+
l := log.New(nil, "cli", "").NewEvent("", "", "", bson.Timestamp{})
785+
defer storage.Close(stg, l)
786+
meta, err = restore.GetPhysRestoreMeta(o.restore, stg, l)
784787
if err != nil && meta == nil {
785788
return nil, errors.Wrap(err, "get restore meta")
786789
}

cmd/pbm/status.go

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -603,11 +603,12 @@ func getStorageStat(
603603

604604
bcpsMatchCluster(bcps, ver.VersionString, fcv, shards, rsMap)
605605

606-
stg, err := util.GetStorage(ctx, conn, inf.Me,
607-
log.FromContext(ctx).NewEvent("", "", "", bson.Timestamp{}))
606+
l := log.FromContext(ctx).NewEvent("", "", "", bson.Timestamp{})
607+
stg, err := util.GetStorage(ctx, conn, inf.Me, l)
608608
if err != nil {
609609
return s, errors.Wrap(err, "get storage")
610610
}
611+
defer storage.Close(stg, l)
611612

612613
now, err := topo.GetClusterTime(ctx, conn)
613614
if err != nil {

e2e-tests/cmd/ensure-oplog/main.go

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -272,10 +272,12 @@ func ensureReplsetOplog(ctx context.Context, uri string, from, till bson.Timesta
272272
return errors.Wrap(err, "get config")
273273
}
274274

275-
stg, err := util.StorageFromConfig(&cfg.Storage, "", log.FromContext(ctx).NewDefaultEvent())
275+
l := log.FromContext(ctx).NewDefaultEvent()
276+
stg, err := util.StorageFromConfig(&cfg.Storage, "", l)
276277
if err != nil {
277278
return errors.Wrap(err, "get storage")
278279
}
280+
defer storage.Close(stg, l)
279281

280282
compression := defs.DefaultCompression
281283
compressionLevel := (*int)(nil)

0 commit comments

Comments
 (0)