Merge pull request #9824 from rainzm/ab/template

Install monitor agent for VM
This commit is contained in:
Zexi Li
2021-04-07 21:22:00 +08:00
committed by GitHub
40 changed files with 2459 additions and 167 deletions
+4
View File
@@ -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
+22
View File
@@ -0,0 +1,22 @@
FROM registry.cn-beijing.aliyuncs.com/yunionio/onecloud-base:v0.2
MAINTAINER "Rain Zheng <zhengyu@yunion.com>"
# 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
+1
View File
@@ -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"
@@ -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))
}
+27
View File
@@ -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
)
@@ -12,7 +12,7 @@
// See the License for the specific language governing permissions and
// limitations under the License.
package ansible
package devtool
import (
"fmt"
@@ -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"
+29
View File
@@ -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))
}
@@ -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)
}
@@ -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)
}
+2 -1
View File
@@ -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 (
+3
View File
@@ -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)
+1
View File
@@ -41,6 +41,7 @@ func StartService() {
dbOpts := &opts.DBOptions
baseOpts := &opts.BaseOptions
models.InitPlaybookWorker()
app := common_app.InitApp(baseOpts, false)
InitHandlers(app)
+37
View File
@@ -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
}
+22
View File
@@ -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"
)
+63
View File
@@ -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"`
}
+13
View File
@@ -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"
@@ -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"`
}
+15
View File
@@ -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"
+67
View File
@@ -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
}
+30
View File
@@ -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"
)
+78
View File
@@ -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"`
}
+11
View File
@@ -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
}
+9 -1
View File
@@ -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)
}
+345
View File
@@ -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
}
+178
View File
@@ -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)
}
+143
View File
@@ -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)
}
+4
View File
@@ -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)
+2 -1
View File
@@ -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!!")
+264
View File
@@ -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()))
}
@@ -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)
}
+19 -2
View File
@@ -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)
}
+6
View File
@@ -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 {
@@ -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
}
+15
View File
@@ -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"
+1
View File
@@ -0,0 +1 @@
package devtool // import "yunion.io/x/onecloud/pkg/mcclient/options/devtool"
+66
View File
@@ -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)
}
+81
View File
@@ -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)
}
+278
View File
@@ -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
}
+15 -160
View File
@@ -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)
}