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
This commit is contained in:
Roland Schilter
2018-03-16 12:39:57 -07:00
committed by GitHub
parent da6995bd87
commit 3de06b5a09
7 changed files with 58 additions and 31 deletions
+2 -1
View File
@@ -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)
+7 -6
View File
@@ -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)
+7 -6
View File
@@ -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,
})
}
+7 -10
View File
@@ -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
})
+26 -1
View File
@@ -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) {
+3 -2
View File
@@ -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
+6 -5
View File
@@ -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)