mirror of
https://github.com/yunionio/cloudpods.git
synced 2026-09-01 15:07:17 +08:00
dns: policy, log, db query optimization
Fixes CLOUD-451 TODO - timeout bound for k8s query - k8s client pools
This commit is contained in:
+75
-86
@@ -18,25 +18,27 @@ import (
|
||||
"github.com/mholt/caddy"
|
||||
"github.com/miekg/dns"
|
||||
v1 "k8s.io/api/core/v1"
|
||||
k8serrors "k8s.io/apimachinery/pkg/api/errors"
|
||||
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||||
"k8s.io/apimachinery/pkg/fields"
|
||||
"k8s.io/apimachinery/pkg/labels"
|
||||
"k8s.io/client-go/kubernetes"
|
||||
|
||||
ylog "yunion.io/x/log"
|
||||
"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/util/k8s"
|
||||
"yunion.io/x/pkg/utils"
|
||||
"yunion.io/x/sqlchemy"
|
||||
)
|
||||
|
||||
const (
|
||||
PluginName string = "yunion"
|
||||
|
||||
// defaultTTL to apply to all answers
|
||||
defaultTTL = 10
|
||||
defaultTTL = 10
|
||||
defaultDbMaxOpenConn = 32
|
||||
defaultDbMaxIdleConn = 32
|
||||
)
|
||||
|
||||
var (
|
||||
@@ -74,15 +76,17 @@ func (r *SRegionDNS) initDB(c *caddy.Controller) error {
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
dbConn, err := sql.Open(dialect, sqlStr)
|
||||
sqlDb, err := sql.Open(dialect, sqlStr)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
sqlchemy.SetDB(dbConn)
|
||||
sqlDb.SetMaxOpenConns(defaultDbMaxOpenConn)
|
||||
sqlDb.SetMaxIdleConns(defaultDbMaxIdleConn)
|
||||
sqlchemy.SetDB(sqlDb)
|
||||
db.InitAllManagers()
|
||||
|
||||
c.OnShutdown(func() error {
|
||||
r.CloseDB()
|
||||
sqlchemy.CloseDB()
|
||||
return nil
|
||||
})
|
||||
return nil
|
||||
@@ -103,10 +107,6 @@ func (r *SRegionDNS) initK8s(c *caddy.Controller) {
|
||||
ylog.Infof("Init k8s client success, %d pods in the cluster", len(pods.Items))
|
||||
}
|
||||
|
||||
func (r *SRegionDNS) CloseDB() {
|
||||
sqlchemy.CloseDB()
|
||||
}
|
||||
|
||||
func (r *SRegionDNS) ServeDNS(ctx context.Context, w dns.ResponseWriter, rmsg *dns.Msg) (int, error) {
|
||||
var (
|
||||
records []dns.RR
|
||||
@@ -116,36 +116,26 @@ func (r *SRegionDNS) ServeDNS(ctx context.Context, w dns.ResponseWriter, rmsg *d
|
||||
|
||||
opt := plugin.Options{}
|
||||
state := request.Request{W: w, Req: rmsg, Context: ctx}
|
||||
|
||||
zone := plugin.Zones(r.Zones).Matches(state.Name())
|
||||
|
||||
switch state.QType() {
|
||||
case dns.TypeA:
|
||||
ylog.Debugf("A question: %#v", state)
|
||||
records, err = plugin.A(r, zone, state, nil, opt)
|
||||
case dns.TypeAAAA:
|
||||
ylog.Debugf("AAAA question: %#v", state)
|
||||
// TODO fallthrough to next
|
||||
records, err = plugin.AAAA(r, zone, state, nil, opt)
|
||||
case dns.TypeTXT:
|
||||
ylog.Debugf("TXT question: %#v", state)
|
||||
records, err = plugin.TXT(r, zone, state, opt)
|
||||
case dns.TypeCNAME:
|
||||
ylog.Debugf("CNAME question: %#v", state)
|
||||
records, err = plugin.CNAME(r, zone, state, opt)
|
||||
case dns.TypePTR:
|
||||
ylog.Debugf("PTR question: %#v", state)
|
||||
records, err = plugin.PTR(r, zone, state, opt)
|
||||
case dns.TypeMX:
|
||||
ylog.Debugf("MX question: %#v", state)
|
||||
records, extra, err = plugin.MX(r, zone, state, opt)
|
||||
case dns.TypeSRV:
|
||||
ylog.Debugf("SRV question: %#v", state)
|
||||
records, extra, err = plugin.SRV(r, zone, state, opt)
|
||||
case dns.TypeSOA:
|
||||
ylog.Debugf("SOA question: %#v", state)
|
||||
records, err = plugin.SOA(r, zone, state, opt)
|
||||
case dns.TypeNS:
|
||||
ylog.Debugf("NS question: %#v", state)
|
||||
if state.Name() == zone {
|
||||
records, extra, err = plugin.NS(r, zone, state, opt)
|
||||
break
|
||||
@@ -157,18 +147,19 @@ func (r *SRegionDNS) ServeDNS(ctx context.Context, w dns.ResponseWriter, rmsg *d
|
||||
_, err = plugin.A(r, zone, state, nil, opt)
|
||||
}
|
||||
|
||||
if r.IsNameError(err) {
|
||||
if err == errCallNext {
|
||||
if r.Fall.Through(state.Name()) {
|
||||
return plugin.NextOrFailure(r.Name(), r.Next, ctx, w, rmsg)
|
||||
}
|
||||
return plugin.BackendError(r, zone, dns.RcodeNameError, state, nil /* err */, opt)
|
||||
}
|
||||
if err != nil {
|
||||
return plugin.BackendError(r, zone, dns.RcodeServerFailure, state, err, opt)
|
||||
} else if err == errRefused {
|
||||
return plugin.BackendError(r, zone, dns.RcodeRefused, state, err, opt)
|
||||
} else if err == errNotFound {
|
||||
return plugin.BackendError(r, zone, dns.RcodeNameError, state, err, opt)
|
||||
}
|
||||
|
||||
if len(records) == 0 {
|
||||
return plugin.BackendError(r, zone, dns.RcodeSuccess, state, err, opt)
|
||||
return plugin.BackendError(r, zone, dns.RcodeNameError, state, err, opt)
|
||||
}
|
||||
|
||||
m := new(dns.Msg)
|
||||
@@ -184,13 +175,14 @@ func (r *SRegionDNS) ServeDNS(ctx context.Context, w dns.ResponseWriter, rmsg *d
|
||||
}
|
||||
|
||||
var (
|
||||
errNoItems = errors.New("no items found")
|
||||
errRefused = errors.New("refused the query")
|
||||
errNotFound = errors.New("not found")
|
||||
errCallNext = errors.New("continue to next")
|
||||
)
|
||||
|
||||
// Services implements the ServiceBackend interface
|
||||
func (r *SRegionDNS) Services(state request.Request, exact bool, opt plugin.Options) (services []msg.Service, err error) {
|
||||
switch state.QType() {
|
||||
|
||||
case dns.TypeTXT:
|
||||
t, _ := dnsutil.TrimZone(state.Name(), state.Zone)
|
||||
|
||||
@@ -203,7 +195,6 @@ func (r *SRegionDNS) Services(state request.Request, exact bool, opt plugin.Opti
|
||||
}
|
||||
svc := msg.Service{Text: "0.0.1", TTL: 28800, Key: msg.Path(state.QName(), "coredns")}
|
||||
return []msg.Service{svc}, nil
|
||||
|
||||
case dns.TypeNS:
|
||||
ns := r.nsAddr()
|
||||
svc := msg.Service{Host: ns.A.String(), Key: msg.Path(state.QName(), "coredns")}
|
||||
@@ -219,7 +210,6 @@ func (r *SRegionDNS) Services(state request.Request, exact bool, opt plugin.Opti
|
||||
}
|
||||
|
||||
services, err = r.Records(state, false)
|
||||
ylog.Debugf("Get records: %#v, error: %v", services, err)
|
||||
return
|
||||
}
|
||||
|
||||
@@ -230,7 +220,7 @@ func (r *SRegionDNS) Lookup(state request.Request, name string, typ uint16) (*dn
|
||||
|
||||
// IsNameError implements the ServiceBackend interface
|
||||
func (r *SRegionDNS) IsNameError(err error) bool {
|
||||
return err == errNoItems
|
||||
return err == errCallNext
|
||||
}
|
||||
|
||||
// Records looks up records in region mysql
|
||||
@@ -242,25 +232,22 @@ func (r *SRegionDNS) Records(state request.Request, exact bool) ([]msg.Service,
|
||||
return r.findRecords(req)
|
||||
}
|
||||
|
||||
func (r *SRegionDNS) getHostIpWithName(req *recordRequest) []string {
|
||||
func (r *SRegionDNS) getHostIpWithName(req *recordRequest) string {
|
||||
name := req.QueryName()
|
||||
host, _ := models.HostManager.FetchByName("", name)
|
||||
if host == nil {
|
||||
return nil
|
||||
return ""
|
||||
}
|
||||
ip := host.(*models.SHost).AccessIp
|
||||
if len(ip) == 0 {
|
||||
return nil
|
||||
}
|
||||
return []string{ip}
|
||||
return ip
|
||||
}
|
||||
|
||||
func (r *SRegionDNS) getGuestIpWithName(req *recordRequest) []string {
|
||||
ips := []string{}
|
||||
name := req.QueryName()
|
||||
projectId := req.ProjectId()
|
||||
isExitOnly := req.IsExitOnly()
|
||||
ips = models.GuestManager.GetIpInProjectWithName(projectId, name, isExitOnly)
|
||||
wantOnlyExit := false
|
||||
ips = models.GuestManager.GetIpInProjectWithName(projectId, name, wantOnlyExit)
|
||||
return ips
|
||||
}
|
||||
|
||||
@@ -268,6 +255,9 @@ func (r *SRegionDNS) getK8sServiceBackends(req *recordRequest) ([]string, error)
|
||||
queryInfo := req.GetK8sQueryInfo()
|
||||
pods, err := r.getK8sServicePods(queryInfo.Namespace, queryInfo.ServiceName)
|
||||
if err != nil {
|
||||
if k8serrors.IsNotFound(err) {
|
||||
err = nil
|
||||
}
|
||||
return nil, err
|
||||
}
|
||||
ips := make([]string, 0)
|
||||
@@ -301,10 +291,9 @@ func (r *SRegionDNS) Name() string {
|
||||
return PluginName
|
||||
}
|
||||
|
||||
func (r *SRegionDNS) queryLocalDnsRecords(req *recordRequest) (recs []msg.Service, err error) {
|
||||
func (r *SRegionDNS) queryLocalDnsRecords(req *recordRequest) (recs []msg.Service) {
|
||||
ips := models.DnsRecordManager.QueryDnsIps(req.ProjectId(), req.Name(), req.Type())
|
||||
if len(ips) == 0 {
|
||||
err = errNoItems
|
||||
return
|
||||
}
|
||||
|
||||
@@ -317,12 +306,12 @@ func (r *SRegionDNS) queryLocalDnsRecords(req *recordRequest) (recs []msg.Servic
|
||||
if req.IsSRV() {
|
||||
parts := strings.SplitN(ip.Addr, ":", 2)
|
||||
if len(parts) != 2 {
|
||||
err = fmt.Errorf("Invalid SRV records: %q", ip.Addr)
|
||||
ylog.Errorf("Invalid SRV records: %q", ip.Addr)
|
||||
return
|
||||
}
|
||||
port, e := strconv.Atoi(parts[1])
|
||||
if e != nil {
|
||||
err = e
|
||||
ylog.Errorf("Invalid SRV records: %q", ip.Addr)
|
||||
return
|
||||
}
|
||||
s = msg.Service{Host: parts[0], Port: port, TTL: ttl}
|
||||
@@ -334,17 +323,6 @@ func (r *SRegionDNS) queryLocalDnsRecords(req *recordRequest) (recs []msg.Servic
|
||||
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) isMyDomain(req *recordRequest) bool {
|
||||
zones := []string{fmt.Sprintf("%s.", r.PrimaryZone)}
|
||||
zone := plugin.Zones(zones).Matches(req.state.Name())
|
||||
@@ -354,54 +332,65 @@ func (r *SRegionDNS) isMyDomain(req *recordRequest) bool {
|
||||
return false
|
||||
}
|
||||
|
||||
func (r *SRegionDNS) findRecords(req *recordRequest) (recs []msg.Service, err error) {
|
||||
func (r *SRegionDNS) findRecords(req *recordRequest) ([]msg.Service, error) {
|
||||
// 1. try local dns records table
|
||||
recs, err = r.queryLocalDnsRecords(req)
|
||||
if len(recs) != 0 {
|
||||
return
|
||||
rrs := r.queryLocalDnsRecords(req)
|
||||
if len(rrs) > 0 {
|
||||
return rrs, nil
|
||||
}
|
||||
|
||||
isPlainName := req.IsPlainName()
|
||||
isMyDomain := r.isMyDomain(req)
|
||||
|
||||
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
|
||||
if isPlainName {
|
||||
isCloudIp := req.SrcInCloud()
|
||||
if isCloudIp {
|
||||
ips := r.findInternalRecordIps(req)
|
||||
if len(ips) > 0 {
|
||||
return ips2DnsRecords(ips), nil
|
||||
} else {
|
||||
return nil, errNotFound
|
||||
}
|
||||
} else {
|
||||
return nil, errRefused
|
||||
}
|
||||
} else if isMyDomain {
|
||||
ips := r.findInternalRecordIps(req)
|
||||
if len(ips) > 0 {
|
||||
return ips2DnsRecords(ips), nil
|
||||
} else {
|
||||
return nil, errNotFound
|
||||
}
|
||||
} else {
|
||||
return nil, errCallNext
|
||||
}
|
||||
|
||||
// 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
|
||||
func (r *SRegionDNS) findInternalRecordIps(req *recordRequest) []string {
|
||||
{
|
||||
// 1. try host table
|
||||
ip := r.getHostIpWithName(req)
|
||||
if len(ip) > 0 {
|
||||
return []string{ip}
|
||||
}
|
||||
}
|
||||
// 2. try guest table
|
||||
ip = r.getGuestIpWithName(req)
|
||||
if len(ip) != 0 {
|
||||
return ip, nil
|
||||
{
|
||||
// 2. try guest table
|
||||
ips := r.getGuestIpWithName(req)
|
||||
if len(ips) > 0 {
|
||||
return ips
|
||||
}
|
||||
}
|
||||
|
||||
if !r.IsK8sClientReady() {
|
||||
if r.K8sClient == nil {
|
||||
ylog.Warningf("K8s client not ready, skip it.")
|
||||
return nil, errNoItems
|
||||
return nil
|
||||
}
|
||||
// 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
|
||||
return ips
|
||||
}
|
||||
|
||||
func ips2DnsRecords(ips []string) []msg.Service {
|
||||
|
||||
+27
-22
@@ -7,17 +7,16 @@ import (
|
||||
"github.com/coredns/coredns/request"
|
||||
"github.com/miekg/dns"
|
||||
|
||||
"yunion.io/x/pkg/tristate"
|
||||
|
||||
"yunion.io/x/onecloud/pkg/compute/models"
|
||||
"yunion.io/x/pkg/tristate"
|
||||
)
|
||||
|
||||
type recordRequest struct {
|
||||
state request.Request
|
||||
domainSegs []string
|
||||
guest *models.SGuest
|
||||
host *models.SHost
|
||||
network *models.SNetwork
|
||||
state request.Request
|
||||
domainSegs []string
|
||||
srcProjectId string
|
||||
srcInCloud bool
|
||||
network *models.SNetwork
|
||||
}
|
||||
|
||||
func parseRequest(state request.Request) (r *recordRequest, err error) {
|
||||
@@ -28,9 +27,19 @@ func parseRequest(state request.Request) (r *recordRequest, err error) {
|
||||
domainSegs: segs,
|
||||
}
|
||||
srcIP := r.SrcIP4()
|
||||
r.guest = models.GuestnetworkManager.GetGuestByAddress(srcIP)
|
||||
r.host = models.HostnetworkManager.GetHostByAddress(srcIP)
|
||||
r.network, _ = models.NetworkManager.GetNetworkOfIP(srcIP, "", tristate.None)
|
||||
// NOTE the check on networks_tbl is a hack, we should be more specific
|
||||
// by querying only guest networks, host networks, and others like
|
||||
// loadbalancer network the to come.
|
||||
//
|
||||
// Order matters here, we want to find the srcIP project as accurately
|
||||
// as possible
|
||||
if guest := models.GuestnetworkManager.GetGuestByAddress(srcIP); guest != nil {
|
||||
r.srcProjectId = guest.ProjectId
|
||||
r.srcInCloud = true
|
||||
} else if network, _ := models.NetworkManager.GetNetworkOfIP(srcIP, "", tristate.None); network != nil {
|
||||
r.srcProjectId = network.ProjectId
|
||||
r.srcInCloud = true
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
@@ -41,6 +50,11 @@ func (r recordRequest) Name() string {
|
||||
return name
|
||||
}
|
||||
|
||||
func (r recordRequest) IsPlainName() bool {
|
||||
nl := dns.CountLabel(r.Name())
|
||||
return nl == 1
|
||||
}
|
||||
|
||||
func (r recordRequest) QueryName() string {
|
||||
seps := strings.Split(r.Name(), ".")
|
||||
if len(seps) == 0 {
|
||||
@@ -63,20 +77,11 @@ func (r recordRequest) SrcIP4() string {
|
||||
}
|
||||
|
||||
func (r recordRequest) ProjectId() string {
|
||||
if r.guest != nil {
|
||||
return r.guest.ProjectId
|
||||
}
|
||||
if r.network != nil {
|
||||
return r.network.ProjectId
|
||||
}
|
||||
return ""
|
||||
return r.srcProjectId
|
||||
}
|
||||
|
||||
func (r recordRequest) IsExitOnly() bool {
|
||||
if r.guest == nil {
|
||||
return false
|
||||
}
|
||||
return r.guest.IsExitOnly()
|
||||
func (r recordRequest) SrcInCloud() bool {
|
||||
return r.srcInCloud
|
||||
}
|
||||
|
||||
type K8sQueryInfo struct {
|
||||
|
||||
+2
-10
@@ -8,8 +8,6 @@ import (
|
||||
"github.com/coredns/coredns/plugin/pkg/dnsutil"
|
||||
"github.com/coredns/coredns/request"
|
||||
|
||||
"yunion.io/x/log"
|
||||
|
||||
"yunion.io/x/onecloud/pkg/compute/models"
|
||||
)
|
||||
|
||||
@@ -21,13 +19,7 @@ func (r *SRegionDNS) Reverse(state request.Request, exact bool, opt plugin.Optio
|
||||
return nil, e
|
||||
}
|
||||
records, err := r.getNameForIp(ip, state)
|
||||
if err != nil {
|
||||
log.Errorf("Reverse get name for ip: %v", err)
|
||||
}
|
||||
if len(records) == 0 {
|
||||
return records, errNoItems
|
||||
}
|
||||
return records, nil
|
||||
return records, err
|
||||
}
|
||||
|
||||
func (r *SRegionDNS) getNameForIp(ip string, state request.Request) ([]msg.Service, error) {
|
||||
@@ -53,7 +45,7 @@ func (r *SRegionDNS) getNameForIp(ip string, state request.Request) ([]msg.Servi
|
||||
if guest != nil {
|
||||
return []msg.Service{{Host: r.joinDomain(guest.Name), TTL: defaultTTL}}, nil
|
||||
}
|
||||
return nil, errNoItems
|
||||
return nil, errNotFound
|
||||
}
|
||||
|
||||
func (r *SRegionDNS) joinDomain(name string) string {
|
||||
|
||||
+3
-6
@@ -2,7 +2,6 @@ package dns
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"os"
|
||||
|
||||
"github.com/coredns/coredns/core/dnsserver"
|
||||
"github.com/coredns/coredns/plugin"
|
||||
@@ -20,18 +19,16 @@ func init() {
|
||||
}
|
||||
|
||||
func setup(c *caddy.Controller) error {
|
||||
os.Stderr = os.Stdout
|
||||
|
||||
rDNS, err := regionDNSParse(c)
|
||||
if err != nil {
|
||||
return plugin.Error(PluginName, err)
|
||||
}
|
||||
|
||||
if rDNS.PrimaryZone == "" {
|
||||
return fmt.Errorf("dns_domain must provided")
|
||||
if len(rDNS.PrimaryZone) == 0 {
|
||||
return fmt.Errorf("dns_domain missing")
|
||||
}
|
||||
if !regutils.MatchDomainName(rDNS.PrimaryZone) {
|
||||
return fmt.Errorf("dns_domain %q not match domain format", rDNS.PrimaryZone)
|
||||
return fmt.Errorf("dns_domain %q invalid", rDNS.PrimaryZone)
|
||||
}
|
||||
|
||||
err = rDNS.initDB(c)
|
||||
|
||||
Reference in New Issue
Block a user