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)