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 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/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/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/ansible/devtoolcronjob.go b/cmd/climc/shell/devtool/devtoolcronjob.go similarity index 99% rename from cmd/climc/shell/ansible/devtoolcronjob.go rename to cmd/climc/shell/devtool/devtoolcronjob.go index a03fcd4d6f..38eb8dfb90 100644 --- a/cmd/climc/shell/ansible/devtoolcronjob.go +++ b/cmd/climc/shell/devtool/devtoolcronjob.go @@ -12,7 +12,7 @@ // See the License for the specific language governing permissions and // limitations under the License. -package ansible +package devtool import ( "fmt" diff --git a/cmd/climc/shell/ansible/devtooltemplate.go b/cmd/climc/shell/devtool/devtooltemplate.go similarity index 99% rename from cmd/climc/shell/ansible/devtooltemplate.go rename to cmd/climc/shell/devtool/devtooltemplate.go index 92290feac0..9c6a1d351a 100644 --- a/cmd/climc/shell/ansible/devtooltemplate.go +++ b/cmd/climc/shell/devtool/devtooltemplate.go @@ -12,7 +12,7 @@ // See the License for the specific language governing permissions and // limitations under the License. -package ansible +package devtool import ( "yunion.io/x/jsonutils" 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/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/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/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) } 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..f56dc85dcf --- /dev/null +++ b/pkg/devtool/tasks/apply_script_task.go @@ -0,0 +1,264 @@ +// 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" + "strings" + "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 + } + + // 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 + switch { + case sshable.user != "": + user = sshable.user + case serverDetail.Hypervisor == comapi.HYPERVISOR_KVM: + user = "root" + default: + 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 + } +} + +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 { + 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_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_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/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/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) +} 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) }