Files
1Panel/agent/app/service/backup.go
ssongliu 1ab3da1fab fix: apply timeout to snapshot uploads (#13275)
* fix: apply timeout to snapshot uploads

* fix: honor snapshot upload timeouts for sftp and upyun
2026-07-16 11:21:26 +08:00

713 lines
22 KiB
Go

package service
import (
"bufio"
"context"
"encoding/base64"
"encoding/json"
"fmt"
"os"
"path"
"strconv"
"strings"
"time"
"github.com/1Panel-dev/1Panel/agent/app/dto"
"github.com/1Panel-dev/1Panel/agent/app/model"
"github.com/1Panel-dev/1Panel/agent/app/repo"
"github.com/1Panel-dev/1Panel/agent/app/task"
"github.com/1Panel-dev/1Panel/agent/buserr"
"github.com/1Panel-dev/1Panel/agent/constant"
"github.com/1Panel-dev/1Panel/agent/global"
"github.com/1Panel-dev/1Panel/agent/i18n"
"github.com/1Panel-dev/1Panel/agent/utils/cloud_storage"
"github.com/1Panel-dev/1Panel/agent/utils/cloud_storage/client"
"github.com/1Panel-dev/1Panel/agent/utils/encrypt"
"github.com/1Panel-dev/1Panel/agent/utils/files"
"github.com/jinzhu/copier"
)
type BackupService struct{}
type IBackupService interface {
CheckUsed(name string, isPublic bool) error
LoadBackupOptions() ([]dto.BackupOption, error)
SearchWithPage(search dto.SearchPageWithType) (int64, interface{}, error)
Create(backupDto dto.BackupOperate) error
CheckConn(req dto.BackupOperate) dto.BackupCheckRes
GetBuckets(backupDto dto.ForBuckets) ([]interface{}, error)
Update(req dto.BackupOperate) error
Delete(id uint) error
RefreshToken(req dto.OperateByID) error
GetLocalDir() (string, error)
UploadForRecover(req dto.UploadForRecover) error
MysqlBackup(db dto.CommonBackup) error
PostgresqlBackup(db dto.CommonBackup) error
MongodbBackup(db dto.CommonBackup) error
MysqlRecover(db dto.CommonRecover) error
PostgresqlRecover(db dto.CommonRecover) error
MongodbRecover(db dto.CommonRecover) error
MysqlRecoverByUpload(req dto.CommonRecover) error
PostgresqlRecoverByUpload(req dto.CommonRecover) error
MongodbRecoverByUpload(req dto.CommonRecover) error
RedisBackup(db dto.CommonBackup) error
RedisRecover(db dto.CommonRecover) error
WebsiteBackup(db dto.CommonBackup) error
WebsiteRecover(req dto.CommonRecover) error
AppBackup(db dto.CommonBackup) (*model.BackupRecord, error)
AppRecover(req dto.CommonRecover) error
ContainerBackup(req dto.CommonBackup) error
ContainerRecover(req dto.CommonRecover) error
ComposeBackup(req dto.CommonBackup) error
ComposeRecover(req dto.CommonRecover) error
}
func NewIBackupService() IBackupService {
return &BackupService{}
}
func (u *BackupService) GetLocalDir() (string, error) {
account, err := backupRepo.Get(repo.WithByType(constant.Local))
if err != nil {
return "", err
}
return account.BackupPath, nil
}
func (u *BackupService) SearchWithPage(req dto.SearchPageWithType) (int64, interface{}, error) {
options := []repo.DBOption{repo.WithOrderDesc("created_at")}
if len(req.Type) != 0 {
options = append(options, repo.WithByType(req.Type))
}
if len(req.Info) != 0 {
options = append(options, repo.WithByType(req.Info))
}
count, accounts, err := backupRepo.Page(req.Page, req.PageSize, options...)
if err != nil {
return 0, nil, err
}
var data []dto.BackupInfo
for _, account := range accounts {
var item dto.BackupInfo
if err := copier.Copy(&item, &account); err != nil {
global.LOG.Errorf("copy backup account to dto backup info failed, err: %v", err)
}
if item.Type != constant.Sftp && item.Type != constant.Local {
item.BackupPath = path.Join("/", strings.TrimPrefix(item.BackupPath, "/"))
}
if !item.RememberAuth {
item.AccessKey = ""
item.Credential = ""
if account.Type == constant.Sftp {
varMap := make(map[string]interface{})
if err := json.Unmarshal([]byte(item.Vars), &varMap); err != nil {
continue
}
delete(varMap, "passPhrase")
itemVars, _ := json.Marshal(varMap)
item.Vars = string(itemVars)
}
} else {
item.AccessKey, _ = encrypt.StringDecryptWithBase64(item.AccessKey)
item.Credential, _ = encrypt.StringDecryptWithBase64(item.Credential)
}
if account.Type == constant.OneDrive || account.Type == constant.ALIYUN || account.Type == constant.GoogleDrive {
varMap := make(map[string]interface{})
if err := json.Unmarshal([]byte(item.Vars), &varMap); err != nil {
continue
}
delete(varMap, "refresh_token")
delete(varMap, "drive_id")
itemVars, _ := json.Marshal(varMap)
item.Vars = string(itemVars)
}
data = append(data, item)
}
return count, data, nil
}
func (u *BackupService) CheckConn(req dto.BackupOperate) dto.BackupCheckRes {
var res dto.BackupCheckRes
var backup model.BackupAccount
if err := copier.Copy(&backup, &req); err != nil {
res.Msg = i18n.GetMsgWithDetail("ErrStructTransform", err.Error())
return res
}
itemAccessKey, err := base64.StdEncoding.DecodeString(backup.AccessKey)
if err != nil {
res.Msg = err.Error()
return res
}
backup.AccessKey = string(itemAccessKey)
itemCredential, err := base64.StdEncoding.DecodeString(backup.Credential)
if err != nil {
res.Msg = err.Error()
return res
}
backup.Credential = string(itemCredential)
if req.Type == constant.OneDrive || req.Type == constant.GoogleDrive {
refreshToken, err := loadRefreshTokenByCode(&backup)
if err != nil {
res.Msg = err.Error()
return res
}
res.Token = base64.StdEncoding.EncodeToString([]byte(refreshToken))
}
isOk, err := u.checkBackupConn(&backup)
if err != nil {
res.Msg = err.Error()
return res
}
res.IsOk = isOk
return res
}
func (u *BackupService) Create(req dto.BackupOperate) error {
if req.Type == constant.Local {
return buserr.New("ErrBackupLocalCreate")
}
if req.Type != constant.Sftp {
req.BackupPath = strings.TrimPrefix(req.BackupPath, "/")
}
backup, _ := backupRepo.Get(repo.WithByName(req.Name))
if backup.ID != 0 {
return buserr.New("ErrRecordExist")
}
if err := copier.Copy(&backup, &req); err != nil {
return buserr.WithDetail("ErrStructTransform", err.Error(), nil)
}
itemAccessKey, err := base64.StdEncoding.DecodeString(backup.AccessKey)
if err != nil {
return err
}
backup.AccessKey = string(itemAccessKey)
itemCredential, err := base64.StdEncoding.DecodeString(backup.Credential)
if err != nil {
return err
}
backup.Credential = string(itemCredential)
backup.AccessKey, err = encrypt.StringEncrypt(backup.AccessKey)
if err != nil {
return err
}
backup.Credential, err = encrypt.StringEncrypt(backup.Credential)
if err != nil {
return err
}
if err := backupRepo.Create(&backup); err != nil {
return err
}
return nil
}
func (u *BackupService) GetBuckets(req dto.ForBuckets) ([]interface{}, error) {
itemAccessKey, err := base64.StdEncoding.DecodeString(req.AccessKey)
if err != nil {
return nil, err
}
req.AccessKey = string(itemAccessKey)
itemCredential, err := base64.StdEncoding.DecodeString(req.Credential)
if err != nil {
return nil, err
}
req.Credential = string(itemCredential)
varMap := make(map[string]interface{})
if err := json.Unmarshal([]byte(req.Vars), &varMap); err != nil {
return nil, err
}
switch req.Type {
case constant.Sftp, constant.WebDAV:
varMap["username"] = req.AccessKey
varMap["password"] = req.Credential
case constant.OSS, constant.S3, constant.MinIo, constant.Cos, constant.Kodo:
varMap["accessKey"] = req.AccessKey
varMap["secretKey"] = req.Credential
}
client, err := cloud_storage.NewCloudStorageClient(req.Type, varMap)
if err != nil {
return nil, err
}
return client.ListBuckets()
}
func (u *BackupService) Delete(id uint) error {
backup, _ := backupRepo.Get(repo.WithByID(id))
if backup.ID == 0 {
return buserr.New("ErrRecordNotFound")
}
if backup.Type == constant.Local {
return buserr.New("ErrBackupLocalDelete")
}
if err := u.CheckUsed(backup.Name, false); err != nil {
return err
}
return backupRepo.Delete(repo.WithByID(id))
}
func (u *BackupService) Update(req dto.BackupOperate) error {
backup, _ := backupRepo.Get(repo.WithByID(req.ID))
if backup.ID == 0 {
return buserr.New("ErrRecordNotFound")
}
if req.Type != constant.Sftp && req.Type != constant.Local && req.BackupPath != "/" {
req.BackupPath = strings.TrimPrefix(req.BackupPath, "/")
}
var newBackup model.BackupAccount
if err := copier.Copy(&newBackup, &req); err != nil {
return buserr.WithDetail("ErrStructTransform", err.Error(), nil)
}
itemAccessKey, err := base64.StdEncoding.DecodeString(newBackup.AccessKey)
if err != nil {
return err
}
newBackup.AccessKey = string(itemAccessKey)
itemCredential, err := base64.StdEncoding.DecodeString(newBackup.Credential)
if err != nil {
return err
}
newBackup.Credential = string(itemCredential)
if backup.Type == constant.Local {
if err := changeLocalBackup(backup.BackupPath, newBackup.BackupPath); err != nil {
return err
}
global.Dir.LocalBackupDir = newBackup.BackupPath
}
if backup.Type != constant.Local {
newBackup.AccessKey, err = encrypt.StringEncrypt(newBackup.AccessKey)
if err != nil {
return err
}
newBackup.Credential, err = encrypt.StringEncrypt(newBackup.Credential)
if err != nil {
return err
}
}
newBackup.ID = backup.ID
newBackup.CreatedAt = backup.CreatedAt
newBackup.UpdatedAt = backup.UpdatedAt
if err := backupRepo.Save(&newBackup); err != nil {
return err
}
return nil
}
func (u *BackupService) RefreshToken(req dto.OperateByID) error {
backup, _ := backupRepo.Get(repo.WithByID(req.ID))
if backup.ID == 0 {
return buserr.New("ErrRecordNotFound")
}
varMap := make(map[string]interface{})
if err := json.Unmarshal([]byte(backup.Vars), &varMap); err != nil {
return fmt.Errorf("failed to refresh %s - %s token, please retry, err: %v", backup.Type, backup.Name, err)
}
var (
refreshToken string
err error
)
switch backup.Type {
case constant.OneDrive:
refreshToken, err = client.RefreshToken("refresh_token", "refreshToken", varMap)
case constant.ALIYUN:
refreshToken, err = client.RefreshALIToken(varMap)
}
if err != nil {
varMap["refresh_status"] = constant.StatusFailed
varMap["refresh_msg"] = err.Error()
return fmt.Errorf("failed to refresh %s-%s token, please retry, err: %v", backup.Type, backup.Name, err)
}
varMap["refresh_status"] = constant.StatusSuccess
varMap["refresh_time"] = time.Now().Format(constant.DateTimeLayout)
varMap["refresh_token"] = refreshToken
varsItem, _ := json.Marshal(varMap)
backup.Vars = string(varsItem)
return backupRepo.Save(&backup)
}
func (u *BackupService) UploadForRecover(req dto.UploadForRecover) error {
fileOp := files.NewFileOp()
if !fileOp.Stat(req.TargetDir) {
if err := fileOp.CreateDir(req.TargetDir, constant.DirPerm); err != nil {
return err
}
}
return fileOp.Copy(req.FilePath, req.TargetDir)
}
func (u *BackupService) checkBackupConn(backup *model.BackupAccount) (bool, error) {
client, err := newClient(backup, false)
if err != nil {
return false, err
}
fileItem := path.Join(global.Dir.BaseDir, "1panel/tmp/test/1panel")
if _, err := os.Stat(path.Dir(fileItem)); err != nil && os.IsNotExist(err) {
if err = os.MkdirAll(path.Dir(fileItem), os.ModePerm); err != nil {
return false, err
}
}
file, err := os.OpenFile(fileItem, os.O_WRONLY|os.O_CREATE, constant.FilePerm)
if err != nil {
return false, err
}
defer file.Close()
write := bufio.NewWriter(file)
_, _ = write.WriteString("1Panel 备份账号测试文件。\n")
_, _ = write.WriteString("1Panel 備份賬號測試文件。\n")
_, _ = write.WriteString("1Panel Backs up account test files.\n")
_, _ = write.WriteString("1Panelアカウントのテストファイルをバックアップします。\n")
write.Flush()
targetPath := path.Join(backup.BackupPath, "test/1panel")
if backup.Type != constant.Sftp && backup.Type != constant.Local && targetPath != "/" {
targetPath = strings.TrimPrefix(targetPath, "/")
}
if _, err := client.Upload(context.Background(), fileItem, targetPath); err != nil {
return false, err
}
_, _ = client.Delete(path.Join(backup.BackupPath, "test/1panel"))
return true, nil
}
func (u *BackupService) LoadBackupOptions() ([]dto.BackupOption, error) {
accounts, err := backupRepo.List(repo.WithOrderDesc("created_at"))
if err != nil {
return nil, err
}
var data []dto.BackupOption
for _, account := range accounts {
var item dto.BackupOption
if err := copier.Copy(&item, &account); err != nil {
global.LOG.Errorf("copy backup account to dto backup info failed, err: %v", err)
}
data = append(data, item)
}
return data, nil
}
func (u *BackupService) CheckUsed(name string, isPublic bool) error {
account, _ := backupRepo.Get(repo.WithByName(name), backupRepo.WithByPublic(isPublic))
if account.ID == 0 {
return nil
}
cronjobs, _ := cronjobRepo.List()
for _, job := range cronjobs {
if job.DownloadAccountID == account.ID {
return buserr.New("ErrBackupInUsed")
}
ids := strings.Split(job.SourceAccountIDs, ",")
for _, idItem := range ids {
if idItem == fmt.Sprintf("%v", account.ID) {
return buserr.New("ErrBackupInUsed")
}
}
}
return nil
}
func NewBackupClientWithID(id uint) (*model.BackupAccount, cloud_storage.CloudStorageClient, error) {
account, _ := backupRepo.Get(repo.WithByID(id))
backClient, err := newClient(&account, true)
if err != nil {
return nil, nil, err
}
return &account, backClient, nil
}
type backupClientHelper struct {
id uint
accountType string
name string
backupPath string
client cloud_storage.CloudStorageClient
isOk bool
hasBackup bool
message string
}
func NewBackupClientMap(ids []string) map[string]backupClientHelper {
return NewBackupClientMapWithContext(context.Background(), ids)
}
func NewBackupClientMapWithContext(ctx context.Context, ids []string) map[string]backupClientHelper {
var accounts []model.BackupAccount
var idItems []uint
for i := 0; i < len(ids); i++ {
item, _ := strconv.Atoi(ids[i])
idItems = append(idItems, uint(item))
}
accounts, _ = backupRepo.List(repo.WithByIDs(idItems))
clientMap := make(map[string]backupClientHelper)
for _, item := range accounts {
backClient, err := newClientWithContext(ctx, &item, true)
itemHelper := backupClientHelper{
client: backClient,
name: item.Name,
backupPath: item.BackupPath,
accountType: item.Type,
id: item.ID,
isOk: err == nil,
}
if err != nil {
itemHelper.message = err.Error()
}
clientMap[fmt.Sprintf("%v", item.ID)] = itemHelper
}
return clientMap
}
func uploadWithMap(taskItem task.Task, accountMap map[string]backupClientHelper, src, dst, accountIDs string, downloadAccountID, retry uint, cleanOnFailure bool) error {
return uploadWithMapWithContext(context.Background(), taskItem, accountMap, src, dst, accountIDs, downloadAccountID, retry, cleanOnFailure, true)
}
func uploadWithMapWithContext(ctx context.Context, taskItem task.Task, accountMap map[string]backupClientHelper, src, dst, accountIDs string, downloadAccountID, retry uint, cleanOnFailure, removeSrc bool) error {
accounts := strings.Split(accountIDs, ",")
for _, account := range accounts {
if len(account) == 0 {
continue
}
itemBackup, ok := accountMap[account]
if !ok {
continue
}
if itemBackup.hasBackup {
continue
}
if !itemBackup.isOk {
taskItem.LogFailed(i18n.GetMsgWithDetail("LoadBackupFailed", itemBackup.message))
continue
}
name := itemBackup.name
if itemBackup.name == "localhost" {
name = i18n.GetMsgByKey("Localhost")
}
taskItem.LogStart(i18n.GetMsgWithMap("UploadFile", map[string]interface{}{
"file": path.Join(itemBackup.backupPath, dst),
"backup": name,
}))
for i := 0; i < int(retry)+1; i++ {
_, err := itemBackup.client.Upload(ctx, src, path.Join(itemBackup.backupPath, dst))
taskItem.LogWithStatus(i18n.GetMsgByKey("Upload"), err)
if err != nil {
if account == fmt.Sprintf("%d", downloadAccountID) {
if cleanOnFailure {
cleanupCronjobBackupArtifacts(accountMap, src, dst)
}
return err
}
} else {
break
}
}
itemBackup.hasBackup = true
accountMap[account] = itemBackup
}
if removeSrc {
os.RemoveAll(src)
}
return nil
}
func cleanupCronjobBackupArtifacts(accountMap map[string]backupClientHelper, src, dst string) {
if err := os.RemoveAll(src); err != nil {
global.LOG.Errorf("remove failed local cronjob backup file %s failed, err: %v", src, err)
}
for _, account := range accountMap {
if !account.isOk {
continue
}
if _, err := account.client.Delete(path.Join(account.backupPath, dst)); err != nil {
global.LOG.Errorf("remove failed cronjob backup file %s failed, err: %v", dst, err)
}
}
}
func markBackupFailed(recordID uint, backupErr error) {
_ = backupRepo.UpdateRecordByMap(recordID, map[string]interface{}{"status": constant.StatusFailed, "message": backupErr.Error()})
record, err := backupRepo.GetRecord(repo.WithByID(recordID))
if err != nil || record.ID == 0 {
global.LOG.Errorf("load failed backup record %d for cleanup failed, err: %v", recordID, err)
return
}
filePath := path.Join(record.FileDir, record.FileName)
if err := os.Remove(path.Join(global.Dir.LocalBackupDir, filePath)); err != nil && !os.IsNotExist(err) {
global.LOG.Errorf("remove failed local backup file %s failed, err: %v", filePath, err)
}
cleaned := make(map[string]struct{})
for _, accountID := range strings.Split(record.SourceAccountIDs, ",") {
if accountID == "" {
continue
}
if _, ok := cleaned[accountID]; ok {
continue
}
cleaned[accountID] = struct{}{}
id, err := strconv.Atoi(accountID)
if err != nil {
global.LOG.Errorf("parse backup account %s for failed backup cleanup failed, err: %v", accountID, err)
continue
}
account, storageClient, err := NewBackupClientWithID(uint(id))
if err != nil {
global.LOG.Errorf("new backup client for failed backup cleanup failed, err: %v", err)
continue
}
if _, err := storageClient.Delete(path.Join(account.BackupPath, filePath)); err != nil {
global.LOG.Errorf("remove failed backup file %s failed, err: %v", filePath, err)
}
}
}
func newClient(account *model.BackupAccount, isEncrypt bool) (cloud_storage.CloudStorageClient, error) {
return newClientWithContext(context.Background(), account, isEncrypt)
}
func newClientWithContext(ctx context.Context, account *model.BackupAccount, isEncrypt bool) (cloud_storage.CloudStorageClient, error) {
varMap := make(map[string]interface{})
if len(account.Vars) != 0 {
if err := json.Unmarshal([]byte(account.Vars), &varMap); err != nil {
return nil, err
}
}
varMap["bucket"] = account.Bucket
varMap["backupPath"] = account.BackupPath
if isEncrypt {
account.AccessKey, _ = encrypt.StringDecrypt(account.AccessKey)
account.Credential, _ = encrypt.StringDecrypt(account.Credential)
}
switch account.Type {
case constant.Sftp, constant.WebDAV:
varMap["username"] = account.AccessKey
varMap["password"] = account.Credential
case constant.OSS, constant.S3, constant.MinIo, constant.Cos, constant.Kodo:
varMap["accessKey"] = account.AccessKey
varMap["secretKey"] = account.Credential
case constant.UPYUN:
varMap["operator"] = account.AccessKey
varMap["password"] = account.Credential
}
client, err := cloud_storage.NewCloudStorageClientWithContext(ctx, account.Type, varMap)
if err != nil {
return nil, err
}
return client, nil
}
func loadRefreshTokenByCode(backup *model.BackupAccount) (string, error) {
varMap := make(map[string]interface{})
if err := json.Unmarshal([]byte(backup.Vars), &varMap); err != nil {
return "", fmt.Errorf("unmarshal backup vars failed, err: %v", err)
}
if token, ok := varMap["refresh_token"]; ok && len(token.(string)) != 0 {
return "", nil
}
refreshToken := ""
var err error
switch backup.Type {
case constant.GoogleDrive:
refreshToken, err = client.RefreshGoogleToken("authorization_code", "refreshToken", varMap)
if err != nil {
return "", err
}
case constant.OneDrive:
refreshToken, err = client.RefreshToken("authorization_code", "refreshToken", varMap)
if err != nil {
return "", err
}
}
if backup.Type != constant.ALIYUN {
varMap["refresh_token"] = refreshToken
}
itemVars, _ := json.Marshal(varMap)
backup.Vars = string(itemVars)
return refreshToken, nil
}
func loadBackupNamesByID(accountIDs string, downloadID uint) ([]string, string, error) {
accountIDList := strings.Split(accountIDs, ",")
var ids []uint
for _, item := range accountIDList {
if len(item) != 0 {
itemID, _ := strconv.Atoi(item)
ids = append(ids, uint(itemID))
}
}
list, err := backupRepo.List(repo.WithByIDs(ids))
if err != nil {
return nil, "", err
}
var accounts []string
var downloadAccount string
for _, item := range list {
accounts = append(accounts, item.Name)
if item.ID == downloadID {
downloadAccount = item.Name
}
}
return accounts, downloadAccount, nil
}
func changeLocalBackup(oldPath, newPath string) error {
fileOp := files.NewFileOp()
if fileOp.Stat(path.Join(oldPath, "app")) {
if err := fileOp.CopyDir(path.Join(oldPath, "app"), newPath); err != nil {
return err
}
}
if fileOp.Stat(path.Join(oldPath, "database")) {
if err := fileOp.CopyDir(path.Join(oldPath, "database"), newPath); err != nil {
return err
}
}
if fileOp.Stat(path.Join(oldPath, "directory")) {
if err := fileOp.CopyDir(path.Join(oldPath, "directory"), newPath); err != nil {
return err
}
}
if fileOp.Stat(path.Join(oldPath, "system_snapshot")) {
if err := fileOp.CopyDir(path.Join(oldPath, "system_snapshot"), newPath); err != nil {
return err
}
}
if fileOp.Stat(path.Join(oldPath, "website")) {
if err := fileOp.CopyDir(path.Join(oldPath, "website"), newPath); err != nil {
return err
}
}
if fileOp.Stat(path.Join(oldPath, "log")) {
if err := fileOp.CopyDir(path.Join(oldPath, "log"), newPath); err != nil {
return err
}
}
if fileOp.Stat(path.Join(oldPath, "master")) {
if err := fileOp.CopyDir(path.Join(oldPath, "master"), newPath); err != nil {
return err
}
}
_ = fileOp.RmRf(path.Join(oldPath, "app"))
_ = fileOp.RmRf(path.Join(oldPath, "database"))
_ = fileOp.RmRf(path.Join(oldPath, "directory"))
_ = fileOp.RmRf(path.Join(oldPath, "system_snapshot"))
_ = fileOp.RmRf(path.Join(oldPath, "website"))
_ = fileOp.RmRf(path.Join(oldPath, "log"))
_ = fileOp.RmRf(path.Join(oldPath, "master"))
return nil
}