diff --git a/probe/kubernetes/client.go b/probe/kubernetes/client.go index 0cba08dfa..352ba855d 100644 --- a/probe/kubernetes/client.go +++ b/probe/kubernetes/client.go @@ -55,6 +55,7 @@ type Client interface { WalkStorageClasses(f func(StorageClass) error) error WalkVolumeSnapshots(f func(VolumeSnapshot) error) error WalkVolumeSnapshotData(f func(VolumeSnapshotData) error) error + WalkJobs(f func(Job) error) error WatchPods(f func(Event, Pod)) @@ -456,6 +457,16 @@ func (c *client) WalkVolumeSnapshotData(f func(VolumeSnapshotData) error) error return nil } +func (c *client) WalkJobs(f func(Job) error) error { + for _, m := range c.jobStore.List() { + job := m.(*apibatchv1.Job) + if err := f(NewJob(job)); err != nil { + return err + } + } + return nil +} + func (c *client) CloneVolumeSnapshot(namespaceID, volumeSnapshotID, persistentVolumeClaimID, capacity string) error { var scName string var claimSize string diff --git a/probe/kubernetes/controls.go b/probe/kubernetes/controls.go index cf39dfed8..71d4f1c93 100644 --- a/probe/kubernetes/controls.go +++ b/probe/kubernetes/controls.go @@ -91,6 +91,10 @@ func (r *Reporter) describeStoragelass(req xfer.Request, storageClassID string) return r.describe(req, "", storageClassID, ResourceMap["StorageClass"], apimeta.RESTMapping{}) } +func (r *Reporter) describeJob(req xfer.Request, namespaceID, jobID string) xfer.Response { + return r.describe(req, namespaceID, jobID, ResourceMap["Job"], apimeta.RESTMapping{}) +} + func (r *Reporter) describeVolumeSnapshot(req xfer.Request, namespaceID, volumeSnapshotID, _, _ string) xfer.Response { restMapping := apimeta.RESTMapping{ Resource: schema.GroupVersionResource{ @@ -204,6 +208,8 @@ func (r *Reporter) Describe() func(xfer.Request) xfer.Response { f = r.CaptureVolumeSnapshot(r.describeVolumeSnapshot) case "": f = r.CaptureVolumeSnapshotData(r.describeVolumeSnapshotData) + case "": + f = r.CaptureJob(r.describeJob) default: return xfer.ResponseErrorf("Node not found: %s", req.NodeID) } @@ -447,6 +453,27 @@ func (r *Reporter) CaptureVolumeSnapshotData(f func(xfer.Request, string) xfer.R } } +// CaptureJob is exported for testing +func (r *Reporter) CaptureJob(f func(xfer.Request, string, string) xfer.Response) func(xfer.Request) xfer.Response { + return func(req xfer.Request) xfer.Response { + uid, ok := report.ParseJobNodeID(req.NodeID) + if !ok { + return xfer.ResponseErrorf("Invalid ID: %s", req.NodeID) + } + var job Job + r.client.WalkJobs(func(c Job) error { + if c.UID() == uid { + job = c + } + return nil + }) + if job == nil { + return xfer.ResponseErrorf("Job not found: %s", uid) + } + return f(req, job.Namespace(), job.Name()) + } +} + // ScaleUp is the control to scale up a deployment func (r *Reporter) ScaleUp(req xfer.Request, namespace, id string) xfer.Response { return xfer.ResponseError(r.client.ScaleUp(namespace, id)) diff --git a/probe/kubernetes/job.go b/probe/kubernetes/job.go new file mode 100644 index 000000000..5d611233f --- /dev/null +++ b/probe/kubernetes/job.go @@ -0,0 +1,47 @@ +package kubernetes + +import ( + batchv1 "k8s.io/api/batch/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/labels" + + "github.com/weaveworks/scope/report" +) + +//Job represents a Kubernetes job +type Job interface { + Meta + Selector() (labels.Selector, error) + GetNode(probeID string) report.Node +} + +type job struct { + *batchv1.Job + Meta +} + +// NewJob creates a new job. +func NewJob(j *batchv1.Job) Job { + return &job{ + Job: j, + Meta: meta{j.ObjectMeta}, + } +} + +func (j *job) Selector() (labels.Selector, error) { + selector, err := metav1.LabelSelectorAsSelector(j.Spec.Selector) + if err != nil { + return nil, err + } + return selector, nil +} + +func (j *job) GetNode(probeID string) report.Node { + latests := map[string]string{ + NodeType: "Job", + report.ControlProbeID: probeID, + } + return j.MetaNode(report.MakeJobNodeID(j.UID())). + WithLatests(latests). + WithLatestActiveControls(Describe) +} diff --git a/probe/kubernetes/reporter.go b/probe/kubernetes/reporter.go index 32583d975..ff2021105 100644 --- a/probe/kubernetes/reporter.go +++ b/probe/kubernetes/reporter.go @@ -145,6 +145,16 @@ var ( VolumeSnapshotName: {ID: VolumeSnapshotName, Label: "Volume snapshot", From: report.FromLatest, Priority: 3}, } + JobMetadataTemplates = report.MetadataTemplates{ + NodeType: {ID: NodeType, Label: "Type", From: report.FromLatest, Priority: 1}, + Name: {ID: Name, Label: "Name", From: report.FromLatest, Priority: 2}, + Namespace: {ID: Namespace, Label: "Namespace", From: report.FromLatest, Priority: 3}, + Created: {ID: Created, Label: "Created", From: report.FromLatest, Datatype: report.DateTime, Priority: 4}, + report.Pod: {ID: report.Pod, Label: "# Pods", From: report.FromCounters, Datatype: report.Number, Priority: 5}, + } + + JobMetricTemplates = PodMetricTemplates + TableTemplates = report.TableTemplates{ LabelPrefix: { ID: LabelPrefix, @@ -313,7 +323,11 @@ func (r *Reporter) Report() (report.Report, error) { if err != nil { return result, err } - podTopology, err := r.podTopology(services, deployments, daemonSets, statefulSets, cronJobs) + jobTopology, jobs, err := r.jobTopology() + if err != nil { + return result, err + } + podTopology, err := r.podTopology(services, deployments, daemonSets, statefulSets, cronJobs, jobs) if err != nil { return result, err } @@ -341,6 +355,7 @@ func (r *Reporter) Report() (report.Report, error) { if err != nil { return result, err } + result.Pod = result.Pod.Merge(podTopology) result.Service = result.Service.Merge(serviceTopology) result.DaemonSet = result.DaemonSet.Merge(daemonSetTopology) @@ -353,6 +368,7 @@ func (r *Reporter) Report() (report.Report, error) { result.StorageClass = result.StorageClass.Merge(storageClassTopology) result.VolumeSnapshot = result.VolumeSnapshot.Merge(volumeSnapshotTopology) result.VolumeSnapshotData = result.VolumeSnapshotData.Merge(volumeSnapshotDataTopology) + result.Job = result.Job.Merge(jobTopology) return result, nil } @@ -525,6 +541,21 @@ func (r *Reporter) volumeSnapshotDataTopology() (report.Topology, []VolumeSnapsh return result, volumeSnapshotData, err } +func (r *Reporter) jobTopology() (report.Topology, []Job, error) { + jobs := []Job{} + result := report.MakeTopology(). + WithMetadataTemplates(JobMetadataTemplates). + WithMetricTemplates(JobMetricTemplates). + WithTableTemplates(TableTemplates) + result.Controls.AddControl(DescribeControl) + err := r.client.WalkJobs(func(c Job) error { + result.AddNode(c.GetNode(r.probeID)) + jobs = append(jobs, c) + return nil + }) + return result, jobs, err +} + type labelledChild interface { Labels() map[string]string AddParent(string, string) @@ -540,7 +571,7 @@ func match(namespace string, selector labels.Selector, topology, id string) func } } -func (r *Reporter) podTopology(services []Service, deployments []Deployment, daemonSets []DaemonSet, statefulSets []StatefulSet, cronJobs []CronJob) (report.Topology, error) { +func (r *Reporter) podTopology(services []Service, deployments []Deployment, daemonSets []DaemonSet, statefulSets []StatefulSet, cronJobs []CronJob, jobs []Job) (report.Topology, error) { var ( pods = report.MakeTopology(). WithMetadataTemplates(PodMetadataTemplates). @@ -619,6 +650,18 @@ func (r *Reporter) podTopology(services []Service, deployments []Deployment, dae report.MakeCronJobNodeID(cronJob.UID()), )) } + for _, job := range jobs { + selector, err := job.Selector() + if err != nil { + return pods, err + } + selectors = append(selectors, match( + job.Namespace(), + selector, + report.Job, + report.MakeJobNodeID(job.UID()), + )) + } } var localPodUIDs map[string]struct{} diff --git a/probe/kubernetes/reporter_test.go b/probe/kubernetes/reporter_test.go index 0487244ce..d742fb79d 100644 --- a/probe/kubernetes/reporter_test.go +++ b/probe/kubernetes/reporter_test.go @@ -172,6 +172,9 @@ func (c *mockClient) WalkVolumeSnapshots(f func(kubernetes.VolumeSnapshot) error func (c *mockClient) WalkVolumeSnapshotData(f func(kubernetes.VolumeSnapshotData) error) error { return nil } +func (c *mockClient) WalkJobs(f func(kubernetes.Job) error) error { + return nil +} func (*mockClient) WatchPods(func(kubernetes.Event, kubernetes.Pod)) {} func (c *mockClient) GetLogs(namespaceID, podName string, _ []string) (io.ReadCloser, error) { r, ok := c.logs[namespaceID+";"+podName] diff --git a/render/detailed/summary.go b/render/detailed/summary.go index a6c24c95d..6f42f79f8 100644 --- a/render/detailed/summary.go +++ b/render/detailed/summary.go @@ -78,6 +78,7 @@ var renderers = map[string]func(BasicNodeSummary, report.Node) BasicNodeSummary{ report.DaemonSet: podGroupNodeSummary, report.StatefulSet: podGroupNodeSummary, report.CronJob: podGroupNodeSummary, + report.Job: podGroupNodeSummary, report.ECSTask: ecsTaskNodeSummary, report.ECSService: ecsServiceNodeSummary, report.SwarmService: swarmServiceNodeSummary, @@ -101,6 +102,7 @@ var primaryAPITopology = map[string]string{ report.DaemonSet: "kube-controllers", report.StatefulSet: "kube-controllers", report.CronJob: "kube-controllers", + report.Job: "kube-controllers", report.Service: "services", report.ECSTask: "ecs-tasks", report.ECSService: "ecs-services", @@ -321,6 +323,7 @@ var podGroupNodeTypeName = map[string]string{ report.DaemonSet: "DaemonSet", report.StatefulSet: "StatefulSet", report.CronJob: "CronJob", + report.Job: "Job", } func podGroupNodeSummary(base BasicNodeSummary, n report.Node) BasicNodeSummary { diff --git a/render/pod.go b/render/pod.go index 7d8e489a2..548f65a90 100644 --- a/render/pod.go +++ b/render/pod.go @@ -30,6 +30,7 @@ func renderKubernetesTopologies(rpt report.Report) bool { &rpt.PersistentVolume, &rpt.PersistentVolumeClaim, &rpt.StorageClass, + &rpt.Job, } for _, t := range topologies { if len(t.Nodes) > 0 { @@ -112,7 +113,7 @@ var PodServiceRenderer = ConditionalRenderer(renderKubernetesTopologies, // not memoised var KubeControllerRenderer = ConditionalRenderer(renderKubernetesTopologies, renderParents( - report.Pod, []string{report.Deployment, report.DaemonSet, report.StatefulSet, report.CronJob}, UnmanagedID, + report.Pod, []string{report.Deployment, report.DaemonSet, report.StatefulSet, report.CronJob, report.Job}, UnmanagedID, PodRenderer, ), ) diff --git a/render/selectors.go b/render/selectors.go index 22856fdfd..f09c027be 100644 --- a/render/selectors.go +++ b/render/selectors.go @@ -30,6 +30,7 @@ var ( SelectDaemonSet = TopologySelector(report.DaemonSet) SelectStatefulSet = TopologySelector(report.StatefulSet) SelectCronJob = TopologySelector(report.CronJob) + SelectJob = TopologySelector(report.Job) SelectECSTask = TopologySelector(report.ECSTask) SelectECSService = TopologySelector(report.ECSService) SelectSwarmService = TopologySelector(report.SwarmService) diff --git a/report/id.go b/report/id.go index be893a214..78db8e43c 100644 --- a/report/id.go +++ b/report/id.go @@ -140,6 +140,12 @@ var ( // ParseCronJobNodeID parses a cronjob node ID ParseCronJobNodeID = parseSingleComponentID("cronjob") + // MakeJobNodeID produces a job node ID from its composite parts. + MakeJobNodeID = makeSingleComponentID("job") + + // ParseJobNodeID parses a job node ID + ParseJobNodeID = parseSingleComponentID("job") + // MakeNamespaceNodeID produces a namespace node ID from its composite parts. MakeNamespaceNodeID = makeSingleComponentID("namespace") diff --git a/report/report.go b/report/report.go index 63a9b2d46..52b6fddbe 100644 --- a/report/report.go +++ b/report/report.go @@ -33,6 +33,7 @@ const ( StorageClass = "storage_class" VolumeSnapshot = "volume_snapshot" VolumeSnapshotData = "volume_snapshot_data" + Job = "job" // Shapes used for different nodes Circle = "circle" @@ -76,6 +77,7 @@ var topologyNames = []string{ StorageClass, VolumeSnapshot, VolumeSnapshotData, + Job, } // Report is the core data type. It's produced by probes, and consumed and @@ -182,6 +184,9 @@ type Report struct { // VolumeSnapshotData represent all Kubernetes Volume Snapshot Data on hosts running probes. VolumeSnapshotData Topology + // Job represent all Kubernetes Job on hosts running probes. + Job Topology + DNS DNSRecords // Sampling data for this report. @@ -297,6 +302,10 @@ func MakeReport() Report { WithTag(Camera). WithLabel("volume snapshot data", "volume snapshot data"), + Job: MakeTopology(). + WithShape(StorageSheet). + WithLabel("job", "jobs"), + DNS: DNSRecords{}, Sampling: Sampling{}, @@ -413,6 +422,8 @@ func (r *Report) topology(name string) *Topology { return &r.VolumeSnapshot case VolumeSnapshotData: return &r.VolumeSnapshotData + case Job: + return &r.Job } return nil }