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 }