mirror of
https://github.com/vmware-tanzu/velero.git
synced 2026-09-25 01:14:18 +00:00
* Update CRDs and CLI to support in-place restore (#10038) Update CRDs(Restore, DataDownload, PodVolumeRestore) and restore create CLI to support in-place restore Signed-off-by: Wenkai Yin(尹文开) <yinw@vmware.com> * Update Kopia(filesystem) uploader to support incremental and deleteExtraFile during restore (#10066) Update Kopia(filesystem) uploader to support incremental and deleteExtraFile during restore Signed-off-by: Wenkai Yin(尹文开) <yinw@vmware.com> * Update Restore Exposer and PVC CSI to support in-place restore (#10104) 1. Update Restore Exposer to support exposing with existing PV for in-place restore 2. Update PVC CSI RIA to continue the restore process for in-place restore Signed-off-by: Wenkai Yin(尹文开) <yinw@vmware.com> * Update Block uploader to support increase restore (#10244) Update Block uploader to support increase restore Signed-off-by: Wenkai Yin(尹文开) <yinw@vmware.com> * Update Exposer to recreate the target PV if the volume mode is different with the restore PVC (#10257) Update Exposer to recreate the target PV if the volume mode is different with t he restore PVC Signed-off-by: Wenkai Yin(尹文开) <yinw@vmware.com> * Preserve PVC selected-node annotation via carrier annotation for in-place restore For in-place volume data restore, the existing PVC is deleted and recreated. For StorageClasses with the WaitForFirstConsumer volume binding mode, losing the volume.kubernetes.io/selected-node annotation could let the scheduler place the recreated workload Pod in a different zone than the original PV, leaving it stuck in ContainerCreating. Instead of relying on RestoreItemAction execution order (the generic PVC RIA unconditionally strips the selected-node annotation), the PVC CSI RIA now captures the annotation from the existing PVC right before deleting it and carries it on the target PVC via the Velero-internal restore.velero.io/inplace-restore-selected-node annotation. The restore engine translates the carrier back to the Kubernetes annotation after all RestoreItemActions have run and always strips the carrier so it never lands on the cluster. This makes the behavior independent of RIA ordering: the Kubernetes annotation is stripped by default on every path (including when the target PVC does not exist and Velero falls back to provisioning a new PVC), and preservation only happens when the CSI RIA explicitly captured a value from the existing PVC. Signed-off-by: chlins <chlins.zhang@gmail.com> * Update the control path to make the in-place incremental restore with block data mover work E2E (#10410) Update the control path to make the in-place incremental restore with block data mover work E2E Signed-off-by: Wenkai Yin(尹文开) <yinw@vmware.com> --------- Signed-off-by: Wenkai Yin(尹文开) <yinw@vmware.com> Signed-off-by: chlins <chlins.zhang@gmail.com> Co-authored-by: chlins <chlins.zhang@gmail.com>
323 lines
11 KiB
Go
323 lines
11 KiB
Go
/*
|
|
Copyright The Velero Contributors.
|
|
|
|
Licensed under the Apache License, Version 2.0 (the "License");
|
|
you may not use this file except in compliance with the License.
|
|
You may obtain a copy of the License at
|
|
|
|
http://www.apache.org/licenses/LICENSE-2.0
|
|
|
|
Unless required by applicable law or agreed to in writing, software
|
|
distributed under the License is distributed on an "AS IS" BASIS,
|
|
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
|
See the License for the specific language governing permissions and
|
|
limitations under the License.
|
|
*/
|
|
|
|
package block
|
|
|
|
import (
|
|
"context"
|
|
"io"
|
|
"maps"
|
|
"path/filepath"
|
|
"time"
|
|
|
|
"github.com/cockroachdb/errors"
|
|
"github.com/sirupsen/logrus"
|
|
|
|
"github.com/vmware-tanzu/velero/pkg/cbtservice"
|
|
"github.com/vmware-tanzu/velero/pkg/repository/udmrepo"
|
|
"github.com/vmware-tanzu/velero/pkg/uploader"
|
|
"github.com/vmware-tanzu/velero/pkg/uploader/cbt"
|
|
)
|
|
|
|
var openBlockDeviceFunc = openBlockDevice
|
|
|
|
type parentBackupInfo struct {
|
|
parentObject udmrepo.ID
|
|
changeID string
|
|
volumeID string
|
|
}
|
|
|
|
// Backup backup specific sourcePath and update progress
|
|
func Backup(ctx context.Context, blkUp Uploader, repoWriter udmrepo.BackupRepo, sourcePath string, realSource string, cbtSource cbtservice.SourceInfo,
|
|
forceFull bool, parentSnapshot string, cbtService cbtservice.Service, uploaderCfg map[string]string, tags map[string]string, log logrus.FieldLogger) (uploader.SnapshotInfo, bool, error) {
|
|
if blkUp == nil {
|
|
return uploader.SnapshotInfo{}, false, errors.New("get empty block uploader")
|
|
}
|
|
|
|
source, err := filepath.Abs(sourcePath)
|
|
if err != nil {
|
|
return uploader.SnapshotInfo{}, false, errors.Wrapf(err, "invalid source path %s", sourcePath)
|
|
}
|
|
|
|
source = filepath.Clean(source)
|
|
|
|
sourceInfo := sourceInfo{
|
|
realSource: filepath.Clean(realSource),
|
|
}
|
|
|
|
if realSource == "" {
|
|
sourceInfo.realSource = source
|
|
}
|
|
|
|
sourceInfo.dev, err = openBlockDeviceFunc(source, true)
|
|
if err != nil {
|
|
return uploader.SnapshotInfo{}, false, errors.Wrapf(err, "error opening block device %s", source)
|
|
}
|
|
|
|
defer sourceInfo.dev.Close()
|
|
|
|
sourceInfo.size, err = sourceInfo.dev.Seek(0, io.SeekEnd)
|
|
if err != nil {
|
|
return uploader.SnapshotInfo{}, false, errors.Wrapf(err, "error getting length of block device %s", source)
|
|
}
|
|
|
|
_, err = sourceInfo.dev.Seek(0, io.SeekStart)
|
|
if err != nil {
|
|
return uploader.SnapshotInfo{}, false, errors.Wrapf(err, "error reset pos of block device %s", source)
|
|
}
|
|
|
|
snapID, backupSize, err := snapshotSource(ctx, repoWriter, blkUp, sourceInfo, forceFull, parentSnapshot, cbtSource, cbtService, tags, uploaderCfg, log, "Block Uploader")
|
|
snapshotInfo := uploader.SnapshotInfo{
|
|
ID: snapID,
|
|
Size: sourceInfo.size,
|
|
IncrementalSize: backupSize,
|
|
}
|
|
|
|
return snapshotInfo, false, err
|
|
}
|
|
|
|
func snapshotSource(
|
|
ctx context.Context,
|
|
rep udmrepo.BackupRepo,
|
|
u Uploader,
|
|
source sourceInfo,
|
|
forceFull bool,
|
|
parentSnapshot string,
|
|
cbtSource cbtservice.SourceInfo,
|
|
cbtService cbtservice.Service,
|
|
snapshotTags map[string]string,
|
|
uploaderCfg map[string]string,
|
|
log logrus.FieldLogger,
|
|
description string,
|
|
) (string, int64, error) {
|
|
log.Info("Start to snapshot...")
|
|
snapshotStartTime := time.Now()
|
|
|
|
parentBackup := getParentBackupInfo(ctx, rep, forceFull, parentSnapshot, cbtSource.VolumeID, source.realSource, snapshotTags, log)
|
|
|
|
bitmap := cbt.NewBitmap(blockSize, uint64(source.size), cbtSource.Snapshot, parentBackup.changeID, parentBackup.volumeID)
|
|
|
|
err := cbt.SetBitmapOrFull(ctx, cbtService, bitmap)
|
|
if err != nil {
|
|
parentBackup.parentObject = ""
|
|
log.WithError(err).Warnf("Failed to create CBT with source %v, fallback to real full backup", cbtSource)
|
|
}
|
|
|
|
snap, backupSize, err := u.Backup(source, parentBackup.parentObject, bitmap.Iterator(), uploaderCfg)
|
|
if err != nil {
|
|
return "", 0, errors.Wrapf(err, "Failed to run uploader backup for si %v", source)
|
|
}
|
|
|
|
if snap.Tags == nil {
|
|
snap.Tags = make(map[string]string)
|
|
}
|
|
|
|
snap.Tags[uploader.CBTChangeIDTag] = cbtSource.ChangeID
|
|
snap.Tags[uploader.CBTVolumeIDTag] = cbtSource.VolumeID
|
|
if snapshotTags != nil {
|
|
maps.Copy(snap.Tags, snapshotTags)
|
|
}
|
|
|
|
snap.Description = description
|
|
|
|
snapID, err := rep.SaveSnapshot(ctx, snap)
|
|
if err != nil {
|
|
return "", 0, errors.Wrapf(err, "Failed to save snapshot %v", snap)
|
|
}
|
|
|
|
if err = rep.Flush(ctx); err != nil {
|
|
return "", 0, errors.Wrapf(err, "Failed to flush repository")
|
|
}
|
|
|
|
log.Infof("Created snapshot with root %v and ID %v in %v", snap.RootObject, snapID, time.Since(snapshotStartTime).Truncate(time.Second))
|
|
|
|
return string(snapID), backupSize, nil
|
|
}
|
|
|
|
func getParentBackupInfo(ctx context.Context, rep udmrepo.BackupRepo, forceFull bool, parentSnapshot string, volumeID string, realSource string, snapshotTags map[string]string, log logrus.FieldLogger) parentBackupInfo {
|
|
var previous *udmrepo.Snapshot
|
|
|
|
// parentID names whichever snapshot ended up being the parent. On the discovery
|
|
// branch the parentSnapshot parameter is empty by definition, so logging it there
|
|
// produces messages that describe a decision without naming the object it was about.
|
|
parentID := parentSnapshot
|
|
|
|
if !forceFull {
|
|
if parentSnapshot != "" {
|
|
snap, err := rep.GetSnapshot(ctx, udmrepo.ID(parentSnapshot))
|
|
if err != nil {
|
|
log.WithError(err).Warn("Failed to load previous snapshot, fallback to full backup")
|
|
} else {
|
|
previous = &snap
|
|
log.Infof("Using provided parent snapshot %s", parentSnapshot)
|
|
}
|
|
} else {
|
|
log.Infof("Searching for parent snapshot")
|
|
|
|
snap, err := findPreviousSnapshot(ctx, rep, realSource, snapshotTags, nil, log)
|
|
if err != nil {
|
|
log.WithError(err).Warn("Failed to search previous snapshot, fallback to full backup")
|
|
} else {
|
|
previous = &snap
|
|
parentID = string(snap.RootObject.ID)
|
|
log.Infof("Using previous snapshot %s", snap.RootObject.ID)
|
|
}
|
|
}
|
|
} else {
|
|
log.Info("Forcing full snapshot")
|
|
}
|
|
|
|
parentInfo := parentBackupInfo{}
|
|
if previous != nil {
|
|
if previous.Tags == nil {
|
|
log.Warnf("No tag from parent snapshot %s, fallback to full backup", parentID)
|
|
} else if previous.Tags[uploader.CBTChangeIDTag] == "" {
|
|
log.Warnf("No ChangeID tag from parent snapshot %s, fallback to full backup", parentID)
|
|
} else if previous.Tags[uploader.CBTVolumeIDTag] == "" {
|
|
log.Warnf("No VolumeID tag from parent snapshot %s, fallback to full backup", parentID)
|
|
} else if previous.Tags[uploader.CBTVolumeIDTag] != volumeID {
|
|
log.Warnf("VolumeID %s from parent snapshot %s is not expected as %s, fallback to full backup", previous.Tags[uploader.CBTVolumeIDTag], parentID, volumeID)
|
|
} else if obj, err := loadObjectFromSnapshot(ctx, rep, previous); err != nil {
|
|
log.WithError(err).Warnf("Failed to load object from parent snapshot %s, fallback to full backup", parentID)
|
|
} else {
|
|
parentInfo.parentObject = obj
|
|
parentInfo.changeID = previous.Tags[uploader.CBTChangeIDTag]
|
|
parentInfo.volumeID = previous.Tags[uploader.CBTVolumeIDTag]
|
|
|
|
log.Infof("Using parent snapshot %s, start time %v, end time %v, description %s", parentID, previous.StartTime, previous.EndTime, previous.Description)
|
|
}
|
|
}
|
|
|
|
return parentInfo
|
|
}
|
|
|
|
// Restore restore specific sourcePath with given snapshotID and update progress
|
|
func Restore(ctx context.Context, blkUp Uploader, rep udmrepo.BackupRepo, snapshotID, dest string, incremental bool, cbtSource cbtservice.SourceInfo, cbtService cbtservice.Service, uploaderCfg map[string]string, log logrus.FieldLogger) (int64, error) {
|
|
log.Info("Start to restore...")
|
|
|
|
snapshot, err := rep.GetSnapshot(ctx, udmrepo.ID(snapshotID))
|
|
if err != nil {
|
|
return 0, errors.Wrapf(err, "Unable to load snapshot %v", snapshotID)
|
|
}
|
|
log.Infof("Restore from snapshot %s, incremental %v, cbt source %v, description %s, created time %v, tags %v", snapshotID, incremental, cbtSource, snapshot.Description, snapshot.EndTime, snapshot.Tags)
|
|
|
|
var volumeSnapshot, changeID, volumeID string
|
|
if incremental {
|
|
if snapshot.Tags == nil {
|
|
log.Warnf("No tag from snapshot %s, fallback to full restore", snapshotID)
|
|
incremental = false
|
|
} else if snapshot.Tags[uploader.CBTChangeIDTag] == "" {
|
|
log.Warnf("No ChangeID tag from snapshot %s, fallback to full restore", snapshotID)
|
|
incremental = false
|
|
} else if snapshot.Tags[uploader.CBTVolumeIDTag] == "" {
|
|
log.Warnf("No VolumeID tag from snapshot %s, fallback to full restore", snapshotID)
|
|
incremental = false
|
|
} else if snapshot.Tags[uploader.CBTVolumeIDTag] != cbtSource.VolumeID {
|
|
log.Warnf("VolumeID %s from snapshot %s is not expected as %s, fallback to full restore", snapshot.Tags[uploader.CBTVolumeIDTag], snapshotID, cbtSource.VolumeID)
|
|
incremental = false
|
|
} else {
|
|
volumeSnapshot = cbtSource.Snapshot
|
|
changeID = snapshot.Tags[uploader.CBTChangeIDTag]
|
|
volumeID = snapshot.Tags[uploader.CBTVolumeIDTag]
|
|
}
|
|
}
|
|
|
|
bitmap := cbt.NewBitmap(blockSize, uint64(snapshot.TotalSize), volumeSnapshot, changeID, volumeID)
|
|
if incremental {
|
|
if err = cbt.SetBitmapOrFull(ctx, cbtService, bitmap); err != nil {
|
|
log.WithError(err).Warnf("Failed to create CBT with source %v, fallback to full restore", cbtSource)
|
|
}
|
|
} else {
|
|
bitmap.SetFull()
|
|
}
|
|
|
|
destPath, err := filepath.Abs(dest)
|
|
if err != nil {
|
|
return 0, errors.Wrapf(err, "invalid dest path '%s'", dest)
|
|
}
|
|
|
|
destPath = filepath.Clean(destPath)
|
|
|
|
destDev, err := openBlockDeviceFunc(destPath, false)
|
|
if err != nil {
|
|
return 0, errors.Wrapf(err, "error opening block device '%s'", destPath)
|
|
}
|
|
|
|
defer destDev.Close()
|
|
|
|
destSize, err := destDev.Seek(0, io.SeekEnd)
|
|
if err != nil {
|
|
return 0, errors.Wrapf(err, "error getting length of block device %s", dest)
|
|
}
|
|
|
|
_, err = destDev.Seek(0, io.SeekStart)
|
|
if err != nil {
|
|
return 0, errors.Wrapf(err, "error reset pos of block device %s", dest)
|
|
}
|
|
|
|
_, totalSize, err := blkUp.Restore(snapshot, destInfo{dev: destDev, path: destPath, size: destSize}, bitmap.Iterator(), uploaderCfg)
|
|
if err != nil {
|
|
return 0, errors.Wrapf(err, "error restoring to block dev %s", destPath)
|
|
}
|
|
|
|
return totalSize, nil
|
|
}
|
|
|
|
func findPreviousSnapshot(ctx context.Context, rep udmrepo.BackupRepo, path string, snapshotTags map[string]string, noLaterThan *time.Time, log logrus.FieldLogger) (udmrepo.Snapshot, error) {
|
|
snaps, err := rep.ListSnapshot(ctx, path)
|
|
if err != nil {
|
|
return udmrepo.Snapshot{}, errors.Wrapf(err, "error list snapshots for %s", path)
|
|
}
|
|
|
|
var previous *udmrepo.Snapshot
|
|
|
|
for _, snap := range snaps {
|
|
log.Debugf("Found one snapshot %s, start time %v, tags %v", snap.RootObject.ID, snap.StartTime, snap.Tags)
|
|
|
|
requester, found := snap.Tags[uploader.SnapshotRequesterTag]
|
|
if !found {
|
|
continue
|
|
}
|
|
|
|
if requester != snapshotTags[uploader.SnapshotRequesterTag] {
|
|
continue
|
|
}
|
|
|
|
uploaderName, found := snap.Tags[uploader.SnapshotUploaderTag]
|
|
if !found {
|
|
continue
|
|
}
|
|
|
|
if uploaderName != uploader.BlockType {
|
|
continue
|
|
}
|
|
|
|
if noLaterThan != nil && snap.StartTime.After(*noLaterThan) {
|
|
continue
|
|
}
|
|
|
|
if previous == nil || snap.StartTime.After(previous.StartTime) {
|
|
previous = &snap
|
|
}
|
|
}
|
|
|
|
if previous == nil {
|
|
return udmrepo.Snapshot{}, errors.Errorf("no matching snapshot found for source %s", path)
|
|
}
|
|
|
|
return *previous, nil
|
|
}
|