diff --git a/cmd/climc/shell/compute/schedulers.go b/cmd/climc/shell/compute/schedulers.go index 70667b439e..52e0fc4adc 100644 --- a/cmd/climc/shell/compute/schedulers.go +++ b/cmd/climc/shell/compute/schedulers.go @@ -111,10 +111,11 @@ func init() { type SchedulerCleanCacheOptions struct { HostId string `help:"ID of host" short-token:"h"` SessionId string `help:"Session id" short-token:"s"` + Sync bool `help:"Sync do clean sched cache"` } R(&SchedulerCleanCacheOptions{}, "scheduler-clean-cache", "Clean scheduler hosts cache", func(s *mcclient.ClientSession, args *SchedulerCleanCacheOptions) error { - err := modules.SchedManager.CleanCache(s, args.HostId, args.SessionId) + err := modules.SchedManager.CleanCache(s, args.HostId, args.SessionId, args.Sync) if err != nil { return err } diff --git a/cmd/climc/shell/compute/servers.go b/cmd/climc/shell/compute/servers.go index 6b1337f172..14405520cf 100644 --- a/cmd/climc/shell/compute/servers.go +++ b/cmd/climc/shell/compute/servers.go @@ -298,6 +298,12 @@ func init() { return nil }) + R(&options.ServerIdsOptions{}, "server-reconcile-backup", "Reconcile backup server", func(s *mcclient.ClientSession, opts *options.ServerIdsOptions) error { + ret := modules.Servers.BatchPerformAction(s, opts.ID, "reconcile-backup", nil) + printBatchResults(ret, modules.Servers.GetColumns(s)) + return nil + }) + R(&options.ServerIdsOptions{}, "server-create-backup", "Create backup guest", func(s *mcclient.ClientSession, opts *options.ServerIdsOptions) error { ret := modules.Servers.BatchPerformAction(s, opts.ID, "create-backup", nil) printBatchResults(ret, modules.Servers.GetColumns(s)) diff --git a/pkg/apis/compute/api.go b/pkg/apis/compute/api.go index a83244605d..ab0ad3da3c 100644 --- a/pkg/apis/compute/api.go +++ b/pkg/apis/compute/api.go @@ -289,7 +289,7 @@ type ServerConfigs struct { // 虚拟机高可用(创建备机) // default: false - // requried: false + // required: false Backup bool `json:"backup"` // 创建虚拟机数量 diff --git a/pkg/cloudcommon/db/db_dispatcher.go b/pkg/cloudcommon/db/db_dispatcher.go index 52e2eed4a5..5be5c41011 100644 --- a/pkg/cloudcommon/db/db_dispatcher.go +++ b/pkg/cloudcommon/db/db_dispatcher.go @@ -1484,6 +1484,9 @@ func (dispatcher *DBModelDispatcher) PerformAction(ctx context.Context, idStr st lockman.LockObject(ctx, model) defer lockman.ReleaseObject(ctx, model) + if err := model.PreCheckPerformAction(ctx, userCred, action, query, data); err != nil { + return nil, err + } return objectPerformAction(dispatcher, model, reflect.ValueOf(model), ctx, userCred, action, query, data) } diff --git a/pkg/cloudcommon/db/interface.go b/pkg/cloudcommon/db/interface.go index 0fbb30e592..85e4e31528 100644 --- a/pkg/cloudcommon/db/interface.go +++ b/pkg/cloudcommon/db/interface.go @@ -162,6 +162,7 @@ type IModel interface { // allow perform action AllowPerformAction(ctx context.Context, userCred mcclient.TokenCredential, action string, query jsonutils.JSONObject, data jsonutils.JSONObject) bool PerformAction(ctx context.Context, userCred mcclient.TokenCredential, action string, query jsonutils.JSONObject, data jsonutils.JSONObject) (jsonutils.JSONObject, error) + PreCheckPerformAction(ctx context.Context, userCred mcclient.TokenCredential, action string, query jsonutils.JSONObject, data jsonutils.JSONObject) error // update hooks ValidateUpdateCondition(ctx context.Context) error diff --git a/pkg/cloudcommon/db/modelbase.go b/pkg/cloudcommon/db/modelbase.go index b0d7b8d1f3..ec8a3372f7 100644 --- a/pkg/cloudcommon/db/modelbase.go +++ b/pkg/cloudcommon/db/modelbase.go @@ -497,6 +497,13 @@ func (model *SModelBase) PerformAction(ctx context.Context, userCred mcclient.To return nil, httperrors.NewActionNotFoundError("Action %s not found", action) } +func (model *SModelBase) PreCheckPerformAction( + ctx context.Context, userCred mcclient.TokenCredential, + action string, query jsonutils.JSONObject, data jsonutils.JSONObject, +) error { + return nil +} + // update hooks func (model *SModelBase) AllowUpdateItem(ctx context.Context, userCred mcclient.TokenCredential) bool { return false diff --git a/pkg/cloudcommon/db/taskman/tasks.go b/pkg/cloudcommon/db/taskman/tasks.go index 2391a894e2..4df261c97b 100644 --- a/pkg/cloudcommon/db/taskman/tasks.go +++ b/pkg/cloudcommon/db/taskman/tasks.go @@ -125,6 +125,13 @@ func (manager *STaskManager) PerformAction(ctx context.Context, userCred mcclien return resp, nil } +func (manager *STask) PreCheckPerformAction( + ctx context.Context, userCred mcclient.TokenCredential, + action string, query jsonutils.JSONObject, data jsonutils.JSONObject, +) error { + return nil +} + func (self *STask) GetOwnerId() mcclient.IIdentityProvider { owner := db.SOwnerId{DomainId: self.UserCred.GetProjectDomainId(), Domain: self.UserCred.GetProjectDomain(), ProjectId: self.UserCred.GetProjectId(), Project: self.UserCred.GetProjectName()} diff --git a/pkg/compute/models/guest_actions.go b/pkg/compute/models/guest_actions.go index 0b4c14ea9d..c08d3bfbf9 100644 --- a/pkg/compute/models/guest_actions.go +++ b/pkg/compute/models/guest_actions.go @@ -82,6 +82,27 @@ func (self *SGuest) GetDetailsVnc(ctx context.Context, userCred mcclient.TokenCr } } +func (self *SGuest) PreCheckPerformAction( + ctx context.Context, userCred mcclient.TokenCredential, + action string, query jsonutils.JSONObject, data jsonutils.JSONObject, +) error { + if self.Hypervisor == api.HYPERVISOR_KVM { + host := self.GetHost() + if (host.HostStatus == api.HOST_OFFLINE || !host.Enabled.Bool()) && + utils.IsInStringArray(action, + []string{ + "start", "restart", "stop", "reset", "rebuild-root", + "change-config", "instance-snapshot", "snapshot-and-clone", + "attach-isolated-device", "detach-isolated-deivce", + "insert-iso", "eject-iso", "deploy", "create-backup", + }) { + return httperrors.NewInvalidStatusError( + "host status %s and enabled %v, can't do server %s", host.HostStatus, host.Enabled.Bool(), action) + } + } + return nil +} + func (self *SGuest) AllowPerformMonitor(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, @@ -3316,7 +3337,10 @@ func (self *SGuest) guestDisksStorageTypeIsShared() bool { return true } -func (self *SGuest) PerformCreateBackup(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) (jsonutils.JSONObject, error) { +func (self *SGuest) PerformCreateBackup( + ctx context.Context, userCred mcclient.TokenCredential, + query jsonutils.JSONObject, data jsonutils.JSONObject, +) (jsonutils.JSONObject, error) { if len(self.BackupHostId) > 0 { return nil, httperrors.NewBadRequestError("Already have backup server") } @@ -3336,7 +3360,12 @@ func (self *SGuest) PerformCreateBackup(ctx context.Context, userCred mcclient.T if hasSnapshot { return nil, httperrors.NewBadRequestError("Cannot create backup with snapshot") } + return self.StartGuestCreateBackupTask(ctx, userCred, "", data) +} +func (self *SGuest) StartGuestCreateBackupTask( + ctx context.Context, userCred mcclient.TokenCredential, parentTaskId string, data jsonutils.JSONObject, +) (jsonutils.JSONObject, error) { req := self.getGuestBackupResourceRequirements(ctx, userCred) keys, err := self.GetQuotaKeys() if err != nil { @@ -3350,7 +3379,7 @@ func (self *SGuest) PerformCreateBackup(ctx context.Context, userCred mcclient.T params := data.(*jsonutils.JSONDict) params.Set("guest_status", jsonutils.NewString(self.Status)) - task, err := taskman.TaskManager.NewTask(ctx, "GuestCreateBackupTask", self, userCred, params, "", "", &req) + task, err := taskman.TaskManager.NewTask(ctx, "GuestCreateBackupTask", self, userCred, params, parentTaskId, "", &req) if err != nil { quotas.CancelPendingUsage(ctx, userCred, &req, &req, false) log.Errorln(err) @@ -3379,6 +3408,7 @@ func (self *SGuest) PerformDeleteBackup(ctx context.Context, userCred mcclient.T taskData := jsonutils.NewDict() taskData.Set("purge", jsonutils.NewBool(jsonutils.QueryBoolean(data, "purge", false))) + taskData.Set("create", jsonutils.NewBool(jsonutils.QueryBoolean(data, "create", false))) self.SetStatus(userCred, api.VM_DELETING_BACKUP, "delete backup server") if task, err := taskman.TaskManager.NewTask( ctx, "GuestDeleteBackupTask", self, userCred, taskData, "", "", nil); err != nil { @@ -3414,6 +3444,44 @@ func (self *SGuest) StartCreateBackup(ctx context.Context, userCred mcclient.Tok return nil } +func (self *SGuest) AllowPerformReconcileBackup(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) bool { + return db.IsAdminAllowPerform(userCred, self, "reconcile-backup") +} + +func (self *SGuest) PerformReconcileBackup(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) (jsonutils.JSONObject, error) { + switchBackup := self.GetMetadata("switch_backup", userCred) + createBackup := self.GetMetadata("create_backup", userCred) + if len(switchBackup) == 0 && len(createBackup) == 0 { + return nil, httperrors.NewBadRequestError("guest doesn't need reconcile backup") + } + if len(switchBackup) > 0 { + data := jsonutils.NewDict() + data.Set("purge_backup", jsonutils.JSONTrue) + return self.PerformSwitchToBackup(ctx, userCred, nil, data) + } else { + return nil, self.StartReconcileBackup(ctx, userCred) + } +} + +func (self *SGuest) StartReconcileBackup(ctx context.Context, userCred mcclient.TokenCredential) error { + if len(self.BackupHostId) > 0 { + data := jsonutils.NewDict() + data.Set("purge", jsonutils.JSONTrue) + data.Set("create", jsonutils.JSONTrue) + _, err := self.PerformDeleteBackup(ctx, userCred, nil, data) + if err != nil { + return err + } + } else { + params := jsonutils.NewDict() + params.Set("reconcile_backup", jsonutils.JSONTrue) + if _, err := self.StartGuestCreateBackupTask(ctx, userCred, "", params); err != nil { + return err + } + } + return nil +} + func (self *SGuest) AllowPerformSetExtraOption(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) bool { return db.IsAdminAllowPerform(userCred, self, "set-extra-option") } diff --git a/pkg/compute/models/guests.go b/pkg/compute/models/guests.go index 4c4f0d9969..679c9b981f 100644 --- a/pkg/compute/models/guests.go +++ b/pkg/compute/models/guests.go @@ -4472,6 +4472,113 @@ func (manager *SGuestManager) DeleteExpiredPostpaidServers(ctx context.Context, } } +func (self *SGuestManager) ReconcileBackupGuests(ctx context.Context, userCred mcclient.TokenCredential, isStart bool) { + self.switchBackupGuests(ctx, userCred) + self.createBackupGuests(ctx, userCred) +} + +func (self *SGuestManager) switchBackupGuests(ctx context.Context, userCred mcclient.TokenCredential) { + q := self.Query() + metaDataQuery := db.Metadata.Query().Startswith("id", "server::"). + Equals("key", "switch_backup").IsNotEmpty("value").GroupBy("id") + metaDataQuery.AppendField(sqlchemy.SubStr("guest_id", metaDataQuery.Field("id"), len("server::")+1, 0)) + subQ := metaDataQuery.SubQuery() + q = q.Join(subQ, sqlchemy.Equals(q.Field("id"), subQ.Field("guest_id"))) + + guests := make([]SGuest, 0) + err := db.FetchModelObjects(GuestManager, q, &guests) + if err != nil { + log.Errorf("ReconcileBackupGuests failed fetch guests %s", err) + return + } + for i := 0; i < len(guests); i++ { + val := guests[i].GetMetadataJson("switch_backup", userCred) + t, err := val.GetTime() + if err != nil { + log.Errorf("failed get time from metadata switch_backup %s", err) + continue + } + if time.Now().After(t) { + data := jsonutils.NewDict() + data.Set("purge_backup", jsonutils.JSONTrue) + _, err := guests[i].PerformSwitchToBackup(ctx, userCred, nil, data) + if err != nil { + db.OpsLog.LogEvent( + &guests[i], db.ACT_SWITCH_FAILED, fmt.Sprintf("switchBackupGuests on reconcile_backup: %s", err), userCred, + ) + logclient.AddSimpleActionLog( + &guests[i], logclient.ACT_SWITCH_TO_BACKUP, + fmt.Sprintf("switchBackupGuests on reconcile_backup: %s", err), userCred, false, + ) + } + } + } +} + +func (self *SGuestManager) createBackupGuests(ctx context.Context, userCred mcclient.TokenCredential) { + q := self.Query() + metaDataQuery := db.Metadata.Query().Startswith("id", "server::"). + Equals("key", "create_backup").IsNotEmpty("value").GroupBy("id") + metaDataQuery.AppendField(sqlchemy.SubStr("guest_id", metaDataQuery.Field("id"), len("server::")+1, 0)) + subQ := metaDataQuery.SubQuery() + q = q.Join(subQ, sqlchemy.Equals(q.Field("id"), subQ.Field("guest_id"))) + + guests := make([]SGuest, 0) + err := db.FetchModelObjects(GuestManager, q, &guests) + if err != nil { + log.Errorf("ReconcileBackupGuests failed fetch guests %s", err) + return + } + for i := 0; i < len(guests); i++ { + val := guests[i].GetMetadataJson("create_backup", userCred) + t, err := val.GetTime() + if err != nil { + log.Errorf("failed get time from metadata create_backup %s", err) + continue + } + if time.Now().After(t) { + if len(guests[i].BackupHostId) > 0 { + data := jsonutils.NewDict() + data.Set("purge", jsonutils.JSONTrue) + data.Set("create", jsonutils.JSONTrue) + _, err := guests[i].PerformDeleteBackup(ctx, userCred, nil, data) + if err != nil { + db.OpsLog.LogEvent( + &guests[i], db.ACT_DELETE_BACKUP_FAILED, + fmt.Sprintf("PerformDeleteBackup on ReconcileBackupGuests: %s", err), userCred, + ) + logclient.AddSimpleActionLog( + &guests[i], logclient.ACT_DELETE_BACKUP, + fmt.Sprintf("PerformDeleteBackup on ReconcileBackupGuests: %s", err), userCred, false, + ) + } + } else { + params := jsonutils.NewDict() + params.Set("reconcile_backup", jsonutils.JSONTrue) + if _, err := guests[i].StartGuestCreateBackupTask(ctx, userCred, "", params); err != nil { + db.OpsLog.LogEvent( + &guests[i], db.ACT_CREATE_BACKUP_FAILED, + fmt.Sprintf("StartGuestCreateBackupTask on ReconcileBackupGuests: %s", err), userCred, + ) + logclient.AddSimpleActionLog( + &guests[i], logclient.ACT_CREATE_BACKUP, + fmt.Sprintf("StartGuestCreateBackupTask on ReconcileBackupGuests: %s", err), userCred, false, + ) + } + } + } + } +} + +func (self *SGuest) isInReconcile(userCred mcclient.TokenCredential) bool { + switchBackup := self.GetMetadata("switch_backup", userCred) + createBackup := self.GetMetadata("create_backup", userCred) + if len(switchBackup) > 0 || len(createBackup) > 0 { + return true + } + return false +} + func (self *SGuest) GetEip() (*SElasticip, error) { return ElasticipManager.getEipForInstance(api.EIP_ASSOCIATE_TYPE_SERVER, self.Id) } diff --git a/pkg/compute/models/host_health.go b/pkg/compute/models/host_health.go index 3f0a7bee44..69f95e2227 100644 --- a/pkg/compute/models/host_health.go +++ b/pkg/compute/models/host_health.go @@ -34,6 +34,8 @@ type SHostHealthChecker struct { cli *etcd.SEtcdClient // time of wait host reconnect timeout time.Duration + // hosts chan + hc map[string]chan struct{} } func hostKey(hostId string) string { @@ -44,7 +46,11 @@ func InitHostHealthChecker(cli *etcd.SEtcdClient, timeout int) *SHostHealthCheck if hostHealthChecker != nil { return hostHealthChecker } - hostHealthChecker = &SHostHealthChecker{cli, time.Duration(timeout) * time.Second} + hostHealthChecker = &SHostHealthChecker{ + cli: cli, + timeout: time.Duration(timeout) * time.Second, + hc: make(map[string]chan struct{}), + } return hostHealthChecker } @@ -70,27 +76,31 @@ func (h *SHostHealthChecker) startHealthCheck(ctx context.Context) { } func (h *SHostHealthChecker) startWatcher(ctx context.Context, hostId string) { - log.Debugf("Start watch host %s", hostId) + log.Infof("Start watch host %s", hostId) var ( ch chan struct{} key = hostKey(hostId) ) + _, err := h.cli.Get(ctx, key) if err == etcd.ErrNoSuchKey { - log.Errorf("No such key %s", hostId) + log.Warningf("No such key %s", hostId) ch = make(chan struct{}) go func() { select { case <-time.NewTimer(h.timeout).C: h.onHostUnhealthy(ctx, hostId) - case <-ch: + case <-h.hc[hostId]: h.startWatcher(ctx, hostId) case <-ctx.Done(): log.Infof("exit watch host %s", hostId) } }() } - h.cli.Watch(ctx, key, h.onHostOnline(hostId, ch), h.onHostOffline(hostId)) + if _, ok := h.hc[hostId]; !ok { + h.hc[hostId] = ch + } + h.cli.Watch(ctx, key, h.onHostOnline(hostId), h.onHostOffline(hostId)) } func (h *SHostHealthChecker) onHostUnhealthy(ctx context.Context, hostId string) { @@ -102,11 +112,11 @@ func (h *SHostHealthChecker) onHostUnhealthy(ctx context.Context, hostId string) } } -func (h *SHostHealthChecker) onHostOnline(hostId string, ch chan struct{}) etcd.TEtcdCreateEventFunc { +func (h *SHostHealthChecker) onHostOnline(hostId string) etcd.TEtcdCreateEventFunc { return func(key, value []byte) { - log.Debugf("Got host online %s", hostId) - if ch != nil { - close(ch) + log.Infof("Got host online %s", hostId) + if h.hc[hostId] != nil { + h.hc[hostId] <- struct{}{} } } } @@ -114,10 +124,14 @@ func (h *SHostHealthChecker) onHostOnline(hostId string, ch chan struct{}) etcd. func (h *SHostHealthChecker) onHostOffline(hostId string) etcd.TEtcdModifyEventFunc { return func(key, oldvalue, value []byte) { log.Warningf("host %s disconnect with etcd", hostId) - host := HostManager.FetchHostById(hostId) - if host.EnableHealthCheck == true { - h.startWatcher(context.Background(), hostId) - } + go func() { + select { + case <-time.NewTimer(h.timeout).C: + h.onHostUnhealthy(context.Background(), hostId) + case <-h.hc[hostId]: + h.startWatcher(context.Background(), hostId) + } + }() } } @@ -127,6 +141,7 @@ func (h *SHostHealthChecker) WatchHost(ctx context.Context, hostId string) { } func (h *SHostHealthChecker) UnwatchHost(ctx context.Context, hostId string) { - log.Debugf("Unwatch host %s", hostId) + log.Infof("Unwatch host %s", hostId) h.cli.Unwatch(hostKey(hostId)) + delete(h.hc, hostId) } diff --git a/pkg/compute/models/hosts.go b/pkg/compute/models/hosts.go index 851943e009..2a66cae2bf 100644 --- a/pkg/compute/models/hosts.go +++ b/pkg/compute/models/hosts.go @@ -207,7 +207,6 @@ func (manager *SHostManager) ListItemFilter( if err != nil { return nil, errors.Wrap(err, "SManagedResourceBaseManager.ListItemFilter") } - q, err = manager.SExternalizedResourceBaseManager.ListItemFilter(ctx, q, userCred, query.ExternalizedResourceBaseListInput) if err != nil { return nil, errors.Wrap(err, "SExternalizedResourceBaseManager.ListItemFilter") @@ -1032,7 +1031,7 @@ func (self *SHostManager) ClearSchedDescCache(hostId string) error { func (self *SHostManager) ClearSchedDescSessionCache(hostId, sessionId string) error { s := auth.GetAdminSession(context.Background(), options.Options.Region, "") - return modules.SchedManager.CleanCache(s, hostId, sessionId) + return modules.SchedManager.CleanCache(s, hostId, sessionId, false) } func (self *SHost) ClearSchedDescCache() error { @@ -1043,6 +1042,20 @@ func (self *SHost) ClearSchedDescSessionCache(sessionId string) error { return HostManager.ClearSchedDescSessionCache(self.Id, sessionId) } +// sync clear sched desc on scheduler side +func (self *SHostManager) SyncClearSchedDescSessionCache(hostId, sessionId string) error { + s := auth.GetAdminSession(context.Background(), options.Options.Region, "") + return modules.SchedManager.CleanCache(s, hostId, sessionId, true) +} + +func (self *SHost) SyncCleanSchedDescCache() error { + return self.SyncClearSchedDescSessionCache("") +} + +func (self *SHost) SyncClearSchedDescSessionCache(sessionId string) error { + return HostManager.SyncClearSchedDescSessionCache(self.Id, sessionId) +} + func (self *SHost) AllowGetDetailsSpec(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject) bool { return db.IsAdminAllowGetSpec(userCred, self, "spec") } @@ -1401,6 +1414,28 @@ func (self *SHost) GetGuests() []SGuest { return guests } +func (self *SHost) GetGuestsMasterOnThisHost() []SGuest { + q := self.GetGuestsQuery().IsNotEmpty("backup_host_id") + guests := make([]SGuest, 0) + err := db.FetchModelObjects(GuestManager, q, &guests) + if err != nil { + log.Errorf("GetGuests %s", err) + return nil + } + return guests +} + +func (self *SHost) GetGuestsBackupOnThisHost() []SGuest { + q := GuestManager.Query().Equals("backup_host_id", self.Id) + guests := make([]SGuest, 0) + err := db.FetchModelObjects(GuestManager, q, &guests) + if err != nil { + log.Errorf("GetGuests %s", err) + return nil + } + return guests +} + func (self *SHost) GetGuestCount() (int, error) { q := self.GetGuestsQuery() return q.CountWithError() @@ -4885,17 +4920,69 @@ func (host *SHost) PerformHostMaintenance(ctx context.Context, userCred mcclient func (host *SHost) OnHostDown(ctx context.Context, userCred mcclient.TokenCredential) { log.Errorf("watched host down %s", host.Id) db.OpsLog.LogEvent(host, db.ACT_HOST_DOWN, "", userCred) + if _, err := host.SaveCleanUpdates(func() error { + host.EnableHealthCheck = false + host.HostStatus = api.HOST_OFFLINE + return nil + }); err != nil { + log.Errorf("update host %s failed %s", host.Id, err) + } + host.SyncCleanSchedDescCache() + host.switchWithBackup(ctx, userCred) + host.migrateOnHostDown(ctx, userCred) +} + +func (host *SHost) switchWithBackup(ctx context.Context, userCred mcclient.TokenCredential) { + guests := host.GetGuestsMasterOnThisHost() + for i := 0; i < len(guests); i++ { + if guests[i].isInReconcile(userCred) { + log.Warningf("guest %s is in reconcile", guests[i].GetName()) + continue + } + data := jsonutils.NewDict() + data.Set("purge_backup", jsonutils.JSONTrue) + _, err := guests[i].PerformSwitchToBackup(ctx, userCred, nil, data) + if err != nil { + db.OpsLog.LogEvent( + &guests[i], db.ACT_SWITCH_FAILED, fmt.Sprintf("PerformSwitchToBackup on host down: %s", err), userCred, + ) + logclient.AddSimpleActionLog( + &guests[i], logclient.ACT_SWITCH_TO_BACKUP, + fmt.Sprintf("PerformSwitchToBackup on host down: %s", err), userCred, false, + ) + } else { + guests[i].SetMetadata(ctx, "origin_status", guests[i].Status, userCred) + } + } + + guests2 := host.GetGuestsBackupOnThisHost() + for i := 0; i < len(guests2); i++ { + if guests2[i].isInReconcile(userCred) { + log.Warningf("guest %s is in reconcile", guests[i].GetName()) + continue + } + data := jsonutils.NewDict() + data.Set("purge", jsonutils.JSONTrue) + data.Set("create", jsonutils.JSONTrue) + _, err := guests2[i].PerformDeleteBackup(ctx, userCred, nil, data) + if err != nil { + db.OpsLog.LogEvent( + &guests2[i], db.ACT_DELETE_BACKUP_FAILED, fmt.Sprintf("PerformDeleteBackup on host down: %s", err), userCred, + ) + logclient.AddSimpleActionLog( + &guests2[i], logclient.ACT_DELETE_BACKUP, + fmt.Sprintf("PerformDeleteBackup on host down: %s", err), userCred, false, + ) + } + } +} + +func (host *SHost) migrateOnHostDown(ctx context.Context, userCred mcclient.TokenCredential) { if host.GetMetadata("__auto_migrate_on_host_down", nil) == "enable" { if err := host.MigrateSharedStorageServers(ctx, userCred); err != nil { db.OpsLog.LogEvent(host, db.ACT_HOST_DOWN, fmt.Sprintf("migrate servers failed %s", err), userCred) } } - if _, err := host.SaveUpdates(func() error { - host.EnableHealthCheck = false - return nil - }); err != nil { - log.Errorf("update host %s failed %s", host.Id, err) - } } func (host *SHost) MigrateSharedStorageServers(ctx context.Context, userCred mcclient.TokenCredential) error { diff --git a/pkg/compute/options/options.go b/pkg/compute/options/options.go index 026c413a2c..92a3b2ca20 100644 --- a/pkg/compute/options/options.go +++ b/pkg/compute/options/options.go @@ -138,6 +138,8 @@ type ComputeOptions struct { ScheduledTaskQueueSize int `help:"the maximum number of scheduled tasks that are being executed simultaneously" default:"100"` + ReconcileGuestBackupIntervalSeconds int `help:"interval reconcile guest bakcups" default:"30"` + SCapabilityOptions SASControllerOptions common_options.CommonOptions diff --git a/pkg/compute/service/service.go b/pkg/compute/service/service.go index ac8364ba3e..b35e126c46 100644 --- a/pkg/compute/service/service.go +++ b/pkg/compute/service/service.go @@ -136,6 +136,7 @@ func StartService() { cron.AddJobAtIntervalsWithStartRun("CalculateInfrasQuotaUsages", time.Duration(opts.CalculateQuotaUsageIntervalSeconds)*time.Second, models.InfrasQuotaManager.CalculateQuotaUsages, true) cron.AddJobAtIntervalsWithStartRun("AutoSyncCloudaccountTask", time.Duration(opts.CloudAutoSyncIntervalSeconds)*time.Second, models.CloudaccountManager.AutoSyncCloudaccountTask, true) + cron.AddJobAtIntervalsWithStartRun("ReconcileBackupGuests", time.Duration(opts.ReconcileGuestBackupIntervalSeconds)*time.Second, models.GuestManager.ReconcileBackupGuests, true) cron.AddJobEveryFewHour("AutoDiskSnapshot", 1, 5, 0, models.DiskManager.AutoDiskSnapshot, false) cron.AddJobEveryFewHour("SnapshotsCleanup", 1, 35, 0, models.SnapshotManager.CleanupSnapshots, false) diff --git a/pkg/compute/tasks/guest_backup_tasks.go b/pkg/compute/tasks/guest_backup_tasks.go index 157a76396f..e4160140a1 100644 --- a/pkg/compute/tasks/guest_backup_tasks.go +++ b/pkg/compute/tasks/guest_backup_tasks.go @@ -17,8 +17,10 @@ package tasks import ( "context" "fmt" + "time" "yunion.io/x/jsonutils" + "yunion.io/x/log" "yunion.io/x/pkg/utils" api "yunion.io/x/onecloud/pkg/apis/compute" @@ -96,6 +98,7 @@ func (self *GuestSwitchToBackupTask) OnBackupGuestStoped(ctx context.Context, gu } func (self *GuestSwitchToBackupTask) OnNewMasterStarted(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { + guest.RemoveMetadata(ctx, "origin_status", self.UserCred) self.OnComplete(ctx, guest, nil) } @@ -103,22 +106,88 @@ func (self *GuestSwitchToBackupTask) OnFail(ctx context.Context, guest *models.S guest.SetStatus(self.UserCred, api.VM_SWITCH_TO_BACKUP_FAILED, reason) db.OpsLog.LogEvent(guest, db.ACT_SWITCH_FAILED, reason, self.UserCred) logclient.AddActionLogWithContext(ctx, guest, logclient.ACT_SWITCH_TO_BACKUP, reason, self.UserCred, false) + self.SetSwitchFiledGuestMetadata(ctx, guest) self.SetStageFailed(ctx, reason) } +func (self *GuestSwitchToBackupTask) SetSwitchFiledGuestMetadata(ctx context.Context, guest *models.SGuest) { + if res := guest.GetMetadata("switch_backup", self.UserCred); len(res) == 0 { + guest.SetMetadata( + ctx, "switch_backup", jsonutils.NewTimeString(time.Now().Add(time.Minute*1)), self.UserCred) + guest.SetMetadata(ctx, "switch_backup_count", jsonutils.NewInt(1), self.UserCred) + } else { + count := guest.GetMetadataJson("switch_backup_count", self.UserCred) + cnt, _ := count.Int() + cnt += 1 + dur := cnt + if dur > 15 { + dur = 15 + } + + guest.SetMetadata(ctx, "switch_backup", + jsonutils.NewTimeString(time.Now().Add(time.Minute*time.Duration(dur))), self.UserCred) + guest.SetMetadata(ctx, "switch_backup_count", jsonutils.NewInt(cnt), self.UserCred) + } +} + +func (self *GuestSwitchToBackupTask) CleanGuestMetadata(ctx context.Context, guest *models.SGuest) { + if res := guest.GetMetadata("switch_backup", self.UserCred); len(res) > 0 { + guest.RemoveMetadata(ctx, "switch_backup", self.UserCred) + guest.RemoveMetadata(ctx, "switch_backup_count", self.UserCred) + } +} + func (self *GuestSwitchToBackupTask) OnComplete(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { self.SetStageComplete(ctx, nil) } func (self *GuestSwitchToBackupTask) OnSwitched(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { - guest.SetMetadata(ctx, "__mirror_job_status", "", self.UserCred) + if err := guest.SetMetadata(ctx, "__mirror_job_status", "", self.UserCred); err != nil { + self.OnSwitchedFailed(ctx, guest, jsonutils.NewString("guest set metadata failed")) + return + } + // switched bakcup success + self.CleanGuestMetadata(ctx, guest) + + // add a new backup server + self.SetStage("OnCreatedBackup", nil) + if jsonutils.QueryBoolean(self.Params, "purge_backup", false) || + jsonutils.QueryBoolean(self.Params, "delete_backup", false) { + params := jsonutils.NewDict() + params.Set("reconcile_backup", jsonutils.JSONTrue) + if _, err := guest.StartGuestCreateBackupTask(ctx, self.UserCred, self.Id, params); err != nil { + log.Errorf("guest start create backup failed %s", err) + self.failedStartCreateBackupTask(ctx, guest) + } + } else { + self.OnCreatedBackup(ctx, guest, nil) + } +} + +func (self *GuestSwitchToBackupTask) failedStartCreateBackupTask(ctx context.Context, guest *models.SGuest) { + guest.SetMetadata(ctx, "create_backup", jsonutils.NewTimeString(time.Now().Add(time.Minute*1)), self.UserCred) + guest.SetMetadata(ctx, "create_backup_count", jsonutils.NewInt(1), self.UserCred) + self.OnCreatedBackup(ctx, guest, nil) +} + +func (self *GuestSwitchToBackupTask) OnCreatedBackup(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { oldStatus, _ := self.Params.GetString("old_status") - if utils.IsInStringArray(oldStatus, api.VM_RUNNING_STATUS) { - self.SetStage("OnNewMasterStarted", nil) + originStatus := guest.GetMetadata("origin_status", self.UserCred) + if (!utils.IsInStringArray(guest.Status, api.VM_RUNNING_STATUS) && utils.IsInStringArray(oldStatus, api.VM_RUNNING_STATUS)) || + utils.IsInStringArray(originStatus, api.VM_RUNNING_STATUS) { + self.SetStage("OnNewMastqerStarted", nil) guest.StartGueststartTask(ctx, self.UserCred, nil, self.GetTaskId()) } else { self.OnComplete(ctx, guest, nil) } + if len(originStatus) > 0 { + guest.RemoveMetadata(ctx, "origin_status", self.UserCred) + } +} + +func (self *GuestSwitchToBackupTask) OnCreatedBackupFailed(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { + log.Errorf("failed create backup %s", data) + self.OnCreatedBackup(ctx, guest, data) } func (self *GuestSwitchToBackupTask) OnSwitchedFailed(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { @@ -335,6 +404,12 @@ func (self *GuestCreateBackupTask) OnSyncToBackupFailed(ctx context.Context, gue func (self *GuestCreateBackupTask) TaskCompleted(ctx context.Context, guest *models.SGuest, reason string) { db.OpsLog.LogEvent(guest, db.ACT_CREATE_BACKUP, reason, self.UserCred) + + if res := guest.GetMetadata("create_backup", self.UserCred); len(res) > 0 { + guest.RemoveMetadata(ctx, "create_backup", self.UserCred) + guest.RemoveMetadata(ctx, "create_backup_count", self.UserCred) + } + logclient.AddActionLogWithContext(ctx, guest, logclient.ACT_CREATE_BACKUP, reason, self.UserCred, true) self.SetStageComplete(ctx, nil) guest.StartSyncstatus(ctx, self.UserCred, "") @@ -342,6 +417,25 @@ func (self *GuestCreateBackupTask) TaskCompleted(ctx context.Context, guest *mod func (self *GuestCreateBackupTask) TaskFailed(ctx context.Context, guest *models.SGuest, reason string) { guest.SetStatus(self.UserCred, api.VM_BACKUP_CREATE_FAILED, reason) + if jsonutils.QueryBoolean(self.Params, "reconcile_backup", false) { + if res := guest.GetMetadata("create_backup", self.UserCred); len(res) == 0 { + guest.SetMetadata( + ctx, "create_backup", jsonutils.NewTimeString(time.Now().Add(time.Minute*1)), self.UserCred) + guest.SetMetadata(ctx, "create_backup_count", jsonutils.NewInt(1), self.UserCred) + } else { + count := guest.GetMetadataJson("create_backup_count", self.UserCred) + cnt, _ := count.Int() + cnt += 1 + dur := cnt + if dur > 15 { + dur = 15 + } + + guest.SetMetadata(ctx, "create_backup", + jsonutils.NewTimeString(time.Now().Add(time.Minute*time.Duration(dur))), self.UserCred) + guest.SetMetadata(ctx, "create_backup_count", jsonutils.NewInt(cnt), self.UserCred) + } + } db.OpsLog.LogEvent(guest, db.ACT_CREATE_BACKUP_FAILED, reason, self.UserCred) logclient.AddActionLogWithContext(ctx, guest, logclient.ACT_CREATE_BACKUP, reason, self.UserCred, false) self.SetStageFailed(ctx, reason) diff --git a/pkg/compute/tasks/guest_delete_backup_task.go b/pkg/compute/tasks/guest_delete_backup_task.go index 5c574c24a5..a88e566029 100644 --- a/pkg/compute/tasks/guest_delete_backup_task.go +++ b/pkg/compute/tasks/guest_delete_backup_task.go @@ -70,16 +70,63 @@ func (self *GuestDeleteBackupTask) StartDeleteBackupOnHost(ctx context.Context, taskData.Set("purge", jsonutils.NewBool(jsonutils.QueryBoolean(self.Params, "purge", false))) taskData.Set("host_id", jsonutils.NewString(guest.BackupHostId)) taskData.Set("failed_status", jsonutils.NewString(compute.VM_BACKUP_DELETE_FAILED)) + + self.SetStage("OnDeleteOnHost", nil) if task, err := taskman.TaskManager.NewTask( - ctx, "GuestDeleteOnHostTask", guest, self.UserCred, taskData, "", "", nil); err != nil { + ctx, "GuestDeleteOnHostTask", guest, self.UserCred, taskData, self.GetId(), "", nil); err != nil { self.OnFail(ctx, guest, err.Error()) return } else { task.ScheduleRun(nil) - self.SetStageComplete(ctx, nil) } } func (self *GuestDeleteBackupTask) OnCancelBlockJobsFailed(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { self.OnFail(ctx, guest, data.String()) } + +func (self *GuestDeleteBackupTask) OnDeleteOnHost(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { + if jsonutils.QueryBoolean(self.Params, "create", false) { + self.OnDeleteBackupComplete(ctx, guest, data) + self.SetStage("OnCreateNewBackup", nil) + + params := jsonutils.NewDict() + params.Set("reconcile_backup", jsonutils.JSONTrue) + _, err := guest.StartGuestCreateBackupTask(ctx, self.UserCred, self.GetId(), params) + if err != nil { + self.onCreateNewBackupFailed(ctx, guest, err.Error()) + } + } else { + self.TaskComplete(ctx, guest, data) + } +} + +func (self *GuestDeleteBackupTask) OnCreateNewBackup(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { + logclient.AddActionLogWithContext(ctx, guest, logclient.ACT_CREATE_BACKUP, "", self.UserCred, false) + db.OpsLog.LogEvent(guest, db.ACT_CREATE_BACKUP, "", self.UserCred) + self.SetStageComplete(ctx, nil) +} + +func (self *GuestDeleteBackupTask) OnCreateNewBackupFailed(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { + self.onCreateNewBackupFailed(ctx, guest, data.String()) +} + +func (self *GuestDeleteBackupTask) onCreateNewBackupFailed(ctx context.Context, guest *models.SGuest, reason string) { + logclient.AddActionLogWithContext(ctx, guest, logclient.ACT_CREATE_BACKUP, reason, self.UserCred, false) + db.OpsLog.LogEvent(guest, db.ACT_CREATE_BACKUP_FAILED, reason, self.UserCred) + self.SetStageFailed(ctx, reason) +} + +func (self *GuestDeleteBackupTask) OnDeleteOnHostFailed(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { + self.OnFail(ctx, guest, data.String()) +} + +func (self *GuestDeleteBackupTask) OnDeleteBackupComplete(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { + logclient.AddActionLogWithContext(ctx, guest, logclient.ACT_DELETE_BACKUP, "", self.UserCred, true) + db.OpsLog.LogEvent(guest, db.ACT_DELETE_BACKUP, "", self.UserCred) +} + +func (self *GuestDeleteBackupTask) TaskComplete(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { + self.OnDeleteBackupComplete(ctx, guest, data) + self.SetStageComplete(ctx, nil) +} diff --git a/pkg/compute/tasks/storage_cache_image_task.go b/pkg/compute/tasks/storage_cache_image_task.go index ed3379fc8d..e8267ec106 100644 --- a/pkg/compute/tasks/storage_cache_image_task.go +++ b/pkg/compute/tasks/storage_cache_image_task.go @@ -63,13 +63,22 @@ func (self *StorageCacheImageTask) OnRelinquishLeastUsedCachedImageComplete(ctx db.OpsLog.LogEvent(storageCache, db.ACT_CACHING_IMAGE, imageId, self.UserCred) - self.SetStage("on_image_cache_complete", nil) + self.SetStage("OnImageCacheComplete", nil) - host, _ := storageCache.GetHost() - err := host.GetHostDriver().CheckAndSetCacheImage(ctx, host, storageCache, self) + host, err := storageCache.GetHost() if err != nil { errData := taskman.Error2TaskData(err) self.OnImageCacheCompleteFailed(ctx, storageCache, errData) + } else if host != nil { + err := host.GetHostDriver().CheckAndSetCacheImage(ctx, host, storageCache, self) + if err != nil { + errData := taskman.Error2TaskData(err) + self.OnImageCacheCompleteFailed(ctx, storageCache, errData) + } + } else { + // host is nil, and err is nil + reason := jsonutils.NewString("storage cache failed get host") + self.OnImageCacheCompleteFailed(ctx, storageCache, reason) } } diff --git a/pkg/hostman/guestman/guestman.go b/pkg/hostman/guestman/guestman.go index 8da679b9df..637e404946 100644 --- a/pkg/hostman/guestman/guestman.go +++ b/pkg/hostman/guestman/guestman.go @@ -72,6 +72,10 @@ type SGuestManager struct { GuestStartWorker *appsrv.SWorkerManager isLoaded bool + + // dirty servers chan + dirtyServers []*SKVMGuestInstance + dirtyServersChan chan struct{} } func NewGuestManager(host hostutils.IHost, serversPath string) *SGuestManager { @@ -86,6 +90,8 @@ func NewGuestManager(host hostutils.IHost, serversPath string) *SGuestManager { manager.StartCpusetBalancer() manager.LoadExistingGuests() manager.host.StartDHCPServer() + manager.dirtyServersChan = make(chan struct{}) + manager.dirtyServers = make([]*SKVMGuestInstance, 0) return manager } @@ -111,7 +117,7 @@ func (m *SGuestManager) SaveServer(sid string, s *SKVMGuestInstance) { m.Servers.Store(sid, s) } -func (m *SGuestManager) Bootstrap() { +func (m *SGuestManager) Bootstrap() chan struct{} { if m.isLoaded || len(m.ServersPath) == 0 { log.Errorln("Guestman bootstrap has been called!!!!!") } else { @@ -123,6 +129,7 @@ func (m *SGuestManager) Bootstrap() { m.OnLoadExistingGuestsComplete() } } + return m.dirtyServersChan } func (m *SGuestManager) VerifyExistingGuests(pendingDelete bool) { @@ -169,7 +176,7 @@ func (m *SGuestManager) OnVerifyExistingGuestsSucc(servers []jsonutils.JSONObjec } else { for id, server := range m.CandidateServers { m.UnknownServers.Store(id, server) - go m.RequestVerifyDirtyServer(server) + m.dirtyServers = append(m.dirtyServers, server) log.Errorf("Server %s not found on this host", server.GetName()) m.RemoveCandidateServer(server) } @@ -189,14 +196,26 @@ func (m *SGuestManager) OnLoadExistingGuestsComplete() { log.Infof("Load existing guests complete...") err := m.host.PutHostOnline() if err != nil { - log.Errorln(err) + log.Fatalf("put host online failed %s", err) } + go m.verifyDirtyServers() + if !options.HostOptions.EnableCpuBinding { m.ClenaupCpuset() } } +func (m *SGuestManager) verifyDirtyServers() { + select { + case <-m.dirtyServersChan: + } + for i := 0; i < len(m.dirtyServers); i++ { + go m.RequestVerifyDirtyServer(m.dirtyServers[i]) + } + m.dirtyServers = nil +} + func (m *SGuestManager) ClenaupCpuset() { m.Servers.Range(func(k, v interface{}) bool { guest := v.(*SKVMGuestInstance) diff --git a/pkg/hostman/host_services.go b/pkg/hostman/host_services.go index 8ff92a9e3e..a6ce693345 100644 --- a/pkg/hostman/host_services.go +++ b/pkg/hostman/host_services.go @@ -91,6 +91,7 @@ func (host *SHostService) RunService() { log.Fatalf(err.Error()) } + var guestChan chan struct{} guestman.Init(hostInstance, options.HostOptions.ServersPath) app_common.InitAuth(&options.HostOptions.CommonOptions, func() { log.Infof("Auth complete!!") @@ -100,7 +101,7 @@ func (host *SHostService) RunService() { } hostInstance.StartRegister(2, func() { - guestman.GetGuestManager().Bootstrap() + guestChan = guestman.GetGuestManager().Bootstrap() // hostmetrics after guestmanager bootstrap hostmetrics.Init() hostmetrics.Start() @@ -119,6 +120,7 @@ func (host *SHostService) RunService() { "CleanRecycleDiskFiles", 1, 3, 0, 0, storageman.CleanRecycleDiskfiles, false) cronManager.Start() + close(guestChan) app_common.ServeForeverWithCleanup(app, &options.HostOptions.BaseOptions, func() { hostinfo.Stop() storageman.Stop() diff --git a/pkg/mcclient/modules/mod_scheduler.go b/pkg/mcclient/modules/mod_scheduler.go index c6c350ad6c..3facd63df7 100644 --- a/pkg/mcclient/modules/mod_scheduler.go +++ b/pkg/mcclient/modules/mod_scheduler.go @@ -205,7 +205,7 @@ func (this *SchedulerManager) HistoryShow(s *mcclient.ClientSession, id string, return modulebase.Post(this.ResourceManager, s, url, params, "history") } -func (this *SchedulerManager) CleanCache(s *mcclient.ClientSession, hostId, sessionId string) error { +func (this *SchedulerManager) CleanCache(s *mcclient.ClientSession, hostId, sessionId string, sync bool) error { url := newSchedURL("clean-cache") if len(hostId) > 0 { url = fmt.Sprintf("%s/%s", url, hostId) @@ -213,6 +213,10 @@ func (this *SchedulerManager) CleanCache(s *mcclient.ClientSession, hostId, sess if len(sessionId) > 0 { url = fmt.Sprintf("%s?session=%s", url, sessionId) } + if sync { + url = fmt.Sprintf("%s?sync_clean=true", url) + } + resp, err := modulebase.RawRequest(this.ResourceManager, s, "POST", url, nil, nil) if err != nil { return err diff --git a/pkg/mcclient/modules/mod_servers.go b/pkg/mcclient/modules/mod_servers.go index b364d2234b..876705b6e7 100644 --- a/pkg/mcclient/modules/mod_servers.go +++ b/pkg/mcclient/modules/mod_servers.go @@ -104,7 +104,7 @@ func init() { "Created_at", "Group_name", "Group_id", "Hypervisor", "os_type", "expired_at"}, - []string{"Host", "Tenant", "is_system", "auto_delete_at"})} + []string{"Host", "Tenant", "is_system", "auto_delete_at", "backup_host_name"})} registerCompute(&Servers) } diff --git a/pkg/mcclient/options/servers.go b/pkg/mcclient/options/servers.go index a06f12c86e..2aa2fd941f 100644 --- a/pkg/mcclient/options/servers.go +++ b/pkg/mcclient/options/servers.go @@ -150,9 +150,10 @@ type ServerConfigs struct { Host string `help:"Preferred host where virtual server should be created" json:"prefer_host"` BackupHost string `help:"Perfered host where virtual backup server should be created"` - Hypervisor string `help:"Hypervisor type" choices:"kvm|esxi|baremetal|container|aliyun|azure|qcloud|aws|huawei|openstack|ucloud|zstack|google|ctyun"` - ResourceType string `help:"Resource type" choices:"shared|prepaid|dedicated"` - Backup bool `help:"Create server with backup server"` + Hypervisor string `help:"Hypervisor type" choices:"kvm|esxi|baremetal|container|aliyun|azure|qcloud|aws|huawei|openstack|ucloud|zstack|google|ctyun"` + ResourceType string `help:"Resource type" choices:"shared|prepaid|dedicated"` + Backup bool `help:"Create server with backup server"` + AutoSwitchToBackupOnHostDown bool `help:"Auto switch to backup server on host down"` Schedtag []string `help:"Schedule policy, key = aggregate name, value = require|exclude|prefer|avoid" metavar:""` Disk []string `help:"Disk descriptions" nargs:"+"` diff --git a/pkg/scheduler/handler/handler.go b/pkg/scheduler/handler/handler.go index 5389cddfa7..69cd39ee67 100644 --- a/pkg/scheduler/handler/handler.go +++ b/pkg/scheduler/handler/handler.go @@ -353,6 +353,15 @@ func getSessionId(c *gin.Context) string { return jsonutils.GetAnyString(query, []string{"session", "session_id"}) } +func isSyncCleanSchedCache(c *gin.Context) bool { + query, err := jsonutils.ParseQueryString(c.Request.URL.RawQuery) + if err != nil { + log.Warningf("not found sync clean cache in query") + return false + } + return jsonutils.QueryBoolean(query, "sync_clean", false) +} + func doCleanHostCache(c *gin.Context, hostID string) { sid := getSessionId(c) args, err := newExpireArgsByHostIDs([]string{hostID}, sid) @@ -365,7 +374,7 @@ func doCleanHostCache(c *gin.Context, hostID string) { } func doCleanHostCacheByArgs(c *gin.Context, args *api.ExpireArgs) { - result, err := schedman.Expire(args) + result, err := schedman.Expire(args, isSyncCleanSchedCache(c)) if err != nil { c.AbortWithError(http.StatusBadRequest, err) return diff --git a/pkg/scheduler/manager/expire_queue.go b/pkg/scheduler/manager/expire_queue.go index f8b0a32cf2..3672ea193c 100644 --- a/pkg/scheduler/manager/expire_queue.go +++ b/pkg/scheduler/manager/expire_queue.go @@ -29,12 +29,15 @@ import ( type ExpireManager struct { expireChannel chan *api.ExpireArgs stopCh <-chan struct{} + + mergeLock *sync.Mutex } func NewExpireManager(stopCh <-chan struct{}) *ExpireManager { return &ExpireManager{ expireChannel: make(chan *api.ExpireArgs, o.GetOptions().ExpireQueueMaxLength), stopCh: stopCh, + mergeLock: new(sync.Mutex), } } @@ -42,6 +45,10 @@ func (e *ExpireManager) Add(expireArgs *api.ExpireArgs) { e.expireChannel <- expireArgs } +func (e *ExpireManager) Trigger() { + e.batchMergeExpire() +} + type expireHost struct { Id string SessionId string @@ -56,83 +63,14 @@ func newExpireHost(id string, sid string) *expireHost { func (e *ExpireManager) Run() { t := time.Tick(u.ToDuration(o.GetOptions().ExpireQueueConsumptionPeriod)) - - waitTimeOut := func(wg *sync.WaitGroup, timeout time.Duration) bool { - ch := make(chan struct{}) - go func() { - wg.Wait() - close(ch) - }() - select { - case <-ch: - return true - case <-time.After(timeout): - return false - } - } - - batchMergeExpire := func() { - expireRequestNumber := len(e.expireChannel) - // If the expireRequestNumber then return right now. - if expireRequestNumber <= 0 { - return - } - dirtyHostSets := sets.NewString() - dirtyBaremetalSets := sets.NewString() - dirtyHosts := make([]*expireHost, 0) - dirtyBaremetals := make([]*expireHost, 0) - // Merge all same host. - for i := 0; i < expireRequestNumber; i++ { - expireArgs := <-e.expireChannel - log.V(4).Infof("Get expireArgs from channel: %#v", expireArgs) - dirtyHostSets.Insert(expireArgs.DirtyHosts...) - for _, host := range expireArgs.DirtyHosts { - dirtyHosts = append(dirtyHosts, newExpireHost(host, expireArgs.SessionId)) - } - dirtyBaremetalSets.Insert(expireArgs.DirtyBaremetals...) - for _, baremetal := range expireArgs.DirtyBaremetals { - dirtyBaremetals = append(dirtyBaremetals, newExpireHost(baremetal, expireArgs.SessionId)) - } - } - log.V(4).Infof("batchMergeExpire dirtyHosts: %v, dirtyBaremetals: %v", dirtyHosts, dirtyBaremetals) - wg := &sync.WaitGroup{} - wg.Add(2) - go func() { - defer wg.Done() - //dirtyHosts = notInSession(dirtyHosts, "host") - if len(dirtyHosts) > 0 { - log.V(10).Debugf("CleanDirty Hosts: %v\n", dirtyHosts) - if _, err := schedManager.CandidateManager.Reload("host", dirtyHostSets.List()); err != nil { - log.Errorf("Clean dirty hosts %v: %v", dirtyHosts, err) - } - schedManager.HistoryManager.CancelCandidatesPendingUsage(dirtyHosts) - } - }() - - go func() { - defer wg.Done() - //dirtyBaremetals = notInSession(dirtyBaremetals, "baremetal") - if len(dirtyBaremetals) > 0 { - log.V(10).Debugf("CleanDirty Baremetals: %v\n", dirtyBaremetals) - if _, err := schedManager.CandidateManager.Reload("baremetal", dirtyBaremetalSets.List()); err != nil { - log.Errorf("Clean dirty baremetals %v: %v", dirtyBaremetals, err) - } - schedManager.HistoryManager.CancelCandidatesPendingUsage(dirtyBaremetals) - } - }() - if ok := waitTimeOut(wg, u.ToDuration(o.GetOptions().ExpireQueueConsumptionTimeout)); !ok { - log.Errorln("time out reload data.") - } - } - // Watching the expires. for { select { case <-t: - batchMergeExpire() + e.batchMergeExpire() case <-e.stopCh: // update all the expire before return - batchMergeExpire() + e.batchMergeExpire() close(e.expireChannel) e.expireChannel = nil log.Errorln("expire manager EXIT!") @@ -140,10 +78,80 @@ func (e *ExpireManager) Run() { default: // if expire number is bigger then 80 then update if len(e.expireChannel) >= o.GetOptions().ExpireQueueDealLength { - batchMergeExpire() + e.batchMergeExpire() } else { time.Sleep(1 * time.Second) } } } } + +func (e *ExpireManager) batchMergeExpire() { + e.mergeLock.Lock() + defer e.mergeLock.Unlock() + expireRequestNumber := len(e.expireChannel) + // If the expireRequestNumber then return right now. + if expireRequestNumber <= 0 { + return + } + dirtyHostSets := sets.NewString() + dirtyBaremetalSets := sets.NewString() + dirtyHosts := make([]*expireHost, 0) + dirtyBaremetals := make([]*expireHost, 0) + // Merge all same host. + for i := 0; i < expireRequestNumber; i++ { + expireArgs := <-e.expireChannel + log.V(4).Infof("Get expireArgs from channel: %#v", expireArgs) + dirtyHostSets.Insert(expireArgs.DirtyHosts...) + for _, host := range expireArgs.DirtyHosts { + dirtyHosts = append(dirtyHosts, newExpireHost(host, expireArgs.SessionId)) + } + dirtyBaremetalSets.Insert(expireArgs.DirtyBaremetals...) + for _, baremetal := range expireArgs.DirtyBaremetals { + dirtyBaremetals = append(dirtyBaremetals, newExpireHost(baremetal, expireArgs.SessionId)) + } + } + log.V(4).Infof("batchMergeExpire dirtyHosts: %v, dirtyBaremetals: %v", dirtyHosts, dirtyBaremetals) + wg := &sync.WaitGroup{} + wg.Add(2) + go func() { + defer wg.Done() + //dirtyHosts = notInSession(dirtyHosts, "host") + if len(dirtyHosts) > 0 { + log.V(10).Debugf("CleanDirty Hosts: %v\n", dirtyHosts) + if _, err := schedManager.CandidateManager.Reload("host", dirtyHostSets.List()); err != nil { + log.Errorf("Clean dirty hosts %v: %v", dirtyHosts, err) + } + schedManager.HistoryManager.CancelCandidatesPendingUsage(dirtyHosts) + } + }() + + go func() { + defer wg.Done() + //dirtyBaremetals = notInSession(dirtyBaremetals, "baremetal") + if len(dirtyBaremetals) > 0 { + log.V(10).Debugf("CleanDirty Baremetals: %v\n", dirtyBaremetals) + if _, err := schedManager.CandidateManager.Reload("baremetal", dirtyBaremetalSets.List()); err != nil { + log.Errorf("Clean dirty baremetals %v: %v", dirtyBaremetals, err) + } + schedManager.HistoryManager.CancelCandidatesPendingUsage(dirtyBaremetals) + } + }() + if ok := e.waitTimeOut(wg, u.ToDuration(o.GetOptions().ExpireQueueConsumptionTimeout)); !ok { + log.Errorln("time out reload data.") + } +} + +func (e *ExpireManager) waitTimeOut(wg *sync.WaitGroup, timeout time.Duration) bool { + ch := make(chan struct{}) + go func() { + wg.Wait() + close(ch) + }() + select { + case <-ch: + return true + case <-time.After(timeout): + return false + } +} diff --git a/pkg/scheduler/manager/manager.go b/pkg/scheduler/manager/manager.go index 0877c23d30..80469a9c78 100644 --- a/pkg/scheduler/manager/manager.go +++ b/pkg/scheduler/manager/manager.go @@ -136,8 +136,11 @@ func GetCandidateManager() *data_manager.CandidateManager { return schedManager.CandidateManager } -func Expire(expireArgs *api.ExpireArgs) (*api.ExpireResult, error) { +func Expire(expireArgs *api.ExpireArgs, trigger bool) (*api.ExpireResult, error) { schedManager.ExpireManager.Add(expireArgs) + if trigger { + schedManager.ExpireManager.Trigger() + } return &api.ExpireResult{}, nil }