From a6736da35e7321bc33a9ea861c4f04d36f2c0c2e Mon Sep 17 00:00:00 2001 From: Zexi Li Date: Fri, 31 Aug 2018 19:32:53 +0800 Subject: [PATCH] webterm init --- cmd/climc/shell/k8s/pods.go | 19 +++ cmd/climc/shell/webconsole.go | 35 +++++ cmd/webconsole/main.go | 9 ++ pkg/compute/handlers.go | 2 +- pkg/mcclient/modules/mod_webconsole.go | 65 +++++++++ pkg/mcclient/options/webconsole.go | 24 ++++ pkg/webconsole/command/command.go | 46 ++++++ pkg/webconsole/command/kube_command.go | 127 +++++++++++++++++ pkg/webconsole/command/kube_command_test.go | 51 +++++++ pkg/webconsole/handlers.go | 150 ++++++++++++++++++++ pkg/webconsole/options/options.go | 16 +++ pkg/webconsole/server/server.go | 45 ++++++ pkg/webconsole/server/tty_server.go | 131 +++++++++++++++++ pkg/webconsole/service/service.go | 62 ++++++++ pkg/webconsole/session/pty_session.go | 71 +++++++++ pkg/webconsole/session/session.go | 103 ++++++++++++++ 16 files changed, 955 insertions(+), 1 deletion(-) create mode 100644 cmd/climc/shell/webconsole.go create mode 100644 cmd/webconsole/main.go create mode 100644 pkg/mcclient/modules/mod_webconsole.go create mode 100644 pkg/mcclient/options/webconsole.go create mode 100644 pkg/webconsole/command/command.go create mode 100644 pkg/webconsole/command/kube_command.go create mode 100644 pkg/webconsole/command/kube_command_test.go create mode 100644 pkg/webconsole/handlers.go create mode 100644 pkg/webconsole/options/options.go create mode 100644 pkg/webconsole/server/server.go create mode 100644 pkg/webconsole/server/tty_server.go create mode 100644 pkg/webconsole/service/service.go create mode 100644 pkg/webconsole/session/pty_session.go create mode 100644 pkg/webconsole/session/session.go diff --git a/cmd/climc/shell/k8s/pods.go b/cmd/climc/shell/k8s/pods.go index 5376147feb..5fc854b7c9 100644 --- a/cmd/climc/shell/k8s/pods.go +++ b/cmd/climc/shell/k8s/pods.go @@ -3,6 +3,8 @@ package k8s import ( "fmt" + "yunion.io/x/jsonutils" + "yunion.io/x/onecloud/pkg/mcclient" "yunion.io/x/onecloud/pkg/mcclient/modules/k8s" ) @@ -28,6 +30,23 @@ func initPod() { return nil }) + type getOpt struct { + resourceGetOptions + } + R(&getOpt{}, cmdN("show"), "Get pod details", func(s *mcclient.ClientSession, args *getOpt) error { + id := args.NAME + params := args.ClusterParams() + if args.Namespace != "" { + params.Add(jsonutils.NewString(args.Namespace), "namespace") + } + ret, err := k8s.Pods.Get(s, id, params) + if err != nil { + return err + } + printObjectYAML(ret) + return nil + }) + type deleteOpt struct { resourceGetOptions } diff --git a/cmd/climc/shell/webconsole.go b/cmd/climc/shell/webconsole.go new file mode 100644 index 0000000000..1d4a21289d --- /dev/null +++ b/cmd/climc/shell/webconsole.go @@ -0,0 +1,35 @@ +package shell + +import ( + "yunion.io/x/onecloud/pkg/mcclient" + "yunion.io/x/onecloud/pkg/mcclient/modules" + o "yunion.io/x/onecloud/pkg/mcclient/options" +) + +func init() { + R(&o.PodShellOptions{}, "webconsole-k8s-pod", "Show TTY console of given pod", func(s *mcclient.ClientSession, args *o.PodShellOptions) error { + params, err := args.Params() + if err != nil { + return err + } + ret, err := modules.WebConsole.DoK8sShellConnect(s, args.NAME, params) + if err != nil { + return err + } + printObject(ret) + return nil + }) + + R(&o.PodLogOptoins{}, "webconsole-k8s-pod-log", "Get logs of given pod", func(s *mcclient.ClientSession, args *o.PodLogOptoins) error { + params, err := args.Params() + if err != nil { + return err + } + ret, err := modules.WebConsole.DoK8sLogConnect(s, args.NAME, params) + if err != nil { + return err + } + printObject(ret) + return nil + }) +} diff --git a/cmd/webconsole/main.go b/cmd/webconsole/main.go new file mode 100644 index 0000000000..46c760e115 --- /dev/null +++ b/cmd/webconsole/main.go @@ -0,0 +1,9 @@ +package main + +import ( + "yunion.io/x/onecloud/pkg/webconsole/service" +) + +func main() { + service.StartService() +} diff --git a/pkg/compute/handlers.go b/pkg/compute/handlers.go index 2cf6dc6115..439e415b09 100644 --- a/pkg/compute/handlers.go +++ b/pkg/compute/handlers.go @@ -2,9 +2,9 @@ package compute import ( "yunion.io/x/log" + "yunion.io/x/onecloud/pkg/appsrv" "yunion.io/x/onecloud/pkg/appsrv/dispatcher" - "yunion.io/x/onecloud/pkg/cloudcommon/db" "yunion.io/x/onecloud/pkg/cloudcommon/db/quotas" "yunion.io/x/onecloud/pkg/cloudcommon/db/taskman" diff --git a/pkg/mcclient/modules/mod_webconsole.go b/pkg/mcclient/modules/mod_webconsole.go new file mode 100644 index 0000000000..6ff9db80ab --- /dev/null +++ b/pkg/mcclient/modules/mod_webconsole.go @@ -0,0 +1,65 @@ +package modules + +import ( + "fmt" + + "yunion.io/x/jsonutils" + "yunion.io/x/onecloud/pkg/mcclient" +) + +var ( + WebConsole WebConsoleManager +) + +func init() { + WebConsole = WebConsoleManager{NewWebConsoleManager()} +} + +type WebConsoleManager struct { + ResourceManager +} + +func NewWebConsoleManager() ResourceManager { + return ResourceManager{BaseManager: BaseManager{serviceType: "webconsole"}, + Keyword: "webconsole", KeywordPlural: "webconsole"} +} + +func (m WebConsoleManager) DoConnect( + s *mcclient.ClientSession, + connType, id, action string, + params jsonutils.JSONObject, +) (jsonutils.JSONObject, error) { + if len(connType) == 0 { + return nil, fmt.Errorf("Empty connection resource type") + } + url := fmt.Sprintf("/webconsole/%s", connType) + if id != "" { + url = fmt.Sprintf("%s/%s", url, id) + } + if action != "" { + url = fmt.Sprintf("%s/%s", url, action) + } + return m._post(s, url, params, "webconsole") +} + +func (m WebConsoleManager) DoK8sConnect( + s *mcclient.ClientSession, + id, action string, + params jsonutils.JSONObject, +) (jsonutils.JSONObject, error) { + return m.DoConnect(s, "k8s", id, action, params) +} + +func (m WebConsoleManager) DoK8sShellConnect( + s *mcclient.ClientSession, + id string, params jsonutils.JSONObject, +) (jsonutils.JSONObject, error) { + return m.DoK8sConnect(s, id, "shell", params) +} + +func (m WebConsoleManager) DoK8sLogConnect( + s *mcclient.ClientSession, + id string, params jsonutils.JSONObject, +) (jsonutils.JSONObject, error) { + return m.DoK8sConnect(s, id, "log", params) +} diff --git a/pkg/mcclient/options/webconsole.go b/pkg/mcclient/options/webconsole.go new file mode 100644 index 0000000000..ad339fa6a8 --- /dev/null +++ b/pkg/mcclient/options/webconsole.go @@ -0,0 +1,24 @@ +package options + +import ( + "yunion.io/x/jsonutils" +) + +type PodBaseOptions struct { + NAME string `help:"Name of k8s pod to connect"` + Namespace string `help:"Namespace of this pod"` + Container string `help:"Container in this pod"` + Cluster string `default:"$K8S_CLUSTER|default" help:"Kubernetes cluster name"` +} + +func (opt *PodBaseOptions) Params() (*jsonutils.JSONDict, error) { + return StructToParams(opt) +} + +type PodShellOptions struct { + PodBaseOptions +} + +type PodLogOptoins struct { + PodBaseOptions +} diff --git a/pkg/webconsole/command/command.go b/pkg/webconsole/command/command.go new file mode 100644 index 0000000000..35ed12f53d --- /dev/null +++ b/pkg/webconsole/command/command.go @@ -0,0 +1,46 @@ +package command + +import ( + "os/exec" + + "yunion.io/x/log" +) + +const ( + PROTOCOL_TTY string = "tty" + //PROTOCOL_VNC string = "vnc" +) + +type ICommand interface { + GetProtocol() string + GetCommand() *exec.Cmd + Cleanup() error +} + +type BaseCommand struct { + name string + args []string +} + +func NewBaseCommand(name string, args ...string) *BaseCommand { + return &BaseCommand{ + name: name, + args: args, + } +} + +func (c *BaseCommand) AppendArgs(args ...string) *BaseCommand { + for _, arg := range args { + c.args = append(c.args, arg) + } + return c +} + +func (c BaseCommand) GetCommand() *exec.Cmd { + return exec.Command(c.name, c.args...) +} + +func (c BaseCommand) Cleanup() error { + log.Infof("BaseCommand Cleanup do nothing") + return nil +} diff --git a/pkg/webconsole/command/kube_command.go b/pkg/webconsole/command/kube_command.go new file mode 100644 index 0000000000..34a9087a6d --- /dev/null +++ b/pkg/webconsole/command/kube_command.go @@ -0,0 +1,127 @@ +package command + +import ( + "fmt" + "os" + "os/exec" + + "yunion.io/x/log" +) + +type Kubectl struct { + *BaseCommand + kubeconfig string +} + +func NewKubectlCommand(kubeconfig, namespace string) *Kubectl { + name := "kubectl" + if len(namespace) == 0 { + namespace = "default" + } + cmd := NewBaseCommand(name, "--namespace", namespace) + return &Kubectl{ + BaseCommand: cmd, + kubeconfig: kubeconfig, + } +} + +func (c *Kubectl) GetCommand() *exec.Cmd { + cmd := c.BaseCommand.GetCommand() + cmd.Env = append(cmd.Env, fmt.Sprintf("KUBECONFIG=%s", c.kubeconfig)) + return cmd +} + +func (c Kubectl) GetProtocol() string { + return PROTOCOL_TTY +} + +func (c *Kubectl) Cleanup() error { + log.Debugf("Remove temp kubeconfig file: %s", c.kubeconfig) + return os.Remove(c.kubeconfig) +} + +type KubectlExec struct { + *Kubectl +} + +func (c *Kubectl) Exec() *KubectlExec { + // Execute a command in a container + cmd := &KubectlExec{ + Kubectl: c, + } + cmd.AppendArgs("exec") + return cmd +} + +func (c *KubectlExec) Stdin() *KubectlExec { + // -i: Pass stdin to the container + c.AppendArgs("-i") + return c +} + +func (c *KubectlExec) TTY() *KubectlExec { + // -t: Stdin is a TTY + c.AppendArgs("-t") + return c +} + +func (c *KubectlExec) Container(name string) *KubectlExec { + if len(name) == 0 { + return c + } + // -c: Container name. If ommitted, the first container in the pod will be chosen + c.AppendArgs("-c", name) + return c +} + +func (c *KubectlExec) Pod(name string) *KubectlExec { + // Pod name + c.AppendArgs(name) + return c +} + +func (c *KubectlExec) Command(cmd string, args ...string) *KubectlExec { + c.AppendArgs("--", cmd) + c.AppendArgs(args...) + return c +} + +func NewPodBashCommand(kubeconfig, namespace, pod, container string) ICommand { + return NewKubectlCommand(kubeconfig, namespace).Exec(). + Stdin(). + TTY(). + Pod(pod). + Container(container). + Command("bash", "-i", "-l") +} + +type KubectlLog struct { + *Kubectl +} + +func (c *Kubectl) Logs() *KubectlLog { + // Print the logs for a container in a pod + cmd := &KubectlLog{ + Kubectl: c, + } + cmd.AppendArgs("logs") + return cmd +} + +func (c *KubectlLog) Follow() *KubectlLog { + // -f: Specify if the logs should be streamed + c.AppendArgs("-f") + return c +} + +func (c *KubectlLog) Pod(name string) *KubectlLog { + // Pod name + c.AppendArgs(name) + return c +} + +func NewPodLogCommand(kubeconfig, namespace, pod, container string) ICommand { + return NewKubectlCommand(kubeconfig, namespace).Logs(). + Follow(). + Pod(pod) +} diff --git a/pkg/webconsole/command/kube_command_test.go b/pkg/webconsole/command/kube_command_test.go new file mode 100644 index 0000000000..eda46f031b --- /dev/null +++ b/pkg/webconsole/command/kube_command_test.go @@ -0,0 +1,51 @@ +package command + +import ( + "reflect" + "strings" + "testing" +) + +func TestKubectlExec_Command(t *testing.T) { + type fields struct { + Kubectl *Kubectl + } + type args struct { + cmd string + args []string + } + tests := []struct { + name string + fields fields + args args + want string + }{ + { + name: "bash command", + fields: fields{ + Kubectl: NewKubectlCommand("system"), + }, + args: args{ + cmd: "bash", + args: []string{"-il"}, + }, + want: "kubectl --namespace system exec -i -t Pod1 -c Container1 -- bash -il", + }, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + c := tt.fields.Kubectl.Exec(). + Stdin().TTY(). + Pod("Pod1"). + Container("Container1") + cmd := c.Command(tt.args.cmd, tt.args.args...) + name := cmd.name + args := cmd.args + got := []string{name} + got = append(got, args...) + if !reflect.DeepEqual(strings.Join(got, " "), tt.want) { + t.Errorf("KubectlExec.Command() = %v, want %v", got, tt.want) + } + }) + } +} diff --git a/pkg/webconsole/handlers.go b/pkg/webconsole/handlers.go new file mode 100644 index 0000000000..bc89034c0d --- /dev/null +++ b/pkg/webconsole/handlers.go @@ -0,0 +1,150 @@ +package webconsole + +import ( + "context" + "fmt" + "io/ioutil" + "net/http" + + "yunion.io/x/jsonutils" + "yunion.io/x/log" + "yunion.io/x/pkg/util/sets" + + "yunion.io/x/onecloud/pkg/appctx" + "yunion.io/x/onecloud/pkg/appsrv" + "yunion.io/x/onecloud/pkg/httperrors" + "yunion.io/x/onecloud/pkg/mcclient/auth" + "yunion.io/x/onecloud/pkg/mcclient/modules/k8s" + "yunion.io/x/onecloud/pkg/webconsole/command" + o "yunion.io/x/onecloud/pkg/webconsole/options" + "yunion.io/x/onecloud/pkg/webconsole/session" +) + +const ( + ApiPathPrefix = "/webconsole/" + ConnectPathPrefix = "/connect/" +) + +func InitHandlers(app *appsrv.Application) { + app.AddHandler("POST", ApiPathPrefix+"k8s//shell", auth.Authenticate(handleK8sShell)) + app.AddHandler("POST", ApiPathPrefix+"k8s//log", auth.Authenticate(handleK8sLog)) +} + +func fetchEnv(ctx context.Context, w http.ResponseWriter, r *http.Request) (map[string]string, jsonutils.JSONObject, jsonutils.JSONObject) { + params := appctx.AppContextParams(ctx) + query, e := jsonutils.ParseQueryString(r.URL.RawQuery) + if e != nil { + log.Errorf("Parse query string %q failed: %v", r.URL.RawQuery, e) + } + var body jsonutils.JSONObject = nil + if sets.NewString("PUT", "POST", "DELETE", "PATCH").Has(r.Method) { + body, e = appsrv.FetchJSON(r) + if e != nil { + log.Errorf("Failed to decode JSON request body: %v", e) + } + } + return params, query, body +} + +type K8sEnv struct { + Cluster string + Namespace string + Pod string + Container string + KubeConfig string +} + +func fetchK8sEnv(ctx context.Context, w http.ResponseWriter, r *http.Request) (*K8sEnv, error) { + params, _, body := fetchEnv(ctx, w, r) + cluster, _ := body.GetString("cluster") + if cluster == "" { + cluster = "default" + } + namespace, _ := body.GetString("namespace") + if namespace == "" { + namespace = "default" + } + podName := params[""] + container, _ := body.GetString("container") + adminSession := auth.GetAdminSession(o.Options.Region, "") + + query := jsonutils.NewDict() + query.Add(jsonutils.NewString(namespace), "namespace") + query.Add(jsonutils.NewString(cluster), "cluster") + obj, err := k8s.Pods.Get(adminSession, podName, query) + if err != nil { + return nil, err + } + if obj == nil { + return nil, httperrors.NewNotFoundError("Not found pod %q", podName) + } + + data := jsonutils.NewDict() + data.Add(jsonutils.JSONTrue, "directly") + ret, err := k8s.Clusters.PerformAction(adminSession, cluster, "generate-kubeconfig", data) + if err != nil { + return nil, err + } + conf, err := ret.GetString("kubeconfig") + if err != nil { + return nil, httperrors.NewNotFoundError("Not found cluster %q kubeconfig", cluster) + } + f, err := ioutil.TempFile("", "kubeconfig-") + if err != nil { + return nil, fmt.Errorf("Save kubeconfig error: %v", err) + } + f.WriteString(conf) + defer f.Close() + + return &K8sEnv{ + Cluster: cluster, + Namespace: namespace, + Pod: podName, + Container: container, + KubeConfig: f.Name(), + }, nil +} + +func handleK8sCommand( + ctx context.Context, + w http.ResponseWriter, + r *http.Request, + cmdFactory func(kubeconfig, namespace, pod, container string) command.ICommand, +) { + env, err := fetchK8sEnv(ctx, w, r) + if err != nil { + httperrors.GeneralServerError(w, err) + return + } + + cmd := cmdFactory(env.KubeConfig, env.Namespace, env.Pod, env.Container) + cmdSession, err := session.Manager.Save(cmd) + if err != nil { + httperrors.GeneralServerError(w, err) + return + } + data := jsonutils.NewDict() + url, err := cmdSession.GetConnectUrl() + if err != nil { + httperrors.GeneralServerError(w, err) + return + } + data.Add(jsonutils.NewString(url), "url") + sendJSON(w, data) +} + +func handleK8sShell(ctx context.Context, w http.ResponseWriter, r *http.Request) { + handleK8sCommand(ctx, w, r, command.NewPodBashCommand) +} + +func handleK8sLog(ctx context.Context, w http.ResponseWriter, r *http.Request) { + handleK8sCommand(ctx, w, r, command.NewPodLogCommand) +} + +func sendJSON(w http.ResponseWriter, body jsonutils.JSONObject) { + ret := jsonutils.NewDict() + if body != nil { + ret.Add(body, "webconsole") + } + appsrv.SendJSON(w, ret) +} diff --git a/pkg/webconsole/options/options.go b/pkg/webconsole/options/options.go new file mode 100644 index 0000000000..eeaab8bf9b --- /dev/null +++ b/pkg/webconsole/options/options.go @@ -0,0 +1,16 @@ +package options + +import ( + "yunion.io/x/onecloud/pkg/cloudcommon" +) + +var ( + Options WebConsoleOptions +) + +type WebConsoleOptions struct { + cloudcommon.Options + + FrontendUrl string `help:"Frontend url to display web console page" default:"http://127.0.0.1:8899"` + TtyStaticPath string `help:"TTY static HTML render pages" default:"./tty"` +} diff --git a/pkg/webconsole/server/server.go b/pkg/webconsole/server/server.go new file mode 100644 index 0000000000..04796caf70 --- /dev/null +++ b/pkg/webconsole/server/server.go @@ -0,0 +1,45 @@ +package server + +import ( + "fmt" + "net/http" + + "yunion.io/x/jsonutils" + "yunion.io/x/log" + + "yunion.io/x/onecloud/pkg/httperrors" + "yunion.io/x/onecloud/pkg/webconsole/session" +) + +type ConnectionServer struct { +} + +func NewConnectionServer() *ConnectionServer { + return &ConnectionServer{} +} + +func (s *ConnectionServer) ServeHTTP(w http.ResponseWriter, req *http.Request) { + query, err := jsonutils.ParseQueryString(req.URL.RawQuery) + if err != nil { + httperrors.GeneralServerError(w, err) + return + } + log.Debugf("[connection] Get query: %v", query) + accessToken, _ := query.GetString("access_token") + if accessToken == "" { + httperrors.BadRequestError(w, fmt.Sprintf("Empty access_token")) + return + } + sessionObj, ok := session.Manager.Get(accessToken) + if !ok { + log.Warningf("Not found session by token: %q", accessToken) + httperrors.NotFoundError(w, fmt.Sprintf("Not found session")) + return + } + ttyServer, err := NewTTYServer(sessionObj) + if err != nil { + httperrors.GeneralServerError(w, fmt.Errorf("New TTY error: %v", err)) + return + } + ttyServer.ServeHTTP(w, req) +} diff --git a/pkg/webconsole/server/tty_server.go b/pkg/webconsole/server/tty_server.go new file mode 100644 index 0000000000..e9b8f6c9d5 --- /dev/null +++ b/pkg/webconsole/server/tty_server.go @@ -0,0 +1,131 @@ +package server + +import ( + "os/exec" + "strings" + + socketio "github.com/googollee/go-socket.io" + "github.com/kr/pty" + + "yunion.io/x/jsonutils" + "yunion.io/x/log" + + "yunion.io/x/onecloud/pkg/webconsole/session" +) + +const ( + ON_CONNECTION = "connection" + ON_DISCONNECTION = "disconnection" + ON_ERROR = "error" + + OUTPUT_EVENT = "output" + INPUT_EVENT = "input" + RESIZE_EVENT = "resize" + + COMMAND_QUERY = "command" + ARGS_QUERY = "args" +) + +type TTYServer struct { + *socketio.Server +} + +func NewTTYServer(s *session.SSession) (*TTYServer, error) { + socketioServer, err := socketio.NewServer(nil) + if err != nil { + return nil, err + } + server := &TTYServer{ + Server: socketioServer, + } + server.initEventHandler(s) + return server, nil +} + +func fetchCommand(so socketio.Socket) (*exec.Cmd, error) { + req := so.Request() + query, err := jsonutils.ParseQueryString(req.URL.RawQuery) + if err != nil { + return nil, err + } + cmd, err := query.GetString(COMMAND_QUERY) + if err != nil { + return nil, err + } + args := []string{} + argsStr, _ := query.GetString(ARGS_QUERY) + if len(argsStr) != 0 { + args = strings.Split(argsStr, ",") + } + return exec.Command(cmd, args...), nil +} + +func (server *TTYServer) initEventHandler(s *session.SSession) { + server.On(ON_CONNECTION, func(so socketio.Socket) error { + log.Infof("[%q] On connection", so.Id()) + p, err := session.NewPty(s) + if err != nil { + log.Errorf("Create Pty error: %v", err) + return err + } + initSocketHandler(so, p) + return nil + }) +} + +func initSocketHandler(so socketio.Socket, p *session.Pty) { + // handle read + go func() { + buf := make([]byte, 1024) + for { + n, err := p.Pty.Read(buf) + if err != nil { + log.Errorf("Failed to read from pty master: %v", err) + cleanUp(so, p) + return + } + so.Emit(OUTPUT_EVENT, string(buf[0:n])) + } + }() + + // handle write + so.On(INPUT_EVENT, func(data string) { + p.Pty.Write([]byte(data)) + }) + + // handle resize + so.On(RESIZE_EVENT, func(colRow []uint16) { + if len(colRow) != 2 { + log.Errorf("Invalid window size: %v", colRow) + cleanUp(so, p) + return + } + //size, err := pty.GetsizeFull(p.Pty) + //if err != nil { + //log.Errorf("Get pty window size error: %v", err) + //return + //} + newSize := pty.Winsize{ + Cols: colRow[0], + Rows: colRow[1], + } + p.Resize(&newSize) + }) + + // handle disconnection + so.On(ON_DISCONNECTION, func(msg string) { + log.Infof("[%s] closed: %s", so.Id(), msg) + cleanUp(so, p) + }) + + // handle error + so.On(ON_ERROR, func(err error) { + log.Errorf("[%s] on error: %v", so.Id(), err) + cleanUp(so, p) + }) +} + +func cleanUp(so socketio.Socket, p *session.Pty) { + so.Disconnect() + p.Stop() +} diff --git a/pkg/webconsole/service/service.go b/pkg/webconsole/service/service.go new file mode 100644 index 0000000000..9bf7667f03 --- /dev/null +++ b/pkg/webconsole/service/service.go @@ -0,0 +1,62 @@ +package service + +import ( + "fmt" + "net" + "net/http" + "net/url" + "os" + "strconv" + + "github.com/gorilla/mux" + + "yunion.io/x/log" + + "yunion.io/x/onecloud/pkg/cloudcommon" + "yunion.io/x/onecloud/pkg/webconsole" + "yunion.io/x/onecloud/pkg/webconsole/command" + o "yunion.io/x/onecloud/pkg/webconsole/options" + "yunion.io/x/onecloud/pkg/webconsole/server" +) + +func StartService() { + cloudcommon.ParseOptions(&o.Options, &o.Options.Options, os.Args, "webconsole.conf") + + if o.Options.FrontendUrl == "" { + log.Fatalf("--frontend-url must specified") + } + _, err := url.Parse(o.Options.FrontendUrl) + if err != nil { + log.Fatalf("invalid --frontend-url %s", o.Options.FrontendUrl) + } + + cloudcommon.InitAuth(&o.Options.Options, func() { + log.Infof("Auth complete") + }) + start() +} + +func start() { + app := cloudcommon.InitApp(&o.Options.Options) + webconsole.InitHandlers(app) + + root := mux.NewRouter() + root.UseEncodedPath() + + // api handler + root.PathPrefix(webconsole.ApiPathPrefix).Handler(app) + + // websocket related console handler + root.Handle(webconsole.ConnectPathPrefix, server.NewConnectionServer()) + + p1 := fmt.Sprintf("/%s/", command.PROTOCOL_TTY) + // static file handler + root.PathPrefix(p1).Handler(http.FileServer(http.Dir(o.Options.TtyStaticPath))) + + addr := net.JoinHostPort(o.Options.Address, strconv.Itoa(o.Options.Port)) + log.Infof("Start listen on %s", addr) + err := http.ListenAndServe(addr, root) + if err != nil { + log.Fatalf("%v", err) + } +} diff --git a/pkg/webconsole/session/pty_session.go b/pkg/webconsole/session/pty_session.go new file mode 100644 index 0000000000..c77f7f5600 --- /dev/null +++ b/pkg/webconsole/session/pty_session.go @@ -0,0 +1,71 @@ +package session + +import ( + "os" + "os/exec" + "os/signal" + "syscall" + + "github.com/kr/pty" + + "yunion.io/x/log" +) + +type Pty struct { + Session *SSession + Cmd *exec.Cmd + Pty *os.File + sizeCh chan os.Signal + size *pty.Winsize +} + +func NewPty(session *SSession) (p *Pty, err error) { + cmd := session.GetCommand() + p = &Pty{ + Session: session, + Cmd: cmd, + } + p.Pty, err = pty.Start(p.Cmd) + if err != nil { + return + } + p.sizeCh = make(chan os.Signal, 1) + p.size = &pty.Winsize{} + p.startResizeMonitor() + signal.Notify(p.sizeCh, syscall.SIGWINCH) // Initail resize + return +} + +func (p *Pty) startResizeMonitor() { + go func() { + for range p.sizeCh { + if err := pty.Setsize(p.Pty, p.size); err != nil { + log.Errorf("Resize pty error: %v", err) + } else { + log.Debugf("Resize pty to %#v, cmd: %#v", p.size, p.Cmd) + } + } + }() +} + +func (p *Pty) Resize(size *pty.Winsize) { + p.size = size + p.sizeCh <- syscall.SIGWINCH +} + +func (p *Pty) Stop() { + var err error + err = p.Pty.Close() + if err != nil { + log.Errorf("Close PTY error: %v", err) + } + err = p.Cmd.Process.Signal(os.Kill) + if err != nil { + log.Errorf("Kill command process error: %v", err) + } + err = p.Cmd.Wait() + if err != nil { + log.Errorf("Wait command error: %v", err) + } + p.Session.Close() +} diff --git a/pkg/webconsole/session/session.go b/pkg/webconsole/session/session.go new file mode 100644 index 0000000000..c13f352ec8 --- /dev/null +++ b/pkg/webconsole/session/session.go @@ -0,0 +1,103 @@ +package session + +import ( + "fmt" + "math/rand" + "net/url" + "sync" + + "github.com/golang-plus/uuid" + + "yunion.io/x/log" + "yunion.io/x/pkg/utils" + + "yunion.io/x/onecloud/pkg/mcclient/auth" + "yunion.io/x/onecloud/pkg/webconsole/command" + o "yunion.io/x/onecloud/pkg/webconsole/options" +) + +var ( + Manager *SSessionManager + AES_KEY string +) + +func init() { + Manager = NewSessionManager() + AES_KEY = fmt.Sprintf("webconsole-%f", rand.Float32()) +} + +type SSessionManager struct { + *sync.Map +} + +func NewSessionManager() *SSessionManager { + s := &SSessionManager{ + Map: &sync.Map{}, + } + return s +} + +func (man *SSessionManager) Save(Command command.ICommand) (session *SSession, err error) { + key, err := uuid.NewV4() + if err != nil { + return + } + idStr := key.String() + token, err := utils.EncryptAESBase64Url(AES_KEY, idStr) + if err != nil { + return + } + session = &SSession{ + Id: idStr, + ICommand: Command, + AccessToken: token, + } + man.Store(idStr, session) + return +} + +func (man *SSessionManager) Get(accessToken string) (*SSession, bool) { + id, err := utils.DescryptAESBase64Url(AES_KEY, accessToken) + if err != nil { + log.Errorf("DescryptAESBase64Url error: %v", err) + return nil, false + } + obj, ok := man.Load(id) + if !ok { + return nil, false + } + return obj.(*SSession), true +} + +type SSession struct { + command.ICommand + Id string + AccessToken string +} + +func (s SSession) GetConnectUrl() (string, error) { + FrontendUrl := o.Options.FrontendUrl + endpointUrl, err := auth.AdminCredential().GetServiceURL("webconsole", o.Options.Region, "", "internalURL") + if err != nil { + return "", err + } + if endpointUrl == "" { + return "", fmt.Errorf("Not found service URL for webconsole") + } + u, _ := url.Parse(FrontendUrl) + params := url.Values{ + "api_server": {endpointUrl}, + "access_token": {s.AccessToken}, + } + u.Path = fmt.Sprintf("%s/", s.GetProtocol()) + u.RawQuery = params.Encode() + return u.String(), nil +} + +func (s *SSession) Close() error { + if err := s.ICommand.Cleanup(); err != nil { + log.Errorf("Clean up command error: %v", err) + } + Manager.Delete(s.Id) + return nil +}