fix(region,host): support limit guest nic traffic (#17852)

This commit is contained in:
wanyaoqi
2023-08-29 08:40:20 +08:00
committed by GitHub
parent 1b81d19a82
commit 11fe9bf6bc
23 changed files with 700 additions and 58 deletions
+2
View File
@@ -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))
+8 -6
View File
@@ -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"`
+3
View File
@@ -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"
+23 -16
View File
@@ -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
}
+6
View File
@@ -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"`
}
+12
View File
@@ -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 {
+3
View File
@@ -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"
)
+8
View File
@@ -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
}
+24
View File
@@ -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
}
+48
View File
@@ -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) {
+2
View File
@@ -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
+40 -6
View File
@@ -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()
+29 -19
View File
@@ -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,
+33
View File
@@ -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)
}
@@ -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)
}
@@ -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/<sid>/%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 {
+136
View File
@@ -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() {
+4
View File
@@ -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 {
+33
View File
@@ -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,
+20 -3
View File
@@ -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
+100 -8
View File
@@ -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{} {
+11
View File
@@ -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"`
+2
View File
@@ -258,4 +258,6 @@ const (
ACT_PROGRESS = "progress"
ACT_ADD_BASTION_SERVER = "add_bastion_server"
ACT_SYNC_TRAFFIC_LIMIT = "sync_traffic_limit"
)