From 4250b147546bf4c1486af2a0e72d80bb13fb02e6 Mon Sep 17 00:00:00 2001 From: rainzm Date: Sat, 9 Jan 2021 16:47:51 +0800 Subject: [PATCH 1/5] feat(docker): inject playbook and telegraf installation packages --- build/docker/Dockerfile.ansibleserver | 4 ++++ build/docker/Dockerfile.file-repo | 22 ++++++++++++++++++++++ 2 files changed, 26 insertions(+) create mode 100644 build/docker/Dockerfile.file-repo diff --git a/build/docker/Dockerfile.ansibleserver b/build/docker/Dockerfile.ansibleserver index fd7b3d727f..3c9b00af6f 100644 --- a/build/docker/Dockerfile.ansibleserver +++ b/build/docker/Dockerfile.ansibleserver @@ -1,3 +1,7 @@ FROM registry.cn-beijing.aliyuncs.com/yunionio/ansibleserver-base:v1.0.2 +# install playbook and telegraf install pkg +COPY --from=registry.cn-beijing.aliyuncs.com/yunionio/file-repo:v0.1.0 /opt/yunion/playbook /opt/yunion/playbook +COPY --from=registry.cn-beijing.aliyuncs.com/yunionio/file-repo:v0.1.0 /opt/yunion/ansible-install-pkg /opt/yunion/ansible-install-pkg + ADD ./_output/alpine-build/bin/ansibleserver /opt/yunion/bin/ansibleserver diff --git a/build/docker/Dockerfile.file-repo b/build/docker/Dockerfile.file-repo new file mode 100644 index 0000000000..27969a44f5 --- /dev/null +++ b/build/docker/Dockerfile.file-repo @@ -0,0 +1,22 @@ +FROM registry.cn-beijing.aliyuncs.com/yunionio/onecloud-base:v0.2 + +MAINTAINER "Rain Zheng " + +# openssh-client, for ansible ssh connection +# git, ca-certificates, for fetching ansible roles +RUN set -x \ + && apk update \ + && apk add git \ + && rm -rf /var/cache/apk/* + +# install default playbook and install pkg +Run mkdir -p /opt/yunion/ansible-install-pkg +Run wget https://dl.influxdata.com/telegraf/releases/telegraf-1.17.0-1.x86_64.rpm -P /opt/yunion/ansible-install-pkg +Run wget https://dl.influxdata.com/telegraf/releases/telegraf-1.17.0-1.aarch64.rpm -P /opt/yunion/ansible-install-pkg +Run wget https://dl.influxdata.com/telegraf/releases/telegraf-1.17.0_windows_amd64.zip -P /opt/yunion/ansible-install-pkg +Run wget https://dl.influxdata.com/telegraf/releases/telegraf_1.17.0-1_amd64.deb -P /opt/yunion/ansible-install-pkg +Run wget https://dl.influxdata.com/telegraf/releases/telegraf_1.17.0-1_arm64.deb -P /opt/yunion/ansible-install-pkg + +Run mkdir -p /opt/yunion/playbook +Run mkdir /opt/yunion/playbook/monitor-agent +Run git clone https://github.com/yunionio/monitor-agent.git /opt/yunion/playbook/monitor-agent --recurse-submodules From f8f99365a2cfc6c5a56b1a351fd712014e7c2e4f Mon Sep 17 00:00:00 2001 From: rainzm Date: Sat, 9 Jan 2021 16:54:57 +0800 Subject: [PATCH 2/5] feat(ansibleserver): add reference and instance Reference is a reference to a local or remote playbook without certain inventory and its location is determined by playbookPath. Instance is a complete playbook that can be executed directly, it is filled with inventory and config. Instance ensures that the playbook is executed successfully as far as possible without exceeding the maximum number of attempts --- .../ansible/ansibleplaybook_reference.go | 34 +++ cmd/climc/shell/ansible/devtoolcronjob.go | 222 -------------- cmd/climc/shell/ansible/devtooltemplate.go | 154 ---------- .../models/ansibleplaybook_instance.go | 238 +++++++++++++++ .../models/ansibleplaybookreference.go | 130 ++++++++ pkg/ansibleserver/options/options.go | 3 +- pkg/ansibleserver/service/handlers.go | 3 + pkg/ansibleserver/service/service.go | 1 + pkg/apis/ansible/ansible.go | 37 +++ pkg/apis/ansible/const.go | 22 ++ pkg/apis/ansible/zz_generated.model.go | 63 ++++ pkg/apis/ansibleserver/doc.go | 13 + pkg/apis/ansibleserver/zz_generated.model.go | 61 ++++ .../modules/mod_ansibleplaybook_reference.go | 51 ++++ pkg/mcclient/modules/mod_tasks.go | 6 + .../ansible/ansibleplaybook_reference.go | 92 ++++++ pkg/mcclient/options/ansible/doc.go | 15 + pkg/util/ansiblev2/offline_session.go | 81 +++++ pkg/util/ansiblev2/playbook_session.go | 278 ++++++++++++++++++ pkg/util/ansiblev2/session.go | 175 +---------- 20 files changed, 1142 insertions(+), 537 deletions(-) create mode 100644 cmd/climc/shell/ansible/ansibleplaybook_reference.go delete mode 100644 cmd/climc/shell/ansible/devtoolcronjob.go delete mode 100644 cmd/climc/shell/ansible/devtooltemplate.go create mode 100644 pkg/ansibleserver/models/ansibleplaybook_instance.go create mode 100644 pkg/ansibleserver/models/ansibleplaybookreference.go create mode 100644 pkg/apis/ansible/const.go create mode 100644 pkg/apis/ansible/zz_generated.model.go create mode 100644 pkg/apis/ansibleserver/doc.go create mode 100644 pkg/apis/ansibleserver/zz_generated.model.go create mode 100644 pkg/mcclient/modules/mod_ansibleplaybook_reference.go create mode 100644 pkg/mcclient/options/ansible/ansibleplaybook_reference.go create mode 100644 pkg/mcclient/options/ansible/doc.go create mode 100644 pkg/util/ansiblev2/offline_session.go create mode 100644 pkg/util/ansiblev2/playbook_session.go diff --git a/cmd/climc/shell/ansible/ansibleplaybook_reference.go b/cmd/climc/shell/ansible/ansibleplaybook_reference.go new file mode 100644 index 0000000000..a2f7d85b15 --- /dev/null +++ b/cmd/climc/shell/ansible/ansibleplaybook_reference.go @@ -0,0 +1,34 @@ +// 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 ansible + +import ( + "yunion.io/x/onecloud/cmd/climc/shell" + "yunion.io/x/onecloud/pkg/mcclient/modules" + options "yunion.io/x/onecloud/pkg/mcclient/options/ansible" +) + +func init() { + cmd := shell.NewResourceCmd(&modules.AnsiblePlaybookReference).WithKeyword("ansibleplaybook-reference") + cmd.List(new(options.APRListOptions)) + cmd.Show(new(options.APROptions)) + cmd.Perform("run", new(options.APRRunOptions)) + cmd.Perform("stop", new(options.APRStopOptions)) + + cmd1 := shell.NewResourceCmd(&modules.AnsiblePlaybookInstance).WithKeyword("ansibleplaybook-instance") + cmd1.List(new(options.APIListOptions)) + cmd1.Show(new(options.APIOptions)) + cmd1.Perform("run", new(options.APIOptions)) +} diff --git a/cmd/climc/shell/ansible/devtoolcronjob.go b/cmd/climc/shell/ansible/devtoolcronjob.go deleted file mode 100644 index a03fcd4d6f..0000000000 --- a/cmd/climc/shell/ansible/devtoolcronjob.go +++ /dev/null @@ -1,222 +0,0 @@ -// 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 ansible - -import ( - "fmt" - - "yunion.io/x/jsonutils" - "yunion.io/x/log" - - "yunion.io/x/onecloud/pkg/mcclient" - "yunion.io/x/onecloud/pkg/mcclient/modulebase" - "yunion.io/x/onecloud/pkg/mcclient/modules" - "yunion.io/x/onecloud/pkg/mcclient/options" -) - -func paramValidator(param jsonutils.JSONObject) (bool, error) { - // TODO 移到 server 端 - log.Infof("paramValidator: param: %+v", param) - interval, err := param.Int("interval") - if err != nil { - return true, nil - } - day, err := param.Int("day") - if err != nil { - return true, nil - } - if interval == 0 && day == 0 { - return false, fmt.Errorf("interval and day can not be 0 at the same time") - } - return true, nil -} - -func init() { - type CronjobListOptions struct { - options.BaseListOptions - Name string `help:"cloud region ID or Name" json:"-"` - } - - R(&CronjobListOptions{}, "devtoolcronjob-list", "List Devtool Cronjobs", func(s *mcclient.ClientSession, args *CronjobListOptions) error { - params, err := options.ListStructToParams(args) - if err != nil { - return err - } - var result *modulebase.ListResult - result, err = modules.DevToolCronjobs.List(s, params) - printList(result, modules.DevToolCronjobs.GetColumns(s)) - return nil - }) - - type CronjobCreateOptions struct { - NAME string `help:"Ansible Playbook ID or Name" json:"-"` - Day int `help:"Cronjob runs at given day" default:"0"` - Hour int `help:"Cronjob runs at given hour" default:"0"` - Min int `help:"Cronjob runs at given min" default:"0"` - Sec int `help:"Cronjob runs at given sec" default:"0"` - Interval int64 `help:"Cronjob runs at given interval" default:"0"` - Start bool `help:"start job when created" default:"false"` - Enabled bool `help:"Set job status enabled" default:"false"` - } - R( - &CronjobCreateOptions{}, - "devtoolcronjob-create", - "Create a cronjob repo component", - func(s *mcclient.ClientSession, args *CronjobCreateOptions) error { - result, err := modules.AnsiblePlaybooks.Get(s, args.NAME, nil) - if err != nil { - return err - } - ansiblePlaybookName, err := result.GetString("name") - if err != nil { - return err - } - - ansiblePlaybookID, err := result.GetString("id") - if err != nil { - return err - } - - params := jsonutils.NewDict() - params.Add(jsonutils.NewString(ansiblePlaybookName), "name") - - params.Add(jsonutils.NewString(ansiblePlaybookID), "ansible_playbook_id") - - if args.Start { - params.Add(jsonutils.JSONTrue, "start") - - } - if args.Enabled { - params.Add(jsonutils.JSONTrue, "enabled") - } else if args.Interval > 0 { - params.Add(jsonutils.NewInt(int64(args.Interval)), "interval") - } else { - params.Add(jsonutils.NewInt(int64(0)), "interval") - params.Add(jsonutils.NewInt(int64(args.Day)), "day") - params.Add(jsonutils.NewInt(int64(args.Hour)), "hour") - params.Add(jsonutils.NewInt(int64(args.Min)), "min") - params.Add(jsonutils.NewInt(int64(args.Sec)), "sec") - } - ok, err := paramValidator(params) - if err != nil || !ok { - log.Infof("paramValidator error %s", err) - return err - } - cronjob, err := modules.DevToolCronjobs.Create(s, params) - if err != nil { - log.Errorf("modules.DevToolCronjobs.Create error %s", err) - return err - } - printObject(cronjob) - return nil - }, - ) - - type DevToolCronjobShowOptions struct { - ID string `help:"ID or Name of the DevToolCronjob to show"` - } - R(&DevToolCronjobShowOptions{}, "devtoolcronjob-show", "Show cronjob details", func(s *mcclient.ClientSession, args *DevToolCronjobShowOptions) error { - result, err := modules.DevToolCronjobs.Get(s, args.ID, nil) - if err != nil { - return err - } - printObject(result) - return nil - }) - - type DevToolCronjobUpdateOptions struct { - ID string `help:"ID or Name of DevToolCronjob to update"` - Day int `help:"Cronjob runs at given day" json:"-" default:"-1"` - Hour int `help:"Cronjob runs at given hour" json:"-" default:"-1"` - Min int `help:"Cronjob runs at given min" default:"-1"` - Sec int `help:"Cronjob runs at given sec" default:"-1"` - Interval int `help:"Cronjob runs at given interval" default:"-1"` - Start bool `help:"start job when created"` - Stop bool `help:"start job when created"` - Enable bool `help:"Set job status enabled"` - Disable bool `help:"Set job status enabled"` - } - R(&DevToolCronjobUpdateOptions{}, "devtoolcronjob-update", "Update DevToolCronjob", func(s *mcclient.ClientSession, args *DevToolCronjobUpdateOptions) error { - result, err := modules.DevToolCronjobs.Get(s, args.ID, nil) - if err != nil { - return err - } - - params := jsonutils.NewDict() - interval, _ := result.Int("interval") - day, _ := result.Int("day") - params.Add(jsonutils.NewString(args.ID), "id") - params.Add(jsonutils.NewInt(int64(interval)), "interval") - params.Add(jsonutils.NewInt(int64(day)), "day") - - log.Infof("DevToolCronjobUpdateOptions args: %+v", args) - if args.Interval >= 0 { - params.Add(jsonutils.NewInt(int64(args.Interval)), "interval") - if args.Interval > 0 { - params.Add(jsonutils.NewInt(0), "day") - } - } else if args.Day >= 0 { - params.Add(jsonutils.NewInt(int64(args.Day)), "day") - if args.Day > 0 { - params.Add(jsonutils.NewInt(0), "interval") - } - } - if args.Hour >= 0 { - params.Add(jsonutils.NewInt(int64(args.Hour)), "hour") - } - if args.Min >= 0 { - params.Add(jsonutils.NewInt(int64(args.Min)), "min") - } - if args.Sec >= 0 { - params.Add(jsonutils.NewInt(int64(args.Sec)), "sec") - } - - ok, err := paramValidator(params) - if err != nil || !ok { - return err - } - - if args.Start && args.Stop { - return fmt.Errorf("can not set job start and stop at the same time") - } else if args.Start { - params.Add(jsonutils.JSONTrue, "start") - } else if args.Stop { - params.Add(jsonutils.JSONFalse, "start") - } - if args.Enable && args.Disable { - return fmt.Errorf("can not set job enabled and disabled at the same time") - } else if args.Enable { - params.Add(jsonutils.JSONTrue, "enabled") - } else if args.Disable { - params.Add(jsonutils.JSONFalse, "enabled") - } - - result, err = modules.DevToolCronjobs.Update(s, args.ID, params) - if err != nil { - return err - } - printObject(result) - return nil - }) - - R(&DevToolCronjobShowOptions{}, "devtoolcronjob-delete", "Delete DevToolCronjob", func(s *mcclient.ClientSession, args *DevToolCronjobShowOptions) error { - result, err := modules.DevToolCronjobs.Delete(s, args.ID, nil) - if err != nil { - return err - } - printObject(result) - return nil - }) -} diff --git a/cmd/climc/shell/ansible/devtooltemplate.go b/cmd/climc/shell/ansible/devtooltemplate.go deleted file mode 100644 index 92290feac0..0000000000 --- a/cmd/climc/shell/ansible/devtooltemplate.go +++ /dev/null @@ -1,154 +0,0 @@ -// 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 ansible - -import ( - "yunion.io/x/jsonutils" - "yunion.io/x/log" - - "yunion.io/x/onecloud/pkg/mcclient" - "yunion.io/x/onecloud/pkg/mcclient/modulebase" - "yunion.io/x/onecloud/pkg/mcclient/modules" - "yunion.io/x/onecloud/pkg/mcclient/options" -) - -func init() { - - printAnsiblePlaybookObject := func(obj jsonutils.JSONObject) { - dict := obj.(*jsonutils.JSONDict) - pbJson, err := dict.Get("playbook") - if err != nil { - printObject(obj) - return - } - pbStr := pbJson.YAMLString() - dict.Set("playbook", jsonutils.NewString(pbStr)) - printObject(obj) - } - - type TemplateListOptions struct { - options.BaseListOptions - Name string `help:"cloud region ID or Name" json:"-"` - } - - R(&TemplateListOptions{}, "devtooltemplate-list", "List Devtool Templates", func(s *mcclient.ClientSession, args *TemplateListOptions) error { - params, err := options.ListStructToParams(args) - if err != nil { - return err - } - var result *modulebase.ListResult - result, err = modules.DevToolTemplates.List(s, params) - printList(result, modules.DevToolTemplates.GetColumns(s)) - return nil - }) - - R( - &options.DevtoolTemplateCreateOptions{}, - "devtooltemplate-create", - "Create a template repo component", - func(s *mcclient.ClientSession, opts *options.DevtoolTemplateCreateOptions) error { - - params, err := opts.Params() - if err != nil { - return err - } - log.Infof("devtool playbook create opts: %+v", params) - apb, err := modules.DevToolTemplates.Create(s, params) - if err != nil { - return err - } - printAnsiblePlaybookObject(apb) - return nil - }, - ) - - R( - &options.DevtoolTemplateIdOptions{}, - "devtooltemplate-show", - "Show devtool template", - func(s *mcclient.ClientSession, opts *options.DevtoolTemplateIdOptions) error { - apb, err := modules.DevToolTemplates.Get(s, opts.ID, nil) - if err != nil { - return err - } - printAnsiblePlaybookObject(apb) - return nil - }, - ) - - R( - &options.DevtoolTemplateBindingOptions{}, - "devtooltemplate-bind", - "Binding devtool template to a host/vm", - func(s *mcclient.ClientSession, opts *options.DevtoolTemplateBindingOptions) error { - params := jsonutils.NewDict() - params.Set("server_id", jsonutils.NewString(opts.ServerID)) - apb, err := modules.DevToolTemplates.PerformAction(s, opts.ID, "bind", params) - if err != nil { - return err - } - printAnsiblePlaybookObject(apb) - return nil - }, - ) - - R( - &options.DevtoolTemplateBindingOptions{}, - "devtooltemplate-unbind", - "UnBinding devtool template to a host/vm", - func(s *mcclient.ClientSession, opts *options.DevtoolTemplateBindingOptions) error { - params := jsonutils.NewDict() - params.Set("server_id", jsonutils.NewString(opts.ServerID)) - apb, err := modules.DevToolTemplates.PerformAction(s, opts.ID, "unbind", params) - if err != nil { - return err - } - printAnsiblePlaybookObject(apb) - return nil - }, - ) - - R( - &options.DevtoolTemplateIdOptions{}, - "devtooltemplate-delete", - "Delete devtool template", - func(s *mcclient.ClientSession, opts *options.DevtoolTemplateIdOptions) error { - apb, err := modules.DevToolTemplates.Delete(s, opts.ID, nil) - if err != nil { - return err - } - printAnsiblePlaybookObject(apb) - return nil - }, - ) - - R( - &options.DevtoolTemplateUpdateOptions{}, - "devtooltemplate-update", - "Update ansible playbook", - func(s *mcclient.ClientSession, opts *options.DevtoolTemplateUpdateOptions) error { - params, err := opts.Params() - if err != nil { - return err - } - apb, err := modules.DevToolTemplates.Update(s, opts.ID, params) - if err != nil { - return err - } - printAnsiblePlaybookObject(apb) - return nil - }, - ) -} diff --git a/pkg/ansibleserver/models/ansibleplaybook_instance.go b/pkg/ansibleserver/models/ansibleplaybook_instance.go new file mode 100644 index 0000000000..7dce2be710 --- /dev/null +++ b/pkg/ansibleserver/models/ansibleplaybook_instance.go @@ -0,0 +1,238 @@ +// 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 models + +import ( + "context" + "sync" + "time" + + "yunion.io/x/jsonutils" + "yunion.io/x/log" + "yunion.io/x/pkg/errors" + "yunion.io/x/sqlchemy" + + "yunion.io/x/onecloud/pkg/ansibleserver/options" + api "yunion.io/x/onecloud/pkg/apis/ansible" + "yunion.io/x/onecloud/pkg/appctx" + "yunion.io/x/onecloud/pkg/cloudcommon/db" + "yunion.io/x/onecloud/pkg/cloudcommon/workmanager" + "yunion.io/x/onecloud/pkg/mcclient" + "yunion.io/x/onecloud/pkg/mcclient/auth" + "yunion.io/x/onecloud/pkg/mcclient/modules" + "yunion.io/x/onecloud/pkg/util/ansible" + "yunion.io/x/onecloud/pkg/util/ansiblev2" +) + +type SAnsiblePlaybookInstance struct { + db.SStatusStandaloneResourceBase + ReferenceId string `width:"36" nullable:"false" get:"user" list:"user"` + Inventory string `length:"text" nullable:"false" get:"user" list:"user"` + Params jsonutils.JSONObject + Output string `length:"medium" get:"user" list:"user"` + StartTime time.Time `list:"user" get:"user"` + EndTime time.Time `list:"user" get:"user"` +} + +type SAnsiblePlaybookInstanceManager struct { + db.SStatusStandaloneResourceBaseManager + + sessions ansible.SessionManager + sessionsMux *sync.Mutex +} + +var AnsiblePlaybookInstanceManager *SAnsiblePlaybookInstanceManager + +func init() { + AnsiblePlaybookInstanceManager = &SAnsiblePlaybookInstanceManager{ + SStatusStandaloneResourceBaseManager: db.NewStatusStandaloneResourceBaseManager( + SAnsiblePlaybookInstance{}, + "ansibleplaybook_instance_tbl", + "ansibleplaybookinstance", + "ansibleplaybookinstances", + ), + sessions: ansible.SessionManager{}, + } + AnsiblePlaybookInstanceManager.SetVirtualObject(AnsiblePlaybookInstanceManager) +} + +func (aim *SAnsiblePlaybookInstanceManager) ListItemFilter(ctx context.Context, q *sqlchemy.SQuery, userCred mcclient.TokenCredential, input api.AnsiblePlaybookInstanceListInput) (*sqlchemy.SQuery, error) { + if len(input.AnsiblePlayboookReferenceId) > 0 { + q = q.Equals("reference_id", input.AnsiblePlayboookReferenceId) + } + return aim.SStatusStandaloneResourceBaseManager.ListItemFilter(ctx, q, userCred, input.StatusStandaloneResourceListInput) +} + +func (aim *SAnsiblePlaybookInstanceManager) createInstance(ctx context.Context, referenceId string, host api.AnsibleHost, params jsonutils.JSONObject) (*SAnsiblePlaybookInstance, error) { + // build inventory + inv := ansiblev2.NewInventory() + vars := map[string]interface{}{ + "ansible_user": host.User, + "ansible_host": host.IP, + "ansible_port": host.Port, + } + h := ansiblev2.NewHost() + h.Vars = vars + inv.SetHost(host.Name, h) + ai := &SAnsiblePlaybookInstance{ + ReferenceId: referenceId, + Params: params, + Inventory: inv.String(), + } + err := aim.TableSpec().Insert(ctx, ai) + if err != nil { + return nil, errors.Wrapf(err, "unable to create AnsiblePlaybookInstance") + } + ai.SetModelManager(aim, ai) + return ai, nil +} + +func (ai *SAnsiblePlaybookInstance) AllowPerformRun(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject) bool { + return true +} + +func (ai *SAnsiblePlaybookInstance) PerformRun(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, input jsonutils.JSONObject) (jsonutils.JSONObject, error) { + return nil, ai.runPlaybook(ctx, userCred, nil) +} + +func (ai *SAnsiblePlaybookInstance) runPlaybook(ctx context.Context, userCred mcclient.TokenCredential, ar *SAnsiblePlaybookReference) error { + man := AnsiblePlaybookInstanceManager + if man.sessions.Has(ai.Id) { + return errors.Error("playbook is already running") + } + + if ar == nil { + obj, err := AnsiblePlaybookReferenceManager.FetchById(ai.ReferenceId) + if err != nil { + return errors.Wrapf(err, "unable to fetch ansibleplaybook reference %s", ai.ReferenceId) + } + ar = obj.(*SAnsiblePlaybookReference) + } + var ( + privateKey string + err error + ) + if privateKey, err = modules.Sshkeypairs.FetchPrivateKey(ctx, userCred); err != nil { + return err + } + _, err = db.Update(ai, func() error { + ai.StartTime = time.Now() + ai.EndTime = time.Time{} + ai.Output = "" + ai.Status = api.AnsiblePlaybookStatusRunning + return nil + }) + if err != nil { + return errors.Wrap(err, "unable to update ansibleplaybookinstance") + } + + convertJO := func(o jsonutils.JSONObject) map[string]interface{} { + ret := make(map[string]interface{}) + o.Unmarshal(&ret) + return ret + } + // merge configs + params := make(map[string]interface{}) + for k, v := range convertJO(ar.DefaultParams) { + params[k] = v + } + for k, v := range convertJO(ai.Params) { + params[k] = v + } + sess := ansiblev2.NewOfflineSession(). + Inventory(ai.Inventory). + PrivateKey(privateKey). + Configs(params). + PlaybookPath(ar.PlaybookPath). + OutputWriter(&ansiblePlaybookOutputWriter{ai}). + KeepTmpdir(options.Options.KeepTmpdir) + + man.sessions.Add(ai.Id, sess) + + // NOTE host state check? run only on online hosts and running guests, skip others + run := func(ctx context.Context, data interface{}) (jsonutils.JSONObject, error) { + defer func() { + man.sessions.Remove(ai.Id) + }() + runErr := man.sessions.Run(ai.Id) + // TODO: try to close local forwarding? + + _, err := db.Update(ai, func() error { + err := man.sessions.Err(ai.Id) + if err != nil { + ai.Status = api.AnsiblePlaybookStatusCanceled + } else if runErr != nil { + log.Warningf("playbook %s(%s) failed: %v", ai.Name, ai.Id, runErr) + ai.Status = api.AnsiblePlaybookStatusFailed + } else { + ai.Status = api.AnsiblePlaybookStatusSucceeded + } + ai.EndTime = time.Now() + return nil + }) + if err != nil { + log.Errorf("updating ansible playbook failed: %v", err) + } + return nil, runErr + } + PlaybookWorker.DelayTask(ctx, run, nil) + return nil +} + +func (ai *SAnsiblePlaybookInstance) stopPlaybook(ctx context.Context, userCred mcclient.TokenCredential) error { + man := AnsiblePlaybookInstanceManager + if !man.sessions.Has(ai.Id) { + return errors.Error("playbook is not running") + } + // the playbook will be removed from session map in runPlaybook() on return from run + man.sessions.Stop(ai.Id) + return nil +} + +func (ai *SAnsiblePlaybookInstance) getMaxOutputLength() int { + return OutputMaxBytes +} + +func (ai *SAnsiblePlaybookInstance) getOutput() string { + return ai.Output +} + +func (ai *SAnsiblePlaybookInstance) setOutput(s string) { + ai.Output = s +} + +var PlaybookWorker *workmanager.SWorkManager + +func taskFailed(ctx context.Context, reason string) { + if taskId := ctx.Value(appctx.APP_CONTEXT_KEY_TASK_ID); taskId != nil { + session := auth.GetAdminSessionWithInternal(ctx, "", "") + modules.TaskFailed(&modules.DevtoolTasks, session, taskId.(string), reason) + } else { + log.Warningf("Reqeuest task failed missing task id, with reason: %s", reason) + } +} + +func taskCompleted(ctx context.Context, data jsonutils.JSONObject) { + if taskId := ctx.Value(appctx.APP_CONTEXT_KEY_TASK_ID); taskId != nil { + session := auth.GetAdminSessionWithInternal(ctx, "", "") + modules.TaskComplete(&modules.DevtoolTasks, session, taskId.(string), data) + } else { + log.Warningf("Reqeuest task failed missing task id, with data: %v", data) + } +} + +func InitPlaybookWorker() { + PlaybookWorker = workmanager.NewWorkManger(taskFailed, taskCompleted, options.Options.PlaybookWorkerCount) +} diff --git a/pkg/ansibleserver/models/ansibleplaybookreference.go b/pkg/ansibleserver/models/ansibleplaybookreference.go new file mode 100644 index 0000000000..7b1f07caa0 --- /dev/null +++ b/pkg/ansibleserver/models/ansibleplaybookreference.go @@ -0,0 +1,130 @@ +// 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 models + +import ( + "context" + "os" + + "yunion.io/x/jsonutils" + "yunion.io/x/pkg/errors" + + api "yunion.io/x/onecloud/pkg/apis/ansible" + "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/rbacutils" +) + +type SAnsiblePlaybookReference struct { + db.SSharableVirtualResourceBase + + PlaybookPath string `length:"text" nullable:"false" create:"required" get:"user" list:"user"` + Method string `width:"8" nullable:"false" default:"offline" get:"user" list:"user"` + DefaultParams jsonutils.JSONObject `get:"user" list:"user"` +} + +type SAnsiblePlaybookReferenceManager struct { + db.SSharableVirtualResourceBaseManager +} + +var AnsiblePlaybookReferenceManager *SAnsiblePlaybookReferenceManager + +func init() { + AnsiblePlaybookReferenceManager = &SAnsiblePlaybookReferenceManager{ + SSharableVirtualResourceBaseManager: db.NewSharableVirtualResourceBaseManager( + SAnsiblePlaybookReference{}, + "ansibleplaybook_reference_tbl", + "ansibleplaybookreference", + "ansibleplaybookreferences", + ), + } + AnsiblePlaybookReferenceManager.SetVirtualObject(AnsiblePlaybookReferenceManager) +} + +func (arm *SAnsiblePlaybookReferenceManager) ResourceScope() rbacutils.TRbacScope { + return rbacutils.ScopeSystem +} + +var ( + monitorAgent = "monitor agent" + monitorAgentId = "monitoragent" +) + +func (arm *SAnsiblePlaybookReferenceManager) ValidateCreateData(ctx context.Context, userCred mcclient.TokenCredential, ownerId mcclient.IIdentityProvider, query jsonutils.JSONObject, input api.AnsiblePlaybookReferenceCreateInput) (api.AnsiblePlaybookReferenceCreateInput, error) { + + if input.Method != api.APReferenceMethodOffline { + return input, httperrors.NewInputParameterError("unkown Method %q", input.Method) + } + if !arm.checkOfflinePath(input.PlaybookPath) { + return input, httperrors.NewInputParameterError("non-existent path: %q", input.PlaybookPath) + } + return input, nil +} + +func (arm *SAnsiblePlaybookReferenceManager) checkOfflinePath(path string) bool { + _, err := os.Stat(path) + if err != nil { + if os.IsExist(err) { + return true + } + return false + } + return true +} + +func (ar *SAnsiblePlaybookReference) CustomizeCreate(ctx context.Context, userCred mcclient.TokenCredential, ownerId mcclient.IIdentityProvider, query jsonutils.JSONObject, data jsonutils.JSONObject) error { + if data.Contains("playbook_params") { + params, _ := data.Get("playbook_params") + ar.DefaultParams = params + } + ar.Status = api.APReferenceStatusReady + return nil +} + +func (ar *SAnsiblePlaybookReference) ValidateUpdateData(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, input api.AnsiblePlaybookReferenceUpdateInput) (api.AnsiblePlaybookReferenceUpdateInput, error) { + return input, httperrors.NewForbiddenError("prohibited operation") +} + +func (ar *SAnsiblePlaybookReference) AllowPerformRun(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject) bool { + return ar.IsOwner(userCred) || db.IsAdminAllowPerform(userCred, ar, "run") +} + +func (ar *SAnsiblePlaybookReference) PerformRun(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, input api.AnsiblePlaybookReferenceRunInput) (api.AnsiblePlaybookReferenceRunOutput, error) { + output := api.AnsiblePlaybookReferenceRunOutput{} + ai, err := AnsiblePlaybookInstanceManager.createInstance(ctx, ar.Id, input.Host, input.Args) + if err != nil { + return output, errors.Wrap(err, "unable to create instance") + } + output.AnsiblePlaybookInstanceId = ai.Id + err = ai.runPlaybook(ctx, userCred, ar) + if err != nil { + return output, errors.Wrap(err, "unable to runPlaybook") + } + return output, nil +} + +func (ar *SAnsiblePlaybookReference) AllowPerformStop(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject) bool { + return ar.IsOwner(userCred) || db.IsAdminAllowPerform(userCred, ar, "stop") +} + +func (ar *SAnsiblePlaybookReference) PerformStop(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, input api.AnsiblePlaybookReferenceStopInput) (jsonutils.JSONObject, error) { + obj, err := AnsiblePlaybookInstanceManager.FetchById(input.AnsiblePlaybookInstanceId) + if err != nil { + return nil, errors.Wrap(err, "unable to fetch ansibleplaybookinstance") + } + ai := obj.(*SAnsiblePlaybookInstance) + return nil, ai.stopPlaybook(ctx, userCred) +} diff --git a/pkg/ansibleserver/options/options.go b/pkg/ansibleserver/options/options.go index d1caacf9e7..1042995129 100644 --- a/pkg/ansibleserver/options/options.go +++ b/pkg/ansibleserver/options/options.go @@ -19,7 +19,8 @@ import common_options "yunion.io/x/onecloud/pkg/cloudcommon/options" type AnsibleServerOptions struct { common_options.CommonOptions common_options.DBOptions - KeepTmpdir bool `help:"Whether to save the tmp directory" json:"keep_tmpdir"` + KeepTmpdir bool `help:"Whether to save the tmp directory" json:"keep_tmpdir"` + PlaybookWorkerCount int `help:"count of worker to run playbook" default:"5" json:"playbook_worker_count"` } var ( diff --git a/pkg/ansibleserver/service/handlers.go b/pkg/ansibleserver/service/handlers.go index 6e09153ea6..aed35209ec 100644 --- a/pkg/ansibleserver/service/handlers.go +++ b/pkg/ansibleserver/service/handlers.go @@ -42,9 +42,12 @@ func InitHandlers(app *appsrv.Application) { db.RegisterModelManager(db.Metadata) db.RegisterModelManager(db.UserCacheManager) db.RegisterModelManager(db.TenantCacheManager) + db.RegisterModelManager(db.SharedResourceManager) for _, manager := range []db.IModelManager{ models.AnsiblePlaybookManager, models.AnsiblePlaybookV2Manager, + models.AnsiblePlaybookReferenceManager, + models.AnsiblePlaybookInstanceManager, } { db.RegisterModelManager(manager) handler := db.NewModelHandler(manager) diff --git a/pkg/ansibleserver/service/service.go b/pkg/ansibleserver/service/service.go index 5fddf1c196..4c891dcafe 100644 --- a/pkg/ansibleserver/service/service.go +++ b/pkg/ansibleserver/service/service.go @@ -41,6 +41,7 @@ func StartService() { dbOpts := &opts.DBOptions baseOpts := &opts.BaseOptions + models.InitPlaybookWorker() app := common_app.InitApp(baseOpts, false) InitHandlers(app) diff --git a/pkg/apis/ansible/ansible.go b/pkg/apis/ansible/ansible.go index d4090361c4..d1b5a61004 100644 --- a/pkg/apis/ansible/ansible.go +++ b/pkg/apis/ansible/ansible.go @@ -15,6 +15,8 @@ package ansible import ( + "yunion.io/x/jsonutils" + "yunion.io/x/onecloud/pkg/apis" "yunion.io/x/onecloud/pkg/util/ansible" ) @@ -27,3 +29,38 @@ type AnsiblePlaybookCreateInput struct { } type AnsiblePlaybookUpdateInput AnsiblePlaybookCreateInput + +type AnsibleHost struct { + User string `json:"user"` + IP string `json:"ip"` + Port int `json:"port"` + Name string `json:"name"` +} + +type AnsiblePlaybookReferenceCreateInput struct { + apis.SharableVirtualResourceCreateInput + + SAnsiblePlaybookReference + PlaybookParams map[string]interface{} `json:"playbook_params"` +} + +type AnsiblePlaybookReferenceUpdateInput struct { +} + +type AnsiblePlaybookReferenceRunInput struct { + Host AnsibleHost + Args jsonutils.JSONObject +} + +type AnsiblePlaybookReferenceRunOutput struct { + AnsiblePlaybookInstanceId string +} + +type AnsiblePlaybookReferenceStopInput struct { + AnsiblePlaybookInstanceId string +} + +type AnsiblePlaybookInstanceListInput struct { + apis.StatusStandaloneResourceListInput + AnsiblePlayboookReferenceId string +} diff --git a/pkg/apis/ansible/const.go b/pkg/apis/ansible/const.go new file mode 100644 index 0000000000..bf27bda61a --- /dev/null +++ b/pkg/apis/ansible/const.go @@ -0,0 +1,22 @@ +// 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 ansible + +const ( + APReferenceMethodOffline = "offline" + APReferenceMethodOnline = "online" + + APReferenceStatusReady = "ready" +) diff --git a/pkg/apis/ansible/zz_generated.model.go b/pkg/apis/ansible/zz_generated.model.go new file mode 100644 index 0000000000..dbbcc89340 --- /dev/null +++ b/pkg/apis/ansible/zz_generated.model.go @@ -0,0 +1,63 @@ +// 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. + +// Code generated by model-api-gen. DO NOT EDIT. + +package ansible + +import ( + time "time" + + "yunion.io/x/onecloud/pkg/apis" + ansible "yunion.io/x/onecloud/pkg/util/ansible" +) + +// SAnsiblePlaybook is an autogenerated struct via yunion.io/x/onecloud/pkg/ansibleserver/models.SAnsiblePlaybook. +type SAnsiblePlaybook struct { + apis.SVirtualResourceBase + Playbook *ansible.Playbook `json:"playbook"` + Output string `json:"output"` + StartTime time.Time `json:"start_time"` + EndTime time.Time `json:"end_time"` +} + +// SAnsiblePlaybookInstance is an autogenerated struct via yunion.io/x/onecloud/pkg/ansibleserver/models.SAnsiblePlaybookInstance. +type SAnsiblePlaybookInstance struct { + apis.SStatusStandaloneResourceBase + ReferenceId string `json:"reference_id"` + ProxyEndpoingId string `json:"proxy_endpoing_id"` + LocalForwardId string `json:"local_forward_id"` + Proxy string `json:"proxy"` + Output string `json:"output"` + StartTime time.Time `json:"start_time"` + EndTime time.Time `json:"end_time"` +} + +// SAnsiblePlaybookReference is an autogenerated struct via yunion.io/x/onecloud/pkg/ansibleserver/models.SAnsiblePlaybookReference. +type SAnsiblePlaybookReference struct { + apis.SVirtualResourceBase + PlaybookPath string `json:"playbook_path"` + Method string `json:"method"` +} + +// SAnsiblePlaybookV2 is an autogenerated struct via yunion.io/x/onecloud/pkg/ansibleserver/models.SAnsiblePlaybookV2. +type SAnsiblePlaybookV2 struct { + apis.SVirtualResourceBase + Playbook string `json:"playbook"` + Inventory string `json:"inventory"` + Requirements string `json:"requirements"` + Files string `json:"files"` + Output string `json:"output"` + StartTime time.Time `json:"start_time"` + EndTime time.Time `json:"end_time"` + CreatorMark string `json:"creator_mark"` +} diff --git a/pkg/apis/ansibleserver/doc.go b/pkg/apis/ansibleserver/doc.go new file mode 100644 index 0000000000..56843ce6de --- /dev/null +++ b/pkg/apis/ansibleserver/doc.go @@ -0,0 +1,13 @@ +// 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 ansibleserver // import "yunion.io/x/onecloud/pkg/apis/ansibleserver" diff --git a/pkg/apis/ansibleserver/zz_generated.model.go b/pkg/apis/ansibleserver/zz_generated.model.go new file mode 100644 index 0000000000..fe4acb94ed --- /dev/null +++ b/pkg/apis/ansibleserver/zz_generated.model.go @@ -0,0 +1,61 @@ +// 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. + +// Code generated by model-api-gen. DO NOT EDIT. + +package ansibleserver + +import ( + time "time" + + "yunion.io/x/onecloud/pkg/apis" + ansible "yunion.io/x/onecloud/pkg/util/ansible" +) + +// SAnsiblePlaybook is an autogenerated struct via yunion.io/x/onecloud/pkg/ansibleserver/models.SAnsiblePlaybook. +type SAnsiblePlaybook struct { + apis.SVirtualResourceBase + Playbook *ansible.Playbook `json:"playbook"` + Output string `json:"output"` + StartTime time.Time `json:"start_time"` + EndTime time.Time `json:"end_time"` +} + +// SAnsiblePlaybookInstance is an autogenerated struct via yunion.io/x/onecloud/pkg/ansibleserver/models.SAnsiblePlaybookInstance. +type SAnsiblePlaybookInstance struct { + apis.SStatusStandaloneResourceBase + ReferenceId string `json:"reference_id"` + Inventory string `json:"inventory"` + Output string `json:"output"` + StartTime time.Time `json:"start_time"` + EndTime time.Time `json:"end_time"` +} + +// SAnsiblePlaybookReference is an autogenerated struct via yunion.io/x/onecloud/pkg/ansibleserver/models.SAnsiblePlaybookReference. +type SAnsiblePlaybookReference struct { + apis.SSharableVirtualResourceBase + PlaybookPath string `json:"playbook_path"` + Method string `json:"method"` +} + +// SAnsiblePlaybookV2 is an autogenerated struct via yunion.io/x/onecloud/pkg/ansibleserver/models.SAnsiblePlaybookV2. +type SAnsiblePlaybookV2 struct { + apis.SVirtualResourceBase + Playbook string `json:"playbook"` + Inventory string `json:"inventory"` + Requirements string `json:"requirements"` + Files string `json:"files"` + Output string `json:"output"` + StartTime time.Time `json:"start_time"` + EndTime time.Time `json:"end_time"` + CreatorMark string `json:"creator_mark"` +} diff --git a/pkg/mcclient/modules/mod_ansibleplaybook_reference.go b/pkg/mcclient/modules/mod_ansibleplaybook_reference.go new file mode 100644 index 0000000000..eea143d444 --- /dev/null +++ b/pkg/mcclient/modules/mod_ansibleplaybook_reference.go @@ -0,0 +1,51 @@ +// 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 modules + +import "yunion.io/x/onecloud/pkg/mcclient/modulebase" + +var ( + AnsiblePlaybookReference modulebase.ResourceManager + AnsiblePlaybookInstance modulebase.ResourceManager +) + +func init() { + AnsiblePlaybookReference = NewAnsibleManager( + "ansibleplaybookreference", + "ansibleplaybookreferences", + []string{ + "Id", + "Name", + "Playbook_Path", + "Default_Params", + "Method", + }, + []string{}, + ) + AnsiblePlaybookInstance = NewAnsibleManager( + "ansibleplaybookinstance", + "ansibleplaybookinstances", + []string{ + "Id", + "Status", + "Start_Time", + "End_Time", + "Output", + }, + []string{}, + ) + registerV2(&AnsiblePlaybookReference) + registerV2(&AnsiblePlaybookInstance) +} diff --git a/pkg/mcclient/modules/mod_tasks.go b/pkg/mcclient/modules/mod_tasks.go index c86a9cb2a8..8ac45ffa20 100644 --- a/pkg/mcclient/modules/mod_tasks.go +++ b/pkg/mcclient/modules/mod_tasks.go @@ -28,6 +28,7 @@ var ( Tasks modulebase.ResourceManager ComputeTasks ComputeTasksManager + DevtoolTasks modulebase.ResourceManager ) type ComputeTasksManager struct { @@ -45,6 +46,11 @@ func init() { []string{"Id", "Obj_name", "Obj_Id", "Task_name", "Stage", "Created_at"}), } registerCompute(&ComputeTasks) + + DevtoolTasks = NewDevtoolManager("task", "tasks", + []string{}, + []string{"Id", "Obj_name", "Obj_Id", "Task_name", "Stage", "Created_at"}, + ) } type ITaskResourceManager interface { diff --git a/pkg/mcclient/options/ansible/ansibleplaybook_reference.go b/pkg/mcclient/options/ansible/ansibleplaybook_reference.go new file mode 100644 index 0000000000..db62caa916 --- /dev/null +++ b/pkg/mcclient/options/ansible/ansibleplaybook_reference.go @@ -0,0 +1,92 @@ +// 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 ansible + +import ( + "yunion.io/x/jsonutils" + + "yunion.io/x/onecloud/pkg/mcclient/options" +) + +type APRListOptions struct { + options.BaseListOptions +} + +func (ao *APRListOptions) Params() (jsonutils.JSONObject, error) { + return options.ListStructToParams(ao) +} + +type APROptions struct { + ID string `help:"id or name of ansible playbook reference"` +} + +func (ao *APROptions) GetId() string { + return ao.ID +} + +func (ao *APROptions) Params() (jsonutils.JSONObject, error) { + return nil, nil +} + +type aprRunOptions struct { + ServerName string + ServerIp string + ServerUser string + Args map[string]interface{} + ProxyEndpoingId string +} + +type APRRunOptions struct { + APROptions + aprRunOptions +} + +func (ao *APRRunOptions) Params() (jsonutils.JSONObject, error) { + return jsonutils.Marshal(ao.aprRunOptions), nil +} + +type aprStopOptions struct { + AnsiblePlaybookInstanceId string +} + +type APRStopOptions struct { + APROptions + aprStopOptions +} + +func (ao *APRStopOptions) Params() (jsonutils.JSONObject, error) { + return jsonutils.Marshal(ao.aprStopOptions), nil +} + +type APIListOptions struct { + options.BaseListOptions + AnsiblePlayboookReferenceId string +} + +func (ao *APIListOptions) Params() (jsonutils.JSONObject, error) { + return options.ListStructToParams(ao) +} + +type APIOptions struct { + ID string +} + +func (ao *APIOptions) GetId() string { + return ao.ID +} + +func (ao *APIOptions) Params() (jsonutils.JSONObject, error) { + return nil, nil +} diff --git a/pkg/mcclient/options/ansible/doc.go b/pkg/mcclient/options/ansible/doc.go new file mode 100644 index 0000000000..54f6ff9e29 --- /dev/null +++ b/pkg/mcclient/options/ansible/doc.go @@ -0,0 +1,15 @@ +// 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 ansible // import "yunion.io/x/onecloud/pkg/mcclient/options/ansible" diff --git a/pkg/util/ansiblev2/offline_session.go b/pkg/util/ansiblev2/offline_session.go new file mode 100644 index 0000000000..040e7bea18 --- /dev/null +++ b/pkg/util/ansiblev2/offline_session.go @@ -0,0 +1,81 @@ +// 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 ansiblev2 + +import ( + "context" + "io" +) + +type OfflineSession struct { + PlaybookSessionBase + + playbookPath string + proxy string + user string + hostIp string + hostName string + configs map[string]interface{} +} + +func NewOfflineSession() *OfflineSession { + sess := &OfflineSession{ + PlaybookSessionBase: NewPlaybookSessionBase(), + configs: map[string]interface{}{}, + } + return sess +} + +func (sess *OfflineSession) PrivateKey(s string) *OfflineSession { + sess.privateKey = s + return sess +} + +func (sess *OfflineSession) PlaybookPath(s string) *OfflineSession { + sess.playbookPath = s + return sess +} + +func (sess *OfflineSession) Inventory(s string) *OfflineSession { + sess.inventory = s + return sess +} + +func (sess *OfflineSession) OutputWriter(w io.Writer) *OfflineSession { + sess.outputWriter = w + return sess +} + +func (sess *OfflineSession) KeepTmpdir(keep bool) *OfflineSession { + sess.keepTmpdir = keep + return sess +} + +func (sess *OfflineSession) Configs(configs map[string]interface{}) *OfflineSession { + sess.configs = configs + return sess +} + +func (sess *OfflineSession) GetPlaybookPath() string { + return sess.playbookPath +} + +func (sess *OfflineSession) GetConfigs() map[string]interface{} { + return sess.configs +} + +func (sess *OfflineSession) Run(ctx context.Context) (err error) { + return runnable{sess}.Run(ctx) +} diff --git a/pkg/util/ansiblev2/playbook_session.go b/pkg/util/ansiblev2/playbook_session.go new file mode 100644 index 0000000000..d8d2bbd9e2 --- /dev/null +++ b/pkg/util/ansiblev2/playbook_session.go @@ -0,0 +1,278 @@ +// 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 ansiblev2 + +import ( + "context" + "io" + "io/ioutil" + "os" + "os/exec" + "path/filepath" + "sync" + + "github.com/go-yaml/yaml" + + "yunion.io/x/pkg/errors" + yerrors "yunion.io/x/pkg/util/errors" +) + +type IPlaybookSession interface { + GetPrivateKey() string + GetPlaybook() string + GetPlaybookPath() string + GetInventory() string + IsKeepTmpdir() bool + GetConfigs() map[string]interface{} + GetRequirements() string + GetFiles() map[string][]byte + GetOutputWriter() io.Writer + CheckAndSetRunning() bool + SetStopped() +} + +type runnable struct { + IPlaybookSession +} + +func (r runnable) Run(ctx context.Context) (err error) { + var ( + tmpdir string + ) + + has := r.CheckAndSetRunning() + if has { + return errors.Errorf("playbook is already running") + } + defer r.SetStopped() + + // make tmpdir + tmpdir, err = ioutil.TempDir("", "onecloud-ansiblev2") + if err != nil { + err = errors.Wrap(err, "making tmp dir") + return + } + defer func() { + if r.IsKeepTmpdir() { + return + } + if err1 := os.RemoveAll(tmpdir); err1 != nil { + err = errors.Wrapf(err1, "removing %q", tmpdir) + } + }() + + // write out inventory + inventory := filepath.Join(tmpdir, "inventory") + err = ioutil.WriteFile(inventory, []byte(r.GetInventory()), os.FileMode(0600)) + if err != nil { + err = errors.Wrapf(err, "writing inventory %s", inventory) + return + } + + // write out playbook + playbook := r.GetPlaybookPath() + if len(playbook) == 0 { + playbook := filepath.Join(tmpdir, "playbook") + err = ioutil.WriteFile(playbook, []byte(r.GetPlaybook()), os.FileMode(0600)) + if err != nil { + err = errors.Wrapf(err, "writing playbook %s", playbook) + return + } + } + + // write out private key + var privateKey string + if len(r.GetPrivateKey()) > 0 { + privateKey = filepath.Join(tmpdir, "private_key") + err = ioutil.WriteFile(privateKey, []byte(r.GetPrivateKey()), os.FileMode(0600)) + if err != nil { + err = errors.Wrapf(err, "writing private key %s", privateKey) + return + } + } + + // write out requirements + var requirements string + if len(r.GetRequirements()) > 0 { + requirements = filepath.Join(tmpdir, "requirements.yml") + err = ioutil.WriteFile(requirements, []byte(r.GetRequirements()), os.FileMode(0600)) + if err != nil { + err = errors.Wrapf(err, "writing requirements %s", requirements) + return + } + } + + // write out files + for name, content := range r.GetFiles() { + path := filepath.Join(tmpdir, name) + dir := filepath.Dir(path) + err = os.MkdirAll(dir, os.FileMode(0700)) + if err != nil { + err = errors.Wrapf(err, "mkdir -p %s", dir) + return + } + err = ioutil.WriteFile(path, content, os.FileMode(0600)) + if err != nil { + err = errors.Wrapf(err, "writing file %s", name) + return + } + } + + // write out configs + var config string + if r.GetConfigs() != nil { + yml, err := yaml.Marshal(r.GetConfigs()) + if err != nil { + return errors.Wrap(err, "unable to marshal map to yaml") + } + config = filepath.Join(tmpdir, "config") + err = ioutil.WriteFile(config, yml, os.FileMode(0600)) + if err != nil { + return errors.Wrapf(err, "unable to write config to file %s", config) + } + } + + // run modules one by one + var errs []error + defer func() { + if len(errs) > 0 { + err = yerrors.NewAggregate(errs) + } + }() + + // install required roles + if len(requirements) > 0 { + args := []string{ + "install", "-r", requirements, "-p", tmpdir, + } + cmd := exec.CommandContext(ctx, "ansible-galaxy", args...) + stdout, _ := cmd.StdoutPipe() + stderr, _ := cmd.StderrPipe() + if err1 := cmd.Start(); err1 != nil { + errs = append(errs, errors.Wrap(err1, "start ansible-galaxy install roles")) + return + } + // Mix stdout, stderr + + if writer := r.GetOutputWriter(); writer != nil { + go io.Copy(writer, stdout) + go io.Copy(writer, stderr) + } + if err1 := cmd.Wait(); err1 != nil { + errs = append(errs, errors.Wrap(err1, "wait ansible-galaxy install roles")) + } + } + + // run playbook + { + args := []string{ + "--inventory", inventory, + } + if config != "" { + args = append(args, "-e", "@"+config) + } + if privateKey != "" { + args = append(args, "--private-key", privateKey) + } + args = append(args, playbook) + cmd := exec.CommandContext(ctx, "ansible-playbook", args...) + cmd.Dir = tmpdir + cmd.Env = os.Environ() + cmd.Env = append(cmd.Env, "ANSIBLE_HOST_KEY_CHECKING=False") + stdout, _ := cmd.StdoutPipe() + stderr, _ := cmd.StderrPipe() + if err1 := cmd.Start(); err1 != nil { + errs = append(errs, errors.Wrapf(err1, "start playbook %s", playbook)) + return + } + // Mix stdout, stderr + if writer := r.GetOutputWriter(); writer != nil { + go io.Copy(writer, stdout) + go io.Copy(writer, stderr) + } + if err1 := cmd.Wait(); err1 != nil { + errs = append(errs, errors.Wrapf(err1, "wait playbook %s", playbook)) + } + } + return nil +} + +type PlaybookSessionBase struct { + privateKey string + + inventory string + outputWriter io.Writer + stateMux *sync.Mutex + isRunning bool + keepTmpdir bool +} + +func NewPlaybookSessionBase() PlaybookSessionBase { + return PlaybookSessionBase{ + stateMux: &sync.Mutex{}, + } +} + +func (pb *PlaybookSessionBase) GetPrivateKey() string { + return pb.privateKey +} + +func (pb *PlaybookSessionBase) IsKeepTmpdir() bool { + return pb.keepTmpdir +} + +func (pb *PlaybookSessionBase) GetOutputWriter() io.Writer { + return pb.outputWriter +} + +func (pb *PlaybookSessionBase) CheckAndSetRunning() bool { + pb.stateMux.Lock() + if pb.isRunning { + return true + } + pb.isRunning = true + pb.stateMux.Unlock() + return false +} + +func (pb *PlaybookSessionBase) SetStopped() { + pb.stateMux.Lock() + pb.isRunning = false + pb.stateMux.Unlock() +} + +func (pb *PlaybookSessionBase) GetPlaybook() string { + return "" +} + +func (pb *PlaybookSessionBase) GetPlaybookPath() string { + return "" +} + +func (pb *PlaybookSessionBase) GetInventory() string { + return pb.inventory +} + +func (pb *PlaybookSessionBase) GetConfigs() map[string]interface{} { + return nil +} + +func (pb *PlaybookSessionBase) GetRequirements() string { + return "" +} + +func (pb *PlaybookSessionBase) GetFiles() map[string][]byte { + return nil +} diff --git a/pkg/util/ansiblev2/session.go b/pkg/util/ansiblev2/session.go index 412b23dec9..26cd6cb8ae 100644 --- a/pkg/util/ansiblev2/session.go +++ b/pkg/util/ansiblev2/session.go @@ -17,34 +17,20 @@ package ansiblev2 import ( "context" "io" - "io/ioutil" - "os" - "os/exec" - "path/filepath" - "sync" - - "github.com/pkg/errors" - - yerrors "yunion.io/x/pkg/util/errors" ) type Session struct { - privateKey string + PlaybookSessionBase + playbook string - inventory string requirements string files map[string][]byte - - outputWriter io.Writer - stateMux *sync.Mutex - isRunning bool - keepTmpdir bool } func NewSession() *Session { sess := &Session{ - stateMux: &sync.Mutex{}, - files: map[string][]byte{}, + PlaybookSessionBase: NewPlaybookSessionBase(), + files: map[string][]byte{}, } return sess } @@ -95,149 +81,18 @@ func (sess *Session) KeepTmpdir(keep bool) *Session { return sess } -func (sess *Session) Run(ctx context.Context) (err error) { - var ( - tmpdir string - ) +func (sess *Session) GetPlaybook() string { + return sess.playbook +} - sess.stateMux.Lock() - if sess.isRunning { - return errors.Errorf("playbook is already running") - } - sess.isRunning = true - sess.stateMux.Unlock() - defer func() { - sess.stateMux.Lock() - sess.isRunning = false - sess.stateMux.Unlock() - }() +func (sess *Session) GetRequirements() string { + return sess.requirements +} - // make tmpdir - tmpdir, err = ioutil.TempDir("", "onecloud-ansiblev2") - if err != nil { - err = errors.WithMessage(err, "making tmp dir") - return - } - defer func() { - if sess.keepTmpdir { - return - } - if err1 := os.RemoveAll(tmpdir); err1 != nil { - err = errors.WithMessagef(err1, "removing %q", tmpdir) - } - }() +func (sess *Session) GetFile() map[string][]byte { + return sess.files +} - // write out inventory - inventory := filepath.Join(tmpdir, "inventory") - err = ioutil.WriteFile(inventory, []byte(sess.inventory), os.FileMode(0600)) - if err != nil { - err = errors.WithMessagef(err, "writing inventory %s", inventory) - return - } - - // write out playbook - playbook := filepath.Join(tmpdir, "playbook") - err = ioutil.WriteFile(playbook, []byte(sess.playbook), os.FileMode(0600)) - if err != nil { - err = errors.WithMessagef(err, "writing playbook %s", playbook) - return - } - - // write out private key - var privateKey string - if len(sess.privateKey) > 0 { - privateKey = filepath.Join(tmpdir, "private_key") - err = ioutil.WriteFile(privateKey, []byte(sess.privateKey), os.FileMode(0600)) - if err != nil { - err = errors.WithMessagef(err, "writing private key %s", privateKey) - return - } - } - - // write out requirements - var requirements string - if len(sess.requirements) > 0 { - requirements = filepath.Join(tmpdir, "requirements.yml") - err = ioutil.WriteFile(requirements, []byte(sess.requirements), os.FileMode(0600)) - if err != nil { - err = errors.WithMessagef(err, "writing requirements %s", requirements) - return - } - } - - // write out files - for name, content := range sess.files { - path := filepath.Join(tmpdir, name) - dir := filepath.Dir(path) - err = os.MkdirAll(dir, os.FileMode(0700)) - if err != nil { - err = errors.WithMessagef(err, "mkdir -p %s", dir) - return - } - err = ioutil.WriteFile(path, content, os.FileMode(0600)) - if err != nil { - err = errors.WithMessagef(err, "writing file %s", name) - return - } - } - - // run modules one by one - var errs []error - defer func() { - if len(errs) > 0 { - err = yerrors.NewAggregate(errs) - } - }() - - // install required roles - if len(requirements) > 0 { - args := []string{ - "install", "-r", requirements, "-p", tmpdir, - } - cmd := exec.CommandContext(ctx, "ansible-galaxy", args...) - stdout, _ := cmd.StdoutPipe() - stderr, _ := cmd.StderrPipe() - if err1 := cmd.Start(); err1 != nil { - errs = append(errs, errors.WithMessage(err1, "start ansible-galaxy install roles")) - return - } - // Mix stdout, stderr - if sess.outputWriter != nil { - go io.Copy(sess.outputWriter, stdout) - go io.Copy(sess.outputWriter, stderr) - } - if err1 := cmd.Wait(); err1 != nil { - errs = append(errs, errors.WithMessage(err1, "wait ansible-galaxy install roles")) - } - } - - // run playbook - { - args := []string{ - "--inventory", inventory, - } - if privateKey != "" { - args = append(args, "--private-key", privateKey) - } - args = append(args, playbook) - cmd := exec.CommandContext(ctx, "ansible-playbook", args...) - cmd.Dir = tmpdir - cmd.Env = os.Environ() - cmd.Env = append(cmd.Env, "ANSIBLE_HOST_KEY_CHECKING=False") - stdout, _ := cmd.StdoutPipe() - stderr, _ := cmd.StderrPipe() - if err1 := cmd.Start(); err1 != nil { - errs = append(errs, errors.WithMessagef(err1, "start playbook %s", playbook)) - return - } - // Mix stdout, stderr - if sess.outputWriter != nil { - go io.Copy(sess.outputWriter, stdout) - go io.Copy(sess.outputWriter, stderr) - } - if err1 := cmd.Wait(); err1 != nil { - errs = append(errs, errors.WithMessagef(err1, "wait playbook %s", playbook)) - } - } - return nil +func (sess *Session) Run(ctx context.Context) (err error) { + return runnable{sess}.Run(ctx) } From f7add553034331104ceeafbeee1ddc32965c75be Mon Sep 17 00:00:00 2001 From: rainzm Date: Sat, 9 Jan 2021 17:03:49 +0800 Subject: [PATCH 3/5] feat(cloudcommon): set TaskNotifyUrl in GetTaskRequestHeader To support cross-service task callbacks, you must set TaskNotifyUrl, which is automatically added to TaskRequestHeader here --- pkg/cloudcommon/db/taskman/tasks.go | 11 +++++++++++ pkg/compute/service/service.go | 10 +++++++++- 2 files changed, 20 insertions(+), 1 deletion(-) diff --git a/pkg/cloudcommon/db/taskman/tasks.go b/pkg/cloudcommon/db/taskman/tasks.go index b10355f9a4..c9aec3443e 100644 --- a/pkg/cloudcommon/db/taskman/tasks.go +++ b/pkg/cloudcommon/db/taskman/tasks.go @@ -19,6 +19,7 @@ import ( "database/sql" "fmt" "net/http" + "path/filepath" "reflect" "runtime/debug" "strconv" @@ -826,9 +827,19 @@ func (task *STask) GetTaskRequestHeader() http.Header { } header := mcclient.GetTokenHeaders(userCred) header.Set(mcclient.TASK_ID, task.GetTaskId()) + if len(serviceUrl) > 0 { + notifyUrl := filepath.Join(serviceUrl, "tasks", task.GetTaskId()) + header.Set(mcclient.TASK_NOTIFY_URL, notifyUrl) + } return header } +var serviceUrl string + +func SetServiceUrl(url string) { + serviceUrl = url +} + func (task *STask) GetStartTime() time.Time { return task.CreatedAt } diff --git a/pkg/compute/service/service.go b/pkg/compute/service/service.go index cc602ce9c1..21768fa04f 100644 --- a/pkg/compute/service/service.go +++ b/pkg/compute/service/service.go @@ -27,10 +27,12 @@ import ( "yunion.io/x/pkg/errors" api "yunion.io/x/onecloud/pkg/apis/compute" + "yunion.io/x/onecloud/pkg/apis/identity" "yunion.io/x/onecloud/pkg/cloudcommon" app_common "yunion.io/x/onecloud/pkg/cloudcommon/app" "yunion.io/x/onecloud/pkg/cloudcommon/cronman" "yunion.io/x/onecloud/pkg/cloudcommon/db" + "yunion.io/x/onecloud/pkg/cloudcommon/db/taskman" "yunion.io/x/onecloud/pkg/cloudcommon/elect" "yunion.io/x/onecloud/pkg/cloudcommon/etcd" common_options "yunion.io/x/onecloud/pkg/cloudcommon/options" @@ -44,6 +46,7 @@ import ( _ "yunion.io/x/onecloud/pkg/compute/tasks" "yunion.io/x/onecloud/pkg/controller/autoscaling" "yunion.io/x/onecloud/pkg/httperrors" + "yunion.io/x/onecloud/pkg/mcclient/auth" "yunion.io/x/onecloud/pkg/multicloud/esxi" _ "yunion.io/x/onecloud/pkg/multicloud/loader" ) @@ -64,8 +67,13 @@ func StartService() { log.Infof("Auth complete!!") }) common_options.StartOptionManager(opts, opts.ConfigSyncPeriodSeconds, api.SERVICE_TYPE, api.SERVICE_VERSION, options.OnOptionsChange) + serviceUrl, err := auth.GetServiceURL(api.SERVICE_TYPE, opts.Region, "", identity.EndpointInterfaceInternal) + if err != nil { + log.Fatalf("unable to get service url: %v", err) + } + taskman.SetServiceUrl(serviceUrl) - err := esxi.InitEsxiConfig(opts.EsxiOptions) + err = esxi.InitEsxiConfig(opts.EsxiOptions) if err != nil { log.Fatalf("unable to init esxi configs: %v", err) } From 1cc917e98a168d1cc5c4240be4ceb87340bdcd2b Mon Sep 17 00:00:00 2001 From: rainzm Date: Sat, 9 Jan 2021 17:01:28 +0800 Subject: [PATCH 4/5] feat(devtool): add script Script in devtool is a program or configuration that can be applied to the target host, currently only supports ansible playbook --- cmd/climc/main.go | 1 + cmd/climc/shell/devtool/common.go | 27 ++ cmd/climc/shell/devtool/devtoolcronjob.go | 222 +++++++++++++ cmd/climc/shell/devtool/devtooltemplate.go | 154 +++++++++ cmd/climc/shell/devtool/script.go | 29 ++ pkg/apis/devtool/doc.go | 15 + pkg/apis/devtool/script.go | 67 ++++ pkg/apis/devtool/script_const.go | 30 ++ pkg/apis/devtool/zz_generated.model.go | 78 +++++ pkg/devtool/models/script.go | 345 +++++++++++++++++++++ pkg/devtool/models/script_apply.go | 178 +++++++++++ pkg/devtool/models/script_apply_record.go | 143 +++++++++ pkg/devtool/service/handler.go | 4 + pkg/devtool/service/service.go | 3 +- pkg/devtool/tasks/apply_script_task.go | 218 +++++++++++++ pkg/mcclient/modules/mod_devtools.go | 21 +- pkg/mcclient/options/devtool/doc.go | 1 + pkg/mcclient/options/devtool/script.go | 66 ++++ 18 files changed, 1599 insertions(+), 3 deletions(-) create mode 100644 cmd/climc/shell/devtool/common.go create mode 100644 cmd/climc/shell/devtool/devtoolcronjob.go create mode 100644 cmd/climc/shell/devtool/devtooltemplate.go create mode 100644 cmd/climc/shell/devtool/script.go create mode 100644 pkg/apis/devtool/doc.go create mode 100644 pkg/apis/devtool/script.go create mode 100644 pkg/apis/devtool/script_const.go create mode 100644 pkg/apis/devtool/zz_generated.model.go create mode 100644 pkg/devtool/models/script.go create mode 100644 pkg/devtool/models/script_apply.go create mode 100644 pkg/devtool/models/script_apply_record.go create mode 100644 pkg/devtool/tasks/apply_script_task.go create mode 100644 pkg/mcclient/options/devtool/doc.go create mode 100644 pkg/mcclient/options/devtool/script.go diff --git a/cmd/climc/main.go b/cmd/climc/main.go index 640bce3f1e..de12cef78f 100644 --- a/cmd/climc/main.go +++ b/cmd/climc/main.go @@ -23,6 +23,7 @@ import ( _ "yunion.io/x/onecloud/cmd/climc/shell/cloudnet" _ "yunion.io/x/onecloud/cmd/climc/shell/cloudproxy" _ "yunion.io/x/onecloud/cmd/climc/shell/compute" + _ "yunion.io/x/onecloud/cmd/climc/shell/devtool" _ "yunion.io/x/onecloud/cmd/climc/shell/etcd" _ "yunion.io/x/onecloud/cmd/climc/shell/events" _ "yunion.io/x/onecloud/cmd/climc/shell/identity" diff --git a/cmd/climc/shell/devtool/common.go b/cmd/climc/shell/devtool/common.go new file mode 100644 index 0000000000..f16ff7021b --- /dev/null +++ b/cmd/climc/shell/devtool/common.go @@ -0,0 +1,27 @@ +// 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 devtool + +import ( + "yunion.io/x/onecloud/cmd/climc/shell" + "yunion.io/x/onecloud/pkg/util/printutils" +) + +var ( + R = shell.R + printList = printutils.PrintJSONList + printObject = printutils.PrintJSONObject + printBatchResults = printutils.PrintJSONBatchResults +) diff --git a/cmd/climc/shell/devtool/devtoolcronjob.go b/cmd/climc/shell/devtool/devtoolcronjob.go new file mode 100644 index 0000000000..38eb8dfb90 --- /dev/null +++ b/cmd/climc/shell/devtool/devtoolcronjob.go @@ -0,0 +1,222 @@ +// 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 devtool + +import ( + "fmt" + + "yunion.io/x/jsonutils" + "yunion.io/x/log" + + "yunion.io/x/onecloud/pkg/mcclient" + "yunion.io/x/onecloud/pkg/mcclient/modulebase" + "yunion.io/x/onecloud/pkg/mcclient/modules" + "yunion.io/x/onecloud/pkg/mcclient/options" +) + +func paramValidator(param jsonutils.JSONObject) (bool, error) { + // TODO 移到 server 端 + log.Infof("paramValidator: param: %+v", param) + interval, err := param.Int("interval") + if err != nil { + return true, nil + } + day, err := param.Int("day") + if err != nil { + return true, nil + } + if interval == 0 && day == 0 { + return false, fmt.Errorf("interval and day can not be 0 at the same time") + } + return true, nil +} + +func init() { + type CronjobListOptions struct { + options.BaseListOptions + Name string `help:"cloud region ID or Name" json:"-"` + } + + R(&CronjobListOptions{}, "devtoolcronjob-list", "List Devtool Cronjobs", func(s *mcclient.ClientSession, args *CronjobListOptions) error { + params, err := options.ListStructToParams(args) + if err != nil { + return err + } + var result *modulebase.ListResult + result, err = modules.DevToolCronjobs.List(s, params) + printList(result, modules.DevToolCronjobs.GetColumns(s)) + return nil + }) + + type CronjobCreateOptions struct { + NAME string `help:"Ansible Playbook ID or Name" json:"-"` + Day int `help:"Cronjob runs at given day" default:"0"` + Hour int `help:"Cronjob runs at given hour" default:"0"` + Min int `help:"Cronjob runs at given min" default:"0"` + Sec int `help:"Cronjob runs at given sec" default:"0"` + Interval int64 `help:"Cronjob runs at given interval" default:"0"` + Start bool `help:"start job when created" default:"false"` + Enabled bool `help:"Set job status enabled" default:"false"` + } + R( + &CronjobCreateOptions{}, + "devtoolcronjob-create", + "Create a cronjob repo component", + func(s *mcclient.ClientSession, args *CronjobCreateOptions) error { + result, err := modules.AnsiblePlaybooks.Get(s, args.NAME, nil) + if err != nil { + return err + } + ansiblePlaybookName, err := result.GetString("name") + if err != nil { + return err + } + + ansiblePlaybookID, err := result.GetString("id") + if err != nil { + return err + } + + params := jsonutils.NewDict() + params.Add(jsonutils.NewString(ansiblePlaybookName), "name") + + params.Add(jsonutils.NewString(ansiblePlaybookID), "ansible_playbook_id") + + if args.Start { + params.Add(jsonutils.JSONTrue, "start") + + } + if args.Enabled { + params.Add(jsonutils.JSONTrue, "enabled") + } else if args.Interval > 0 { + params.Add(jsonutils.NewInt(int64(args.Interval)), "interval") + } else { + params.Add(jsonutils.NewInt(int64(0)), "interval") + params.Add(jsonutils.NewInt(int64(args.Day)), "day") + params.Add(jsonutils.NewInt(int64(args.Hour)), "hour") + params.Add(jsonutils.NewInt(int64(args.Min)), "min") + params.Add(jsonutils.NewInt(int64(args.Sec)), "sec") + } + ok, err := paramValidator(params) + if err != nil || !ok { + log.Infof("paramValidator error %s", err) + return err + } + cronjob, err := modules.DevToolCronjobs.Create(s, params) + if err != nil { + log.Errorf("modules.DevToolCronjobs.Create error %s", err) + return err + } + printObject(cronjob) + return nil + }, + ) + + type DevToolCronjobShowOptions struct { + ID string `help:"ID or Name of the DevToolCronjob to show"` + } + R(&DevToolCronjobShowOptions{}, "devtoolcronjob-show", "Show cronjob details", func(s *mcclient.ClientSession, args *DevToolCronjobShowOptions) error { + result, err := modules.DevToolCronjobs.Get(s, args.ID, nil) + if err != nil { + return err + } + printObject(result) + return nil + }) + + type DevToolCronjobUpdateOptions struct { + ID string `help:"ID or Name of DevToolCronjob to update"` + Day int `help:"Cronjob runs at given day" json:"-" default:"-1"` + Hour int `help:"Cronjob runs at given hour" json:"-" default:"-1"` + Min int `help:"Cronjob runs at given min" default:"-1"` + Sec int `help:"Cronjob runs at given sec" default:"-1"` + Interval int `help:"Cronjob runs at given interval" default:"-1"` + Start bool `help:"start job when created"` + Stop bool `help:"start job when created"` + Enable bool `help:"Set job status enabled"` + Disable bool `help:"Set job status enabled"` + } + R(&DevToolCronjobUpdateOptions{}, "devtoolcronjob-update", "Update DevToolCronjob", func(s *mcclient.ClientSession, args *DevToolCronjobUpdateOptions) error { + result, err := modules.DevToolCronjobs.Get(s, args.ID, nil) + if err != nil { + return err + } + + params := jsonutils.NewDict() + interval, _ := result.Int("interval") + day, _ := result.Int("day") + params.Add(jsonutils.NewString(args.ID), "id") + params.Add(jsonutils.NewInt(int64(interval)), "interval") + params.Add(jsonutils.NewInt(int64(day)), "day") + + log.Infof("DevToolCronjobUpdateOptions args: %+v", args) + if args.Interval >= 0 { + params.Add(jsonutils.NewInt(int64(args.Interval)), "interval") + if args.Interval > 0 { + params.Add(jsonutils.NewInt(0), "day") + } + } else if args.Day >= 0 { + params.Add(jsonutils.NewInt(int64(args.Day)), "day") + if args.Day > 0 { + params.Add(jsonutils.NewInt(0), "interval") + } + } + if args.Hour >= 0 { + params.Add(jsonutils.NewInt(int64(args.Hour)), "hour") + } + if args.Min >= 0 { + params.Add(jsonutils.NewInt(int64(args.Min)), "min") + } + if args.Sec >= 0 { + params.Add(jsonutils.NewInt(int64(args.Sec)), "sec") + } + + ok, err := paramValidator(params) + if err != nil || !ok { + return err + } + + if args.Start && args.Stop { + return fmt.Errorf("can not set job start and stop at the same time") + } else if args.Start { + params.Add(jsonutils.JSONTrue, "start") + } else if args.Stop { + params.Add(jsonutils.JSONFalse, "start") + } + if args.Enable && args.Disable { + return fmt.Errorf("can not set job enabled and disabled at the same time") + } else if args.Enable { + params.Add(jsonutils.JSONTrue, "enabled") + } else if args.Disable { + params.Add(jsonutils.JSONFalse, "enabled") + } + + result, err = modules.DevToolCronjobs.Update(s, args.ID, params) + if err != nil { + return err + } + printObject(result) + return nil + }) + + R(&DevToolCronjobShowOptions{}, "devtoolcronjob-delete", "Delete DevToolCronjob", func(s *mcclient.ClientSession, args *DevToolCronjobShowOptions) error { + result, err := modules.DevToolCronjobs.Delete(s, args.ID, nil) + if err != nil { + return err + } + printObject(result) + return nil + }) +} diff --git a/cmd/climc/shell/devtool/devtooltemplate.go b/cmd/climc/shell/devtool/devtooltemplate.go new file mode 100644 index 0000000000..9c6a1d351a --- /dev/null +++ b/cmd/climc/shell/devtool/devtooltemplate.go @@ -0,0 +1,154 @@ +// 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 devtool + +import ( + "yunion.io/x/jsonutils" + "yunion.io/x/log" + + "yunion.io/x/onecloud/pkg/mcclient" + "yunion.io/x/onecloud/pkg/mcclient/modulebase" + "yunion.io/x/onecloud/pkg/mcclient/modules" + "yunion.io/x/onecloud/pkg/mcclient/options" +) + +func init() { + + printAnsiblePlaybookObject := func(obj jsonutils.JSONObject) { + dict := obj.(*jsonutils.JSONDict) + pbJson, err := dict.Get("playbook") + if err != nil { + printObject(obj) + return + } + pbStr := pbJson.YAMLString() + dict.Set("playbook", jsonutils.NewString(pbStr)) + printObject(obj) + } + + type TemplateListOptions struct { + options.BaseListOptions + Name string `help:"cloud region ID or Name" json:"-"` + } + + R(&TemplateListOptions{}, "devtooltemplate-list", "List Devtool Templates", func(s *mcclient.ClientSession, args *TemplateListOptions) error { + params, err := options.ListStructToParams(args) + if err != nil { + return err + } + var result *modulebase.ListResult + result, err = modules.DevToolTemplates.List(s, params) + printList(result, modules.DevToolTemplates.GetColumns(s)) + return nil + }) + + R( + &options.DevtoolTemplateCreateOptions{}, + "devtooltemplate-create", + "Create a template repo component", + func(s *mcclient.ClientSession, opts *options.DevtoolTemplateCreateOptions) error { + + params, err := opts.Params() + if err != nil { + return err + } + log.Infof("devtool playbook create opts: %+v", params) + apb, err := modules.DevToolTemplates.Create(s, params) + if err != nil { + return err + } + printAnsiblePlaybookObject(apb) + return nil + }, + ) + + R( + &options.DevtoolTemplateIdOptions{}, + "devtooltemplate-show", + "Show devtool template", + func(s *mcclient.ClientSession, opts *options.DevtoolTemplateIdOptions) error { + apb, err := modules.DevToolTemplates.Get(s, opts.ID, nil) + if err != nil { + return err + } + printAnsiblePlaybookObject(apb) + return nil + }, + ) + + R( + &options.DevtoolTemplateBindingOptions{}, + "devtooltemplate-bind", + "Binding devtool template to a host/vm", + func(s *mcclient.ClientSession, opts *options.DevtoolTemplateBindingOptions) error { + params := jsonutils.NewDict() + params.Set("server_id", jsonutils.NewString(opts.ServerID)) + apb, err := modules.DevToolTemplates.PerformAction(s, opts.ID, "bind", params) + if err != nil { + return err + } + printAnsiblePlaybookObject(apb) + return nil + }, + ) + + R( + &options.DevtoolTemplateBindingOptions{}, + "devtooltemplate-unbind", + "UnBinding devtool template to a host/vm", + func(s *mcclient.ClientSession, opts *options.DevtoolTemplateBindingOptions) error { + params := jsonutils.NewDict() + params.Set("server_id", jsonutils.NewString(opts.ServerID)) + apb, err := modules.DevToolTemplates.PerformAction(s, opts.ID, "unbind", params) + if err != nil { + return err + } + printAnsiblePlaybookObject(apb) + return nil + }, + ) + + R( + &options.DevtoolTemplateIdOptions{}, + "devtooltemplate-delete", + "Delete devtool template", + func(s *mcclient.ClientSession, opts *options.DevtoolTemplateIdOptions) error { + apb, err := modules.DevToolTemplates.Delete(s, opts.ID, nil) + if err != nil { + return err + } + printAnsiblePlaybookObject(apb) + return nil + }, + ) + + R( + &options.DevtoolTemplateUpdateOptions{}, + "devtooltemplate-update", + "Update ansible playbook", + func(s *mcclient.ClientSession, opts *options.DevtoolTemplateUpdateOptions) error { + params, err := opts.Params() + if err != nil { + return err + } + apb, err := modules.DevToolTemplates.Update(s, opts.ID, params) + if err != nil { + return err + } + printAnsiblePlaybookObject(apb) + return nil + }, + ) +} diff --git a/cmd/climc/shell/devtool/script.go b/cmd/climc/shell/devtool/script.go new file mode 100644 index 0000000000..4a50a12a5a --- /dev/null +++ b/cmd/climc/shell/devtool/script.go @@ -0,0 +1,29 @@ +// 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 devtool + +import ( + "yunion.io/x/onecloud/cmd/climc/shell" + "yunion.io/x/onecloud/pkg/mcclient/modules" + options "yunion.io/x/onecloud/pkg/mcclient/options/devtool" +) + +func init() { + cmd := shell.NewResourceCmd(&modules.DevToolScripts).WithKeyword("devtool-script") + cmd.List(new(options.ScriptListOptions)) + cmd.Perform("apply", new(options.ScriptApplyOptions)) + cmd1 := shell.NewResourceCmd(&modules.DevToolScriptApplyRecords).WithKeyword("devtool-script-record") + cmd1.List(new(options.ScriptApplyRecordListOptions)) +} diff --git a/pkg/apis/devtool/doc.go b/pkg/apis/devtool/doc.go new file mode 100644 index 0000000000..c3a15f4646 --- /dev/null +++ b/pkg/apis/devtool/doc.go @@ -0,0 +1,15 @@ +// 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 devtool // import "yunion.io/x/onecloud/pkg/apis/devtool" diff --git a/pkg/apis/devtool/script.go b/pkg/apis/devtool/script.go new file mode 100644 index 0000000000..0a60baf0c8 --- /dev/null +++ b/pkg/apis/devtool/script.go @@ -0,0 +1,67 @@ +// 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 devtool + +import "yunion.io/x/onecloud/pkg/apis" + +type ScriptApplyInput struct { + // description: server id + // required: true + // example: b48c5c84-9952-4394-8ca9-c3b84e946a03 + ServerID string + // description: whether to use eip first + // example: true + EipFirst bool + // description: Id of proxyEndpoint + // example: cf1d1a0f-9b9d-4629-8036-af3ed87c0821 + ProxyEndpointId string + // description: whether to automatically select proxy endpoint + AutoChooseProxyEndpoint bool +} + +type ScriptApplyOutput struct { + // description: Instantiation of script apply + // example: cf1d1a0f-9b9d-4629-8036-af3ed87c0821 + ScriptApplyId string +} + +type ScriptApplyRecoredListInput struct { + apis.StatusStandaloneResourceListInput + // description: Id of Script + // example: cc2e2ba6-e33d-4be3-8e2d-4d2aa843dd03 + ScriptId string +} + +type ScriptCreateInput struct { + apis.SharableVirtualResourceCreateInput + // description: Id or Name of ansible playbook reference + // example: cf1d1a0f-9b9d-4629-8036-af3ed87c0821 + PlaybookReference string + // description: The script may fail to execute, MaxTryTime represents the maximum number of attempts to execute + MaxTryTimes int +} + +type ScriptDetails struct { + apis.SharableVirtualResourceDetails + SScript + ApplyInfos []SApplyInfo +} + +type SApplyInfo struct { + ServerId string + EipFirst bool + ProxyEndpointId string + TryTimes int +} diff --git a/pkg/apis/devtool/script_const.go b/pkg/apis/devtool/script_const.go new file mode 100644 index 0000000000..631e8c001e --- /dev/null +++ b/pkg/apis/devtool/script_const.go @@ -0,0 +1,30 @@ +// 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 devtool + +const ( + SCRIPT_APPLY_STATUS_APPLYING = "applying" + SCRIPT_APPLY_STATUS_APPLY_FAILED = "apply_failed" + SCRIPT_APPLY_STATUS_READY = "ready" + + SCRIPT_APPLY_RECORD_APPLYING = "applying" + SCRIPT_APPLY_RECORD_SUCCEED = "succeed" + SCRIPT_APPLY_RECORD_FAILED = "failed" + + SCRIPT_NAME = "monitor agent" + SERVICE_TYPE = "devtool" + + SCRIPT_STATUS_READY = "ready" +) diff --git a/pkg/apis/devtool/zz_generated.model.go b/pkg/apis/devtool/zz_generated.model.go new file mode 100644 index 0000000000..342a508cd5 --- /dev/null +++ b/pkg/apis/devtool/zz_generated.model.go @@ -0,0 +1,78 @@ +// 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. + +// Code generated by model-api-gen. DO NOT EDIT. + +package devtool + +import ( + time "time" + + "yunion.io/x/onecloud/pkg/apis" + ansible "yunion.io/x/onecloud/pkg/util/ansible" +) + +// SCronjob is an autogenerated struct via yunion.io/x/onecloud/pkg/devtool/models.SCronjob. +type SCronjob struct { + SVSCronjob + AnsiblePlaybookID string `json:"ansible_playbook_id"` + TemplateID string `json:"template_id"` + ServerID string `json:"server_id"` + apis.SVirtualResourceBase +} + +// SDevtoolTemplate is an autogenerated struct via yunion.io/x/onecloud/pkg/devtool/models.SDevtoolTemplate. +type SDevtoolTemplate struct { + SVSCronjob + Playbook *ansible.Playbook `json:"playbook"` + apis.SVirtualResourceBase +} + +// SScript is an autogenerated struct via yunion.io/x/onecloud/pkg/devtool/models.SScript. +type SScript struct { + apis.SVirtualResourceBase + // remote + Type string `json:"type"` + PlaybookReference string `json:"playbook_reference"` + MaxTryTimes int `json:"max_try_times"` +} + +// SScriptApply is an autogenerated struct via yunion.io/x/onecloud/pkg/devtool/models.SScriptApply. +type SScriptApply struct { + apis.SStatusStandaloneResourceBase + ScriptId string `json:"script_id"` + GuestId string `json:"guest_id"` + EipFirst *bool `json:"eip_first,omitempty"` + ProxyEndpointId string `json:"proxy_endpoint_id"` + TryTimes int `json:"try_times"` +} + +// SScriptApplyRecord is an autogenerated struct via yunion.io/x/onecloud/pkg/devtool/models.SScriptApplyRecord. +type SScriptApplyRecord struct { + apis.SStatusStandaloneResourceBase + ScriptId string `json:"script_id"` + ServerId string `json:"server_id"` + StartTime time.Time `json:"start_time"` + EndTime time.Time `json:"end_time"` + Reason string `json:"reason"` +} + +// SVSCronjob is an autogenerated struct via yunion.io/x/onecloud/pkg/devtool/models.SVSCronjob. +type SVSCronjob struct { + Day int `json:"day"` + Hour int `json:"hour"` + Min int `json:"min"` + Sec int `json:"sec"` + Interval int64 `json:"interval"` + Start bool `json:"start"` + Enabled bool `json:"enabled"` +} diff --git a/pkg/devtool/models/script.go b/pkg/devtool/models/script.go new file mode 100644 index 0000000000..e433bfc690 --- /dev/null +++ b/pkg/devtool/models/script.go @@ -0,0 +1,345 @@ +// 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 models + +import ( + "context" + "fmt" + "net/url" + "sync" + + "github.com/coredns/coredns/plugin/pkg/log" + + "yunion.io/x/jsonutils" + "yunion.io/x/pkg/errors" + "yunion.io/x/pkg/util/sets" + + proxy_api "yunion.io/x/onecloud/pkg/apis/cloudproxy" + comapi "yunion.io/x/onecloud/pkg/apis/compute" + api "yunion.io/x/onecloud/pkg/apis/devtool" + "yunion.io/x/onecloud/pkg/cloudcommon/db" + "yunion.io/x/onecloud/pkg/httperrors" + "yunion.io/x/onecloud/pkg/mcclient" + "yunion.io/x/onecloud/pkg/mcclient/auth" + "yunion.io/x/onecloud/pkg/mcclient/modules" + "yunion.io/x/onecloud/pkg/mcclient/modules/cloudproxy" + "yunion.io/x/onecloud/pkg/util/httputils" + "yunion.io/x/onecloud/pkg/util/stringutils2" +) + +type SScript struct { + db.SSharableVirtualResourceBase + // remote + Type string `width:"16" nullable:"false"` + PlaybookReferenceId string `width:"128" nullable:"false"` + MaxTryTimes int `default:"1"` +} + +type SScriptManager struct { + db.SSharableVirtualResourceBaseManager +} + +var ScriptManager *SScriptManager + +func init() { + ScriptManager = &SScriptManager{ + SSharableVirtualResourceBaseManager: db.NewSharableVirtualResourceBaseManager( + SScript{}, + "script_tbl", + "script", + "scripts", + ), + } + ScriptManager.SetVirtualObject(ScriptManager) + registerArgGenerator(MonitorAgent, getArgs) +} + +type argGenerator func(ctx context.Context, input api.ScriptApplyInput, details *comapi.ServerDetails) (map[string]interface{}, error) + +var argGenerators = &sync.Map{} + +func registerArgGenerator(name string, ag argGenerator) { + argGenerators.Store(name, ag) +} + +func getArgGenerator(name string) (argGenerator, bool) { + v, ok := argGenerators.Load(name) + if !ok { + return nil, ok + } + return v.(argGenerator), ok +} + +func convertInfluxdbUrl(ctx context.Context, pUrl string, endpointId string) (string, error) { + session := auth.AdminSessionWithInternal(ctx, "", "", "") + filter := jsonutils.NewDict() + filter.Set("proxy_endpoint_id", jsonutils.NewString(endpointId)) + filter.Set("opaque", jsonutils.NewString(pUrl)) + filter.Set("scope", jsonutils.NewString("system")) + lr, err := cloudproxy.Forwards.List(session, filter) + if err != nil { + return "", errors.Wrap(err, "failed to list forward") + } + var port int64 + if len(lr.Data) > 0 { + port, _ = lr.Data[0].Int("bind_port") + } else { + rUrl, err := url.Parse(pUrl) + if err != nil { + return "", errors.Wrap(err, "invalid influxdbUrl?") + } + // create one + createP := jsonutils.NewDict() + createP.Set("proxy_endpoint", jsonutils.NewString(endpointId)) + createP.Set("type", jsonutils.NewString(proxy_api.FORWARD_TYPE_REMOTE)) + createP.Set("remote_addr", jsonutils.NewString(rUrl.Hostname())) + createP.Set("remote_port", jsonutils.NewString(rUrl.Port())) + createP.Set("generate_name", jsonutils.NewString("influxdb proxy")) + createP.Set("opaque", jsonutils.NewString(pUrl)) + forward, err := cloudproxy.Forwards.Create(session, createP) + if err != nil { + return "", errors.Wrapf(err, "unable to create forward with create params %s", createP.String()) + } + port, _ = forward.Int("bind_port") + } + // fetch proxy_endpoint address + ep, err := cloudproxy.ProxyEndpoints.Get(session, endpointId, nil) + if err != nil { + return "", errors.Wrapf(err, "unable to get proxy endpoint %s", endpointId) + } + address, _ := ep.GetString("intranet_ip_addr") + return fmt.Sprintf("https://%s:%d", address, port), nil +} + +func getArgs(ctx context.Context, input api.ScriptApplyInput, detail *comapi.ServerDetails) (map[string]interface{}, error) { + influxdbUrl, err := getInfluxdbUrl(ctx) + if err != nil { + return nil, errors.Wrap(err, "unable to get influxdbUrl") + } + // convert influxdbUrl + if len(input.ProxyEndpointId) > 0 { + influxdbUrl, err = convertInfluxdbUrl(ctx, influxdbUrl, input.ProxyEndpointId) + if err != nil { + return nil, errors.Wrapf(err, "unable to convertInfluxdbUrl %s", influxdbUrl) + } + } + vmId := detail.Id + tenantId := detail.ProjectId + domainId := detail.DomainId + ret := map[string]interface{}{ + "influxdb_url": influxdbUrl, + "influxdb_name": "telegraf", + "onecloud_vm_id": vmId, + "onecloud_tenant_id": tenantId, + "onecloud_domain_id": domainId, + } + return ret, nil +} + +var influxdbUrl string + +func getInfluxdbUrl(ctx context.Context) (string, error) { + if len(influxdbUrl) > 0 { + return influxdbUrl, nil + } + session := auth.GetAdminSession(ctx, "", "") + params := jsonutils.NewDict() + params.Set("interface", jsonutils.NewString("public")) + params.Set("service", jsonutils.NewString("influxdb")) + ret, err := modules.EndpointsV3.List(session, params) + if err != nil { + return "", err + } + if len(ret.Data) == 0 { + return "", fmt.Errorf("no sucn endpoint with 'internal' interface and 'influxdb' service") + } + url, _ := ret.Data[0].GetString("url") + return url, nil +} + +var MonitorAgent = "monitor agent" + +func (sm *SScriptManager) InitializeData() error { + q := sm.Query().Equals("playbook_reference", MonitorAgent) + n, err := q.CountWithError() + if err != nil { + return err + } + if n > 0 { + return nil + } + s := SScript{ + PlaybookReferenceId: MonitorAgent, + } + s.ProjectId = "system" + s.IsPublic = true + s.PublicScope = "system" + err = sm.TableSpec().Insert(context.Background(), &s) + if err != nil { + return err + } + return nil +} + +func (sm *SScriptManager) ValidateCreateData(ctx context.Context, userCred mcclient.TokenCredential, ownerId mcclient.IIdentityProvider, query jsonutils.JSONObject, input api.ScriptCreateInput) (api.ScriptCreateInput, error) { + // check ansible playbook reference + session := auth.GetSessionWithInternal(ctx, userCred, "", "") + pr, err := modules.AnsiblePlaybookReference.Get(session, input.PlaybookReference, nil) + if err != nil { + return input, errors.Wrapf(err, "unable to get AnsiblePlaybookReference %q", input.PlaybookReference) + } + id, _ := pr.GetString("id") + input.PlaybookReference = id + return input, nil +} + +func (s *SScript) CustomizeCreate(ctx context.Context, userCred mcclient.TokenCredential, ownerId mcclient.IIdentityProvider, query jsonutils.JSONObject, data jsonutils.JSONObject) error { + s.Status = api.SCRIPT_STATUS_READY + s.PlaybookReferenceId, _ = data.GetString("playbook_reference") + return nil +} + +func (sm *SScriptManager) FetchCustomizeColumns(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, objs []interface{}, fields stringutils2.SSortedStrings, isList bool) []api.ScriptDetails { + vDetails := sm.SSharableVirtualResourceBaseManager.FetchCustomizeColumns(ctx, userCred, query, objs, fields, isList) + details := make([]api.ScriptDetails, len(objs)) + for i := range details { + details[i].SharableVirtualResourceDetails = vDetails[i] + script := objs[i].(*SScript) + ais, err := script.ApplyInfos() + if err != nil { + log.Errorf("unable to get ApplyInfos of script %s: %v", script.Id, err) + } + details[i].ApplyInfos = ais + } + return details +} + +func (s *SScript) ApplyInfos() ([]api.SApplyInfo, error) { + q := ScriptApplyManager.Query().Equals("script_id", s.Id) + var sa []SScriptApply + err := db.FetchModelObjects(ScriptApplyManager, q, &sa) + if err != nil { + return nil, err + } + ai := make([]api.SApplyInfo, len(sa)) + for i := range ai { + ai[i].ServerId = sa[i].GuestId + ai[i].EipFirst = sa[i].EipFirst.Bool() + ai[i].ProxyEndpointId = sa[i].ProxyEndpointId + ai[i].TryTimes = sa[i].TryTimes + } + return ai, nil +} + +func (s *SScript) AllowPerformApply(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject) bool { + return s.IsOwner(userCred) || db.IsAdminAllowPerform(userCred, s, "apply") +} + +func (s *SScript) PerformApply(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, input api.ScriptApplyInput) (api.ScriptApplyOutput, error) { + output := api.ScriptApplyOutput{} + serverInfo, err := s.checkServer(ctx, userCred, input.ServerID) + if err != nil { + return output, err + } + // select proxyEndpoint automatically + if len(input.ProxyEndpointId) == 0 && input.AutoChooseProxyEndpoint { + var proxyEndpointId string + // find suitable proxyEndpoint + // network first + session := auth.GetAdminSession(ctx, "", "") + for _, netId := range serverInfo.NetworkIds { + filter := jsonutils.NewDict() + filter.Set("network_id", jsonutils.NewString(netId)) + lr, err := cloudproxy.ProxyEndpoints.List(session, filter) + if err != nil { + return output, errors.Wrapf(err, "unable to list proxy endpoint in network %q", netId) + } + if len(lr.Data) == 0 { + continue + } + proxyEndpointId, _ = lr.Data[0].GetString("id") + break + } + if len(proxyEndpointId) == 0 { + filter := jsonutils.NewDict() + filter.Set("vpc_id", jsonutils.NewString(serverInfo.VpcId)) + lr, err := cloudproxy.ProxyEndpoints.List(session, filter) + if err != nil { + return output, errors.Wrapf(err, "unable to list proxy endpoint in vpc %q", serverInfo.VpcId) + } + if len(lr.Data) > 0 { + // TODO Choose strictly + proxyEndpointId, _ = lr.Data[0].GetString("id") + } + } + if len(proxyEndpointId) == 0 { + return output, httperrors.NewInputParameterError("can't find suitable proxy endpoint for server %s, please connect with admin to create one", serverInfo.serverDetails.Name) + } + input.ProxyEndpointId = proxyEndpointId + } + ag, _ := getArgGenerator(MonitorAgent) + args, err := ag(ctx, input, serverInfo.serverDetails) + if err != nil { + return output, errors.Wrapf(err, "unable to get args of server %s", serverInfo.ServerId) + } + sa, err := ScriptApplyManager.createScriptApply(ctx, s.Id, serverInfo.ServerId, input.ProxyEndpointId, input.EipFirst, args) + if err != nil { + return output, errors.Wrapf(err, "unable to apply script to server %s", serverInfo.ServerId) + } + err = sa.StartApply(ctx, userCred) + if err != nil { + return output, errors.Wrapf(err, "unable to apply script to server %s", serverInfo.ServerId) + } + output.ScriptApplyId = sa.Id + return output, nil +} + +type sServerInfo struct { + ServerId string + VpcId string + NetworkIds []string + serverDetails *comapi.ServerDetails +} + +func (s *SScript) checkServer(ctx context.Context, userCred mcclient.TokenCredential, serverId string) (sServerInfo, error) { + session := auth.GetSessionWithInternal(ctx, userCred, "", "") + // check server + data, err := modules.Servers.Get(session, serverId, nil) + if err != nil { + if httputils.ErrorCode(err) == 404 { + return sServerInfo{}, httperrors.NewInputParameterError("no such server %s", serverId) + } + return sServerInfo{}, fmt.Errorf("unable to get server %s: %s", serverId, httputils.ErrorMsg(err)) + } + info := sServerInfo{} + var serverDetails comapi.ServerDetails + err = data.Unmarshal(&serverDetails) + if err != nil { + return info, errors.Wrap(err, "unable to unmarshal serverDetails") + } + if serverDetails.Status != comapi.VM_RUNNING { + return info, httperrors.NewInputParameterError("can only apply scripts to %s server", comapi.VM_RUNNING) + } + info.serverDetails = &serverDetails + info.ServerId = serverDetails.Id + + networkIds := sets.NewString() + for _, nic := range serverDetails.Nics { + networkIds.Insert(nic.NetworkId) + info.VpcId = nic.VpcId + } + info.NetworkIds = networkIds.UnsortedList() + return info, nil +} diff --git a/pkg/devtool/models/script_apply.go b/pkg/devtool/models/script_apply.go new file mode 100644 index 0000000000..7a6045d4e5 --- /dev/null +++ b/pkg/devtool/models/script_apply.go @@ -0,0 +1,178 @@ +// 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 models + +import ( + "context" + "fmt" + "sync" + + "yunion.io/x/jsonutils" + "yunion.io/x/pkg/errors" + "yunion.io/x/pkg/tristate" + "yunion.io/x/pkg/util/sets" + + api "yunion.io/x/onecloud/pkg/apis/devtool" + "yunion.io/x/onecloud/pkg/cloudcommon/db" + "yunion.io/x/onecloud/pkg/cloudcommon/db/taskman" + "yunion.io/x/onecloud/pkg/mcclient" +) + +type SScriptApply struct { + db.SStatusStandaloneResourceBase + ScriptId string `width:"36" nullable:"false" index:"true"` + GuestId string `width:"36" nullable:"false" index:"true"` + EipFirst tristate.TriState + // + Args jsonutils.JSONObject + ProxyEndpointId string `width:"36" nullable:"false"` + TryTimes int +} + +type SScriptApplyManager struct { + db.SStatusStandaloneResourceBaseManager + Session *sScriptApplySession +} + +var ScriptApplyManager *SScriptApplyManager + +func init() { + ScriptApplyManager = &SScriptApplyManager{ + SStatusStandaloneResourceBaseManager: db.NewStatusStandaloneResourceBaseManager( + SScriptApply{}, + "scriptapply_tbl", + "scriptapply", + "scirptapplys", + ), + Session: newScriptApplySession(), + } + ScriptApplyManager.SetVirtualObject(ScriptApplyManager) +} + +func (sam *SScriptApplyManager) createScriptApply(ctx context.Context, scriptId, guestId, proxyEndpointId string, eipFirst bool, args map[string]interface{}) (*SScriptApply, error) { + sa := &SScriptApply{ + ScriptId: scriptId, + GuestId: guestId, + EipFirst: tristate.NewFromBool(eipFirst), + ProxyEndpointId: proxyEndpointId, + Args: jsonutils.Marshal(args), + } + err := ScriptApplyManager.TableSpec().Insert(ctx, sa) + sa.SetModelManager(ScriptApplyManager, sa) + return sa, err +} + +func (sa *SScriptApply) StartApply(ctx context.Context, userCred mcclient.TokenCredential) (err error) { + if ok := ScriptApplyManager.Session.CheckAndSet(sa.Id); !ok { + return fmt.Errorf("script %s is applying to server %s", sa.ScriptId, sa.GuestId) + } + defer func() { + if err != nil { + ScriptApplyManager.Session.Remove(sa.Id) + } + }() + // check try times + script, err := sa.Script() + if err != nil { + return err + } + if sa.TryTimes >= script.MaxTryTimes { + return fmt.Errorf("The times to try has exceeded the maximum times %d setted by the script", script.MaxTryTimes) + } + _, err = db.Update(sa, func() error { + sa.TryTimes += 1 + sa.Status = api.SCRIPT_APPLY_STATUS_APPLYING + return nil + }) + if err != nil { + return errors.Wrap(err, "unable to update scriptapply") + } + + err = sa.startApplyScriptTask(ctx, userCred, "") + if err != nil { + f := false + _, err = ScriptApplyRecordManager.createRecordWithResult(ctx, sa.ScriptId, sa.GuestId, &f, fmt.Sprintf("unabel to start ApplyScriptTask: %v", err)) + if err != nil { + return errors.Wrap(err, "unable to record") + } + ScriptApplyManager.Session.Remove(sa.Id) + } + return nil +} + +func (sa *SScriptApply) startApplyScriptTask(ctx context.Context, userCred mcclient.TokenCredential, parentTaskId string) error { + task, err := taskman.TaskManager.NewTask(ctx, "ApplyScriptTask", sa, userCred, nil, "", parentTaskId) + if err != nil { + return err + } + task.ScheduleRun(nil) + return nil +} + +func (sa *SScriptApply) StopApply(userCred mcclient.TokenCredential, record *SScriptApplyRecord, success bool, reason string) error { + var status string + if success { + status = api.SCRIPT_APPLY_STATUS_READY + if record != nil { + record.Succeed(reason) + } + } else { + status = api.SCRIPT_APPLY_RECORD_FAILED + if record != nil { + record.Fail(reason) + } + } + sa.SetStatus(userCred, status, "") + ScriptApplyManager.Session.Remove(sa.Id) + return nil +} + +func (sa *SScriptApply) Script() (*SScript, error) { + obj, err := ScriptManager.FetchById(sa.ScriptId) + if err != nil { + return nil, errors.Wrapf(err, "unable to fetch Script %s", sa.Id) + } + s := obj.(*SScript) + s.SetModelManager(ScriptManager, s) + return s, nil +} + +type sScriptApplySession struct { + mux *sync.Mutex + applyingOnes sets.String +} + +func newScriptApplySession() *sScriptApplySession { + return &sScriptApplySession{ + mux: &sync.Mutex{}, + applyingOnes: sets.NewString(), + } +} + +func (sas *sScriptApplySession) CheckAndSet(id string) bool { + sas.mux.Lock() + defer sas.mux.Unlock() + if sas.applyingOnes.Has(id) { + return false + } + sas.applyingOnes.Insert(id) + return true +} + +func (sas *sScriptApplySession) Remove(id string) { + sas.mux.Lock() + defer sas.mux.Unlock() + sas.applyingOnes.Delete(id) +} diff --git a/pkg/devtool/models/script_apply_record.go b/pkg/devtool/models/script_apply_record.go new file mode 100644 index 0000000000..2d820a1972 --- /dev/null +++ b/pkg/devtool/models/script_apply_record.go @@ -0,0 +1,143 @@ +// 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 models + +import ( + "context" + "time" + + "yunion.io/x/jsonutils" + "yunion.io/x/sqlchemy" + + api "yunion.io/x/onecloud/pkg/apis/devtool" + "yunion.io/x/onecloud/pkg/cloudcommon/db" + "yunion.io/x/onecloud/pkg/mcclient" + "yunion.io/x/onecloud/pkg/util/rbacutils" +) + +type SScriptApplyRecord struct { + db.SStatusStandaloneResourceBase + ScriptId string `width:"36" charset:"ascii" nullable:"true" list:"user" index:"true"` + ServerId string `width:"36" charset:"ascii" nullable:"true" list:"user"` + StartTime time.Time `list:"user"` + EndTime time.Time `list:"user"` + Reason string `list:"user"` +} + +type SScriptApplyRecordManager struct { + db.SStatusStandaloneResourceBaseManager +} + +var ScriptApplyRecordManager *SScriptApplyRecordManager + +func init() { + ScriptApplyRecordManager = &SScriptApplyRecordManager{ + SStatusStandaloneResourceBaseManager: db.NewStatusStandaloneResourceBaseManager( + SScriptApplyRecord{}, + "scriptapplyrecord_tbl", + "scriptapplyrecord", + "scriptapplyrecords", + ), + } + ScriptApplyRecordManager.SetVirtualObject(ScriptApplyRecordManager) +} + +func (sarm *SScriptApplyRecordManager) ListItemFilter(ctx context.Context, q *sqlchemy.SQuery, userCred mcclient.TokenCredential, input api.ScriptApplyRecoredListInput) (*sqlchemy.SQuery, error) { + q, err := sarm.SStatusStandaloneResourceBaseManager.ListItemFilter(ctx, q, userCred, input.StatusStandaloneResourceListInput) + if err != nil { + return q, err + } + if len(input.ScriptId) > 0 { + q = q.Equals("script_id", input.ScriptId) + } + return q, nil +} + +func (sarm *SScriptApplyRecordManager) CreateRecord(ctx context.Context, scriptId, serverId string) (*SScriptApplyRecord, error) { + return sarm.createRecordWithResult(ctx, scriptId, serverId, nil, "") +} + +func (sarm *SScriptApplyRecordManager) createRecordWithResult(ctx context.Context, scriptId, serverId string, success *bool, reason string) (*SScriptApplyRecord, error) { + now := time.Now() + sar := &SScriptApplyRecord{ + StartTime: now, + ScriptId: scriptId, + ServerId: serverId, + } + if success == nil { + sar.Status = api.SCRIPT_APPLY_RECORD_APPLYING + } else if *success { + sar.Status = api.SCRIPT_APPLY_RECORD_SUCCEED + } else if !*success { + sar.Status = api.SCRIPT_APPLY_RECORD_FAILED + } + sar.Reason = reason + err := sarm.TableSpec().Insert(ctx, sar) + if err != nil { + return nil, err + } + sar.SetModelManager(sarm, sar) + return sar, nil +} + +func (sarm *SScriptApplyRecordManager) NamespaceScope() rbacutils.TRbacScope { + return rbacutils.ScopeProject +} + +func (sarm *SScriptApplyRecordManager) ResourceScope() rbacutils.TRbacScope { + return rbacutils.ScopeProject +} + +func (sarm *SScriptApplyRecordManager) FileterByOwner(q *sqlchemy.SQuery, owner mcclient.IIdentityProvider, scope rbacutils.TRbacScope) *sqlchemy.SQuery { + if owner != nil { + switch scope { + case rbacutils.ScopeProject, rbacutils.ScopeDomain: + scriptQ := ScriptManager.Query("id", "domain_id").SubQuery() + q = q.Join(scriptQ, sqlchemy.Equals(q.Field("script_id"), scriptQ.Field("id"))) + q = q.Filter(sqlchemy.Equals(scriptQ.Field("domain_id"), owner.GetProjectDomainId())) + } + } + return q +} + +func (sarm *SScriptApplyRecordManager) FetchOwnerId(ctx context.Context, data jsonutils.JSONObject) (mcclient.IIdentityProvider, error) { + return db.FetchDomainInfo(ctx, data) +} + +func (sar *SScriptApplyRecord) GetOwnerId() mcclient.IIdentityProvider { + obj, _ := ScriptManager.FetchById(sar.ScriptId) + if obj == nil { + return nil + } + return obj.GetOwnerId() +} + +func (sar *SScriptApplyRecord) SetResult(status, reason string) error { + _, err := db.Update(sar, func() error { + sar.Status = status + sar.Reason = reason + sar.EndTime = time.Now() + return nil + }) + return err +} + +func (sar *SScriptApplyRecord) Fail(reason string) error { + return sar.SetResult(api.SCRIPT_APPLY_RECORD_FAILED, reason) +} + +func (sar *SScriptApplyRecord) Succeed(reason string) error { + return sar.SetResult(api.SCRIPT_APPLY_RECORD_SUCCEED, reason) +} diff --git a/pkg/devtool/service/handler.go b/pkg/devtool/service/handler.go index b4a645e4fe..303e17b969 100644 --- a/pkg/devtool/service/handler.go +++ b/pkg/devtool/service/handler.go @@ -30,6 +30,7 @@ func InitHandlers(app *appsrv.Application) { taskman.TaskManager, taskman.SubTaskManager, taskman.TaskObjectManager, + db.SharedResourceManager, db.UserCacheManager, db.TenantCacheManager, } { @@ -42,6 +43,9 @@ func InitHandlers(app *appsrv.Application) { models.CronjobManager, models.DevtoolTemplateManager, + models.ScriptManager, + models.ScriptApplyManager, + models.ScriptApplyRecordManager, } { db.RegisterModelManager(manager) handler := db.NewModelHandler(manager) diff --git a/pkg/devtool/service/service.go b/pkg/devtool/service/service.go index 29b1fa6f70..61cc5f9f8a 100644 --- a/pkg/devtool/service/service.go +++ b/pkg/devtool/service/service.go @@ -21,6 +21,7 @@ import ( "yunion.io/x/log" + api "yunion.io/x/onecloud/pkg/apis/devtool" "yunion.io/x/onecloud/pkg/cloudcommon" app_common "yunion.io/x/onecloud/pkg/cloudcommon/app" "yunion.io/x/onecloud/pkg/cloudcommon/db" @@ -36,7 +37,7 @@ func StartService() { commonOpts := &opts.CommonOptions dbOpts := &options.Options.DBOptions baseOpts := &opts.BaseOptions - common_options.ParseOptions(opts, os.Args, "devtool.conf", "devtool") + common_options.ParseOptions(opts, os.Args, "devtool.conf", api.SERVICE_TYPE) app_common.InitAuth(commonOpts, func() { log.Infof("Auth complete!!") diff --git a/pkg/devtool/tasks/apply_script_task.go b/pkg/devtool/tasks/apply_script_task.go new file mode 100644 index 0000000000..31161cc98e --- /dev/null +++ b/pkg/devtool/tasks/apply_script_task.go @@ -0,0 +1,218 @@ +// 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" + "net" + "time" + + "yunion.io/x/jsonutils" + "yunion.io/x/log" + "yunion.io/x/pkg/errors" + + ansible_api "yunion.io/x/onecloud/pkg/apis/ansible" + cloudproxy_api "yunion.io/x/onecloud/pkg/apis/cloudproxy" + comapi "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/devtool/models" + "yunion.io/x/onecloud/pkg/mcclient" + "yunion.io/x/onecloud/pkg/mcclient/auth" + "yunion.io/x/onecloud/pkg/mcclient/modules" + "yunion.io/x/onecloud/pkg/mcclient/modules/cloudproxy" +) + +type ApplyScriptTask struct { + taskman.STask +} + +func init() { + taskman.RegisterTask(ApplyScriptTask{}) +} + +func (self *ApplyScriptTask) taskFailed(ctx context.Context, sa *models.SScriptApply, sar *models.SScriptApplyRecord, err error) { + err = sa.StopApply(self.UserCred, sar, false, err.Error()) + if err != nil { + log.Errorf("unable to StopApply script %s to server %s", sa.ScriptId, sa.GuestId) + self.SetStageFailed(ctx, jsonutils.NewString(err.Error())) + return + } + // restart + err = sa.StartApply(ctx, self.UserCred) + if err != nil { + log.Errorf("unable to StartApply script %s to server %s", sa.ScriptId, sa.GuestId) + } + self.SetStageFailed(ctx, jsonutils.NewString(err.Error())) +} + +func (self *ApplyScriptTask) taskSuccess(ctx context.Context, sa *models.SScriptApply, sar *models.SScriptApplyRecord) { + err := sa.StopApply(self.UserCred, sar, true, "") + if err != nil { + log.Errorf("unable to StopApply script %s to server %s", sa.ScriptId, sa.GuestId) + self.SetStageComplete(ctx, nil) + } +} + +func (self *ApplyScriptTask) OnInit(ctx context.Context, obj db.IStandaloneModel, body jsonutils.JSONObject) { + sa := obj.(*models.SScriptApply) + // create record + sar, err := models.ScriptApplyRecordManager.CreateRecord(ctx, sa.ScriptId, sa.GuestId) + if err != nil { + self.taskFailed(ctx, sa, nil, err) + return + } + s, err := sa.Script() + if err != nil { + self.taskFailed(ctx, sa, sar, err) + return + } + session := auth.GetAdminSession(ctx, "", "") + params := jsonutils.NewDict() + params.Set("details", jsonutils.JSONTrue) + data, err := modules.Servers.GetById(session, sa.GuestId, params) + if err != nil { + self.taskFailed(ctx, sa, sar, errors.Wrapf(err, "unable to fetch server %s", sa.GuestId)) + return + } + var serverDetail comapi.ServerDetails + err = data.Unmarshal(&serverDetail) + if err != nil { + self.taskFailed(ctx, sa, sar, errors.Wrapf(err, "unable to unmarshal %q to ServerDetails", data)) + return + } + // make sure user + var user string + if serverDetail.Hypervisor == comapi.HYPERVISOR_KVM { + user = "root" + } else { + user = "cloudroot" + } + // create local forward + createP := jsonutils.NewDict() + createP.Set("type", jsonutils.NewString(cloudproxy_api.FORWARD_TYPE_LOCAL)) + createP.Set("remote_port", jsonutils.NewInt(22)) + createP.Set("server_id", jsonutils.NewString(serverDetail.Id)) + + forward, err := cloudproxy.Forwards.PerformClassAction(session, "create-from-server", createP) + if err != nil { + self.taskFailed(ctx, sa, sar, errors.Wrapf(err, "fail to create local forward from server %q", serverDetail.Id)) + return + } + + port, _ := forward.Int("bind_port") + forwardId, _ := forward.GetString("id") + agentId, _ := forward.GetString("proxy_agent_id") + agent, err := cloudproxy.ProxyAgents.Get(session, agentId, nil) + if err != nil { + self.clearLocalForward(session, forwardId) + self.taskFailed(ctx, sa, sar, errors.Wrapf(err, "fail to get proxy agent %q", agentId)) + return + } + address, _ := agent.GetString("advertise_addr") + host := ansible_api.AnsibleHost{ + User: user, + IP: address, + Port: int(port), + Name: serverDetail.Name, + } + params = jsonutils.NewDict() + params.Set("args", sa.Args) + params.Set("host", jsonutils.Marshal(host)) + // fetch ansible playbook reference id + updateData := jsonutils.NewDict() + updateData.Set("script_apply_record_id", jsonutils.NewString(sar.GetId())) + updateData.Set("proxy_forward_id", jsonutils.NewString(forwardId)) + + // check proxy forward + if ok := self.ensureLocalForwardWork(address, int(port)); !ok { + self.clearLocalForward(session, forwardId) + self.taskFailed(ctx, sa, sar, errors.Wrapf(err, "The created local forward is actually not usable")) + return + } + self.SetStage("OnAnsiblePlaybookComplete", updateData) + // Inject Task Header + session.Header = self.GetTaskRequestHeader() + _, err = modules.AnsiblePlaybookReference.PerformAction(session, s.PlaybookReferenceId, "run", params) + if err != nil { + self.clearLocalForward(session, forwardId) + self.taskFailed(ctx, sa, sar, errors.Wrapf(err, "can't run ansible playbook reference %s", s.PlaybookReferenceId)) + return + } +} + +func (self *ApplyScriptTask) clearLocalForward(s *mcclient.ClientSession, forwardId string) { + _, err := cloudproxy.Forwards.Delete(s, forwardId, nil) + if err != nil { + log.Errorf("unable to delete proxy forward %s", forwardId) + } +} + +func (self *ApplyScriptTask) ensureLocalForwardWork(host string, port int) bool { + maxWaitTimes, wt := 10, 1*time.Second + waitTimes := 1 + address := fmt.Sprintf("%s:%d", host, port) + for waitTimes < maxWaitTimes { + _, err := net.DialTimeout("tcp", address, 1*time.Second) + if err == nil { + return true + } + time.Sleep(wt) + waitTimes += 1 + wt += 1 * time.Second + } + return false +} + +func mapStringSlice(f func(string) string, a []string) []string { + for i := range a { + a[i] = f(a[i]) + } + return a +} + +func (self *ApplyScriptTask) OnAnsiblePlaybookComplete(ctx context.Context, obj db.IStandaloneModel, body jsonutils.JSONObject) { + // try to delete local forward + session := auth.GetAdminSession(ctx, "", "") + forwardId, _ := self.Params.GetString("proxy_forward_id") + self.clearLocalForward(session, forwardId) + sa := obj.(*models.SScriptApply) + sarId, _ := self.Params.GetString("script_apply_record_id") + osar, err := models.ScriptApplyRecordManager.FetchById(sarId) + if err != nil { + log.Errorf("unable to fetch script apply record %s", sarId) + self.taskSuccess(ctx, sa, nil) + } + self.taskSuccess(ctx, sa, osar.(*models.SScriptApplyRecord)) +} + +func (self *ApplyScriptTask) OnAnsiblePlaybookCompleteFailed(ctx context.Context, obj db.IStandaloneModel, body jsonutils.JSONObject) { + // try to delete local forward + session := auth.GetAdminSession(ctx, "", "") + forwardId, _ := self.Params.GetString("proxy_forward_id") + _, err := cloudproxy.Forwards.Delete(session, forwardId, nil) + if err != nil { + log.Errorf("unable to delete proxy forward %s", forwardId) + } + sa := obj.(*models.SScriptApply) + sarId, _ := self.Params.GetString("script_apply_record_id") + osar, err := models.ScriptApplyRecordManager.FetchById(sarId) + if err != nil { + log.Errorf("unable to fetch script apply record %s", sarId) + self.taskSuccess(ctx, sa, nil) + } + self.taskFailed(ctx, sa, osar.(*models.SScriptApplyRecord), errors.Error(body.String())) +} diff --git a/pkg/mcclient/modules/mod_devtools.go b/pkg/mcclient/modules/mod_devtools.go index f5d16d71b9..47efe8c752 100644 --- a/pkg/mcclient/modules/mod_devtools.go +++ b/pkg/mcclient/modules/mod_devtools.go @@ -19,8 +19,10 @@ import ( ) var ( - DevToolCronjobs modulebase.ResourceManager - DevToolTemplates modulebase.ResourceManager + DevToolCronjobs modulebase.ResourceManager + DevToolTemplates modulebase.ResourceManager + DevToolScripts modulebase.ResourceManager + DevToolScriptApplyRecords modulebase.ResourceManager ) func init() { @@ -40,4 +42,19 @@ func init() { []string{"is_system"}, ) registerCompute(&DevToolTemplates) + + DevToolScripts = NewDevtoolManager( + "script", + "scripts", + []string{"Id", "Name", "Type", "Playbook_Reference", "Max_Try_Times"}, + []string{}, + ) + registerCompute(&DevToolScripts) + DevToolScriptApplyRecords = NewDevtoolManager( + "scriptapplyrecord", + "scriptapplyrecords", + []string{"Script_Id", "Server_Id", "Start_Time", "End_Time", "Reason", "Status"}, + []string{}, + ) + registerCompute(&DevToolScriptApplyRecords) } diff --git a/pkg/mcclient/options/devtool/doc.go b/pkg/mcclient/options/devtool/doc.go new file mode 100644 index 0000000000..3493988073 --- /dev/null +++ b/pkg/mcclient/options/devtool/doc.go @@ -0,0 +1 @@ +package devtool // import "yunion.io/x/onecloud/pkg/mcclient/options/devtool" diff --git a/pkg/mcclient/options/devtool/script.go b/pkg/mcclient/options/devtool/script.go new file mode 100644 index 0000000000..0911f523d4 --- /dev/null +++ b/pkg/mcclient/options/devtool/script.go @@ -0,0 +1,66 @@ +// 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 devtool + +import ( + "yunion.io/x/jsonutils" + + "yunion.io/x/onecloud/pkg/mcclient/options" +) + +type ScriptListOptions struct { + options.BaseListOptions +} + +func (so *ScriptListOptions) Params() (jsonutils.JSONObject, error) { + return options.ListStructToParams(so) +} + +type ScriptOptions struct { + ID string `help:"id or name of script"` +} + +func (so *ScriptOptions) GetId() string { + return so.ID +} + +func (so *ScriptOptions) Params() (jsonutils.JSONObject, error) { + return nil, nil +} + +type SscriptApplyOptions struct { + SERVERID string `help:"server id" json:"server_id"` + EipFirst bool `help:"whether to use eip first"` + ProxyEndpointId string `help:"proxy endpoint id"` + AutoChooseProxyEndpoint bool `help:"automatically choose proxy endpoint"` +} + +type ScriptApplyOptions struct { + ScriptOptions + SscriptApplyOptions +} + +func (so *ScriptApplyOptions) Params() (jsonutils.JSONObject, error) { + return jsonutils.Marshal(so.SscriptApplyOptions), nil +} + +type ScriptApplyRecordListOptions struct { + options.BaseListOptions + ScriptId string +} + +func (so *ScriptApplyRecordListOptions) Params() (jsonutils.JSONObject, error) { + return options.ListStructToParams(so) +} From 7e452712eef3ffcc5b55a860df6697ab931f9295 Mon Sep 17 00:00:00 2001 From: rainzm Date: Wed, 7 Apr 2021 15:20:56 +0800 Subject: [PATCH 5/5] feat(devtool): check sshable before applying ansible playbook for server --- pkg/devtool/tasks/apply_script_task.go | 50 ++++++++++++++++++++++++-- 1 file changed, 48 insertions(+), 2 deletions(-) diff --git a/pkg/devtool/tasks/apply_script_task.go b/pkg/devtool/tasks/apply_script_task.go index 31161cc98e..f56dc85dcf 100644 --- a/pkg/devtool/tasks/apply_script_task.go +++ b/pkg/devtool/tasks/apply_script_task.go @@ -18,6 +18,7 @@ import ( "context" "fmt" "net" + "strings" "time" "yunion.io/x/jsonutils" @@ -94,11 +95,25 @@ func (self *ApplyScriptTask) OnInit(ctx context.Context, obj db.IStandaloneModel self.taskFailed(ctx, sa, sar, errors.Wrapf(err, "unable to unmarshal %q to ServerDetails", data)) return } + + // check sshable + sshable, err := self.checkSshable(session, serverDetail.Id) + if err != nil { + self.taskFailed(ctx, sa, sar, err) + return + } + if !sshable.ok { + self.taskFailed(ctx, sa, sar, fmt.Errorf("server %s is not sshable: %s", serverDetail.Id, sshable.reason)) + return + } // make sure user var user string - if serverDetail.Hypervisor == comapi.HYPERVISOR_KVM { + switch { + case sshable.user != "": + user = sshable.user + case serverDetail.Hypervisor == comapi.HYPERVISOR_KVM: user = "root" - } else { + default: user = "cloudroot" } // create local forward @@ -154,6 +169,37 @@ func (self *ApplyScriptTask) OnInit(ctx context.Context, obj db.IStandaloneModel } } +type sSSHable struct { + user string + ok bool + reason string +} + +func (self *ApplyScriptTask) checkSshable(session *mcclient.ClientSession, serverId string) (sSSHable, error) { + data, err := modules.Servers.GetSpecific(session, serverId, "sshable", nil) + if err != nil { + return sSSHable{}, errors.Wrapf(err, "unable to get sshable info of server %s", serverId) + } + methodTrieds, _ := data.GetArray("method_tried") + sshable := sSSHable{} + reasons := make([]string, 0, len(methodTrieds)) + for _, methodTried := range methodTrieds { + ok, _ := methodTried.Bool("sshable") + if ok { + sshable.ok = true + break + } + reason, _ := methodTried.GetString("reason") + reasons = append(reasons, reason) + } + if !sshable.ok { + sshable.reason = strings.Join(reasons, "; ") + } else { + sshable.user, _ = data.GetString("user") + } + return sshable, nil +} + func (self *ApplyScriptTask) clearLocalForward(s *mcclient.ClientSession, forwardId string) { _, err := cloudproxy.Forwards.Delete(s, forwardId, nil) if err != nil {