From 217f5a65ef449dd034b78c01bf86dc7a0dba0496 Mon Sep 17 00:00:00 2001 From: Zexi Li Date: Thu, 2 Aug 2018 20:11:39 +0800 Subject: [PATCH] - k8s client util - dns config sample --- .../root/etc/yunion/region-dns.conf.sample | 11 + cmd/region-dns/main.go | 42 +-- pkg/compute/models/dnsrecords.go | 4 +- pkg/compute/models/hostnetworks.go | 63 +++-- pkg/compute/models/hosts.go | 17 +- pkg/compute/models/networks.go | 15 +- pkg/compute/models/usage.go | 3 +- pkg/dns/dns.go | 242 +++++++++++------- pkg/dns/parse.go | 53 +++- pkg/dns/setup.go | 10 +- pkg/util/k8s/k8s.go | 64 +++++ 11 files changed, 335 insertions(+), 189 deletions(-) create mode 100644 build/region-dns/root/etc/yunion/region-dns.conf.sample create mode 100644 pkg/util/k8s/k8s.go diff --git a/build/region-dns/root/etc/yunion/region-dns.conf.sample b/build/region-dns/root/etc/yunion/region-dns.conf.sample new file mode 100644 index 0000000000..167be6e98e --- /dev/null +++ b/build/region-dns/root/etc/yunion/region-dns.conf.sample @@ -0,0 +1,11 @@ +.:53 { + cache 30 + whoami + yunion . { + sql_connection mysql+pymysql://yunioncloud:AxuaId26ZhfPXWgr@10.168.222.183:3306/yunioncloud?charset=utf8 + kube_config /home/lzx/.kube/config-hq + fallthrough . + } + proxy . 114.114.114.114:53 8.8.8.8:53 + log +} diff --git a/cmd/region-dns/main.go b/cmd/region-dns/main.go index 22b490ce8d..b5dc26adce 100644 --- a/cmd/region-dns/main.go +++ b/cmd/region-dns/main.go @@ -1,39 +1,21 @@ package main import ( - _ "github.com/coredns/coredns/plugin/auto" - _ "github.com/coredns/coredns/plugin/autopath" - _ "github.com/coredns/coredns/plugin/bind" _ "github.com/coredns/coredns/plugin/cache" _ "github.com/coredns/coredns/plugin/chaos" _ "github.com/coredns/coredns/plugin/debug" - _ "github.com/coredns/coredns/plugin/dnssec" - _ "github.com/coredns/coredns/plugin/dnstap" - _ "github.com/coredns/coredns/plugin/erratic" _ "github.com/coredns/coredns/plugin/errors" - _ "github.com/coredns/coredns/plugin/etcd" - _ "github.com/coredns/coredns/plugin/federation" _ "github.com/coredns/coredns/plugin/file" _ "github.com/coredns/coredns/plugin/forward" _ "github.com/coredns/coredns/plugin/health" _ "github.com/coredns/coredns/plugin/hosts" - _ "github.com/coredns/coredns/plugin/kubernetes" - _ "github.com/coredns/coredns/plugin/loadbalance" _ "github.com/coredns/coredns/plugin/log" _ "github.com/coredns/coredns/plugin/metrics" _ "github.com/coredns/coredns/plugin/nsid" - _ "github.com/coredns/coredns/plugin/pprof" _ "github.com/coredns/coredns/plugin/proxy" _ "github.com/coredns/coredns/plugin/reload" - _ "github.com/coredns/coredns/plugin/rewrite" - _ "github.com/coredns/coredns/plugin/root" - _ "github.com/coredns/coredns/plugin/route53" - _ "github.com/coredns/coredns/plugin/secondary" - _ "github.com/coredns/coredns/plugin/template" - _ "github.com/coredns/coredns/plugin/tls" _ "github.com/coredns/coredns/plugin/trace" _ "github.com/coredns/coredns/plugin/whoami" - _ "github.com/mholt/caddy/onevent" _ "github.com/mholt/caddy/startupshutdown" "github.com/coredns/coredns/core/dnsserver" @@ -43,43 +25,21 @@ import ( ) var directives = []string{ - "tls", - "reload", - "nsid", - "root", - "bind", "debug", "trace", "health", - "pprof", - "prometheus", "errors", "log", - "dnstap", "chaos", - "loadbalance", "cache", - "rewrite", - "dnssec", - "autopath", - "template", "hosts", - "route53", - "federation", - "kubernetes", "file", - "auto", - "secondary", - "etcd", - "redis", + "yunion", "forward", "proxy", - "erratic", "whoami", - "on", "startup", "shutdown", - "yunion", } func init() { diff --git a/pkg/compute/models/dnsrecords.go b/pkg/compute/models/dnsrecords.go index c906a9d1fd..7d8bb67237 100644 --- a/pkg/compute/models/dnsrecords.go +++ b/pkg/compute/models/dnsrecords.go @@ -7,12 +7,12 @@ import ( "github.com/yunionio/jsonutils" "github.com/yunionio/log" - "github.com/yunionio/onecloud/pkg/mcclient" - "github.com/yunionio/onecloud/pkg/httperrors" "github.com/yunionio/pkg/util/regutils" "github.com/yunionio/sqlchemy" "github.com/yunionio/onecloud/pkg/cloudcommon/db" + "github.com/yunionio/onecloud/pkg/httperrors" + "github.com/yunionio/onecloud/pkg/mcclient" ) type SDnsRecordManager struct { diff --git a/pkg/compute/models/hostnetworks.go b/pkg/compute/models/hostnetworks.go index d2bdc21a28..e5aa52dafa 100644 --- a/pkg/compute/models/hostnetworks.go +++ b/pkg/compute/models/hostnetworks.go @@ -4,9 +4,10 @@ import ( "context" "github.com/yunionio/jsonutils" - "github.com/yunionio/onecloud/pkg/mcclient" + "github.com/yunionio/sqlchemy" "github.com/yunionio/onecloud/pkg/cloudcommon/db" + "github.com/yunionio/onecloud/pkg/mcclient" ) type SHostnetworkManager struct { @@ -31,27 +32,27 @@ type SHostnetwork struct { MacAddr string `width:"18" charset:"ascii" list:"admin"` // Column(VARCHAR(18, charset='ascii')) } -func (joint *SHostnetwork) Master() db.IStandaloneModel { - return db.JointMaster(joint) +func (bn *SHostnetwork) Master() db.IStandaloneModel { + return db.JointMaster(bn) } -func (joint *SHostnetwork) Slave() db.IStandaloneModel { - return db.JointSlave(joint) +func (bn *SHostnetwork) Slave() db.IStandaloneModel { + return db.JointSlave(bn) } -func (self *SHostnetwork) GetCustomizeColumns(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject) *jsonutils.JSONDict { - extra := self.SHostJointsBase.GetCustomizeColumns(ctx, userCred, query) - extra = db.JointModelExtra(self, extra) - netif := self.GetNetInterface() +func (bn *SHostnetwork) GetCustomizeColumns(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject) *jsonutils.JSONDict { + extra := bn.SHostJointsBase.GetCustomizeColumns(ctx, userCred, query) + extra = db.JointModelExtra(bn, extra) + netif := bn.GetNetInterface() if netif != nil { extra.Add(jsonutils.NewString(netif.NicType), "nic_type") } return extra } -func (self *SHostnetwork) GetExtraDetails(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject) *jsonutils.JSONDict { - extra := self.SHostJointsBase.GetExtraDetails(ctx, userCred, query) - return db.JointModelExtra(self, extra) +func (bn *SHostnetwork) GetExtraDetails(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject) *jsonutils.JSONDict { + extra := bn.SHostJointsBase.GetExtraDetails(ctx, userCred, query) + return db.JointModelExtra(bn, extra) } func (bn *SHostnetwork) GetHost() *SHost { @@ -75,10 +76,40 @@ func (bn *SHostnetwork) GetNetInterface() *SNetInterface { return netIf } -func (self *SHostnetwork) Delete(ctx context.Context, userCred mcclient.TokenCredential) error { - return db.DeleteModel(ctx, userCred, self) +func (bn *SHostnetwork) Delete(ctx context.Context, userCred mcclient.TokenCredential) error { + return db.DeleteModel(ctx, userCred, bn) } -func (self *SHostnetwork) Detach(ctx context.Context, userCred mcclient.TokenCredential) error { - return db.DetachJoint(ctx, userCred, self) +func (bn *SHostnetwork) Detach(ctx context.Context, userCred mcclient.TokenCredential) error { + return db.DetachJoint(ctx, userCred, bn) +} + +func (man *SHostnetworkManager) QueryByAddress(addr string) *sqlchemy.SQuery { + q := HostnetworkManager.Query() + return q.Filter(sqlchemy.Equals(q.Field("ip_addr"), addr)) +} + +func (man *SHostnetworkManager) GetHostNetworkByAddress(addr string) *SHostnetwork { + network := SHostnetwork{} + err := man.QueryByAddress(addr).First(&network) + if err != nil { + return &network + } + return nil +} + +func (man *SHostnetworkManager) GetNetworkByAddress(addr string) *SNetwork { + net := man.GetHostNetworkByAddress(addr) + if net == nil { + return nil + } + return net.GetNetwork() +} + +func (man *SHostnetworkManager) GetHostByAddress(addr string) *SHost { + net := man.GetHostNetworkByAddress(addr) + if net == nil { + return nil + } + return net.GetHost() } diff --git a/pkg/compute/models/hosts.go b/pkg/compute/models/hosts.go index 32dcdf3d80..a784d4375c 100644 --- a/pkg/compute/models/hosts.go +++ b/pkg/compute/models/hosts.go @@ -9,13 +9,6 @@ import ( "github.com/yunionio/jsonutils" "github.com/yunionio/log" - "github.com/yunionio/onecloud/pkg/mcclient" - "github.com/yunionio/onecloud/pkg/mcclient/auth" - "github.com/yunionio/onecloud/pkg/mcclient/modules" - "github.com/yunionio/onecloud/pkg/cloudcommon/db" - "github.com/yunionio/onecloud/pkg/cloudprovider" - "github.com/yunionio/onecloud/pkg/compute/options" - "github.com/yunionio/onecloud/pkg/httperrors" "github.com/yunionio/pkg/tristate" "github.com/yunionio/pkg/util/compare" "github.com/yunionio/pkg/util/netutils" @@ -23,6 +16,14 @@ import ( "github.com/yunionio/pkg/util/sysutils" "github.com/yunionio/pkg/utils" "github.com/yunionio/sqlchemy" + + "github.com/yunionio/onecloud/pkg/cloudcommon/db" + "github.com/yunionio/onecloud/pkg/cloudprovider" + "github.com/yunionio/onecloud/pkg/compute/options" + "github.com/yunionio/onecloud/pkg/httperrors" + "github.com/yunionio/onecloud/pkg/mcclient" + "github.com/yunionio/onecloud/pkg/mcclient/auth" + "github.com/yunionio/onecloud/pkg/mcclient/modules" ) const ( @@ -1109,7 +1110,7 @@ func (self *SHost) SyncHostVMs(ctx context.Context, userCred mcclient.TokenCrede } func (self *SHost) getNetworkOfIPOnHost(ipAddr string) (*SNetwork, error) { - net, err := NetworkManager.getNetworkOfIP(ipAddr, "", tristate.None) + net, err := NetworkManager.GetNetworkOfIP(ipAddr, "", tristate.None) if err != nil { return nil, err } diff --git a/pkg/compute/models/networks.go b/pkg/compute/models/networks.go index 0fbfea6808..c2e530d8e7 100644 --- a/pkg/compute/models/networks.go +++ b/pkg/compute/models/networks.go @@ -8,12 +8,6 @@ import ( "github.com/yunionio/jsonutils" "github.com/yunionio/log" - "github.com/yunionio/onecloud/pkg/mcclient" - "github.com/yunionio/onecloud/pkg/cloudcommon/db" - "github.com/yunionio/onecloud/pkg/cloudcommon/db/taskman" - "github.com/yunionio/onecloud/pkg/cloudprovider" - "github.com/yunionio/onecloud/pkg/compute/options" - "github.com/yunionio/onecloud/pkg/httperrors" "github.com/yunionio/pkg/tristate" "github.com/yunionio/pkg/util/compare" "github.com/yunionio/pkg/util/fileutils" @@ -21,6 +15,13 @@ import ( "github.com/yunionio/pkg/util/regutils" "github.com/yunionio/pkg/utils" "github.com/yunionio/sqlchemy" + + "github.com/yunionio/onecloud/pkg/cloudcommon/db" + "github.com/yunionio/onecloud/pkg/cloudcommon/db/taskman" + "github.com/yunionio/onecloud/pkg/cloudprovider" + "github.com/yunionio/onecloud/pkg/compute/options" + "github.com/yunionio/onecloud/pkg/httperrors" + "github.com/yunionio/onecloud/pkg/mcclient" ) const ( @@ -507,7 +508,7 @@ func (self *SNetwork) isAddressUsed(address string) bool { return false } -func (manager *SNetworkManager) getNetworkOfIP(ipAddr string, serverType string, isPublic tristate.TriState) (*SNetwork, error) { +func (manager *SNetworkManager) GetNetworkOfIP(ipAddr string, serverType string, isPublic tristate.TriState) (*SNetwork, error) { address, err := netutils.NewIPV4Addr(ipAddr) if err != nil { return nil, err diff --git a/pkg/compute/models/usage.go b/pkg/compute/models/usage.go index 28c619a0cd..27d05a2228 100644 --- a/pkg/compute/models/usage.go +++ b/pkg/compute/models/usage.go @@ -2,8 +2,9 @@ package models import ( "github.com/yunionio/log" - "github.com/yunionio/onecloud/pkg/cloudcommon/db" "github.com/yunionio/sqlchemy" + + "github.com/yunionio/onecloud/pkg/cloudcommon/db" ) func AttachUsageQuery( diff --git a/pkg/dns/dns.go b/pkg/dns/dns.go index f8d6782c12..321fc76eb4 100644 --- a/pkg/dns/dns.go +++ b/pkg/dns/dns.go @@ -5,25 +5,29 @@ import ( "database/sql" "errors" "fmt" - //"os" "github.com/coredns/coredns/plugin" "github.com/coredns/coredns/plugin/etcd/msg" "github.com/coredns/coredns/plugin/pkg/dnsutil" "github.com/coredns/coredns/plugin/pkg/fall" + "github.com/coredns/coredns/plugin/pkg/upstream" "github.com/coredns/coredns/request" _ "github.com/go-sql-driver/mysql" - //clog "github.com/coredns/coredns/plugin/log" - "github.com/coredns/coredns/plugin/pkg/upstream" - //"github.com/coredns/coredns/plugin/pkg/replacer" "github.com/mholt/caddy" "github.com/miekg/dns" + v1 "k8s.io/api/core/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/fields" + "k8s.io/apimachinery/pkg/labels" + "k8s.io/client-go/kubernetes" + ylog "github.com/yunionio/log" "github.com/yunionio/pkg/utils" "github.com/yunionio/sqlchemy" "github.com/yunionio/onecloud/pkg/cloudcommon/db" "github.com/yunionio/onecloud/pkg/compute/models" + "github.com/yunionio/onecloud/pkg/util/k8s" ) const ( @@ -51,6 +55,7 @@ type SRegionDNS struct { Upstream upstream.Upstream SqlConnection string K8sConfigFile string + K8sClient *kubernetes.Clientset } func New() *SRegionDNS { @@ -77,8 +82,19 @@ func (r *SRegionDNS) initDB(c *caddy.Controller) error { return nil } -func (r *SRegionDNS) initK8s(c *caddy.Controller) error { - return nil +func (r *SRegionDNS) initK8s(c *caddy.Controller) { + cli, err := k8s.NewClientByFile(r.K8sConfigFile, nil) + if err != nil { + ylog.Errorf("Init kubernetes client error: %v", err) + return + } + r.K8sClient = cli + pods, err := cli.CoreV1().Pods("").List(metav1.ListOptions{}) + if err != nil { + ylog.Errorf("Get all pods in kubernetes cluster error: %v", err) + return + } + ylog.Infof("Init k8s client success, %d pods in the cluster", len(pods.Items)) } func (r *SRegionDNS) CloseDB() { @@ -86,11 +102,6 @@ func (r *SRegionDNS) CloseDB() { } func (r *SRegionDNS) ServeDNS(ctx context.Context, w dns.ResponseWriter, rmsg *dns.Msg) (int, error) { - //rrw := dnstest.NewRecorder(w) - //rep := replacer.New(r, rrw, corelog.CommonLogEmptyValue) - //log.Infof("%v", rep.Replace(format)) - //fmt.Fprintln(output, rep.Replace(format)) - //count := models.DnsRecordManager.QueryDns("", "drone.yunion.io") var ( records []dns.RR extra []dns.RR @@ -102,45 +113,45 @@ func (r *SRegionDNS) ServeDNS(ctx context.Context, w dns.ResponseWriter, rmsg *d //isMyDomain := true zone := plugin.Zones(r.Zones).Matches(state.Name()) - if zone == "" { - //isMyDomain = false - return plugin.NextOrFailure(r.Name(), r.Next, ctx, w, rmsg) - } + //if zone == "" { + ////isMyDomain = false + //return plugin.NextOrFailure(r.Name(), r.Next, ctx, w, rmsg) + //} switch state.QType() { case dns.TypeA: - ylog.Debugf("===A question: %#v", state) + ylog.Debugf("A question: %#v", state) records, err = plugin.A(r, zone, state, nil, opt) case dns.TypeAAAA: - ylog.Debugf("===AAAA question: %#v", state) + ylog.Debugf("AAAA question: %#v", state) records, err = plugin.AAAA(r, zone, state, nil, opt) case dns.TypeTXT: - ylog.Debugf("===TXT question: %#v", state) + ylog.Debugf("TXT question: %#v", state) records, err = plugin.TXT(r, zone, state, opt) case dns.TypeCNAME: - ylog.Debugf("===CNAME question: %#v", state) + ylog.Debugf("CNAME question: %#v", state) records, err = plugin.CNAME(r, zone, state, opt) case dns.TypePTR: - ylog.Debugf("===PTR question: %#v", state) + ylog.Debugf("PTR question: %#v", state) records, err = plugin.PTR(r, zone, state, opt) case dns.TypeMX: - ylog.Debugf("===MX question: %#v", state) + ylog.Debugf("MX question: %#v", state) records, extra, err = plugin.MX(r, zone, state, opt) case dns.TypeSRV: - ylog.Debugf("===SRV question: %#v", state) + ylog.Debugf("SRV question: %#v", state) records, extra, err = plugin.SRV(r, zone, state, opt) case dns.TypeSOA: - ylog.Debugf("===SOA question: %#v", state) + ylog.Debugf("SOA question: %#v", state) records, err = plugin.SOA(r, zone, state, opt) case dns.TypeNS: - ylog.Debugf("===NS question: %#v", state) + ylog.Debugf("NS question: %#v", state) if state.Name() == zone { records, extra, err = plugin.NS(r, zone, state, opt) break } fallthrough default: - ylog.Infof("=== not processed state: %#v", state) + ylog.Warningf("Not processed state: %#v", state) // Do a fake A lookup, so we can distinguish between NODATA and NXDOMAIN _, err = plugin.A(r, zone, state, nil, opt) } @@ -152,7 +163,6 @@ func (r *SRegionDNS) ServeDNS(ctx context.Context, w dns.ResponseWriter, rmsg *d return plugin.BackendError(r, zone, dns.RcodeNameError, state, nil /* err */, opt) } if err != nil { - //return dns.RcodeServerFailure, err return plugin.BackendError(r, zone, dns.RcodeServerFailure, state, err, opt) } @@ -178,11 +188,6 @@ var ( // Services implements the ServiceBackend interface func (r *SRegionDNS) Services(state request.Request, exact bool, opt plugin.Options) (services []msg.Service, err error) { - //services, err = r.Records(state, exact) - //if err != nil { - //return - //} - //services = msg.Group switch state.QType() { case dns.TypeTXT: @@ -213,6 +218,7 @@ func (r *SRegionDNS) Services(state request.Request, exact bool, opt plugin.Opti } s, e := r.Records(state, false) + ylog.Debugf("Get records: %#v, error: %v", s, e) // SRV is not yet implemented, so remove those records. if state.QType() != dns.TypeSRV { @@ -252,79 +258,66 @@ func (r *SRegionDNS) Records(state request.Request, exact bool) ([]msg.Service, return r.findRecords(req) } +func (r *SRegionDNS) getHostIpWithName(req *recordRequest) []string { + name := req.QueryName() + host, _ := models.HostManager.FetchByName("", name) + if host == nil { + return nil + } + ip := host.(*models.SHost).AccessIp + if len(ip) == 0 { + return nil + } + return []string{ip} +} + func (r *SRegionDNS) getGuestIpWithName(req *recordRequest) []string { ips := []string{} - name := req.GuestName() + name := req.QueryName() projectId := req.ProjectId() isExitOnly := req.IsExitOnly() - ylog.Warningf("===========args: %q, %q, %v", projectId, name, isExitOnly) ips = models.GuestManager.GetIpInProjectWithName(projectId, name, isExitOnly) return ips } +func (r *SRegionDNS) getK8sServiceBackends(req *recordRequest) ([]string, error) { + queryInfo := req.GetK8sQueryInfo() + pods, err := r.getK8sServicePods(queryInfo.Namespace, queryInfo.ServiceName) + if err != nil { + return nil, err + } + ips := make([]string, 0) + for _, pod := range pods { + ip := pod.Status.PodIP + if len(ip) != 0 { + ips = append(ips, ip) + } + } + return ips, nil +} + +func (r *SRegionDNS) getK8sServicePods(namespace, name string) ([]v1.Pod, error) { + cli := r.K8sClient + svc, err := cli.CoreV1().Services(namespace).Get(name, metav1.GetOptions{}) + if err != nil { + return nil, err + } + labelSelector := labels.SelectorFromSet(svc.Spec.Selector) + pods, err := cli.CoreV1().Pods(namespace).List(metav1.ListOptions{ + LabelSelector: labelSelector.String(), + FieldSelector: fields.Everything().String(), + }) + if err != nil { + return nil, err + } + return pods.Items, nil +} + func (r *SRegionDNS) Name() string { return PluginName } -type SGuestInfo struct { - *models.SGuest -} - -func (g *SGuestInfo) GetProjectId() string { - return g.SGuest.ProjectId -} - -func (g *SGuestInfo) GetGuestId() string { - return g.SGuest.Id -} - -func (g *SGuestInfo) IsExitOnly() bool { - return g.SGuest.IsExitOnly() -} - -func NewGuestInfoByAddress(address string) *SGuestInfo { - guest := models.GuestnetworkManager.GetGuestByAddress(address) - if guest == nil { - return nil - } - return &SGuestInfo{SGuest: guest} -} - -func (r *SRegionDNS) findRecords(req *recordRequest) (recs []msg.Service, err error) { - isMyDomain := false - zone := plugin.Zones(r.Zones).Matches(req.state.Name()) - if zone != "" { - isMyDomain = true - } - //isPrivateAddr := false - if isMyDomain { - return r.findInternalRecords(req) - } - return r.findExternalRecords(req) -} - -func (r *SRegionDNS) findInternalRecords(req *recordRequest) ([]msg.Service, error) { - ylog.Debugf("=======findInternalRecords, srcip: %q, projectId: %q", req.SrcIP4(), req.ProjectId()) - if req.ProjectId() == "" { - return nil, nil - } - // first, try dns records - svcs, _ := r.findLocalRecords(req) - if len(svcs) != 0 { - return svcs, nil - } - // second, try guest table - ips := r.getGuestIpWithName(req) - return ips2DnsRecords(ips), nil -} - -func (r *SRegionDNS) findExternalRecords(req *recordRequest) ([]msg.Service, error) { - srcIP := req.SrcIP4() - ylog.Debugf("Get client ip: %q", srcIP) - return r.findLocalRecords(req) -} - -func (r *SRegionDNS) findLocalRecords(req *recordRequest) (recs []msg.Service, err error) { +func (r *SRegionDNS) queryLocalDnsRecords(req *recordRequest) (recs []msg.Service, err error) { ips := models.DnsRecordManager.QueryDnsIps(req.ProjectId(), req.Name(), req.Type()) if len(ips) == 0 { err = errNoItems @@ -337,6 +330,71 @@ func (r *SRegionDNS) findLocalRecords(req *recordRequest) (recs []msg.Service, e return } +func (r *SRegionDNS) IsCloudNetworkIp(req *recordRequest) bool { + if req.network != nil { + return true + } + return false +} + +func (r *SRegionDNS) IsK8sClientReady() bool { + return r.K8sClient != nil +} + +func (r *SRegionDNS) findRecords(req *recordRequest) (recs []msg.Service, err error) { + // 1. try local dns records table + recs, err = r.queryLocalDnsRecords(req) + if len(recs) != 0 { + return + } + + isMyDomain := false + zone := plugin.Zones(r.Zones).Matches(req.state.Name()) + if zone != "" { + isMyDomain = true + } + + isCloudIp := r.IsCloudNetworkIp(req) + + // 2. not my domain and src ip not in cloud network table + // query from upstream + if !isMyDomain && !isCloudIp { + err = errNoItems + return + } + + // 3. internal query + ips, err := r.findInternalRecordIps(req) + return ips2DnsRecords(ips), err +} + +func (r *SRegionDNS) findInternalRecordIps(req *recordRequest) ([]string, error) { + // 1. try host table + ip := r.getHostIpWithName(req) + if len(ip) != 0 { + return ip, nil + } + // 2. try guest table + ip = r.getGuestIpWithName(req) + if len(ip) != 0 { + return ip, nil + } + + if !r.IsK8sClientReady() { + ylog.Warningf("K8s client not ready, skip it.") + return nil, errNoItems + } + // 3. try k8s service backends + ips, err := r.getK8sServiceBackends(req) + if len(ips) != 0 { + return ips, nil + } + if err != nil { + ylog.Errorf("Get k8s service backends error: %v", err) + } + return nil, errNoItems +} + func ips2DnsRecords(ips []string) []msg.Service { recs := make([]msg.Service, 0) for _, ip := range ips { diff --git a/pkg/dns/parse.go b/pkg/dns/parse.go index 82bb791f41..92a7f586d5 100644 --- a/pkg/dns/parse.go +++ b/pkg/dns/parse.go @@ -7,14 +7,17 @@ import ( "github.com/coredns/coredns/request" "github.com/miekg/dns" - "github.com/yunionio/log" + "github.com/yunionio/pkg/tristate" + + "github.com/yunionio/onecloud/pkg/compute/models" ) type recordRequest struct { state request.Request domainSegs []string - projectId string - guestInfo *SGuestInfo + guest *models.SGuest + host *models.SHost + network *models.SNetwork } func parseRequest(state request.Request) (r *recordRequest, err error) { @@ -25,20 +28,20 @@ func parseRequest(state request.Request) (r *recordRequest, err error) { domainSegs: segs, } srcIP := r.SrcIP4() - guestInfo := NewGuestInfoByAddress(srcIP) - r.guestInfo = guestInfo + r.guest = models.GuestnetworkManager.GetGuestByAddress(srcIP) + r.host = models.HostnetworkManager.GetHostByAddress(srcIP) + r.network, _ = models.NetworkManager.GetNetworkOfIP(srcIP, "", tristate.None) return } func (r recordRequest) Name() string { //fullName, _ := dnsutil.TrimZone(r.state.Name(), "") name := r.state.Name() - log.Errorf("==name: %q", name) name = strings.TrimSuffix(name, ".") return name } -func (r recordRequest) GuestName() string { +func (r recordRequest) QueryName() string { seps := strings.Split(r.Name(), ".") if len(seps) == 0 { return "" @@ -52,20 +55,44 @@ func (r recordRequest) Type() string { func (r recordRequest) SrcIP4() string { ip := r.state.IP() - log.Debugf("Source ip: %q, guestName: %q", ip, r.GuestName()) return ip } func (r recordRequest) ProjectId() string { - if r.guestInfo == nil { - return "" + if r.guest != nil { + return r.guest.ProjectId } - return r.guestInfo.GetProjectId() + if r.network != nil { + return r.network.ProjectId + } + return "" } func (r recordRequest) IsExitOnly() bool { - if r.guestInfo == nil { + if r.guest == nil { return false } - return r.guestInfo.IsExitOnly() + return r.guest.IsExitOnly() +} + +type K8sQueryInfo struct { + ServiceName string + Namespace string +} + +func (r recordRequest) GetK8sQueryInfo() K8sQueryInfo { + parts := strings.SplitN(r.Name(), ".", 3) + var svcName string + var namespace string + if len(parts) >= 2 { + svcName = parts[0] + namespace = parts[1] + } else { + svcName = parts[0] + namespace = "default" + } + return K8sQueryInfo{ + ServiceName: svcName, + Namespace: namespace, + } } diff --git a/pkg/dns/setup.go b/pkg/dns/setup.go index 6915c9de89..04e6ea00c7 100644 --- a/pkg/dns/setup.go +++ b/pkg/dns/setup.go @@ -7,8 +7,6 @@ import ( "github.com/coredns/coredns/plugin" "github.com/coredns/coredns/plugin/pkg/upstream" "github.com/mholt/caddy" - - "github.com/yunionio/log" ) func init() { @@ -25,17 +23,13 @@ func setup(c *caddy.Controller) error { if err != nil { return plugin.Error(PluginName, err) } - log.Infof("regionDNSParse succ: %#v", rDNS) err = rDNS.initDB(c) if err != nil { return plugin.Error(PluginName, err) } - err = rDNS.initK8s(c) - if err != nil { - return plugin.Error(PluginName, err) - } + rDNS.initK8s(c) dnsserver.GetConfig(c).AddPlugin(func(next plugin.Handler) plugin.Handler { rDNS.Next = next @@ -60,11 +54,9 @@ func parseConfig(c *caddy.Controller) (*SRegionDNS, error) { for i, str := range rDNS.Zones { rDNS.Zones[i] = plugin.Host(str).Normalize() } - log.Warningf("==zones: %v", rDNS.Zones) if c.NextBlock() { for { - log.Printf("===val: %v", c.Val()) switch c.Val() { case "fallthrough": rDNS.Fall.SetZonesFromArgs(c.RemainingArgs()) diff --git a/pkg/util/k8s/k8s.go b/pkg/util/k8s/k8s.go new file mode 100644 index 0000000000..eaf8284dc5 --- /dev/null +++ b/pkg/util/k8s/k8s.go @@ -0,0 +1,64 @@ +package k8s + +import ( + "fmt" + "net/http" + "time" + + "k8s.io/client-go/kubernetes" + "k8s.io/client-go/rest" + "k8s.io/client-go/tools/clientcmd" +) + +const ( + K8sWrapTransportTimeout = 30 +) + +type WrapTransport func(rt http.RoundTripper) http.RoundTripper + +func NewClientByFile(kubeConfigPath string, k8sWrapTransport WrapTransport) (*kubernetes.Clientset, error) { + config, err := clientcmd.BuildConfigFromFlags("", kubeConfigPath) + if err != nil { + return nil, err + } + return kubernetes.NewForConfig(setConfigField(config, k8sWrapTransport)) +} + +func GetK8sClientConfig(kubeConfig []byte) (*rest.Config, error) { + var config *rest.Config + var err error + if kubeConfig != nil { + apiconfig, err := clientcmd.Load(kubeConfig) + if err != nil { + return nil, err + } + + clientConfig := clientcmd.NewDefaultClientConfig(*apiconfig, &clientcmd.ConfigOverrides{}) + config, err = clientConfig.ClientConfig() + if err != nil { + return nil, err + } + } else { + return nil, fmt.Errorf("kubeconfig value is nil") + } + if err != nil { + return nil, fmt.Errorf("create kubernetes config failed: %v", err) + } + return config, nil +} + +func NewClientByContent(kubeConfig []byte, k8sWrapTransport WrapTransport) (*kubernetes.Clientset, error) { + config, err := GetK8sClientConfig(kubeConfig) + if err != nil { + return nil, fmt.Errorf("Create kubernetes config: %v", err) + } + return kubernetes.NewForConfig(setConfigField(config, k8sWrapTransport)) +} + +func setConfigField(c *rest.Config, tr WrapTransport) *rest.Config { + if tr != nil { + c.WrapTransport = tr + } + c.Timeout = time.Second * time.Duration(K8sWrapTransportTimeout) + return c +}