From 19fae5bc7d9c58a8a981c0929f92bc8768e464b8 Mon Sep 17 00:00:00 2001 From: cimaol Date: Tue, 10 Mar 2020 08:03:34 +0000 Subject: [PATCH] Removed kubelet port flag. Node name now always need from env/flag --- probe/kubernetes/kubelet.go | 37 --------- probe/kubernetes/kubelet_pods.json | 127 ----------------------------- probe/kubernetes/kubelet_test.go | 48 ----------- probe/kubernetes/reporter.go | 27 ++---- probe/kubernetes/reporter_test.go | 22 +---- prog/main.go | 2 - prog/probe.go | 3 +- 7 files changed, 9 insertions(+), 257 deletions(-) delete mode 100644 probe/kubernetes/kubelet.go delete mode 100644 probe/kubernetes/kubelet_pods.json delete mode 100644 probe/kubernetes/kubelet_test.go diff --git a/probe/kubernetes/kubelet.go b/probe/kubernetes/kubelet.go deleted file mode 100644 index 74720056f..000000000 --- a/probe/kubernetes/kubelet.go +++ /dev/null @@ -1,37 +0,0 @@ -package kubernetes - -import ( - "fmt" - "net/http" - - "github.com/ugorji/go/codec" -) - -// Intentionally not using the full kubernetes library DS -// to make parsing faster and more tolerant to schema changes -type podList struct { - Items []struct { - Metadata struct { - UID string `json:"uid"` - } `json:"metadata"` - } `json:"items"` -} - -// GetLocalPodUIDs obtains the UID of the pods run locally (it's just exported for testing) -var GetLocalPodUIDs = func(kubeletHost string) (map[string]struct{}, error) { - url := fmt.Sprintf("http://%s/pods/", kubeletHost) - resp, err := http.Get(url) - if err != nil { - return nil, err - } - defer resp.Body.Close() - var localPods podList - if err := codec.NewDecoder(resp.Body, &codec.JsonHandle{}).Decode(&localPods); err != nil { - return nil, err - } - result := make(map[string]struct{}, len(localPods.Items)) - for _, pod := range localPods.Items { - result[pod.Metadata.UID] = struct{}{} - } - return result, nil -} diff --git a/probe/kubernetes/kubelet_pods.json b/probe/kubernetes/kubelet_pods.json deleted file mode 100644 index beda44287..000000000 --- a/probe/kubernetes/kubelet_pods.json +++ /dev/null @@ -1,127 +0,0 @@ -{ - "kind": "PodList", - "apiVersion": "v1", - "metadata": {}, - "items": [ - { - "metadata": { - "name": "configs-db-289919966-phz6s", - "generateName": "configs-db-289919966-", - "namespace": "default", - "selfLink": "/api/v1/namespaces/default/pods/configs-db-289919966-phz6s", - "uid": "af1b5325-d8cf-11e6-84fa-0800278a0c83", - "resourceVersion": "16439", - "creationTimestamp": "2017-01-12T14:01:58Z", - "labels": { - "name": "configs-db", - "pod-template-hash": "289919966" - }, - "annotations": { - "kubernetes.io/config.seen": "2017-01-16T13:01:12.824003566Z", - "kubernetes.io/config.source": "api", - "kubernetes.io/created-by": "{\"kind\":\"SerializedReference\",\"apiVersion\":\"v1\",\"reference\":{\"kind\":\"ReplicaSet\",\"namespace\":\"default\",\"name\":\"configs-db-289919966\",\"uid\":\"af1ad88a-d8cf-11e6-84fa-0800278a0c83\",\"apiVersion\":\"extensions\",\"resourceVersion\":\"16139\"}}\n", - "prometheus.io.scrape": "false" - }, - "ownerReferences": [ - { - "apiVersion": "extensions/v1beta1", - "kind": "ReplicaSet", - "name": "configs-db-289919966", - "uid": "af1ad88a-d8cf-11e6-84fa-0800278a0c83", - "controller": true - } - ] - }, - "spec": { - "volumes": [ - { - "name": "default-token-bzbt2", - "secret": { - "secretName": "default-token-bzbt2", - "defaultMode": 420 - } - } - ], - "containers": [ - { - "name": "configs-db", - "image": "quay.io/weaveworks/configs-db", - "ports": [ - { - "containerPort": 5432, - "protocol": "TCP" - } - ], - "resources": {}, - "volumeMounts": [ - { - "name": "default-token-bzbt2", - "readOnly": true, - "mountPath": "/var/run/secrets/kubernetes.io/serviceaccount" - } - ], - "terminationMessagePath": "/dev/termination-log", - "imagePullPolicy": "IfNotPresent" - } - ], - "restartPolicy": "Always", - "terminationGracePeriodSeconds": 30, - "dnsPolicy": "ClusterFirst", - "serviceAccountName": "default", - "serviceAccount": "default", - "nodeName": "minikube", - "securityContext": {} - }, - "status": { - "phase": "Running", - "conditions": [ - { - "type": "Initialized", - "status": "True", - "lastProbeTime": null, - "lastTransitionTime": "2017-01-12T14:01:59Z" - }, - { - "type": "Ready", - "status": "True", - "lastProbeTime": null, - "lastTransitionTime": "2017-01-16T13:01:32Z" - }, - { - "type": "PodScheduled", - "status": "True", - "lastProbeTime": null, - "lastTransitionTime": "2017-01-12T14:01:58Z" - } - ], - "hostIP": "192.168.99.100", - "podIP": "172.17.0.23", - "startTime": "2017-01-12T14:01:59Z", - "containerStatuses": [ - { - "name": "configs-db", - "state": { - "running": { - "startedAt": "2017-01-16T13:01:25Z" - } - }, - "lastState": { - "terminated": { - "exitCode": 0, - "reason": "Completed", - "startedAt": "2017-01-12T14:02:00Z", - "finishedAt": "2017-01-16T11:54:52Z", - "containerID": "docker://1e8bb78d93fbf1395d301e81bc30e3416c934e6180d0684f8576e8383c41d070" - } - }, - "ready": true, - "restartCount": 1, - "image": "quay.io/weaveworks/configs-db", - "imageID": "docker://sha256:f6d024621defbdb2e9124da2f4fef1ab02d9d329405722c5c519d96e5a17eeb7", - "containerID": "docker://9a5091295c7e6616237162ffefb89012e04f9d26d8e1b732b6a8a77b239e6052" - } - ] - } - } - ] -} diff --git a/probe/kubernetes/kubelet_test.go b/probe/kubernetes/kubelet_test.go deleted file mode 100644 index 31a6471c8..000000000 --- a/probe/kubernetes/kubelet_test.go +++ /dev/null @@ -1,48 +0,0 @@ -package kubernetes_test - -import ( - "io/ioutil" - "net/http" - "net/http/httptest" - "net/url" - "testing" - - "github.com/weaveworks/scope/probe/kubernetes" -) - -const kubeletPodsJSONFile = "kubelet_pods.json" - -// obtained with jq .items[].metadata.uid < kubelet_pods.json -var expectedPodUIDs = []string{ - "af1b5325-d8cf-11e6-84fa-0800278a0c83", -} - -func TestGetLocalPodUIDs(t *testing.T) { - server := httptest.NewServer(http.HandlerFunc( - func(w http.ResponseWriter, r *http.Request) { - if r.URL.Path != "/pods/" { - t.Fatalf("unexpected path: %s", r.URL.Path) - } - b, err := ioutil.ReadFile(kubeletPodsJSONFile) - if err != nil { - t.Fatalf("unexpected error reading json file: %v", err) - } - w.Write(b) - }, - )) - defer server.Close() - - serverURL, _ := url.Parse(server.URL) - uids, err := kubernetes.GetLocalPodUIDs(serverURL.Host) - if err != nil { - t.Fatalf("unexpected error: %v", err) - } - if len(expectedPodUIDs) != len(uids) { - t.Errorf("nnexpected length in pod UIDs (%d): expected %d", len(uids), len(expectedPodUIDs)) - } - for _, expectedUID := range expectedPodUIDs { - if _, ok := uids[expectedUID]; !ok { - t.Errorf("uid not found: %s", expectedUID) - } - } -} diff --git a/probe/kubernetes/reporter.go b/probe/kubernetes/reporter.go index fadf602f0..98468720c 100644 --- a/probe/kubernetes/reporter.go +++ b/probe/kubernetes/reporter.go @@ -2,10 +2,8 @@ package kubernetes import ( "fmt" - "k8s.io/apimachinery/pkg/labels" - log "github.com/sirupsen/logrus" "github.com/weaveworks/common/mtime" "github.com/weaveworks/scope/probe" "github.com/weaveworks/scope/probe/controls" @@ -195,11 +193,10 @@ type Reporter struct { hostID string handlerRegistry *controls.HandlerRegistry nodeName string - kubeletPort uint } // NewReporter makes a new Reporter -func NewReporter(client Client, pipes controls.PipeClient, probeID string, hostID string, probe *probe.Probe, handlerRegistry *controls.HandlerRegistry, nodeName string, kubeletPort uint) *Reporter { +func NewReporter(client Client, pipes controls.PipeClient, probeID string, hostID string, probe *probe.Probe, handlerRegistry *controls.HandlerRegistry, nodeName string) *Reporter { reporter := &Reporter{ client: client, pipes: pipes, @@ -208,7 +205,6 @@ func NewReporter(client Client, pipes controls.PipeClient, probeID string, hostI hostID: hostID, handlerRegistry: handlerRegistry, nodeName: nodeName, - kubeletPort: kubeletPort, } reporter.registerControls() client.WatchPods(reporter.podEvent) @@ -658,26 +654,13 @@ func (r *Reporter) podTopology(services []Service, deployments []Deployment, dae )) } } - - var localPodUIDs map[string]struct{} - if r.nodeName == "" && r.kubeletPort != 0 { - // We don't know the node name: fall back to obtaining the local pods from kubelet - var err error - localPodUIDs, err = GetLocalPodUIDs(fmt.Sprintf("127.0.0.1:%d", r.kubeletPort)) - if err != nil { - log.Warnf("No node name and cannot obtain local pods, reporting all (which may impact performance): %v", err) - } + // filter out non-local pods: we only want to report local ones for performance reasons. + if r.nodeName == "" { + return pods, fmt.Errorf("pod topology failure: no node name given for reporter") } err := r.client.WalkPods(func(p Pod) error { - // filter out non-local pods: we only want to report local ones for performance reasons. - if r.nodeName != "" { - if p.NodeName() != r.nodeName { + if p.NodeName() != r.nodeName { return nil - } - } else if localPodUIDs != nil { - if _, ok := localPodUIDs[p.UID()]; !ok { - return nil - } } for _, selector := range selectors { selector(p) diff --git a/probe/kubernetes/reporter_test.go b/probe/kubernetes/reporter_test.go index 34a284eb6..3393e9b2d 100644 --- a/probe/kubernetes/reporter_test.go +++ b/probe/kubernetes/reporter_test.go @@ -219,21 +219,11 @@ func (c mockPipeClient) PipeClose(appID, id string) error { } func TestReporter(t *testing.T) { - oldGetNodeName := kubernetes.GetLocalPodUIDs - defer func() { kubernetes.GetLocalPodUIDs = oldGetNodeName }() - kubernetes.GetLocalPodUIDs = func(string) (map[string]struct{}, error) { - uids := map[string]struct{}{ - pod1UID: {}, - pod2UID: {}, - } - return uids, nil - } - pod1ID := report.MakePodNodeID(pod1UID) pod2ID := report.MakePodNodeID(pod2UID) serviceID := report.MakeServiceNodeID(serviceUID) hr := controls.NewDefaultHandlerRegistry() - rpt, _ := kubernetes.NewReporter(newMockClient(), nil, "probe-id", "foo", nil, hr, "", 0).Report() + rpt, _ := kubernetes.NewReporter(newMockClient(), nil, "probe-id", "foo", nil, hr, nodeName).Report() // Reporter should have added the following pods for _, pod := range []struct { @@ -337,7 +327,7 @@ func BenchmarkReporter(b *testing.B) { } mockK8s.deployments = append(mockK8s.deployments, kubernetes.NewDeployment(&deployment)) } - reporter := kubernetes.NewReporter(mockK8s, nil, "probe-id", "foo", nil, hr, nodeName, 0) + reporter := kubernetes.NewReporter(mockK8s, nil, "probe-id", "foo", nil, hr, nodeName) b.ResetTimer() for i := 0; i < b.N; i++ { @@ -400,16 +390,10 @@ type callbackReadCloser struct { func (c *callbackReadCloser) Close() error { return c.close() } func TestReporterGetLogs(t *testing.T) { - oldGetNodeName := kubernetes.GetLocalPodUIDs - defer func() { kubernetes.GetLocalPodUIDs = oldGetNodeName }() - kubernetes.GetLocalPodUIDs = func(string) (map[string]struct{}, error) { - return map[string]struct{}{}, nil - } - client := newMockClient() pipes := mockPipeClient{} hr := controls.NewDefaultHandlerRegistry() - reporter := kubernetes.NewReporter(client, pipes, "", "", nil, hr, "", 0) + reporter := kubernetes.NewReporter(client, pipes, "", "", nil, hr, nodeName) // Should error on invalid IDs { diff --git a/prog/main.go b/prog/main.go index e44755228..3e73196b9 100644 --- a/prog/main.go +++ b/prog/main.go @@ -131,7 +131,6 @@ type probeFlags struct { kubernetesRole string kubernetesNodeName string kubernetesClientConfig kubernetes.ClientConfig - kubernetesKubeletPort uint ecsEnabled bool ecsCacheSize int @@ -343,7 +342,6 @@ func setupFlags(flags *flags) { flag.StringVar(&flags.probe.kubernetesClientConfig.User, "probe.kubernetes.user", "", "The name of the kubeconfig user to use") flag.StringVar(&flags.probe.kubernetesClientConfig.Username, "probe.kubernetes.username", "", "Username for basic authentication to the API server") flag.StringVar(&flags.probe.kubernetesNodeName, "probe.kubernetes.node-name", "", "Name of this node, for filtering pods") - flag.UintVar(&flags.probe.kubernetesKubeletPort, "probe.kubernetes.kubelet-port", 10255, "Node-local TCP port for contacting kubelet (zero to disable)") // AWS ECS flag.BoolVar(&flags.probe.ecsEnabled, "probe.ecs", false, "Collect ecs-related attributes for containers on this node") diff --git a/prog/probe.go b/prog/probe.go index e69ff4984..549a5716c 100644 --- a/prog/probe.go +++ b/prog/probe.go @@ -135,7 +135,6 @@ func probeMain(flags probeFlags, targets []appclient.Target) { case kubernetesRoleHost: flags.kubernetesEnabled = true case kubernetesRoleCluster: - flags.kubernetesKubeletPort = 0 flags.kubernetesEnabled = true flags.spyProcs = false flags.procEnabled = false @@ -319,7 +318,7 @@ func probeMain(flags probeFlags, targets []appclient.Target) { if flags.kubernetesEnabled && flags.kubernetesRole != kubernetesRoleHost { if client, err := kubernetes.NewClient(flags.kubernetesClientConfig); err == nil { defer client.Stop() - reporter := kubernetes.NewReporter(client, clients, probeID, hostID, p, handlerRegistry, flags.kubernetesNodeName, flags.kubernetesKubeletPort) + reporter := kubernetes.NewReporter(client, clients, probeID, hostID, p, handlerRegistry, flags.kubernetesNodeName) defer reporter.Stop() p.AddReporter(reporter) } else {