diff --git a/app/api_topologies.go b/app/api_topologies.go index e9bcb2a67..7f0e205fc 100644 --- a/app/api_topologies.go +++ b/app/api_topologies.go @@ -16,23 +16,26 @@ 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 ( topologyRegistry = MakeRegistry() - k8sPseudoFilter = APITopologyOptionGroup{ + unmanagedFilter = APITopologyOptionGroup{ ID: "pseudo", Default: "hide", Options: []APITopologyOption{ @@ -42,134 +45,6 @@ var ( } ) -func init() { - AddInitialTopologiesToRegistry(topologyRegistry) -} - -// AddInitialTopologiesToRegistry does the initial setup for a Registry. -// This is needed for testing. -func AddInitialTopologiesToRegistry(registry *Registry) { - containerFilters := []APITopologyOptionGroup{ - { - 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}}, - }, - { - ID: "stopped", - Default: "running", - Options: []APITopologyOption{ - {Value: "stopped", Label: "Stopped containers", filter: render.IsStopped, filterPseudo: false}, - {Value: "running", Label: "Running containers", filter: render.IsRunning, filterPseudo: false}, - {Value: "both", Label: "Both", filter: nil, filterPseudo: false}, - }, - }, - { - ID: "pseudo", - Default: "hide", - Options: []APITopologyOption{ - {Value: "show", Label: "Show Uncontained", filter: nil, filterPseudo: false}, - {Value: "hide", Label: "Hide Uncontained", filter: render.IsNotPseudo, filterPseudo: true}, - }, - }, - } - - unconnectedFilter := []APITopologyOptionGroup{ - { - ID: "unconnected", - Default: "hide", - Options: []APITopologyOption{ - // Show the user why there are filtered nodes in this view. - // Don't give them the option to show those nodes. - {Value: "hide", Label: "Unconnected nodes hidden", filter: nil, filterPseudo: false}, - }, - }, - } - - // Topology option labels should tell the current state. The first item must - // be the verb to get to that state - registry.Add( - APITopologyDesc{ - id: processesTopologyDescID, - renderer: render.FilterUnconnected(render.ProcessWithContainerNameRenderer), - Name: "Processes", - Rank: 1, - Options: unconnectedFilter, - HideIfEmpty: true, - }, - APITopologyDesc{ - id: processesByNameTopologyDescID, - parent: "processes", - renderer: render.FilterUnconnected(render.ProcessNameRenderer), - Name: "by name", - Options: unconnectedFilter, - HideIfEmpty: true, - }, - APITopologyDesc{ - id: containersTopologyDescID, - renderer: render.ContainerWithImageNameRenderer, - Name: "Containers", - Rank: 2, - Options: containerFilters, - }, - APITopologyDesc{ - id: containersByHostnameTopologyDescID, - parent: "containers", - renderer: render.ContainerHostnameRenderer, - Name: "by DNS name", - Options: containerFilters, - }, - APITopologyDesc{ - id: containersByImageTopologyDescID, - parent: "containers", - renderer: render.ContainerImageRenderer, - Name: "by image", - Options: containerFilters, - }, - APITopologyDesc{ - id: podsTopologyDescID, - renderer: render.PodRenderer, - Name: "Pods", - Rank: 3, - HideIfEmpty: true, - }, - APITopologyDesc{ - id: replicaSetsTopologyDescID, - parent: "pods", - renderer: render.ReplicaSetRenderer, - Name: "replica sets", - HideIfEmpty: true, - }, - APITopologyDesc{ - id: deploymentsTopologyDescID, - parent: "pods", - renderer: render.DeploymentRenderer, - Name: "deployments", - HideIfEmpty: true, - }, - APITopologyDesc{ - id: servicesTopologyDescID, - parent: "pods", - renderer: render.PodServiceRenderer, - Name: "services", - HideIfEmpty: true, - }, - APITopologyDesc{ - id: hostsTopologyDescID, - renderer: render.HostRenderer, - Name: "Hosts", - Rank: 4, - }, - APITopologyDesc{ - id: "weave", - parent: "hosts", - renderer: render.WeaveRenderer, - Name: "weave net", - HideIfEmpty: true, - }, - ) -} - // kubernetesFilters generates the current kubernetes filters based on the // available k8s topologies. func kubernetesFilters(namespaces ...string) APITopologyOptionGroup { @@ -206,9 +81,9 @@ 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, + kubernetesFilters(ns...), unmanagedFilter, }) } } @@ -237,10 +112,148 @@ type Registry struct { // MakeRegistry returns a new Registry func MakeRegistry() *Registry { - newRegistry := &Registry{ + registry := &Registry{ items: map[string]APITopologyDesc{}, } - return newRegistry + containerFilters := []APITopologyOptionGroup{ + { + 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}}, + }, + { + ID: "stopped", + Default: "running", + Options: []APITopologyOption{ + {Value: "stopped", Label: "Stopped containers", filter: render.IsStopped, filterPseudo: false}, + {Value: "running", Label: "Running containers", filter: render.IsRunning, filterPseudo: false}, + {Value: "both", Label: "Both", filter: nil, filterPseudo: false}, + }, + }, + { + ID: "pseudo", + Default: "hide", + Options: []APITopologyOption{ + {Value: "show", Label: "Show Uncontained", filter: nil, filterPseudo: false}, + {Value: "hide", Label: "Hide Uncontained", filter: render.IsNotPseudo, filterPseudo: true}, + }, + }, + } + + unconnectedFilter := []APITopologyOptionGroup{ + { + ID: "unconnected", + Default: "hide", + Options: []APITopologyOption{ + // Show the user why there are filtered nodes in this view. + // Don't give them the option to show those nodes. + {Value: "hide", Label: "Unconnected nodes hidden", filter: nil, filterPseudo: false}, + }, + }, + } + + // Topology option labels should tell the current state. The first item must + // be the verb to get to that state + registry.Add( + APITopologyDesc{ + id: processesID, + renderer: render.FilterUnconnected(render.ProcessWithContainerNameRenderer), + Name: "Processes", + Rank: 1, + Options: unconnectedFilter, + HideIfEmpty: true, + }, + APITopologyDesc{ + id: processesByNameID, + parent: processesID, + renderer: render.FilterUnconnected(render.ProcessNameRenderer), + Name: "by name", + Options: unconnectedFilter, + HideIfEmpty: true, + }, + APITopologyDesc{ + id: containersID, + renderer: render.ContainerWithImageNameRenderer, + Name: "Containers", + Rank: 2, + Options: containerFilters, + }, + APITopologyDesc{ + id: containersByHostnameID, + parent: containersID, + renderer: render.ContainerHostnameRenderer, + Name: "by DNS name", + Options: containerFilters, + }, + APITopologyDesc{ + id: containersByImageID, + parent: containersID, + renderer: render.ContainerImageRenderer, + Name: "by image", + Options: containerFilters, + }, + APITopologyDesc{ + id: podsID, + renderer: render.PodRenderer, + Name: "Pods", + Rank: 3, + HideIfEmpty: true, + }, + APITopologyDesc{ + id: replicaSetsID, + parent: podsID, + renderer: render.ReplicaSetRenderer, + Name: "replica sets", + HideIfEmpty: true, + }, + APITopologyDesc{ + id: deploymentsID, + parent: podsID, + renderer: render.DeploymentRenderer, + Name: "deployments", + HideIfEmpty: true, + }, + APITopologyDesc{ + id: servicesID, + parent: podsID, + renderer: render.PodServiceRenderer, + Name: "services", + HideIfEmpty: true, + }, + APITopologyDesc{ + id: ecsTasksID, + renderer: render.ECSTaskRenderer, + Name: "Tasks", + Rank: 3, + Options: []APITopologyOptionGroup{unmanagedFilter}, + HideIfEmpty: true, + }, + APITopologyDesc{ + id: ecsServicesID, + parent: ecsTasksID, + renderer: render.ECSServiceRenderer, + Name: "services", + Options: []APITopologyOptionGroup{unmanagedFilter}, + HideIfEmpty: true, + }, + APITopologyDesc{ + id: hostsID, + renderer: render.HostRenderer, + Name: "Hosts", + Rank: 4, + }, + APITopologyDesc{ + id: weaveID, + parent: hostsID, + renderer: render.WeaveRenderer, + Name: "Weave Net", + }, + ) + + return registry } // APITopologyDesc is returned in a list by the /api/topology handler. @@ -297,7 +310,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/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/app/merger.go b/app/merger.go index be4d149b2..94231baf8 100644 --- a/app/merger.go +++ b/app/merger.go @@ -53,9 +53,10 @@ func (smartMerger) Merge(reports []report.Report) report.Report { c <- r } for ; l > 1; l-- { - go func(left, right report.Report) { + left, right := <-c, <-c + go func() { c <- left.Merge(right) - }(<-c, <-c) + }() } return <-c } 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/awsecs/client.go b/probe/awsecs/client.go new file mode 100644 index 000000000..8774544af --- /dev/null +++ b/probe/awsecs/client.go @@ -0,0 +1,150 @@ +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 +} + +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() + + 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 +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{} + + 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) + serviceArns := page.ServiceArns + go func() { + defer group.Done() + + resp, err := c.client.DescribeServices(&ecs.DescribeServicesInput{ + Cluster: &c.cluster, + Services: serviceArns, + }) + if err != nil { + log.Warnf("Error describing some ECS services, ECS service report may be incomplete: %v", err) + return + } + + for _, failure := range resp.Failures { + 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 { + results[*service.ServiceName] = service + } + lock.Unlock() + }() + return true + }, + ) + group.Wait() + + if err != nil { + log.Warnf("Error listing ECS services, ECS service report may be incomplete: %v", err) + } + return results +} + +func (c ecsClient) getTasks(taskArns []string) (map[string]*ecs.Task, 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.Warnf("Failed to describe ECS task %s, ECS service report may be incomplete: %s", failure.Arn, failure.Reason) + } + + results := make(map[string]*ecs.Task, len(resp.Tasks)) + for _, task := range resp.Tasks { + results[*task.TaskArn] = task + } + return results, nil +} + +// returns a map from task ARNs to service names +func (c ecsClient) getInfo(taskArns []string) (ecsInfo, error) { + servicesChan := make(chan map[string]*ecs.Service) + go func() { + servicesChan <- c.getServices() + }() + + // do these two fetches in parallel + tasks, err := c.getTasks(taskArns) + services := <-servicesChan + + if err != nil { + return ecsInfo{}, err + } + + 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 serviceName, ok := deploymentMap[*task.StartedBy]; ok { + taskServiceMap[taskArn] = serviceName + } + } + + return ecsInfo{services: services, tasks: tasks, taskServiceMap: taskServiceMap}, nil +} diff --git a/probe/awsecs/reporter.go b/probe/awsecs/reporter.go new file mode 100644 index 000000000..9a6e2321c --- /dev/null +++ b/probe/awsecs/reporter.go @@ -0,0 +1,168 @@ +package awsecs + +import ( + "fmt" + "time" + + log "github.com/Sirupsen/logrus" + "github.com/weaveworks/scope/probe/docker" + "github.com/weaveworks/scope/report" +) + +// TaskFamily is the key that stores the task family of an ECS Task +const ( + Cluster = "ecs_cluster" + CreatedAt = "ecs_created_at" + TaskFamily = "ecs_task_family" + ServiceDesiredCount = "ecs_service_desired_count" + ServiceRunningCount = "ecs_service_running_count" +) + +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 Tasks", From: report.FromLatest, Priority: 2, Datatype: "number"}, + ServiceRunningCount: {ID: ServiceRunningCount, Label: "Running Tasks", 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]*taskLabelInfo { + results := map[string]map[string]*taskLabelInfo{} + log.Debug("scanning for ECS containers") + for nodeID, node := range rpt.Container.Nodes { + + 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]*taskLabelInfo{} + results[cluster] = taskMap + } + + task, ok := taskMap[taskArn] + if !ok { + task = &taskLabelInfo{containerIDs: []string{}, family: family} + taskMap[taskArn] = task + } + + task.containerIDs = append(task.containerIDs, nodeID) + } + log.Debug("Got ECS container info: %v", results) + return results +} + +// Reporter implements Tagger, Reporter +type Reporter struct { +} + +// Tag needed for Tagger +func (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) + } + + ecsInfo, err := client.getInfo(taskArns) + if err != nil { + return rpt, err + } + + // 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{ + Cluster: cluster, + ServiceDesiredCount: fmt.Sprintf("%d", *service.DesiredCount), + ServiceRunningCount: fmt.Sprintf("%d", *service.RunningCount), + })) + } + 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, + 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 := ecsInfo.taskServiceMap[taskArn]; ok { + 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) + } else { + log.Warnf("Got task info for non-existent container %v, this shouldn't be able to happen", containerID) + } + } + } + + } + + return rpt, nil +} + +// 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/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)} 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), } { diff --git a/prog/main.go b/prog/main.go index 956549cea..0a679970c 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..dd49881c7 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" @@ -94,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() @@ -201,6 +205,12 @@ func probeMain(flags probeFlags, targets []appclient.Target) { } } + if flags.ecsEnabled { + reporter := awsecs.Reporter{} + p.AddReporter(reporter) + p.AddTagger(reporter) + } + if flags.weaveEnabled { client := weave.NewClient(sanitize.URL("http://", 6784, "")(flags.weaveAddr)) weave, err := overlay.NewWeave(hostID, client, dockerEndpoint) 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..10ba0e75a --- /dev/null +++ b/render/ecs.go @@ -0,0 +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 = ConditionalRenderer(renderECSTopologies, + ApplyDecorators( + MakeMap( + PropagateSingleMetrics(report.Container), + MakeReduce( + MakeMap( + MapContainer2ECSTask, + ContainerWithImageNameRenderer, + ), + SelectECSTask, + ), + ), + ), +) + +// ECSServiceRenderer is a Renderer for Amazon ECS services. +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 +} 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 c72aeb832..890ab5be1 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. @@ -147,7 +158,17 @@ func MakeReport() Report { WithShape(Heptagon). WithLabel("replica set", "replica sets"), - Overlay: MakeTopology(), + Overlay: MakeTopology(). + WithShape(Circle). + WithLabel("peer", "peers"), + + ECSTask: MakeTopology(). + WithShape(Heptagon). + WithLabel("task", "tasks"), + + ECSService: MakeTopology(). + WithShape(Heptagon). + WithLabel("service", "services"), Sampling: Sampling{}, Window: 0, @@ -169,6 +190,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 +213,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 +244,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 +261,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 }