From 3de06b5a0963549308bd28a296063ec5c5e963d7 Mon Sep 17 00:00:00 2001 From: Roland Schilter Date: Fri, 16 Mar 2018 12:39:57 -0700 Subject: [PATCH] Support controls in more k8s topologies (#3110) * Refactor: func has already access to probeID * Allow controls for more k8s topologies Controls for nodes generally need to know about the probe that is in control of them. This PR appends the probe ID info to k8s topologies CronJob, DaemonSet, Service, and StatefulSet. Therefore allowing plugins to append controls. * Remove superfluous error check * Add some tests to verify controls allowance --- probe/host/reporter_test.go | 3 ++- probe/kubernetes/cronjob.go | 13 +++++++------ probe/kubernetes/daemonset.go | 13 +++++++------ probe/kubernetes/reporter.go | 17 +++++++---------- probe/kubernetes/reporter_test.go | 27 ++++++++++++++++++++++++++- probe/kubernetes/service.go | 5 +++-- probe/kubernetes/statefulset.go | 11 ++++++----- 7 files changed, 58 insertions(+), 31 deletions(-) diff --git a/probe/host/reporter_test.go b/probe/host/reporter_test.go index c3084c378..2fd6d4b7e 100644 --- a/probe/host/reporter_test.go +++ b/probe/host/reporter_test.go @@ -57,7 +57,7 @@ func TestReporter(t *testing.T) { host.GetLocalNetworks = func() ([]*net.IPNet, error) { return []*net.IPNet{ipnet}, nil } hr := controls.NewDefaultHandlerRegistry() - rpt, err := host.NewReporter(hostID, hostname, "", "", nil, hr).Report() + rpt, err := host.NewReporter(hostID, hostname, "probe-id", "", nil, hr).Report() if err != nil { t.Fatal(err) } @@ -77,6 +77,7 @@ func TestReporter(t *testing.T) { {host.OS, runtime.GOOS}, {host.Uptime, uptime}, {host.KernelVersion, kernel}, + {report.ControlProbeID, "probe-id"}, } { if have, ok := node.Latest.Lookup(tuple.key); !ok || have != tuple.want { t.Errorf("Expected %s %q, got %q", tuple.key, tuple.want, have) diff --git a/probe/kubernetes/cronjob.go b/probe/kubernetes/cronjob.go index 5b8177b5e..abdeb421c 100644 --- a/probe/kubernetes/cronjob.go +++ b/probe/kubernetes/cronjob.go @@ -26,7 +26,7 @@ const ( type CronJob interface { Meta Selectors() ([]labels.Selector, error) - GetNode() report.Node + GetNode(probeID string) report.Node } type cronJob struct { @@ -74,12 +74,13 @@ func (cj *cronJob) Selectors() ([]labels.Selector, error) { return selectors, nil } -func (cj *cronJob) GetNode() report.Node { +func (cj *cronJob) GetNode(probeID string) report.Node { latest := map[string]string{ - NodeType: "CronJob", - Schedule: cj.Spec.Schedule, - Suspended: fmt.Sprint(cj.Spec.Suspend != nil && *cj.Spec.Suspend), // nil -> false - ActiveJobs: fmt.Sprint(len(cj.jobs)), + NodeType: "CronJob", + Schedule: cj.Spec.Schedule, + Suspended: fmt.Sprint(cj.Spec.Suspend != nil && *cj.Spec.Suspend), // nil -> false + ActiveJobs: fmt.Sprint(len(cj.jobs)), + report.ControlProbeID: probeID, } if cj.Status.LastScheduleTime != nil { latest[LastScheduled] = cj.Status.LastScheduleTime.Format(time.RFC3339Nano) diff --git a/probe/kubernetes/daemonset.go b/probe/kubernetes/daemonset.go index 4bba18d2a..285f155c3 100644 --- a/probe/kubernetes/daemonset.go +++ b/probe/kubernetes/daemonset.go @@ -19,7 +19,7 @@ const ( type DaemonSet interface { Meta Selector() (labels.Selector, error) - GetNode() report.Node + GetNode(probeID string) report.Node } type daemonSet struct { @@ -43,11 +43,12 @@ func (d *daemonSet) Selector() (labels.Selector, error) { return selector, nil } -func (d *daemonSet) GetNode() report.Node { +func (d *daemonSet) GetNode(probeID string) report.Node { return d.MetaNode(report.MakeDaemonSetNodeID(d.UID())).WithLatests(map[string]string{ - DesiredReplicas: fmt.Sprint(d.Status.DesiredNumberScheduled), - Replicas: fmt.Sprint(d.Status.CurrentNumberScheduled), - MisscheduledReplicas: fmt.Sprint(d.Status.NumberMisscheduled), - NodeType: "DaemonSet", + DesiredReplicas: fmt.Sprint(d.Status.DesiredNumberScheduled), + Replicas: fmt.Sprint(d.Status.CurrentNumberScheduled), + MisscheduledReplicas: fmt.Sprint(d.Status.NumberMisscheduled), + NodeType: "DaemonSet", + report.ControlProbeID: probeID, }) } diff --git a/probe/kubernetes/reporter.go b/probe/kubernetes/reporter.go index c4264e0e4..b46d6a0aa 100644 --- a/probe/kubernetes/reporter.go +++ b/probe/kubernetes/reporter.go @@ -235,9 +235,6 @@ func (r *Reporter) Report() (report.Report, error) { return result, err } hostTopology := r.hostTopology(services) - if err != nil { - return result, err - } daemonSetTopology, daemonSets, err := r.daemonSetTopology() if err != nil { return result, err @@ -250,7 +247,7 @@ func (r *Reporter) Report() (report.Report, error) { if err != nil { return result, err } - deploymentTopology, deployments, err := r.deploymentTopology(r.probeID) + deploymentTopology, deployments, err := r.deploymentTopology() if err != nil { return result, err } @@ -282,7 +279,7 @@ func (r *Reporter) serviceTopology() (report.Topology, []Service, error) { services = []Service{} ) err := r.client.WalkServices(func(s Service) error { - result.AddNode(s.GetNode()) + result.AddNode(s.GetNode(r.probeID)) services = append(services, s) return nil }) @@ -319,7 +316,7 @@ func (r *Reporter) hostTopology(services []Service) report.Topology { return t } -func (r *Reporter) deploymentTopology(probeID string) (report.Topology, []Deployment, error) { +func (r *Reporter) deploymentTopology() (report.Topology, []Deployment, error) { var ( result = report.MakeTopology(). WithMetadataTemplates(DeploymentMetadataTemplates). @@ -330,7 +327,7 @@ func (r *Reporter) deploymentTopology(probeID string) (report.Topology, []Deploy result.Controls.AddControls(ScalingControls) err := r.client.WalkDeployments(func(d Deployment) error { - result.AddNode(d.GetNode(probeID)) + result.AddNode(d.GetNode(r.probeID)) deployments = append(deployments, d) return nil }) @@ -344,7 +341,7 @@ func (r *Reporter) daemonSetTopology() (report.Topology, []DaemonSet, error) { WithMetricTemplates(DaemonSetMetricTemplates). WithTableTemplates(TableTemplates) err := r.client.WalkDaemonSets(func(d DaemonSet) error { - result.AddNode(d.GetNode()) + result.AddNode(d.GetNode(r.probeID)) daemonSets = append(daemonSets, d) return nil }) @@ -358,7 +355,7 @@ func (r *Reporter) statefulSetTopology() (report.Topology, []StatefulSet, error) WithMetricTemplates(StatefulSetMetricTemplates). WithTableTemplates(TableTemplates) err := r.client.WalkStatefulSets(func(s StatefulSet) error { - result.AddNode(s.GetNode()) + result.AddNode(s.GetNode(r.probeID)) statefulSets = append(statefulSets, s) return nil }) @@ -372,7 +369,7 @@ func (r *Reporter) cronJobTopology() (report.Topology, []CronJob, error) { WithMetricTemplates(CronJobMetricTemplates). WithTableTemplates(TableTemplates) err := r.client.WalkCronJobs(func(c CronJob) error { - result.AddNode(c.GetNode()) + result.AddNode(c.GetNode(r.probeID)) cronJobs = append(cronJobs, c) return nil }) diff --git a/probe/kubernetes/reporter_test.go b/probe/kubernetes/reporter_test.go index e2ca252bf..6be9dd5ad 100644 --- a/probe/kubernetes/reporter_test.go +++ b/probe/kubernetes/reporter_test.go @@ -194,7 +194,7 @@ func TestReporter(t *testing.T) { pod2ID := report.MakePodNodeID(pod2UID) serviceID := report.MakeServiceNodeID(serviceUID) hr := controls.NewDefaultHandlerRegistry() - rpt, _ := kubernetes.NewReporter(newMockClient(), nil, "", "foo", nil, hr, "", 0).Report() + rpt, _ := kubernetes.NewReporter(newMockClient(), nil, "probe-id", "foo", nil, hr, "", 0).Report() // Reporter should have added the following pods for _, pod := range []struct { @@ -246,6 +246,31 @@ func TestReporter(t *testing.T) { } } } + + // Reporter should allow controls for k8s topologies by providing a probe ID + { + for _, topologyName := range []string{ + report.Container, + report.CronJob, + report.DaemonSet, + report.Deployment, + report.Pod, + report.Service, + report.StatefulSet, + } { + topology, ok := rpt.Topology(topologyName) + if !ok { + // TODO: this mock report doesn't have nodes for all the topologies yet, so don't fail for now. + // t.Errorf("Expected report to have nodes in topology %q, but none found", topology) + } + for _, n := range topology.Nodes { + if probeID, ok := n.Latest.Lookup(report.ControlProbeID); !ok || probeID != "probe-id" { + t.Errorf("Expected node %q to have probeID, but not found", n.ID) + } + } + } + } + } func TestTagger(t *testing.T) { diff --git a/probe/kubernetes/service.go b/probe/kubernetes/service.go index dea2a931a..8251db4e0 100644 --- a/probe/kubernetes/service.go +++ b/probe/kubernetes/service.go @@ -17,7 +17,7 @@ const ( // Service represents a Kubernetes service type Service interface { Meta - GetNode() report.Node + GetNode(probeID string) report.Node Selector() labels.Selector ClusterIP() string } @@ -47,10 +47,11 @@ func servicePortString(p apiv1.ServicePort) string { return fmt.Sprintf("%d:%d/%s", p.Port, p.NodePort, p.Protocol) } -func (s *service) GetNode() report.Node { +func (s *service) GetNode(probeID string) report.Node { latest := map[string]string{ IP: s.Spec.ClusterIP, Type: string(s.Spec.Type), + report.ControlProbeID: probeID, } if s.Spec.LoadBalancerIP != "" { latest[PublicIP] = s.Spec.LoadBalancerIP diff --git a/probe/kubernetes/statefulset.go b/probe/kubernetes/statefulset.go index 3d3bc973b..7cd5ae01e 100644 --- a/probe/kubernetes/statefulset.go +++ b/probe/kubernetes/statefulset.go @@ -14,7 +14,7 @@ import ( type StatefulSet interface { Meta Selector() (labels.Selector, error) - GetNode() report.Node + GetNode(probeID string) report.Node } type statefulSet struct { @@ -38,15 +38,16 @@ func (s *statefulSet) Selector() (labels.Selector, error) { return selector, nil } -func (s *statefulSet) GetNode() report.Node { +func (s *statefulSet) GetNode(probeID string) report.Node { desiredReplicas := 1 if s.Spec.Replicas != nil { desiredReplicas = int(*s.Spec.Replicas) } latests := map[string]string{ - NodeType: "StatefulSet", - DesiredReplicas: fmt.Sprint(desiredReplicas), - Replicas: fmt.Sprint(s.Status.Replicas), + NodeType: "StatefulSet", + DesiredReplicas: fmt.Sprint(desiredReplicas), + Replicas: fmt.Sprint(s.Status.Replicas), + report.ControlProbeID: probeID, } if s.Status.ObservedGeneration != nil { latests[ObservedGeneration] = fmt.Sprint(*s.Status.ObservedGeneration)