mirror of
https://github.com/weaveworks/scope.git
synced 2026-07-21 06:20:31 +00:00
Removed kubelet port flag. Node name now always need from env/flag
This commit is contained in:
@@ -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
|
||||
}
|
||||
@@ -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"
|
||||
}
|
||||
]
|
||||
}
|
||||
}
|
||||
]
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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)
|
||||
|
||||
@@ -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
|
||||
{
|
||||
|
||||
@@ -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")
|
||||
|
||||
@@ -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 {
|
||||
|
||||
Reference in New Issue
Block a user