From 7d93e2cfe7cba928c0e1a8c5dfa09368ca481e5e Mon Sep 17 00:00:00 2001 From: Mike Lang Date: Fri, 18 Nov 2016 18:06:13 -0800 Subject: [PATCH 01/19] merger: Pass pointers, not structs just in one function, where two of them are passed at once. This was causing errors because it was too large to fit the stack. --- app/merger.go | 11 ++++++----- 1 file changed, 6 insertions(+), 5 deletions(-) diff --git a/app/merger.go b/app/merger.go index be4d149b2..ab0f9379f 100644 --- a/app/merger.go +++ b/app/merger.go @@ -48,14 +48,15 @@ func (smartMerger) Merge(reports []report.Report) report.Report { case 1: return reports[0] } - c := make(chan report.Report, l) + c := make(chan *report.Report, l) for _, r := range reports { - c <- r + c <- &r } for ; l > 1; l-- { - go func(left, right report.Report) { - c <- left.Merge(right) + go func(left, right *report.Report) { + r := left.Merge(*right) + c <- &r }(<-c, <-c) } - return <-c + return *<-c } From db23e64e9caef76272890e858f10e39b577783f7 Mon Sep 17 00:00:00 2001 From: Mike Lang Date: Tue, 22 Nov 2016 15:12:14 -0800 Subject: [PATCH 02/19] Add new topologies ECSTask and ECSService --- report/report.go | 27 +++++++++++++++++++++++++++ 1 file changed, 27 insertions(+) diff --git a/report/report.go b/report/report.go index c72aeb832..932756a8d 100644 --- a/report/report.go +++ b/report/report.go @@ -22,6 +22,8 @@ const ( ContainerImage = "container_image" Host = "host" Overlay = "overlay" + ECSService = "ecs_service" + ECSTask = "ecs_task" // Shapes used for different nodes Circle = "circle" @@ -81,6 +83,15 @@ type Report struct { // probes with each published report. Edges are not present. Host Topology + // ECS Task nodes are AWS ECS tasks, which represent a group of containers. + // Metadata is limited for now, more to come later. Edges are not present. + ECSTask Topology + + // ECS Service nodes are AWS ECS services, which represent a specification for a + // desired count of tasks with a task definition template. + // Metadata is limited for now, more to come later. Edges are not present. + ECSService Topology + // Overlay nodes are active peers in any software-defined network that's // overlaid on the infrastructure. The information is scraped by polling // their status endpoints. Edges could be present, but aren't currently. @@ -149,6 +160,14 @@ func MakeReport() Report { Overlay: MakeTopology(), + ECSTask: MakeTopology(). + WithShape(Heptagon). + WithLabel("task", "tasks"), + + ECSService: MakeTopology(). + WithShape(Heptagon). + WithLabel("service", "services"), + Sampling: Sampling{}, Window: 0, Plugins: xfer.MakePluginSpecs(), @@ -169,6 +188,8 @@ func (r Report) Copy() Report { Deployment: r.Deployment.Copy(), ReplicaSet: r.ReplicaSet.Copy(), Overlay: r.Overlay.Copy(), + ECSTask: r.ECSTask.Copy(), + ECSService: r.ECSService.Copy(), Sampling: r.Sampling, Window: r.Window, Plugins: r.Plugins.Copy(), @@ -190,6 +211,8 @@ func (r Report) Merge(other Report) Report { Deployment: r.Deployment.Merge(other.Deployment), ReplicaSet: r.ReplicaSet.Merge(other.ReplicaSet), Overlay: r.Overlay.Merge(other.Overlay), + ECSTask: r.ECSTask.Merge(other.ECSTask), + ECSService: r.ECSService.Merge(other.ECSService), Sampling: r.Sampling.Merge(other.Sampling), Window: r.Window + other.Window, Plugins: r.Plugins.Merge(other.Plugins), @@ -219,6 +242,8 @@ func (r *Report) WalkTopologies(f func(*Topology)) { f(&r.ReplicaSet) f(&r.Host) f(&r.Overlay) + f(&r.ECSTask) + f(&r.ECSService) } // Topology gets a topology by name @@ -234,6 +259,8 @@ func (r Report) Topology(name string) (Topology, bool) { ReplicaSet: r.ReplicaSet, Host: r.Host, Overlay: r.Overlay, + ECSTask: r.ECSTask, + ECSService: r.ECSService, }[name] return t, ok } From 511f6dad6a9e85aef68cea45b3503259b6e9885b Mon Sep 17 00:00:00 2001 From: Mike Lang Date: Tue, 22 Nov 2016 15:12:57 -0800 Subject: [PATCH 03/19] Add report tagger for populating ECS topologies --- probe/awsecs/client.go | 138 +++++++++++++++++++++++++++++++++++++++ probe/awsecs/reporter.go | 118 +++++++++++++++++++++++++++++++++ 2 files changed, 256 insertions(+) create mode 100644 probe/awsecs/client.go create mode 100644 probe/awsecs/reporter.go diff --git a/probe/awsecs/client.go b/probe/awsecs/client.go new file mode 100644 index 000000000..bf5b1a38d --- /dev/null +++ b/probe/awsecs/client.go @@ -0,0 +1,138 @@ +package awsecs + +import ( + "sync" + + log "github.com/Sirupsen/logrus" + "github.com/aws/aws-sdk-go/aws" + "github.com/aws/aws-sdk-go/aws/ec2metadata" + "github.com/aws/aws-sdk-go/aws/session" + "github.com/aws/aws-sdk-go/service/ecs" +) + +// a wrapper around an AWS client that makes all the needed calls and just exposes the final results +type ecsClient struct { + client *ecs.ECS + cluster string +} + +func newClient(cluster string) (*ecsClient, error) { + sess := session.New() + + region, err := ec2metadata.New(sess).Region() + if err != nil { + return nil, err + } + + return &ecsClient{ + client: ecs.New(sess, &aws.Config{Region: aws.String(region)}), + cluster: cluster, + }, nil +} + +// returns a map from deployment ids to service names +// cannot fail as it will attempt to deliver partial results, though that may end up being no results +func (c ecsClient) getDeploymentMap() map[string]string { + results := make(map[string]string) + lock := sync.Mutex{} // lock mediates access to results + + group := sync.WaitGroup{} + + err := c.client.ListServicesPages( + &ecs.ListServicesInput{Cluster: &c.cluster}, + func(page *ecs.ListServicesOutput, lastPage bool) bool { + // describe each page of 10 (the max for one describe command) concurrently + group.Add(1) + go func() { + defer group.Done() + + resp, err := c.client.DescribeServices(&ecs.DescribeServicesInput{ + Cluster: &c.cluster, + Services: page.ServiceArns, + }) + if err != nil { + // rather than trying to propogate errors up, just log a warning here + log.Warnf("Error describing some ECS services, ECS service report may be incomplete: %v", err) + return + } + + for _, failure := range resp.Failures { + // log the failures but still continue with what succeeded + log.Warnf("Failed to describe ECS service %s, ECS service report may be incomplete: %s", failure.Arn, failure.Reason) + } + + lock.Lock() + for _, service := range resp.Services { + for _, deployment := range service.Deployments { + results[*deployment.Id] = *service.ServiceName + } + } + lock.Unlock() + }() + return true + }, + ) + group.Wait() + + if err != nil { + // We want to still return partial results if we have any, so just log a warning + log.Warnf("Error listing ECS services, ECS service report may be incomplete: %v", err) + } + return results +} + +// returns a map from task ARNs to deployment ids +func (c ecsClient) getTaskDeployments(taskArns []string) (map[string]string, error) { + taskPtrs := make([]*string, len(taskArns)) + for i := range taskArns { + taskPtrs[i] = &taskArns[i] + } + + // You'd think there's a limit on how many tasks can be described here, + // but the docs don't mention anything. + resp, err := c.client.DescribeTasks(&ecs.DescribeTasksInput{ + Cluster: &c.cluster, + Tasks: taskPtrs, + }) + if err != nil { + return nil, err + } + + for _, failure := range resp.Failures { + // log the failures but still continue with what succeeded + log.Warnf("Failed to describe ECS task %s, ECS service report may be incomplete: %s", failure.Arn, failure.Reason) + } + + results := make(map[string]string) + for _, task := range resp.Tasks { + results[*task.TaskArn] = *task.StartedBy + } + return results, nil +} + +// returns a map from task ARNs to service names +func (c ecsClient) getTaskServices(taskArns []string) (map[string]string, error) { + deploymentMapChan := make(chan map[string]string) + go func() { + deploymentMapChan <- c.getDeploymentMap() + }() + + // do these two fetches in parallel + taskDeployments, err := c.getTaskDeployments(taskArns) + deploymentMap := <-deploymentMapChan + + if err != nil { + return nil, err + } + + results := make(map[string]string) + for taskArn, depID := range taskDeployments { + // Note not all tasks map to a deployment, or we could otherwise mismatch due to races. + // It's safe to just ignore all these cases and consider them "non-service" tasks. + if service, ok := deploymentMap[depID]; ok { + results[taskArn] = service + } + } + + return results, nil +} diff --git a/probe/awsecs/reporter.go b/probe/awsecs/reporter.go new file mode 100644 index 000000000..3c7ef70f7 --- /dev/null +++ b/probe/awsecs/reporter.go @@ -0,0 +1,118 @@ +package awsecs + +import ( + "fmt" + + log "github.com/Sirupsen/logrus" + "github.com/weaveworks/scope/probe/docker" + "github.com/weaveworks/scope/report" +) + +type taskInfo struct { + containerIDs []string + family string +} + +// return map from cluster to map of task arns to task infos +func getLabelInfo(rpt report.Report) map[string]map[string]*taskInfo { + results := make(map[string]map[string]*taskInfo) + log.Debug("scanning for ECS containers") + for nodeID, node := range rpt.Container.Nodes { + + taskArn, taskArnOk := node.Latest.Lookup(docker.LabelPrefix + "com.amazonaws.ecs.task-arn") + cluster, clusterOk := node.Latest.Lookup(docker.LabelPrefix + "com.amazonaws.ecs.cluster") + family, familyOk := node.Latest.Lookup(docker.LabelPrefix + "com.amazonaws.ecs.task-definition-family") + + if taskArnOk && clusterOk && familyOk { + taskMap, ok := results[cluster] + if !ok { + taskMap = make(map[string]*taskInfo) + results[cluster] = taskMap + } + + task, ok := taskMap[taskArn] + if !ok { + task = &taskInfo{containerIDs: make([]string, 0), family: family} + taskMap[taskArn] = task + } + + task.containerIDs = append(task.containerIDs, nodeID) + } + } + log.Debug("Got ECS container info: %v", results) + return results +} + +// implements Tagger +type Reporter struct { +} + +// Tag needed for Tagger +func (r Reporter) Tag(rpt report.Report) (report.Report, error) { + rpt = rpt.Copy() + + clusterMap := getLabelInfo(rpt) + + for cluster, taskMap := range clusterMap { + log.Debugf("Fetching ECS info for cluster %v with %v tasks", cluster, len(taskMap)) + + client, err := newClient(cluster) + if err != nil { + return rpt, err + } + + taskArns := make([]string, 0, len(taskMap)) + for taskArn := range taskMap { + taskArns = append(taskArns, taskArn) + } + + taskServices, err := client.getTaskServices(taskArns) + if err != nil { + return rpt, err + } + + // Create all the services first + unique := make(map[string]bool) + for _, serviceName := range taskServices { + if !unique[serviceName] { + rpt.ECSService = rpt.ECSService.AddNode(report.MakeNode(serviceNodeID(serviceName))) + unique[serviceName] = true + } + } + log.Debugf("Created %v ECS service nodes", len(taskServices)) + + for taskArn, info := range taskMap { + + // new task node + node := report.MakeNodeWith(taskNodeID(taskArn), map[string]string{"family": info.family}) + + rpt.ECSTask = rpt.ECSTask.AddNode(node) + + for _, containerID := range info.containerIDs { + // TODO set task node as parent of container + log.Debugf("task %v has container %v", taskArn, containerID) + } + + if serviceName, ok := taskServices[taskArn]; ok { + // TODO set service node as parent of task node + log.Debugf("service %v has task %v", serviceName, taskArn) + } + } + + } + + return rpt, nil +} + +// Name needed for Tagger +func (r Reporter) Name() string { + return "awsecs" +} + +func serviceNodeID(id string) string { + return fmt.Sprintf("%s;ECSService", id) +} + +func taskNodeID(id string) string { + return fmt.Sprintf("%s;ECSTask", id) +} From 88499b4e9dec601249eae9b6cdb0b2946e73a935 Mon Sep 17 00:00:00 2001 From: Mike Lang Date: Tue, 22 Nov 2016 15:13:36 -0800 Subject: [PATCH 04/19] Add --probe.ecs flag to enable running the ECS probe tagger --- prog/main.go | 5 +++++ prog/probe.go | 5 +++++ 2 files changed, 10 insertions(+) diff --git a/prog/main.go b/prog/main.go index 956549cea..2f8999968 100644 --- a/prog/main.go +++ b/prog/main.go @@ -104,6 +104,8 @@ type probeFlags struct { kubernetesEnabled bool kubernetesConfig kubernetes.ClientConfig + ecsEnabled bool + weaveEnabled bool weaveAddr string weaveHostname string @@ -282,6 +284,9 @@ func main() { flag.StringVar(&flags.probe.kubernetesConfig.User, "probe.kubernetes.user", "", "The name of the kubeconfig user to use") flag.StringVar(&flags.probe.kubernetesConfig.Username, "probe.kubernetes.username", "", "Username for basic authentication to the API server") + // AWS ECS + flag.BoolVar(&flags.probe.ecsEnabled, "probe.ecs", false, "collect ecs-related attributes for containers on this node") + // Weave flag.StringVar(&flags.probe.weaveAddr, "probe.weave.addr", "127.0.0.1:6784", "IP address & port of the Weave router") flag.StringVar(&flags.probe.weaveHostname, "probe.weave.hostname", app.DefaultHostname, "Hostname to lookup in WeaveDNS") diff --git a/prog/probe.go b/prog/probe.go index 5d638ea81..2cb26bab4 100644 --- a/prog/probe.go +++ b/prog/probe.go @@ -23,6 +23,7 @@ import ( "github.com/weaveworks/scope/common/xfer" "github.com/weaveworks/scope/probe" "github.com/weaveworks/scope/probe/appclient" + "github.com/weaveworks/scope/probe/awsecs" "github.com/weaveworks/scope/probe/controls" "github.com/weaveworks/scope/probe/docker" "github.com/weaveworks/scope/probe/endpoint" @@ -201,6 +202,10 @@ func probeMain(flags probeFlags, targets []appclient.Target) { } } + if flags.ecsEnabled { + p.AddTagger(awsecs.Reporter{}) + } + if flags.weaveEnabled { client := weave.NewClient(sanitize.URL("http://", 6784, "")(flags.weaveAddr)) weave, err := overlay.NewWeave(hostID, client, dockerEndpoint) From a2d329dee77a62e21f1fbbf9281048627666dcad Mon Sep 17 00:00:00 2001 From: Mike Lang Date: Tue, 22 Nov 2016 15:37:21 -0800 Subject: [PATCH 05/19] ecs reporter: Associate containers with ECS tasks and services as parents --- probe/awsecs/reporter.go | 19 ++++++++++++------- 1 file changed, 12 insertions(+), 7 deletions(-) diff --git a/probe/awsecs/reporter.go b/probe/awsecs/reporter.go index 3c7ef70f7..895bf8e8e 100644 --- a/probe/awsecs/reporter.go +++ b/probe/awsecs/reporter.go @@ -85,18 +85,23 @@ func (r Reporter) Tag(rpt report.Report) (report.Report, error) { // new task node node := report.MakeNodeWith(taskNodeID(taskArn), map[string]string{"family": info.family}) - rpt.ECSTask = rpt.ECSTask.AddNode(node) - for _, containerID := range info.containerIDs { - // TODO set task node as parent of container - log.Debugf("task %v has container %v", taskArn, containerID) + // parents sets to merge into all matching container nodes + parentsSets := report.MakeSets() + parentsSets.Add(report.ECSTask, report.MakeStringSet(taskNodeID(taskArn))) + if serviceName, ok := taskServices[taskArn]; ok { + parentsSets.Add(report.ECSService, report.MakeStringSet(serviceNodeID(serviceName))) } - if serviceName, ok := taskServices[taskArn]; ok { - // TODO set service node as parent of task node - log.Debugf("service %v has task %v", serviceName, taskArn) + for _, containerID := range info.containerIDs { + if containerNode, ok := rpt.Container.Nodes[containerID]; ok { + rpt.Container.Nodes[containerID] = containerNode.WithParents(parentsSets) + } else { + log.Warnf("Got task info for non-existent container %v, this shouldn't be able to happen", containerID) + } } + } } From 78775bbdb835437169ba265f317f45383664ebb4 Mon Sep 17 00:00:00 2001 From: Alfonso Acosta Date: Wed, 23 Nov 2016 16:08:46 +0000 Subject: [PATCH 06/19] Initial rendering for ECS (not working yet) --- app/api_topologies.go | 86 +++++++++++++++++++++++--------------- probe/awsecs/reporter.go | 25 +++++------ render/detailed/summary.go | 14 +++++++ render/ecs.go | 19 +++++++++ render/selectors.go | 2 + report/id.go | 12 ++++++ report/report.go | 4 +- 7 files changed, 114 insertions(+), 48 deletions(-) create mode 100644 render/ecs.go diff --git a/app/api_topologies.go b/app/api_topologies.go index e9bcb2a67..848cee933 100644 --- a/app/api_topologies.go +++ b/app/api_topologies.go @@ -16,18 +16,21 @@ import ( ) const ( - apiTopologyURL = "/api/topology/" - processesTopologyDescID = "processes" - processesByNameTopologyDescID = "processes-by-name" - containerLabelFiltersGroupID = "container_label_filters_group" - containersTopologyDescID = "containers" - containersByHostnameTopologyDescID = "containers-by-hostname" - containersByImageTopologyDescID = "containers-by-image" - podsTopologyDescID = "pods" - replicaSetsTopologyDescID = "replica-sets" - deploymentsTopologyDescID = "deployments" - servicesTopologyDescID = "services" - hostsTopologyDescID = "hosts" + apiTopologyURL = "/api/topology/" + processesID = "processes" + processesByNameID = "processes-by-name" + containerLabelFiltersGroupID = "container-label-filters-group" + containersID = "containers" + containersByHostnameID = "containers-by-hostname" + containersByImageID = "containers-by-image" + podsID = "pods" + replicaSetsID = "replica-sets" + deploymentsID = "deployments" + servicesID = "services" + hostsID = "hosts" + weaveID = "weave" + ecsTasksID = "ecs-tasks" + ecsServicesID = "ecs-services" ) var ( @@ -53,7 +56,10 @@ func AddInitialTopologiesToRegistry(registry *Registry) { { ID: containerLabelFiltersGroupID, Default: "application", - Options: []APITopologyOption{{Value: "all", Label: "All", filter: nil, filterPseudo: false}, {Value: "system", Label: "System Containers", filter: render.IsSystem, filterPseudo: false}, {Value: "notsystem", Label: "Application Containers", filter: render.IsApplication, filterPseudo: false}}, + Options: []APITopologyOption{ + {Value: "all", Label: "All", filter: nil, filterPseudo: false}, + {Value: "system", Label: "System Containers", filter: render.IsSystem, filterPseudo: false}, + {Value: "notsystem", Label: "Application Containers", filter: render.IsApplication, filterPseudo: false}}, }, { ID: "stopped", @@ -90,7 +96,7 @@ func AddInitialTopologiesToRegistry(registry *Registry) { // be the verb to get to that state registry.Add( APITopologyDesc{ - id: processesTopologyDescID, + id: processesID, renderer: render.FilterUnconnected(render.ProcessWithContainerNameRenderer), Name: "Processes", Rank: 1, @@ -98,71 +104,85 @@ func AddInitialTopologiesToRegistry(registry *Registry) { HideIfEmpty: true, }, APITopologyDesc{ - id: processesByNameTopologyDescID, - parent: "processes", + id: processesByNameID, + parent: processesID, renderer: render.FilterUnconnected(render.ProcessNameRenderer), Name: "by name", Options: unconnectedFilter, HideIfEmpty: true, }, APITopologyDesc{ - id: containersTopologyDescID, + id: containersID, renderer: render.ContainerWithImageNameRenderer, Name: "Containers", Rank: 2, Options: containerFilters, }, APITopologyDesc{ - id: containersByHostnameTopologyDescID, - parent: "containers", + id: containersByHostnameID, + parent: containersID, renderer: render.ContainerHostnameRenderer, Name: "by DNS name", Options: containerFilters, }, APITopologyDesc{ - id: containersByImageTopologyDescID, - parent: "containers", + id: containersByImageID, + parent: containersID, renderer: render.ContainerImageRenderer, Name: "by image", Options: containerFilters, }, APITopologyDesc{ - id: podsTopologyDescID, + id: podsID, renderer: render.PodRenderer, Name: "Pods", Rank: 3, HideIfEmpty: true, }, APITopologyDesc{ - id: replicaSetsTopologyDescID, - parent: "pods", + id: replicaSetsID, + parent: podsID, renderer: render.ReplicaSetRenderer, Name: "replica sets", HideIfEmpty: true, }, APITopologyDesc{ - id: deploymentsTopologyDescID, - parent: "pods", + id: deploymentsID, + parent: podsID, renderer: render.DeploymentRenderer, Name: "deployments", HideIfEmpty: true, }, APITopologyDesc{ - id: servicesTopologyDescID, - parent: "pods", + id: servicesID, + parent: podsID, renderer: render.PodServiceRenderer, Name: "services", HideIfEmpty: true, }, APITopologyDesc{ - id: hostsTopologyDescID, + id: ecsTasksID, + renderer: render.ECSTaskRenderer, + Name: "Tasks", + Rank: 3, + HideIfEmpty: true, + }, + APITopologyDesc{ + id: ecsServicesID, + parent: ecsTasksID, + renderer: render.ECSServiceRenderer, + Name: "services", + HideIfEmpty: true, + }, + APITopologyDesc{ + id: hostsID, renderer: render.HostRenderer, Name: "Hosts", Rank: 4, }, APITopologyDesc{ - id: "weave", - parent: "hosts", + id: weaveID, + parent: hostsID, renderer: render.WeaveRenderer, Name: "weave net", HideIfEmpty: true, @@ -206,7 +226,7 @@ func updateFilters(rpt report.Report, topologies []APITopologyDesc) []APITopolog } sort.Strings(ns) for i, t := range topologies { - if t.id == podsTopologyDescID || t.id == servicesTopologyDescID || t.id == deploymentsTopologyDescID || t.id == replicaSetsTopologyDescID { + if t.id == podsID || t.id == servicesID || t.id == deploymentsID || t.id == replicaSetsID { topologies[i] = updateTopologyFilters(t, []APITopologyOptionGroup{ kubernetesFilters(ns...), k8sPseudoFilter, }) @@ -297,7 +317,7 @@ func AddContainerFilters(newFilters ...APITopologyOption) { func (r *Registry) AddContainerFilters(newFilters ...APITopologyOption) { r.Lock() defer r.Unlock() - for _, key := range []string{containersTopologyDescID, containersByHostnameTopologyDescID, containersByImageTopologyDescID} { + for _, key := range []string{containersID, containersByHostnameID, containersByImageID} { for i := range r.items[key].Options { if r.items[key].Options[i].ID == containerLabelFiltersGroupID { r.items[key].Options[i].Options = append(r.items[key].Options[i].Options, newFilters...) diff --git a/probe/awsecs/reporter.go b/probe/awsecs/reporter.go index 895bf8e8e..3932151ff 100644 --- a/probe/awsecs/reporter.go +++ b/probe/awsecs/reporter.go @@ -1,13 +1,15 @@ package awsecs import ( - "fmt" - log "github.com/Sirupsen/logrus" "github.com/weaveworks/scope/probe/docker" "github.com/weaveworks/scope/report" ) +const ( + TaskFamily = "ecs_task_family" +) + type taskInfo struct { containerIDs []string family string @@ -75,7 +77,8 @@ func (r Reporter) Tag(rpt report.Report) (report.Report, error) { unique := make(map[string]bool) for _, serviceName := range taskServices { if !unique[serviceName] { - rpt.ECSService = rpt.ECSService.AddNode(report.MakeNode(serviceNodeID(serviceName))) + serviceID := report.MakeECSServiceNodeID(serviceName) + rpt.ECSService = rpt.ECSService.AddNode(report.MakeNode(serviceID)) unique[serviceName] = true } } @@ -84,14 +87,16 @@ func (r Reporter) Tag(rpt report.Report) (report.Report, error) { for taskArn, info := range taskMap { // new task node - node := report.MakeNodeWith(taskNodeID(taskArn), map[string]string{"family": info.family}) + taskID := report.MakeECSTaskNodeID(taskArn) + node := report.MakeNodeWith(taskID, map[string]string{TaskFamily: info.family}) rpt.ECSTask = rpt.ECSTask.AddNode(node) // parents sets to merge into all matching container nodes parentsSets := report.MakeSets() - parentsSets.Add(report.ECSTask, report.MakeStringSet(taskNodeID(taskArn))) + parentsSets.Add(report.ECSTask, report.MakeStringSet(taskID)) if serviceName, ok := taskServices[taskArn]; ok { - parentsSets.Add(report.ECSService, report.MakeStringSet(serviceNodeID(serviceName))) + serviceID := report.MakeECSServiceNodeID(serviceName) + parentsSets.Add(report.ECSService, report.MakeStringSet(serviceID)) } for _, containerID := range info.containerIDs { @@ -113,11 +118,3 @@ func (r Reporter) Tag(rpt report.Report) (report.Report, error) { func (r Reporter) Name() string { return "awsecs" } - -func serviceNodeID(id string) string { - return fmt.Sprintf("%s;ECSService", id) -} - -func taskNodeID(id string) string { - return fmt.Sprintf("%s;ECSTask", id) -} diff --git a/render/detailed/summary.go b/render/detailed/summary.go index caa2d0a3e..e4ad49e24 100644 --- a/render/detailed/summary.go +++ b/render/detailed/summary.go @@ -4,6 +4,7 @@ import ( "fmt" "strings" + "github.com/weaveworks/scope/probe/awsecs" "github.com/weaveworks/scope/probe/docker" "github.com/weaveworks/scope/probe/endpoint" "github.com/weaveworks/scope/probe/host" @@ -68,6 +69,8 @@ var renderers = map[string]func(NodeSummary, report.Node) (NodeSummary, bool){ report.Service: podGroupNodeSummary, report.Deployment: podGroupNodeSummary, report.ReplicaSet: podGroupNodeSummary, + report.ECSTask: ecsTaskNodeSummary, + report.ECSService: ecsServiceNodeSummary, report.Host: hostNodeSummary, report.Overlay: weaveNodeSummary, } @@ -240,6 +243,17 @@ func podGroupNodeSummary(base NodeSummary, n report.Node) (NodeSummary, bool) { return base, true } +func ecsTaskNodeSummary(base NodeSummary, n report.Node) (NodeSummary, bool) { + base.Label, _ = n.Latest.Lookup(awsecs.TaskFamily) + return base, true +} + +func ecsServiceNodeSummary(base NodeSummary, n report.Node) (NodeSummary, bool) { + base.Label, _ = report.ParseECSServiceNodeID(n.ID) + base.Stack = true + return base, true +} + func hostNodeSummary(base NodeSummary, n report.Node) (NodeSummary, bool) { var ( hostname, _ = n.Latest.Lookup(host.HostName) diff --git a/render/ecs.go b/render/ecs.go new file mode 100644 index 000000000..595772b93 --- /dev/null +++ b/render/ecs.go @@ -0,0 +1,19 @@ +package render + +import ( + "github.com/weaveworks/scope/report" +) + +// ECSTaskRenderer is a Renderer for Amazon ECS tasks. +var ECSTaskRenderer = MakeFilter( + // TODO + func(n report.Node) bool { return true }, + SelectECSTask, +) + +// ECSServiceRenderer is a Renderer for Amazon ECS services. +var ECSServiceRenderer = MakeFilter( + // TODO + func(n report.Node) bool { return true }, + SelectECSService, +) diff --git a/render/selectors.go b/render/selectors.go index 286bc29ee..b4894b9c7 100644 --- a/render/selectors.go +++ b/render/selectors.go @@ -31,5 +31,7 @@ var ( SelectService = TopologySelector(report.Service) SelectDeployment = TopologySelector(report.Deployment) SelectReplicaSet = TopologySelector(report.ReplicaSet) + SelectECSTask = TopologySelector(report.ECSTask) + SelectECSService = TopologySelector(report.ECSService) SelectOverlay = TopologySelector(report.Overlay) ) diff --git a/report/id.go b/report/id.go index 0f5466e80..66a5b3418 100644 --- a/report/id.go +++ b/report/id.go @@ -119,6 +119,18 @@ var ( // ParseReplicaSetNodeID parses a replica set node ID ParseReplicaSetNodeID = parseSingleComponentID("replica_set") + + // MakeECSTaskNodeID produces a replica set node ID from its composite parts. + MakeECSTaskNodeID = makeSingleComponentID("ecs_task") + + // ParseECSTaskNodeID parses a replica set node ID + ParseECSTaskNodeID = parseSingleComponentID("ecs_task") + + // MakeECSServiceNodeID produces a replica set node ID from its composite parts. + MakeECSServiceNodeID = makeSingleComponentID("ecs_service") + + // ParseECSServiceNodeID parses a replica set node ID + ParseECSServiceNodeID = parseSingleComponentID("ecs_service") ) // makeSingleComponentID makes a single-component node id encoder diff --git a/report/report.go b/report/report.go index 932756a8d..890ab5be1 100644 --- a/report/report.go +++ b/report/report.go @@ -158,7 +158,9 @@ func MakeReport() Report { WithShape(Heptagon). WithLabel("replica set", "replica sets"), - Overlay: MakeTopology(), + Overlay: MakeTopology(). + WithShape(Circle). + WithLabel("peer", "peers"), ECSTask: MakeTopology(). WithShape(Heptagon). From 90c8b6eeed58626686f31832aa3024d2859bb918 Mon Sep 17 00:00:00 2001 From: Alfonso Acosta Date: Wed, 23 Nov 2016 18:35:25 +0000 Subject: [PATCH 07/19] Add ECS topologies to tagger --- probe/topology_tagger.go | 2 ++ 1 file changed, 2 insertions(+) diff --git a/probe/topology_tagger.go b/probe/topology_tagger.go index 354913387..cf41896d1 100644 --- a/probe/topology_tagger.go +++ b/probe/topology_tagger.go @@ -23,6 +23,8 @@ func (topologyTagger) Tag(r report.Report) (report.Report, error) { report.ContainerImage: &(r.ContainerImage), report.Pod: &(r.Pod), report.Service: &(r.Service), + report.ECSTask: &(r.ECSTask), + report.ECSService: &(r.ECSService), report.Host: &(r.Host), report.Overlay: &(r.Overlay), } { From f5b7b5bec2337b1e751d8135a4afdb16d3e59317 Mon Sep 17 00:00:00 2001 From: Mike Lang Date: Wed, 23 Nov 2016 20:46:23 -0800 Subject: [PATCH 08/19] ecs reporter: Fix a bug where parents weren't actually set --- probe/awsecs/reporter.go | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/probe/awsecs/reporter.go b/probe/awsecs/reporter.go index 3932151ff..bfbbc0c15 100644 --- a/probe/awsecs/reporter.go +++ b/probe/awsecs/reporter.go @@ -93,10 +93,10 @@ func (r Reporter) Tag(rpt report.Report) (report.Report, error) { // parents sets to merge into all matching container nodes parentsSets := report.MakeSets() - parentsSets.Add(report.ECSTask, report.MakeStringSet(taskID)) + parentsSets = parentsSets.Add(report.ECSTask, report.MakeStringSet(taskID)) if serviceName, ok := taskServices[taskArn]; ok { serviceID := report.MakeECSServiceNodeID(serviceName) - parentsSets.Add(report.ECSService, report.MakeStringSet(serviceID)) + parentsSets = parentsSets.Add(report.ECSService, report.MakeStringSet(serviceID)) } for _, containerID := range info.containerIDs { From b53de4317d3f8b7768346991565890bb5d707334 Mon Sep 17 00:00:00 2001 From: Mike Lang Date: Wed, 23 Nov 2016 20:58:14 -0800 Subject: [PATCH 09/19] Appease linter spellcheck and required comments and required comment formatting --- probe/awsecs/client.go | 2 +- probe/awsecs/reporter.go | 3 ++- 2 files changed, 3 insertions(+), 2 deletions(-) diff --git a/probe/awsecs/client.go b/probe/awsecs/client.go index bf5b1a38d..8d56a4fd3 100644 --- a/probe/awsecs/client.go +++ b/probe/awsecs/client.go @@ -51,7 +51,7 @@ func (c ecsClient) getDeploymentMap() map[string]string { Services: page.ServiceArns, }) if err != nil { - // rather than trying to propogate errors up, just log a warning here + // rather than trying to propagate errors up, just log a warning here log.Warnf("Error describing some ECS services, ECS service report may be incomplete: %v", err) return } diff --git a/probe/awsecs/reporter.go b/probe/awsecs/reporter.go index bfbbc0c15..8752a07b5 100644 --- a/probe/awsecs/reporter.go +++ b/probe/awsecs/reporter.go @@ -6,6 +6,7 @@ import ( "github.com/weaveworks/scope/report" ) +// TaskFamily is the key that stores the task family of an ECS Task const ( TaskFamily = "ecs_task_family" ) @@ -45,7 +46,7 @@ func getLabelInfo(rpt report.Report) map[string]map[string]*taskInfo { return results } -// implements Tagger +// Reporter implements Tagger type Reporter struct { } From f8b1f71f06e4aa9b89f225a32955dd12614f2296 Mon Sep 17 00:00:00 2001 From: Alfonso Acosta Date: Thu, 24 Nov 2016 12:18:19 +0000 Subject: [PATCH 10/19] Finish rendering for ECS topologies --- app/api_topologies.go | 6 ++-- render/ecs.go | 83 ++++++++++++++++++++++++++++++++++++++----- 2 files changed, 79 insertions(+), 10 deletions(-) diff --git a/app/api_topologies.go b/app/api_topologies.go index 848cee933..557c8f61c 100644 --- a/app/api_topologies.go +++ b/app/api_topologies.go @@ -35,7 +35,7 @@ const ( var ( topologyRegistry = MakeRegistry() - k8sPseudoFilter = APITopologyOptionGroup{ + unmanagedFilter = APITopologyOptionGroup{ ID: "pseudo", Default: "hide", Options: []APITopologyOption{ @@ -165,6 +165,7 @@ func AddInitialTopologiesToRegistry(registry *Registry) { renderer: render.ECSTaskRenderer, Name: "Tasks", Rank: 3, + Options: []APITopologyOptionGroup{unmanagedFilter}, HideIfEmpty: true, }, APITopologyDesc{ @@ -172,6 +173,7 @@ func AddInitialTopologiesToRegistry(registry *Registry) { parent: ecsTasksID, renderer: render.ECSServiceRenderer, Name: "services", + Options: []APITopologyOptionGroup{unmanagedFilter}, HideIfEmpty: true, }, APITopologyDesc{ @@ -228,7 +230,7 @@ func updateFilters(rpt report.Report, topologies []APITopologyDesc) []APITopolog for i, t := range topologies { if t.id == podsID || t.id == servicesID || t.id == deploymentsID || t.id == replicaSetsID { topologies[i] = updateTopologyFilters(t, []APITopologyOptionGroup{ - kubernetesFilters(ns...), k8sPseudoFilter, + kubernetesFilters(ns...), unmanagedFilter, }) } } diff --git a/render/ecs.go b/render/ecs.go index 595772b93..10ba0e75a 100644 --- a/render/ecs.go +++ b/render/ecs.go @@ -1,19 +1,86 @@ package render import ( + "strings" + + "github.com/weaveworks/scope/probe/docker" "github.com/weaveworks/scope/report" ) // ECSTaskRenderer is a Renderer for Amazon ECS tasks. -var ECSTaskRenderer = MakeFilter( - // TODO - func(n report.Node) bool { return true }, - SelectECSTask, +var ECSTaskRenderer = ConditionalRenderer(renderECSTopologies, + ApplyDecorators( + MakeMap( + PropagateSingleMetrics(report.Container), + MakeReduce( + MakeMap( + MapContainer2ECSTask, + ContainerWithImageNameRenderer, + ), + SelectECSTask, + ), + ), + ), ) // ECSServiceRenderer is a Renderer for Amazon ECS services. -var ECSServiceRenderer = MakeFilter( - // TODO - func(n report.Node) bool { return true }, - SelectECSService, +var ECSServiceRenderer = ConditionalRenderer(renderECSTopologies, + ApplyDecorators( + MakeMap( + PropagateSingleMetrics(report.ECSTask), + MakeReduce( + MakeMap( + Map2Parent(report.ECSService), + ECSTaskRenderer, + ), + SelectECSService, + ), + ), + ), ) + +// MapContainer2ECSTask maps container Nodes to ECS Task +// Nodes. +// +// If this function is given a node without an ECS Task parent +// (including other pseudo nodes), it will produce an "Unmanaged" +// pseudo node. +// +// TODO: worth merging with MapContainer2Pod? +func MapContainer2ECSTask(n report.Node, _ report.Networks) report.Nodes { + // Uncontained becomes unmanaged in the tasks view + if strings.HasPrefix(n.ID, MakePseudoNodeID(UncontainedID)) { + id := MakePseudoNodeID(UnmanagedID, report.ExtractHostID(n)) + node := NewDerivedPseudoNode(id, n) + return report.Nodes{id: node} + } + + // Propagate all pseudo nodes + if n.Topology == Pseudo { + return report.Nodes{n.ID: n} + } + + // Ignore non-running containers + if state, ok := n.Latest.Lookup(docker.ContainerState); ok && state != docker.StateRunning { + return report.Nodes{} + } + + taskIDSet, ok := n.Parents.Lookup(report.ECSTask) + if !ok || len(taskIDSet) == 0 { + id := MakePseudoNodeID(UnmanagedID, report.ExtractHostID(n)) + node := NewDerivedPseudoNode(id, n) + return report.Nodes{id: node} + } + nodeID := taskIDSet[0] + node := NewDerivedNode(nodeID, n).WithTopology(report.ECSTask) + // Propagate parent service + if serviceIDSet, ok := n.Parents.Lookup(report.ECSService); ok { + node = node.WithParents(report.MakeSets().Add(report.ECSService, serviceIDSet)) + } + node.Counters = node.Counters.Add(n.Topology, 1) + return report.Nodes{nodeID: node} +} + +func renderECSTopologies(rpt report.Report) bool { + return len(rpt.ECSTask.Nodes)+len(rpt.ECSService.Nodes) >= 1 +} From ab1d2d2c6ddaaa4b7ca25309a41bcfbebd9868f1 Mon Sep 17 00:00:00 2001 From: Alfonso Acosta Date: Thu, 24 Nov 2016 12:20:32 +0000 Subject: [PATCH 11/19] Add checkpoint flag for ECS --- prog/probe.go | 3 +++ 1 file changed, 3 insertions(+) diff --git a/prog/probe.go b/prog/probe.go index 2cb26bab4..66a9e36eb 100644 --- a/prog/probe.go +++ b/prog/probe.go @@ -95,6 +95,9 @@ func probeMain(flags probeFlags, targets []appclient.Target) { if flags.kubernetesEnabled { checkpointFlags["kubernetes_enabled"] = "true" } + if flags.ecsEnabled { + checkpointFlags["ecs_enabled"] = "true" + } go check(checkpointFlags) handlerRegistry := controls.NewDefaultHandlerRegistry() From 8f2e3e7d9bc04663b51482b122948359043b83ea Mon Sep 17 00:00:00 2001 From: Mike Lang Date: Thu, 24 Nov 2016 14:06:24 -0800 Subject: [PATCH 12/19] merger: Fix a pointer bug that trashed the merge process Turns out that when iterating in go, &loop_var is the same address every time --- app/merger.go | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/app/merger.go b/app/merger.go index ab0f9379f..a6e23f56e 100644 --- a/app/merger.go +++ b/app/merger.go @@ -49,8 +49,8 @@ func (smartMerger) Merge(reports []report.Report) report.Report { return reports[0] } c := make(chan *report.Report, l) - for _, r := range reports { - c <- &r + for i := range reports { + c <- &reports[i] } for ; l > 1; l-- { go func(left, right *report.Report) { From d20381d30b0b2e975d6729d48d099b9325937131 Mon Sep 17 00:00:00 2001 From: Mike Lang Date: Sat, 26 Nov 2016 21:13:30 -0800 Subject: [PATCH 13/19] merger: Pass reports via closure, instead of by reference in args This is an alternate way of solving the same problem as 4007a902a264e5ff2c3be6b269ade515c9c1c145, but in a nicer way. Compared to using pointers, this approach more obviously preserves the original behaviour, and is arguably more readable than the original code. Credit to @rade for this approach. --- app/merger.go | 16 ++++++++-------- 1 file changed, 8 insertions(+), 8 deletions(-) diff --git a/app/merger.go b/app/merger.go index a6e23f56e..94231baf8 100644 --- a/app/merger.go +++ b/app/merger.go @@ -48,15 +48,15 @@ func (smartMerger) Merge(reports []report.Report) report.Report { case 1: return reports[0] } - c := make(chan *report.Report, l) - for i := range reports { - c <- &reports[i] + c := make(chan report.Report, l) + for _, r := range reports { + c <- r } for ; l > 1; l-- { - go func(left, right *report.Report) { - r := left.Merge(*right) - c <- &r - }(<-c, <-c) + left, right := <-c, <-c + go func() { + c <- left.Merge(right) + }() } - return *<-c + return <-c } From 747577f4144415de521f430a04799fa4ecb9c45c Mon Sep 17 00:00:00 2001 From: Alfonso Acosta Date: Mon, 28 Nov 2016 16:32:57 +0000 Subject: [PATCH 14/19] Review feedback --- app/api_topologies.go | 9 ++++----- 1 file changed, 4 insertions(+), 5 deletions(-) diff --git a/app/api_topologies.go b/app/api_topologies.go index 557c8f61c..f8b5cb7b4 100644 --- a/app/api_topologies.go +++ b/app/api_topologies.go @@ -183,11 +183,10 @@ func AddInitialTopologiesToRegistry(registry *Registry) { Rank: 4, }, APITopologyDesc{ - id: weaveID, - parent: hostsID, - renderer: render.WeaveRenderer, - Name: "weave net", - HideIfEmpty: true, + id: weaveID, + parent: hostsID, + renderer: render.WeaveRenderer, + Name: "Weave Net", }, ) } From b06fee8c0fa4da7899ca3ff0434553eb355b1d7b Mon Sep 17 00:00:00 2001 From: Mike Lang Date: Mon, 28 Nov 2016 08:55:10 -0800 Subject: [PATCH 15/19] Review feedback --- probe/awsecs/client.go | 15 +++++------- probe/awsecs/reporter.go | 53 +++++++++++++++++++++++----------------- prog/main.go | 2 +- 3 files changed, 37 insertions(+), 33 deletions(-) diff --git a/probe/awsecs/client.go b/probe/awsecs/client.go index 8d56a4fd3..cac336603 100644 --- a/probe/awsecs/client.go +++ b/probe/awsecs/client.go @@ -33,7 +33,7 @@ func newClient(cluster string) (*ecsClient, error) { // returns a map from deployment ids to service names // cannot fail as it will attempt to deliver partial results, though that may end up being no results func (c ecsClient) getDeploymentMap() map[string]string { - results := make(map[string]string) + results := map[string]string{} lock := sync.Mutex{} // lock mediates access to results group := sync.WaitGroup{} @@ -43,21 +43,20 @@ func (c ecsClient) getDeploymentMap() map[string]string { func(page *ecs.ListServicesOutput, lastPage bool) bool { // describe each page of 10 (the max for one describe command) concurrently group.Add(1) + serviceArns := page.ServiceArns go func() { defer group.Done() resp, err := c.client.DescribeServices(&ecs.DescribeServicesInput{ Cluster: &c.cluster, - Services: page.ServiceArns, + Services: serviceArns, }) if err != nil { - // rather than trying to propagate errors up, just log a warning here log.Warnf("Error describing some ECS services, ECS service report may be incomplete: %v", err) return } for _, failure := range resp.Failures { - // log the failures but still continue with what succeeded log.Warnf("Failed to describe ECS service %s, ECS service report may be incomplete: %s", failure.Arn, failure.Reason) } @@ -75,7 +74,6 @@ func (c ecsClient) getDeploymentMap() map[string]string { group.Wait() if err != nil { - // We want to still return partial results if we have any, so just log a warning log.Warnf("Error listing ECS services, ECS service report may be incomplete: %v", err) } return results @@ -99,11 +97,10 @@ func (c ecsClient) getTaskDeployments(taskArns []string) (map[string]string, err } for _, failure := range resp.Failures { - // log the failures but still continue with what succeeded log.Warnf("Failed to describe ECS task %s, ECS service report may be incomplete: %s", failure.Arn, failure.Reason) } - results := make(map[string]string) + results := make(map[string]string, len(resp.Tasks)) for _, task := range resp.Tasks { results[*task.TaskArn] = *task.StartedBy } @@ -112,7 +109,7 @@ func (c ecsClient) getTaskDeployments(taskArns []string) (map[string]string, err // returns a map from task ARNs to service names func (c ecsClient) getTaskServices(taskArns []string) (map[string]string, error) { - deploymentMapChan := make(chan map[string]string) + deploymentMapChan := chan map[string]string{} go func() { deploymentMapChan <- c.getDeploymentMap() }() @@ -125,7 +122,7 @@ func (c ecsClient) getTaskServices(taskArns []string) (map[string]string, error) return nil, err } - results := make(map[string]string) + results := map[string]string{} for taskArn, depID := range taskDeployments { // Note not all tasks map to a deployment, or we could otherwise mismatch due to races. // It's safe to just ignore all these cases and consider them "non-service" tasks. diff --git a/probe/awsecs/reporter.go b/probe/awsecs/reporter.go index 8752a07b5..8537c4611 100644 --- a/probe/awsecs/reporter.go +++ b/probe/awsecs/reporter.go @@ -18,29 +18,38 @@ type taskInfo struct { // return map from cluster to map of task arns to task infos func getLabelInfo(rpt report.Report) map[string]map[string]*taskInfo { - results := make(map[string]map[string]*taskInfo) + results := map[string]map[string]*taskInfo{} log.Debug("scanning for ECS containers") for nodeID, node := range rpt.Container.Nodes { - taskArn, taskArnOk := node.Latest.Lookup(docker.LabelPrefix + "com.amazonaws.ecs.task-arn") - cluster, clusterOk := node.Latest.Lookup(docker.LabelPrefix + "com.amazonaws.ecs.cluster") - family, familyOk := node.Latest.Lookup(docker.LabelPrefix + "com.amazonaws.ecs.task-definition-family") - - if taskArnOk && clusterOk && familyOk { - taskMap, ok := results[cluster] - if !ok { - taskMap = make(map[string]*taskInfo) - results[cluster] = taskMap - } - - task, ok := taskMap[taskArn] - if !ok { - task = &taskInfo{containerIDs: make([]string, 0), family: family} - taskMap[taskArn] = task - } - - task.containerIDs = append(task.containerIDs, nodeID) + taskArn, ok := node.Latest.Lookup(docker.LabelPrefix + "com.amazonaws.ecs.task-arn") + if !ok { + continue } + + cluster, ok := node.Latest.Lookup(docker.LabelPrefix + "com.amazonaws.ecs.cluster") + if !ok { + continue + } + + family, ok := node.Latest.Lookup(docker.LabelPrefix + "com.amazonaws.ecs.task-definition-family") + if !ok { + continue + } + + taskMap, ok := results[cluster] + if !ok { + taskMap = map[string]*taskInfo{} + results[cluster] = taskMap + } + + task, ok := taskMap[taskArn] + if !ok { + task = &taskInfo{containerIDs: []string{}, family: family} + taskMap[taskArn] = task + } + + task.containerIDs = append(task.containerIDs, nodeID) } log.Debug("Got ECS container info: %v", results) return results @@ -51,7 +60,7 @@ type Reporter struct { } // Tag needed for Tagger -func (r Reporter) Tag(rpt report.Report) (report.Report, error) { +func (Reporter) Tag(rpt report.Report) (report.Report, error) { rpt = rpt.Copy() clusterMap := getLabelInfo(rpt) @@ -75,7 +84,7 @@ func (r Reporter) Tag(rpt report.Report) (report.Report, error) { } // Create all the services first - unique := make(map[string]bool) + unique := map[string]bool{} for _, serviceName := range taskServices { if !unique[serviceName] { serviceID := report.MakeECSServiceNodeID(serviceName) @@ -99,7 +108,6 @@ func (r Reporter) Tag(rpt report.Report) (report.Report, error) { serviceID := report.MakeECSServiceNodeID(serviceName) parentsSets = parentsSets.Add(report.ECSService, report.MakeStringSet(serviceID)) } - for _, containerID := range info.containerIDs { if containerNode, ok := rpt.Container.Nodes[containerID]; ok { rpt.Container.Nodes[containerID] = containerNode.WithParents(parentsSets) @@ -107,7 +115,6 @@ func (r Reporter) Tag(rpt report.Report) (report.Report, error) { log.Warnf("Got task info for non-existent container %v, this shouldn't be able to happen", containerID) } } - } } diff --git a/prog/main.go b/prog/main.go index 2f8999968..0a679970c 100644 --- a/prog/main.go +++ b/prog/main.go @@ -285,7 +285,7 @@ func main() { flag.StringVar(&flags.probe.kubernetesConfig.Username, "probe.kubernetes.username", "", "Username for basic authentication to the API server") // AWS ECS - flag.BoolVar(&flags.probe.ecsEnabled, "probe.ecs", false, "collect ecs-related attributes for containers on this node") + flag.BoolVar(&flags.probe.ecsEnabled, "probe.ecs", false, "Collect ecs-related attributes for containers on this node") // Weave flag.StringVar(&flags.probe.weaveAddr, "probe.weave.addr", "127.0.0.1:6784", "IP address & port of the Weave router") From 9a10e9650d8756953994ce2cfcb236c23d6cd9b5 Mon Sep 17 00:00:00 2001 From: Mike Lang Date: Mon, 28 Nov 2016 09:16:19 -0800 Subject: [PATCH 16/19] Fix the one instance of "make" that is actually apparently required --- probe/awsecs/client.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/probe/awsecs/client.go b/probe/awsecs/client.go index cac336603..d016362f3 100644 --- a/probe/awsecs/client.go +++ b/probe/awsecs/client.go @@ -109,7 +109,7 @@ func (c ecsClient) getTaskDeployments(taskArns []string) (map[string]string, err // returns a map from task ARNs to service names func (c ecsClient) getTaskServices(taskArns []string) (map[string]string, error) { - deploymentMapChan := chan map[string]string{} + deploymentMapChan := make(chan map[string]string) go func() { deploymentMapChan <- c.getDeploymentMap() }() From 003ef6b4eadf8683a453b721ed48843c20dbbb46 Mon Sep 17 00:00:00 2001 From: Mike Lang Date: Mon, 28 Nov 2016 13:32:24 -0800 Subject: [PATCH 17/19] Add some basic metadata to ECS nodes --- probe/awsecs/client.go | 55 ++++++++++++++++++++++------------- probe/awsecs/reporter.go | 63 +++++++++++++++++++++++++++++----------- 2 files changed, 81 insertions(+), 37 deletions(-) diff --git a/probe/awsecs/client.go b/probe/awsecs/client.go index d016362f3..8774544af 100644 --- a/probe/awsecs/client.go +++ b/probe/awsecs/client.go @@ -16,6 +16,12 @@ type ecsClient struct { cluster string } +type ecsInfo struct { + tasks map[string]*ecs.Task + services map[string]*ecs.Service + taskServiceMap map[string]string +} + func newClient(cluster string) (*ecsClient, error) { sess := session.New() @@ -31,9 +37,19 @@ func newClient(cluster string) (*ecsClient, error) { } // returns a map from deployment ids to service names -// cannot fail as it will attempt to deliver partial results, though that may end up being no results -func (c ecsClient) getDeploymentMap() map[string]string { +func (c ecsClient) getDeploymentMap(services map[string]*ecs.Service) map[string]string { results := map[string]string{} + for serviceName, service := range services { + for _, deployment := range service.Deployments { + results[*deployment.Id] = serviceName + } + } + return results +} + +// cannot fail as it will attempt to deliver partial results, though that may end up being no results +func (c ecsClient) getServices() map[string]*ecs.Service { + results := map[string]*ecs.Service{} lock := sync.Mutex{} // lock mediates access to results group := sync.WaitGroup{} @@ -62,9 +78,7 @@ func (c ecsClient) getDeploymentMap() map[string]string { lock.Lock() for _, service := range resp.Services { - for _, deployment := range service.Deployments { - results[*deployment.Id] = *service.ServiceName - } + results[*service.ServiceName] = service } lock.Unlock() }() @@ -79,8 +93,7 @@ func (c ecsClient) getDeploymentMap() map[string]string { return results } -// returns a map from task ARNs to deployment ids -func (c ecsClient) getTaskDeployments(taskArns []string) (map[string]string, error) { +func (c ecsClient) getTasks(taskArns []string) (map[string]*ecs.Task, error) { taskPtrs := make([]*string, len(taskArns)) for i := range taskArns { taskPtrs[i] = &taskArns[i] @@ -100,36 +113,38 @@ func (c ecsClient) getTaskDeployments(taskArns []string) (map[string]string, err log.Warnf("Failed to describe ECS task %s, ECS service report may be incomplete: %s", failure.Arn, failure.Reason) } - results := make(map[string]string, len(resp.Tasks)) + results := make(map[string]*ecs.Task, len(resp.Tasks)) for _, task := range resp.Tasks { - results[*task.TaskArn] = *task.StartedBy + results[*task.TaskArn] = task } return results, nil } // returns a map from task ARNs to service names -func (c ecsClient) getTaskServices(taskArns []string) (map[string]string, error) { - deploymentMapChan := make(chan map[string]string) +func (c ecsClient) getInfo(taskArns []string) (ecsInfo, error) { + servicesChan := make(chan map[string]*ecs.Service) go func() { - deploymentMapChan <- c.getDeploymentMap() + servicesChan <- c.getServices() }() // do these two fetches in parallel - taskDeployments, err := c.getTaskDeployments(taskArns) - deploymentMap := <-deploymentMapChan + tasks, err := c.getTasks(taskArns) + services := <-servicesChan if err != nil { - return nil, err + return ecsInfo{}, err } - results := map[string]string{} - for taskArn, depID := range taskDeployments { + deploymentMap := c.getDeploymentMap(services) + + taskServiceMap := map[string]string{} + for taskArn, task := range tasks { // Note not all tasks map to a deployment, or we could otherwise mismatch due to races. // It's safe to just ignore all these cases and consider them "non-service" tasks. - if service, ok := deploymentMap[depID]; ok { - results[taskArn] = service + if serviceName, ok := deploymentMap[*task.StartedBy]; ok { + taskServiceMap[taskArn] = serviceName } } - return results, nil + return ecsInfo{services: services, tasks: tasks, taskServiceMap: taskServiceMap}, nil } diff --git a/probe/awsecs/reporter.go b/probe/awsecs/reporter.go index 8537c4611..6a4c1fd27 100644 --- a/probe/awsecs/reporter.go +++ b/probe/awsecs/reporter.go @@ -1,6 +1,8 @@ package awsecs import ( + "time" + log "github.com/Sirupsen/logrus" "github.com/weaveworks/scope/probe/docker" "github.com/weaveworks/scope/report" @@ -8,17 +10,35 @@ import ( // TaskFamily is the key that stores the task family of an ECS Task const ( - TaskFamily = "ecs_task_family" + Cluster = "ecs_cluster" + CreatedAt = "ecs_created_at" + TaskFamily = "ecs_task_family" + ServiceDesiredCount = "ecs_service_desired_count" + ServiceRunningCount = "ecs_service_running_count" ) -type taskInfo struct { +var ( + taskMetadata = report.MetadataTemplates{ + Cluster: {ID: Cluster, Label: "Cluster", From: report.FromLatest, Priority: 0}, + CreatedAt: {ID: CreatedAt, Label: "Created At", From: report.FromLatest, Priority: 1, Datatype: "datetime"}, + TaskFamily: {ID: TaskFamily, Label: "Family", From: report.FromLatest, Priority: 2}, + } + serviceMetadata = report.MetadataTemplates{ + Cluster: {ID: Cluster, Label: "Cluster", From: report.FromLatest, Priority: 0}, + CreatedAt: {ID: CreatedAt, Label: "Created At", From: report.FromLatest, Priority: 1, Datatype: "datetime"}, + ServiceDesiredCount: {ID: ServiceDesiredCount, Label: "Desired Task Count", From: report.FromLatest, Priority: 2, Datatype: "number"}, + ServiceRunningCount: {ID: ServiceRunningCount, Label: "Running Task Count", From: report.FromLatest, Priority: 3, Datatype: "number"}, + } +) + +type taskLabelInfo struct { containerIDs []string family string } // return map from cluster to map of task arns to task infos -func getLabelInfo(rpt report.Report) map[string]map[string]*taskInfo { - results := map[string]map[string]*taskInfo{} +func getLabelInfo(rpt report.Report) map[string]map[string]*taskLabelInfo { + results := map[string]map[string]*taskLabelInfo{} log.Debug("scanning for ECS containers") for nodeID, node := range rpt.Container.Nodes { @@ -39,13 +59,13 @@ func getLabelInfo(rpt report.Report) map[string]map[string]*taskInfo { taskMap, ok := results[cluster] if !ok { - taskMap = map[string]*taskInfo{} + taskMap = map[string]*taskLabelInfo{} results[cluster] = taskMap } task, ok := taskMap[taskArn] if !ok { - task = &taskInfo{containerIDs: []string{}, family: family} + task = &taskLabelInfo{containerIDs: []string{}, family: family} taskMap[taskArn] = task } @@ -78,33 +98,42 @@ func (Reporter) Tag(rpt report.Report) (report.Report, error) { taskArns = append(taskArns, taskArn) } - taskServices, err := client.getTaskServices(taskArns) + ecsInfo, err := client.getInfo(taskArns) if err != nil { return rpt, err } // Create all the services first - unique := map[string]bool{} - for _, serviceName := range taskServices { - if !unique[serviceName] { - serviceID := report.MakeECSServiceNodeID(serviceName) - rpt.ECSService = rpt.ECSService.AddNode(report.MakeNode(serviceID)) - unique[serviceName] = true - } + for serviceName, service := range ecsInfo.services { + serviceID := report.MakeECSServiceNodeID(serviceName) + rpt.ECSService = rpt.ECSService.AddNode(report.MakeNodeWith(serviceID, map[string]string{ // TODO add task metadata + Cluster: cluster, + ServiceDesiredCount: string(*service.DesiredCount), + ServiceRunningCount: string(*service.RunningCount), + })) } - log.Debugf("Created %v ECS service nodes", len(taskServices)) + log.Debugf("Created %v ECS service nodes", len(ecsInfo.services)) for taskArn, info := range taskMap { + task, ok := ecsInfo.tasks[taskArn] + if !ok { + // can happen due to partial failures, just skip it + continue + } // new task node taskID := report.MakeECSTaskNodeID(taskArn) - node := report.MakeNodeWith(taskID, map[string]string{TaskFamily: info.family}) + node := report.MakeNodeWith(taskID, map[string]string{ + TaskFamily: info.family, + Cluster: cluster, + CreatedAt: task.CreatedAt.Format(time.RFC3339Nano), + }) rpt.ECSTask = rpt.ECSTask.AddNode(node) // parents sets to merge into all matching container nodes parentsSets := report.MakeSets() parentsSets = parentsSets.Add(report.ECSTask, report.MakeStringSet(taskID)) - if serviceName, ok := taskServices[taskArn]; ok { + if serviceName, ok := ecsInfo.taskServiceMap[taskArn]; ok { serviceID := report.MakeECSServiceNodeID(serviceName) parentsSets = parentsSets.Add(report.ECSService, report.MakeStringSet(serviceID)) } From 9c7282231f78797647f4d8991cbade9f55abbe89 Mon Sep 17 00:00:00 2001 From: Alfonso Acosta Date: Tue, 29 Nov 2016 01:36:35 +0000 Subject: [PATCH 18/19] Fix tests Also, refactor some tests and MakeRegistry in api_topologies --- app/api_topologies.go | 146 +++++++++----------- app/api_topologies_test.go | 31 +++-- probe/appclient/app_client_internal_test.go | 4 + probe/probe_internal_test.go | 2 + 4 files changed, 95 insertions(+), 88 deletions(-) diff --git a/app/api_topologies.go b/app/api_topologies.go index f8b5cb7b4..7f0e205fc 100644 --- a/app/api_topologies.go +++ b/app/api_topologies.go @@ -45,13 +45,76 @@ var ( } ) -func init() { - AddInitialTopologiesToRegistry(topologyRegistry) +// kubernetesFilters generates the current kubernetes filters based on the +// available k8s topologies. +func kubernetesFilters(namespaces ...string) APITopologyOptionGroup { + options := APITopologyOptionGroup{ID: "namespace", Default: "all"} + for _, namespace := range namespaces { + if namespace == "default" { + options.Default = namespace + } + options.Options = append(options.Options, APITopologyOption{ + Value: namespace, Label: namespace, filter: render.IsNamespace(namespace), filterPseudo: false, + }) + } + options.Options = append(options.Options, APITopologyOption{Value: "all", Label: "All Namespaces", filter: nil, filterPseudo: false}) + return options } -// AddInitialTopologiesToRegistry does the initial setup for a Registry. -// This is needed for testing. -func AddInitialTopologiesToRegistry(registry *Registry) { +// updateFilters updates the available filters based on the current report. +// Currently only kubernetes changes. +func updateFilters(rpt report.Report, topologies []APITopologyDesc) []APITopologyDesc { + namespaces := map[string]struct{}{} + for _, t := range []report.Topology{rpt.Pod, rpt.Service, rpt.Deployment, rpt.ReplicaSet} { + for _, n := range t.Nodes { + if state, ok := n.Latest.Lookup(kubernetes.State); ok && state == kubernetes.StateDeleted { + continue + } + if namespace, ok := n.Latest.Lookup(kubernetes.Namespace); ok { + namespaces[namespace] = struct{}{} + } + } + } + var ns []string + for namespace := range namespaces { + ns = append(ns, namespace) + } + sort.Strings(ns) + for i, t := range topologies { + if t.id == podsID || t.id == servicesID || t.id == deploymentsID || t.id == replicaSetsID { + topologies[i] = updateTopologyFilters(t, []APITopologyOptionGroup{ + kubernetesFilters(ns...), unmanagedFilter, + }) + } + } + return topologies +} + +// updateTopologyFilters recursively sets the options on a topology description +func updateTopologyFilters(t APITopologyDesc, options []APITopologyOptionGroup) APITopologyDesc { + t.Options = options + for i, sub := range t.SubTopologies { + t.SubTopologies[i] = updateTopologyFilters(sub, options) + } + return t +} + +// MakeAPITopologyOption provides an external interface to the package for creating an APITopologyOption. +func MakeAPITopologyOption(value string, label string, filterFunc render.FilterFunc, pseudo bool) APITopologyOption { + return APITopologyOption{Value: value, Label: label, filter: filterFunc, filterPseudo: pseudo} +} + +// Registry is a threadsafe store of the available topologies +type Registry struct { + sync.RWMutex + items map[string]APITopologyDesc +} + +// MakeRegistry returns a new Registry +func MakeRegistry() *Registry { + registry := &Registry{ + items: map[string]APITopologyDesc{}, + } containerFilters := []APITopologyOptionGroup{ { ID: containerLabelFiltersGroupID, @@ -189,79 +252,8 @@ func AddInitialTopologiesToRegistry(registry *Registry) { Name: "Weave Net", }, ) -} -// kubernetesFilters generates the current kubernetes filters based on the -// available k8s topologies. -func kubernetesFilters(namespaces ...string) APITopologyOptionGroup { - options := APITopologyOptionGroup{ID: "namespace", Default: "all"} - for _, namespace := range namespaces { - if namespace == "default" { - options.Default = namespace - } - options.Options = append(options.Options, APITopologyOption{ - Value: namespace, Label: namespace, filter: render.IsNamespace(namespace), filterPseudo: false, - }) - } - options.Options = append(options.Options, APITopologyOption{Value: "all", Label: "All Namespaces", filter: nil, filterPseudo: false}) - return options -} - -// updateFilters updates the available filters based on the current report. -// Currently only kubernetes changes. -func updateFilters(rpt report.Report, topologies []APITopologyDesc) []APITopologyDesc { - namespaces := map[string]struct{}{} - for _, t := range []report.Topology{rpt.Pod, rpt.Service, rpt.Deployment, rpt.ReplicaSet} { - for _, n := range t.Nodes { - if state, ok := n.Latest.Lookup(kubernetes.State); ok && state == kubernetes.StateDeleted { - continue - } - if namespace, ok := n.Latest.Lookup(kubernetes.Namespace); ok { - namespaces[namespace] = struct{}{} - } - } - } - var ns []string - for namespace := range namespaces { - ns = append(ns, namespace) - } - sort.Strings(ns) - for i, t := range topologies { - if t.id == podsID || t.id == servicesID || t.id == deploymentsID || t.id == replicaSetsID { - topologies[i] = updateTopologyFilters(t, []APITopologyOptionGroup{ - kubernetesFilters(ns...), unmanagedFilter, - }) - } - } - return topologies -} - -// updateTopologyFilters recursively sets the options on a topology description -func updateTopologyFilters(t APITopologyDesc, options []APITopologyOptionGroup) APITopologyDesc { - t.Options = options - for i, sub := range t.SubTopologies { - t.SubTopologies[i] = updateTopologyFilters(sub, options) - } - return t -} - -// MakeAPITopologyOption provides an external interface to the package for creating an APITopologyOption. -func MakeAPITopologyOption(value string, label string, filterFunc render.FilterFunc, pseudo bool) APITopologyOption { - return APITopologyOption{Value: value, Label: label, filter: filterFunc, filterPseudo: pseudo} -} - -// Registry is a threadsafe store of the available topologies -type Registry struct { - sync.RWMutex - items map[string]APITopologyDesc -} - -// MakeRegistry returns a new Registry -func MakeRegistry() *Registry { - newRegistry := &Registry{ - items: map[string]APITopologyDesc{}, - } - return newRegistry + return registry } // APITopologyDesc is returned in a list by the /api/topology handler. diff --git a/app/api_topologies_test.go b/app/api_topologies_test.go index ea23b7945..63e214662 100644 --- a/app/api_topologies_test.go +++ b/app/api_topologies_test.go @@ -20,7 +20,7 @@ import ( ) const ( - containerLabelFiltersGroupID = "container_label_filters_group" + containerLabelFiltersGroupID = "container-label-filters-group" customAPITopologyOptionFilterID = "containerLabelFilter0" ) @@ -35,7 +35,7 @@ func TestAPITopology(t *testing.T) { if err := decoder.Decode(&topologies); err != nil { t.Fatalf("JSON parse error: %s", err) } - equals(t, 4, len(topologies)) + equals(t, 5, len(topologies)) for _, topology := range topologies { is200(t, ts, topology.URL) @@ -44,6 +44,11 @@ func TestAPITopology(t *testing.T) { is200(t, ts, subTopology.URL) } + // TODO: add ECS nodes in report fixture + if topology.Name == "Tasks" { + continue + } + if have := topology.Stats.EdgeCount; have <= 0 { t.Errorf("EdgeCount isn't positive for %s: %d", topology.Name, have) } @@ -79,8 +84,9 @@ func TestContainerLabelFilterExclude(t *testing.T) { // all containers but the excluded container should be present for key := range topologySummaries { - if report.MakeContainerNodeID(fixture.ServerContainerNodeID) == key { - t.Errorf("TestAPITopologyNegativeContainerLabelFilter Failed. Expected to not find " + report.MakeContainerNodeID(fixture.ServerContainerNodeID) + " in report") + id := report.MakeContainerNodeID(fixture.ServerContainerNodeID) + if id == key { + t.Errorf("Didn't expect to find %q in report", id) } } } @@ -89,14 +95,17 @@ func getTestContainerLabelFilterTopologySummary(t *testing.T, exclude bool) (det ts := topologyServer() defer ts.Close() - topologyRegistry := app.MakeRegistry() - app.AddInitialTopologiesToRegistry(topologyRegistry) - + var ( + topologyRegistry = app.MakeRegistry() + filter render.FilterFunc + ) if exclude == true { - topologyRegistry.AddContainerFilters(app.MakeAPITopologyOption(customAPITopologyOptionFilterID, "title", render.DoesNotHaveLabel(fixture.TestLabelKey2, fixture.ApplicationLabelValue2), false)) + filter = render.DoesNotHaveLabel(fixture.TestLabelKey2, fixture.ApplicationLabelValue2) } else { - topologyRegistry.AddContainerFilters(app.MakeAPITopologyOption(customAPITopologyOptionFilterID, "title", render.HasLabel(fixture.TestLabelKey1, fixture.ApplicationLabelValue1), false)) + filter = render.HasLabel(fixture.TestLabelKey1, fixture.ApplicationLabelValue1) } + option := app.MakeAPITopologyOption(customAPITopologyOptionFilterID, "title", filter, false) + topologyRegistry.AddContainerFilters(option) urlvalues := url.Values{} urlvalues.Set(containerLabelFiltersGroupID, customAPITopologyOptionFilterID) @@ -123,7 +132,7 @@ func TestAPITopologyAddsKubernetes(t *testing.T) { if err := decoder.Decode(&topologies); err != nil { t.Fatalf("JSON parse error: %s", err) } - equals(t, 4, len(topologies)) + equals(t, 5, len(topologies)) // Enable the kubernetes topologies rpt := report.MakeReport() @@ -157,7 +166,7 @@ func TestAPITopologyAddsKubernetes(t *testing.T) { if err := decoder.Decode(&topologies); err != nil { t.Fatalf("JSON parse error: %s", err) } - equals(t, 4, len(topologies)) + equals(t, 5, len(topologies)) found := false for _, topology := range topologies { diff --git a/probe/appclient/app_client_internal_test.go b/probe/appclient/app_client_internal_test.go index e29a5bc75..79df715e2 100644 --- a/probe/appclient/app_client_internal_test.go +++ b/probe/appclient/app_client_internal_test.go @@ -83,6 +83,8 @@ func TestAppClientPublish(t *testing.T) { rpt.ReplicaSet = report.MakeTopology() rpt.Host = report.MakeTopology() rpt.Overlay = report.MakeTopology() + rpt.ECSTask = report.MakeTopology() + rpt.ECSService = report.MakeTopology() rpt.Endpoint.Controls = nil rpt.Process.Controls = nil rpt.Container.Controls = nil @@ -93,6 +95,8 @@ func TestAppClientPublish(t *testing.T) { rpt.ReplicaSet.Controls = nil rpt.Host.Controls = nil rpt.Overlay.Controls = nil + rpt.ECSTask.Controls = nil + rpt.ECSService.Controls = nil s := dummyServer(t, token, id, version, rpt, done) defer s.Close() diff --git a/probe/probe_internal_test.go b/probe/probe_internal_test.go index 55cd52e68..0202ad36c 100644 --- a/probe/probe_internal_test.go +++ b/probe/probe_internal_test.go @@ -91,6 +91,8 @@ func TestProbe(t *testing.T) { want.ReplicaSet.Controls = nil want.Host.Controls = nil want.Overlay.Controls = nil + want.ECSTask.Controls = nil + want.ECSService.Controls = nil want.Endpoint.AddNode(node) pub := mockPublisher{make(chan report.Report, 10)} From d0caee47486e5be08d9f5fac88704f529ac759aa Mon Sep 17 00:00:00 2001 From: Mike Lang Date: Tue, 29 Nov 2016 07:04:59 -0800 Subject: [PATCH 19/19] Add some basic metadata to the ECS task/service details panels --- probe/awsecs/reporter.go | 25 ++++++++++++++++++------- prog/probe.go | 4 +++- 2 files changed, 21 insertions(+), 8 deletions(-) diff --git a/probe/awsecs/reporter.go b/probe/awsecs/reporter.go index 6a4c1fd27..9a6e2321c 100644 --- a/probe/awsecs/reporter.go +++ b/probe/awsecs/reporter.go @@ -1,6 +1,7 @@ package awsecs import ( + "fmt" "time" log "github.com/Sirupsen/logrus" @@ -26,8 +27,8 @@ var ( serviceMetadata = report.MetadataTemplates{ Cluster: {ID: Cluster, Label: "Cluster", From: report.FromLatest, Priority: 0}, CreatedAt: {ID: CreatedAt, Label: "Created At", From: report.FromLatest, Priority: 1, Datatype: "datetime"}, - ServiceDesiredCount: {ID: ServiceDesiredCount, Label: "Desired Task Count", From: report.FromLatest, Priority: 2, Datatype: "number"}, - ServiceRunningCount: {ID: ServiceRunningCount, Label: "Running Task Count", From: report.FromLatest, Priority: 3, Datatype: "number"}, + ServiceDesiredCount: {ID: ServiceDesiredCount, Label: "Desired Tasks", From: report.FromLatest, Priority: 2, Datatype: "number"}, + ServiceRunningCount: {ID: ServiceRunningCount, Label: "Running Tasks", From: report.FromLatest, Priority: 3, Datatype: "number"}, } ) @@ -75,7 +76,7 @@ func getLabelInfo(rpt report.Report) map[string]map[string]*taskLabelInfo { return results } -// Reporter implements Tagger +// Reporter implements Tagger, Reporter type Reporter struct { } @@ -106,10 +107,10 @@ func (Reporter) Tag(rpt report.Report) (report.Report, error) { // Create all the services first for serviceName, service := range ecsInfo.services { serviceID := report.MakeECSServiceNodeID(serviceName) - rpt.ECSService = rpt.ECSService.AddNode(report.MakeNodeWith(serviceID, map[string]string{ // TODO add task metadata + rpt.ECSService = rpt.ECSService.AddNode(report.MakeNodeWith(serviceID, map[string]string{ Cluster: cluster, - ServiceDesiredCount: string(*service.DesiredCount), - ServiceRunningCount: string(*service.RunningCount), + ServiceDesiredCount: fmt.Sprintf("%d", *service.DesiredCount), + ServiceRunningCount: fmt.Sprintf("%d", *service.RunningCount), })) } log.Debugf("Created %v ECS service nodes", len(ecsInfo.services)) @@ -151,7 +152,17 @@ func (Reporter) Tag(rpt report.Report) (report.Report, error) { return rpt, nil } -// Name needed for Tagger +// Report needed for Reporter +func (Reporter) Report() (report.Report, error) { + result := report.MakeReport() + taskTopology := report.MakeTopology().WithMetadataTemplates(taskMetadata) + result.ECSTask = result.ECSTask.Merge(taskTopology) + serviceTopology := report.MakeTopology().WithMetadataTemplates(serviceMetadata) + result.ECSService = result.ECSService.Merge(serviceTopology) + return result, nil +} + +// Name needed for Tagger, Reporter func (r Reporter) Name() string { return "awsecs" } diff --git a/prog/probe.go b/prog/probe.go index 66a9e36eb..dd49881c7 100644 --- a/prog/probe.go +++ b/prog/probe.go @@ -206,7 +206,9 @@ func probeMain(flags probeFlags, targets []appclient.Target) { } if flags.ecsEnabled { - p.AddTagger(awsecs.Reporter{}) + reporter := awsecs.Reporter{} + p.AddReporter(reporter) + p.AddTagger(reporter) } if flags.weaveEnabled {