Files
worker/checks/cdns/master_task.go
Gleb Tv 2c7a0236da feat: publish standalone worker
Separate worker packaging and service lifecycle from the control plane.
2026-07-13 17:55:14 +03:00

124 строки
3.6 KiB
Go

package cdns
import (
"fmt"
"time"
"github.com/miekg/dns"
)
// Results provides functionality.
type Results map[string]nameServer
func masterTask(zone string, nameservers map[string]nameServer) (uint, uint, bool, Results) {
var numRequests uint
success := true
addressChannel := make(chan DNSreply)
soaChannel := make(chan SOAreply)
numNS := uint(0)
numAddrNS := uint(0)
results := make(Results)
for name := range nameservers {
if !v6only {
go localQuery(addressChannel, name, dns.TypeA)
}
if !v4only {
go localQuery(addressChannel, name, dns.TypeAAAA)
}
numNS++
}
if v6only || v4only {
numRequests = numNS
} else {
numRequests = numNS * 2
}
for i := uint(0); i < numRequests; i++ {
addrResult := <-addressChannel
addrFamily := "IPv6"
if addrResult.qtype == dns.TypeA {
addrFamily = "IPv4"
}
if addrResult.r == nil {
// TODO We may have different globalErrMsg is it
// works with IPv4 but not IPv6 (it should not happen but it does)
nameservers[addrResult.qname] = nameServer{
name: addrResult.qname,
ips: nil,
globalErrMsg: fmt.Sprintf("Cannot get the %s address: %s", addrFamily, addrResult.err),
}
success = false
} else {
if addrResult.r.Rcode != dns.RcodeSuccess {
nameservers[addrResult.qname] = nameServer{
name: addrResult.qname,
ips: nil,
globalErrMsg: fmt.Sprintf("Cannot get the %s address: %s", addrFamily, dns.RcodeToString[addrResult.r.Rcode]),
}
success = false
} else {
for j := range addrResult.r.Answer {
ansa := addrResult.r.Answer[j]
var ns string
switch a := ansa.(type) {
case *dns.A:
ns = a.A.String()
existing := nameservers[addrResult.qname]
nameservers[addrResult.qname] = nameServer{name: addrResult.qname, ips: append(existing.ips, ns)}
numAddrNS++
go soaQuery(soaChannel, zone, addrResult.qname, ns)
case *dns.AAAA:
ns = a.AAAA.String()
existing2 := nameservers[addrResult.qname]
nameservers[addrResult.qname] = nameServer{name: addrResult.qname, ips: append(existing2.ips, ns)}
numAddrNS++
go soaQuery(soaChannel, zone, addrResult.qname, ns)
}
}
}
}
}
for i := uint(0); i < numAddrNS; i++ {
if debug {
fmt.Printf("DEBUG Getting result for ns #%d/%d\n", i+1, numAddrNS)
}
soaResult := <-soaChannel
_, present := results[soaResult.name]
if !present {
results[soaResult.name] = nameServer{
name: soaResult.name,
ips: make([]string, 0),
success: make([]bool, 0),
errMsg: make([]string, 0),
serial: make([]uint32, 0),
rtts: make([]time.Duration, 0),
}
}
if !soaResult.retrieved {
results[soaResult.name] = nameServer{
name: soaResult.name,
ips: append(results[soaResult.name].ips, soaResult.address),
success: append(results[soaResult.name].success, false),
errMsg: append(results[soaResult.name].errMsg, soaResult.msg),
serial: append(results[soaResult.name].serial, 0),
rtts: append(results[soaResult.name].rtts, soaResult.rtt),
}
success = false
} else {
results[soaResult.name] = nameServer{
name: soaResult.name,
ips: append(results[soaResult.name].ips, soaResult.address),
success: append(results[soaResult.name].success, true),
errMsg: append(results[soaResult.name].errMsg, ""),
serial: append(results[soaResult.name].serial, soaResult.serial),
rtts: append(results[soaResult.name].rtts, soaResult.rtt),
}
}
}
for name := range nameservers {
if nameservers[name].ips == nil {
results[name] = nameservers[name]
}
}
return numNS, numAddrNS, success, results
}