From 7c4caed98457eb89d64a6e819d42e5674ab67b56 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=B1=88=E8=BD=A9?= Date: Thu, 18 Oct 2018 12:03:09 +0800 Subject: [PATCH 01/23] =?UTF-8?q?=E5=A4=84=E7=90=86=E5=8D=B8=E8=BD=BD?= =?UTF-8?q?=E7=A3=81=E7=9B=98=E5=A4=B1=E8=B4=A5=E5=90=8E=E6=83=85=E5=86=B5?= =?UTF-8?q?:=20=E9=87=8D=E6=96=B0=E6=8C=82=E8=BD=BD=E5=8E=9F=E6=9C=89?= =?UTF-8?q?=E7=A3=81=E7=9B=98=E5=88=B0=E8=99=9A=E6=8B=9F=E6=9C=BA=E4=B8=8A?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- pkg/compute/models/guests.go | 4 +++ pkg/compute/tasks/guest_detach_disk_task.go | 37 ++++++++++++++++++++- 2 files changed, 40 insertions(+), 1 deletion(-) diff --git a/pkg/compute/models/guests.go b/pkg/compute/models/guests.go index 2c8f06fe80..c52d7beaff 100644 --- a/pkg/compute/models/guests.go +++ b/pkg/compute/models/guests.go @@ -1669,6 +1669,10 @@ func (self *SGuest) getMaxDiskIndex() int8 { return int8(len(guestdisks)) } +func (self *SGuest) AttachDisk(disk *SDisk, userCred mcclient.TokenCredential, driver string, cache string, mountpoint string) error { + return self.attach2Disk(disk, userCred, driver, cache, mountpoint) +} + func (self *SGuest) attach2Disk(disk *SDisk, userCred mcclient.TokenCredential, driver string, cache string, mountpoint string) error { if self.isAttach2Disk(disk) { return fmt.Errorf("Guest has been attached to disk") diff --git a/pkg/compute/tasks/guest_detach_disk_task.go b/pkg/compute/tasks/guest_detach_disk_task.go index 0a6653750f..10deaa7799 100644 --- a/pkg/compute/tasks/guest_detach_disk_task.go +++ b/pkg/compute/tasks/guest_detach_disk_task.go @@ -35,6 +35,14 @@ func (self *GuestDetachDiskTask) OnInit(ctx context.Context, obj db.IStandaloneM return } + guestdisks := disk.GetGuestdisks() + if len(guestdisks) > 0 { + guestdisk := guestdisks[0] + self.Params.Add(jsonutils.NewString(guestdisk.Driver), "driver") + self.Params.Add(jsonutils.NewString(guestdisk.CacheMode), "cache") + self.Params.Add(jsonutils.NewString(guestdisk.Mountpoint), "mountpoint") + } + guest.DetachDisk(ctx, disk, self.UserCred) if disk.Status == models.DISK_INIT { self.OnSyncConfigComplete(ctx, guest, nil) @@ -94,9 +102,36 @@ func (self *GuestDetachDiskTask) OnSyncConfigComplete(ctx context.Context, guest } } +func (self *GuestDetachDiskTask) OnSyncConfigCompleteFailed(ctx context.Context, obj db.IStandaloneModel, resion jsonutils.JSONObject) { + guest := obj.(*models.SGuest) + driver, _ := self.Params.GetString("driver") + cache, _ := self.Params.GetString("cache") + mountpoint, _ := self.Params.GetString("mountpoint") + diskId, _ := self.Params.GetString("disk_id") + objDisk, err := models.DiskManager.FetchById(diskId) + if err != nil { + self.OnTaskFail(ctx, guest, err) + return + } + disk := objDisk.(*models.SDisk) + db.OpsLog.LogEvent(disk, db.ACT_DETACH, resion.String(), self.UserCred) + err = guest.AttachDisk(disk, self.UserCred, driver, cache, mountpoint) + if err != nil { + self.OnTaskFail(ctx, guest, err) + return + } +} + func (self *GuestDetachDiskTask) OnTaskFail(ctx context.Context, guest *models.SGuest, err error) { + diskId, _ := self.Params.GetString("disk_id") + objDisk, err := models.DiskManager.FetchById(diskId) + if err != nil { + return + } + disk := objDisk.(*models.SDisk) + disk.SetStatus(self.UserCred, models.DISK_READY, err.Error()) self.SetStageFailed(ctx, err.Error()) - log.Errorf("Guest %s GuestDetachDiskTask failed %s", guest.Id, err.Error()) + log.Errorf("Guest %s disk %s GuestDetachDiskTask failed %s ", guest.Id, disk.Name, err.Error()) } func (self *GuestDetachDiskTask) OnDiskDeleteComplete(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) { From c96849e3343879e60925d2ab92bcfbbb6f43dff1 Mon Sep 17 00:00:00 2001 From: Yousong Zhou Date: Thu, 18 Oct 2018 07:14:39 +0000 Subject: [PATCH 02/23] guests: allow change nic bandwidth to 0 To restore the default no limit behavior --- pkg/compute/models/guestnetworks.go | 2 +- pkg/compute/models/guests.go | 4 ++-- 2 files changed, 3 insertions(+), 3 deletions(-) diff --git a/pkg/compute/models/guestnetworks.go b/pkg/compute/models/guestnetworks.go index b2cc9b7df5..bd17515f73 100644 --- a/pkg/compute/models/guestnetworks.go +++ b/pkg/compute/models/guestnetworks.go @@ -122,7 +122,7 @@ func (manager *SGuestnetworkManager) newGuestNetwork(ctx context.Context, userCr driver = "virtio" } gn.Driver = driver - if bwLimit > 0 { + if bwLimit >= 0 { gn.BwLimit = bwLimit } diff --git a/pkg/compute/models/guests.go b/pkg/compute/models/guests.go index 2c8f06fe80..09940bd3bc 100644 --- a/pkg/compute/models/guests.go +++ b/pkg/compute/models/guests.go @@ -2881,8 +2881,8 @@ func (self *SGuest) PerformChangeBandwidth(ctx context.Context, userCred mcclien return nil, httperrors.NewBadRequestError("Index Not fount or out of NIC index") } bandwidth, err := data.Int("bandwidth") - if err != nil || bandwidth <= 0 { - return nil, httperrors.NewBadRequestError("Bandwidth must be larger than 0") + if err != nil || bandwidth < 0 { + return nil, httperrors.NewBadRequestError("Bandwidth must be non-negative") } guestnic := &guestnics[index] if guestnic.BwLimit != int(bandwidth) { From 15a5c91cf17537ad51a7e2b81bb18577913715ee Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=B1=88=E8=BD=A9?= Date: Thu, 18 Oct 2018 12:11:39 +0800 Subject: [PATCH 03/23] =?UTF-8?q?=E5=8D=B8=E8=BD=BD=E7=A3=81=E7=9B=98?= =?UTF-8?q?=E6=97=B6=EF=BC=8C=E7=A3=81=E7=9B=98=E7=8A=B6=E6=80=81=E8=B7=9F?= =?UTF-8?q?=E9=9A=8F=E5=8F=98=E5=8C=96?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- pkg/compute/tasks/guest_detach_disk_task.go | 37 +++++++++++---------- 1 file changed, 19 insertions(+), 18 deletions(-) diff --git a/pkg/compute/tasks/guest_detach_disk_task.go b/pkg/compute/tasks/guest_detach_disk_task.go index 10deaa7799..abbe01ea0e 100644 --- a/pkg/compute/tasks/guest_detach_disk_task.go +++ b/pkg/compute/tasks/guest_detach_disk_task.go @@ -26,12 +26,12 @@ func (self *GuestDetachDiskTask) OnInit(ctx context.Context, obj db.IStandaloneM diskId, _ := self.Params.GetString("disk_id") objDisk, err := models.DiskManager.FetchById(diskId) if err != nil { - self.OnTaskFail(ctx, guest, err) + self.OnTaskFail(ctx, guest, nil, err) return } disk := objDisk.(*models.SDisk) if disk == nil { - self.OnTaskFail(ctx, guest, fmt.Errorf("Connot find disk %s", diskId)) + self.OnTaskFail(ctx, guest, nil, fmt.Errorf("Connot find disk %s", diskId)) return } @@ -48,6 +48,8 @@ func (self *GuestDetachDiskTask) OnInit(ctx context.Context, obj db.IStandaloneM self.OnSyncConfigComplete(ctx, guest, nil) return } + disk.SetStatus(self.UserCred, models.DISK_DETACHING, "Disk detach") + host := guest.GetHost() purge := false if host != nil && host.Status == models.HOST_DISABLED && jsonutils.QueryBoolean(self.Params, "purge", false) { @@ -55,13 +57,12 @@ func (self *GuestDetachDiskTask) OnInit(ctx context.Context, obj db.IStandaloneM } detachStatus, err := guest.GetDriver().GetDetachDiskStatus() if err != nil { - self.OnTaskFail(ctx, guest, err) + self.OnTaskFail(ctx, guest, disk, err) return } if utils.IsInStringArray(guest.Status, detachStatus) && !purge { self.SetStage("on_sync_config_complete", nil) guest.GetDriver().RequestDetachDisk(ctx, guest, self) - disk.SetStatus(self.UserCred, models.DISK_READY, "Disk detach") } else { self.OnSyncConfigComplete(ctx, guest, nil) } @@ -71,14 +72,15 @@ func (self *GuestDetachDiskTask) OnSyncConfigComplete(ctx context.Context, guest diskId, _ := self.Params.GetString("disk_id") objDisk, err := models.DiskManager.FetchById(diskId) if err != nil { - self.OnTaskFail(ctx, guest, err) + self.OnTaskFail(ctx, guest, nil, err) return } disk := objDisk.(*models.SDisk) if disk == nil { - self.OnTaskFail(ctx, guest, fmt.Errorf("Connot find disk %s", diskId)) + self.OnTaskFail(ctx, guest, nil, fmt.Errorf("Connot find disk %s", diskId)) return } + disk.SetDiskReady(ctx, self.UserCred, "") keepDisk := jsonutils.QueryBoolean(self.Params, "keep_disk", true) host := guest.GetHost() purge := false @@ -89,12 +91,14 @@ func (self *GuestDetachDiskTask) OnSyncConfigComplete(ctx context.Context, guest db.OpsLog.LogEvent(disk, db.ACT_DELETE, "", self.UserCred) disk.RealDelete(ctx, self.UserCred) self.SetStageComplete(ctx, nil) - } else if (disk.Status == models.DISK_READY || !keepDisk) && disk.GetGuestDiskCount() == 0 && disk.AutoDelete { + return + } + if !keepDisk && disk.GetGuestDiskCount() == 0 && disk.AutoDelete { self.SetStage("on_disk_delete_complete", nil) db.OpsLog.LogEvent(disk, db.ACT_DELETE, "", self.UserCred) err := guest.GetDriver().RequestDeleteDetachedDisk(ctx, disk, self, purge) if err != nil { - self.OnTaskFail(ctx, guest, err) + self.OnTaskFail(ctx, guest, disk, err) return } } else { @@ -110,28 +114,25 @@ func (self *GuestDetachDiskTask) OnSyncConfigCompleteFailed(ctx context.Context, diskId, _ := self.Params.GetString("disk_id") objDisk, err := models.DiskManager.FetchById(diskId) if err != nil { - self.OnTaskFail(ctx, guest, err) + self.OnTaskFail(ctx, guest, nil, err) return } disk := objDisk.(*models.SDisk) db.OpsLog.LogEvent(disk, db.ACT_DETACH, resion.String(), self.UserCred) + disk.SetDiskReady(ctx, self.UserCred, "") err = guest.AttachDisk(disk, self.UserCred, driver, cache, mountpoint) if err != nil { - self.OnTaskFail(ctx, guest, err) + self.OnTaskFail(ctx, guest, disk, err) return } } -func (self *GuestDetachDiskTask) OnTaskFail(ctx context.Context, guest *models.SGuest, err error) { - diskId, _ := self.Params.GetString("disk_id") - objDisk, err := models.DiskManager.FetchById(diskId) - if err != nil { - return +func (self *GuestDetachDiskTask) OnTaskFail(ctx context.Context, guest *models.SGuest, disk *models.SDisk, err error) { + if disk != nil { + disk.SetDiskReady(ctx, self.UserCred, "") } - disk := objDisk.(*models.SDisk) - disk.SetStatus(self.UserCred, models.DISK_READY, err.Error()) self.SetStageFailed(ctx, err.Error()) - log.Errorf("Guest %s disk %s GuestDetachDiskTask failed %s ", guest.Id, disk.Name, err.Error()) + log.Errorf("Guest %s GuestDetachDiskTask failed %s", guest.Id, err.Error()) } func (self *GuestDetachDiskTask) OnDiskDeleteComplete(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) { From cb92fb4e671783eb42701027a6dbd3487ed85a80 Mon Sep 17 00:00:00 2001 From: Qiu Jian Date: Thu, 18 Oct 2018 18:06:15 +0800 Subject: [PATCH 04/23] Remove local terraform code --- pkg/terraform/provider/doc.go | 1 - pkg/terraform/provider/provider.go | 1 - 2 files changed, 2 deletions(-) delete mode 100644 pkg/terraform/provider/doc.go delete mode 100644 pkg/terraform/provider/provider.go diff --git a/pkg/terraform/provider/doc.go b/pkg/terraform/provider/doc.go deleted file mode 100644 index f32b27848a..0000000000 --- a/pkg/terraform/provider/doc.go +++ /dev/null @@ -1 +0,0 @@ -package provider // import "yunion.io/x/onecloud/pkg/terraform/provider" diff --git a/pkg/terraform/provider/provider.go b/pkg/terraform/provider/provider.go deleted file mode 100644 index 4f504f6688..0000000000 --- a/pkg/terraform/provider/provider.go +++ /dev/null @@ -1 +0,0 @@ -package provider From 0101570c8686d4698ff5b0a8507408f12a9e36a8 Mon Sep 17 00:00:00 2001 From: Zexi Li Date: Thu, 18 Oct 2018 19:31:59 +0800 Subject: [PATCH 05/23] fix: scheduler start auth interface --- cmd/scheduler/app/server.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/cmd/scheduler/app/server.go b/cmd/scheduler/app/server.go index e3a8f93be8..7a3ff45fc2 100644 --- a/cmd/scheduler/app/server.go +++ b/cmd/scheduler/app/server.go @@ -64,7 +64,7 @@ func Run(s *SchedulerServer) error { debug := o.GetOptions().LogLevel == "debug" - auth.AsyncInit(s.AuthInfo, debug, true, startSched) + auth.AsyncInit(s.AuthInfo, debug, true, "", "", startSched) return startHTTP(s) } From 83f09321cab7a3b12b67ef83e4bf8db27f8d9bf5 Mon Sep 17 00:00:00 2001 From: Zexi Li Date: Fri, 19 Oct 2018 12:07:57 +0800 Subject: [PATCH 06/23] fix: kvm find vnc port index out of range --- pkg/compute/guestdrivers/kvm.go | 9 ++++----- 1 file changed, 4 insertions(+), 5 deletions(-) diff --git a/pkg/compute/guestdrivers/kvm.go b/pkg/compute/guestdrivers/kvm.go index 4b43f4f79d..fa50ef7e95 100644 --- a/pkg/compute/guestdrivers/kvm.go +++ b/pkg/compute/guestdrivers/kvm.go @@ -3,8 +3,8 @@ package guestdrivers import ( "context" "fmt" - "regexp" "strconv" + "strings" "yunion.io/x/jsonutils" "yunion.io/x/log" @@ -43,10 +43,9 @@ func (self *SKVMGuestDriver) OnDeleteGuestFinalCleanup(ctx context.Context, gues } func findVNCPort(results string) int { - reg := regexp.MustCompile(`(\d+\.\d+\.\d+\.\d+):([\d]+)`) - finds := reg.FindStringSubmatch(results) - log.Debugf("finds=%s", finds) - port, _ := strconv.Atoi(finds[2]) + vncInfo := strings.Split(results, "\n") + addrParts := strings.Split(vncInfo[1], ":") + port, _ := strconv.Atoi(addrParts[len(addrParts)-1]) return port } From db7183976b2202a0e4958e6be8df9b5d2fd10d1f Mon Sep 17 00:00:00 2001 From: Qiu Jian Date: Sat, 20 Oct 2018 13:14:38 +0800 Subject: [PATCH 07/23] =?UTF-8?q?=E4=BF=AE=E6=AD=A3=EF=BC=9A=E5=88=A0?= =?UTF-8?q?=E9=99=A4=E4=B8=BB=E6=9C=BA=E6=97=B6=EF=BC=8C=E5=9C=A8=E5=81=9C?= =?UTF-8?q?=E6=AD=A2=E4=B8=BB=E6=9C=BA=E6=97=B6=EF=BC=8C=E5=A6=82=E6=9E=9C?= =?UTF-8?q?=E5=AE=BF=E4=B8=BB=E6=9C=BA=E4=B8=8D=E5=93=8D=E5=BA=94=EF=BC=8C?= =?UTF-8?q?=E5=BA=94=E8=AF=A5=E5=AF=BC=E8=87=B4=E5=88=A0=E9=99=A4=E5=A4=B1?= =?UTF-8?q?=E8=B4=A5?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- pkg/compute/guestdrivers/virtualization.go | 12 ++++++------ pkg/compute/tasks/guest_delete_task.go | 6 +++++- 2 files changed, 11 insertions(+), 7 deletions(-) diff --git a/pkg/compute/guestdrivers/virtualization.go b/pkg/compute/guestdrivers/virtualization.go index 2ed28aca95..04c1307f88 100644 --- a/pkg/compute/guestdrivers/virtualization.go +++ b/pkg/compute/guestdrivers/virtualization.go @@ -115,12 +115,12 @@ func (self *SVirtualizedGuestDriver) StartGuestSyncstatusTask(guest *models.SGue } func (self *SVirtualizedGuestDriver) RequestStopGuestForDelete(ctx context.Context, guest *models.SGuest, task taskman.ITask) error { - guestStatus, _ := task.GetParams().GetString("guest_status") - if guestStatus == models.VM_RUNNING { - host := guest.GetHost() - if host != nil && host.Enabled && host.HostStatus == models.HOST_ONLINE && !jsonutils.QueryBoolean(task.GetParams(), "purge", false) { - return guest.StartGuestStopTask(ctx, task.GetUserCred(), true, task.GetTaskId()) - } + host := guest.GetHost() + if host != nil && host.Enabled && host.HostStatus == models.HOST_ONLINE { + return guest.StartGuestStopTask(ctx, task.GetUserCred(), true, task.GetTaskId()) + } + if !jsonutils.QueryBoolean(task.GetParams(), "purge", false) { + return fmt.Errorf("fail to contact host") } task.ScheduleRun(nil) return nil diff --git a/pkg/compute/tasks/guest_delete_task.go b/pkg/compute/tasks/guest_delete_task.go index 5b5a8e5d51..01c6ac6a3d 100644 --- a/pkg/compute/tasks/guest_delete_task.go +++ b/pkg/compute/tasks/guest_delete_task.go @@ -26,7 +26,11 @@ func init() { func (self *GuestDeleteTask) OnInit(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) { guest := obj.(*models.SGuest) self.SetStage("on_guest_stop_complete", nil) - guest.GetDriver().RequestStopGuestForDelete(ctx, guest, self) + err := guest.GetDriver().RequestStopGuestForDelete(ctx, guest, self) + if err != nil { + errMsg := jsonutils.NewString(err.Error()) + self.OnGuestStopCompleteFailed(ctx, obj, errMsg) + } } func (self *GuestDeleteTask) OnGuestStopComplete(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) { From 12ae5b6aca3af606b47e58e8d9afa1345006324e Mon Sep 17 00:00:00 2001 From: Qiu Jian Date: Sat, 20 Oct 2018 16:47:26 +0800 Subject: [PATCH 08/23] =?UTF-8?q?=E6=9B=B4=E6=96=B0vendor?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- Gopkg.lock | 8 +-- vendor/yunion.io/x/jsonutils/compond.go | 13 ++++ vendor/yunion.io/x/jsonutils/currency.go | 30 +++++++++ vendor/yunion.io/x/jsonutils/interface.go | 39 +++++++++++ vendor/yunion.io/x/jsonutils/jsonutils.go | 2 + vendor/yunion.io/x/jsonutils/unmarshal.go | 3 +- vendor/yunion.io/x/jsonutils/yamlutils.go | 64 +++++++++++-------- .../yunion.io/x/pkg/util/regutils/regutils.go | 12 ++++ 8 files changed, 140 insertions(+), 31 deletions(-) create mode 100644 vendor/yunion.io/x/jsonutils/compond.go create mode 100644 vendor/yunion.io/x/jsonutils/currency.go create mode 100644 vendor/yunion.io/x/jsonutils/interface.go diff --git a/Gopkg.lock b/Gopkg.lock index bcd51fb9d7..6a367673e1 100644 --- a/Gopkg.lock +++ b/Gopkg.lock @@ -1223,11 +1223,11 @@ [[projects]] branch = "master" - digest = "1:49ffc35ec8d3f7789393cd132acd359e8ac1f5d38c7a2b91c484b041840c62c0" + digest = "1:54554b3c72f4fcbd3c27c5137f9acd9d153d7e7150d77cd509af50c75fcb41b9" name = "yunion.io/x/jsonutils" packages = ["."] pruneopts = "UT" - revision = "d1290e94d4753c1748fc7c89f472a523cc0a5c08" + revision = "7079aada4c7e9a37e4c3b183bddb0ef684e038f2" [[projects]] branch = "master" @@ -1242,7 +1242,7 @@ [[projects]] branch = "master" - digest = "1:211ecaed7d1d87e5d216b598c3cd5b7c47977b73a06308e1b39c08fd03c9c417" + digest = "1:201bff9d99a538dd71fd8e55bbbfda8a301fd89676d0668a2cedc574e1aeaa03" name = "yunion.io/x/pkg" packages = [ "gotypes", @@ -1276,7 +1276,7 @@ "utils", ] pruneopts = "UT" - revision = "9d246215b0a167b153bcdbc7d75031ced7138cf8" + revision = "122c7d4ce76b4611b63d86af8e9bb2610ab55c20" [[projects]] branch = "master" diff --git a/vendor/yunion.io/x/jsonutils/compond.go b/vendor/yunion.io/x/jsonutils/compond.go new file mode 100644 index 0000000000..1ffc930c35 --- /dev/null +++ b/vendor/yunion.io/x/jsonutils/compond.go @@ -0,0 +1,13 @@ +package jsonutils + +func (val *JSONValue) isCompond() bool { + return false +} + +func (val *JSONDict) isCompond() bool { + return true +} + +func (val *JSONArray) isCompond() bool { + return true +} diff --git a/vendor/yunion.io/x/jsonutils/currency.go b/vendor/yunion.io/x/jsonutils/currency.go new file mode 100644 index 0000000000..0f2d2b77f8 --- /dev/null +++ b/vendor/yunion.io/x/jsonutils/currency.go @@ -0,0 +1,30 @@ +package jsonutils + +import ( + "fmt" + "strings" + "yunion.io/x/pkg/util/regutils" +) + +func normalizeUSCurrency(currency string) string { + return strings.Replace(currency, ",", "", -1) +} + +func normalizeEUCurrency(currency string) string { + commaPos := strings.IndexByte(currency, ',') + if commaPos >= 0 { + return fmt.Sprintf("%s.%s", strings.Replace(currency[:commaPos], ".", "", -1), currency[commaPos+1:]) + } else { + return strings.Replace(currency, ".", "", -1) + } +} + +func normalizeCurrencyString(currency string) string { + if regutils.MatchUSCurrency(currency) { + return normalizeUSCurrency(currency) + } + if regutils.MatchEUCurrency(currency) { + return normalizeEUCurrency(currency) + } + return currency +} diff --git a/vendor/yunion.io/x/jsonutils/interface.go b/vendor/yunion.io/x/jsonutils/interface.go new file mode 100644 index 0000000000..3ee19a10f6 --- /dev/null +++ b/vendor/yunion.io/x/jsonutils/interface.go @@ -0,0 +1,39 @@ +package jsonutils + +func (self *JSONValue) Interface() interface{} { + return nil +} + +func (self *JSONBool) Interface() interface{} { + return self.data +} + +func (self *JSONInt) Interface() interface{} { + return self.data +} + +func (self *JSONFloat) Interface() interface{} { + return self.data +} + +func (self *JSONString) Interface() interface{} { + return self.data +} + +func (self *JSONArray) Interface() interface{} { + ret := make([]interface{}, len(self.data)) + for i := 0; i < len(self.data); i += 1 { + ret[i] = self.data[i].Interface() + } + return ret +} + +func (self *JSONDict) Interface() interface{} { + mapping := make(map[string]interface{}) + + for k, v := range self.data { + mapping[k] = v.Interface() + } + + return mapping +} diff --git a/vendor/yunion.io/x/jsonutils/jsonutils.go b/vendor/yunion.io/x/jsonutils/jsonutils.go index d3321d1caf..2b9f289cfe 100644 --- a/vendor/yunion.io/x/jsonutils/jsonutils.go +++ b/vendor/yunion.io/x/jsonutils/jsonutils.go @@ -65,6 +65,8 @@ type JSONObject interface { Equals(obj JSONObject) bool unmarshalValue(val reflect.Value) error // IsZero() bool + Interface() interface{} + isCompond() bool } type JSONValue struct { diff --git a/vendor/yunion.io/x/jsonutils/unmarshal.go b/vendor/yunion.io/x/jsonutils/unmarshal.go index 79a7a2ac6a..9a1015f3da 100644 --- a/vendor/yunion.io/x/jsonutils/unmarshal.go +++ b/vendor/yunion.io/x/jsonutils/unmarshal.go @@ -321,7 +321,7 @@ func (this *JSONString) unmarshalValue(val reflect.Value) error { } val.SetInt(intVal) case reflect.Float32, reflect.Float64: - floatVal, err := strconv.ParseFloat(this.data, 64) + floatVal, err := strconv.ParseFloat(normalizeCurrencyString(this.data), 64) if err != nil { return err } @@ -348,6 +348,7 @@ func (this *JSONArray) unmarshalValue(val reflect.Value) error { if this.data != nil { array.Add(this.data...) } + val.Set(reflect.ValueOf(array)) return nil case JSONArrayPtrType, JSONObjectType: val.Set(reflect.ValueOf(this)) diff --git a/vendor/yunion.io/x/jsonutils/yamlutils.go b/vendor/yunion.io/x/jsonutils/yamlutils.go index 1daf0a63a3..fd2b978882 100644 --- a/vendor/yunion.io/x/jsonutils/yamlutils.go +++ b/vendor/yunion.io/x/jsonutils/yamlutils.go @@ -61,43 +61,49 @@ func parseYAMLDict(lines []string) (map[string]JSONObject, error) { } else { key := lines[i][0:keypos] val := strings.Trim(lines[i][keypos+1:], " ") + if len(val) > 0 && val != "|" { - o, e := Parse([]byte(val)) - if e != nil { - return dict, e - } else { - dict[key] = o - } + dict[key] = NewString(val) i++ } else { + sublines := make([]string, 0) j := i + 1 for j < len(lines) && len(strings.Trim(lines[j], " ")) == 0 { + sublines = append(sublines, "") j++ } - if j >= len(lines) || lines[j][0] != ' ' { - return dict, fmt.Errorf("Illformat") - } - indent := 0 - for indent < len(lines[j]) && lines[j][indent] == ' ' { - indent++ - } - sublines := make([]string, 0) - for j < len(lines) { - if indent >= len(lines[j]) && len(strings.Trim(lines[j], " ")) == 0 { - j++ - } else if indent < len(lines[j]) && len(strings.Trim(lines[j][:indent], " ")) == 0 { - sublines = append(sublines, lines[j][indent:]) - j++ - } else { - break + if j < len(lines) { + if lines[j][0] != ' ' { + return dict, fmt.Errorf("Illformat") + } + + indent := 0 + for indent < len(lines[j]) && lines[j][indent] == ' ' { + indent++ + } + + for j < len(lines) { + if indent >= len(lines[j]) && len(strings.Trim(lines[j], " ")) == 0 { + sublines = append(sublines, "") + j++ + } else if indent < len(lines[j]) && len(strings.Trim(lines[j][:indent], " ")) == 0 { + sublines = append(sublines, lines[j][indent:]) + j++ + } else { + break + } } } - o, e := parseYAMLLines(sublines) - if e != nil { - return dict, e + if val == "|" { + dict[key] = NewString(strings.Join(sublines, "\n")) } else { + o, e := parseYAMLLines(sublines) + if e != nil { + return dict, e + } dict[key] = o } + i = j } } @@ -192,8 +198,14 @@ func (this *JSONDict) yamlLines() []string { var ret = make([]string, 0) for _, key := range this.SortedKeys() { val := this.data[key] + if val.IsZero() { + switch val.(type) { + case *JSONString, *JSONDict, *JSONArray, *JSONValue: + continue + } + } lines := val.yamlLines() - if len(lines) == 1 { + if !val.isCompond() && len(lines) == 1 { ret = append(ret, fmt.Sprintf("%s: %s", key, lines[0])) } else { switch val.(type) { diff --git a/vendor/yunion.io/x/pkg/util/regutils/regutils.go b/vendor/yunion.io/x/pkg/util/regutils/regutils.go index 5ae4abb988..19fbb24f0f 100644 --- a/vendor/yunion.io/x/pkg/util/regutils/regutils.go +++ b/vendor/yunion.io/x/pkg/util/regutils/regutils.go @@ -32,6 +32,8 @@ var RFC2882_TIME_REG *regexp.Regexp var EMAIL_REG *regexp.Regexp var CHINA_MOBILE_REG *regexp.Regexp var FS_FORMAT_REG *regexp.Regexp +var US_CURRENCY_REG *regexp.Regexp +var EU_CURRENCY_REG *regexp.Regexp func init() { FUNCTION_REG = regexp.MustCompile(`^\w+\(.*\)$`) @@ -62,6 +64,8 @@ func init() { EMAIL_REG = regexp.MustCompile(`^[A-Za-z0-9._%+-]+@[A-Za-z0-9.-]+\.[A-Za-z]{2,4}$`) CHINA_MOBILE_REG = regexp.MustCompile(`^1[0-9-]{10}$`) FS_FORMAT_REG = regexp.MustCompile(`^(ext|fat|hfs|xfs|swap|ntfs|reiserfs|ufs|btrfs)`) + US_CURRENCY_REG = regexp.MustCompile(`^(\d{0,3}|((\d{1,3},)+\d{3}))(\.\d*)?$`) + EU_CURRENCY_REG = regexp.MustCompile(`^(\d{0,3}|((\d{1,3}\.)+\d{3}))(,\d*)?$`) } func MatchFunction(str string) bool { @@ -179,3 +183,11 @@ func MatchMobile(str string) bool { func MatchFS(str string) bool { return FS_FORMAT_REG.MatchString(str) } + +func MatchUSCurrency(str string) bool { + return US_CURRENCY_REG.MatchString(str) +} + +func MatchEUCurrency(str string) bool { + return EU_CURRENCY_REG.MatchString(str) +} From 0dbc93725dc5ff30b4e0cd836f69e694f552c2b1 Mon Sep 17 00:00:00 2001 From: Qiu Jian Date: Sat, 20 Oct 2018 22:32:43 +0800 Subject: [PATCH 09/23] =?UTF-8?q?=E4=BF=AE=E6=AD=A3=EF=BC=9A=E5=88=A9?= =?UTF-8?q?=E7=94=A8context.Timeout=EF=BC=8Cappsrv=E6=94=AF=E6=8C=81proces?= =?UTF-8?q?sTimeout=EF=BC=8C=E9=81=BF=E5=85=8Dyunionapi=E7=94=B1=E4=BA=8E?= =?UTF-8?q?=E5=90=8E=E7=AB=AF=E5=8D=A1=E4=BD=8F=E5=AF=BC=E8=87=B4yunionapi?= =?UTF-8?q?=E5=8D=A1=E4=BD=8F?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- pkg/appsrv/appsrv.go | 117 +++++++++-------------------------- pkg/appsrv/response.go | 87 ++++++++++++++++++++++++++ pkg/httperrors/errors.go | 5 ++ pkg/httperrors/httperrors.go | 4 ++ 4 files changed, 126 insertions(+), 87 deletions(-) create mode 100644 pkg/appsrv/response.go diff --git a/pkg/appsrv/appsrv.go b/pkg/appsrv/appsrv.go index 314bfc185f..f434bd44b4 100644 --- a/pkg/appsrv/appsrv.go +++ b/pkg/appsrv/appsrv.go @@ -10,79 +10,15 @@ import ( "time" "yunion.io/x/log" - "yunion.io/x/onecloud/pkg/appctx" - "yunion.io/x/onecloud/pkg/proxy" "yunion.io/x/pkg/trace" "yunion.io/x/pkg/utils" + + "yunion.io/x/onecloud/pkg/appctx" + "yunion.io/x/onecloud/pkg/httperrors" + "yunion.io/x/onecloud/pkg/proxy" + "yunion.io/x/onecloud/pkg/util/httputils" ) -type responseWriterResponse struct { - count int - err error -} - -type responseWriterChannel struct { - backend http.ResponseWriter - bodyChan chan []byte - bodyResp chan responseWriterResponse - statusChan chan int - statusResp chan bool -} - -func newResponseWriterChannel(backend http.ResponseWriter) responseWriterChannel { - return responseWriterChannel{backend: backend, - bodyChan: make(chan []byte), - bodyResp: make(chan responseWriterResponse), - statusChan: make(chan int), - statusResp: make(chan bool)} -} - -func (w *responseWriterChannel) Header() http.Header { - return w.backend.Header() -} - -func (w *responseWriterChannel) Write(bytes []byte) (int, error) { - w.bodyChan <- bytes - v := <-w.bodyResp - return v.count, v.err -} - -func (w *responseWriterChannel) WriteHeader(status int) { - w.statusChan <- status - <-w.statusResp -} - -func (w *responseWriterChannel) wait() { - stop := false - for !stop { - select { - case bytes, more := <-w.bodyChan: - // log.Print("Recive body ", len(bytes), " more ", more) - if more { - c, e := w.backend.Write(bytes) - w.bodyResp <- responseWriterResponse{count: c, err: e} - } else { - stop = true - } - case status, more := <-w.statusChan: - // log.Print("Recive status ", status, " more ", more) - if more { - w.backend.WriteHeader(status) - w.statusResp <- true - } else { - stop = true - } - } - } -} - -func (w *responseWriterChannel) closeChannels() { - close(w.bodyChan) - close(w.bodyResp) - close(w.statusChan) - close(w.statusResp) -} - type Application struct { name string context context.Context @@ -106,6 +42,7 @@ const ( DEFAULT_READ_TIMEOUT = 0 DEFAULT_READ_HEADER_TIMEOUT = 10 * time.Second DEFAULT_WRITE_TIMEOUT = 0 + DEFAULT_PROCESS_TIMEOUT = 15 * time.Second ) func NewApplication(name string, connMax int) *Application { @@ -118,7 +55,9 @@ func NewApplication(name string, connMax int) *Application { idleTimeout: DEFAULT_IDLE_TIMEOUT, readTimeout: DEFAULT_READ_TIMEOUT, readHeaderTimeout: DEFAULT_READ_HEADER_TIMEOUT, - writeTimeout: DEFAULT_WRITE_TIMEOUT} + writeTimeout: DEFAULT_WRITE_TIMEOUT, + processTimeout: DEFAULT_PROCESS_TIMEOUT, + } app.SetContext(appctx.APP_CONTEXT_KEY_APP, &app) app.SetContext(appctx.APP_CONTEXT_KEY_APPNAME, app.name) @@ -174,9 +113,6 @@ func (app *Application) AddHandler(method string, prefix string, handler func(co func (app *Application) AddHandler2(method string, prefix string, handler func(context.Context, http.ResponseWriter, *http.Request), metadata map[string]interface{}, name string, tags map[string]string) { log.Debugf("%s - %s", method, prefix) segs := SplitPath(prefix) - // for i := len(this.middlewares) - 1; i >= 0; i -= 1 { - // handler = this.middlewares[i](handler) - // } e := app.getRoot(method).Add(segs, newHandlerInfo(method, segs, handler, metadata, name, tags)) if e != nil { log.Fatalf("Fail to register %s %s: %s", method, prefix, e) @@ -254,9 +190,9 @@ func (app *Application) defaultHandle(w http.ResponseWriter, r *http.Request, ri if ok { fw := newResponseWriterChannel(w) errChan := make(chan interface{}) + ctx, cancel := context.WithTimeout(app.context, app.processTimeout) + defer cancel() app.session.Run(func() { - ctx, cancel := context.WithCancel(app.context) - defer cancel() defer fw.closeChannels() if ctx.Err() == nil { ctx = context.WithValue(ctx, appctx.APP_CONTEXT_KEY_REQUEST_ID, rid) @@ -266,25 +202,32 @@ func (app *Application) defaultHandle(w http.ResponseWriter, r *http.Request, ri if hand.metadata != nil { ctx = context.WithValue(ctx, appctx.APP_CONTEXT_KEY_METADATA, hand.metadata) } - span := trace.StartServerTrace(w, r, hand.GetName(params), app.GetName(), hand.GetTags()) - ctx = context.WithValue(ctx, appctx.APP_CONTEXT_KEY_TRACE, span) - hand.handler(ctx, &fw, r) - span.EndTrace() + func() { + span := trace.StartServerTrace(w, r, hand.GetName(params), app.GetName(), hand.GetTags()) + defer span.EndTrace() + ctx = context.WithValue(ctx, appctx.APP_CONTEXT_KEY_TRACE, span) + hand.handler(ctx, &fw, r) + }() } // otherwise, the task has been timeout }, errChan) - fw.wait() - runerr := WaitChannel(errChan) - if runerr != nil { - http.Error(w, fmt.Sprintf("Internal error: %s", runerr), http.StatusInternalServerError) + runErr := fw.wait(ctx, errChan) + if runErr != nil { + switch runErr.(type) { + case *httputils.JSONClientError: + je := runErr.(*httputils.JSONClientError) + httperrors.GeneralServerError(w, je) + default: + httperrors.InternalServerError(w, "Internal server error") + } } return hand } else { - log.Printf("Invalid handler for %s", r.URL) - http.Error(w, "Invalid handler", 500) + log.Errorf("Invalid handler for %s", r.URL) + httperrors.InternalServerError(w, "Invalid handler %s", r.URL) } } else if !isCors { - log.Printf("Handler not found") - http.NotFound(w, r) + log.Errorf("Handler not found") + httperrors.NotFoundError(w, "Handler not found") } return nil } diff --git a/pkg/appsrv/response.go b/pkg/appsrv/response.go new file mode 100644 index 0000000000..c2a4263700 --- /dev/null +++ b/pkg/appsrv/response.go @@ -0,0 +1,87 @@ +package appsrv + +import ( + "context" + "net/http" + + "yunion.io/x/onecloud/pkg/httperrors" +) + +type responseWriterResponse struct { + count int + err error +} + +type responseWriterChannel struct { + backend http.ResponseWriter + bodyChan chan []byte + bodyResp chan responseWriterResponse + statusChan chan int + statusResp chan bool +} + +func newResponseWriterChannel(backend http.ResponseWriter) responseWriterChannel { + return responseWriterChannel{backend: backend, + bodyChan: make(chan []byte), + bodyResp: make(chan responseWriterResponse), + statusChan: make(chan int), + statusResp: make(chan bool)} +} + +func (w *responseWriterChannel) Header() http.Header { + return w.backend.Header() +} + +func (w *responseWriterChannel) Write(bytes []byte) (int, error) { + w.bodyChan <- bytes + v := <-w.bodyResp + return v.count, v.err +} + +func (w *responseWriterChannel) WriteHeader(status int) { + w.statusChan <- status + <-w.statusResp +} + +func (w *responseWriterChannel) wait(ctx context.Context, errChan chan interface{}) interface{} { + var err interface{} + stop := false + for !stop { + select { + case <-ctx.Done(): + // ctx deadline reached, timeout + err = httperrors.NewTimeoutError("request process timeout") + stop = true + case e, more := <-errChan: + if more { + err = e + } else { + stop = true + } + case bytes, more := <-w.bodyChan: + // log.Print("Recive body ", len(bytes), " more ", more) + if more { + c, e := w.backend.Write(bytes) + w.bodyResp <- responseWriterResponse{count: c, err: e} + } else { + stop = true + } + case status, more := <-w.statusChan: + // log.Print("Recive status ", status, " more ", more) + if more { + w.backend.WriteHeader(status) + w.statusResp <- true + } else { + stop = true + } + } + } + return err +} + +func (w *responseWriterChannel) closeChannels() { + close(w.bodyChan) + close(w.bodyResp) + close(w.statusChan) + close(w.statusResp) +} diff --git a/pkg/httperrors/errors.go b/pkg/httperrors/errors.go index 4e4b226e03..90927a6a10 100644 --- a/pkg/httperrors/errors.go +++ b/pkg/httperrors/errors.go @@ -185,6 +185,11 @@ func NewRequireLicenseError(msg string, params ...interface{}) *httputils.JSONCl return NewJsonClientError(402, "RequireLicenseError", msg, err) } +func NewTimeoutError(msg string, params ...interface{}) *httputils.JSONClientError { + msg, err := errorMessage(msg, params...) + return NewJsonClientError(504, "TimeoutError", msg, err) +} + func NewGeneralError(err error) *httputils.JSONClientError { switch err.(type) { case *httputils.JSONClientError: diff --git a/pkg/httperrors/httperrors.go b/pkg/httperrors/httperrors.go index b9be4abbaf..e23e18c018 100644 --- a/pkg/httperrors/httperrors.go +++ b/pkg/httperrors/httperrors.go @@ -94,3 +94,7 @@ func TenantNotFoundError(w http.ResponseWriter, msg string, params ...interface{ func OutOfQuotaError(w http.ResponseWriter, msg string, params ...interface{}) { JsonClientError(w, NewOutOfQuotaError(msg, params...)) } + +func TimeoutError(w http.ResponseWriter, msg string, params ...interface{}) { + JsonClientError(w, NewTimeoutError(msg, params...)) +} From 4d1d4e56fd55f8e10a14206feff6ef93e7ca8d15 Mon Sep 17 00:00:00 2001 From: Qiu Jian Date: Mon, 22 Oct 2018 01:38:52 +0800 Subject: [PATCH 10/23] =?UTF-8?q?=E4=BF=AE=E6=AD=A3=EF=BC=9Aworker=20timeo?= =?UTF-8?q?ut=E4=B9=8B=E5=90=8E=E9=9C=80=E8=A6=81=E4=BB=8Eworker=20manager?= =?UTF-8?q?=E4=B8=AD=E7=A7=BB=E9=99=A4=EF=BC=8C=E4=B8=8D=E5=86=8D=E5=8D=A0?= =?UTF-8?q?=E7=94=A8worker=E7=9A=84=E6=95=B0=E9=87=8F=E9=99=90=E9=A2=9D?= =?UTF-8?q?=EF=BC=8C=E5=90=A6=E5=88=99timeout=E4=B9=8B=E5=90=8E=EF=BC=8C?= =?UTF-8?q?=E5=B9=B6=E4=B8=8D=E8=83=BD=E8=A7=A3=E5=86=B3=E6=95=B4=E4=B8=AA?= =?UTF-8?q?=E6=9C=8D=E5=8A=A1=E8=A2=AB=E5=8D=A1=E4=BD=8F=E7=9A=84=E9=97=AE?= =?UTF-8?q?=E9=A2=98?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- pkg/appsrv/appsrv.go | 9 +- pkg/appsrv/response.go | 10 +- pkg/appsrv/workers.go | 194 ++++++++++++++---- pkg/appsrv/workers_test.go | 19 +- pkg/cloudcommon/db/taskman/handler.go | 4 +- pkg/cloudcommon/db/taskman/localtaskworker.go | 4 +- pkg/util/logclient/logclient.go | 4 +- 7 files changed, 185 insertions(+), 59 deletions(-) diff --git a/pkg/appsrv/appsrv.go b/pkg/appsrv/appsrv.go index f434bd44b4..c31936ceb7 100644 --- a/pkg/appsrv/appsrv.go +++ b/pkg/appsrv/appsrv.go @@ -22,7 +22,7 @@ import ( type Application struct { name string context context.Context - session *WorkerManager + session *SWorkerManager roots map[string]*RadixNode rootLock *sync.Mutex connMax int @@ -42,7 +42,7 @@ const ( DEFAULT_READ_TIMEOUT = 0 DEFAULT_READ_HEADER_TIMEOUT = 10 * time.Second DEFAULT_WRITE_TIMEOUT = 0 - DEFAULT_PROCESS_TIMEOUT = 15 * time.Second + DEFAULT_PROCESS_TIMEOUT = 15 * time.Millisecond ) func NewApplication(name string, connMax int) *Application { @@ -189,6 +189,7 @@ func (app *Application) defaultHandle(w http.ResponseWriter, r *http.Request, ri hand, ok := handler.(*handlerInfo) if ok { fw := newResponseWriterChannel(w) + worker := make(chan *SWorker) errChan := make(chan interface{}) ctx, cancel := context.WithTimeout(app.context, app.processTimeout) defer cancel() @@ -209,8 +210,8 @@ func (app *Application) defaultHandle(w http.ResponseWriter, r *http.Request, ri hand.handler(ctx, &fw, r) }() } // otherwise, the task has been timeout - }, errChan) - runErr := fw.wait(ctx, errChan) + }, worker, errChan) + runErr := fw.wait(ctx, worker, errChan) if runErr != nil { switch runErr.(type) { case *httputils.JSONClientError: diff --git a/pkg/appsrv/response.go b/pkg/appsrv/response.go index c2a4263700..6bf18eea4c 100644 --- a/pkg/appsrv/response.go +++ b/pkg/appsrv/response.go @@ -4,6 +4,8 @@ import ( "context" "net/http" + "yunion.io/x/log" + "yunion.io/x/onecloud/pkg/httperrors" ) @@ -43,13 +45,19 @@ func (w *responseWriterChannel) WriteHeader(status int) { <-w.statusResp } -func (w *responseWriterChannel) wait(ctx context.Context, errChan chan interface{}) interface{} { +func (w *responseWriterChannel) wait(ctx context.Context, workerChan chan *SWorker, errChan chan interface{}) interface{} { var err interface{} + var worker *SWorker stop := false for !stop { select { + case worker = <-workerChan: + log.Infof("request is being handled by worker %s", worker) case <-ctx.Done(): // ctx deadline reached, timeout + if worker != nil { + worker.Detach("timeout") + } err = httperrors.NewTimeoutError("request process timeout") stop = true case e, more := <-errChan: diff --git a/pkg/appsrv/workers.go b/pkg/appsrv/workers.go index 675853be3d..bffed2f678 100644 --- a/pkg/appsrv/workers.go +++ b/pkg/appsrv/workers.go @@ -4,61 +4,164 @@ import ( "runtime/debug" "sync" + "container/list" + "fmt" "yunion.io/x/log" ) -type WorkerManager struct { - name string - queue *Ring - workerCount int - backlog int - activeWorker int - workerLock *sync.Mutex - workerId uint64 +const ( + WORKER_STATE_ACTIVE = 0 + WORKER_STATE_DETACH = 1 +) + +var isDebug = false + +func enableDebug() { + isDebug = true } -func NewWorkerManager(name string, workerCount int, backlog int) *WorkerManager { - manager := WorkerManager{name: name, - queue: NewRing(workerCount * backlog), - workerCount: workerCount, - backlog: backlog, - activeWorker: 0, - workerLock: &sync.Mutex{}, - workerId: 0} +type SWorker struct { + id uint64 + state int + container *list.Element + manager *SWorkerManager +} + +func newWorker(id uint64, manager *SWorkerManager) *SWorker { + return &SWorker{ + id: id, + state: WORKER_STATE_ACTIVE, + container: nil, + manager: manager, + } +} + +func (worker *SWorker) isDetached() bool { + worker.manager.workerLock.Lock() + defer worker.manager.workerLock.Unlock() + + return worker.state == WORKER_STATE_DETACH +} + +func (worker *SWorker) run() { + for { + if worker.isDetached() { + if isDebug { + log.Debugf("deteched worker %s, no need to pick up new job", worker) + } + break + } + req := worker.manager.queue.Pop() + if req != nil { + task := req.(*sWorkerTask) + if task.worker != nil { + task.worker <- worker + } + if isDebug { + log.Debugf("start exec task on worker %s", worker) + } + execCallback(task) + if isDebug { + log.Debugf("end exec task on worker %s", worker) + } + } else { + if isDebug { + log.Debugf("no more job, exit worker %s", worker) + } + break + } + } + worker.manager.removeWorker(worker) +} + +func (worker *SWorker) Detach(reason string) { + worker.manager.workerLock.Lock() + defer worker.manager.workerLock.Unlock() + + worker.state = WORKER_STATE_DETACH + worker.manager.activeWorker.removeWithLock(worker) + worker.manager.detachedWorker.addWithLock(worker) + + log.Warningf("detach worker %s due to reason %s", worker, reason) +} + +func (worker *SWorker) String() string { + return fmt.Sprintf("#%d(%d)", worker.id, worker.state) +} + +type SWorkerList struct { + list *list.List +} + +func newWorkerList() SWorkerList { + return SWorkerList{ + list: list.New(), + } +} + +func (wl *SWorkerList) addWithLock(worker *SWorker) { + ele := wl.list.PushBack(worker) + worker.container = ele +} + +func (wl *SWorkerList) removeWithLock(worker *SWorker) { + wl.list.Remove(worker.container) + worker.container = nil +} + +func (wl *SWorkerList) size() int { + return wl.list.Len() +} + +type SWorkerManager struct { + name string + queue *Ring + workerCount int + backlog int + activeWorker SWorkerList + detachedWorker SWorkerList + workerLock *sync.Mutex + workerId uint64 +} + +func NewWorkerManager(name string, workerCount int, backlog int) *SWorkerManager { + manager := SWorkerManager{name: name, + queue: NewRing(workerCount * backlog), + workerCount: workerCount, + backlog: backlog, + activeWorker: newWorkerList(), + detachedWorker: newWorkerList(), + workerLock: &sync.Mutex{}, + workerId: 0} return &manager } -type workerTask struct { - task func() - err chan interface{} +type sWorkerTask struct { + task func() + worker chan *SWorker + err chan interface{} } -func (wm *WorkerManager) Run(task func(), err chan interface{}) bool { - ret := wm.queue.Push(&workerTask{task: task, err: err}) +func (wm *SWorkerManager) Run(task func(), worker chan *SWorker, err chan interface{}) bool { + ret := wm.queue.Push(&sWorkerTask{task: task, worker: worker, err: err}) if ret { wm.schedule() } return ret } -func (wm *WorkerManager) workerRun(id uint64) { - //log.Println("Start worker", id) - defer wm.decActiveWorker() - req := wm.queue.Pop() - for req != nil { - wm.execCallback(req.(*workerTask)) - req = wm.queue.Pop() - } - //log.Println("End worker", id) -} - -func (wm *WorkerManager) decActiveWorker() { +func (wm *SWorkerManager) removeWorker(worker *SWorker) { wm.workerLock.Lock() defer wm.workerLock.Unlock() - wm.activeWorker -= 1 + + if worker.state == WORKER_STATE_ACTIVE { + wm.activeWorker.removeWithLock(worker) + } else { + wm.detachedWorker.removeWithLock(worker) + } } -func (wm *WorkerManager) execCallback(task *workerTask) { +func execCallback(task *sWorkerTask) { defer func() { if r := recover(); r != nil { log.Errorf("WorkerManager exec callback error: %s", r) @@ -75,16 +178,29 @@ func (wm *WorkerManager) execCallback(task *workerTask) { } } -func (wm *WorkerManager) schedule() { +func (wm *SWorkerManager) schedule() { wm.workerLock.Lock() defer wm.workerLock.Unlock() - if wm.activeWorker < wm.workerCount && wm.queue.Size() > 0 { - wm.activeWorker += 1 + + if wm.activeWorker.size() < wm.workerCount && wm.queue.Size() > 0 { wm.workerId += 1 - go wm.workerRun(wm.workerId) + worker := newWorker(wm.workerId, wm) + wm.activeWorker.addWithLock(worker) + if isDebug { + log.Debugf("no enough worker, add new worker %s", worker) + } + go worker.run() } } +func (wm *SWorkerManager) ActiveWorkerCount() int { + return wm.activeWorker.size() +} + +func (wm *SWorkerManager) DetachedWorkerCount() int { + return wm.detachedWorker.size() +} + func WaitChannel(ch chan interface{}) interface{} { var ret interface{} stop := false diff --git a/pkg/appsrv/workers_test.go b/pkg/appsrv/workers_test.go index 69f4877b9f..bc1527ab37 100644 --- a/pkg/appsrv/workers_test.go +++ b/pkg/appsrv/workers_test.go @@ -6,22 +6,22 @@ import ( ) func TestWorkerManager(t *testing.T) { + enableDebug() startTime := time.Now() - end := make(chan int) + // end := make(chan int) wm := NewWorkerManager("testwm", 2, 10) counter := 0 for i := 0; i < 10; i += 1 { wm.Run(func() { counter += 1 time.Sleep(1 * time.Second) - if counter >= i { - end <- 1 - } - }, nil) + }, nil, nil) + } + for wm.ActiveWorkerCount() != 0 { + time.Sleep(time.Second) } - <-end if time.Since(startTime) < 5*time.Second { - t.Error("Increct timing") + t.Error("Incorrect timing") } } @@ -30,7 +30,7 @@ func TestWorkerManagerError(t *testing.T) { err := make(chan interface{}) wm.Run(func() { panic("Panic inside worker") - }, err) + }, nil, err) e := WaitChannel(err) if e == nil { t.Error("Panic not captured") @@ -38,9 +38,10 @@ func TestWorkerManagerError(t *testing.T) { err = make(chan interface{}) wm.Run(func() { time.Sleep(1 * time.Second) - }, err) + }, nil, err) e = WaitChannel(err) if e != nil { t.Error("Should no error") } + } diff --git a/pkg/cloudcommon/db/taskman/handler.go b/pkg/cloudcommon/db/taskman/handler.go index beb9de938b..23dc6dff15 100644 --- a/pkg/cloudcommon/db/taskman/handler.go +++ b/pkg/cloudcommon/db/taskman/handler.go @@ -8,7 +8,7 @@ import ( "yunion.io/x/onecloud/pkg/cloudcommon/db" ) -var taskWorkMan *appsrv.WorkerManager +var taskWorkMan *appsrv.SWorkerManager func init() { taskWorkMan = appsrv.NewWorkerManager("TaskWorkerManager", 4, 100) @@ -22,5 +22,5 @@ func AddTaskHandler(prefix string, app *appsrv.Application) { func runTask(taskId string, data jsonutils.JSONObject) { taskWorkMan.Run(func() { TaskManager.execTask(taskId, data) - }, nil) + }, nil, nil) } diff --git a/pkg/cloudcommon/db/taskman/localtaskworker.go b/pkg/cloudcommon/db/taskman/localtaskworker.go index f60e8b146f..57eaf00fb5 100644 --- a/pkg/cloudcommon/db/taskman/localtaskworker.go +++ b/pkg/cloudcommon/db/taskman/localtaskworker.go @@ -9,7 +9,7 @@ import ( "yunion.io/x/onecloud/pkg/appsrv" ) -var localTaskWorkerMan *appsrv.WorkerManager +var localTaskWorkerMan *appsrv.SWorkerManager func init() { localTaskWorkerMan = appsrv.NewWorkerManager("LocalTaskWorkerManager", 4, 10) @@ -42,5 +42,5 @@ func LocalTaskRun(task ITask, proc func() (jsonutils.JSONObject, error)) { task.ScheduleRun(data) } - }, nil) + }, nil, nil) } diff --git a/pkg/util/logclient/logclient.go b/pkg/util/logclient/logclient.go index 0d8a541788..176d55f9b9 100644 --- a/pkg/util/logclient/logclient.go +++ b/pkg/util/logclient/logclient.go @@ -60,7 +60,7 @@ const ( // golang 不支持 const 的string array, http://t.cn/EzAvbw8 var BLACK_LIST_OBJ_TYPE = []string{"parameter"} -var logclientWorkerMan *appsrv.WorkerManager +var logclientWorkerMan *appsrv.SWorkerManager func init() { logclientWorkerMan = appsrv.NewWorkerManager("LogClientWorkerManager", 1, 50) @@ -119,5 +119,5 @@ func AddActionLog(model IObject, action string, iNotes interface{}, userCred mcc if err != nil { log.Errorf("create action log failed %s", err) } - }, nil) + }, nil, nil) } From 2d7e575e6f54c1ff50e657ac30c91a6c8db41ca3 Mon Sep 17 00:00:00 2001 From: Qiu Jian Date: Mon, 22 Oct 2018 01:43:31 +0800 Subject: [PATCH 11/23] minor fixes --- pkg/appsrv/appsrv.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pkg/appsrv/appsrv.go b/pkg/appsrv/appsrv.go index c31936ceb7..c57af5a03b 100644 --- a/pkg/appsrv/appsrv.go +++ b/pkg/appsrv/appsrv.go @@ -42,7 +42,7 @@ const ( DEFAULT_READ_TIMEOUT = 0 DEFAULT_READ_HEADER_TIMEOUT = 10 * time.Second DEFAULT_WRITE_TIMEOUT = 0 - DEFAULT_PROCESS_TIMEOUT = 15 * time.Millisecond + DEFAULT_PROCESS_TIMEOUT = 15 * time.Second ) func NewApplication(name string, connMax int) *Application { From 9cbd8891a9aa491b4a647a1d92b2121b28801268 Mon Sep 17 00:00:00 2001 From: Qiu Jian Date: Mon, 22 Oct 2018 09:53:50 +0800 Subject: [PATCH 12/23] =?UTF-8?q?=E5=A2=9E=E5=8A=A0=EF=BC=9AworkerManager?= =?UTF-8?q?=E7=BB=9F=E8=AE=A1=E6=95=B0=E6=8D=AE=E8=8E=B7=E5=8F=96API?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- pkg/appsrv/appsrv.go | 3 ++- pkg/appsrv/workers.go | 43 +++++++++++++++++++++++++++++++++++++++++-- 2 files changed, 43 insertions(+), 3 deletions(-) diff --git a/pkg/appsrv/appsrv.go b/pkg/appsrv/appsrv.go index c57af5a03b..38c17a7604 100644 --- a/pkg/appsrv/appsrv.go +++ b/pkg/appsrv/appsrv.go @@ -49,7 +49,7 @@ func NewApplication(name string, connMax int) *Application { app := Application{name: name, context: context.Background(), connMax: connMax, - session: NewWorkerManager("sessionMan", connMax, DEFAULT_BACKLOG), + session: NewWorkerManager("HttpRequestWorkerManager", connMax, DEFAULT_BACKLOG), roots: make(map[string]*RadixNode), rootLock: &sync.Mutex{}, idleTimeout: DEFAULT_IDLE_TIMEOUT, @@ -239,6 +239,7 @@ func (app *Application) addDefaultHandler() { app.AddHandler("POST", "/ping", PingHandler) app.AddHandler("GET", "/ping", PingHandler) // app.AddHandler("OPTIONS", "/", CORSHandler) + app.AddHandler("GET", "/worker_stats", WorkerStatsHandler) } func timeoutHandle(h http.Handler) http.HandlerFunc { diff --git a/pkg/appsrv/workers.go b/pkg/appsrv/workers.go index bffed2f678..ed3b940eee 100644 --- a/pkg/appsrv/workers.go +++ b/pkg/appsrv/workers.go @@ -1,11 +1,14 @@ package appsrv import ( + "container/list" + "context" + "fmt" + "net/http" "runtime/debug" "sync" - "container/list" - "fmt" + "yunion.io/x/jsonutils" "yunion.io/x/log" ) @@ -20,6 +23,12 @@ func enableDebug() { isDebug = true } +var workerManagers []*SWorkerManager + +func init() { + workerManagers = make([]*SWorkerManager, 0) +} + type SWorker struct { id uint64 state int @@ -133,6 +142,8 @@ func NewWorkerManager(name string, workerCount int, backlog int) *SWorkerManager detachedWorker: newWorkerList(), workerLock: &sync.Mutex{}, workerId: 0} + + workerManagers = append(workerManagers, &manager) return &manager } @@ -201,6 +212,34 @@ func (wm *SWorkerManager) DetachedWorkerCount() int { return wm.detachedWorker.size() } +type SWorkerManagerStates struct { + Name string + Backlog int + ActiveWorkerCnt int + DetachWorkerCnt int +} + +func (wm *SWorkerManager) getState() SWorkerManagerStates { + state := SWorkerManagerStates{} + + state.Name = wm.name + state.Backlog = wm.queue.Size() + state.ActiveWorkerCnt = wm.activeWorker.size() + state.DetachWorkerCnt = wm.detachedWorker.size() + + return state +} + +func WorkerStatsHandler(ctx context.Context, w http.ResponseWriter, r *http.Request) { + stats := make([]SWorkerManagerStates, 0) + for i := 0; i < len(workerManagers); i += 1 { + stats = append(stats, workerManagers[i].getState()) + } + result := jsonutils.NewDict() + result.Add(jsonutils.Marshal(&stats), "workers") + fmt.Fprintf(w, result.String()) +} + func WaitChannel(ch chan interface{}) interface{} { var ret interface{} stop := false From fc61a8f9c5d69995df4fd6a31bada8a552c9b789 Mon Sep 17 00:00:00 2001 From: Qiu Jian Date: Mon, 22 Oct 2018 10:02:40 +0800 Subject: [PATCH 13/23] add max_worker state --- pkg/appsrv/workers.go | 2 ++ 1 file changed, 2 insertions(+) diff --git a/pkg/appsrv/workers.go b/pkg/appsrv/workers.go index ed3b940eee..45c48b9a6e 100644 --- a/pkg/appsrv/workers.go +++ b/pkg/appsrv/workers.go @@ -215,6 +215,7 @@ func (wm *SWorkerManager) DetachedWorkerCount() int { type SWorkerManagerStates struct { Name string Backlog int + MaxWorkerCnt int ActiveWorkerCnt int DetachWorkerCnt int } @@ -224,6 +225,7 @@ func (wm *SWorkerManager) getState() SWorkerManagerStates { state.Name = wm.name state.Backlog = wm.queue.Size() + state.MaxWorkerCnt = wm.workerCount state.ActiveWorkerCnt = wm.activeWorker.size() state.DetachWorkerCnt = wm.detachedWorker.size() From b27be13cc391f4869ffcbeaf41f46ea642295882 Mon Sep 17 00:00:00 2001 From: Qiu Jian Date: Mon, 22 Oct 2018 11:24:54 +0800 Subject: [PATCH 14/23] =?UTF-8?q?=E4=BF=AE=E6=AD=A3=EF=BC=9A=E8=A1=A5?= =?UTF-8?q?=E5=85=852.2.0=20purge=20eip/snapshot=E6=8E=A5=E5=8F=A3?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- cmd/climc/shell/elasticips.go | 12 ++++++++++++ cmd/climc/shell/snapshots.go | 13 +++++++++++++ pkg/compute/models/cloudproviders.go | 9 +++++++++ pkg/compute/models/elasticips.go | 19 +++++++++++++++++++ pkg/compute/models/snapshots.go | 19 +++++++++++++++++++ 5 files changed, 72 insertions(+) diff --git a/cmd/climc/shell/elasticips.go b/cmd/climc/shell/elasticips.go index 054283940c..74e6e0f281 100644 --- a/cmd/climc/shell/elasticips.go +++ b/cmd/climc/shell/elasticips.go @@ -187,4 +187,16 @@ func init() { return nil }) + type EipPurgeOptions struct { + ID string `help:"ID or name of EIP"` + } + R(&EipPurgeOptions{}, "eip-purge", "Purge EIP db records", func(s *mcclient.ClientSession, args *EipPurgeOptions) error { + result, err := modules.Elasticips.PerformAction(s, args.ID, "purge", nil) + if err != nil { + return err + } + printObject(result) + return nil + }) + } diff --git a/cmd/climc/shell/snapshots.go b/cmd/climc/shell/snapshots.go index 3be5a0ad2b..d809a1226e 100644 --- a/cmd/climc/shell/snapshots.go +++ b/cmd/climc/shell/snapshots.go @@ -84,4 +84,17 @@ func init() { printObject(result) return nil }) + + type SnapshotPurgeOptions struct { + ID string `help:"ID or name of Snapshot"` + } + R(&SnapshotPurgeOptions{}, "snapshot-purge", "Purge Snapshot db records", func(s *mcclient.ClientSession, args *SnapshotPurgeOptions) error { + result, err := modules.Snapshots.PerformAction(s, args.ID, "purge", nil) + if err != nil { + return err + } + printObject(result) + return nil + }) + } diff --git a/pkg/compute/models/cloudproviders.go b/pkg/compute/models/cloudproviders.go index 31afc0984f..d3e9b0e5bf 100644 --- a/pkg/compute/models/cloudproviders.go +++ b/pkg/compute/models/cloudproviders.go @@ -96,6 +96,10 @@ func (self *SCloudprovider) getEipCount() int { return ElasticipManager.Query().Equals("manager_id", self.Id).Count() } +func (self *SCloudprovider) getSnapshotCount() int { + return SnapshotManager.Query().Equals("manager_id", self.Id).Count() +} + func (self *SCloudprovider) ValidateUpdateData(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data *jsonutils.JSONDict) (*jsonutils.JSONDict, error) { return self.SEnabledStatusStandaloneResourceBase.ValidateUpdateData(ctx, userCred, query, data) } @@ -425,6 +429,7 @@ type SCloudproviderUsage struct { StorageCount int StorageCacheCount int EipCount int + SnapshotCount int } func (usage *SCloudproviderUsage) isEmpty() bool { @@ -443,6 +448,9 @@ func (usage *SCloudproviderUsage) isEmpty() bool { if usage.EipCount > 0 { return false } + if usage.SnapshotCount > 0 { + return false + } return true } @@ -454,6 +462,7 @@ func (self *SCloudprovider) getUsage() *SCloudproviderUsage { usage.StorageCount = self.getStorageCount() usage.StorageCacheCount = self.getStoragecacheCount() usage.EipCount = self.getEipCount() + usage.SnapshotCount = self.getSnapshotCount() return &usage } diff --git a/pkg/compute/models/elasticips.go b/pkg/compute/models/elasticips.go index 5c5ad97792..a164e35ada 100644 --- a/pkg/compute/models/elasticips.go +++ b/pkg/compute/models/elasticips.go @@ -799,3 +799,22 @@ func (manager *SElasticipManager) TotalCount(projectId string, rangeObj db.IStan usage.EIPUsedCount = q3.Count() return usage } + +func (self *SElasticip) AllowPerformPurge(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) bool { + return userCred.IsSystemAdmin() +} + +func (self *SElasticip) PerformPurge(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) (jsonutils.JSONObject, error) { + err := self.ValidateDeleteCondition(ctx) + if err != nil { + return nil, err + } + provider := self.GetCloudprovider() + if provider != nil { + if provider.Enabled { + return nil, httperrors.NewInvalidStatusError("Cannot purge elastic_ip on enabled cloud provider") + } + } + err = self.RealDelete(ctx, userCred) + return nil, err +} diff --git a/pkg/compute/models/snapshots.go b/pkg/compute/models/snapshots.go index 3367203ee4..5d5073710c 100644 --- a/pkg/compute/models/snapshots.go +++ b/pkg/compute/models/snapshots.go @@ -530,3 +530,22 @@ func (self *SSnapshot) GetISnapshotRegion() (cloudprovider.ICloudRegion, error) } return provider.GetIRegionById(region.GetExternalId()) } + +func (self *SSnapshot) AllowPerformPurge(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) bool { + return userCred.IsSystemAdmin() +} + +func (self *SSnapshot) PerformPurge(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) (jsonutils.JSONObject, error) { + err := self.ValidateDeleteCondition(ctx) + if err != nil { + return nil, err + } + provider := self.GetCloudprovider() + if provider != nil { + if provider.Enabled { + return nil, httperrors.NewInvalidStatusError("Cannot purge snapshot on enabled cloud provider") + } + } + err = self.RealDelete(ctx, userCred) + return nil, err +} From 46a39b56124ee5a6b3f318f86aaf13d4820ec985 Mon Sep 17 00:00:00 2001 From: Qiu Jian Date: Mon, 22 Oct 2018 12:04:55 +0800 Subject: [PATCH 15/23] allow snapshot-list --manager --- cmd/climc/shell/snapshots.go | 5 +++++ cmd/climc/shell/vpcs.go | 2 +- pkg/compute/models/snapshots.go | 12 ++++++++++++ pkg/compute/models/storages.go | 6 +++++- 4 files changed, 23 insertions(+), 2 deletions(-) diff --git a/cmd/climc/shell/snapshots.go b/cmd/climc/shell/snapshots.go index d809a1226e..967abb7387 100644 --- a/cmd/climc/shell/snapshots.go +++ b/cmd/climc/shell/snapshots.go @@ -16,6 +16,8 @@ func init() { Share bool `help:"Show shared snapshots"` DiskType string `help: "Filter by disk type" choices:"sys|data"` Provider string `help: "Cloud provider" choices:"Aliyun|VMware|Azure"` + + Manager string `help:"Show snapshots belongs to a specific cloud provider"` } R(&SnapshotsListOptions{}, "snapshot-list", "Show snapshots", func(s *mcclient.ClientSession, args *SnapshotsListOptions) error { params, err := args.BaseListOptions.Params() @@ -41,6 +43,9 @@ func init() { if len(args.Provider) > 0 { params.Add(jsonutils.NewString(args.Provider), "provider") } + if len(args.Manager) > 0 { + params.Add(jsonutils.NewString(args.Manager), "manager") + } result, err := modules.Snapshots.List(s, params) if err != nil { return err diff --git a/cmd/climc/shell/vpcs.go b/cmd/climc/shell/vpcs.go index 547a775dd0..79517f2880 100644 --- a/cmd/climc/shell/vpcs.go +++ b/cmd/climc/shell/vpcs.go @@ -11,7 +11,7 @@ func init() { type VpcListOptions struct { options.BaseListOptions Region string `help:"ID or Name of region"` - Manager string `help:"Show regions belongs to the cloud provider"` + Manager string `help:"Show vpcs belongs to the cloud provider"` } R(&VpcListOptions{}, "vpc-list", "List VPCs", func(s *mcclient.ClientSession, args *VpcListOptions) error { var params *jsonutils.JSONDict diff --git a/pkg/compute/models/snapshots.go b/pkg/compute/models/snapshots.go index 5d5073710c..34fe6a4945 100644 --- a/pkg/compute/models/snapshots.go +++ b/pkg/compute/models/snapshots.go @@ -112,6 +112,18 @@ func (manager *SSnapshotManager) ListItemFilter(ctx context.Context, q *sqlchemy sq := cloudproviderTbl.Query(cloudproviderTbl.Field("id")).Equals("provider", provider) q = q.In("manager_id", sq) } + + if managerStr := jsonutils.GetAnyString(query, []string{"manager", "manager_id"}); len(managerStr) > 0 { + managerObj, err := CloudproviderManager.FetchByIdOrName("", managerStr) + if err != nil { + if err == sql.ErrNoRows { + return nil, httperrors.NewNotFoundError("manager %s not found", managerStr) + } + return nil, httperrors.NewGeneralError(err) + } + q = q.Equals("manager_id", managerObj.GetId()) + } + return q, nil } diff --git a/pkg/compute/models/storages.go b/pkg/compute/models/storages.go index 060c266c62..256321c33c 100644 --- a/pkg/compute/models/storages.go +++ b/pkg/compute/models/storages.go @@ -83,7 +83,7 @@ func (manager *SStorageManager) GetContextManager() []db.IModelManager { } func (self *SStorage) ValidateDeleteCondition(ctx context.Context) error { - if self.GetHostCount() > 0 || self.GetDiskCount() > 0 { + if self.GetHostCount() > 0 || self.GetDiskCount() > 0 || self.GetSnapshotCount() > 0 { return httperrors.NewNotEmptyError("Not an empty storage provider") } return self.SEnabledStatusStandaloneResourceBase.ValidateDeleteCondition(ctx) @@ -97,6 +97,10 @@ func (self *SStorage) GetDiskCount() int { return DiskManager.Query().Equals("storage_id", self.Id).Count() } +func (self *SStorage) GetSnapshotCount() int { + return SnapshotManager.Query().Equals("storage_id", self.Id).Count() +} + func (self *SStorage) IsLocal() bool { return self.StorageType == STORAGE_LOCAL || self.StorageType == STORAGE_BAREMETAL } From 3d7723499d79cba98f9ec541d91b5835cff63387 Mon Sep 17 00:00:00 2001 From: Zhang Dongliang Date: Mon, 22 Oct 2018 15:20:56 +0800 Subject: [PATCH 16/23] =?UTF-8?q?=E4=BA=91=E6=9C=8D=E5=8A=A1=E5=99=A8?= =?UTF-8?q?=E5=85=B3=E8=81=94=E5=AE=89=E5=85=A8=E7=BB=84=E5=90=8E=EF=BC=8C?= =?UTF-8?q?=E6=93=8D=E4=BD=9C=E6=97=A5=E5=BF=97=E8=AE=B0=E5=BD=95=E5=86=85?= =?UTF-8?q?=E5=AE=B9=E5=BA=94=E8=AF=A5=E6=98=AF=E2=80=9C=E5=85=B3=E8=81=94?= =?UTF-8?q?=E5=AE=89=E5=85=A8=E7=BB=84=E2=80=9D?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- pkg/compute/models/guests.go | 7 +++++++ pkg/util/logclient/logclient.go | 1 + 2 files changed, 8 insertions(+) diff --git a/pkg/compute/models/guests.go b/pkg/compute/models/guests.go index 922c21b4b3..5dc37f6eec 100644 --- a/pkg/compute/models/guests.go +++ b/pkg/compute/models/guests.go @@ -2451,24 +2451,31 @@ func (self *SGuest) PerformRevokeSecgroup(ctx context.Context, userCred mcclient func (self *SGuest) PerformAssignSecgroup(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) (jsonutils.JSONObject, error) { if !utils.IsInStringArray(self.Status, []string{VM_READY, VM_RUNNING, VM_SUSPEND}) { + logclient.AddActionLog(self, logclient.ACT_VM_ASSIGNSECGROUP, "Cannot assign security rules in status "+self.Status, userCred, false) return nil, httperrors.NewInputParameterError("Cannot assign security rules in status %s", self.Status) } else { if secgrp, err := data.GetString("secgrp"); err != nil { + logclient.AddActionLog(self, logclient.ACT_VM_ASSIGNSECGROUP, err, userCred, false) return nil, err } else if sg, err := SecurityGroupManager.FetchByIdOrName(userCred.GetProjectId(), secgrp); err != nil { + msg := fmt.Sprintf("SecurityGroup %s not found", secgrp) + logclient.AddActionLog(self, logclient.ACT_VM_ASSIGNSECGROUP, msg, userCred, false) return nil, httperrors.NewNotFoundError("SecurityGroup %s not found", secgrp) } else { if _, err := self.GetModelManager().TableSpec().Update(self, func() error { self.SecgrpId = sg.GetId() return nil }); err != nil { + logclient.AddActionLog(self, logclient.ACT_VM_ASSIGNSECGROUP, err, userCred, false) return nil, err } if err := self.StartSyncTask(ctx, userCred, true, ""); err != nil { + logclient.AddActionLog(self, logclient.ACT_VM_ASSIGNSECGROUP, err, userCred, false) return nil, err } } } + logclient.AddActionLog(self, logclient.ACT_VM_ASSIGNSECGROUP, nil, userCred, true) return nil, nil } diff --git a/pkg/util/logclient/logclient.go b/pkg/util/logclient/logclient.go index 176d55f9b9..76796053e3 100644 --- a/pkg/util/logclient/logclient.go +++ b/pkg/util/logclient/logclient.go @@ -55,6 +55,7 @@ const ( ACT_VM_SYNC_CONF = "同步配置" ACT_VM_SYNC_STATUS = "同步状态" ACT_VM_UNBIND_KEYPAIR = "解绑密钥" + ACT_VM_ASSIGNSECGROUP = "关联安全组" ) // golang 不支持 const 的string array, http://t.cn/EzAvbw8 From 2adec7a7926fea8ebf0a4531cb73608f989ed88b Mon Sep 17 00:00:00 2001 From: Qiu Jian Date: Mon, 22 Oct 2018 15:21:42 +0800 Subject: [PATCH 17/23] =?UTF-8?q?=E4=BF=AE=E6=AD=A3=EF=BC=9A=E5=88=A0?= =?UTF-8?q?=E9=99=A4guest=E6=97=B6=E5=80=99=E6=B2=A1=E6=9C=89=E6=AD=A3?= =?UTF-8?q?=E7=A1=AE=E6=B8=85=E7=90=86eip=E7=9A=84=E6=95=B0=E6=8D=AE?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- pkg/compute/models/elasticips.go | 8 ++++++++ pkg/compute/models/guests.go | 20 ++++++++++++++++---- pkg/compute/models/storagecachedimages.go | 4 ++-- pkg/compute/models/storagecaches.go | 2 +- 4 files changed, 27 insertions(+), 7 deletions(-) diff --git a/pkg/compute/models/elasticips.go b/pkg/compute/models/elasticips.go index a164e35ada..69db4e1e1a 100644 --- a/pkg/compute/models/elasticips.go +++ b/pkg/compute/models/elasticips.go @@ -818,3 +818,11 @@ func (self *SElasticip) PerformPurge(ctx context.Context, userCred mcclient.Toke err = self.RealDelete(ctx, userCred) return nil, err } + +func (self *SElasticip) DoPendingDelete(ctx context.Context, userCred mcclient.TokenCredential) { + if self.Mode == EIP_MODE_INSTANCE_PUBLICIP { + self.SVirtualResourceBase.DoPendingDelete(ctx, userCred) + return + } + self.Dissociate(ctx, userCred) +} \ No newline at end of file diff --git a/pkg/compute/models/guests.go b/pkg/compute/models/guests.go index 922c21b4b3..85f3e3c909 100644 --- a/pkg/compute/models/guests.go +++ b/pkg/compute/models/guests.go @@ -3086,6 +3086,10 @@ func (self *SGuest) StartChangeConfigTask(ctx context.Context, userCred mcclient } func (self *SGuest) DoPendingDelete(ctx context.Context, userCred mcclient.TokenCredential) { + eip, _ := self.GetEip() + if eip != nil { + eip.DoPendingDelete(ctx, userCred) + } for _, guestdisk := range self.GetDisks() { disk := guestdisk.GetDisk() storage := disk.GetStorage() @@ -4479,10 +4483,18 @@ func (self *SGuest) DeleteEip(ctx context.Context, userCred mcclient.TokenCreden if eip == nil { return nil } - err = eip.Delete(ctx, userCred) - if err != nil { - log.Errorf("Delete eip fail %s", err) - return err + if eip.Mode == EIP_MODE_INSTANCE_PUBLICIP { + err = eip.RealDelete(ctx, userCred) + if err != nil { + log.Errorf("Delete eip on delete server fail %s", err) + return err + } + } else { + err = eip.Dissociate(ctx, userCred) + if err != nil { + log.Errorf("Dissociate eip on delete server fail %s", err) + return err + } } return nil } diff --git a/pkg/compute/models/storagecachedimages.go b/pkg/compute/models/storagecachedimages.go index 0bafb4624b..30d24a13cb 100644 --- a/pkg/compute/models/storagecachedimages.go +++ b/pkg/compute/models/storagecachedimages.go @@ -225,7 +225,7 @@ func (self *SStoragecachedimage) isDownloadSessionExpire() bool { } } -func (self *SStoragecachedimage) markDeleting(ctx context.Context, userCred mcclient.TokenCredential) error { +func (self *SStoragecachedimage) markDeleting(ctx context.Context, userCred mcclient.TokenCredential, isForce bool) error { err := self.ValidateDeleteCondition(ctx) if err != nil { return err @@ -237,7 +237,7 @@ func (self *SStoragecachedimage) markDeleting(ctx context.Context, userCred mccl lockman.LockJointObject(ctx, cache, image) defer lockman.ReleaseJointObject(ctx, cache, image) - if utils.IsInStringArray(self.Status, []string{CACHED_IMAGE_STATUS_READY, CACHED_IMAGE_STATUS_DELETING}) { + if !isForce && ! utils.IsInStringArray(self.Status, []string{CACHED_IMAGE_STATUS_READY, CACHED_IMAGE_STATUS_DELETING}) { return httperrors.NewInvalidStatusError("Cannot uncache in status %s", self.Status) } _, err = self.GetModelManager().TableSpec().Update(self, func() error { diff --git a/pkg/compute/models/storagecaches.go b/pkg/compute/models/storagecaches.go index 04546e8429..5e46d08c31 100644 --- a/pkg/compute/models/storagecaches.go +++ b/pkg/compute/models/storagecaches.go @@ -315,7 +315,7 @@ func (self *SStoragecache) PerformUncacheImage(ctx context.Context, userCred mcc return nil, err } - err = scimg.markDeleting(ctx, userCred) + err = scimg.markDeleting(ctx, userCred, isForce) if err != nil { return nil, httperrors.NewInvalidStatusError("Fail to mark cache status: %s", err) } From fdcb0148c3f08ce0c2b0f285e2284d2b46970186 Mon Sep 17 00:00:00 2001 From: Zexi Li Date: Mon, 22 Oct 2018 19:49:59 +0800 Subject: [PATCH 18/23] fix: server list and details no is_gpu field --- pkg/compute/models/guests.go | 17 +++++++++++++++++ 1 file changed, 17 insertions(+) diff --git a/pkg/compute/models/guests.go b/pkg/compute/models/guests.go index 9c952822f5..c3b42f105e 100644 --- a/pkg/compute/models/guests.go +++ b/pkg/compute/models/guests.go @@ -914,6 +914,12 @@ func (self *SGuest) GetCustomizeColumns(ctx context.Context, userCred mcclient.T extra.Add(jsonutils.NewString(timeutils.FullIsoTime(pendingDeletedAt)), "auto_delete_at") } + isGpu := jsonutils.JSONFalse + if self.isGpu() { + isGpu = jsonutils.JSONTrue + } + extra.Add(isGpu, "is_gpu") + return self.moreExtraInfo(extra) } @@ -965,6 +971,13 @@ func (self *SGuest) GetExtraDetails(ctx context.Context, userCred mcclient.Token } extra.Add(jsonutils.NewString(self.getAdminSecurityRules()), "admin_security_rules") } + + isGpu := jsonutils.JSONFalse + if self.isGpu() { + isGpu = jsonutils.JSONTrue + } + extra.Add(isGpu, "is_gpu") + return self.moreExtraInfo(extra) } @@ -1166,6 +1179,10 @@ func (self *SGuest) getAdminSecurityRules() string { } } +func (self *SGuest) isGpu() bool { + return len(self.GetIsolatedDevices()) != 0 +} + func (self *SGuest) GetIsolatedDevices() []SIsolatedDevice { return IsolatedDeviceManager.findAttachedDevicesOfGuest(self) } From a3a5259505a5c254b969955b6b70433fe8db7869 Mon Sep 17 00:00:00 2001 From: wanyaoqi Date: Tue, 23 Oct 2018 22:07:04 +0800 Subject: [PATCH 19/23] add disk type in snapshots --- pkg/cloudprovider/resources.go | 1 + pkg/compute/models/snapshots.go | 8 +++++--- pkg/util/aliyun/snapshot.go | 4 ++++ pkg/util/azure/snapshot.go | 4 ++++ 4 files changed, 14 insertions(+), 3 deletions(-) diff --git a/pkg/cloudprovider/resources.go b/pkg/cloudprovider/resources.go index f93abdc326..0935ffb38f 100644 --- a/pkg/cloudprovider/resources.go +++ b/pkg/cloudprovider/resources.go @@ -252,6 +252,7 @@ type ICloudSnapshot interface { GetManagerId() string GetSize() int32 GetDiskId() string + GetDiskType() string Delete() error GetRegionId() string } diff --git a/pkg/compute/models/snapshots.go b/pkg/compute/models/snapshots.go index 34fe6a4945..9d2445cd08 100644 --- a/pkg/compute/models/snapshots.go +++ b/pkg/compute/models/snapshots.go @@ -47,6 +47,7 @@ type SSnapshot struct { Size int `nullable:"false" list:"user"` // MB OutOfChain bool `nullable:"false" default:"false" index:"true" list:"admin"` FakeDeleted bool `nullable:"false" default:"false" index:"true"` + DiskType string `width:"32" charset:"ascii" nullable:"true" list:"user"` CloudregionId string `width:"36" charset:"ascii" nullable:"true" list:"user"` } @@ -144,8 +145,7 @@ func (self *SSnapshot) getMoreDetails(extra *jsonutils.JSONDict) *jsonutils.JSON } disk, _ := self.GetDisk() if disk != nil { - extra.Add(jsonutils.NewString(disk.DiskType), "disk_type") - + // extra.Add(jsonutils.NewString(disk.DiskType), "disk_type") guests := disk.GetGuests() if len(guests) == 1 { extra.Add(jsonutils.NewString(guests[0].Name), "guest") @@ -271,6 +271,7 @@ func (self *SSnapshotManager) CreateSnapshot(ctx context.Context, userCred mccli snapshot.DiskId = disk.Id snapshot.StorageId = disk.StorageId snapshot.Size = disk.DiskSize + snapshot.DiskType = disk.DiskType snapshot.Location = location snapshot.CreatedBy = createdBy snapshot.Name = name @@ -428,10 +429,10 @@ func totalSnapshotCount(projectId string) int { return count } -// Only sync snapshot status func (self *SSnapshot) SyncWithCloudSnapshot(userCred mcclient.TokenCredential, ext cloudprovider.ICloudSnapshot) error { _, err := self.GetModelManager().TableSpec().Update(self, func() error { self.Status = ext.GetStatus() + self.DiskType = ext.GetDiskType() return nil }) if err != nil { @@ -456,6 +457,7 @@ func (manager *SSnapshotManager) newFromCloudSnapshot(userCred mcclient.TokenCre } } + snapshot.DiskType = extSnapshot.GetDiskType() snapshot.Size = int(extSnapshot.GetSize()) * 1024 snapshot.ManagerId = extSnapshot.GetManagerId() snapshot.CloudregionId = region.Id diff --git a/pkg/util/aliyun/snapshot.go b/pkg/util/aliyun/snapshot.go index 3e896607a0..fc1ee257bc 100644 --- a/pkg/util/aliyun/snapshot.go +++ b/pkg/util/aliyun/snapshot.go @@ -64,6 +64,10 @@ func (self *SSnapshot) GetDiskId() string { return self.SourceDiskId } +func (self *SSnapshot) GetDiskType() string { + return self.SourceDiskType +} + func (self *SSnapshot) Refresh() error { if snapshots, total, err := self.region.GetSnapshots("", "", "", []string{self.SnapshotId}, 0, 1); err != nil { return err diff --git a/pkg/util/azure/snapshot.go b/pkg/util/azure/snapshot.go index 246e9a7387..050970c0c5 100644 --- a/pkg/util/azure/snapshot.go +++ b/pkg/util/azure/snapshot.go @@ -165,3 +165,7 @@ func (self *SSnapshot) GetManagerId() string { func (self *SSnapshot) GetRegionId() string { return self.region.GetId() } + +func (self *SSnapshot) GetDiskType() string { + return "" +} From 532c41e5a8d65e8fbe293e47e6af8d0f211cb55e Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=B1=88=E8=BD=A9?= Date: Tue, 23 Oct 2018 22:15:03 +0800 Subject: [PATCH 20/23] =?UTF-8?q?=E9=81=BF=E5=85=8D=E6=9C=89=E9=87=8D?= =?UTF-8?q?=E5=A4=8Dip=E6=97=B6=E8=A6=86=E7=9B=96=E9=83=A8=E5=88=86guest?= =?UTF-8?q?=20ip?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- pkg/compute/models/guestnetworks.go | 6 +++--- pkg/compute/models/guests.go | 2 +- 2 files changed, 4 insertions(+), 4 deletions(-) diff --git a/pkg/compute/models/guestnetworks.go b/pkg/compute/models/guestnetworks.go index b2cc9b7df5..4fe7a290a0 100644 --- a/pkg/compute/models/guestnetworks.go +++ b/pkg/compute/models/guestnetworks.go @@ -338,10 +338,10 @@ func (manager *SGuestnetworkManager) DeleteGuestNics(ctx context.Context, guest return nil } -func (manager *SGuestnetworkManager) getGuestNicByIP(ip string) (*SGuestnetwork, error) { +func (manager *SGuestnetworkManager) getGuestNicByIP(ip string, networkId string) (*SGuestnetwork, error) { gn := SGuestnetwork{} q := manager.Query() - q = q.Equals("ip_addr", ip) + q = q.Equals("ip_addr", ip).Equals("network_id", networkId) err := q.First(&gn) if err != nil { if err != sql.ErrNoRows { @@ -578,4 +578,4 @@ func (manager *SGuestnetworkManager) getRecentlyReleasedIPAddresses(networkId st } } return ret -} \ No newline at end of file +} diff --git a/pkg/compute/models/guests.go b/pkg/compute/models/guests.go index 2c8f06fe80..2ea38009ac 100644 --- a/pkg/compute/models/guests.go +++ b/pkg/compute/models/guests.go @@ -1635,7 +1635,7 @@ func (self *SGuest) SyncVMNics(ctx context.Context, userCred mcclient.TokenCrede continue // cannot determine which network it attached to } // check if the IP has been occupied, if yes, release the IP - gn, err := GuestnetworkManager.getGuestNicByIP(add.nic.GetIP()) + gn, err := GuestnetworkManager.getGuestNicByIP(add.nic.GetIP(), add.net.Id) if err != nil { result.AddError(err) continue From a5be2a406ae5089711d71d9167157b8f753851e6 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=B1=88=E8=BD=A9?= Date: Wed, 24 Oct 2018 10:33:17 +0800 Subject: [PATCH 21/23] =?UTF-8?q?=E9=81=BF=E5=85=8D=E5=9B=A0=E4=B8=AD?= =?UTF-8?q?=E6=96=87=E9=97=AE=E9=A2=98=E5=AF=BC=E8=87=B4=E5=90=8C=E6=AD=A5?= =?UTF-8?q?=E6=97=B6=E5=87=BA=E7=8E=B0=E7=9B=B8=E5=90=8C=E7=9A=84Azure?= =?UTF-8?q?=E8=B5=84=E6=BA=90?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- pkg/cloudcommon/db/standalone.go | 2 +- pkg/compute/models/storagecachedimages.go | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/pkg/cloudcommon/db/standalone.go b/pkg/cloudcommon/db/standalone.go index d7d2b74a14..b0f72c8d85 100644 --- a/pkg/cloudcommon/db/standalone.go +++ b/pkg/cloudcommon/db/standalone.go @@ -20,7 +20,7 @@ type SStandaloneResourceBase struct { Id string `width:"128" charset:"ascii" primary:"true" list:"user"` Name string `width:"128" charset:"utf8" nullable:"false" index:"true" list:"user" update:"user" create:"required"` - ExternalId string `width:"256" charset:"ascii" index:"true" list:"admin" create:"admin_optional"` + ExternalId string `width:"256" charset:"utf8" index:"true" list:"admin" create:"admin_optional"` Description string `width:"256" charset:"utf8" get:"user" list:"user" update:"user" create:"optional"` diff --git a/pkg/compute/models/storagecachedimages.go b/pkg/compute/models/storagecachedimages.go index 0bafb4624b..c0e458185c 100644 --- a/pkg/compute/models/storagecachedimages.go +++ b/pkg/compute/models/storagecachedimages.go @@ -57,7 +57,7 @@ type SStoragecachedimage struct { StoragecacheId string `width:"36" charset:"ascii" nullable:"false" list:"admin" create:"admin_required" key_index:"true"` CachedimageId string `width:"36" charset:"ascii" nullable:"false" list:"admin" create:"admin_required" key_index:"true"` - ExternalId string `width:"256" charset:"ascii" nullable:"false" get:"admin"` + ExternalId string `width:"256" charset:"utf8" nullable:"false" get:"admin"` Status string `width:"32" charset:"ascii" nullable:"false" default:"init" list:"admin" update:"admin" create:"admin_required"` // = Column(VARCHAR(32, charset='ascii'), nullable=False, Path string `width:"256" charset:"utf8" nullable:"true" list:"admin" update:"admin" create:"admin_optional"` // = Column(VARCHAR(256, charset='utf8'), nullable=True) From 350f4e4db95231990b0671ffbacd9f064bbea5a2 Mon Sep 17 00:00:00 2001 From: Qiu Jian Date: Wed, 24 Oct 2018 12:00:53 +0800 Subject: [PATCH 22/23] =?UTF-8?q?=E4=BF=AE=E6=AD=A3=EF=BC=9A=E4=BF=9D?= =?UTF-8?q?=E6=8A=A4=E7=B3=BB=E7=BB=9F=E9=BB=98=E8=AE=A4=E5=88=9B=E5=BB=BA?= =?UTF-8?q?=E7=9A=84default=E8=B5=84=E6=BA=90=EF=BC=8C=E4=B8=8D=E5=85=81?= =?UTF-8?q?=E8=AE=B8=E5=88=A0=E9=99=A4?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- pkg/compute/models/cloudregions.go | 3 +++ pkg/compute/models/secgroups.go | 11 +++++++++++ pkg/compute/models/vpcs.go | 3 +++ pkg/httperrors/errors.go | 5 +++++ pkg/httperrors/httperrors.go | 4 ++++ 5 files changed, 26 insertions(+) diff --git a/pkg/compute/models/cloudregions.go b/pkg/compute/models/cloudregions.go index 3d0e2e9a53..7b84f1ca42 100644 --- a/pkg/compute/models/cloudregions.go +++ b/pkg/compute/models/cloudregions.go @@ -57,6 +57,9 @@ func (self *SCloudregion) ValidateDeleteCondition(ctx context.Context) error { if self.GetZoneCount() > 0 || self.GetVpcCount() > 0 { return httperrors.NewNotEmptyError("not empty cloud region") } + if self.Id == "default" { + return httperrors.NewProtectedResourceError("not allow to delete default cloud region") + } return self.SEnabledStatusStandaloneResourceBase.ValidateDeleteCondition(ctx) } diff --git a/pkg/compute/models/secgroups.go b/pkg/compute/models/secgroups.go index bdba6dddf3..68fcd19623 100644 --- a/pkg/compute/models/secgroups.go +++ b/pkg/compute/models/secgroups.go @@ -367,3 +367,14 @@ func (manager *SSecurityGroupManager) InitializeData() error { } return nil } + +func (self *SSecurityGroup) ValidateDeleteCondition(ctx context.Context) error { + cnt := self.GetGuestsCount() + if cnt > 0 { + return httperrors.NewNotEmptyError("the security group is in use") + } + if self.Id == "default" { + return httperrors.NewProtectedResourceError("not allow to delete default security group") + } + return self.SSharableVirtualResourceBase.ValidateDeleteCondition(ctx) +} diff --git a/pkg/compute/models/vpcs.go b/pkg/compute/models/vpcs.go index 67fd0cbc82..2ec8e1bd24 100644 --- a/pkg/compute/models/vpcs.go +++ b/pkg/compute/models/vpcs.go @@ -78,6 +78,9 @@ func (self *SVpc) ValidateDeleteCondition(ctx context.Context) error { if self.GetNetworkCount() > 0 { return httperrors.NewNotEmptyError("VPC not empty") } + if self.Id == "default" { + return httperrors.NewProtectedResourceError("not allow to delete default vpc") + } return self.SEnabledStatusStandaloneResourceBase.ValidateDeleteCondition(ctx) } diff --git a/pkg/httperrors/errors.go b/pkg/httperrors/errors.go index 662b59ec0e..da247ea7a8 100644 --- a/pkg/httperrors/errors.go +++ b/pkg/httperrors/errors.go @@ -188,3 +188,8 @@ func NewGeneralError(err error) *httputils.JSONClientError { return NewInternalServerError(err.Error()) } } + +func NewProtectedResourceError(msg string, params ...interface{}) *httputils.JSONClientError { + msg, err := errorMessage(msg, params) + return NewJsonClientError(403, "ProtectedResourceError(", msg, err) +} diff --git a/pkg/httperrors/httperrors.go b/pkg/httperrors/httperrors.go index 84dc0d3216..4974d2f01e 100644 --- a/pkg/httperrors/httperrors.go +++ b/pkg/httperrors/httperrors.go @@ -90,3 +90,7 @@ func TenantNotFoundError(w http.ResponseWriter, msg string, params ...interface{ func OutOfQuotaError(w http.ResponseWriter, msg string, params ...interface{}) { JsonClientError(w, NewOutOfQuotaError(msg, params...)) } + +func ProtectedResourceError(w http.ResponseWriter, msg string, params ...interface{}) { + JsonClientError(w, NewProtectedResourceError(msg, params...)) +} From 5d7bdace41aeb5493f9ab10b730029325dd55478 Mon Sep 17 00:00:00 2001 From: Qiu Jian Date: Thu, 25 Oct 2018 00:04:04 +0800 Subject: [PATCH 23/23] =?UTF-8?q?=E4=BF=AE=E6=AD=A3=EF=BC=9A=E6=97=A0?= =?UTF-8?q?=E6=B3=95=E5=88=A0=E9=99=A4sched=5Ffail=E7=9A=84=E4=B8=BB?= =?UTF-8?q?=E6=9C=BA?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- pkg/compute/guestdrivers/virtualization.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pkg/compute/guestdrivers/virtualization.go b/pkg/compute/guestdrivers/virtualization.go index 04c1307f88..c0db1e612b 100644 --- a/pkg/compute/guestdrivers/virtualization.go +++ b/pkg/compute/guestdrivers/virtualization.go @@ -119,7 +119,7 @@ func (self *SVirtualizedGuestDriver) RequestStopGuestForDelete(ctx context.Conte if host != nil && host.Enabled && host.HostStatus == models.HOST_ONLINE { return guest.StartGuestStopTask(ctx, task.GetUserCred(), true, task.GetTaskId()) } - if !jsonutils.QueryBoolean(task.GetParams(), "purge", false) { + if host != nil && !jsonutils.QueryBoolean(task.GetParams(), "purge", false) { return fmt.Errorf("fail to contact host") } task.ScheduleRun(nil)