From 93a42678f8c6bb472b7ed67375757e37a4fa91b4 Mon Sep 17 00:00:00 2001 From: yunlzheng Date: Fri, 6 Mar 2020 18:36:34 +0800 Subject: [PATCH 01/12] add cluster interface --- pkg/kt/cluster/types.go | 9 +++++++++ pkg/kt/command/connect.go | 36 +++++++++++++++++------------------- pkg/kt/connect/outbound.go | 16 +++++++++------- 3 files changed, 35 insertions(+), 26 deletions(-) create mode 100644 pkg/kt/cluster/types.go diff --git a/pkg/kt/cluster/types.go b/pkg/kt/cluster/types.go new file mode 100644 index 0000000..9536d55 --- /dev/null +++ b/pkg/kt/cluster/types.go @@ -0,0 +1,9 @@ +package cluster + +// KubernetesInterface kubernetes interface +type KubernetesInterface interface { +} + +// Kubernetes implements KubernetesInterface +type Kubernetes struct { +} diff --git a/pkg/kt/command/connect.go b/pkg/kt/command/connect.go index 5433f35..e8ee10c 100644 --- a/pkg/kt/command/connect.go +++ b/pkg/kt/command/connect.go @@ -4,12 +4,12 @@ import ( "fmt" "strings" + "github.com/alibaba/kt-connect/pkg/kt/connect" "github.com/rs/zerolog" "github.com/rs/zerolog/log" "github.com/urfave/cli" "github.com/alibaba/kt-connect/pkg/kt/cluster" - "github.com/alibaba/kt-connect/pkg/kt/connect" "github.com/alibaba/kt-connect/pkg/kt/options" "github.com/alibaba/kt-connect/pkg/kt/util" ) @@ -83,8 +83,9 @@ func (action *Action) Connect(options *options.DaemonOptions) (err error) { if err != nil { return } + log.Info().Msgf("Connect Start At %d", pid) - factory := connect.Connect{Options: options} + clientSet, err := cluster.GetKubernetesClient(options.KubeConfig) if err != nil { return @@ -99,23 +100,8 @@ func (action *Action) Connect(options *options.DaemonOptions) (err error) { workload := fmt.Sprintf("kt-connect-daemon-%s", strings.ToLower(util.RandomString(5))) options.RuntimeOptions.Shadow = workload - labels := map[string]string{ - "kt": workload, - "kt-component": "connect", - "control-by": "kt", - } - - for k, v := range util.String2Map(options.Labels) { - labels[k] = v - } - endPointIP, podName, err := cluster.CreateShadow( - clientSet, - workload, - labels, - options.Namespace, - options.Image, - ) + clientSet, workload, labels(workload, options), options.Namespace, options.Image) if err != nil { return @@ -126,7 +112,7 @@ func (action *Action) Connect(options *options.DaemonOptions) (err error) { return } - err = factory.StartConnect(podName, endPointIP, cidrs, options.Debug) + err = connect.StartConnect(podName, endPointIP, cidrs, options) if err != nil { return } @@ -135,3 +121,15 @@ func (action *Action) Connect(options *options.DaemonOptions) (err error) { log.Info().Msgf("Terminal Signal is %s", s) return } + +func labels(workload string, options *options.DaemonOptions) map[string]string { + labels := map[string]string{ + "kt": workload, + "kt-component": "connect", + "control-by": "kt", + } + for k, v := range util.String2Map(options.Labels) { + labels[k] = v + } + return labels +} diff --git a/pkg/kt/connect/outbound.go b/pkg/kt/connect/outbound.go index 58f0128..8eb56dd 100644 --- a/pkg/kt/connect/outbound.go +++ b/pkg/kt/connect/outbound.go @@ -5,6 +5,8 @@ import ( "io/ioutil" "time" + "github.com/alibaba/kt-connect/pkg/kt/options" + "github.com/alibaba/kt-connect/pkg/kt/exec" "github.com/alibaba/kt-connect/pkg/kt/exec/kubectl" "github.com/alibaba/kt-connect/pkg/kt/exec/ssh" @@ -14,24 +16,24 @@ import ( ) // StartConnect start vpn connection -func (c *Connect) StartConnect(name, podIP string, cidrs []string, debug bool) (err error) { +func StartConnect(name, podIP string, cidrs []string, options *options.DaemonOptions) (err error) { err = util.PrepareSSHPrivateKey() if err != nil { return } - err = exec.BackgroundRun(kubectl.PortForward(c.Options.KubeConfig, c.Options.Namespace, name, c.Options.ConnectOptions.SSHPort), "port-forward", debug) + err = exec.BackgroundRun(kubectl.PortForward(options.KubeConfig, options.Namespace, name, options.ConnectOptions.SSHPort), "port-forward", options.Debug) if err != nil { return } time.Sleep(time.Duration(5) * time.Second) - if c.Options.ConnectOptions.Method == "socks5" { + if options.ConnectOptions.Method == "socks5" { log.Info().Msgf("==============================================================") - log.Info().Msgf("Start SOCKS5 Proxy: export http_proxy=socks5://127.0.0.1:%d", c.Options.ConnectOptions.Socke5Proxy) + log.Info().Msgf("Start SOCKS5 Proxy: export http_proxy=socks5://127.0.0.1:%d", options.ConnectOptions.Socke5Proxy) log.Info().Msgf("==============================================================") - _ = ioutil.WriteFile(".jvmrc", []byte(fmt.Sprintf("-DsocksProxyHost=127.0.0.1\n-DsocksProxyPort=%d", c.Options.ConnectOptions.Socke5Proxy)), 0644) - err = exec.BackgroundRun(ssh.DynamicForwardLocalRequestToRemote("127.0.0.1", c.Options.ConnectOptions.SSHPort, c.Options.ConnectOptions.Socke5Proxy), "vpn(ssh)", debug) + _ = ioutil.WriteFile(".jvmrc", []byte(fmt.Sprintf("-DsocksProxyHost=127.0.0.1\n-DsocksProxyPort=%d", options.ConnectOptions.Socke5Proxy)), 0644) + err = exec.BackgroundRun(ssh.DynamicForwardLocalRequestToRemote("127.0.0.1", options.ConnectOptions.SSHPort, options.ConnectOptions.Socke5Proxy), "vpn(ssh)", options.Debug) } else { - err = exec.BackgroundRun(sshuttle.SSHUttle("127.0.0.1", c.Options.ConnectOptions.SSHPort, podIP, c.Options.ConnectOptions.DisableDNS, cidrs, debug), "vpn(sshuttle)", debug) + err = exec.BackgroundRun(sshuttle.SSHUttle("127.0.0.1", options.ConnectOptions.SSHPort, podIP, options.ConnectOptions.DisableDNS, cidrs, options.Debug), "vpn(sshuttle)", options.Debug) } if err != nil { return From 74941ad53b31209ad64861d1707c54728743651a Mon Sep 17 00:00:00 2001 From: yunlzheng Date: Fri, 6 Mar 2020 18:47:51 +0800 Subject: [PATCH 02/12] =?UTF-8?q?improve:=20=F0=9F=A7=AA=20clean=20code?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- pkg/kt/command/exchange.go | 3 +-- pkg/kt/command/mesh.go | 3 +-- pkg/kt/connect/exchange.go | 6 +++--- pkg/kt/connect/mesh.go | 6 +++--- pkg/kt/connect/types.go | 10 ---------- 5 files changed, 8 insertions(+), 20 deletions(-) delete mode 100644 pkg/kt/connect/types.go diff --git a/pkg/kt/command/exchange.go b/pkg/kt/command/exchange.go index c9f021f..225d830 100644 --- a/pkg/kt/command/exchange.go +++ b/pkg/kt/command/exchange.go @@ -68,8 +68,7 @@ func (action *Action) Exchange(exchange string, options *options.DaemonOptions) options.RuntimeOptions.Origin = exchange options.RuntimeOptions.Replicas = *replicas - factory := connect.Connect{} - _, err = factory.Exchange(options, origin, clientset, util.String2Map(options.Labels)) + _, err = connect.Exchange(options, origin, clientset, util.String2Map(options.Labels)) if err != nil { return err } diff --git a/pkg/kt/command/mesh.go b/pkg/kt/command/mesh.go index c98b7e5..6780339 100644 --- a/pkg/kt/command/mesh.go +++ b/pkg/kt/command/mesh.go @@ -57,8 +57,7 @@ func (action *Action) Mesh(mesh string, options *options.DaemonOptions) error { return err } - factory := connect.Connect{} - _, err = factory.Mesh(mesh, options, clientset, util.String2Map(options.Labels)) + _, err = connect.Mesh(mesh, options, clientset, util.String2Map(options.Labels)) if err != nil { return err diff --git a/pkg/kt/connect/exchange.go b/pkg/kt/connect/exchange.go index bf8b0e5..8a4ca79 100644 --- a/pkg/kt/connect/exchange.go +++ b/pkg/kt/connect/exchange.go @@ -13,9 +13,9 @@ import ( ) // Exchange exchange request to local -func (c *Connect) Exchange(options *options.DaemonOptions, origin *v1.Deployment, clientset *kubernetes.Clientset, labels map[string]string) (workload string, err error) { +func Exchange(options *options.DaemonOptions, origin *v1.Deployment, clientset *kubernetes.Clientset, labels map[string]string) (workload string, err error) { workload = origin.GetObjectMeta().GetName() + "-kt-" + strings.ToLower(util.RandomString(5)) - podIP, podName, err := c.createExchangeShadow(origin, options.Namespace, workload, clientset, labels, options.Image) + podIP, podName, err := createExchangeShadow(origin, options.Namespace, workload, clientset, labels, options.Image) options.RuntimeOptions.Shadow = workload down := int32(0) scaleTo(origin, options.Namespace, clientset, &down) @@ -38,7 +38,7 @@ func scaleTo(deployment *v1.Deployment, namespace string, clientset *kubernetes. return nil } -func (c *Connect) createExchangeShadow( +func createExchangeShadow( origin *v1.Deployment, namespace string, workload string, diff --git a/pkg/kt/connect/mesh.go b/pkg/kt/connect/mesh.go index 17b727e..f8ef87b 100644 --- a/pkg/kt/connect/mesh.go +++ b/pkg/kt/connect/mesh.go @@ -13,8 +13,8 @@ import ( ) // Mesh prepare swap deployment -func (c *Connect) Mesh(swap string, options *options.DaemonOptions, clientset *kubernetes.Clientset, labels map[string]string) (workload string, err error) { - workload, podIP, podName, err := c.createMeshShadown(swap, clientset, labels, options.Namespace, options.Image) +func Mesh(swap string, options *options.DaemonOptions, clientset *kubernetes.Clientset, labels map[string]string) (workload string, err error) { + workload, podIP, podName, err := createMeshShadown(swap, clientset, labels, options.Namespace, options.Image) if err != nil { return } @@ -23,7 +23,7 @@ func (c *Connect) Mesh(swap string, options *options.DaemonOptions, clientset *k return } -func (c *Connect) createMeshShadown( +func createMeshShadown( swap string, clientset *kubernetes.Clientset, extraLabels map[string]string, diff --git a/pkg/kt/connect/types.go b/pkg/kt/connect/types.go deleted file mode 100644 index db3c73a..0000000 --- a/pkg/kt/connect/types.go +++ /dev/null @@ -1,10 +0,0 @@ -package connect - -import ( - "github.com/alibaba/kt-connect/pkg/kt/options" -) - -// Connect VPN connect interface -type Connect struct { - Options *options.DaemonOptions -} From 77c71e04fcbcfb19067ae024545fb25c58534cf9 Mon Sep 17 00:00:00 2001 From: yunlzheng Date: Fri, 6 Mar 2020 19:40:38 +0800 Subject: [PATCH 03/12] improve: #89 refactor command connect use kubernetes interface --- pkg/kt/cluster/kubernetes.go | 23 ------ pkg/kt/cluster/types.go | 151 +++++++++++++++++++++++++++++++++++ pkg/kt/command/connect.go | 16 ++-- pkg/kt/util/kubernetes.go | 94 ---------------------- 4 files changed, 161 insertions(+), 123 deletions(-) delete mode 100644 pkg/kt/util/kubernetes.go diff --git a/pkg/kt/cluster/kubernetes.go b/pkg/kt/cluster/kubernetes.go index f4407d2..e2c7983 100644 --- a/pkg/kt/cluster/kubernetes.go +++ b/pkg/kt/cluster/kubernetes.go @@ -18,10 +18,6 @@ import ( "k8s.io/client-go/tools/clientcmd" ) -// Signal structure -type Signal struct { -} - // GetKubernetesClient get Kubernetes client from config func GetKubernetesClient(kubeConfig string) (clientset *kubernetes.Clientset, err error) { config, err := clientcmd.BuildConfigFromFlags("", kubeConfig) @@ -237,22 +233,3 @@ func generatorDeployment(namespace, name string, labels map[string]string, image }, } } - -// LocalHosts LocalHosts -func LocalHosts(clientset *kubernetes.Clientset, namespace string) (hosts map[string]string) { - serviceListener, err := clusterWatcher.ServiceListener(clientset) - if err != nil { - return - } - - services, err := serviceListener.Services(namespace).List(labels.Everything()) - if err != nil { - return - } - - hosts = map[string]string{} - for _, service := range services { - hosts[service.ObjectMeta.Name] = service.Spec.ClusterIP - } - return -} diff --git a/pkg/kt/cluster/types.go b/pkg/kt/cluster/types.go index 9536d55..59cb54f 100644 --- a/pkg/kt/cluster/types.go +++ b/pkg/kt/cluster/types.go @@ -1,9 +1,160 @@ package cluster +import ( + "fmt" + "strings" + + clusterWatcher "github.com/alibaba/kt-connect/pkg/apiserver/cluster" + mapset "github.com/deckarep/golang-set" + "github.com/rs/zerolog/log" + coreV1 "k8s.io/api/core/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/labels" + "k8s.io/client-go/kubernetes" + v1 "k8s.io/client-go/listers/core/v1" +) + +// KubernetesFactory kubernetes factory +type KubernetesFactory struct { +} + +// Create kubernetes instance +func (f *KubernetesFactory) Create(kubeConfig string) (kubernetes Kubernetes, err error) { + clientSet, err := GetKubernetesClient(kubeConfig) + if err != nil { + return + } + serviceListener, err := clusterWatcher.ServiceListener(clientSet) + podListener, err := clusterWatcher.PodListener(clientSet) + if err != nil { + return + } + kubernetes = Kubernetes{ + Clientset: clientSet, + ServiceListener: serviceListener, + PodListener: podListener, + } + return +} + // KubernetesInterface kubernetes interface type KubernetesInterface interface { + ServiceHosts(namespace string) (hosts map[string]string) + ClusterCrids(podCIDR string) (cidrs []string, err error) + CreateShadow(name, namespace, image string, labels map[string]string) (podIP, podName string, err error) } // Kubernetes implements KubernetesInterface type Kubernetes struct { + Clientset *kubernetes.Clientset + ServiceListener v1.ServiceLister + PodListener v1.PodLister +} + +// CreateShadow create shadow +func (k *Kubernetes) CreateShadow(name, namespace, image string, labels map[string]string) (podIP, podName string, err error) { + return CreateShadow(k.Clientset, name, labels, namespace, image) +} + +// ClusterCrids get cluster cirds +func (k *Kubernetes) ClusterCrids(podCIDR string) (cidrs []string, err error) { + serviceList, err := k.ServiceListener.List(labels.Everything()) + if err != nil { + return + } + + cidrs, err = getPodCirds(k.Clientset, podCIDR) + if err != nil { + return + } + + serviceCird, err := getServiceCird(serviceList) + if err != nil { + return + } + cidrs = append(cidrs, serviceCird...) + return +} + +// ServiceHosts get service dns map +func (k *Kubernetes) ServiceHosts(namespace string) (hosts map[string]string) { + services, err := k.ServiceListener.Services(namespace).List(labels.Everything()) + if err != nil { + return + } + hosts = map[string]string{} + for _, service := range services { + hosts[service.ObjectMeta.Name] = service.Spec.ClusterIP + } + return +} + +func getPodCirds(clientset *kubernetes.Clientset, podCIDR string) (cidrs []string, err error) { + cidrs = []string{} + + if len(podCIDR) != 0 { + cidrs = append(cidrs, podCIDR) + return + } + + nodeList, err := clientset.CoreV1().Nodes().List(metav1.ListOptions{}) + + if err != nil { + log.Printf("Fails to get node info of cluster") + return nil, err + } + + for _, node := range nodeList.Items { + if node.Spec.PodCIDR != "" && len(node.Spec.PodCIDR) != 0 { + cidrs = append(cidrs, node.Spec.PodCIDR) + } + } + + if len(cidrs) == 0 { + samples, err2 := getPodCirdByInstance(clientset) + if err2 != nil { + err = err2 + return + } + for _, sample := range samples.ToSlice() { + cidrs = append(cidrs, fmt.Sprint(sample)) + } + } + + return +} + +func getPodCirdByInstance(clientset *kubernetes.Clientset) (samples mapset.Set, err error) { + log.Info().Msgf("Fail to get pod cidr from node.Spec.PODCIDR, try to get with pod sample") + podList, err := clientset.CoreV1().Pods("").List(metav1.ListOptions{}) + if err != nil { + log.Printf("Fails to get service info of cluster") + return + } + + samples = mapset.NewSet() + for _, pod := range podList.Items { + if pod.Status.PodIP != "" && pod.Status.PodIP != "None" { + samples.Add(getCirdFromSample(pod.Status.PodIP)) + } + } + return +} + +func getServiceCird(serviceList []*coreV1.Service) (cidr []string, err error) { + samples := mapset.NewSet() + for _, service := range serviceList { + if service.Spec.ClusterIP != "" && service.Spec.ClusterIP != "None" { + samples.Add(getCirdFromSample(service.Spec.ClusterIP)) + } + } + + for _, sample := range samples.ToSlice() { + cidr = append(cidr, fmt.Sprint(sample)) + } + return +} + +func getCirdFromSample(sample string) string { + return strings.Join(append(strings.Split(sample, ".")[:2], []string{"0", "0"}...), ".") + "/16" } diff --git a/pkg/kt/command/connect.go b/pkg/kt/command/connect.go index e8ee10c..3f5ecec 100644 --- a/pkg/kt/command/connect.go +++ b/pkg/kt/command/connect.go @@ -86,28 +86,32 @@ func (action *Action) Connect(options *options.DaemonOptions) (err error) { log.Info().Msgf("Connect Start At %d", pid) - clientSet, err := cluster.GetKubernetesClient(options.KubeConfig) + factory := cluster.KubernetesFactory{} + kubernetes, err := factory.Create(options.KubeConfig) if err != nil { return } if options.ConnectOptions.Dump2Hosts { - hosts := cluster.LocalHosts(clientSet, options.Namespace) + hosts := kubernetes.ServiceHosts(options.Namespace) util.DumpHosts(hosts) options.ConnectOptions.Hosts = hosts } workload := fmt.Sprintf("kt-connect-daemon-%s", strings.ToLower(util.RandomString(5))) - options.RuntimeOptions.Shadow = workload - endPointIP, podName, err := cluster.CreateShadow( - clientSet, workload, labels(workload, options), options.Namespace, options.Image) + endPointIP, podName, err := kubernetes.CreateShadow( + workload, options.Namespace, options.Image, labels(workload, options), + ) if err != nil { return } - cidrs, err := util.GetCirds(clientSet, options.ConnectOptions.CIDR) + // record shadow name will clean up terminal + options.RuntimeOptions.Shadow = workload + + cidrs, err := kubernetes.ClusterCrids(options.ConnectOptions.CIDR) if err != nil { return } diff --git a/pkg/kt/util/kubernetes.go b/pkg/kt/util/kubernetes.go deleted file mode 100644 index fdcaa2b..0000000 --- a/pkg/kt/util/kubernetes.go +++ /dev/null @@ -1,94 +0,0 @@ -package util - -import ( - "fmt" - "strings" - - "github.com/deckarep/golang-set" - "github.com/rs/zerolog/log" - metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" - "k8s.io/client-go/kubernetes" -) - -// GetCirds Get kubernetes cluster resource crids -func GetCirds(clientset *kubernetes.Clientset, podCIDR string) (cidrs []string, err error) { - cidrs, err = getPodCirds(clientset, podCIDR) - if err != nil { - return - } - serviceCird, err := getServiceCird(clientset) - if err != nil { - return - } - cidrs = append(cidrs, serviceCird...) - return -} - -func getPodCirds(clientset *kubernetes.Clientset, podCIDR string) (cidrs []string, err error) { - cidrs = []string{} - - if len(podCIDR) != 0 { - cidrs = append(cidrs, podCIDR) - return - } - - nodeList, err := clientset.CoreV1().Nodes().List(metav1.ListOptions{}) - - if err != nil { - log.Printf("Fails to get node info of cluster") - return nil, err - } - - for _, node := range nodeList.Items { - if node.Spec.PodCIDR != "" && len(node.Spec.PodCIDR) != 0 { - cidrs = append(cidrs, node.Spec.PodCIDR) - } - } - - if len(cidrs) == 0 { - log.Info().Msgf("Fail to get pod cidr from node.Spec.PODCIDR, try to get with pod sample") - podList, err2 := clientset.CoreV1().Pods("").List(metav1.ListOptions{}) - if err2 != nil { - log.Printf("Fails to get service info of cluster") - return - } - - samples := mapset.NewSet() - for _, pod := range podList.Items { - if pod.Status.PodIP != "" && pod.Status.PodIP != "None" { - samples.Add(getCirdFromSample(pod.Status.PodIP)) - } - } - - for _, sample := range samples.ToSlice() { - cidrs = append(cidrs, fmt.Sprint(sample)) - } - } - - return -} - -func getServiceCird(clientset *kubernetes.Clientset) (cidr []string, err error) { - serviceList, err := clientset.CoreV1().Services("").List(metav1.ListOptions{}) - if err != nil { - log.Printf("Fails to get service info of cluster") - return cidr, err - } - - samples := mapset.NewSet() - for _, service := range serviceList.Items { - if service.Spec.ClusterIP != "" && service.Spec.ClusterIP != "None" { - samples.Add(getCirdFromSample(service.Spec.ClusterIP)) - } - } - - for _, sample := range samples.ToSlice() { - cidr = append(cidr, fmt.Sprint(sample)) - } - - return -} - -func getCirdFromSample(sample string) string { - return strings.Join(append(strings.Split(sample, ".")[:2], []string{"0", "0"}...), ".") + "/16" -} From 41d985465d13cc9c9d30875ce97fc51b03a52282 Mon Sep 17 00:00:00 2001 From: yunlzheng Date: Fri, 6 Mar 2020 20:13:42 +0800 Subject: [PATCH 04/12] feat: #89 refactor command exchange use kubernetes interface --- pkg/kt/cluster/types.go | 24 ++++++++++++++ pkg/kt/command/exchange.go | 47 ++++++++++++++++++++++------ pkg/kt/connect/exchange.go | 64 -------------------------------------- 3 files changed, 62 insertions(+), 73 deletions(-) diff --git a/pkg/kt/cluster/types.go b/pkg/kt/cluster/types.go index 59cb54f..6838612 100644 --- a/pkg/kt/cluster/types.go +++ b/pkg/kt/cluster/types.go @@ -7,6 +7,7 @@ import ( clusterWatcher "github.com/alibaba/kt-connect/pkg/apiserver/cluster" mapset "github.com/deckarep/golang-set" "github.com/rs/zerolog/log" + appV1 "k8s.io/api/apps/v1" coreV1 "k8s.io/api/core/v1" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/labels" @@ -39,6 +40,8 @@ func (f *KubernetesFactory) Create(kubeConfig string) (kubernetes Kubernetes, er // KubernetesInterface kubernetes interface type KubernetesInterface interface { + Deployment(name, namespace string) (deployment appV1.Deployment, err error) + Scale(name, namespace string, replicas *int32) (err error) ServiceHosts(namespace string) (hosts map[string]string) ClusterCrids(podCIDR string) (cidrs []string, err error) CreateShadow(name, namespace, image string, labels map[string]string) (podIP, podName string, err error) @@ -51,6 +54,27 @@ type Kubernetes struct { PodListener v1.PodLister } +// Scale scale deployment to +func (k *Kubernetes) Scale(deployment *appV1.Deployment, replicas *int32) (err error) { + log.Printf("scale deployment %s to %d\n", deployment.GetObjectMeta().GetName(), *replicas) + client := k.Clientset.AppsV1().Deployments(deployment.GetObjectMeta().GetNamespace()) + deployment.Spec.Replicas = replicas + + d, err := client.Update(deployment) + if err != nil { + log.Printf("%s Fails scale deployment %s to %d\n", err.Error(), deployment.GetObjectMeta().GetName(), *replicas) + return + } + log.Printf(" * %s (%d replicas) success", d.Name, *d.Spec.Replicas) + return +} + +// Deployment get deployment +func (k *Kubernetes) Deployment(name, namespace string) (deployment *appV1.Deployment, err error) { + deployment, err = k.Clientset.AppsV1().Deployments(namespace).Get(name, metav1.GetOptions{}) + return +} + // CreateShadow create shadow func (k *Kubernetes) CreateShadow(name, namespace, image string, labels map[string]string) (podIP, podName string, err error) { return CreateShadow(k.Clientset, name, labels, namespace, image) diff --git a/pkg/kt/command/exchange.go b/pkg/kt/command/exchange.go index 225d830..1badd83 100644 --- a/pkg/kt/command/exchange.go +++ b/pkg/kt/command/exchange.go @@ -2,6 +2,9 @@ package command import ( "errors" + "strings" + + v1 "k8s.io/api/apps/v1" "github.com/rs/zerolog" "github.com/rs/zerolog/log" @@ -11,8 +14,6 @@ import ( "github.com/alibaba/kt-connect/pkg/kt/options" "github.com/alibaba/kt-connect/pkg/kt/util" "github.com/urfave/cli" - - metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" ) // newExchangeCommand return new exchange command @@ -52,29 +53,57 @@ func (action *Action) Exchange(exchange string, options *options.DaemonOptions) ch := SetUpCloseHandler(options) checkConnectRunning(options.RuntimeOptions.PidFile) - clientset, err := cluster.GetKubernetesClient(options.KubeConfig) + + factory := cluster.KubernetesFactory{} + kubernetes, err := factory.Create(options.KubeConfig) if err != nil { return err } - origin, err := clientset.AppsV1().Deployments(options.Namespace).Get(exchange, metav1.GetOptions{}) + app, err := kubernetes.Deployment(exchange, options.Namespace) if err != nil { return err } - replicas := origin.Spec.Replicas + // record context inorder to remove after command exit + options.RuntimeOptions.Origin = app.GetObjectMeta().GetName() + options.RuntimeOptions.Replicas = *app.Spec.Replicas - // Prepare context inorder to remove after command exit - options.RuntimeOptions.Origin = exchange - options.RuntimeOptions.Replicas = *replicas + workload := app.GetObjectMeta().GetName() + "-kt-" + strings.ToLower(util.RandomString(5)) + podIP, podName, err := kubernetes.CreateShadow( + workload, options.Namespace, options.Image, getExchangeLables(options.Labels, workload, app)) + log.Info().Msgf("create exchange shadow %s in namespace %s", workload, options.Namespace) - _, err = connect.Exchange(options, origin, clientset, util.String2Map(options.Labels)) if err != nil { return err } + // record data + options.RuntimeOptions.Shadow = workload + + down := int32(0) + kubernetes.Scale(app, &down) + + connect.RemotePortForward(options.ExchangeOptions.Expose, options.KubeConfig, options.Namespace, podName, podIP, options.Debug) + s := <-ch log.Info().Msgf("Terminal Signal is %s", s) return nil } + +func getExchangeLables(customLabels string, workload string, origin *v1.Deployment) map[string]string { + labels := map[string]string{ + "kt": workload, + "kt-component": "exchange", + "control-by": "kt", + } + for k, v := range origin.Spec.Selector.MatchLabels { + labels[k] = v + } + // extra labels must be applied after origin labels + for k, v := range util.String2Map(customLabels) { + labels[k] = v + } + return labels +} diff --git a/pkg/kt/connect/exchange.go b/pkg/kt/connect/exchange.go index 8a4ca79..9c5ea44 100644 --- a/pkg/kt/connect/exchange.go +++ b/pkg/kt/connect/exchange.go @@ -1,65 +1 @@ package connect - -import ( - "strings" - - "github.com/rs/zerolog/log" - - "github.com/alibaba/kt-connect/pkg/kt/cluster" - "github.com/alibaba/kt-connect/pkg/kt/options" - "github.com/alibaba/kt-connect/pkg/kt/util" - v1 "k8s.io/api/apps/v1" - "k8s.io/client-go/kubernetes" -) - -// Exchange exchange request to local -func Exchange(options *options.DaemonOptions, origin *v1.Deployment, clientset *kubernetes.Clientset, labels map[string]string) (workload string, err error) { - workload = origin.GetObjectMeta().GetName() + "-kt-" + strings.ToLower(util.RandomString(5)) - podIP, podName, err := createExchangeShadow(origin, options.Namespace, workload, clientset, labels, options.Image) - options.RuntimeOptions.Shadow = workload - down := int32(0) - scaleTo(origin, options.Namespace, clientset, &down) - RemotePortForward(options.ExchangeOptions.Expose, options.KubeConfig, options.Namespace, podName, podIP, options.Debug) - return -} - -//ScaleTo Scale -func scaleTo(deployment *v1.Deployment, namespace string, clientset *kubernetes.Clientset, replicas *int32) (err error) { - log.Printf("Try Scale deployment %s to %d\n", deployment.GetObjectMeta().GetName(), *replicas) - client := clientset.AppsV1().Deployments(namespace) - deployment.Spec.Replicas = replicas - - d, err := client.Update(deployment) - if err != nil { - log.Printf("%s Fails scale deployment %s to %d\n", err.Error(), deployment.GetObjectMeta().GetName(), *replicas) - return err - } - log.Printf(" * %s (%d replicas) success", d.Name, *d.Spec.Replicas) - return nil -} - -func createExchangeShadow( - origin *v1.Deployment, - namespace string, - workload string, - clientset *kubernetes.Clientset, - extraLabels map[string]string, - image string, -) (podIP string, podName string, err error) { - log.Info().Msgf("Create Exchange shadow %s in namespace %s", workload, namespace) - labels := map[string]string{ - "kt": workload, - "kt-component": "exchange", - "control-by": "kt", - } - for k, v := range origin.Spec.Selector.MatchLabels { - labels[k] = v - } - // extra labels must be applied after origin labels - for k, v := range extraLabels { - labels[k] = v - } - - podIP, podName, err = cluster.CreateShadow(clientset, workload, labels, namespace, image) - return -} From ccac33c6e7119c170b796ba5e04e3f1c483fa961 Mon Sep 17 00:00:00 2001 From: yunlzheng Date: Fri, 6 Mar 2020 20:26:28 +0800 Subject: [PATCH 05/12] =?UTF-8?q?improve:=20=F0=9F=A7=AA=20refactor=20comm?= =?UTF-8?q?and=20exchange=20use=20kubernetes=20interface?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- pkg/kt/cluster/helper.go | 93 ++++++++++++++++++++++++ pkg/kt/cluster/kubernetes.go | 97 ++++++++++++++++++------- pkg/kt/cluster/types.go | 137 ----------------------------------- 3 files changed, 165 insertions(+), 162 deletions(-) create mode 100644 pkg/kt/cluster/helper.go diff --git a/pkg/kt/cluster/helper.go b/pkg/kt/cluster/helper.go new file mode 100644 index 0000000..32387d7 --- /dev/null +++ b/pkg/kt/cluster/helper.go @@ -0,0 +1,93 @@ +package cluster + +import ( + "fmt" + "strings" + + mapset "github.com/deckarep/golang-set" + "github.com/rs/zerolog/log" + coreV1 "k8s.io/api/core/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/client-go/kubernetes" + "k8s.io/client-go/tools/clientcmd" +) + +// GetKubernetesClient get Kubernetes client from config +func GetKubernetesClient(kubeConfig string) (clientset *kubernetes.Clientset, err error) { + config, err := clientcmd.BuildConfigFromFlags("", kubeConfig) + if err != nil { + return nil, err + } + clientset, err = kubernetes.NewForConfig(config) + return +} + +func getPodCirds(clientset *kubernetes.Clientset, podCIDR string) (cidrs []string, err error) { + cidrs = []string{} + + if len(podCIDR) != 0 { + cidrs = append(cidrs, podCIDR) + return + } + + nodeList, err := clientset.CoreV1().Nodes().List(metav1.ListOptions{}) + + if err != nil { + log.Printf("Fails to get node info of cluster") + return nil, err + } + + for _, node := range nodeList.Items { + if node.Spec.PodCIDR != "" && len(node.Spec.PodCIDR) != 0 { + cidrs = append(cidrs, node.Spec.PodCIDR) + } + } + + if len(cidrs) == 0 { + samples, err2 := getPodCirdByInstance(clientset) + if err2 != nil { + err = err2 + return + } + for _, sample := range samples.ToSlice() { + cidrs = append(cidrs, fmt.Sprint(sample)) + } + } + + return +} + +func getPodCirdByInstance(clientset *kubernetes.Clientset) (samples mapset.Set, err error) { + log.Info().Msgf("Fail to get pod cidr from node.Spec.PODCIDR, try to get with pod sample") + podList, err := clientset.CoreV1().Pods("").List(metav1.ListOptions{}) + if err != nil { + log.Printf("Fails to get service info of cluster") + return + } + + samples = mapset.NewSet() + for _, pod := range podList.Items { + if pod.Status.PodIP != "" && pod.Status.PodIP != "None" { + samples.Add(getCirdFromSample(pod.Status.PodIP)) + } + } + return +} + +func getServiceCird(serviceList []*coreV1.Service) (cidr []string, err error) { + samples := mapset.NewSet() + for _, service := range serviceList { + if service.Spec.ClusterIP != "" && service.Spec.ClusterIP != "None" { + samples.Add(getCirdFromSample(service.Spec.ClusterIP)) + } + } + + for _, sample := range samples.ToSlice() { + cidr = append(cidr, fmt.Sprint(sample)) + } + return +} + +func getCirdFromSample(sample string) string { + return strings.Join(append(strings.Split(sample, ".")[:2], []string{"0", "0"}...), ".") + "/16" +} diff --git a/pkg/kt/cluster/kubernetes.go b/pkg/kt/cluster/kubernetes.go index e2c7983..6bdb1dd 100644 --- a/pkg/kt/cluster/kubernetes.go +++ b/pkg/kt/cluster/kubernetes.go @@ -8,30 +8,77 @@ import ( clusterWatcher "github.com/alibaba/kt-connect/pkg/apiserver/cluster" "github.com/alibaba/kt-connect/pkg/kt/util" "github.com/rs/zerolog/log" - appsv1 "k8s.io/api/apps/v1" - apiv1 "k8s.io/api/core/v1" + appV1 "k8s.io/api/apps/v1" v1 "k8s.io/api/core/v1" - metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + metaV1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/labels" "k8s.io/apimachinery/pkg/selection" "k8s.io/client-go/kubernetes" - "k8s.io/client-go/tools/clientcmd" ) -// GetKubernetesClient get Kubernetes client from config -func GetKubernetesClient(kubeConfig string) (clientset *kubernetes.Clientset, err error) { - config, err := clientcmd.BuildConfigFromFlags("", kubeConfig) +// Scale scale deployment to +func (k *Kubernetes) Scale(deployment *appV1.Deployment, replicas *int32) (err error) { + log.Printf("scale deployment %s to %d\n", deployment.GetObjectMeta().GetName(), *replicas) + client := k.Clientset.AppsV1().Deployments(deployment.GetObjectMeta().GetNamespace()) + deployment.Spec.Replicas = replicas + + d, err := client.Update(deployment) if err != nil { - return nil, err + log.Printf("%s Fails scale deployment %s to %d\n", err.Error(), deployment.GetObjectMeta().GetName(), *replicas) + return + } + log.Printf(" * %s (%d replicas) success", d.Name, *d.Spec.Replicas) + return +} + +// Deployment get deployment +func (k *Kubernetes) Deployment(name, namespace string) (deployment *appV1.Deployment, err error) { + deployment, err = k.Clientset.AppsV1().Deployments(namespace).Get(name, metaV1.GetOptions{}) + return +} + +// CreateShadow create shadow +func (k *Kubernetes) CreateShadow(name, namespace, image string, labels map[string]string) (podIP, podName string, err error) { + return CreateShadow(k.Clientset, name, labels, namespace, image) +} + +// ClusterCrids get cluster cirds +func (k *Kubernetes) ClusterCrids(podCIDR string) (cidrs []string, err error) { + serviceList, err := k.ServiceListener.List(labels.Everything()) + if err != nil { + return + } + + cidrs, err = getPodCirds(k.Clientset, podCIDR) + if err != nil { + return + } + + serviceCird, err := getServiceCird(serviceList) + if err != nil { + return + } + cidrs = append(cidrs, serviceCird...) + return +} + +// ServiceHosts get service dns map +func (k *Kubernetes) ServiceHosts(namespace string) (hosts map[string]string) { + services, err := k.ServiceListener.Services(namespace).List(labels.Everything()) + if err != nil { + return + } + hosts = map[string]string{} + for _, service := range services { + hosts[service.ObjectMeta.Name] = service.Spec.ClusterIP } - clientset, err = kubernetes.NewForConfig(config) return } // ScaleTo scale app func ScaleTo(clientSet *kubernetes.Clientset, namespace, name string, replicas int32) (err error) { client := clientSet.AppsV1().Deployments(namespace) - deployment, err := client.Get(name, metav1.GetOptions{}) + deployment, err := client.Get(name, metaV1.GetOptions{}) if err != nil { return } @@ -51,8 +98,8 @@ func ScaleTo(clientSet *kubernetes.Clientset, namespace, name string, replicas i // RemoveShadow remove shadow from cluster func RemoveShadow(client *kubernetes.Clientset, namespace, name string) { deploymentsClient := client.AppsV1().Deployments(namespace) - deletePolicy := metav1.DeletePropagationBackground - err := deploymentsClient.Delete(name, &metav1.DeleteOptions{ + deletePolicy := metaV1.DeletePropagationBackground + err := deploymentsClient.Delete(name, &metaV1.DeleteOptions{ PropagationPolicy: &deletePolicy, }) if err != nil { @@ -78,7 +125,7 @@ func RemoveService( clientset *kubernetes.Clientset, ) (err error) { client := clientset.CoreV1().Services(namespace) - return client.Delete(name, &metav1.DeleteOptions{}) + return client.Delete(name, &metaV1.DeleteOptions{}) } // CreateShadow create shadow @@ -112,12 +159,12 @@ func CreateShadow( return } -func waitPodReadyUsingInformer(namespace, name string, clientset *kubernetes.Clientset) (pod apiv1.Pod, err error) { +func waitPodReadyUsingInformer(namespace, name string, clientset *kubernetes.Clientset) (pod v1.Pod, err error) { podListener, err := clusterWatcher.PodListener(clientset) if err != nil { return } - pod = apiv1.Pod{} + pod = v1.Pod{} podLabels := labels.NewSelector() log.Info().Msgf("pod label: kt=%s", name) labelKeys := []string{ @@ -191,7 +238,7 @@ func generateService(name, namespace string, labels map[string]string, port int) }) return &v1.Service{ - ObjectMeta: metav1.ObjectMeta{ + ObjectMeta: metaV1.ObjectMeta{ Name: name, Namespace: namespace, Labels: labels, @@ -205,23 +252,23 @@ func generateService(name, namespace string, labels map[string]string, port int) } -func generatorDeployment(namespace, name string, labels map[string]string, image string) *appsv1.Deployment { - return &appsv1.Deployment{ - ObjectMeta: metav1.ObjectMeta{ +func generatorDeployment(namespace, name string, labels map[string]string, image string) *appV1.Deployment { + return &appV1.Deployment{ + ObjectMeta: metaV1.ObjectMeta{ Name: name, Namespace: namespace, Labels: labels, }, - Spec: appsv1.DeploymentSpec{ - Selector: &metav1.LabelSelector{ + Spec: appV1.DeploymentSpec{ + Selector: &metaV1.LabelSelector{ MatchLabels: labels, }, - Template: apiv1.PodTemplateSpec{ - ObjectMeta: metav1.ObjectMeta{ + Template: v1.PodTemplateSpec{ + ObjectMeta: metaV1.ObjectMeta{ Labels: labels, }, - Spec: apiv1.PodSpec{ - Containers: []apiv1.Container{ + Spec: v1.PodSpec{ + Containers: []v1.Container{ { Name: "standalone", Image: image, diff --git a/pkg/kt/cluster/types.go b/pkg/kt/cluster/types.go index 6838612..dbd3597 100644 --- a/pkg/kt/cluster/types.go +++ b/pkg/kt/cluster/types.go @@ -1,16 +1,8 @@ package cluster import ( - "fmt" - "strings" - clusterWatcher "github.com/alibaba/kt-connect/pkg/apiserver/cluster" - mapset "github.com/deckarep/golang-set" - "github.com/rs/zerolog/log" appV1 "k8s.io/api/apps/v1" - coreV1 "k8s.io/api/core/v1" - metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" - "k8s.io/apimachinery/pkg/labels" "k8s.io/client-go/kubernetes" v1 "k8s.io/client-go/listers/core/v1" ) @@ -53,132 +45,3 @@ type Kubernetes struct { ServiceListener v1.ServiceLister PodListener v1.PodLister } - -// Scale scale deployment to -func (k *Kubernetes) Scale(deployment *appV1.Deployment, replicas *int32) (err error) { - log.Printf("scale deployment %s to %d\n", deployment.GetObjectMeta().GetName(), *replicas) - client := k.Clientset.AppsV1().Deployments(deployment.GetObjectMeta().GetNamespace()) - deployment.Spec.Replicas = replicas - - d, err := client.Update(deployment) - if err != nil { - log.Printf("%s Fails scale deployment %s to %d\n", err.Error(), deployment.GetObjectMeta().GetName(), *replicas) - return - } - log.Printf(" * %s (%d replicas) success", d.Name, *d.Spec.Replicas) - return -} - -// Deployment get deployment -func (k *Kubernetes) Deployment(name, namespace string) (deployment *appV1.Deployment, err error) { - deployment, err = k.Clientset.AppsV1().Deployments(namespace).Get(name, metav1.GetOptions{}) - return -} - -// CreateShadow create shadow -func (k *Kubernetes) CreateShadow(name, namespace, image string, labels map[string]string) (podIP, podName string, err error) { - return CreateShadow(k.Clientset, name, labels, namespace, image) -} - -// ClusterCrids get cluster cirds -func (k *Kubernetes) ClusterCrids(podCIDR string) (cidrs []string, err error) { - serviceList, err := k.ServiceListener.List(labels.Everything()) - if err != nil { - return - } - - cidrs, err = getPodCirds(k.Clientset, podCIDR) - if err != nil { - return - } - - serviceCird, err := getServiceCird(serviceList) - if err != nil { - return - } - cidrs = append(cidrs, serviceCird...) - return -} - -// ServiceHosts get service dns map -func (k *Kubernetes) ServiceHosts(namespace string) (hosts map[string]string) { - services, err := k.ServiceListener.Services(namespace).List(labels.Everything()) - if err != nil { - return - } - hosts = map[string]string{} - for _, service := range services { - hosts[service.ObjectMeta.Name] = service.Spec.ClusterIP - } - return -} - -func getPodCirds(clientset *kubernetes.Clientset, podCIDR string) (cidrs []string, err error) { - cidrs = []string{} - - if len(podCIDR) != 0 { - cidrs = append(cidrs, podCIDR) - return - } - - nodeList, err := clientset.CoreV1().Nodes().List(metav1.ListOptions{}) - - if err != nil { - log.Printf("Fails to get node info of cluster") - return nil, err - } - - for _, node := range nodeList.Items { - if node.Spec.PodCIDR != "" && len(node.Spec.PodCIDR) != 0 { - cidrs = append(cidrs, node.Spec.PodCIDR) - } - } - - if len(cidrs) == 0 { - samples, err2 := getPodCirdByInstance(clientset) - if err2 != nil { - err = err2 - return - } - for _, sample := range samples.ToSlice() { - cidrs = append(cidrs, fmt.Sprint(sample)) - } - } - - return -} - -func getPodCirdByInstance(clientset *kubernetes.Clientset) (samples mapset.Set, err error) { - log.Info().Msgf("Fail to get pod cidr from node.Spec.PODCIDR, try to get with pod sample") - podList, err := clientset.CoreV1().Pods("").List(metav1.ListOptions{}) - if err != nil { - log.Printf("Fails to get service info of cluster") - return - } - - samples = mapset.NewSet() - for _, pod := range podList.Items { - if pod.Status.PodIP != "" && pod.Status.PodIP != "None" { - samples.Add(getCirdFromSample(pod.Status.PodIP)) - } - } - return -} - -func getServiceCird(serviceList []*coreV1.Service) (cidr []string, err error) { - samples := mapset.NewSet() - for _, service := range serviceList { - if service.Spec.ClusterIP != "" && service.Spec.ClusterIP != "None" { - samples.Add(getCirdFromSample(service.Spec.ClusterIP)) - } - } - - for _, sample := range samples.ToSlice() { - cidr = append(cidr, fmt.Sprint(sample)) - } - return -} - -func getCirdFromSample(sample string) string { - return strings.Join(append(strings.Split(sample, ".")[:2], []string{"0", "0"}...), ".") + "/16" -} From d0570da9ee129b5a27a946d5dc5cb3c912949d2a Mon Sep 17 00:00:00 2001 From: yunlzheng Date: Fri, 6 Mar 2020 20:41:38 +0800 Subject: [PATCH 06/12] improve: use kubernetes interface in mesh command --- pkg/kt/command/connect.go | 2 +- pkg/kt/command/mesh.go | 47 ++++++++++++++++++++++++++-- pkg/kt/connect/exchange.go | 1 - pkg/kt/connect/mesh.go | 64 -------------------------------------- pkg/kt/connect/outbound.go | 4 +-- 5 files changed, 48 insertions(+), 70 deletions(-) delete mode 100644 pkg/kt/connect/exchange.go delete mode 100644 pkg/kt/connect/mesh.go diff --git a/pkg/kt/command/connect.go b/pkg/kt/command/connect.go index 3f5ecec..deb88f7 100644 --- a/pkg/kt/command/connect.go +++ b/pkg/kt/command/connect.go @@ -116,7 +116,7 @@ func (action *Action) Connect(options *options.DaemonOptions) (err error) { return } - err = connect.StartConnect(podName, endPointIP, cidrs, options) + err = connect.Connect(podName, endPointIP, cidrs, options) if err != nil { return } diff --git a/pkg/kt/command/mesh.go b/pkg/kt/command/mesh.go index 6780339..032b37f 100644 --- a/pkg/kt/command/mesh.go +++ b/pkg/kt/command/mesh.go @@ -2,17 +2,25 @@ package command import ( "errors" + "strings" "github.com/alibaba/kt-connect/pkg/kt/cluster" "github.com/alibaba/kt-connect/pkg/kt/connect" "github.com/alibaba/kt-connect/pkg/kt/options" "github.com/alibaba/kt-connect/pkg/kt/util" + v1 "k8s.io/api/apps/v1" "github.com/rs/zerolog" "github.com/rs/zerolog/log" "github.com/urfave/cli" ) +// ComponentMesh mesh component +const ComponentMesh = "mesh" + +// KubernetesTool kt sign +const KubernetesTool = "kt" + // newMeshCommand return new mesh command func newMeshCommand(options *options.DaemonOptions, action ActionInterface) cli.Command { return cli.Command{ @@ -52,12 +60,30 @@ func (action *Action) Mesh(mesh string, options *options.DaemonOptions) error { ch := SetUpCloseHandler(options) - clientset, err := cluster.GetKubernetesClient(options.KubeConfig) + factory := cluster.KubernetesFactory{} + kubernetes, err := factory.Create(options.KubeConfig) if err != nil { return err } - _, err = connect.Mesh(mesh, options, clientset, util.String2Map(options.Labels)) + app, err := kubernetes.Deployment(mesh, options.Namespace) + if err != nil { + return err + } + + meshVersion := strings.ToLower(util.RandomString(5)) + workload := app.GetObjectMeta().GetName() + "-kt-" + meshVersion + + labels := getMeshLabels(workload, meshVersion, app, options) + + podIP, podName, err := kubernetes.CreateShadow(workload, options.Namespace, options.Image, labels) + if err != nil { + return err + } + + // record context data + options.RuntimeOptions.Shadow = workload + err = connect.RemotePortForward(options.MeshOptions.Expose, options.KubeConfig, options.Namespace, podName, podIP, options.Debug) if err != nil { return err @@ -68,3 +94,20 @@ func (action *Action) Mesh(mesh string, options *options.DaemonOptions) error { return nil } + +func getMeshLabels(workload string, meshVersion string, app *v1.Deployment, options *options.DaemonOptions) map[string]string { + labels := map[string]string{ + "kt": workload, + "version": meshVersion, + "kt-component": ComponentMesh, + "control-by": KubernetesTool, + } + for k, v := range app.Spec.Selector.MatchLabels { + labels[k] = v + } + // extra labels must be applied after origin labels + for k, v := range util.String2Map(options.Labels) { + labels[k] = v + } + return labels +} diff --git a/pkg/kt/connect/exchange.go b/pkg/kt/connect/exchange.go deleted file mode 100644 index 9c5ea44..0000000 --- a/pkg/kt/connect/exchange.go +++ /dev/null @@ -1 +0,0 @@ -package connect diff --git a/pkg/kt/connect/mesh.go b/pkg/kt/connect/mesh.go deleted file mode 100644 index f8ef87b..0000000 --- a/pkg/kt/connect/mesh.go +++ /dev/null @@ -1,64 +0,0 @@ -package connect - -import ( - "strings" - - "github.com/rs/zerolog/log" - - "github.com/alibaba/kt-connect/pkg/kt/cluster" - "github.com/alibaba/kt-connect/pkg/kt/options" - "github.com/alibaba/kt-connect/pkg/kt/util" - metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" - "k8s.io/client-go/kubernetes" -) - -// Mesh prepare swap deployment -func Mesh(swap string, options *options.DaemonOptions, clientset *kubernetes.Clientset, labels map[string]string) (workload string, err error) { - workload, podIP, podName, err := createMeshShadown(swap, clientset, labels, options.Namespace, options.Image) - if err != nil { - return - } - options.RuntimeOptions.Shadow = workload - err = RemotePortForward(options.MeshOptions.Expose, options.KubeConfig, options.Namespace, podName, podIP, options.Debug) - return -} - -func createMeshShadown( - swap string, - clientset *kubernetes.Clientset, - extraLabels map[string]string, - namespace, image string, -) (shadowName, podIP, podName string, err error) { - deploymentsClient := clientset.AppsV1().Deployments(namespace) - origin, err := deploymentsClient.Get(swap, metav1.GetOptions{}) - if err != nil { - return "", "", "", err - } - - meshVersion := strings.ToLower(util.RandomString(5)) - shadowName = origin.GetObjectMeta().GetName() + "-kt-" + meshVersion - labels := map[string]string{ - "kt": shadowName, - "kt-component": "mesh", - "control-by": "kt", - "version": meshVersion, - } - for k, v := range origin.Spec.Selector.MatchLabels { - labels[k] = v - } - // extra labels must be applied after origin labels - for k, v := range extraLabels { - labels[k] = v - } - - podIP, podName, err = cluster.CreateShadow(clientset, shadowName, labels, namespace, image) - if err != nil { - return "", "", "", err - } - - log.Printf("-----------------------------------------------------------\n") - log.Printf("| Mesh Version '%s' You can update Istio rule |\n", meshVersion) - log.Printf("-----------------------------------------------------------\n") - - return -} diff --git a/pkg/kt/connect/outbound.go b/pkg/kt/connect/outbound.go index 8eb56dd..d62cec9 100644 --- a/pkg/kt/connect/outbound.go +++ b/pkg/kt/connect/outbound.go @@ -15,8 +15,8 @@ import ( "github.com/rs/zerolog/log" ) -// StartConnect start vpn connection -func StartConnect(name, podIP string, cidrs []string, options *options.DaemonOptions) (err error) { +// Connect start vpn connection +func Connect(name, podIP string, cidrs []string, options *options.DaemonOptions) (err error) { err = util.PrepareSSHPrivateKey() if err != nil { return From f9c15068c954739a5706b0a73da52e8836ca02d1 Mon Sep 17 00:00:00 2001 From: yunlzheng Date: Fri, 6 Mar 2020 20:43:02 +0800 Subject: [PATCH 07/12] improve: clean code --- pkg/kt/cluster/types.go | 6 +----- pkg/kt/command/connect.go | 3 +-- pkg/kt/command/exchange.go | 3 +-- pkg/kt/command/mesh.go | 3 +-- 4 files changed, 4 insertions(+), 11 deletions(-) diff --git a/pkg/kt/cluster/types.go b/pkg/kt/cluster/types.go index dbd3597..26513b4 100644 --- a/pkg/kt/cluster/types.go +++ b/pkg/kt/cluster/types.go @@ -7,12 +7,8 @@ import ( v1 "k8s.io/client-go/listers/core/v1" ) -// KubernetesFactory kubernetes factory -type KubernetesFactory struct { -} - // Create kubernetes instance -func (f *KubernetesFactory) Create(kubeConfig string) (kubernetes Kubernetes, err error) { +func Create(kubeConfig string) (kubernetes Kubernetes, err error) { clientSet, err := GetKubernetesClient(kubeConfig) if err != nil { return diff --git a/pkg/kt/command/connect.go b/pkg/kt/command/connect.go index deb88f7..c6c805b 100644 --- a/pkg/kt/command/connect.go +++ b/pkg/kt/command/connect.go @@ -86,8 +86,7 @@ func (action *Action) Connect(options *options.DaemonOptions) (err error) { log.Info().Msgf("Connect Start At %d", pid) - factory := cluster.KubernetesFactory{} - kubernetes, err := factory.Create(options.KubeConfig) + kubernetes, err := cluster.Create(options.KubeConfig) if err != nil { return } diff --git a/pkg/kt/command/exchange.go b/pkg/kt/command/exchange.go index 1badd83..d082d7a 100644 --- a/pkg/kt/command/exchange.go +++ b/pkg/kt/command/exchange.go @@ -54,8 +54,7 @@ func (action *Action) Exchange(exchange string, options *options.DaemonOptions) checkConnectRunning(options.RuntimeOptions.PidFile) - factory := cluster.KubernetesFactory{} - kubernetes, err := factory.Create(options.KubeConfig) + kubernetes, err := cluster.Create(options.KubeConfig) if err != nil { return err } diff --git a/pkg/kt/command/mesh.go b/pkg/kt/command/mesh.go index 032b37f..a0ba25c 100644 --- a/pkg/kt/command/mesh.go +++ b/pkg/kt/command/mesh.go @@ -60,8 +60,7 @@ func (action *Action) Mesh(mesh string, options *options.DaemonOptions) error { ch := SetUpCloseHandler(options) - factory := cluster.KubernetesFactory{} - kubernetes, err := factory.Create(options.KubeConfig) + kubernetes, err := cluster.Create(options.KubeConfig) if err != nil { return err } From a86e313445c7113da4dfa43998cc809a227b1bf9 Mon Sep 17 00:00:00 2001 From: yunlzheng Date: Fri, 6 Mar 2020 21:02:00 +0800 Subject: [PATCH 08/12] improve: create shadow interface --- pkg/kt/command/connect.go | 3 ++- pkg/kt/command/exchange.go | 3 ++- pkg/kt/command/mesh.go | 4 +++- pkg/kt/command/run.go | 7 ++++++- pkg/kt/connect/inbound.go | 7 +++++-- pkg/kt/connect/outbound.go | 14 ++++++++++---- pkg/kt/connect/types.go | 21 +++++++++++++++++++++ 7 files changed, 49 insertions(+), 10 deletions(-) create mode 100644 pkg/kt/connect/types.go diff --git a/pkg/kt/command/connect.go b/pkg/kt/command/connect.go index c6c805b..403621c 100644 --- a/pkg/kt/command/connect.go +++ b/pkg/kt/command/connect.go @@ -115,7 +115,8 @@ func (action *Action) Connect(options *options.DaemonOptions) (err error) { return } - err = connect.Connect(podName, endPointIP, cidrs, options) + shadow := connect.Create(options) + err = shadow.Connect(podName, endPointIP, cidrs) if err != nil { return } diff --git a/pkg/kt/command/exchange.go b/pkg/kt/command/exchange.go index d082d7a..4ba9fc0 100644 --- a/pkg/kt/command/exchange.go +++ b/pkg/kt/command/exchange.go @@ -83,7 +83,8 @@ func (action *Action) Exchange(exchange string, options *options.DaemonOptions) down := int32(0) kubernetes.Scale(app, &down) - connect.RemotePortForward(options.ExchangeOptions.Expose, options.KubeConfig, options.Namespace, podName, podIP, options.Debug) + shadow := connect.Create(options) + shadow.RemotePortForward(options.ExchangeOptions.Expose, podName, podIP) s := <-ch log.Info().Msgf("Terminal Signal is %s", s) diff --git a/pkg/kt/command/mesh.go b/pkg/kt/command/mesh.go index a0ba25c..5a41192 100644 --- a/pkg/kt/command/mesh.go +++ b/pkg/kt/command/mesh.go @@ -82,7 +82,9 @@ func (action *Action) Mesh(mesh string, options *options.DaemonOptions) error { // record context data options.RuntimeOptions.Shadow = workload - err = connect.RemotePortForward(options.MeshOptions.Expose, options.KubeConfig, options.Namespace, podName, podIP, options.Debug) + + shadow := connect.Create(options) + err = shadow.RemotePortForward(options.MeshOptions.Expose, podName, podIP) if err != nil { return err diff --git a/pkg/kt/command/run.go b/pkg/kt/command/run.go index 805c694..237e5a1 100644 --- a/pkg/kt/command/run.go +++ b/pkg/kt/command/run.go @@ -79,7 +79,12 @@ func (action *Action) Run(service string, options *options.DaemonOptions) error } options.RuntimeOptions.Shadow = service - connect.RemotePortForward(strconv.Itoa(options.RunOptions.Port), options.KubeConfig, options.Namespace, podName, podIP, options.Debug) + + shadow := connect.Create(options) + err = shadow.RemotePortForward(strconv.Itoa(options.RunOptions.Port), podName, podIP) + if err != nil { + return err + } log.Info().Msgf("forward remote %s:%v -> 127.0.0.1:%v", podIP, options.RunOptions.Port, options.RunOptions.Port) diff --git a/pkg/kt/connect/inbound.go b/pkg/kt/connect/inbound.go index 0ef3c65..775c98b 100644 --- a/pkg/kt/connect/inbound.go +++ b/pkg/kt/connect/inbound.go @@ -13,7 +13,10 @@ import ( ) // RemotePortForward mapping local port from cluster -func RemotePortForward(expose, kubeconfig, namespace, target, remoteIP string, debug bool) (err error) { +func (s *Shadow) RemotePortForward(expose, podName, remoteIP string) (err error) { + debug := s.Options.Debug + kubeConfig := s.Options.KubeConfig + namespace := s.Options.Namespace log.Info().Msgf("remote %s forward to local %s", remoteIP, expose) localSSHPort, err := strconv.Atoi(util.GetRandomSSHPort(remoteIP)) if err != nil { @@ -22,7 +25,7 @@ func RemotePortForward(expose, kubeconfig, namespace, target, remoteIP string, d var wg sync.WaitGroup wg.Add(1) go func(wg *sync.WaitGroup) { - portforward := kubectl.PortForward(kubeconfig, namespace, target, localSSHPort) + portforward := kubectl.PortForward(kubeConfig, namespace, podName, localSSHPort) err = exec.BackgroundRun(portforward, "exchange port forward to local", debug) wg.Done() }(&wg) diff --git a/pkg/kt/connect/outbound.go b/pkg/kt/connect/outbound.go index d62cec9..21591cc 100644 --- a/pkg/kt/connect/outbound.go +++ b/pkg/kt/connect/outbound.go @@ -5,8 +5,6 @@ import ( "io/ioutil" "time" - "github.com/alibaba/kt-connect/pkg/kt/options" - "github.com/alibaba/kt-connect/pkg/kt/exec" "github.com/alibaba/kt-connect/pkg/kt/exec/kubectl" "github.com/alibaba/kt-connect/pkg/kt/exec/ssh" @@ -16,12 +14,20 @@ import ( ) // Connect start vpn connection -func Connect(name, podIP string, cidrs []string, options *options.DaemonOptions) (err error) { +func (s *Shadow) Connect(name, podIP string, cidrs []string) (err error) { + options := s.Options err = util.PrepareSSHPrivateKey() if err != nil { return } - err = exec.BackgroundRun(kubectl.PortForward(options.KubeConfig, options.Namespace, name, options.ConnectOptions.SSHPort), "port-forward", options.Debug) + err = exec.BackgroundRun( + kubectl.PortForward( + options.KubeConfig, + options.Namespace, + name, + options.ConnectOptions.SSHPort), + "port-forward", + options.Debug) if err != nil { return } diff --git a/pkg/kt/connect/types.go b/pkg/kt/connect/types.go new file mode 100644 index 0000000..6b205a9 --- /dev/null +++ b/pkg/kt/connect/types.go @@ -0,0 +1,21 @@ +package connect + +import "github.com/alibaba/kt-connect/pkg/kt/options" + +// ShadowInterface shadow interface +type ShadowInterface interface { + Connect(name, podIP string, cidrs []string, options *options.DaemonOptions) (err error) + RemotePortForward(expose, kubeconfig, namespace, target, remoteIP string, debug bool) (err error) +} + +// Shadow shadow +type Shadow struct { + Options *options.DaemonOptions +} + +// Create create shadow +func Create(options *options.DaemonOptions) (shadow Shadow) { + shadow = Shadow{ + Options: options, + } +} From e4f8f6e0218ca643a43acf5b5de0c0c6fa6f974c Mon Sep 17 00:00:00 2001 From: yunlzheng Date: Fri, 6 Mar 2020 21:07:55 +0800 Subject: [PATCH 09/12] improve: rename shadow interface --- pkg/kt/command/connect.go | 2 +- pkg/kt/command/exchange.go | 2 +- pkg/kt/command/mesh.go | 2 +- pkg/kt/command/run.go | 2 +- pkg/kt/connect/inbound.go | 12 ++++++------ pkg/kt/connect/outbound.go | 4 ++-- pkg/kt/connect/types.go | 4 ++-- 7 files changed, 14 insertions(+), 14 deletions(-) diff --git a/pkg/kt/command/connect.go b/pkg/kt/command/connect.go index 403621c..70041d9 100644 --- a/pkg/kt/command/connect.go +++ b/pkg/kt/command/connect.go @@ -116,7 +116,7 @@ func (action *Action) Connect(options *options.DaemonOptions) (err error) { } shadow := connect.Create(options) - err = shadow.Connect(podName, endPointIP, cidrs) + err = shadow.Outbound(podName, endPointIP, cidrs) if err != nil { return } diff --git a/pkg/kt/command/exchange.go b/pkg/kt/command/exchange.go index 4ba9fc0..aa22c71 100644 --- a/pkg/kt/command/exchange.go +++ b/pkg/kt/command/exchange.go @@ -84,7 +84,7 @@ func (action *Action) Exchange(exchange string, options *options.DaemonOptions) kubernetes.Scale(app, &down) shadow := connect.Create(options) - shadow.RemotePortForward(options.ExchangeOptions.Expose, podName, podIP) + shadow.Inbound(options.ExchangeOptions.Expose, podName, podIP) s := <-ch log.Info().Msgf("Terminal Signal is %s", s) diff --git a/pkg/kt/command/mesh.go b/pkg/kt/command/mesh.go index 5a41192..edef13e 100644 --- a/pkg/kt/command/mesh.go +++ b/pkg/kt/command/mesh.go @@ -84,7 +84,7 @@ func (action *Action) Mesh(mesh string, options *options.DaemonOptions) error { options.RuntimeOptions.Shadow = workload shadow := connect.Create(options) - err = shadow.RemotePortForward(options.MeshOptions.Expose, podName, podIP) + err = shadow.Inbound(options.MeshOptions.Expose, podName, podIP) if err != nil { return err diff --git a/pkg/kt/command/run.go b/pkg/kt/command/run.go index 237e5a1..f34f4f1 100644 --- a/pkg/kt/command/run.go +++ b/pkg/kt/command/run.go @@ -81,7 +81,7 @@ func (action *Action) Run(service string, options *options.DaemonOptions) error options.RuntimeOptions.Shadow = service shadow := connect.Create(options) - err = shadow.RemotePortForward(strconv.Itoa(options.RunOptions.Port), podName, podIP) + err = shadow.Inbound(strconv.Itoa(options.RunOptions.Port), podName, podIP) if err != nil { return err } diff --git a/pkg/kt/connect/inbound.go b/pkg/kt/connect/inbound.go index 775c98b..32d6f20 100644 --- a/pkg/kt/connect/inbound.go +++ b/pkg/kt/connect/inbound.go @@ -12,12 +12,12 @@ import ( "github.com/rs/zerolog/log" ) -// RemotePortForward mapping local port from cluster -func (s *Shadow) RemotePortForward(expose, podName, remoteIP string) (err error) { +// Inbound mapping local port from cluster +func (s *Shadow) Inbound(exposePort, podName, remoteIP string) (err error) { debug := s.Options.Debug kubeConfig := s.Options.KubeConfig namespace := s.Options.Namespace - log.Info().Msgf("remote %s forward to local %s", remoteIP, expose) + log.Info().Msgf("remote %s forward to local %s", remoteIP, exposePort) localSSHPort, err := strconv.Atoi(util.GetRandomSSHPort(remoteIP)) if err != nil { return @@ -34,9 +34,9 @@ func (s *Shadow) RemotePortForward(expose, podName, remoteIP string) (err error) return } log.Printf("SSH Remote port-forward POD %s 22 to 127.0.0.1:%d starting\n", remoteIP, localSSHPort) - localPort := expose - remotePort := expose - ports := strings.SplitN(expose, ":", 2) + localPort := exposePort + remotePort := exposePort + ports := strings.SplitN(exposePort, ":", 2) if len(ports) > 1 { localPort = ports[1] remotePort = ports[0] diff --git a/pkg/kt/connect/outbound.go b/pkg/kt/connect/outbound.go index 21591cc..516165d 100644 --- a/pkg/kt/connect/outbound.go +++ b/pkg/kt/connect/outbound.go @@ -13,8 +13,8 @@ import ( "github.com/rs/zerolog/log" ) -// Connect start vpn connection -func (s *Shadow) Connect(name, podIP string, cidrs []string) (err error) { +// Outbound start vpn connection +func (s *Shadow) Outbound(name, podIP string, cidrs []string) (err error) { options := s.Options err = util.PrepareSSHPrivateKey() if err != nil { diff --git a/pkg/kt/connect/types.go b/pkg/kt/connect/types.go index 6b205a9..b70c4bc 100644 --- a/pkg/kt/connect/types.go +++ b/pkg/kt/connect/types.go @@ -4,8 +4,8 @@ import "github.com/alibaba/kt-connect/pkg/kt/options" // ShadowInterface shadow interface type ShadowInterface interface { - Connect(name, podIP string, cidrs []string, options *options.DaemonOptions) (err error) - RemotePortForward(expose, kubeconfig, namespace, target, remoteIP string, debug bool) (err error) + Inbound(exposePort, podName, remoteIP string) (err error) + Outbound(name, podIP string, cidrs []string) (err error) } // Shadow shadow From d162f20a0238cccc35b9f2780dbecef64672577c Mon Sep 17 00:00:00 2001 From: yunlzheng Date: Fri, 6 Mar 2020 21:14:15 +0800 Subject: [PATCH 10/12] fixed code error --- pkg/kt/connect/types.go | 1 + 1 file changed, 1 insertion(+) diff --git a/pkg/kt/connect/types.go b/pkg/kt/connect/types.go index b70c4bc..a36c1bd 100644 --- a/pkg/kt/connect/types.go +++ b/pkg/kt/connect/types.go @@ -18,4 +18,5 @@ func Create(options *options.DaemonOptions) (shadow Shadow) { shadow = Shadow{ Options: options, } + return } From 57cbb6c2c184bdf457bf82bd1b0e2165391a995b Mon Sep 17 00:00:00 2001 From: yunlzheng Date: Fri, 6 Mar 2020 22:31:43 +0800 Subject: [PATCH 11/12] improve add test case for util --- pkg/kt/util/strings_test.go | 35 +++++++++++++++++++++++++++++++++++ 1 file changed, 35 insertions(+) create mode 100644 pkg/kt/util/strings_test.go diff --git a/pkg/kt/util/strings_test.go b/pkg/kt/util/strings_test.go new file mode 100644 index 0000000..e4e2d8c --- /dev/null +++ b/pkg/kt/util/strings_test.go @@ -0,0 +1,35 @@ +package util + +import ( + "reflect" + "testing" +) + +func TestString2Map(t *testing.T) { + type args struct { + str string + } + tests := []struct { + name string + args args + want map[string]string + }{ + { + name: "should covert to key value", + args: args{ + str: "k1=v1,k2=v2", + }, + want: map[string]string{ + "k1": "v1", + "k2": "v2", + }, + }, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + if got := String2Map(tt.args.str); !reflect.DeepEqual(got, tt.want) { + t.Errorf("String2Map() = %v, want %v", got, tt.want) + } + }) + } +} From 6526667ef0fb653a9423b53a06be2fb48c3676a5 Mon Sep 17 00:00:00 2001 From: yunlzheng Date: Fri, 6 Mar 2020 22:32:43 +0800 Subject: [PATCH 12/12] improve: remove go master from travis --- .travis.yml | 1 - 1 file changed, 1 deletion(-) diff --git a/.travis.yml b/.travis.yml index b03ea30..ade911b 100644 --- a/.travis.yml +++ b/.travis.yml @@ -2,7 +2,6 @@ language: go go: - 1.13.x - - master services: - docker