diff --git a/cmd/climc/shell/compute/wires.go b/cmd/climc/shell/compute/wires.go index 91648314f8..643f79c9be 100644 --- a/cmd/climc/shell/compute/wires.go +++ b/cmd/climc/shell/compute/wires.go @@ -15,160 +15,21 @@ package compute import ( - "yunion.io/x/jsonutils" - - "yunion.io/x/onecloud/pkg/mcclient" - "yunion.io/x/onecloud/pkg/mcclient/modulebase" + "yunion.io/x/onecloud/cmd/climc/shell" "yunion.io/x/onecloud/pkg/mcclient/modules" "yunion.io/x/onecloud/pkg/mcclient/options" ) func init() { - type WireListOptions struct { - options.BaseListOptions - - Bandwidth *int `help:"List wires by bandwidth"` - - Region string `help:"List wires in region"` - Zone string `help:"list wires in zone" json:"-"` - Vpc string `help:"List wires in vpc"` - Host string `help:"List wires attached to a host"` - } - R(&WireListOptions{}, "wire-list", "List wires", func(s *mcclient.ClientSession, opts *WireListOptions) error { - params, err := options.ListStructToParams(opts) - if err != nil { - return err - } - - var result *modulebase.ListResult - if len(opts.Zone) > 0 { - result, err = modules.Wires.ListInContext(s, params, &modules.Zones, opts.Zone) - } else { - result, err = modules.Wires.List(s, params) - } - if err != nil { - return err - } - printList(result, modules.Wires.GetColumns(s)) - return nil - }) - - type WireUpdateOptions struct { - ID string `help:"ID or Name of zone to update"` - Name string `help:"Name of wire"` - Desc string `metavar:"" help:"Description"` - Bw int64 `help:"Bandwidth in mbps"` - Mtu int64 `help:"mtu in bytes"` - } - R(&WireUpdateOptions{}, "wire-update", "Update wire", func(s *mcclient.ClientSession, args *WireUpdateOptions) error { - params := jsonutils.NewDict() - if len(args.Name) > 0 { - params.Add(jsonutils.NewString(args.Name), "name") - } - if len(args.Desc) > 0 { - params.Add(jsonutils.NewString(args.Desc), "description") - } - if args.Bw > 0 { - params.Add(jsonutils.NewInt(args.Bw), "bandwidth") - } - if args.Mtu > 0 { - params.Add(jsonutils.NewInt(args.Mtu), "mtu") - } - if params.Size() == 0 { - return InvalidUpdateError() - } - result, err := modules.Wires.Update(s, args.ID, params) - if err != nil { - return err - } - printObject(result) - return nil - }) - - type WireCreateOptions struct { - ZONE string `help:"Zone ID or Name"` - Vpc string `help:"VPC ID or Name" default:"default"` - NAME string `help:"Name of wire"` - BW int64 `help:"Bandwidth in mbps"` - Mtu int64 `help:"mtu in bytes"` - Desc string `metavar:"" help:"Description"` - } - R(&WireCreateOptions{}, "wire-create", "Create a wire", func(s *mcclient.ClientSession, args *WireCreateOptions) error { - params := jsonutils.NewDict() - params.Add(jsonutils.NewString(args.NAME), "name") - params.Add(jsonutils.NewInt(args.BW), "bandwidth") - if args.Mtu > 0 { - params.Add(jsonutils.NewInt(args.Mtu), "mtu") - } - if len(args.Vpc) > 0 { - params.Add(jsonutils.NewString(args.Vpc), "vpc") - } - if len(args.Desc) > 0 { - params.Add(jsonutils.NewString(args.Desc), "description") - } - result, err := modules.Wires.CreateInContext(s, params, &modules.Zones, args.ZONE) - if err != nil { - return err - } - printObject(result) - return nil - }) - - type WireShowOptions struct { - ID string `help:"ID or Name of the wire to show"` - } - R(&WireShowOptions{}, "wire-show", "Show wire details", func(s *mcclient.ClientSession, args *WireShowOptions) error { - result, err := modules.Wires.Get(s, args.ID, nil) - if err != nil { - return err - } - printObject(result) - return nil - }) - - R(&WireShowOptions{}, "wire-delete", "Delete wire", func(s *mcclient.ClientSession, args *WireShowOptions) error { - result, err := modules.Wires.Delete(s, args.ID, nil) - if err != nil { - return err - } - printObject(result) - return nil - }) - - type WirePublicOptions struct { - ID string `help:"ID or name of wire" json:"-"` - Scope string `help:"sharing scope" choices:"system|domain"` - SharedDomains []string `help:"share to domains"` - } - R(&WirePublicOptions{}, "wire-public", "Make wire public", func(s *mcclient.ClientSession, args *WirePublicOptions) error { - params := jsonutils.Marshal(args) - result, err := modules.Wires.PerformAction(s, args.ID, "public", params) - if err != nil { - return err - } - printObject(result) - return nil - }) - - type WirePrivateOptions struct { - ID string `help:"ID or name of wire" json:"-"` - } - R(&WirePrivateOptions{}, "wire-private", "Make wire private", func(s *mcclient.ClientSession, args *WirePrivateOptions) error { - params := jsonutils.Marshal(args) - result, err := modules.Wires.PerformAction(s, args.ID, "private", params) - if err != nil { - return err - } - printObject(result) - return nil - }) - - R(&WireShowOptions{}, "wire-change-owner-candidate-domains", "Show candiate domains of a wire for changing owner", func(s *mcclient.ClientSession, args *WireShowOptions) error { - result, err := modules.Wires.GetSpecific(s, args.ID, "change-owner-candidate-domains", nil) - if err != nil { - return err - } - printObject(result) - return nil - }) + cmd := shell.NewResourceCmd(&modules.Wires).WithKeyword("wire") + cmd.List(new(options.WireListOptions)) + cmd.Create(new(options.WireCreateOptions)) + cmd.Update(new(options.WireUpdateOptions)) + cmd.Show(new(options.WireOptions)) + cmd.Delete(new(options.WireOptions)) + cmd.Perform("public", new(options.WirePublicOptions)) + cmd.Perform("private", new(options.WireOptions)) + cmd.Perform("change-owner-candidate-domains", new(options.WireOptions)) + cmd.Perform("merge", new(options.WireMergeOptions)) + cmd.Perform("merge-network", new(options.WireOptions)) } diff --git a/pkg/apis/compute/wire.go b/pkg/apis/compute/wire.go index 4fa9f2dafa..c1fe69d1be 100644 --- a/pkg/apis/compute/wire.go +++ b/pkg/apis/compute/wire.go @@ -90,3 +90,16 @@ type WireListInput struct { Bandwidth *int `json:"bandwidth"` } + +type WireMergeInput struct { + // description: wire id or name to be merged + // required: true + // example: test-wire + Target string `json:"target"` + // description: if merge networks under wire + // required: false + MergeNetwork bool `json:"merge_network"` +} + +type WireMergeNetworkInput struct { +} diff --git a/pkg/apis/compute/wire_const.go b/pkg/apis/compute/wire_const.go new file mode 100644 index 0000000000..4133972f8d --- /dev/null +++ b/pkg/apis/compute/wire_const.go @@ -0,0 +1,21 @@ +// 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 compute + +const ( + WIRE_STATUS_READY = "ready" + WIRE_STATUS_MERGE_NETWORK = "merge_network" + WIRE_STATUS_MERGE_NETWORK_FAILED = "merge_network_failed" +) diff --git a/pkg/cloudcommon/db/statusbase.go b/pkg/cloudcommon/db/statusbase.go index aa6d05bac0..6554fc1d9b 100644 --- a/pkg/cloudcommon/db/statusbase.go +++ b/pkg/cloudcommon/db/statusbase.go @@ -54,6 +54,10 @@ func (model SStatusResourceBase) GetStatus() string { return model.Status } +func StatusBaseSetStatus(model IStatusBaseModel, userCred mcclient.TokenCredential, status string, reason string) error { + return statusBaseSetStatus(model, userCred, status, reason) +} + func statusBaseSetStatus(model IStatusBaseModel, userCred mcclient.TokenCredential, status string, reason string) error { if model.GetStatus() == status { return nil diff --git a/pkg/compute/models/networks.go b/pkg/compute/models/networks.go index ea3c20f687..0b659ecbf1 100644 --- a/pkg/compute/models/networks.go +++ b/pkg/compute/models/networks.go @@ -2277,6 +2277,52 @@ func (self *SNetwork) PerformMerge(ctx context.Context, userCred mcclient.TokenC return nil, err } + startIp, endIp, err := self.CheckInvalidToMerge(ctx, net, nil) + if err != nil { + logclient.AddActionLogWithContext(ctx, self, logclient.ACT_MERGE, err.Error(), userCred, false) + return nil, err + } + + return nil, self.MergeToNetworkAfterCheck(ctx, userCred, net, startIp, endIp) +} + +func (self *SNetwork) MergeToNetworkAfterCheck(ctx context.Context, userCred mcclient.TokenCredential, net *SNetwork, startIp string, endIp string) error { + lockman.LockClass(ctx, NetworkManager, db.GetLockClassKey(NetworkManager, userCred)) + defer lockman.ReleaseClass(ctx, NetworkManager, db.GetLockClassKey(NetworkManager, userCred)) + + _, err := db.Update(net, func() error { + net.GuestIpStart = startIp + net.GuestIpEnd = endIp + return nil + }) + if err != nil { + logclient.AddActionLogWithContext(ctx, self, logclient.ACT_MERGE, err.Error(), userCred, false) + return err + } + + if err := NetworkManager.handleNetworkIdChange(ctx, &networkIdChangeArgs{ + action: logclient.ACT_MERGE, + oldNet: self, + newNet: net, + userCred: userCred, + }); err != nil { + return err + } + + note := map[string]string{"start_ip": startIp, "end_ip": endIp} + db.OpsLog.LogEvent(self, db.ACT_MERGE, note, userCred) + logclient.AddActionLogWithContext(ctx, self, logclient.ACT_MERGE, note, userCred, true) + + if err = self.RealDelete(ctx, userCred); err != nil { + return err + } + note = map[string]string{"network": self.Id} + db.OpsLog.LogEvent(self, db.ACT_DELETE, note, userCred) + logclient.AddActionLogWithContext(ctx, self, logclient.ACT_DELOCATE, note, userCred, true) + return nil +} + +func (self *SNetwork) CheckInvalidToMerge(ctx context.Context, net *SNetwork, allNets []*SNetwork) (string, string, error) { failReason := make([]string, 0) if self.WireId != net.WireId { @@ -2288,11 +2334,13 @@ func (self *SNetwork) PerformMerge(ctx context.Context, userCred mcclient.TokenC if self.VlanId != net.VlanId { failReason = append(failReason, "vlan_id") } + if self.ServerType != net.ServerType { + failReason = append(failReason, "server_type") + } if len(failReason) > 0 { - err = httperrors.NewInputParameterError("Invalid Target Network %s: inconsist %s", input.Target, strings.Join(failReason, ",")) - logclient.AddActionLogWithContext(ctx, self, logclient.ACT_MERGE, err.Error(), userCred, false) - return nil, err + err := httperrors.NewInputParameterError("Invalid Target Network %s: inconsist %s", net.GetId(), strings.Join(failReason, ",")) + return "", "", err } var startIp, endIp string @@ -2301,12 +2349,18 @@ func (self *SNetwork) PerformMerge(ctx context.Context, userCred mcclient.TokenC ipSS, _ := netutils.NewIPV4Addr(self.GuestIpStart) ipSE, _ := netutils.NewIPV4Addr(self.GuestIpEnd) - wireNets := make([]SNetwork, 0) - q := NetworkManager.Query().Equals("wire_id", self.WireId).NotEquals("id", self.Id).NotEquals("id", net.Id) - err = db.FetchModelObjects(NetworkManager, q, &wireNets) - if err != nil && errors.Cause(err) != sql.ErrNoRows { - logclient.AddActionLogWithContext(ctx, self, logclient.ACT_MERGE, err.Error(), userCred, false) - return nil, errors.Wrap(err, "Query nets of same wire") + var wireNets []SNetwork + if allNets == nil { + q := NetworkManager.Query().Equals("wire_id", self.WireId).NotEquals("id", self.Id).NotEquals("id", net.Id) + err := db.FetchModelObjects(NetworkManager, q, &allNets) + if err != nil && errors.Cause(err) != sql.ErrNoRows { + return "", "", errors.Wrap(err, "Query nets of same wire") + } + } else { + wireNets = make([]SNetwork, len(allNets)) + for i := range wireNets { + wireNets[i] = *allNets[i] + } } if ipNE.StepUp() == ipSS || (ipNE.StepUp() < ipSS && !isOverlapNetworks(wireNets, ipNE.StepUp(), ipSS.StepDown())) { @@ -2315,44 +2369,9 @@ func (self *SNetwork) PerformMerge(ctx context.Context, userCred mcclient.TokenC startIp, endIp = self.GuestIpStart, net.GuestIpEnd } else { note := "Incontinuity Network for %s and %s" - logclient.AddActionLogWithContext(ctx, self, logclient.ACT_MERGE, - fmt.Sprintf(note, self.Name, net.Name), userCred, false) - return nil, httperrors.NewBadRequestError(note, self.Name, net.Name) + return "", "", httperrors.NewBadRequestError(note, self.Name, net.Name) } - - lockman.LockClass(ctx, NetworkManager, db.GetLockClassKey(NetworkManager, userCred)) - defer lockman.ReleaseClass(ctx, NetworkManager, db.GetLockClassKey(NetworkManager, userCred)) - - _, err = db.Update(net, func() error { - net.GuestIpStart = startIp - net.GuestIpEnd = endIp - return nil - }) - if err != nil { - logclient.AddActionLogWithContext(ctx, self, logclient.ACT_MERGE, err.Error(), userCred, false) - return nil, err - } - - if err := NetworkManager.handleNetworkIdChange(ctx, &networkIdChangeArgs{ - action: logclient.ACT_MERGE, - oldNet: self, - newNet: net, - userCred: userCred, - }); err != nil { - return nil, err - } - - note := map[string]string{"start_ip": startIp, "end_ip": endIp} - db.OpsLog.LogEvent(self, db.ACT_MERGE, note, userCred) - logclient.AddActionLogWithContext(ctx, self, logclient.ACT_MERGE, note, userCred, true) - - if err = self.RealDelete(ctx, userCred); err != nil { - return nil, err - } - note = map[string]string{"network": self.Id} - db.OpsLog.LogEvent(self, db.ACT_DELETE, note, userCred) - logclient.AddActionLogWithContext(ctx, self, logclient.ACT_DELOCATE, note, userCred, true) - return nil, nil + return startIp, endIp, nil } // 分割IP子网 diff --git a/pkg/compute/models/wire_id_change_handler.go b/pkg/compute/models/wire_id_change_handler.go new file mode 100644 index 0000000000..19e06a9547 --- /dev/null +++ b/pkg/compute/models/wire_id_change_handler.go @@ -0,0 +1,75 @@ +package models + +import ( + "context" + + "yunion.io/x/onecloud/pkg/cloudcommon/db" + "yunion.io/x/pkg/errors" +) + +type wireIdChangeArgs struct { + oldWire *SWire + newWire *SWire +} + +type wireIdChangeHandler interface { + handleWireIdChange(ctx context.Context, args *wireIdChangeArgs) error +} + +func (manager *SHostwireManager) handleWireIdChange(ctx context.Context, args *wireIdChangeArgs) error { + hws := make([]SHostwire, 0, 8) + err := db.FetchModelObjects(manager, manager.Query().Equals("wire_id", args.oldWire.Id), &hws) + if err != nil { + return err + } + for i := range hws { + hw := &hws[i] + _, err := db.Update(hw, func() error { + hw.WireId = args.newWire.Id + return nil + }) + if err != nil { + return errors.Wrapf(err, "unable to update hostwire host %q wire %q", hw.HostId, hw.WireId) + } + + } + return nil +} + +func (manager *SLoadbalancerClusterManager) handleWireIdChange(ctx context.Context, args *wireIdChangeArgs) error { + lcs := make([]SLoadbalancerCluster, 0, 8) + err := db.FetchModelObjects(manager, manager.Query().Equals("wire_id", args.oldWire.Id), &lcs) + if err != nil { + return err + } + for i := range lcs { + lc := &lcs[i] + _, err := db.Update(lc, func() error { + lc.WireId = args.newWire.Id + return nil + }) + if err != nil { + return errors.Wrapf(err, "unable to update loadbalancercluster %q", lc.GetId()) + } + } + return nil +} + +func (manager *SNetworkManager) handleWireIdChange(ctx context.Context, args *wireIdChangeArgs) error { + ns := make([]SNetwork, 0, 8) + err := db.FetchModelObjects(manager, manager.Query().Equals("wire_id", args.oldWire.Id), &ns) + if err != nil { + return err + } + for i := range ns { + n := &ns[i] + _, err := db.Update(n, func() error { + n.WireId = args.newWire.Id + return nil + }) + if err != nil { + return errors.Wrapf(err, "unable to update network %q", n.GetId()) + } + } + return nil +} diff --git a/pkg/compute/models/wires.go b/pkg/compute/models/wires.go index be6d30ee94..f00d3f042e 100644 --- a/pkg/compute/models/wires.go +++ b/pkg/compute/models/wires.go @@ -38,6 +38,7 @@ import ( "yunion.io/x/onecloud/pkg/cloudprovider" "yunion.io/x/onecloud/pkg/httperrors" "yunion.io/x/onecloud/pkg/mcclient" + "yunion.io/x/onecloud/pkg/util/logclient" "yunion.io/x/onecloud/pkg/util/rbacutils" "yunion.io/x/onecloud/pkg/util/stringutils2" ) @@ -45,6 +46,8 @@ import ( type SWireManager struct { db.SInfrasResourceBaseManager db.SExternalizedResourceBaseManager + db.SStatusResourceBaseManager + SManagedResourceBaseManager SVpcResourceBaseManager SZoneResourceBaseManager } @@ -66,7 +69,9 @@ func init() { type SWire struct { db.SInfrasResourceBase db.SExternalizedResourceBase + db.SStatusResourceBase + // SManagedResourceBase SVpcResourceBase `wdith:"36" charset:"ascii" nullable:"false" list:"domain" create:"domain_required" update:""` SZoneResourceBase `width:"36" charset:"ascii" nullable:"true" list:"domain" create:"domain_required" update:""` @@ -139,6 +144,10 @@ func (manager *SWireManager) ValidateCreateData( return input, nil } +func (wire *SWire) SetStatus(userCred mcclient.TokenCredential, status string, reason string) error { + return db.StatusBaseSetStatus(wire, userCred, status, reason) +} + func (wire *SWire) ValidateUpdateData(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, input api.WireUpdateInput) (api.WireUpdateInput, error) { data := jsonutils.Marshal(input).(*jsonutils.JSONDict) keysV := []validators.IValidator{ @@ -704,6 +713,10 @@ func (self *SWire) getNetworkQuery(ownerId mcclient.IIdentityProvider, scope rba return q } +func (self *SWire) GetNetworks(ownerId mcclient.IIdentityProvider, scope rbacutils.TRbacScope) ([]SNetwork, error) { + return self.getNetworks(ownerId, scope) +} + func (self *SWire) getNetworks(ownerId mcclient.IIdentityProvider, scope rbacutils.TRbacScope) ([]SNetwork, error) { q := self.getNetworkQuery(ownerId, scope) nets := make([]SNetwork, 0) @@ -877,17 +890,21 @@ func chooseCandidateNetworksByNetworkType(nets []SNetwork, isExit bool, serverTy func (manager *SWireManager) InitializeData() error { wires := make([]SWire, 0) q := manager.Query() + q.Filter(sqlchemy.OR(sqlchemy.IsEmpty(q.Field("vpc_id")), sqlchemy.IsEmpty(q.Field("status")))) err := db.FetchModelObjects(manager, q, &wires) if err != nil { return err } for _, w := range wires { - if len(w.VpcId) == 0 { - db.Update(&w, func() error { + db.Update(&w, func() error { + if len(w.VpcId) == 0 { w.VpcId = api.DEFAULT_VPC_ID - return nil - }) - } + } + if len(w.Status) == 0 { + w.Status = api.WIRE_STATUS_READY + } + return nil + }) } return nil } @@ -970,6 +987,89 @@ func (manager *SWireManager) GetOnPremiseWireOfIp(ipAddr string) (*SWire, error) } } +func (w *SWire) AllowPerformMergeNetwork(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject) bool { + return w.IsOwner(userCred) || db.IsAdminAllowPerform(userCred, w, "merge-network") +} + +func (w *SWire) PerformMergeNetwork(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, input api.WireMergeNetworkInput) (jsonutils.JSONObject, error) { + return nil, w.StartMergeNetwork(ctx, userCred, "") +} + +func (w *SWire) AllowPerformMerge(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject) bool { + return w.IsOwner(userCred) || db.IsAdminAllowPerform(userCred, w, "merge") +} + +func (w *SWire) PerformMerge(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, input api.WireMergeInput) (ret jsonutils.JSONObject, err error) { + if len(input.Target) == 0 { + return nil, httperrors.NewMissingParameterError("target") + } + defer func() { + if err != nil { + logclient.AddActionLogWithContext(ctx, w, logclient.ACT_MERGE, err.Error(), userCred, false) + } + }() + iw, err := WireManager.FetchByIdOrName(userCred, input.Target) + if err == sql.ErrNoRows { + err = httperrors.NewNotFoundError("Wire %q", input.Target) + return + } + if err != nil { + return + } + + tw := iw.(*SWire) + lockman.LockClass(ctx, WireManager, db.GetLockClassKey(WireManager, userCred)) + defer lockman.ReleaseClass(ctx, WireManager, db.GetLockClassKey(WireManager, userCred)) + + err = WireManager.handleWireIdChange(ctx, &wireIdChangeArgs{ + oldWire: w, + newWire: tw, + }) + if err != nil { + return + } + logclient.AddActionLogWithContext(ctx, w, logclient.ACT_MERGE, "", userCred, true) + if err = w.Delete(ctx, userCred); err != nil { + return nil, err + } + if input.MergeNetwork { + err = w.StartMergeNetwork(ctx, userCred, "") + if err != nil { + return nil, errors.Wrap(err, "unableto StartMergeNetwork") + } + } + return +} + +func (w *SWire) StartMergeNetwork(ctx context.Context, userCred mcclient.TokenCredential, parentTaskId string) error { + task, err := taskman.TaskManager.NewTask(ctx, "NetworksUnderWireMergeTask", w, userCred, nil, parentTaskId, "", nil) + if err != nil { + return err + } + task.ScheduleRun(nil) + return nil +} + +func (wm *SWireManager) handleWireIdChange(ctx context.Context, args *wireIdChangeArgs) error { + handlers := []wireIdChangeHandler{ + HostwireManager, + NetworkManager, + LoadbalancerClusterManager, + } + + errs := []error{} + for _, h := range handlers { + if err := h.handleWireIdChange(ctx, args); err != nil { + errs = append(errs, err) + } + } + if len(errs) > 0 { + err := errors.NewAggregate(errs) + return httperrors.NewGeneralError(err) + } + return nil +} + // 二层网络列表 func (manager *SWireManager) ListItemFilter( ctx context.Context, @@ -1146,6 +1246,7 @@ func (model *SWire) CustomizeCreate(ctx context.Context, userCred mcclient.Token } data.(*jsonutils.JSONDict).Set("public_scope", jsonutils.NewString(model.PublicScope)) } + model.Status = api.WIRE_STATUS_READY return model.SInfrasResourceBase.CustomizeCreate(ctx, userCred, ownerId, query, data) } diff --git a/pkg/compute/tasks/cloud_account_sync_vmware_net.go b/pkg/compute/tasks/cloud_account_sync_vmware_net.go index 392d885afc..dc7273da10 100644 --- a/pkg/compute/tasks/cloud_account_sync_vmware_net.go +++ b/pkg/compute/tasks/cloud_account_sync_vmware_net.go @@ -290,6 +290,7 @@ func (self *CloudAccountSyncVMwareNetworkTask) createWire(ctx context.Context, c wire.Name = wireName wire.DomainId = cloudaccount.GetOwnerId().GetDomainId() wire.Description = desc + wire.Status = api.WIRE_STATUS_READY wire.SetModelManager(models.WireManager, wire) err := models.WireManager.TableSpec().Insert(ctx, wire) if err != nil { diff --git a/pkg/compute/tasks/networks_under_wire_merge_task.go b/pkg/compute/tasks/networks_under_wire_merge_task.go new file mode 100644 index 0000000000..71f7dc9bc3 --- /dev/null +++ b/pkg/compute/tasks/networks_under_wire_merge_task.go @@ -0,0 +1,106 @@ +package tasks + +import ( + "context" + "fmt" + "sort" + + "yunion.io/x/jsonutils" + "yunion.io/x/log" + "yunion.io/x/pkg/util/netutils" + + api "yunion.io/x/onecloud/pkg/apis/compute" + "yunion.io/x/onecloud/pkg/cloudcommon/db" + "yunion.io/x/onecloud/pkg/cloudcommon/db/lockman" + "yunion.io/x/onecloud/pkg/cloudcommon/db/taskman" + "yunion.io/x/onecloud/pkg/compute/models" + "yunion.io/x/onecloud/pkg/util/logclient" + "yunion.io/x/onecloud/pkg/util/rbacutils" +) + +type NetworksUnderWireMergeTask struct { + taskman.STask +} + +func init() { + taskman.RegisterTask(NetworksUnderWireMergeTask{}) +} + +func (self *NetworksUnderWireMergeTask) taskFailed(ctx context.Context, wire *models.SWire, desc string, err error) { + d := jsonutils.NewDict() + d.Set("description", jsonutils.NewString(desc)) + if err != nil { + d.Set("error", jsonutils.NewString(err.Error())) + } + wire.SetStatus(self.UserCred, api.WIRE_STATUS_MERGE_NETWORK_FAILED, d.PrettyString()) + db.OpsLog.LogEvent(wire, db.ACT_MERGE_NETWORK_FAILED, d, self.UserCred) + logclient.AddActionLogWithStartable(self, wire, logclient.ACT_MERGE_NETWORK, d, self.UserCred, false) + self.SetStageFailed(ctx, nil) +} + +func (self *NetworksUnderWireMergeTask) taskSuccess(ctx context.Context, wire *models.SWire, desc string) { + d := jsonutils.NewString(desc) + wire.SetStatus(self.UserCred, api.WIRE_STATUS_READY, "") + db.OpsLog.LogEvent(wire, db.ACT_MERGE_NETWORK, d, self.UserCred) + logclient.AddActionLogWithStartable(self, wire, logclient.ACT_MERGE_NETWORK, d, self.UserCred, true) + self.SetStageComplete(ctx, nil) +} + +type Net struct { + *models.SNetwork + StartIp netutils.IPV4Addr +} + +func (self *NetworksUnderWireMergeTask) OnInit(ctx context.Context, obj db.IStandaloneModel, body jsonutils.JSONObject) { + w := obj.(*models.SWire) + w.SetStatus(self.UserCred, api.WIRE_STATUS_MERGE_NETWORK, "") + + lockman.LockClass(ctx, models.NetworkManager, db.GetLockClassKey(models.NetworkManager, self.UserCred)) + defer lockman.ReleaseClass(ctx, models.NetworkManager, db.GetLockClassKey(models.NetworkManager, self.UserCred)) + networks, err := w.GetNetworks(self.UserCred, rbacutils.ScopeDomain) + if err != nil { + self.taskFailed(ctx, w, "unable to GetNetworks", err) + return + } + if len(networks) <= 1 { + self.taskSuccess(ctx, w, fmt.Sprintf("num of networks under wire is %d", len(networks))) + } + nets := make([]Net, len(networks)) + for i := range nets { + startIp, _ := netutils.NewIPV4Addr(networks[i].GuestIpStart) + nets[i] = Net{ + SNetwork: &networks[i], + StartIp: startIp, + } + } + sort.Slice(nets, func(i, j int) bool { + if nets[i].VlanId == nets[j].VlanId { + return nets[i].StartIp < nets[j].StartIp + } + return nets[i].VlanId < nets[j].VlanId + }) + log.Infof("nets sorted: %s", jsonutils.Marshal(nets)) + for i := 0; i < len(nets)-1; i++ { + if nets[i].VlanId != nets[i+1].VlanId { + continue + } + // preparenets + wireNets := make([]*models.SNetwork, 0, len(nets)-2) + for j := range nets { + if j != i && j != i+1 { + wireNets = append(wireNets, nets[i].SNetwork) + } + } + startIp, endIp, err := nets[i].CheckInvalidToMerge(ctx, nets[i+1].SNetwork, wireNets) + if err != nil { + log.Debugf("unable to merge network %q to %q: %v", nets[i].GetId(), nets[i+1].GetId(), err) + continue + } + err = nets[i].MergeToNetworkAfterCheck(ctx, self.UserCred, nets[i+1].SNetwork, startIp, endIp) + if err != nil { + self.taskFailed(ctx, w, fmt.Sprintf("unable to merge network %q to %q", nets[i].GetId(), nets[i+1].GetId()), err) + return + } + } + self.taskSuccess(ctx, w, "") +} diff --git a/pkg/mcclient/options/wire.go b/pkg/mcclient/options/wire.go new file mode 100644 index 0000000000..2bdba84a90 --- /dev/null +++ b/pkg/mcclient/options/wire.go @@ -0,0 +1,105 @@ +// 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 options + +import "yunion.io/x/jsonutils" + +type WireListOptions struct { + BaseListOptions + + Bandwidth *int `help:"List wires by bandwidth"` + + Region string `help:"List wires in region"` + Zone string `help:"list wires in zone" json:"-"` + Vpc string `help:"List wires in vpc"` + Host string `help:"List wires attached to a host"` +} + +func (wo *WireListOptions) GetContextId() string { + return wo.Vpc +} + +func (wo *WireListOptions) Params() (jsonutils.JSONObject, error) { + return ListStructToParams(wo) +} + +type WireOptions struct { + ID string `help:"Id or Name of wire to update"` +} + +func (wo *WireOptions) GetId() string { + return wo.ID +} + +func (wo *WireOptions) Params() (jsonutils.JSONObject, error) { + return nil, nil +} + +type WireUpdateOptions struct { + WireOptions + SwireupdateOptions +} + +type SwireupdateOptions struct { + Name string `help:"Name of wire" json:"name"` + Desc string `metavar:"" help:"Description" json:"description"` + Bw int64 `help:"Bandwidth in mbps" json:"bandwidth"` + Mtu int64 `help:"mtu in bytes" json:"mtu"` +} + +func (wo *WireUpdateOptions) Params() (jsonutils.JSONObject, error) { + return jsonutils.Marshal(wo.SwireupdateOptions), nil +} + +type WireCreateOptions struct { + ZONE string `help:"Zone ID or Name"` + Vpc string `help:"VPC ID or Name" default:"default"` + NAME string `help:"Name of wire"` + BW int64 `help:"Bandwidth in mbps"` + Mtu int64 `help:"mtu in bytes"` + Desc string `metavar:"" help:"Description"` +} + +func (wo *WireCreateOptions) Params() (jsonutils.JSONObject, error) { + return jsonutils.Marshal(wo), nil +} + +type SwirePublicOptions struct { + Scope string `help:"sharing scope" choices:"system|domain"` + SharedDomains []string `help:"share to domains"` +} + +type WirePublicOptions struct { + WireOptions + SwirePublicOptions +} + +func (wo *WirePublicOptions) Params() (jsonutils.JSONObject, error) { + return jsonutils.Marshal(wo.SwirePublicOptions), nil +} + +type WireMergeOptions struct { + FROM string `help:"ID or name of merge wire from"` + TARGET string `help:"ID or name of merge wire target"` + MergeNetwork bool `help:"whether to merge network under wire"` +} + +func (wo *WireMergeOptions) GetId() string { + return wo.FROM +} + +func (wo *WireMergeOptions) Params() (jsonutils.JSONObject, error) { + return jsonutils.Marshal(wo), nil +}