Store results by pod name instead of pod IP

Signed-off-by: Sachin Kamboj <skamboj1@bloomberg.net>
This commit is contained in:
Sachin Kamboj
2020-04-03 23:22:52 -04:00
parent 86febf8295
commit f0c66f29c7
+59 -37
View File
@@ -29,20 +29,21 @@ import (
// CheckNeighbours queries the kubernetes API server for all other goldpinger pods
// then calls Ping() on each one
func CheckNeighbours(ps *PodSelecter) *models.CheckResults {
return PingAllPods(ps.SelectPods())
func CheckNeighbours() *models.CheckResults {
return PingAllPods(SelectPods())
}
// CheckNeighboursNeighbours queries the kubernetes API server for all other goldpinger
// pods then calls Check() on each one
func CheckNeighboursNeighbours(ps *PodSelecter) *models.CheckAllResults {
return CheckAllPods(ps.SelectPods())
func CheckNeighboursNeighbours() *models.CheckAllResults {
return CheckAllPods(SelectPods())
}
type PingAllPodsResult struct {
podName string
podResult models.PodResult
hostIPv4 strfmt.IPv4
podIP string
podIPv4 strfmt.IPv4
}
func pickPodHostIP(podIP, hostIP string) string {
@@ -70,7 +71,7 @@ func checkDNS() *models.DNSResults {
return &results
}
func PingAllPods(pods map[string]string) *models.CheckResults {
func PingAllPods(pods map[string]*GoldpingerPod) *models.CheckResults {
result := models.CheckResults{}
@@ -78,44 +79,65 @@ func PingAllPods(pods map[string]string) *models.CheckResults {
wg := sync.WaitGroup{}
wg.Add(len(pods))
for podIP, hostIP := range pods {
for _, pod := range pods {
go func(podIP string, hostIP string) {
go func(pod *GoldpingerPod) {
// metrics
CountCall("made", "ping")
timer := GetLabeledPeersCallsTimer("ping", hostIP, podIP)
timer := GetLabeledPeersCallsTimer("ping", pod.HostIP, pod.PodIP)
start := time.Now()
// setup
var channelResult PingAllPodsResult
channelResult.hostIPv4.UnmarshalText([]byte(hostIP))
channelResult.podIP = podIP
channelResult.podName = pod.Name
channelResult.hostIPv4.UnmarshalText([]byte(pod.HostIP))
channelResult.podIPv4.UnmarshalText([]byte(pod.PodIP))
OK := false
var responseTime int64
client, err := getClient(pickPodHostIP(podIP, hostIP))
client, err := getClient(pickPodHostIP(pod.PodIP, pod.HostIP))
if err != nil {
channelResult.podResult = models.PodResult{HostIP: channelResult.hostIPv4, OK: &OK, Error: err.Error(), StatusCode: 500, ResponseTimeMs: responseTime}
channelResult.podIP = hostIP
channelResult.podResult = models.PodResult{
PodIP: channelResult.podIPv4,
HostIP: channelResult.hostIPv4,
OK: &OK,
Error: err.Error(),
StatusCode: 500,
ResponseTimeMs: responseTime,
}
CountError("ping")
} else {
resp, err := client.Operations.Ping(nil)
responseTime = time.Since(start).Nanoseconds() / int64(time.Millisecond)
OK = (err == nil)
if OK {
channelResult.podResult = models.PodResult{HostIP: channelResult.hostIPv4, OK: &OK, Response: resp.Payload, StatusCode: 200, ResponseTimeMs: responseTime}
channelResult.podResult = models.PodResult{
PodIP: channelResult.podIPv4,
HostIP: channelResult.hostIPv4,
OK: &OK,
Response: resp.Payload,
StatusCode: 200,
ResponseTimeMs: responseTime,
}
timer.ObserveDuration()
} else {
channelResult.podResult = models.PodResult{HostIP: channelResult.hostIPv4, OK: &OK, Error: err.Error(), StatusCode: 504, ResponseTimeMs: responseTime}
channelResult.podResult = models.PodResult{
PodIP: channelResult.podIPv4,
HostIP: channelResult.hostIPv4,
OK: &OK,
Error: err.Error(),
StatusCode: 504,
ResponseTimeMs: responseTime,
}
CountError("ping")
}
}
ch <- channelResult
wg.Done()
}(podIP, hostIP)
}(pod)
}
if len(GoldpingerConfig.DnsHosts) > 0 {
result.DNSResults = *checkDNS()
@@ -127,26 +149,25 @@ func PingAllPods(pods map[string]string) *models.CheckResults {
result.PodResults = make(map[string]models.PodResult)
for response := range ch {
var podIPv4 strfmt.IPv4
podIPv4.UnmarshalText([]byte(response.podIP))
if *response.podResult.OK {
counterHealthy++
} else {
counterUnhealthy++
}
result.PodResults[response.podIP] = response.podResult
result.PodResults[response.podName] = response.podResult
}
CountHealthyUnhealthyNodes(counterHealthy, counterUnhealthy)
return &result
}
type CheckServicePodsResult struct {
podName string
checkAllPodResult models.CheckAllPodResult
hostIPv4 strfmt.IPv4
podIP string
podIPv4 strfmt.IPv4
}
func CheckAllPods(pods map[string]string) *models.CheckAllResults {
func CheckAllPods(pods map[string]*GoldpingerPod) *models.CheckAllResults {
result := models.CheckAllResults{Responses: make(map[string]models.CheckAllPodResult)}
@@ -154,28 +175,29 @@ func CheckAllPods(pods map[string]string) *models.CheckAllResults {
wg := sync.WaitGroup{}
wg.Add(len(pods))
for podIP, hostIP := range pods {
for _, pod := range pods {
go func(podIP string, hostIP string) {
go func(pod *GoldpingerPod) {
// stats
CountCall("made", "check")
timer := GetLabeledPeersCallsTimer("check", hostIP, podIP)
timer := GetLabeledPeersCallsTimer("check", pod.HostIP, pod.PodIP)
// setup
var channelResult CheckServicePodsResult
channelResult.hostIPv4.UnmarshalText([]byte(hostIP))
channelResult.podIP = podIP
client, err := getClient(pickPodHostIP(podIP, hostIP))
channelResult.podName = pod.Name
channelResult.hostIPv4.UnmarshalText([]byte(pod.HostIP))
channelResult.podIPv4.UnmarshalText([]byte(pod.PodIP))
client, err := getClient(pickPodHostIP(pod.PodIP, pod.HostIP))
OK := false
if err != nil {
channelResult.checkAllPodResult = models.CheckAllPodResult{
OK: &OK,
PodIP: channelResult.podIPv4,
HostIP: channelResult.hostIPv4,
Error: err.Error(),
}
channelResult.podIP = hostIP
CountError("checkAll")
} else {
resp, err := client.Operations.CheckServicePods(nil)
@@ -183,6 +205,7 @@ func CheckAllPods(pods map[string]string) *models.CheckAllResults {
if OK {
channelResult.checkAllPodResult = models.CheckAllPodResult{
OK: &OK,
PodIP: channelResult.podIPv4,
HostIP: channelResult.hostIPv4,
Response: resp.Payload,
}
@@ -190,6 +213,7 @@ func CheckAllPods(pods map[string]string) *models.CheckAllResults {
} else {
channelResult.checkAllPodResult = models.CheckAllPodResult{
OK: &OK,
PodIP: channelResult.podIPv4,
HostIP: channelResult.hostIPv4,
Error: err.Error(),
}
@@ -199,19 +223,17 @@ func CheckAllPods(pods map[string]string) *models.CheckAllResults {
ch <- channelResult
wg.Done()
}(podIP, hostIP)
}(pod)
}
wg.Wait()
close(ch)
for response := range ch {
var podIPv4 strfmt.IPv4
podIPv4.UnmarshalText([]byte(response.podIP))
result.Responses[response.podIP] = response.checkAllPodResult
result.Responses[response.podName] = response.checkAllPodResult
result.Hosts = append(result.Hosts, &models.CheckAllResultsHostsItems0{
HostIP: response.hostIPv4,
PodIP: podIPv4,
PodName: response.podName,
HostIP: response.hostIPv4,
PodIP: response.podIPv4,
})
if response.checkAllPodResult.Response != nil &&
response.checkAllPodResult.Response.DNSResults != nil {
@@ -222,7 +244,7 @@ func CheckAllPods(pods map[string]string) *models.CheckAllResults {
if result.DNSResults[host] == nil {
result.DNSResults[host] = make(map[string]models.DNSResult)
}
result.DNSResults[host][response.podIP] = response.checkAllPodResult.Response.DNSResults[host]
result.DNSResults[host][response.podName] = response.checkAllPodResult.Response.DNSResults[host]
}
}
}