diff --git a/cmd/climc/shell/dynamicschedtags.go b/cmd/climc/shell/dynamicschedtags.go new file mode 100644 index 0000000000..6f5ddabe47 --- /dev/null +++ b/cmd/climc/shell/dynamicschedtags.go @@ -0,0 +1,123 @@ +package shell + +import ( + "yunion.io/x/jsonutils" + + "yunion.io/x/onecloud/pkg/mcclient" + "yunion.io/x/onecloud/pkg/mcclient/modules" + "yunion.io/x/onecloud/pkg/mcclient/options" +) + +func init() { + type DynamicschedtagListOptions struct { + options.BaseListOptions + } + R(&DynamicschedtagListOptions{}, "dynamic-schedtag-list", "List dynamic schedtag conditions", func(s *mcclient.ClientSession, args *DynamicschedtagListOptions) error { + var params *jsonutils.JSONDict + { + var err error + params, err = args.BaseListOptions.Params() + if err != nil { + return err + } + } + results, err := modules.Dynamicschedtags.List(s, params) + if err != nil { + return err + } + printList(results, modules.Dynamicschedtags.GetColumns(s)) + return nil + }) + + type DynamicschedtagCreateOptions struct { + NAME string `help:"name of the dynamic schedtag"` + SCHEDTAG string `help:"ID or name of schedtag"` + CONDITION string `help:"condition that assign schedtag to hosts"` + Enable bool `help:"create the policy with enabled status"` + Disable bool `help:"create the policy with disabled status"` + } + R(&DynamicschedtagCreateOptions{}, "dynamic-schedtag-create", "create dynamic schedtag", func(s *mcclient.ClientSession, args *DynamicschedtagCreateOptions) error { + params := jsonutils.NewDict() + params.Add(jsonutils.NewString(args.NAME), "name") + params.Add(jsonutils.NewString(args.CONDITION), "condition") + params.Add(jsonutils.NewString(args.SCHEDTAG), "schedtag") + + if args.Enable { + params.Add(jsonutils.JSONTrue, "enabled") + } else if args.Disable { + params.Add(jsonutils.JSONFalse, "enabled") + } + + result, err := modules.Dynamicschedtags.Create(s, params) + if err != nil { + return err + } + printObject(result) + return nil + }) + + type DynamicschedtagUpdateOptions struct { + ID string `help:"ID or name of the dynamic schedtag"` + Name string `help:"new name of the dynamic schedtag"` + SchedTag string `help:"ID or name of schedtag"` + Condition string `help:"condition that assign schedtag to hosts"` + Enable bool `help:"update to enabled"` + Disable bool `help:"update to disabled"` + } + R(&DynamicschedtagUpdateOptions{}, "dynamic-schedtag-update", "update dynamic schedtag", func(s *mcclient.ClientSession, args *DynamicschedtagUpdateOptions) error { + params := jsonutils.NewDict() + if len(args.Name) > 0 { + params.Add(jsonutils.NewString(args.Name), "name") + } + if len(args.Condition) > 0 { + params.Add(jsonutils.NewString(args.Condition), "condition") + } + if len(args.SchedTag) > 0 { + params.Add(jsonutils.NewString(args.SchedTag), "schedtag") + } + if args.Enable { + params.Add(jsonutils.JSONTrue, "enabled") + } else if args.Disable { + params.Add(jsonutils.JSONFalse, "enabled") + } + + if params.Size() == 0 { + return InvalidUpdateError() + } + result, err := modules.Dynamicschedtags.Update(s, args.ID, params) + if err != nil { + return err + } + printObject(result) + return nil + }) + + type DynamicschedtagDeleteOptions struct { + ID string `help:"ID or name of the dynamic schedtag"` + } + R(&DynamicschedtagDeleteOptions{}, "dynamic-schedtag-delete", "delete dynamic schedtag", func(s *mcclient.ClientSession, args *DynamicschedtagDeleteOptions) error { + result, err := modules.Dynamicschedtags.Delete(s, args.ID, nil) + if err != nil { + return err + } + printObject(result) + return nil + }) + + type DynamicSchedtagEvaluateOptions struct { + ID string `help:"ID or name of the sched policy"` + HOST string `help:"ID or name of the host"` + SERVER string `help:"ID or name of the server"` + } + R(&DynamicSchedtagEvaluateOptions{}, "dynamic-schedtag-evaluate", "Evaluate dynamic schedtag condition", func(s *mcclient.ClientSession, args *DynamicSchedtagEvaluateOptions) error { + params := jsonutils.NewDict() + params.Add(jsonutils.NewString(args.HOST), "host") + params.Add(jsonutils.NewString(args.SERVER), "server") + result, err := modules.Dynamicschedtags.PerformAction(s, args.ID, "evaluate", params) + if err != nil { + return err + } + printObject(result) + return nil + }) +} diff --git a/cmd/climc/shell/schedpolicies.go b/cmd/climc/shell/schedpolicies.go new file mode 100644 index 0000000000..fac45935b7 --- /dev/null +++ b/cmd/climc/shell/schedpolicies.go @@ -0,0 +1,128 @@ +package shell + +import ( + "yunion.io/x/jsonutils" + + "yunion.io/x/onecloud/pkg/mcclient" + "yunion.io/x/onecloud/pkg/mcclient/modules" + "yunion.io/x/onecloud/pkg/mcclient/options" +) + +func init() { + type SchedpoliciesListOptions struct { + options.BaseListOptions + } + R(&SchedpoliciesListOptions{}, "sched-policy-list", "List scheduler policies", func(s *mcclient.ClientSession, args *SchedpoliciesListOptions) error { + var params *jsonutils.JSONDict + { + var err error + params, err = args.BaseListOptions.Params() + if err != nil { + return err + } + } + results, err := modules.Schedpolicies.List(s, params) + if err != nil { + return err + } + printList(results, modules.Schedpolicies.GetColumns(s)) + return nil + }) + + type SchedpoliciesCreateOptions struct { + NAME string `help:"name of the sched policy"` + STRATEGY string `help:"strategy for the schedtag" choices:"require|prefer|avoid|exclude"` + SCHEDTAG string `help:"ID or name of schedtag"` + CONDITION string `help:"condition that assign schedtag to hosts"` + Enable bool `help:"create the policy with enabled status"` + Disable bool `help:"create the policy with disabled status"` + } + R(&SchedpoliciesCreateOptions{}, "sched-policy-create", "create a sched policty", func(s *mcclient.ClientSession, args *SchedpoliciesCreateOptions) error { + params := jsonutils.NewDict() + params.Add(jsonutils.NewString(args.NAME), "name") + params.Add(jsonutils.NewString(args.STRATEGY), "strategy") + params.Add(jsonutils.NewString(args.CONDITION), "condition") + params.Add(jsonutils.NewString(args.SCHEDTAG), "schedtag") + + if args.Enable { + params.Add(jsonutils.JSONTrue, "enabled") + } else if args.Disable { + params.Add(jsonutils.JSONFalse, "disabled") + } + + result, err := modules.Schedpolicies.Create(s, params) + if err != nil { + return err + } + printObject(result) + return nil + }) + + type SchedpoliciesUpdateOptions struct { + ID string `help:"ID or name of the sched policy"` + Name string `help:"new name of sched policy"` + Strategy string `help:"schedtag strategy" choices:"require|prefer|avoid|exclude"` + SchedTag string `help:"ID or name of schedtag"` + Condition string `help:"condition that assign schedtag to hosts"` + Enable bool `help:"make the sched policy enabled"` + Disable bool `help:"make the sched policy disabled"` + } + R(&SchedpoliciesUpdateOptions{}, "sched-policy-update", "update a sched policy", func(s *mcclient.ClientSession, args *SchedpoliciesUpdateOptions) error { + params := jsonutils.NewDict() + + if len(args.Name) > 0 { + params.Add(jsonutils.NewString(args.Name), "name") + } + if len(args.Strategy) > 0 { + params.Add(jsonutils.NewString(args.Strategy), "strategy") + } + if len(args.Condition) > 0 { + params.Add(jsonutils.NewString(args.Condition), "condition") + } + if len(args.SchedTag) > 0 { + params.Add(jsonutils.NewString(args.SchedTag), "schedtag") + } + if args.Enable { + params.Add(jsonutils.JSONTrue, "enabled") + } else if args.Disable { + params.Add(jsonutils.JSONFalse, "enabled") + } + + if params.Size() == 0 { + return InvalidUpdateError() + } + result, err := modules.Schedpolicies.Update(s, args.ID, params) + if err != nil { + return err + } + printObject(result) + return nil + }) + + type SchedpoliciesDeleteOptions struct { + ID string `help:"ID or name of the sched policy"` + } + R(&SchedpoliciesDeleteOptions{}, "sched-policy-delete", "delete a sched policy", func(s *mcclient.ClientSession, args *SchedpoliciesDeleteOptions) error { + result, err := modules.Schedpolicies.Delete(s, args.ID, nil) + if err != nil { + return err + } + printObject(result) + return nil + }) + + type SchedpoliciesEvaluateOptions struct { + ID string `help:"ID or name of the sched policy"` + SERVER string `help:"ID or name of the server"` + } + R(&SchedpoliciesEvaluateOptions{}, "sched-policy-evaluate", "Evaluate sched policy", func(s *mcclient.ClientSession, args *SchedpoliciesEvaluateOptions) error { + params := jsonutils.NewDict() + params.Add(jsonutils.NewString(args.SERVER), "server") + result, err := modules.Schedpolicies.PerformAction(s, args.ID, "evaluate", params) + if err != nil { + return err + } + printObject(result) + return nil + }) +} diff --git a/pkg/compute/handlers.go b/pkg/compute/handlers.go index f44a9ac639..5df211ee60 100644 --- a/pkg/compute/handlers.go +++ b/pkg/compute/handlers.go @@ -72,6 +72,9 @@ func InitHandlers(app *appsrv.Application) { models.LoadbalancerCertificateManager, models.LoadbalancerAclManager, models.LoadbalancerAgentManager, + + models.SchedpolicyManager, + models.DynamicschedtagManager, } { db.RegisterModelManager(manager) handler := db.NewModelHandler(manager) diff --git a/pkg/compute/models/dynamicschedtags.go b/pkg/compute/models/dynamicschedtags.go new file mode 100644 index 0000000000..9c3beaa4bb --- /dev/null +++ b/pkg/compute/models/dynamicschedtags.go @@ -0,0 +1,176 @@ +package models + +import ( + "context" + + "yunion.io/x/jsonutils" + "yunion.io/x/log" + + "database/sql" + "yunion.io/x/onecloud/pkg/cloudcommon/db" + "yunion.io/x/onecloud/pkg/httperrors" + "yunion.io/x/onecloud/pkg/mcclient" + "yunion.io/x/onecloud/pkg/util/conditionparser" +) + +type SDynamicschedtagManager struct { + db.SStandaloneResourceBaseManager +} + +var DynamicschedtagManager *SDynamicschedtagManager + +func init() { + DynamicschedtagManager = &SDynamicschedtagManager{SStandaloneResourceBaseManager: db.NewStandaloneResourceBaseManager(SDynamicschedtag{}, "dynamicschedtags_tbl", "dynamicschedtag", "dynamicschedtags")} +} + +// dynamic schedtag is called before scan host candidates, dynamically adding additional schedtag to hosts +// condition examples: +// host.sys_load > 1.5 || host.mem_used_percent > 0.7 => "high_load" +// +type SDynamicschedtag struct { + db.SStandaloneResourceBase + + Condition string `width:"256" charset:"ascii" nullable:"false" list:"user" create:"required"` + SchedtagId string `width:"36" charset:"ascii" nullable:"false" list:"user" create:"required"` + + Enabled bool `nullable:"false" default:"true" create:"optional" list:"user" update:"user"` +} + +func validateDynamicSchedtagInputData(data *jsonutils.JSONDict, create bool) error { + condStr := jsonutils.GetAnyString(data, []string{"condition"}) + if len(condStr) == 0 && create { + return httperrors.NewInputParameterError("empty condition") + } + if len(condStr) > 0 && !conditionparser.IsValid(condStr) { + return httperrors.NewInputParameterError("invalid condition") + } + + schedStr := jsonutils.GetAnyString(data, []string{"schedtag", "schedtag_id"}) + if len(schedStr) == 0 && create { + return httperrors.NewInputParameterError("missing schedtag") + } + if len(schedStr) > 0 { + schedObj, err := SchedtagManager.FetchByIdOrName(nil, schedStr) + if err != nil { + if err == sql.ErrNoRows { + return httperrors.NewResourceNotFoundError("schedtag %s not found", schedStr) + } else { + log.Errorf("fetch schedtag %s fail %s", schedStr, err) + return httperrors.NewGeneralError(err) + } + } + data.Set("schedtag_id", jsonutils.NewString(schedObj.GetId())) + } + + return nil +} + +func (manager *SDynamicschedtagManager) ValidateCreateData(ctx context.Context, userCred mcclient.TokenCredential, ownerProjId string, query jsonutils.JSONObject, data *jsonutils.JSONDict) (*jsonutils.JSONDict, error) { + err := validateDynamicSchedtagInputData(data, true) + if err != nil { + return nil, err + } + return manager.SStandaloneResourceBaseManager.ValidateCreateData(ctx, userCred, ownerProjId, query, data) +} + +func (self *SDynamicschedtag) ValidateUpdateData(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data *jsonutils.JSONDict) (*jsonutils.JSONDict, error) { + err := validateDynamicSchedtagInputData(data, false) + if err != nil { + return nil, err + } + + return self.SStandaloneResourceBase.ValidateUpdateData(ctx, userCred, query, data) +} + +func (self *SDynamicschedtag) getSchedtag() *SSchedtag { + obj, err := SchedtagManager.FetchById(self.SchedtagId) + if err != nil { + log.Errorf("fail to fetch sched tag by id %s", err) + return nil + } + return obj.(*SSchedtag) +} + +func (self *SDynamicschedtag) getMoreColumns(extra *jsonutils.JSONDict) *jsonutils.JSONDict { + schedtag := self.getSchedtag() + if schedtag != nil { + extra.Add(jsonutils.NewString(schedtag.GetName()), "schedtag") + } + return extra +} + +func (self *SDynamicschedtag) GetCustomizeColumns(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject) *jsonutils.JSONDict { + extra := self.SStandaloneResourceBase.GetCustomizeColumns(ctx, userCred, query) + return self.getMoreColumns(extra) +} + +func (self *SDynamicschedtag) GetExtraDetails(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject) *jsonutils.JSONDict { + extra := self.SStandaloneResourceBase.GetExtraDetails(ctx, userCred, query) + return self.getMoreColumns(extra) +} + +func (manager *SDynamicschedtagManager) getAllEnabledDynamicSchedtags() []SDynamicschedtag { + rules := make([]SDynamicschedtag, 0) + + q := DynamicschedtagManager.Query().IsTrue("enabled") + err := db.FetchModelObjects(manager, q, &rules) + if err != nil { + log.Errorf("getAllEnabledDynamicSchedtags fail %s", err) + return nil + } + + return rules +} + +func (self *SDynamicschedtag) AllowPerformEvaluate(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) bool { + return userCred.IsSystemAdmin() +} + +func (self *SDynamicschedtag) PerformEvaluate(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) (jsonutils.JSONObject, error) { + serverStr := jsonutils.GetAnyString(data, []string{"server", "server_id", "guest", "guest_id"}) + serverObj, err := GuestManager.FetchByIdOrName(userCred, serverStr) + if err != nil { + if err == sql.ErrNoRows { + return nil, httperrors.NewResourceNotFoundError("server %s not found", serverStr) + } else { + return nil, httperrors.NewGeneralError(err) + } + } + + server := serverObj.(*SGuest) + srvDesc := server.getSchedDesc() + + hostStr := jsonutils.GetAnyString(data, []string{"host", "host_id"}) + hostObj, err := HostManager.FetchByIdOrName(userCred, hostStr) + if err != nil { + if err == sql.ErrNoRows { + return nil, httperrors.NewResourceNotFoundError("host %s not found", serverStr) + } else { + return nil, httperrors.NewGeneralError(err) + } + } + + host := hostObj.(*SHost) + // TODO: to fill host scheduling information + hostDesc := jsonutils.Marshal(host) + + params := jsonutils.NewDict() + params.Add(srvDesc, "server") + params.Add(hostDesc, "host") + + meet, err := conditionparser.Eval(self.Condition, params) + if err != nil { + return nil, err + } + result := jsonutils.NewDict() + result.Add(srvDesc, "server") + result.Add(hostDesc, "host") + + if meet { + result.Add(jsonutils.JSONTrue, "result") + } else { + result.Add(jsonutils.JSONFalse, "result") + } + return result, nil +} + diff --git a/pkg/compute/models/guests.go b/pkg/compute/models/guests.go index c242d2fb9b..c3c71ae43a 100644 --- a/pkg/compute/models/guests.go +++ b/pkg/compute/models/guests.go @@ -681,11 +681,11 @@ func (manager *SGuestManager) ValidateCreateData(ctx context.Context, userCred m return nil, httperrors.NewInputParameterError("invalid aggregate_strategy") } } - for idx := 0; data.Contains(fmt.Sprintf("srvtag.%d", idx)); idx += 1 { - aggStr, _ := data.GetString(fmt.Sprintf("srvtag.%d", idx)) + for idx := 0; data.Contains(fmt.Sprintf("schedtag.%d", idx)); idx += 1 { + aggStr, _ := data.GetString(fmt.Sprintf("schedtag.%d", idx)) if len(aggStr) > 0 { parts := strings.Split(aggStr, ":") - if len(parts) >= 2 && len(parts) > 0 && len(parts[1]) > 0 { + if len(parts) >= 2 && len(parts[0]) > 0 && len(parts[1]) > 0 { schedtags[parts[0]] = parts[1] } } @@ -4567,3 +4567,35 @@ func (self *SGuest) PerformUserData(ctx context.Context, userCred mcclient.Token } return nil, nil } + +func (self *SGuest) getSchedDesc() jsonutils.JSONObject { + desc := jsonutils.NewDict() + + desc.Add(jsonutils.NewString(self.Id), "id") + desc.Add(jsonutils.NewString(self.Name), "name") + desc.Add(jsonutils.NewInt(int64(self.VmemSize)), "vmem_size") + desc.Add(jsonutils.NewInt(int64(self.VcpuCount)), "vcpu_count") + + gds := self.GetDisks() + if gds != nil { + for i := 0; i < len(gds); i += 1 { + desc.Add(jsonutils.Marshal(gds[i].ToDiskInfo()), fmt.Sprintf("disk.%d", i)) + } + } + + gns := self.GetNetworks() + if gns != nil { + for i := 0; i < len(gns); i += 1 { + desc.Add(jsonutils.NewString(fmt.Sprintf("%s:%s", gns[i].NetworkId, gns[i].IpAddr)), fmt.Sprintf("net.%d", i)) + } + } + + if len(self.HostId) > 0 && regutils.MatchUUID(self.HostId) { + desc.Add(jsonutils.NewString(self.HostId), "host_id") + } + + desc.Add(jsonutils.NewString(self.ProjectId), "owner_tenant_id") + desc.Add(jsonutils.NewString(self.GetHypervisor()), "hypervisor") + + return desc +} diff --git a/pkg/compute/models/schedpolicies.go b/pkg/compute/models/schedpolicies.go new file mode 100644 index 0000000000..67106a5faa --- /dev/null +++ b/pkg/compute/models/schedpolicies.go @@ -0,0 +1,189 @@ +package models + +import ( + "context" + + "yunion.io/x/jsonutils" + + "database/sql" + "yunion.io/x/log" + "yunion.io/x/onecloud/pkg/cloudcommon/db" + "yunion.io/x/onecloud/pkg/httperrors" + "yunion.io/x/onecloud/pkg/mcclient" + "yunion.io/x/onecloud/pkg/util/conditionparser" + "yunion.io/x/pkg/utils" +) + +type SSchedpolicyManager struct { + db.SStandaloneResourceBaseManager + SInfrastructureManager +} + +var SchedpolicyManager *SSchedpolicyManager + +func init() { + SchedpolicyManager = &SSchedpolicyManager{SStandaloneResourceBaseManager: db.NewStandaloneResourceBaseManager(SSchedpolicy{}, "schedpolicies_tbl", "schedpolicy", "schedpolicies")} +} + +// sched policy is called before calling scheduler, add additional preferences for schedtags +type SSchedpolicy struct { + db.SStandaloneResourceBase + SInfrastructure + + Condition string `width:"256" charset:"ascii" nullable:"false" list:"user" create:"required" update:"user"` + SchedtagId string `width:"36" charset:"ascii" nullable:"false" list:"user" create:"required" update:"user"` + Strategy string `width:"32" charset:"ascii" nullable:"false" list:"user" create:"required" update:"user"` + + Enabled bool `nullable:"false" default:"true" create:"optional" list:"user" update:"user"` +} + +func validateSchedpolicyInputData(data *jsonutils.JSONDict, create bool) error { + err := validateDynamicSchedtagInputData(data, create) + if err != nil { + return err + } + + strategyStr := jsonutils.GetAnyString(data, []string{"strategy"}) + if len(strategyStr) == 0 && create { + return httperrors.NewInputParameterError("missing strategy") + } + + if len(strategyStr) > 0 && !utils.IsInStringArray(strategyStr, STRATEGY_LIST) { + return httperrors.NewInputParameterError("invalid strategy %s", strategyStr) + } + + return nil +} + +func (manager *SSchedpolicyManager) ValidateCreateData(ctx context.Context, userCred mcclient.TokenCredential, ownerProjId string, query jsonutils.JSONObject, data *jsonutils.JSONDict) (*jsonutils.JSONDict, error) { + err := validateSchedpolicyInputData(data, true) + if err != nil { + return nil, err + } + + return manager.SStandaloneResourceBaseManager.ValidateCreateData(ctx, userCred, ownerProjId, query, data) +} + +func (self *SSchedpolicy) ValidateUpdateData(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data *jsonutils.JSONDict) (*jsonutils.JSONDict, error) { + err := validateSchedpolicyInputData(data, false) + if err != nil { + return nil, err + } + + return self.SStandaloneResourceBase.ValidateUpdateData(ctx, userCred, query, data) +} + +func (self *SSchedpolicy) getSchedtag() *SSchedtag { + obj, err := SchedtagManager.FetchById(self.SchedtagId) + if err != nil { + log.Errorf("fail to fetch sched tag by id %s", err) + return nil + } + return obj.(*SSchedtag) +} + +func (self *SSchedpolicy) getMoreColumns(extra *jsonutils.JSONDict) *jsonutils.JSONDict { + schedtag := self.getSchedtag() + if schedtag != nil { + extra.Add(jsonutils.NewString(schedtag.GetName()), "schedtag") + } + return extra +} + +func (self *SSchedpolicy) GetCustomizeColumns(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject) *jsonutils.JSONDict { + extra := self.SStandaloneResourceBase.GetCustomizeColumns(ctx, userCred, query) + return self.getMoreColumns(extra) +} + +func (self *SSchedpolicy) GetExtraDetails(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject) *jsonutils.JSONDict { + extra := self.SStandaloneResourceBase.GetExtraDetails(ctx, userCred, query) + return self.getMoreColumns(extra) +} + +func (manager *SSchedpolicyManager) getAllEnabledPolicies() []SSchedpolicy { + policies := make([]SSchedpolicy, 0) + + q := SchedpolicyManager.Query().IsTrue("enabled") + err := db.FetchModelObjects(manager, q, &policies) + if err != nil { + log.Errorf("getAllEnabledPolicies fail %s", err) + return nil + } + + return policies +} + +func (self *SSchedpolicy) AllowPerformEvaluate(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) bool { + return userCred.IsSystemAdmin() +} + +func (self *SSchedpolicy) PerformEvaluate(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) (jsonutils.JSONObject, error) { + serverStr := jsonutils.GetAnyString(data, []string{"server", "server_id", "guest", "guest_id"}) + serverObj, err := GuestManager.FetchByIdOrName(userCred, serverStr) + if err != nil { + if err == sql.ErrNoRows { + return nil, httperrors.NewResourceNotFoundError("server %s not found", serverStr) + } else { + return nil, httperrors.NewGeneralError(err) + } + } + + server := serverObj.(*SGuest) + desc := server.getSchedDesc() + + params := jsonutils.NewDict() + params.Add(desc, "server") + + meet, err := conditionparser.Eval(self.Condition, params) + if err != nil { + return nil, err + } + result := jsonutils.NewDict() + result.Add(desc, "server") + if meet { + result.Add(jsonutils.JSONTrue, "result") + } else { + result.Add(jsonutils.JSONFalse, "result") + } + return result, nil +} + +func ApplySchedPolicies(params *jsonutils.JSONDict) *jsonutils.JSONDict { + policies := SchedpolicyManager.getAllEnabledPolicies() + if policies == nil { + log.Errorf("getAllEnabledPolicies fail") + return params + } + + schedtags := make(map[string]string) + + if params.Contains("aggregate_strategy") { + err := params.Unmarshal(&schedtags, "aggregate_strategy") + if err != nil { + log.Errorf("unmarshall aggregate_strategy fail %s", err) + return params + } + log.Infof("original sched tag %#v", schedtags) + } + + input := jsonutils.NewDict() + input.Add(params, "server") + + for i := 0; i < len(policies); i += 1 { + meet, err := conditionparser.Eval(policies[i].Condition, input) + if err == nil && meet { + st := policies[i].getSchedtag() + if st != nil { + schedtags[st.Name] = policies[i].Strategy + } + } + } + + newSchedtags := jsonutils.Marshal(schedtags) + log.Infof("updated sched tag %s", newSchedtags) + + ret := jsonutils.NewDict() + ret.Add(newSchedtags, "aggregate_strategy") + + return ret +} diff --git a/pkg/compute/models/schedtags.go b/pkg/compute/models/schedtags.go index 79134250e9..bd4721cd23 100644 --- a/pkg/compute/models/schedtags.go +++ b/pkg/compute/models/schedtags.go @@ -30,6 +30,7 @@ var STRATEGY_LIST = []string{STRATEGY_REQUIRE, STRATEGY_EXCLUDE, STRATEGY_PREFER type SSchedtagManager struct { db.SStandaloneResourceBaseManager + SInfrastructureManager } var SchedtagManager *SSchedtagManager @@ -40,6 +41,7 @@ func init() { type SSchedtag struct { db.SStandaloneResourceBase + SInfrastructure DefaultStrategy string `width:"16" charset:"ascii" nullable:"true" default:"" list:"user" update:"admin" create:"admin_optional"` // Column(VARCHAR(16, charset='ascii'), nullable=True, default='') } @@ -49,11 +51,7 @@ func (manager *SSchedtagManager) AllowListItems(ctx context.Context, userCred mc } func (self *SSchedtag) AllowGetDetails(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject) bool { - return userCred.IsSystemAdmin() -} - -func (manager *SSchedtagManager) AllowCreateItem(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) bool { - return userCred.IsSystemAdmin() + return true } func (manager *SSchedtagManager) ValidateSchedtags(userCred mcclient.TokenCredential, schedtags map[string]string) (map[string]string, error) { @@ -113,6 +111,12 @@ func (self *SSchedtag) ValidateDeleteCondition(ctx context.Context) error { if self.GetHostCount() > 0 { return httperrors.NewNotEmptyError("Tag is associated with hosts") } + if self.getDynamicSchedtagCount() > 0 { + return httperrors.NewNotEmptyError("tag has dynamic rules") + } + if self.getSchedPoliciesCount() > 0 { + return httperrors.NewNotEmptyError("tag is associate with sched policies") + } return self.SStandaloneResourceBase.ValidateDeleteCondition(ctx) } @@ -151,8 +155,18 @@ func (self *SSchedtag) GetHostCount() int { return HostschedtagManager.Query().Equals("schedtag_id", self.Id).Count() } +func (self *SSchedtag) getSchedPoliciesCount() int { + return SchedpolicyManager.Query().Equals("schedtag_id", self.Id).Count() +} + +func (self *SSchedtag) getDynamicSchedtagCount() int { + return DynamicschedtagManager.Query().Equals("schedtag_id", self.Id).Count() +} + func (self *SSchedtag) getMoreColumns(extra *jsonutils.JSONDict) *jsonutils.JSONDict { extra.Add(jsonutils.NewInt(int64(self.GetHostCount())), "host_count") + extra.Add(jsonutils.NewInt(int64(self.getDynamicSchedtagCount())), "dynamic_schedtag_count") + extra.Add(jsonutils.NewInt(int64(self.getSchedPoliciesCount())), "schedpolicy_count") return extra } diff --git a/pkg/compute/tasks/guest_batch_create_task.go b/pkg/compute/tasks/guest_batch_create_task.go index 65758a2313..7e18f7b3c6 100644 --- a/pkg/compute/tasks/guest_batch_create_task.go +++ b/pkg/compute/tasks/guest_batch_create_task.go @@ -36,7 +36,10 @@ func (self *GuestBatchCreateTask) OnInit(ctx context.Context, objs []db.IStandal } func (self *GuestBatchCreateTask) startScheduleGuests(ctx context.Context, guests []*models.SGuest) { - self.SetStage("on_guest_schedule_complete", nil) + // log.Infof("%s", self.Params) + schedtags := models.ApplySchedPolicies(self.Params) + + self.SetStage("on_guest_schedule_complete", schedtags) s := auth.GetAdminSession(options.Options.Region, "") results, err := modules.SchedManager.DoSchedule(s, self.Params, len(guests)) diff --git a/pkg/mcclient/modules/mod_dynamicschedtags.go b/pkg/mcclient/modules/mod_dynamicschedtags.go new file mode 100644 index 0000000000..4ef681aa1c --- /dev/null +++ b/pkg/mcclient/modules/mod_dynamicschedtags.go @@ -0,0 +1,14 @@ +package modules + +var ( + Dynamicschedtags ResourceManager +) + +func init() { + Dynamicschedtags = NewComputeManager("dynamicschedtag", "dynamicschedtag", + []string{"ID", "Name", "Description", + "Condition", "Schedtag", "Schedtag_Id", "Enabled"}, + []string{}) + + registerComputeV2(&Dynamicschedtags) +} diff --git a/pkg/mcclient/modules/mod_schedpolicies.go b/pkg/mcclient/modules/mod_schedpolicies.go new file mode 100644 index 0000000000..eda171e212 --- /dev/null +++ b/pkg/mcclient/modules/mod_schedpolicies.go @@ -0,0 +1,14 @@ +package modules + +var ( + Schedpolicies ResourceManager +) + +func init() { + Schedpolicies = NewComputeManager("schedpolicy", "schedpolicies", + []string{"ID", "Name", "Description", + "Condition", "Schedtag", "Schedtag_Id", "Strategy", "Enabled"}, + []string{}) + + registerComputeV2(&Schedpolicies) +} diff --git a/pkg/mcclient/options/servers.go b/pkg/mcclient/options/servers.go index 3d84e591c2..abd4ead59c 100644 --- a/pkg/mcclient/options/servers.go +++ b/pkg/mcclient/options/servers.go @@ -109,7 +109,7 @@ type ServerCreateOptions struct { AutoStart *bool `help:"Auto start server after it is created"` Zone string `help:"Preferred zone where virtual server should be created" json:"prefer_zone"` Host string `help:"Preferred host where virtual server should be created" json:"prefer_host"` - SchedTag []string `help:"Schedule policy, key = aggregate name, value = require|exclude|prefer|avoid" metavar:""` + Schedtag []string `help:"Schedule policy, key = aggregate name, value = require|exclude|prefer|avoid" metavar:""` Deploy []string `help:"Specify deploy files in virtual server file system" json:"-"` Group []string `help:"Group of virtual server"` Project string `help:"'Owner project ID or Name" json:"tenant"` diff --git a/pkg/util/conditionparser/doc.go b/pkg/util/conditionparser/doc.go new file mode 100644 index 0000000000..2378f3cb07 --- /dev/null +++ b/pkg/util/conditionparser/doc.go @@ -0,0 +1 @@ +package conditionparser // import "yunion.io/x/onecloud/pkg/util/conditionparser" diff --git a/pkg/util/conditionparser/parser.go b/pkg/util/conditionparser/parser.go new file mode 100644 index 0000000000..f8a7dba747 --- /dev/null +++ b/pkg/util/conditionparser/parser.go @@ -0,0 +1,574 @@ +package conditionparser + +import ( + "errors" + "go/ast" + "go/parser" + "go/token" + "reflect" + "strconv" + "strings" + + "yunion.io/x/jsonutils" + "yunion.io/x/pkg/utils" + + "fmt" + "yunion.io/x/log" +) + +var ( + ErrInvalidOp = errors.New("invalid operation") + ErrFieldNotFound = errors.New("field not found") + ErrOutOfIndex = errors.New("out of index") +) + +func IsValid(exprStr string) bool { + _, err := parser.ParseExpr(exprStr) + if err != nil { + return false + } + return true +} + +func Eval(exprStr string, input interface{}) (bool, error) { + expr, err := parser.ParseExpr(exprStr) + if err != nil { + log.Errorf("parse expr %s error %s", exprStr, err) + return false, err + } + result, err := eval(expr, input) + if err != nil { + return false, err + } + switch result.(type) { + case bool, *jsonutils.JSONBool: + return getBool(result), nil + case []interface{}, *jsonutils.JSONArray: + arrX := getArray(result) + for i := 0; i < len(arrX); i += 1 { + val := getBool(arrX[i]) + if val { + return true, nil + } + } + } + return false, nil +} + +func eval(expr ast.Expr, input interface{}) (interface{}, error) { + switch expr.(type) { + case *ast.BinaryExpr: + return evalBinary(expr.(*ast.BinaryExpr), input) + case *ast.UnaryExpr: + return evalUnary(expr.(*ast.UnaryExpr), input) + case *ast.SelectorExpr: + return evalSelector(expr.(*ast.SelectorExpr), input) + case *ast.Ident: + return evalIdent(expr.(*ast.Ident), input) + case *ast.ParenExpr: + return evalParen(expr.(*ast.ParenExpr), input) + case *ast.CallExpr: + return evalCall(expr.(*ast.CallExpr), input) + case *ast.IndexExpr: + return evalIndex(expr.(*ast.IndexExpr), input) + case *ast.BasicLit: + return evalBasicLit(expr.(*ast.BasicLit), input) + default: + return nil, ErrInvalidOp + } +} + +func evalBasicLit(expr *ast.BasicLit, input interface{}) (interface{}, error) { + switch expr.Kind { + case token.IDENT: + switch input.(type) { + case *jsonutils.JSONDict: + jsonX := input.(*jsonutils.JSONDict) + return jsonX.Get(expr.Value) + default: + return nil, ErrInvalidOp + } + case token.INT: + return strconv.Atoi(expr.Value) + case token.FLOAT: + return strconv.ParseFloat(expr.Value, 64) + case token.CHAR: + return expr.Value[1 : len(expr.Value)-1][0], nil + case token.STRING: + return expr.Value[1 : len(expr.Value)-1], nil + default: + return nil, ErrInvalidOp + } +} + +func evalIndex(expr *ast.IndexExpr, input interface{}) (interface{}, error) { + X, err := eval(expr.X, input) + if err != nil { + return nil, err + } + indexV, err := eval(expr.Index, input) + if err != nil { + return nil, err + } + var indexI int64 + switch indexV.(type) { + case int, int8, int16, int32, int64, uint, uint8, uint16, uint32, uint64: + indexI = reflect.ValueOf(indexV).Int() + default: + return nil, ErrInvalidOp + } + switch X.(type) { + case []interface{}, *jsonutils.JSONArray: + arrX := getArray(X) + if indexI >= 0 && indexI < int64(len(arrX)) { + return arrX[indexI], nil + } else { + return nil, ErrOutOfIndex + } + case string: + strX := X.(string) + if indexI >= 0 && indexI < int64(len(strX)) { + return strX[indexI], nil + } else { + return nil, ErrOutOfIndex + } + default: + return nil, err + } +} + +func args2Strings(args []interface{}) ([]string, error) { + strs := make([]string, len(args)) + for i := 0; i < len(args); i += 1 { + var ok bool + strs[i], ok = args[i].(string) + if !ok { + return nil, ErrInvalidOp + } + } + return strs, nil +} + +func evalCallInternal(funcV interface{}, args []interface{}) (interface{}, error) { + switch funcV.(type) { + case []interface{}, *jsonutils.JSONArray: + arrFuncV := getArray(funcV) + ret := make([]interface{}, len(arrFuncV)) + for i := 0; i < len(arrFuncV); i += 1 { + reti, err := evalCallInternal(arrFuncV[i], args) + if err != nil { + return nil, err + } + ret[i] = reti + } + return ret, nil + case *funcCaller: + caller := funcV.(*funcCaller) + switch caller.caller.(type) { + case string: + strX := caller.caller.(string) + strArgs, err := args2Strings(args) + if err != nil { + return nil, err + } + switch caller.method { + case "startswith": + if len(strArgs) != 1 { + return nil, ErrInvalidOp + } + return strings.HasPrefix(strX, strArgs[0]), nil + case "endswith": + if len(strArgs) != 1 { + return nil, ErrInvalidOp + } + return strings.HasSuffix(strX, strArgs[0]), nil + case "contains": + if len(strArgs) != 1 { + return nil, ErrInvalidOp + } + return strings.Contains(strX, strArgs[0]), nil + case "in": + return utils.IsInStringArray(strX, strArgs), nil + default: + return nil, ErrInvalidOp + } + default: + return nil, ErrInvalidOp + } + default: + return nil, ErrInvalidOp + } +} + +func evalCall(expr *ast.CallExpr, input interface{}) (interface{}, error) { + funcV, err := eval(expr.Fun, input) + if err != nil { + return nil, err + } + args := make([]interface{}, len(expr.Args)) + for i := 0; i < len(expr.Args); i += 1 { + args[i], err = eval(expr.Args[i], input) + if err != nil { + return nil, err + } + } + return evalCallInternal(funcV, args) +} + +func evalParen(expr *ast.ParenExpr, input interface{}) (interface{}, error) { + return eval(expr.X, input) +} + +func getJSONProperty(json *jsonutils.JSONDict, identStr string) (jsonutils.JSONObject, error) { + if json.Contains(identStr) { + return json.Get(identStr) + } else if json.Contains(fmt.Sprintf("%s.0", identStr)) { + idx := 0 + jsonArray := jsonutils.NewArray() + for { + obj, _ := json.Get(fmt.Sprintf("%s.%d", identStr, idx)) + if obj != nil { + jsonArray.Add(obj) + idx += 1 + } else { + break + } + } + return jsonArray, nil + } else { + return nil, ErrFieldNotFound + } +} + +func evalIdent(expr *ast.Ident, input interface{}) (interface{}, error) { + if expr.Obj == nil || input == nil { + return expr.Name, nil + } else { + switch input.(type) { + case *jsonutils.JSONDict: + json := input.(*jsonutils.JSONDict) + return getJSONProperty(json, expr.Name) + default: + return nil, ErrInvalidOp + } + } +} + +type funcCaller struct { + caller interface{} + method string +} + +func getArray(X interface{}) []interface{} { + switch X.(type) { + case *jsonutils.JSONArray: + arr := X.(*jsonutils.JSONArray) + ret := make([]interface{}, arr.Size()) + for i := 0; i < arr.Size(); i += 1 { + ret[i], _ = arr.GetAt(i) + } + return ret + default: + return X.([]interface{}) + } +} + +func evalSelectorInternal(X interface{}, identStr string) (interface{}, error) { + switch X.(type) { + case *jsonutils.JSONDict: + json := X.(*jsonutils.JSONDict) + return getJSONProperty(json, identStr) + case []interface{}, *jsonutils.JSONArray: + arrX := getArray(X) + ret := make([]interface{}, len(arrX)) + for i := 0; i < len(arrX); i += 1 { + reti, err := evalSelectorInternal(arrX[i], identStr) + if err != nil { + return nil, err + } + ret[i] = reti + } + return ret, nil + case string, *jsonutils.JSONString: + return &funcCaller{caller: getString(X), method: identStr}, nil + default: + return nil, ErrInvalidOp + } +} + +func evalSelector(expr *ast.SelectorExpr, input interface{}) (interface{}, error) { + X, err := eval(expr.X, input) + if err != nil { + return nil, ErrInvalidOp + } + ident, err := evalIdent(expr.Sel, nil) + if err != nil { + return nil, ErrInvalidOp + } + identStr := ident.(string) + + return evalSelectorInternal(X, identStr) +} + +func evalUnaryInternal(X interface{}, op token.Token) (interface{}, error) { + switch X.(type) { + case []interface{}, *jsonutils.JSONArray: + arrX := getArray(X) + ret := make([]interface{}, len(arrX)) + for i := 0; i < len(arrX); i += 1 { + reti, err := evalUnaryInternal(arrX[i], op) + if err != nil { + return nil, err + } + ret[i] = reti + } + return ret, nil + case bool, *jsonutils.JSONBool: + boolX := getBool(X) + switch op { + case token.NOT: + return !boolX, nil + default: + return nil, ErrInvalidOp + } + case int, int8, int16, int32, int64, uint, uint8, uint16, uint32, uint64, *jsonutils.JSONInt: + intX := getInt(X) + switch op { + case token.SUB: + return -intX, nil + default: + return nil, ErrInvalidOp + } + case float32, float64, *jsonutils.JSONFloat: + floatX := getFloat(X) + switch op { + case token.SUB: + return -floatX, nil + default: + return nil, ErrInvalidOp + } + default: + return nil, ErrInvalidOp + } +} + +func evalUnary(expr *ast.UnaryExpr, input interface{}) (interface{}, error) { + X, err := eval(expr.X, input) + if err != nil { + return nil, err + } + op := expr.Op + return evalUnaryInternal(X, op) +} + +func getString(val interface{}) string { + switch val.(type) { + case *jsonutils.JSONString: + jsonVal, _ := val.(*jsonutils.JSONString).GetString() + return jsonVal + default: + return val.(string) + } +} + +func getBool(val interface{}) bool { + switch val.(type) { + case *jsonutils.JSONBool: + jsonVal, _ := val.(*jsonutils.JSONBool).Bool() + return jsonVal + default: + return val.(bool) + } +} + +func getInt(val interface{}) int64 { + switch val.(type) { + case *jsonutils.JSONInt: + jsonVal, _ := val.(*jsonutils.JSONInt).Int() + return jsonVal + default: + return reflect.ValueOf(val).Int() + } +} + +func getFloat(val interface{}) float64 { + switch val.(type) { + case *jsonutils.JSONFloat: + jsonVal, _ := val.(*jsonutils.JSONFloat).Float() + return jsonVal + default: + return reflect.ValueOf(val).Float() + } +} + +func evalBinaryInternal(X, Y interface{}, op token.Token) (interface{}, error) { + switch X.(type) { + case []interface{}, *jsonutils.JSONArray: + arrX := getArray(X) + ret := make([]interface{}, len(arrX)) + for i := 0; i < len(arrX); i += 1 { + reti, err := evalBinaryInternal(arrX[i], Y, op) + if err != nil { + return nil, err + } + ret[i] = reti + } + return ret, nil + } + switch Y.(type) { + case []interface{}, *jsonutils.JSONArray: + arrY := getArray(Y) + ret := make([]interface{}, len(arrY)) + for i := 0; i < len(arrY); i += 1 { + reti, err := evalBinaryInternal(X, arrY[i], op) + if err != nil { + return nil, err + } + ret[i] = reti + } + return ret, nil + } + switch X.(type) { + case bool, *jsonutils.JSONBool: + switch Y.(type) { + case bool, *jsonutils.JSONBool: + boolX := getBool(X) + boolY := getBool(Y) + switch op { + case token.LAND: + return boolX && boolY, nil + case token.LOR: + return boolX || boolY, nil + default: + return nil, ErrInvalidOp + } + default: + return nil, ErrInvalidOp + } + case string, *jsonutils.JSONString: + switch Y.(type) { + case string, *jsonutils.JSONString: + strX := getString(X) + strY := getString(Y) + switch op { + case token.ADD: + return strX + strY, nil + case token.EQL: + return strX == strY, nil + case token.NEQ: + return strX != strY, nil + default: + return nil, ErrInvalidOp + } + default: + return nil, ErrInvalidOp + } + case int, int8, int16, int32, int64, uint, uint8, uint16, uint32, uint64, *jsonutils.JSONInt: + intX := getInt(X) + switch Y.(type) { + case int, int8, int16, int32, int64, uint, uint8, uint16, uint32, uint64, *jsonutils.JSONInt: + intY := getInt(Y) + return evalIntegerOp(intX, intY, op) + case float32, float64, jsonutils.JSONFloat: + floatX := float64(intX) + floatY := getFloat(Y) + return evalFloatOp(floatX, floatY, op) + default: + return nil, ErrInvalidOp + } + case float32, float64, *jsonutils.JSONFloat: + floatX := getFloat(X) + switch Y.(type) { + case int, int8, int16, int32, int64, uint, uint8, uint16, uint32, uint64, *jsonutils.JSONInt: + floatY := float64(getInt(Y)) + return evalFloatOp(floatX, floatY, op) + case float32, float64, *jsonutils.JSONFloat: + floatY := getFloat(Y) + return evalFloatOp(floatX, floatY, op) + default: + return nil, ErrInvalidOp + } + default: + return nil, ErrInvalidOp + } +} + +func evalBinary(bExpr *ast.BinaryExpr, input interface{}) (interface{}, error) { + X, err := eval(bExpr.X, input) + if err != nil { + return nil, err + } + Y, err := eval(bExpr.Y, input) + if err != nil { + return nil, err + } + return evalBinaryInternal(X, Y, bExpr.Op) +} + +func evalIntegerOp(X, Y int64, op token.Token) (interface{}, error) { + switch op { + case token.ADD: + return X + Y, nil + case token.SUB: + return X - Y, nil + case token.MUL: + return X * Y, nil + case token.QUO: + return X / Y, nil + case token.REM: + return X % Y, nil + case token.AND: + return X & Y, nil + case token.OR: + return X | Y, nil + case token.XOR: + return X ^ Y, nil + case token.SHL: + return X << uint64(Y), nil + case token.SHR: + return X >> uint64(Y), nil + case token.AND_NOT: + return X &^ Y, nil + case token.EQL: + return X == Y, nil + case token.LSS: + return X < Y, nil + case token.GTR: + return X > Y, nil + case token.NEQ: + return X != Y, nil + case token.LEQ: + return X <= Y, nil + case token.GEQ: + return X >= Y, nil + default: + return nil, ErrInvalidOp + } +} + +func evalFloatOp(X, Y float64, op token.Token) (interface{}, error) { + switch op { + case token.ADD: + return X + Y, nil + case token.SUB: + return X - Y, nil + case token.MUL: + return X * Y, nil + case token.QUO: + return X / Y, nil + case token.EQL: + return X == Y, nil + case token.LSS: + return X < Y, nil + case token.GTR: + return X > Y, nil + case token.NEQ: + return X != Y, nil + case token.LEQ: + return X <= Y, nil + case token.GEQ: + return X >= Y, nil + default: + return nil, ErrInvalidOp + } +} diff --git a/pkg/util/conditionparser/parser_test.go b/pkg/util/conditionparser/parser_test.go new file mode 100644 index 0000000000..fac612e176 --- /dev/null +++ b/pkg/util/conditionparser/parser_test.go @@ -0,0 +1,147 @@ +package conditionparser + +import ( + "testing" + + "go/parser" + "yunion.io/x/jsonutils" +) + +func TestAst(t *testing.T) { + input := jsonutils.NewDict() + input.Add(jsonutils.NewString("windows"), "server", "os_type") + disk := jsonutils.NewDict() + disk.Add(jsonutils.NewString("ssd"), "medium_type") + disks := jsonutils.NewArray(disk) + input.Add(disks, "server", "disks") + input.Add(disk, "server", "disk.0") + + cases := []struct { + in string + want bool + }{ + {`server.os_type == "windows"`, true}, + {`server.os_type.startswith("window")`, true}, + {`server.disks[0].medium_type == "ssd"`, true}, + {`server.disks[0].medium_type == "hdd"`, false}, + {`server.os_type == "windows" && server.disks[0].medium_type == "ssd"`, true}, + } + + for _, c := range cases { + result, err := Eval(c.in, input) + if err != nil { + t.Errorf("eval expr %s error %s", c.in, err) + return + } + + if result != c.want { + t.Errorf("expect %v got %v", c.want, result) + return + } + } + +} + +func TestIsValid(t *testing.T) { + cases := []struct { + in string + want bool + }{ + {`server.os_type == "windows"`, true}, + {`server.os_type ==`, false}, + {`dfadsfsdf ==`, false}, + {`dfadsfsdf +`, false}, + } + for _, c := range cases { + got := IsValid(c.in) + if got != c.want { + t.Errorf("%s is valid %v got %v", c.in, c.want, got) + } + } +} + +func TestEval2(t *testing.T) { + inputStr := `{"server": {"disable_delete":false, +"disk.0":{"backend":"local","format":"qcow2","image_id":"4b9fa54c-858c-4c2b-8719-27aee120b3cb","image_properties":{"os_arch":"x86_64","os_distribution":"CentOS","os_type":"Linux","os_version":"7.5.1804"},"medium":"hybrid","size":40960}, +"hypervisor":"kvm","keypair_id":"None","name":"testsched","os_type":"Linux","owner_tenant_id":"5d65667d112e47249ae66dbd7bc07030", +"sched_tag.0":"ssd","secgrp_id":"default","vcpu_count":1,"vmem_size":1024}}` + input, err := jsonutils.ParseString(inputStr) + if err != nil { + t.Errorf("fail to parse server json") + return + } + + cases := []struct { + in string + want bool + }{ + {`server.os_type == "Linux"`, true}, + {`server.vmem_size > 2048`, false}, + {`server.hypervisor.in("kvm", "aliyun")`, true}, + {`server.disable_delete`, false}, + {`server.disk[0].backend == "local"`, true}, + {`server.sched_tag[0] != "ssd"`, false}, + {`server.sched_tag[0] == "ssd"`, true}, + // {`server.os_type == "windows" && server.disks[0].medium_type == "ssd"`, true}, + } + + for _, c := range cases { + result, err := Eval(c.in, input) + if err != nil { + t.Errorf("eval expr %s error %s", c.in, err) + return + } + + if result != c.want { + t.Errorf("expect %v got %v", c.want, result) + return + } + } +} + +func TestEval3(t *testing.T) { + exprStr := `server.disk.backend == "local"` + expr, err := parser.ParseExpr(exprStr) + if err != nil { + t.Errorf("parse exprStr fail %s %s", exprStr, err) + return + } + + t.Logf("%s", jsonutils.Marshal(expr)) + + inputStr := `{"server":{ + "disk.0":{"backend": "local", "medium": "hdd"}, + "disk.1":{"backend": "local", "medium": "hdd"}, + "disk.2":{"backend": "rbd", "medium": "ssd"}, + "disk.3":{"backend": "rbd", "medium": "ssd"}} +}` + input, err := jsonutils.ParseString(inputStr) + if err != nil { + t.Errorf("fail to parse server json %s", err) + return + } + + cases := []struct { + in string + want bool + }{ + {`server.disk.backend == "local"`, true}, + {`server.disk.medium == "ssd"`, true}, + {`server.disk.medium == "hdd"`, true}, + {`server.disk.medium == "hybrid"`, false}, + {`server.disk.medium.contains("ssd")`, true}, + } + + for _, c := range cases { + result, err := Eval(c.in, input) + if err != nil { + t.Errorf("eval expr %s error %s", c.in, err) + return + } + + if result != c.want { + t.Errorf("expect %v got %v", c.want, result) + return + } + } +}