From 11fe9bf6bcab32933bb47908db7f1155c73e689d Mon Sep 17 00:00:00 2001 From: wanyaoqi <18528551+wanyaoqi@users.noreply.github.com> Date: Tue, 29 Aug 2023 08:40:20 +0800 Subject: [PATCH] fix(region,host): support limit guest nic traffic (#17852) --- cmd/climc/shell/compute/servers.go | 2 + pkg/apis/compute/api.go | 14 +- pkg/apis/compute/guest_const.go | 3 + pkg/apis/compute/guestnetwork.go | 39 ++--- pkg/apis/compute/guests.go | 6 + pkg/cloudcommon/cmdline/parser.go | 12 ++ pkg/cloudcommon/db/opslog_const.go | 3 + pkg/compute/guestdrivers/base.go | 8 ++ pkg/compute/guestdrivers/kvm.go | 24 ++++ pkg/compute/models/guest_actions.go | 48 +++++++ pkg/compute/models/guestdrivers.go | 2 + pkg/compute/models/guestnetworks.go | 46 +++++- pkg/compute/models/guests.go | 48 ++++--- pkg/compute/models/hosts.go | 33 +++++ .../tasks/guest_sync_nic_traffics_task.go | 126 ++++++++++++++++ .../guestman/guesthandlers/guesthandler.go | 27 ++++ pkg/hostman/guestman/guestman.go | 136 ++++++++++++++++++ pkg/hostman/guestman/qemu-kvm.go | 4 + pkg/hostman/guestman/qemu-kvmhelper.go | 33 +++++ pkg/hostman/guestman/qemu/generate.go | 23 ++- pkg/hostman/hostmetrics/hostmetrics.go | 108 ++++++++++++-- pkg/mcclient/options/compute/servers.go | 11 ++ pkg/util/logclient/consts.go | 2 + 23 files changed, 700 insertions(+), 58 deletions(-) create mode 100644 pkg/compute/tasks/guest_sync_nic_traffics_task.go diff --git a/cmd/climc/shell/compute/servers.go b/cmd/climc/shell/compute/servers.go index 3939a67236..86c7a910ab 100644 --- a/cmd/climc/shell/compute/servers.go +++ b/cmd/climc/shell/compute/servers.go @@ -117,6 +117,8 @@ func init() { cmd.Perform("qga-ping", &options.ServerQgaPing{}) cmd.Perform("set-password", &options.ServerSetPasswordOptions{}) cmd.Perform("set-boot-index", &options.ServerSetBootIndexOptions{}) + cmd.Perform("reset-nic-traffic-limit", &options.ServerNicTrafficLimitOptions{}) + cmd.Perform("set-nic-traffic-limit", &options.ServerNicTrafficLimitOptions{}) cmd.Get("vnc", new(options.ServerVncOptions)) cmd.Get("desc", new(options.ServerIdOptions)) diff --git a/pkg/apis/compute/api.go b/pkg/apis/compute/api.go index ee3b9ad5e4..7cb1a4a45e 100644 --- a/pkg/apis/compute/api.go +++ b/pkg/apis/compute/api.go @@ -75,12 +75,14 @@ type NetworkConfig struct { // 驱动方式 // 若指定镜像的网络驱动方式,此参数会被覆盖 - Driver string `json:"driver"` - BwLimit int `json:"bw_limit"` - Vip bool `json:"vip"` - Reserved bool `json:"reserved"` - NetType string `json:"net_type"` - NumQueues int `json:"num_queues"` + Driver string `json:"driver"` + BwLimit int `json:"bw_limit"` + Vip bool `json:"vip"` + Reserved bool `json:"reserved"` + NetType string `json:"net_type"` + NumQueues int `json:"num_queues"` + RxTrafficLimit int64 `json:"rx_traffic_limit"` + TxTrafficLimit int64 `json:"tx_traffic_limit"` // sriov nic SriovDevice *IsolatedDeviceConfig `json:"sriov_device"` diff --git a/pkg/apis/compute/guest_const.go b/pkg/apis/compute/guest_const.go index 9511774c48..21a3db3eb0 100644 --- a/pkg/apis/compute/guest_const.go +++ b/pkg/apis/compute/guest_const.go @@ -108,6 +108,9 @@ const ( VM_SYNC_CONFIG = compute.VM_SYNC_CONFIG VM_SYNC_FAIL = "sync_fail" + VM_SYNC_TRAFFIC_LIMIT = "sync_traffic_limit" + VM_SYNC_TRAFFIC_LIMIT_FAILED = "sync_traffic_limit_failed" + VM_START_RESIZE_DISK = "start_resize_disk" VM_RESIZE_DISK = "resize_disk" VM_RESIZE_DISK_FAILED = "resize_disk_fail" diff --git a/pkg/apis/compute/guestnetwork.go b/pkg/apis/compute/guestnetwork.go index ceb34a1917..ec86655c16 100644 --- a/pkg/apis/compute/guestnetwork.go +++ b/pkg/apis/compute/guestnetwork.go @@ -80,22 +80,24 @@ type GuestnetworkUpdateInput struct { } type GuestnetworkBaseDesc struct { - Net string `json:"net"` - NetId string `json:"net_id"` - Mac string `json:"mac"` - Virtual bool `json:"virtual"` - Ip string `json:"ip"` - Gateway string `json:"gateway"` - Dns string `json:"dns"` - Domain string `json:"domain"` - Ntp string `json:"ntp"` - Routes jsonutils.JSONObject `json:"routes"` - Ifname string `json:"ifname"` - Masklen int8 `json:"masklen"` - Vlan int `json:"vlan"` - Bw int `json:"bw"` - Mtu int16 `json:"mtu"` - Index int8 `json:"index"` + Net string `json:"net"` + NetId string `json:"net_id"` + Mac string `json:"mac"` + Virtual bool `json:"virtual"` + Ip string `json:"ip"` + Gateway string `json:"gateway"` + Dns string `json:"dns"` + Domain string `json:"domain"` + Ntp string `json:"ntp"` + Routes jsonutils.JSONObject `json:"routes"` + Ifname string `json:"ifname"` + Masklen int8 `json:"masklen"` + Vlan int `json:"vlan"` + Bw int `json:"bw"` + Mtu int16 `json:"mtu"` + Index int8 `json:"index"` + RxTrafficLimit int64 `json:"rx_traffic_limit"` + TxTrafficLimit int64 `json:"tx_traffic_limit"` Bridge string `json:"bridge"` WireId string `json:"wire_id"` @@ -133,3 +135,8 @@ type GuestnetworkJsonDesc struct { LinkUp bool `json:"link_up"` } + +type SNicTrafficRecord struct { + RxTraffic int64 + TxTraffic int64 +} diff --git a/pkg/apis/compute/guests.go b/pkg/apis/compute/guests.go index 803d2435af..a64eb33a7a 100644 --- a/pkg/apis/compute/guests.go +++ b/pkg/apis/compute/guests.go @@ -1016,3 +1016,9 @@ type ServerSetLiveMigrateParamsInput struct { MaxBandwidthMB *int64 DowntimeLimitMS *int64 } + +type ServerNicTrafficLimit struct { + Mac string `json:"mac"` + RxTrafficLimit *int64 `json:"rx_traffic_limit"` + TxTrafficLimit *int64 `json:"tx_traffic_limit"` +} diff --git a/pkg/cloudcommon/cmdline/parser.go b/pkg/cloudcommon/cmdline/parser.go index 86d96c009b..c8a210e5b9 100644 --- a/pkg/cloudcommon/cmdline/parser.go +++ b/pkg/cloudcommon/cmdline/parser.go @@ -277,6 +277,18 @@ func ParseNetworkConfig(desc string, idx int) (*compute.NetworkConfig, error) { netConfig.SriovDevice = &compute.IsolatedDeviceConfig{ Model: p[len("sriov-nic-model="):], } + } else if strings.HasPrefix(p, "rx-traffic-limit=") { + var err error + netConfig.RxTrafficLimit, err = strconv.ParseInt(p[len("rx-traffic-limit="):], 10, 0) + if err != nil { + return nil, errors.Wrap(err, "parse rx-traffic-limit") + } + } else if strings.HasPrefix(p, "tx-traffic-limit=") { + var err error + netConfig.TxTrafficLimit, err = strconv.ParseInt(p[len("tx-traffic-limit="):], 10, 0) + if err != nil { + return nil, errors.Wrap(err, "parse tx-traffic-limit") + } } else if utils.IsInStringArray(p, compute.ALL_NETWORK_TYPES) { netConfig.NetType = p } else { diff --git a/pkg/cloudcommon/db/opslog_const.go b/pkg/cloudcommon/db/opslog_const.go index f98cfe7f4e..e2dd67608b 100644 --- a/pkg/cloudcommon/db/opslog_const.go +++ b/pkg/cloudcommon/db/opslog_const.go @@ -317,6 +317,9 @@ const ( ACT_ENCRYPT_FAIL = "encrypt_fail" ACT_ENCRYPT_DONE = "encrypted" + ACT_SYNC_TRAFFIC_LIMIT = "sync_traffic_limit" + ACT_SYNC_TRAFFIC_LIMIT_FAIL = "sync_traffic_limit_fail" + ACT_BIND = "bind" ACT_UNBIND = "unbind" ) diff --git a/pkg/compute/guestdrivers/base.go b/pkg/compute/guestdrivers/base.go index d880465fff..5108439c0e 100644 --- a/pkg/compute/guestdrivers/base.go +++ b/pkg/compute/guestdrivers/base.go @@ -516,3 +516,11 @@ func (self *SBaseGuestDriver) FetchMonitorUrl(ctx context.Context, guest *models } return influxdbUrl } + +func (self *SBaseGuestDriver) RequestResetNicTrafficLimit(ctx context.Context, task taskman.ITask, host *models.SHost, guest *models.SGuest, input *api.ServerNicTrafficLimit) error { + return httperrors.ErrNotImplemented +} + +func (self *SBaseGuestDriver) RequestSetNicTrafficLimit(ctx context.Context, task taskman.ITask, host *models.SHost, guest *models.SGuest, input *api.ServerNicTrafficLimit) error { + return httperrors.ErrNotImplemented +} diff --git a/pkg/compute/guestdrivers/kvm.go b/pkg/compute/guestdrivers/kvm.go index 7b1483a641..50374b88bc 100644 --- a/pkg/compute/guestdrivers/kvm.go +++ b/pkg/compute/guestdrivers/kvm.go @@ -1081,3 +1081,27 @@ func (self *SKVMGuestDriver) FetchMonitorUrl(ctx context.Context, guest *models. } return self.SVirtualizedGuestDriver.FetchMonitorUrl(ctx, guest) } + +func (self *SKVMGuestDriver) RequestResetNicTrafficLimit(ctx context.Context, task taskman.ITask, host *models.SHost, guest *models.SGuest, input *api.ServerNicTrafficLimit) error { + url := fmt.Sprintf("%s/servers/%s/reset-nic-traffic-limit", host.ManagerUri, guest.Id) + httpClient := httputils.GetDefaultClient() + header := task.GetTaskRequestHeader() + body := jsonutils.Marshal(input) + _, _, err := httputils.JSONRequest(httpClient, ctx, "POST", url, header, body, false) + if err != nil { + return errors.Wrap(err, "host request") + } + return nil +} + +func (self *SKVMGuestDriver) RequestSetNicTrafficLimit(ctx context.Context, task taskman.ITask, host *models.SHost, guest *models.SGuest, input *api.ServerNicTrafficLimit) error { + url := fmt.Sprintf("%s/servers/%s/set-nic-traffic-limit", host.ManagerUri, guest.Id) + httpClient := httputils.GetDefaultClient() + header := task.GetTaskRequestHeader() + body := jsonutils.Marshal(input) + _, _, err := httputils.JSONRequest(httpClient, ctx, "POST", url, header, body, false) + if err != nil { + return errors.Wrap(err, "host request") + } + return nil +} diff --git a/pkg/compute/models/guest_actions.go b/pkg/compute/models/guest_actions.go index 9793784c55..72fc10a17c 100644 --- a/pkg/compute/models/guest_actions.go +++ b/pkg/compute/models/guest_actions.go @@ -5266,6 +5266,54 @@ func (guest *SGuest) StartDeleteGuestSnapshots(ctx context.Context, userCred mcc return nil } +func (self *SGuest) PerformResetNicTrafficLimit(ctx context.Context, userCred mcclient.TokenCredential, + query jsonutils.JSONObject, input *api.ServerNicTrafficLimit) (jsonutils.JSONObject, error) { + + if !utils.IsInStringArray(self.Status, []string{api.VM_READY, api.VM_RUNNING}) { + return nil, httperrors.NewUnsupportOperationError("The guest status need be %s or %s, current is %s", api.VM_READY, api.VM_RUNNING, self.Status) + } + input.Mac = strings.ToLower(input.Mac) + _, err := self.GetGuestnetworkByMac(input.Mac) + if err != nil { + return nil, errors.Wrap(err, "get guest network by mac") + } + + params := jsonutils.Marshal(input).(*jsonutils.JSONDict) + params.Set("old_status", jsonutils.NewString(self.Status)) + self.SetStatus(userCred, api.VM_SYNC_TRAFFIC_LIMIT, "PerformResetNicTrafficLimit") + task, err := taskman.TaskManager.NewTask(ctx, "GuestResetNicTrafficsTask", self, userCred, params, "", "", nil) + if err != nil { + return nil, err + } + task.ScheduleRun(nil) + return nil, nil +} + +func (self *SGuest) PerformSetNicTrafficLimit(ctx context.Context, userCred mcclient.TokenCredential, + query jsonutils.JSONObject, input *api.ServerNicTrafficLimit) (jsonutils.JSONObject, error) { + + if !utils.IsInStringArray(self.Status, []string{api.VM_READY, api.VM_RUNNING}) { + return nil, httperrors.NewUnsupportOperationError("The guest status need be %s or %s, current is %s", api.VM_READY, api.VM_RUNNING, self.Status) + } + if input.RxTrafficLimit == nil && input.TxTrafficLimit == nil { + return nil, httperrors.NewBadRequestError("rx/tx traffic not provider") + } + input.Mac = strings.ToLower(input.Mac) + _, err := self.GetGuestnetworkByMac(input.Mac) + if err != nil { + return nil, errors.Wrap(err, "get guest network by mac") + } + params := jsonutils.Marshal(input).(*jsonutils.JSONDict) + params.Set("old_status", jsonutils.NewString(self.Status)) + self.SetStatus(userCred, api.VM_SYNC_TRAFFIC_LIMIT, "GuestSetNicTrafficsTask") + task, err := taskman.TaskManager.NewTask(ctx, "GuestSetNicTrafficsTask", self, userCred, params, "", "", nil) + if err != nil { + return nil, err + } + task.ScheduleRun(nil) + return nil, nil +} + func (self *SGuest) PerformBindGroups(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) (jsonutils.JSONObject, error) { diff --git a/pkg/compute/models/guestdrivers.go b/pkg/compute/models/guestdrivers.go index a100beced6..d5814beaf7 100644 --- a/pkg/compute/models/guestdrivers.go +++ b/pkg/compute/models/guestdrivers.go @@ -236,6 +236,8 @@ type IGuestDriver interface { RequestQgaCommand(ctx context.Context, userCred mcclient.TokenCredential, body jsonutils.JSONObject, host *SHost, guest *SGuest) (jsonutils.JSONObject, error) FetchMonitorUrl(ctx context.Context, guest *SGuest) string + RequestResetNicTrafficLimit(ctx context.Context, task taskman.ITask, host *SHost, guest *SGuest, input *api.ServerNicTrafficLimit) error + RequestSetNicTrafficLimit(ctx context.Context, task taskman.ITask, host *SHost, guest *SGuest, input *api.ServerNicTrafficLimit) error } var guestDrivers map[string]IGuestDriver diff --git a/pkg/compute/models/guestnetworks.go b/pkg/compute/models/guestnetworks.go index 4259006ff0..0e5f991d7d 100644 --- a/pkg/compute/models/guestnetworks.go +++ b/pkg/compute/models/guestnetworks.go @@ -89,6 +89,12 @@ type SGuestnetwork struct { NumQueues int `nullable:"true" default:"1" list:"user" update:"user"` // 带宽限制,单位mbps BwLimit int `nullable:"false" default:"0" list:"user"` + // 下行流量限制,单位 bytes + RxTrafficLimit int64 `nullable:"false" default:"0" list:"user"` + RxTrafficUsed int64 `nullable:"false" default:"0" list:"user"` + // 上行流量限制,单位 bytes + TxTrafficLimit int64 `nullable:"false" default:"0" list:"user"` + TxTrafficUsed int64 `nullable:"false" default:"0" list:"user"` // 网卡序号 Index int8 `nullable:"false" default:"0" list:"user" update:"user"` // 是否为虚拟接口(无IP) @@ -229,12 +235,14 @@ type newGuestNetworkArgs struct { requireDesignatedIP bool useDesignatedIP bool - ifname string - macAddr string - bwLimit int - nicDriver string - numQueues int - teamWithMac string + ifname string + macAddr string + bwLimit int + nicDriver string + numQueues int + teamWithMac string + rxTrafficLimit int64 + txTrafficLimit int64 virtual bool } @@ -274,6 +282,8 @@ func (manager *SGuestnetworkManager) newGuestNetwork( } gn.Driver = driver gn.NumQueues = numQueues + gn.RxTrafficLimit = args.rxTrafficLimit + gn.TxTrafficLimit = args.txTrafficLimit if bwLimit >= 0 { gn.BwLimit = bwLimit } @@ -579,6 +589,8 @@ func (self *SGuestnetwork) getJsonDesc() *api.GuestnetworkJsonDesc { desc.Masklen = net.GuestIpMask desc.Driver = self.Driver desc.NumQueues = self.NumQueues + desc.RxTrafficLimit = self.RxTrafficLimit + desc.TxTrafficLimit = self.TxTrafficLimit desc.Vlan = net.VlanId desc.Bw = self.getBandwidth() desc.Mtu = self.getMtu(net) @@ -606,6 +618,28 @@ func (self *SGuestnetwork) IsSriovWithoutOffload() bool { return true } +func (self *SGuestnetwork) UpdateNicTrafficUsed(rx, tx int64) error { + _, err := db.Update(self, func() error { + self.RxTrafficUsed = rx + self.TxTrafficUsed = tx + return nil + }) + return err +} + +func (self *SGuestnetwork) UpdateNicTrafficLimit(rx, tx *int64) error { + _, err := db.Update(self, func() error { + if rx != nil { + self.RxTrafficLimit = *rx + } + if tx != nil { + self.TxTrafficLimit = *tx + } + return nil + }) + return err +} + func (manager *SGuestnetworkManager) GetGuestByAddress(address string) *SGuest { networks := manager.TableSpec().Instance() guests := GuestManager.Query() diff --git a/pkg/compute/models/guests.go b/pkg/compute/models/guests.go index 0250fd2f47..619cf722ff 100644 --- a/pkg/compute/models/guests.go +++ b/pkg/compute/models/guests.go @@ -3272,10 +3272,12 @@ type Attach2NetworkArgs struct { RequireDesignatedIP bool UseDesignatedIP bool - BwLimit int - NicDriver string - NumQueues int - NicConfs []SNicConfig + BwLimit int + NicDriver string + NumQueues int + RxTrafficLimit int64 + TxTrafficLimit int64 + NicConfs []SNicConfig Virtual bool @@ -3295,10 +3297,12 @@ func (args *Attach2NetworkArgs) onceArgs(i int) attach2NetworkOnceArgs { requireDesignatedIP: args.RequireDesignatedIP, useDesignatedIP: args.UseDesignatedIP, - bwLimit: args.BwLimit, - nicDriver: args.NicDriver, - numQueues: args.NumQueues, - nicConf: args.NicConfs[i], + bwLimit: args.BwLimit, + nicDriver: args.NicDriver, + numQueues: args.NumQueues, + txTrafficLimit: args.TxTrafficLimit, + rxTrafficLimit: args.RxTrafficLimit, + nicConf: args.NicConfs[i], virtual: args.Virtual, @@ -3327,11 +3331,13 @@ type attach2NetworkOnceArgs struct { requireDesignatedIP bool useDesignatedIP bool - bwLimit int - nicDriver string - numQueues int - nicConf SNicConfig - teamWithMac string + bwLimit int + nicDriver string + numQueues int + nicConf SNicConfig + teamWithMac string + rxTrafficLimit int64 + txTrafficLimit int64 virtual bool @@ -3396,12 +3402,14 @@ func (self *SGuest) attach2NetworkOnce( requireDesignatedIP: args.requireDesignatedIP, useDesignatedIP: args.useDesignatedIP, - ifname: args.nicConf.Ifname, - macAddr: args.nicConf.Mac, - bwLimit: args.bwLimit, - nicDriver: nicDriver, - numQueues: args.numQueues, - teamWithMac: args.teamWithMac, + ifname: args.nicConf.Ifname, + macAddr: args.nicConf.Mac, + bwLimit: args.bwLimit, + nicDriver: nicDriver, + numQueues: args.numQueues, + teamWithMac: args.teamWithMac, + rxTrafficLimit: args.rxTrafficLimit, + txTrafficLimit: args.txTrafficLimit, virtual: args.virtual, } @@ -4143,6 +4151,8 @@ func (self *SGuest) attach2NamedNetworkDesc(ctx context.Context, userCred mcclie NicDriver: netConfig.Driver, NumQueues: netConfig.NumQueues, BwLimit: netConfig.BwLimit, + RxTrafficLimit: netConfig.RxTrafficLimit, + TxTrafficLimit: netConfig.TxTrafficLimit, Virtual: netConfig.Vip, TryReserved: netConfig.Reserved, AllocDir: allocDir, diff --git a/pkg/compute/models/hosts.go b/pkg/compute/models/hosts.go index 5b60bfa033..d3088bef57 100644 --- a/pkg/compute/models/hosts.go +++ b/pkg/compute/models/hosts.go @@ -6363,6 +6363,39 @@ func (hh *SHost) GetPinnedCpusetCores(ctx context.Context, userCred mcclient.Tok return ret, nil } +func (h *SHost) PerformSyncGuestNicTraffics(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) (jsonutils.JSONObject, error) { + guestTraffics, err := data.GetMap() + if err != nil { + return nil, errors.Wrap(err, "get guest traffics") + } + for guestId, nicTraffics := range guestTraffics { + nicTrafficMap := make(map[string]api.SNicTrafficRecord) + err = nicTraffics.Unmarshal(&nicTrafficMap) + if err != nil { + log.Errorf("failed unmarshal guest %s nic traffics %s", guestId, err) + continue + } + + guest := GuestManager.FetchGuestById(guestId) + gns, err := guest.GetNetworks("") + if err != nil { + log.Errorf("failed fetch guest %s networks %s", guestId, err) + continue + } + for i := range gns { + nicTraffic, ok := nicTrafficMap[strconv.Itoa(int(gns[i].Index))] + if !ok { + continue + } + if err = gns[i].UpdateNicTrafficUsed(nicTraffic.RxTraffic, nicTraffic.TxTraffic); err != nil { + log.Errorf("failed update guestnetwork %d traffic used %s", gns[i].RowId, err) + continue + } + } + } + return nil, nil +} + func (h *SHost) GetDetailsAppOptions(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject) (jsonutils.JSONObject, error) { return h.Request(ctx, userCred, httputils.GET, "/app-options", nil, nil) } diff --git a/pkg/compute/tasks/guest_sync_nic_traffics_task.go b/pkg/compute/tasks/guest_sync_nic_traffics_task.go new file mode 100644 index 0000000000..286cfc8f30 --- /dev/null +++ b/pkg/compute/tasks/guest_sync_nic_traffics_task.go @@ -0,0 +1,126 @@ +// 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" + "fmt" + + "yunion.io/x/jsonutils" + + "yunion.io/x/onecloud/pkg/apis/compute" + "yunion.io/x/onecloud/pkg/cloudcommon/db" + "yunion.io/x/onecloud/pkg/cloudcommon/db/taskman" + "yunion.io/x/onecloud/pkg/compute/models" + "yunion.io/x/onecloud/pkg/util/logclient" +) + +type GuestResetNicTrafficsTask struct { + SGuestBaseTask +} + +func init() { + taskman.RegisterTask(GuestResetNicTrafficsTask{}) + taskman.RegisterTask(GuestSetNicTrafficsTask{}) +} + +func (self *GuestResetNicTrafficsTask) taskFailed(ctx context.Context, guest *models.SGuest, reason string) { + guest.SetStatus(self.UserCred, compute.VM_SYNC_TRAFFIC_LIMIT, "PerformResetNicTrafficLimit") + db.OpsLog.LogEvent(guest, db.ACT_SYNC_TRAFFIC_LIMIT_FAIL, reason, self.UserCred) + logclient.AddActionLogWithStartable(self, guest, logclient.ACT_SYNC_TRAFFIC_LIMIT, reason, self.UserCred, false) + self.SetStageFailed(ctx, jsonutils.NewString(reason)) +} + +func (self *GuestResetNicTrafficsTask) OnInit(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) { + guest := obj.(*models.SGuest) + host, err := guest.GetHost() + if err != nil { + self.taskFailed(ctx, guest, fmt.Sprintf("get host %s", err)) + return + } + input := &compute.ServerNicTrafficLimit{} + self.GetParams().Unmarshal(input) + self.SetStage("OnResetNicTrafficLimit", nil) + err = guest.GetDriver().RequestResetNicTrafficLimit(ctx, self, host, guest, input) + if err != nil { + self.taskFailed(ctx, guest, err.Error()) + } +} + +func (self *GuestResetNicTrafficsTask) OnResetNicTrafficLimit(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { + input := &compute.ServerNicTrafficLimit{} + self.GetParams().Unmarshal(input) + gn, _ := guest.GetGuestnetworkByMac(input.Mac) + err := gn.UpdateNicTrafficLimit(input.RxTrafficLimit, input.TxTrafficLimit) + if err != nil { + self.taskFailed(ctx, guest, fmt.Sprintf("failed update guest nic traffic limit %s", err)) + return + } + err = gn.UpdateNicTrafficUsed(0, 0) + if err != nil { + self.taskFailed(ctx, guest, fmt.Sprintf("failed update guest nic traffic used %s", err)) + return + } + + oldStatus, _ := self.Params.GetString("old_status") + guest.SetStatus(self.UserCred, oldStatus, "OnResetNicTrafficLimit") + db.OpsLog.LogEvent(guest, db.ACT_SYNC_TRAFFIC_LIMIT, "OnResetNicTrafficLimit", self.UserCred) + logclient.AddActionLogWithStartable(self, guest, logclient.ACT_SYNC_TRAFFIC_LIMIT, "OnResetNicTrafficLimit", self.UserCred, true) + self.SetStageComplete(ctx, nil) +} + +type GuestSetNicTrafficsTask struct { + SGuestBaseTask +} + +func (self *GuestSetNicTrafficsTask) taskFailed(ctx context.Context, guest *models.SGuest, reason string) { + guest.SetStatus(self.UserCred, compute.VM_SYNC_TRAFFIC_LIMIT, "PerformResetNicTrafficLimit") + db.OpsLog.LogEvent(guest, db.ACT_SYNC_TRAFFIC_LIMIT_FAIL, reason, self.UserCred) + logclient.AddActionLogWithStartable(self, guest, logclient.ACT_SYNC_TRAFFIC_LIMIT, reason, self.UserCred, false) + self.SetStageFailed(ctx, jsonutils.NewString(reason)) +} + +func (self *GuestSetNicTrafficsTask) OnInit(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) { + guest := obj.(*models.SGuest) + host, err := guest.GetHost() + if err != nil { + self.taskFailed(ctx, guest, fmt.Sprintf("get host %s", err)) + return + } + input := &compute.ServerNicTrafficLimit{} + self.GetParams().Unmarshal(input) + self.SetStage("OnSetNicTrafficLimit", nil) + err = guest.GetDriver().RequestSetNicTrafficLimit(ctx, self, host, guest, input) + if err != nil { + self.taskFailed(ctx, guest, err.Error()) + } +} + +func (self *GuestSetNicTrafficsTask) OnSetNicTrafficLimit(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { + input := &compute.ServerNicTrafficLimit{} + self.GetParams().Unmarshal(input) + gn, _ := guest.GetGuestnetworkByMac(input.Mac) + err := gn.UpdateNicTrafficLimit(input.RxTrafficLimit, input.TxTrafficLimit) + if err != nil { + self.taskFailed(ctx, guest, fmt.Sprintf("failed update guest nic traffic limit %s", err)) + return + } + + oldStatus, _ := self.Params.GetString("old_status") + guest.SetStatus(self.UserCred, oldStatus, "OnSetNicTrafficLimit") + db.OpsLog.LogEvent(guest, db.ACT_SYNC_TRAFFIC_LIMIT, "OnSetNicTrafficLimit", self.UserCred) + logclient.AddActionLogWithStartable(self, guest, logclient.ACT_SYNC_TRAFFIC_LIMIT, "OnSetNicTrafficLimit", self.UserCred, true) + self.SetStageComplete(ctx, nil) +} diff --git a/pkg/hostman/guestman/guesthandlers/guesthandler.go b/pkg/hostman/guestman/guesthandlers/guesthandler.go index 94168992f8..8afbee3a13 100644 --- a/pkg/hostman/guestman/guesthandlers/guesthandler.go +++ b/pkg/hostman/guestman/guesthandlers/guesthandler.go @@ -97,6 +97,8 @@ func AddGuestTaskHandler(prefix string, app *appsrv.Application) { "qga-set-password": qgaGuestSetPassword, "qga-guest-ping": qgaGuestPing, "qga-command": qgaCommand, + "reset-nic-traffic-limit": guestResetNicTrafficLimit, + "set-nic-traffic-limit": guestSetNicTrafficLimit, } { app.AddHandler("POST", fmt.Sprintf("%s/%s//%s", prefix, keyWord, action), @@ -805,6 +807,31 @@ func guestMemorySnapshotDelete(ctx context.Context, w http.ResponseWriter, r *ht }) } +func guestResetNicTrafficLimit(ctx context.Context, userCred mcclient.TokenCredential, sid string, body jsonutils.JSONObject) (interface{}, error) { + input := new(computeapi.ServerNicTrafficLimit) + if err := body.Unmarshal(input); err != nil { + return nil, httperrors.NewInputParameterError("failed unmarshal input %s", err) + } + + hostutils.DelayTask(ctx, func(ctx context.Context, params interface{}) (jsonutils.JSONObject, error) { + return nil, guestman.GetGuestManager().ResetGuestNicTrafficLimit(sid, input) + }, nil) + + return nil, nil +} + +func guestSetNicTrafficLimit(ctx context.Context, userCred mcclient.TokenCredential, sid string, body jsonutils.JSONObject) (interface{}, error) { + input := new(computeapi.ServerNicTrafficLimit) + if err := body.Unmarshal(input); err != nil { + return nil, httperrors.NewInputParameterError("failed unmarshal input %s", err) + } + + hostutils.DelayTask(ctx, func(ctx context.Context, params interface{}) (jsonutils.JSONObject, error) { + return nil, guestman.GetGuestManager().SetGuestNicTrafficLimit(sid, input) + }, nil) + return nil, nil +} + func qgaGuestSetPassword(ctx context.Context, userCred mcclient.TokenCredential, sid string, body jsonutils.JSONObject) (interface{}, error) { input := new(hostapi.GuestSetPasswordRequest) if err := body.Unmarshal(input); err != nil { diff --git a/pkg/hostman/guestman/guestman.go b/pkg/hostman/guestman/guestman.go index 1305ce3f29..9c9e8882c5 100644 --- a/pkg/hostman/guestman/guestman.go +++ b/pkg/hostman/guestman/guestman.go @@ -17,12 +17,14 @@ package guestman import ( "bytes" "context" + "encoding/json" "fmt" "io/ioutil" "os" "path" "path/filepath" "runtime/debug" + "strconv" "strings" "sync" "time" @@ -53,6 +55,7 @@ import ( modules "yunion.io/x/onecloud/pkg/mcclient/modules/compute" "yunion.io/x/onecloud/pkg/util/cgrouputils" "yunion.io/x/onecloud/pkg/util/cgrouputils/cpuset" + "yunion.io/x/onecloud/pkg/util/fileutils2" "yunion.io/x/onecloud/pkg/util/netutils2" "yunion.io/x/onecloud/pkg/util/procutils" "yunion.io/x/onecloud/pkg/util/timeutils2" @@ -84,6 +87,9 @@ type SGuestManager struct { ServersLock *sync.Mutex portsInUse *sync.Map + // guests nics traffics lock + TrafficLock *sync.Mutex + GuestStartWorker *appsrv.SWorkerManager isLoaded bool @@ -108,6 +114,7 @@ func NewGuestManager(host hostutils.IHost, serversPath string) *SGuestManager { manager.CandidateServers = make(map[string]*SKVMGuestInstance, 0) manager.UnknownServers = new(sync.Map) manager.ServersLock = &sync.Mutex{} + manager.TrafficLock = &sync.Mutex{} manager.GuestStartWorker = appsrv.NewWorkerManager("GuestStart", 1, appsrv.DEFAULT_BACKLOG, false) manager.cpuSet = NewGuestCpuSetCounter(host.GetHostTopology(), host.GetReservedCpusInfo()) // manager.StartCpusetBalancer() @@ -1438,6 +1445,135 @@ func (m *SGuestManager) RequestVerifyDirtyServer(s *SKVMGuestInstance) { } } +func (m *SGuestManager) ResetGuestNicTrafficLimit(guestId string, input *compute.ServerNicTrafficLimit) error { + guest, ok := m.GetServer(guestId) + if !ok { + return httperrors.NewNotFoundError("guest %s not found", guestId) + } + var nic *desc.SGuestNetwork + for i := range guest.Desc.Nics { + if guest.Desc.Nics[i].Mac == input.Mac { + nic = guest.Desc.Nics[i] + break + } + } + if nic == nil { + return httperrors.NewNotFoundError("guest nic %s not found", input.Mac) + } + m.TrafficLock.Lock() + defer m.TrafficLock.Unlock() + + recordPath := guest.NicTrafficRecordPath() + if fileutils2.Exists(recordPath) { + record, err := m.GetGuestTrafficRecord(guest.Id) + if err != nil { + return errors.Wrap(err, "failed load guest traffic record") + } + if nicRecord, ok := record[strconv.Itoa(int(nic.Index))]; ok { + if nicRecord.TxTraffic >= nic.TxTrafficLimit || nicRecord.RxTraffic >= nic.RxTrafficLimit { + err = guest.SetNicUp(nic) + if err != nil { + return errors.Wrap(err, "set nic up") + } + } + } + delete(record, strconv.Itoa(int(nic.Index))) + if err = m.SaveGuestTrafficRecord(guestId, record); err != nil { + return errors.Wrap(err, "failed save guest traffic record") + } + } + if input.RxTrafficLimit != nil { + nic.RxTrafficLimit = *input.RxTrafficLimit + } + if input.TxTrafficLimit != nil { + nic.TxTrafficLimit = *input.TxTrafficLimit + } + if err := guest.SaveLiveDesc(guest.Desc); err != nil { + return errors.Wrap(err, "guest save desc") + } + return nil +} + +func (m *SGuestManager) SetGuestNicTrafficLimit(guestId string, input *compute.ServerNicTrafficLimit) error { + guest, ok := m.GetServer(guestId) + if !ok { + return httperrors.NewNotFoundError("guest %s not found", guestId) + } + var nic *desc.SGuestNetwork + for i := range guest.Desc.Nics { + if guest.Desc.Nics[i].Mac == input.Mac { + nic = guest.Desc.Nics[i] + break + } + } + if nic == nil { + return httperrors.NewNotFoundError("guest nic %s not found", input.Mac) + } + m.TrafficLock.Lock() + defer m.TrafficLock.Unlock() + if input.RxTrafficLimit != nil { + nic.RxTrafficLimit = *input.RxTrafficLimit + } + if input.TxTrafficLimit != nil { + nic.TxTrafficLimit = *input.TxTrafficLimit + } + if err := guest.SaveLiveDesc(guest.Desc); err != nil { + return errors.Wrap(err, "guest save desc") + } + recordPath := guest.NicTrafficRecordPath() + if fileutils2.Exists(recordPath) { + record, err := m.GetGuestTrafficRecord(guest.Id) + if err != nil { + return errors.Wrap(err, "failed load guest traffic record") + } + if nicRecord, ok := record[strconv.Itoa(int(nic.Index))]; ok { + if nicRecord.TxTraffic < nic.TxTrafficLimit && nicRecord.RxTraffic < nic.RxTrafficLimit { + err = guest.SetNicUp(nic) + if err != nil { + return errors.Wrap(err, "set nic up") + } + } + } + return m.SaveGuestTrafficRecord(guestId, record) + } + return nil +} + +func (m *SGuestManager) SaveGuestTrafficRecord(sid string, record map[string]compute.SNicTrafficRecord) error { + guest, _ := m.GetServer(sid) + recordPath := guest.NicTrafficRecordPath() + v, _ := json.Marshal(record) + return fileutils2.FilePutContents(recordPath, string(v), false) +} + +func (m *SGuestManager) GetGuestTrafficRecord(sid string) (map[string]compute.SNicTrafficRecord, error) { + guest, _ := m.GetServer(sid) + recordPath := guest.NicTrafficRecordPath() + if !fileutils2.Exists(recordPath) { + return nil, nil + } + recordStr, err := ioutil.ReadFile(recordPath) + if err != nil { + return nil, errors.Wrapf(err, "read traffic record %s", recordPath) + } + record := make(map[string]compute.SNicTrafficRecord) + err = json.Unmarshal(recordStr, &record) + if err != nil { + return nil, errors.Wrapf(err, "failed unmarshal traffic record %s", recordPath) + } + return record, nil +} + +func SyncGuestNicsTraffics(guestNicsTraffics map[string]map[string]compute.SNicTrafficRecord) { + session := hostutils.GetComputeSession(context.Background()) + hostId := guestManager.host.GetHostId() + data := jsonutils.Marshal(guestNicsTraffics) + _, err := modules.Hosts.PerformAction(session, hostId, "sync-guest-nic-traffics", data) + if err != nil { + log.Errorf("failed sync-guest-nic-traffics %s", err) + } +} + var guestManager *SGuestManager func Stop() { diff --git a/pkg/hostman/guestman/qemu-kvm.go b/pkg/hostman/guestman/qemu-kvm.go index a243515d4a..b52aa37076 100644 --- a/pkg/hostman/guestman/qemu-kvm.go +++ b/pkg/hostman/guestman/qemu-kvm.go @@ -955,6 +955,10 @@ func (s *SKVMGuestInstance) QgaPath() string { return path.Join(s.HomeDir(), "qga.sock") } +func (s *SKVMGuestInstance) NicTrafficRecordPath() string { + return path.Join(s.HomeDir(), "nic_traffic.json") +} + func (s *SKVMGuestInstance) InitQga() error { guestAgent, err := qga.NewQemuGuestAgent(s.Id, s.QgaPath()) if err != nil { diff --git a/pkg/hostman/guestman/qemu-kvmhelper.go b/pkg/hostman/guestman/qemu-kvmhelper.go index 1e1dde8e2d..3d3992d605 100644 --- a/pkg/hostman/guestman/qemu-kvmhelper.go +++ b/pkg/hostman/guestman/qemu-kvmhelper.go @@ -378,6 +378,11 @@ func (s *SKVMGuestInstance) generateStartScript(data *jsonutils.JSONDict) (strin downscript := s.getNicDownScriptPath(nic) cmd += fmt.Sprintf("%s %s\n", downscript, nic.Ifname) } + traffic, err := guestManager.GetGuestTrafficRecord(s.Id) + if err != nil { + return "", errors.Wrap(err, "get guest traffic record") + } + input.NicTraffics = traffic if input.HugepagesEnabled { cmd += fmt.Sprintf("mkdir -p /dev/hugepages/%s\n", s.Desc.Uuid) @@ -776,6 +781,34 @@ func (s *SKVMGuestInstance) WriteMigrateCerts(certs map[string]string) error { return nil } +func (s *SKVMGuestInstance) SetNicDown(index int8) error { + var nic *desc.SGuestNetwork + for i := range s.Desc.Nics { + if s.Desc.Nics[i].Index == index { + nic = s.Desc.Nics[i] + break + } + } + if nic == nil { + return errors.Errorf("guest %s has no nic index %d", s.GetName(), index) + } + scriptPath := s.getNicDownScriptPath(nic) + out, err := procutils.NewRemoteCommandAsFarAsPossible("bash", scriptPath).Output() + if err != nil { + return errors.Wrapf(err, "failed run nic down script %s", out) + } + return nil +} + +func (s *SKVMGuestInstance) SetNicUp(nic *desc.SGuestNetwork) error { + scriptPath := s.getNicUpScriptPath(nic) + out, err := procutils.NewRemoteCommandAsFarAsPossible("bash", scriptPath).Output() + if err != nil { + return errors.Wrapf(err, "failed run nic down script %s", out) + } + return nil +} + func (s *SKVMGuestInstance) startMemCleaner() error { err := procutils.NewRemoteCommandAsFarAsPossible( options.HostOptions.BinaryMemcleanPath, diff --git a/pkg/hostman/guestman/qemu/generate.go b/pkg/hostman/guestman/qemu/generate.go index 8da0e687cf..0bb1e27c98 100644 --- a/pkg/hostman/guestman/qemu/generate.go +++ b/pkg/hostman/guestman/qemu/generate.go @@ -16,6 +16,7 @@ package qemu import ( "fmt" + "strconv" "strings" "yunion.io/x/pkg/errors" @@ -412,12 +413,25 @@ func generateNicOptions(drvOpt QemuOptions, input *GenerateStartOptionsInput) ([ opts := make([]string, 0) nics := input.GuestDesc.Nics + //input.guest for idx := range nics { if nics[idx].Driver == api.NETWORK_DRIVER_VFIO { continue } + var nicTrafficExceed = false + if input.NicTraffics != nil { + nicTraffic, ok := input.NicTraffics[strconv.Itoa(int(nics[idx].Index))] + if ok { + if nics[idx].TxTrafficLimit > 0 && nicTraffic.TxTraffic > nics[idx].TxTrafficLimit { + nicTrafficExceed = true + } + if nics[idx].RxTrafficLimit > 0 && nicTraffic.RxTraffic > nics[idx].RxTrafficLimit { + nicTrafficExceed = true + } + } + } - netDevOpt, err := getNicNetdevOption(drvOpt, nics[idx], input.IsKVMSupport) + netDevOpt, err := getNicNetdevOption(drvOpt, nics[idx], input.IsKVMSupport, nicTrafficExceed) if err != nil { return nil, errors.Wrapf(err, "getNicNetdevOption %v", nics[idx]) } @@ -430,7 +444,7 @@ func generateNicOptions(drvOpt QemuOptions, input *GenerateStartOptionsInput) ([ return opts, nil } -func getNicNetdevOption(drvOpt QemuOptions, nic *desc.SGuestNetwork, isKVMSupport bool) (string, error) { +func getNicNetdevOption(drvOpt QemuOptions, nic *desc.SGuestNetwork, isKVMSupport bool, nicTrafficExceed bool) (string, error) { if nic.Ifname == "" { return "", errors.Error("ifname is empty") } @@ -450,7 +464,9 @@ func getNicNetdevOption(drvOpt QemuOptions, nic *desc.SGuestNetwork, isKVMSuppor opt += fmt.Sprintf(",queues=%d", nic.NumQueues) } } - opt += fmt.Sprintf(",script=%s", nic.UpscriptPath) + if !nicTrafficExceed { + opt += fmt.Sprintf(",script=%s", nic.UpscriptPath) + } opt += fmt.Sprintf(",downscript=%s", nic.DownscriptPath) return opt, nil } @@ -579,6 +595,7 @@ type GenerateStartOptionsInput struct { GuestDesc *desc.SGuestDesc IsKVMSupport bool + NicTraffics map[string]api.SNicTrafficRecord EnableUUID bool OsName string diff --git a/pkg/hostman/hostmetrics/hostmetrics.go b/pkg/hostman/hostmetrics/hostmetrics.go index 19c6fd2ef1..0de6596357 100644 --- a/pkg/hostman/hostmetrics/hostmetrics.go +++ b/pkg/hostman/hostmetrics/hostmetrics.go @@ -31,6 +31,7 @@ import ( "yunion.io/x/pkg/util/httputils" "yunion.io/x/pkg/util/netutils" + "yunion.io/x/onecloud/pkg/apis/compute" "yunion.io/x/onecloud/pkg/hostman/guestman" "yunion.io/x/onecloud/pkg/hostman/guestman/desc" "yunion.io/x/onecloud/pkg/hostman/hostinfo/hostconsts" @@ -219,12 +220,90 @@ func (s *SGuestMonitorCollector) CollectReportData() (ret string) { reportData[gm.Id] = s.collectGmReport(gm, prevUsage) s.prevPids[gm.Id] = gm.Pid } + s.saveNicTraffics(reportData, gms) s.prevReportData = reportData ret = s.toTelegrafReportData(reportData) return } +func (s *SGuestMonitorCollector) saveNicTraffics(reportData map[string]*GuestMetrics, gms map[string]*SGuestMonitor) { + guestman.GetGuestManager().TrafficLock.Lock() + defer guestman.GetGuestManager().TrafficLock.Unlock() + var guestNicsTraffics = make(map[string]map[string]compute.SNicTrafficRecord) + for guestId, data := range reportData { + gm := gms[guestId] + guestTrafficRecord, err := guestman.GetGuestManager().GetGuestTrafficRecord(gm.Id) + if err != nil { + log.Errorf("failed get guest traffic record %s", err) + continue + } + guestTraffics := make(map[string]compute.SNicTrafficRecord) + for i := range gm.Nics { + if gm.Nics[i].RxTrafficLimit <= 0 && gm.Nics[i].TxTrafficLimit <= 0 { + continue + } + + var nicIo *NetIOMetric + for j := range data.VmNetio { + if gm.Nics[i].Index == int8(data.VmNetio[j].Meta.Index) { + nicIo = data.VmNetio[j] + break + } + } + if nicIo == nil { + log.Warningf("failed found report data for nic %s", gm.Nics[i].Ifname) + continue + } + nicTraffic := compute.SNicTrafficRecord{} + for index, record := range guestTrafficRecord { + if index == strconv.Itoa(nicIo.Meta.Index) { + nicTraffic.RxTraffic += record.RxTraffic + nicTraffic.TxTraffic += record.TxTraffic + } + } + + if gm.Nics[i].RxTrafficLimit > 0 || gm.Nics[i].TxTrafficLimit > 0 { + var nicDown, nicHasBeenSetDown = false, false + if nicTraffic.RxTraffic >= gm.Nics[i].RxTrafficLimit || nicTraffic.RxTraffic >= gm.Nics[i].TxTrafficLimit { + // record traffic excced, nic must has been set down + nicHasBeenSetDown = true + } + if gm.Nics[i].RxTrafficLimit > 0 { + nicTraffic.RxTraffic += int64(nicIo.TimeDiff * nicIo.BPSRecv / 8) + if nicTraffic.RxTraffic >= gm.Nics[i].RxTrafficLimit { + // nic down + nicDown = true + } + } + if gm.Nics[i].TxTrafficLimit > 0 { + nicTraffic.TxTraffic += int64(nicIo.TimeDiff * nicIo.BPSSent / 8) + if nicTraffic.TxTraffic >= gm.Nics[i].TxTrafficLimit { + // nic down + nicDown = true + } + } + guestTraffics[strconv.Itoa(nicIo.Meta.Index)] = nicTraffic + if !nicHasBeenSetDown && nicDown { + log.Infof("guest %s nic %d traffic exceed, set nic down", gm.Id, nicIo.Meta.Index) + gm.SetNicDown(nicIo.Meta.Index) + } + } + } + if len(guestTraffics) == 0 { + continue + } + guestNicsTraffics[gm.Id] = guestTraffics + if err = guestman.GetGuestManager().SaveGuestTrafficRecord(gm.Id, guestTraffics); err != nil { + log.Errorf("failed save guest %s traffic record %v", gm.Id, guestTraffics) + continue + } + } + if len(guestNicsTraffics) > 0 { + guestman.SyncGuestNicsTraffics(guestNicsTraffics) + } +} + func (s *SGuestMonitorCollector) toTelegrafReportData(data map[string]*GuestMetrics) string { ret := []string{} for guestId, report := range data { @@ -368,6 +447,7 @@ func (s *SGuestMonitorCollector) reportNetIo(cur, prev *NetIOMetric) { timeOld := prev.Meta.Uptime diffTime := float64(timeCur - timeOld) + cur.TimeDiff = diffTime if diffTime > 0 { if cur.BytesSent < prev.BytesSent { cur.BPSSent = float64(cur.BytesSent*8) / diffTime @@ -420,6 +500,16 @@ func NewGuestMonitor(name, id string, pid int, nics []*desc.SGuestNetwork, cpuCo return &SGuestMonitor{name, id, pid, nics, cpuCount, ip, proc, "", "", "", "", ""}, nil } +func (m *SGuestMonitor) SetNicDown(index int) { + guest, ok := guestman.GetGuestManager().GetServer(m.Id) + if !ok { + return + } + if err := guest.SetNicDown(int8(index)); err != nil { + log.Errorf("guest %s SetNicDown failed %s", m.Id, err) + } +} + func (m *SGuestMonitor) UpdateVmName(name string) { m.Name = name } @@ -531,14 +621,14 @@ func (m *SGuestMonitor) Netio() []*NetIOMetric { data.Meta.NetId = nic.NetId data.Meta.Uptime, _ = host.Uptime() - data.BytesSent = nicStat.BytesSent - data.BytesRecv = nicStat.BytesRecv - data.PacketsRecv = nicStat.PacketsRecv - data.PacketsSent = nicStat.PacketsSent - data.ErrIn = nicStat.Errin - data.ErrOut = nicStat.Errout - data.DropIn = nicStat.Dropin - data.DropOut = nicStat.Dropout + data.BytesSent = nicStat.BytesRecv + data.BytesRecv = nicStat.BytesSent + data.PacketsRecv = nicStat.PacketsSent + data.PacketsSent = nicStat.PacketsRecv + data.ErrIn = nicStat.Errout + data.ErrOut = nicStat.Errin + data.DropIn = nicStat.Dropout + data.DropOut = nicStat.Dropin res = append(res, data) } return res @@ -561,6 +651,8 @@ type NetIOMetric struct { BPSSent float64 `json:"bps_sent"` PPSRecv float64 `json:"pps_recv"` PPSSent float64 `json:"pps_sent"` + + TimeDiff float64 `json:"-"` } func (n *NetIOMetric) ToMap() map[string]interface{} { diff --git a/pkg/mcclient/options/compute/servers.go b/pkg/mcclient/options/compute/servers.go index 92595249b7..2717045eff 100644 --- a/pkg/mcclient/options/compute/servers.go +++ b/pkg/mcclient/options/compute/servers.go @@ -891,6 +891,17 @@ func (o *ServerSetBootIndexOptions) Params() (jsonutils.JSONObject, error) { return options.StructToParams(o) } +type ServerNicTrafficLimitOptions struct { + ServerIdOptions + MAC string `help:"guest network mac address"` + RxTrafficLimit *int64 `help:" rx traffic limit, unit Byte"` + TxTrafficLimit *int64 `help:" tx traffic limit, unit Byte"` +} + +func (o *ServerNicTrafficLimitOptions) Params() (jsonutils.JSONObject, error) { + return options.StructToParams(o) +} + type ServerSaveImageOptions struct { ServerIdOptions IMAGE string `help:"Image name" json:"name"` diff --git a/pkg/util/logclient/consts.go b/pkg/util/logclient/consts.go index 791a3bb28b..f8ecfbdf68 100644 --- a/pkg/util/logclient/consts.go +++ b/pkg/util/logclient/consts.go @@ -258,4 +258,6 @@ const ( ACT_PROGRESS = "progress" ACT_ADD_BASTION_SERVER = "add_bastion_server" + + ACT_SYNC_TRAFFIC_LIMIT = "sync_traffic_limit" )