diff --git a/pkg/cloudcommon/db/quotas/handler.go b/pkg/cloudcommon/db/quotas/handler.go index d250e7c129..0b43eb10e8 100644 --- a/pkg/cloudcommon/db/quotas/handler.go +++ b/pkg/cloudcommon/db/quotas/handler.go @@ -88,7 +88,7 @@ func AddQuotaHandler(manager *SQuotaBaseManager, prefix string, app *appsrv.Appl auth.Authenticate(checkQuotaHanlder), nil, "check_quota", nil)*/ } -func (manager *SQuotaBaseManager) queryQuota(ctx context.Context, scope rbacutils.TRbacScope, ownerId mcclient.IIdentityProvider, platforma []string) (*jsonutils.JSONDict, IQuota, error) { +func (manager *SQuotaBaseManager) queryQuota(ctx context.Context, scope rbacutils.TRbacScope, ownerId mcclient.IIdentityProvider, platforma []string, refresh bool) (*jsonutils.JSONDict, IQuota, error) { ret := jsonutils.NewDict() quota := manager.newQuota() @@ -101,9 +101,9 @@ func (manager *SQuotaBaseManager) queryQuota(ctx context.Context, scope rbacutil if err != nil { return nil, nil, err } - if usage.IsEmpty() { + if usage.IsEmpty() || refresh { usageChan := make(chan IQuota) - manager.PostUsageJob(scope, ownerId, nil, usageChan, false) + manager.PostUsageJob(scope, ownerId, nil, usageChan, false, true) usage = <-usageChan } @@ -157,7 +157,7 @@ func (manager *SQuotaBaseManager) getQuotaHanlder(ctx context.Context, w http.Re scope = rbacutils.ScopeProject } - quota, _, err := manager.queryQuota(ctx, scope, ownerId, nil) + quota, _, err := manager.queryQuota(ctx, scope, ownerId, nil, true) if err != nil { httperrors.GeneralServerError(w, err) @@ -455,7 +455,7 @@ func (manager *SQuotaBaseManager) listQuotas(ctx context.Context, targetDomainId scope = rbacutils.ScopeDomain } platform := strings.Split(platformStr, nameSeparator) - quota, _, err := manager.queryQuota(ctx, scope, &owner, platform) + quota, _, err := manager.queryQuota(ctx, scope, &owner, platform, false) if err != nil { log.Errorf("query quota for %s fail %s", getMemoryStoreKey(scope, &owner, platform), err) continue @@ -494,7 +494,7 @@ func (manager *SQuotaBaseManager) listQuotas(ctx context.Context, targetDomainId DomainId: targetDomainId, } platform := []string{} - quota, _, err := manager.queryQuota(ctx, scope, &owner, platform) + quota, _, err := manager.queryQuota(ctx, scope, &owner, platform, false) if err != nil { return nil, httperrors.NewInternalServerError("query domain initial quotas %s", err) } diff --git a/pkg/cloudcommon/db/quotas/quotas.go b/pkg/cloudcommon/db/quotas/quotas.go index 445be41f34..a6f974509a 100644 --- a/pkg/cloudcommon/db/quotas/quotas.go +++ b/pkg/cloudcommon/db/quotas/quotas.go @@ -87,7 +87,7 @@ func (manager *SQuotaBaseManager) _cancelPendingUsage(ctx context.Context, userC } // update usage - manager.PostUsageJob(scope, ownerId, platform, nil, false) + manager.PostUsageJob(scope, ownerId, platform, nil, false, false) return err } diff --git a/pkg/cloudcommon/db/quotas/usageworker.go b/pkg/cloudcommon/db/quotas/usageworker.go index 930a44168e..78ae712605 100644 --- a/pkg/cloudcommon/db/quotas/usageworker.go +++ b/pkg/cloudcommon/db/quotas/usageworker.go @@ -34,9 +34,11 @@ import ( ) var ( - usageCalculateWorker = appsrv.NewWorkerManager("usageCalculateWorker", 1, 1024, true) - usageDirtyMap = make(map[string]bool, 0) - usageDirtyMapLock = &sync.Mutex{} + usageCalculateWorker = appsrv.NewWorkerManager("usageCalculateWorker", 1, 1024, true) + realTimeUsageCalculateWorker = appsrv.NewWorkerManager("realTimeUsageCalculateWorker", 1, 1024, true) + + usageDirtyMap = make(map[string]bool, 0) + usageDirtyMapLock = &sync.Mutex{} ) type sUsageCalculateJob struct { @@ -70,11 +72,18 @@ func isDirty(key string) bool { return false } -func (manager *SQuotaBaseManager) PostUsageJob(scope rbacutils.TRbacScope, ownerId mcclient.IIdentityProvider, platform []string, usageChan chan IQuota, cleanEmpty bool) { +func (manager *SQuotaBaseManager) PostUsageJob(scope rbacutils.TRbacScope, ownerId mcclient.IIdentityProvider, platform []string, usageChan chan IQuota, cleanEmpty bool, realTime bool) { key := getMemoryStoreKey(scope, ownerId, platform) setDirty(key) - usageCalculateWorker.Run(func() { + var worker *appsrv.SWorkerManager + if realTime { + worker = realTimeUsageCalculateWorker + } else { + worker = usageCalculateWorker + } + + worker.Run(func() { ctx := context.Background() if !isDirty(key) { @@ -167,6 +176,6 @@ func (manager *SQuotaBaseManager) CalculateQuotaUsages(ctx context.Context, user } platforms := strings.Split(platform, nameSeparator) // log.Debugf("PostUsageJob %s %s %s", scope, owner, platforms) - manager.PostUsageJob(scope, &owner, platforms, nil, true) + manager.PostUsageJob(scope, &owner, platforms, nil, true, false) } } diff --git a/pkg/cloudcommon/db/taskman/interface.go b/pkg/cloudcommon/db/taskman/interface.go index e0b779d665..998e812749 100644 --- a/pkg/cloudcommon/db/taskman/interface.go +++ b/pkg/cloudcommon/db/taskman/interface.go @@ -21,6 +21,7 @@ import ( "yunion.io/x/jsonutils" "yunion.io/x/onecloud/pkg/cloudcommon" + "yunion.io/x/onecloud/pkg/cloudcommon/db/quotas" "yunion.io/x/onecloud/pkg/mcclient" ) @@ -37,4 +38,7 @@ type ITask interface { SetStageComplete(ctx context.Context, data *jsonutils.JSONDict) SetStageFailed(ctx context.Context, reason string) + + GetPendingUsage(quota quotas.IQuota) error + ClearPendingUsage() error } diff --git a/pkg/compute/tasks/disk_batch_create_task.go b/pkg/compute/tasks/disk_batch_create_task.go index cfbe68eb9e..269d2cd82b 100644 --- a/pkg/compute/tasks/disk_batch_create_task.go +++ b/pkg/compute/tasks/disk_batch_create_task.go @@ -48,6 +48,12 @@ func (self *DiskBatchCreateTask) getNeedScheduleDisks(objs []db.IStandaloneModel return toSchedDisks } +func (self *DiskBatchCreateTask) clearPendingUsage(ctx context.Context, disk *models.SDisk) { + input, _ := self.GetCreateInput() + quotaPlatform := models.GetQuotaPlatformID(input.Hypervisor) + ClearTaskPendingUsage(ctx, self, rbacutils.ScopeProject, disk.GetOwnerId(), quotaPlatform) +} + func (self *DiskBatchCreateTask) OnInit(ctx context.Context, objs []db.IStandaloneModel, body jsonutils.JSONObject) { toSchedDisks := self.getNeedScheduleDisks(objs) if len(toSchedDisks) == 0 { @@ -82,6 +88,7 @@ func (self *DiskBatchCreateTask) OnScheduleFailCallback(ctx context.Context, obj self.SSchedTask.OnScheduleFailCallback(ctx, obj, reason) disk := obj.(*models.SDisk) log.Errorf("Schedule disk %s failed", disk.Name) + self.clearPendingUsage(ctx, disk) } func (self *DiskBatchCreateTask) SaveScheduleResult(ctx context.Context, obj IScheduleModel, candidate *schedapi.CandidateResource) { @@ -93,13 +100,14 @@ func (self *DiskBatchCreateTask) SaveScheduleResult(ctx context.Context, obj ISc log.Errorf("GetPendingUsage fail %s", err) } - input, _ := self.GetCreateInput() - quotaPlatform := models.GetQuotaPlatformID(input.Hypervisor) + // input, _ := self.GetCreateInput() + // quotaPlatform := models.GetQuotaPlatformID(input.Hypervisor) - quotaStorage := models.SQuota{Storage: disk.DiskSize} + // quotaStorage := models.SQuota{Storage: disk.DiskSize} onError := func(err error) { - models.QuotaManager.CancelPendingUsage(ctx, self.UserCred, rbacutils.ScopeProject, disk.GetOwnerId(), quotaPlatform, &pendingUsage, "aStorage) + // models.QuotaManager.CancelPendingUsage(ctx, self.UserCred, rbacutils.ScopeProject, disk.GetOwnerId(), quotaPlatform, &pendingUsage, "aStorage) + self.clearPendingUsage(ctx, disk) disk.SetStatus(self.UserCred, api.DISK_ALLOC_FAILED, err.Error()) self.SetStageFailed(ctx, err.Error()) db.OpsLog.LogEvent(disk, db.ACT_ALLOCATE_FAIL, err, self.UserCred) @@ -120,11 +128,12 @@ func (self *DiskBatchCreateTask) SaveScheduleResult(ctx context.Context, obj ISc } err = disk.SetStorageByHost(hostId, diskConfig, storageIds) if err != nil { - models.QuotaManager.CancelPendingUsage(ctx, self.UserCred, rbacutils.ScopeProject, disk.GetOwnerId(), quotaPlatform, &pendingUsage, "aStorage) - disk.SetStatus(self.UserCred, api.DISK_ALLOC_FAILED, err.Error()) - self.SetStageFailed(ctx, err.Error()) - db.OpsLog.LogEvent(disk, db.ACT_ALLOCATE_FAIL, err, self.UserCred) - notifyclient.NotifySystemError(disk.Id, disk.Name, api.DISK_ALLOC_FAILED, err.Error()) + onError(err) + // models.QuotaManager.CancelPendingUsage(ctx, self.UserCred, rbacutils.ScopeProject, disk.GetOwnerId(), quotaPlatform, &pendingUsage, "aStorage) + // disk.SetStatus(self.UserCred, api.DISK_ALLOC_FAILED, err.Error()) + // self.SetStageFailed(ctx, err.Error()) + // db.OpsLog.LogEvent(disk, db.ACT_ALLOCATE_FAIL, err, self.UserCred) + // notifyclient.NotifySystemError(disk.Id, disk.Name, api.DISK_ALLOC_FAILED, err.Error()) return } diff --git a/pkg/compute/tasks/guest_batch_create_task.go b/pkg/compute/tasks/guest_batch_create_task.go index 528e0ccd0d..d454847fe0 100644 --- a/pkg/compute/tasks/guest_batch_create_task.go +++ b/pkg/compute/tasks/guest_batch_create_task.go @@ -44,6 +44,15 @@ func (self *GuestBatchCreateTask) GetCreateInput() (*api.ServerCreateInput, erro return input, err } +func (self *GuestBatchCreateTask) clearPendingUsage(ctx context.Context, guest *models.SGuest) { + platform := make([]string, 0) + input, _ := self.GetCreateInput() + if len(input.Hypervisor) > 0 { + platform = models.GetDriver(input.Hypervisor).GetQuotaPlatformID() + } + ClearTaskPendingUsage(ctx, self, rbacutils.ScopeProject, guest.GetOwnerId(), platform) +} + func (self *GuestBatchCreateTask) OnInit(ctx context.Context, objs []db.IStandaloneModel, body jsonutils.JSONObject) { StartScheduleObjects(ctx, self, objs) } @@ -54,6 +63,7 @@ func (self *GuestBatchCreateTask) OnScheduleFailCallback(ctx context.Context, ob if guest.DisableDelete.IsTrue() { guest.SetDisableDelete(self.UserCred, false) } + self.clearPendingUsage(ctx, guest) } func (self *GuestBatchCreateTask) SaveScheduleResultWithBackup(ctx context.Context, obj IScheduleModel, master, slave *schedapi.CandidateResource) { @@ -175,6 +185,7 @@ func (self *GuestBatchCreateTask) SaveScheduleResult(ctx context.Context, obj IS err = self.allocateGuestOnHost(ctx, guest, candidate) if err != nil { + self.clearPendingUsage(ctx, guest) db.OpsLog.LogEvent(guest, db.ACT_ALLOCATE_FAIL, err, self.UserCred) logclient.AddActionLogWithStartable(self, obj, logclient.ACT_ALLOCATE, err.Error(), self.GetUserCred(), false) notifyclient.NotifySystemError(guest.Id, guest.Name, api.VM_CREATE_FAILED, err.Error()) diff --git a/pkg/compute/tasks/pending_usage.go b/pkg/compute/tasks/pending_usage.go new file mode 100644 index 0000000000..f4feff7313 --- /dev/null +++ b/pkg/compute/tasks/pending_usage.go @@ -0,0 +1,45 @@ +// Copyright 2019 Yunion +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package tasks + +import ( + "context" + + "yunion.io/x/log" + "yunion.io/x/pkg/errors" + + "yunion.io/x/onecloud/pkg/cloudcommon/db/taskman" + "yunion.io/x/onecloud/pkg/compute/models" + "yunion.io/x/onecloud/pkg/mcclient" + "yunion.io/x/onecloud/pkg/util/rbacutils" +) + +func ClearTaskPendingUsage(ctx context.Context, task taskman.ITask, scope rbacutils.TRbacScope, ownerId mcclient.IIdentityProvider, platform []string) error { + pendingUsage := models.SQuota{} + err := task.GetPendingUsage(&pendingUsage) + if err != nil { + log.Errorf("GetPendingUsage fail %s", err) + return errors.Wrap(err, "task.GetPendingUsage") + } + err = models.QuotaManager.CancelPendingUsage(ctx, task.GetUserCred(), scope, ownerId, platform, &pendingUsage, &pendingUsage) + if err != nil { + return errors.Wrap(err, "models.QuotaManager.CancelPendingUsage") + } + err = task.ClearPendingUsage() + if err != nil { + return errors.Wrap(err, "task.ClearPendingUsage") + } + return nil +} diff --git a/pkg/hostman/diskutils/doc.go b/pkg/hostman/diskutils/doc.go index a23c20344e..53cb65d680 100644 --- a/pkg/hostman/diskutils/doc.go +++ b/pkg/hostman/diskutils/doc.go @@ -1 +1,15 @@ +// Copyright 2019 Yunion +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + package diskutils // import "yunion.io/x/onecloud/pkg/hostman/diskutils" diff --git a/pkg/hostman/hostdeployer/apis/doc.go b/pkg/hostman/hostdeployer/apis/doc.go index 577a579cfc..f8b4172da0 100644 --- a/pkg/hostman/hostdeployer/apis/doc.go +++ b/pkg/hostman/hostdeployer/apis/doc.go @@ -1 +1,15 @@ +// Copyright 2019 Yunion +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + package apis // import "yunion.io/x/onecloud/pkg/hostman/hostdeployer/apis" diff --git a/pkg/hostman/hostdeployer/deployclient/doc.go b/pkg/hostman/hostdeployer/deployclient/doc.go index 3778fee40b..9a4ea1ba2f 100644 --- a/pkg/hostman/hostdeployer/deployclient/doc.go +++ b/pkg/hostman/hostdeployer/deployclient/doc.go @@ -1 +1,15 @@ +// Copyright 2019 Yunion +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + package deployclient // import "yunion.io/x/onecloud/pkg/hostman/hostdeployer/deployclient" diff --git a/pkg/hostman/hostdeployer/deployserver/doc.go b/pkg/hostman/hostdeployer/deployserver/doc.go index e1f3eaba12..6c1235787b 100644 --- a/pkg/hostman/hostdeployer/deployserver/doc.go +++ b/pkg/hostman/hostdeployer/deployserver/doc.go @@ -1 +1,15 @@ +// Copyright 2019 Yunion +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + package deployserver // import "yunion.io/x/onecloud/pkg/hostman/hostdeployer/deployserver" diff --git a/pkg/hostman/storageman/storagehandler/doc.go b/pkg/hostman/storageman/storagehandler/doc.go index ffde444e3e..2ab0804706 100644 --- a/pkg/hostman/storageman/storagehandler/doc.go +++ b/pkg/hostman/storageman/storagehandler/doc.go @@ -1 +1,15 @@ +// Copyright 2019 Yunion +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + package storagehandler // import "yunion.io/x/onecloud/pkg/hostman/storageman/storagehandler" diff --git a/pkg/hostman/system_service/host_deployer.go b/pkg/hostman/system_service/host_deployer.go index 7bb656a288..ec17438caa 100644 --- a/pkg/hostman/system_service/host_deployer.go +++ b/pkg/hostman/system_service/host_deployer.go @@ -1,3 +1,17 @@ +// Copyright 2019 Yunion +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + package system_service type SHostDeployer struct { diff --git a/pkg/util/sysutils/kvm_test.go b/pkg/util/sysutils/kvm_test.go index 96365adfaf..3d453b1468 100644 --- a/pkg/util/sysutils/kvm_test.go +++ b/pkg/util/sysutils/kvm_test.go @@ -1,3 +1,17 @@ +// Copyright 2019 Yunion +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + package sysutils import (