diff --git a/pkg/goldpinger/client.go b/pkg/goldpinger/client.go index a830e8b..680126e 100644 --- a/pkg/goldpinger/client.go +++ b/pkg/goldpinger/client.go @@ -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] } } }