From 2bef47d456852436f590c275b0c0194ce22a5784 Mon Sep 17 00:00:00 2001 From: Zexi Li Date: Fri, 2 Nov 2018 15:59:02 +0800 Subject: [PATCH] region-dns: auto k8s config from kube server --- pkg/dns/dns.go | 119 ++++++++++++++++++++++++++++++++++++++--------- pkg/dns/setup.go | 26 +++++++++-- 2 files changed, 120 insertions(+), 25 deletions(-) diff --git a/pkg/dns/dns.go b/pkg/dns/dns.go index 62fcb22f37..44382c02eb 100644 --- a/pkg/dns/dns.go +++ b/pkg/dns/dns.go @@ -7,6 +7,8 @@ import ( "fmt" "strconv" "strings" + "sync" + "time" "github.com/coredns/coredns/plugin" "github.com/coredns/coredns/plugin/etcd/msg" @@ -24,12 +26,17 @@ import ( "k8s.io/apimachinery/pkg/labels" "k8s.io/client-go/kubernetes" + "yunion.io/x/jsonutils" ylog "yunion.io/x/log" - "yunion.io/x/onecloud/pkg/cloudcommon/db" - "yunion.io/x/onecloud/pkg/compute/models" - "yunion.io/x/onecloud/pkg/util/k8s" "yunion.io/x/pkg/utils" "yunion.io/x/sqlchemy" + + "yunion.io/x/onecloud/pkg/cloudcommon/db" + "yunion.io/x/onecloud/pkg/compute/models" + "yunion.io/x/onecloud/pkg/mcclient" + "yunion.io/x/onecloud/pkg/mcclient/auth" + kubeserver "yunion.io/x/onecloud/pkg/mcclient/modules/k8s" + "yunion.io/x/onecloud/pkg/util/k8s" ) const ( @@ -62,12 +69,19 @@ type SRegionDNS struct { PrimaryZone string Upstream upstream.Upstream SqlConnection string - K8sConfigFile string - K8sClient *kubernetes.Clientset + AuthUrl string + AdminProject string + AdminUser string + AdminPassword string + Region string + k8sConfigLock *sync.RWMutex + k8sConfig string } func New() *SRegionDNS { - r := new(SRegionDNS) + r := &SRegionDNS{ + k8sConfigLock: new(sync.RWMutex), + } return r } @@ -92,19 +106,75 @@ func (r *SRegionDNS) initDB(c *caddy.Controller) error { return nil } -func (r *SRegionDNS) initK8s(c *caddy.Controller) { - cli, err := k8s.NewClientByFile(r.K8sConfigFile, nil) +func (r *SRegionDNS) getAdminSession() *mcclient.ClientSession { + return auth.GetAdminSession(r.Region, "") +} + +func (r *SRegionDNS) initAuth() { + authInfo := auth.NewAuthInfo(r.AuthUrl, "", r.AdminUser, r.AdminPassword, r.AdminProject) + auth.Init(authInfo, false, true, "", "") +} + +func (r *SRegionDNS) getKubeClusterConfig() (string, error) { + session := r.getAdminSession() + params := jsonutils.NewDict() + params.Add(jsonutils.JSONTrue, "directly") + ret, err := kubeserver.Clusters.PerformAction(session, "default", "generate-kubeconfig", params) if err != nil { - ylog.Errorf("Init kubernetes client error: %v", err) + return "", err + } + return ret.GetString("kubeconfig") +} + +func (r *SRegionDNS) getK8sConfig() string { + r.k8sConfigLock.RLock() + defer r.k8sConfigLock.RUnlock() + return r.k8sConfig +} + +func (r *SRegionDNS) setK8sConfig(conf string) { + r.k8sConfigLock.Lock() + defer r.k8sConfigLock.Unlock() + r.k8sConfig = conf +} + +func (r *SRegionDNS) isK8sHealthy() bool { + if r.getK8sConfig() == "" { + return false + } + cli, err := r.getK8sClient() + if err != nil { + return false + } + _, err = cli.Discovery().ServerVersion() + if err != nil { + ylog.Errorf("Discovery k8s version: %v", err) + return false + } + return true +} + +func (r *SRegionDNS) startRefreshKubeConfig() { + r.initAuth() + r.refreshKubeConfig() + tick := time.Tick(30 * time.Second) + for { + select { + case <-tick: + r.refreshKubeConfig() + } + } +} + +func (r *SRegionDNS) refreshKubeConfig() { + if r.isK8sHealthy() { return } - r.K8sClient = cli - pods, err := cli.CoreV1().Pods("").List(metav1.ListOptions{}) + kubeConfig, err := r.getKubeClusterConfig() if err != nil { - ylog.Errorf("Get all pods in kubernetes cluster error: %v", err) - return + ylog.Errorf("Get default k8s config from kube server error: %v", err) } - ylog.Infof("Init k8s client success, %d pods in the cluster", len(pods.Items)) + r.setK8sConfig(kubeConfig) } func (r *SRegionDNS) ServeDNS(ctx context.Context, w dns.ResponseWriter, rmsg *dns.Msg) (int, error) { @@ -251,9 +321,9 @@ func (r *SRegionDNS) getGuestIpWithName(req *recordRequest) []string { return ips } -func (r *SRegionDNS) getK8sServiceBackends(req *recordRequest) ([]string, error) { +func getK8sServiceBackends(cli *kubernetes.Clientset, req *recordRequest) ([]string, error) { queryInfo := req.GetK8sQueryInfo() - pods, err := r.getK8sServicePods(queryInfo.Namespace, queryInfo.ServiceName) + pods, err := getK8sServicePods(cli, queryInfo.Namespace, queryInfo.ServiceName) if err != nil { if k8serrors.IsNotFound(err) { err = nil @@ -270,12 +340,16 @@ func (r *SRegionDNS) getK8sServiceBackends(req *recordRequest) ([]string, error) return ips, nil } -func (r *SRegionDNS) getK8sServicePods(namespace, name string) ([]v1.Pod, error) { - cli, err := k8s.NewClientByFile(r.K8sConfigFile, nil) +func (r *SRegionDNS) getK8sClient() (*kubernetes.Clientset, error) { + cli, err := k8s.NewClientByContent([]byte(r.getK8sConfig()), nil) if err != nil { - ylog.Errorf("Init kubernetes client error: %v", err) + ylog.Warningf("Init kubernetes client error: %v", err) return nil, err } + return cli, nil +} + +func getK8sServicePods(cli *kubernetes.Clientset, namespace, name string) ([]v1.Pod, error) { svc, err := cli.CoreV1().Services(namespace).Get(name, metav1.GetOptions{}) if err != nil { return nil, err @@ -403,12 +477,13 @@ func (r *SRegionDNS) findInternalRecordIps(req *recordRequest) []string { } } - if r.K8sClient == nil { - ylog.Warningf("K8s client not ready, skip it.") + k8sCli, err := r.getK8sClient() + if err != nil { + ylog.Warningf("Get k8s client error: %v, skip it.", err) return nil } // 3. try k8s service backends - ips, err := r.getK8sServiceBackends(req) + ips, err := getK8sServiceBackends(k8sCli, req) if err != nil { ylog.Errorf("Get k8s service backends error: %v", err) } diff --git a/pkg/dns/setup.go b/pkg/dns/setup.go index 77a9cb1c4b..f753d80724 100644 --- a/pkg/dns/setup.go +++ b/pkg/dns/setup.go @@ -36,7 +36,7 @@ func setup(c *caddy.Controller) error { return plugin.Error(PluginName, err) } - rDNS.initK8s(c) + go rDNS.startRefreshKubeConfig() dnsserver.GetConfig(c).AddPlugin(func(next plugin.Handler) plugin.Handler { rDNS.Next = next @@ -77,11 +77,31 @@ func parseConfig(c *caddy.Controller) (*SRegionDNS, error) { return nil, c.ArgErr() } rDNS.SqlConnection = c.Val() - case "kube_config": + case "auth_url": if !c.NextArg() { return nil, c.ArgErr() } - rDNS.K8sConfigFile = c.Val() + rDNS.AuthUrl = c.Val() + case "admin_project": + if !c.NextArg() { + return nil, c.ArgErr() + } + rDNS.AdminProject = c.Val() + case "admin_user": + if !c.NextArg() { + return nil, c.ArgErr() + } + rDNS.AdminUser = c.Val() + case "admin_password": + if !c.NextArg() { + return nil, c.ArgErr() + } + rDNS.AdminPassword = c.Val() + case "region": + if !c.NextArg() { + return nil, c.ArgErr() + } + rDNS.Region = c.Val() case "upstream": args := c.RemainingArgs() u, err := upstream.New(args)