region-dns: use k8s.SKubeClusterManager

This commit is contained in:
Zexi Li
2018-11-15 00:26:56 +08:00
parent 33f744e21a
commit 43255fe7e8
3 changed files with 113 additions and 77 deletions
+9 -76
View File
@@ -7,7 +7,6 @@ import (
"fmt"
"strconv"
"strings"
"sync"
"time"
"github.com/coredns/coredns/plugin"
@@ -26,7 +25,6 @@ import (
"k8s.io/apimachinery/pkg/labels"
"k8s.io/client-go/kubernetes"
"yunion.io/x/jsonutils"
ylog "yunion.io/x/log"
"yunion.io/x/pkg/utils"
"yunion.io/x/sqlchemy"
@@ -35,7 +33,6 @@ import (
"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"
)
@@ -74,14 +71,11 @@ type SRegionDNS struct {
AdminUser string
AdminPassword string
Region string
k8sConfigLock *sync.RWMutex
k8sConfig string
K8sManager *k8s.SKubeClusterManager
}
func New() *SRegionDNS {
r := &SRegionDNS{
k8sConfigLock: new(sync.RWMutex),
}
r := &SRegionDNS{}
return r
}
@@ -106,6 +100,12 @@ func (r *SRegionDNS) initDB(c *caddy.Controller) error {
return nil
}
func (r *SRegionDNS) initK8s() {
r.initAuth()
r.K8sManager = k8s.NewKubeClusterManager(r.Region, 30*time.Second)
r.K8sManager.Start()
}
func (r *SRegionDNS) getAdminSession() *mcclient.ClientSession {
return auth.GetAdminSession(r.Region, "")
}
@@ -115,68 +115,6 @@ func (r *SRegionDNS) initAuth() {
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 {
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
}
kubeConfig, err := r.getKubeClusterConfig()
if err != nil {
ylog.Errorf("Get default k8s config from kube server error: %v", err)
}
r.setK8sConfig(kubeConfig)
}
func (r *SRegionDNS) ServeDNS(ctx context.Context, w dns.ResponseWriter, rmsg *dns.Msg) (int, error) {
var (
records []dns.RR
@@ -341,12 +279,7 @@ func getK8sServiceBackends(cli *kubernetes.Clientset, req *recordRequest) ([]str
}
func (r *SRegionDNS) getK8sClient() (*kubernetes.Clientset, error) {
cli, err := k8s.NewClientByContent([]byte(r.getK8sConfig()), nil)
if err != nil {
ylog.Warningf("Init kubernetes client error: %v", err)
return nil, err
}
return cli, nil
return r.K8sManager.GetK8sClient()
}
func getK8sServicePods(cli *kubernetes.Clientset, namespace, name string) ([]v1.Pod, error) {
+1 -1
View File
@@ -36,7 +36,7 @@ func setup(c *caddy.Controller) error {
return plugin.Error(PluginName, err)
}
go rDNS.startRefreshKubeConfig()
go rDNS.initK8s()
dnsserver.GetConfig(c).AddPlugin(func(next plugin.Handler) plugin.Handler {
rDNS.Next = next
+103
View File
@@ -0,0 +1,103 @@
package k8s
import (
"sync"
"time"
"k8s.io/client-go/kubernetes"
"yunion.io/x/jsonutils"
"yunion.io/x/log"
"yunion.io/x/onecloud/pkg/mcclient/auth"
kubeserver "yunion.io/x/onecloud/pkg/mcclient/modules/k8s"
)
type SKubeClusterManager struct {
k8sConfigLock *sync.RWMutex
k8sConfig string
interval time.Duration
region string
}
func NewKubeClusterManager(region string, interval time.Duration) *SKubeClusterManager {
return &SKubeClusterManager{
k8sConfigLock: new(sync.RWMutex),
interval: interval,
region: region,
}
}
func (man *SKubeClusterManager) GetK8sConfig() string {
man.k8sConfigLock.RLock()
defer man.k8sConfigLock.RUnlock()
return man.k8sConfig
}
func (man *SKubeClusterManager) GetK8sClient() (*kubernetes.Clientset, error) {
cli, err := NewClientByContent([]byte(man.GetK8sConfig()), nil)
if err != nil {
log.Warningf("Init kubernetes client error: %v", err)
return nil, err
}
return cli, nil
}
func (man *SKubeClusterManager) Start() {
go man.startRefreshKubeConfig()
}
func (man *SKubeClusterManager) setK8sConfig(conf string) {
man.k8sConfigLock.Lock()
defer man.k8sConfigLock.Unlock()
man.k8sConfig = conf
}
func (man *SKubeClusterManager) isK8sHealthy() bool {
if man.GetK8sConfig() == "" {
return false
}
cli, err := man.GetK8sClient()
if err != nil {
return false
}
_, err = cli.Discovery().ServerVersion()
if err != nil {
log.Errorf("Discovery k8s version: %v", err)
return false
}
return true
}
func (man *SKubeClusterManager) startRefreshKubeConfig() {
man.refreshKubeConfig()
tick := time.Tick(man.interval)
for {
select {
case <-tick:
man.refreshKubeConfig()
}
}
}
func (man *SKubeClusterManager) refreshKubeConfig() {
if man.isK8sHealthy() {
return
}
kubeConfig, err := man.getKubeClusterConfig()
if err != nil {
log.Errorf("Get default k8s config from kube server error: %v", err)
}
man.setK8sConfig(kubeConfig)
}
func (man *SKubeClusterManager) getKubeClusterConfig() (string, error) {
session := auth.GetAdminSession(man.region, "v1")
params := jsonutils.NewDict()
params.Add(jsonutils.JSONTrue, "directly")
ret, err := kubeserver.Clusters.PerformAction(session, "default", "generate-kubeconfig", params)
if err != nil {
return "", err
}
return ret.GetString("kubeconfig")
}