Merge pull request #2026 from weaveworks/mike/ecs/initial

Add ECS views
This commit is contained in:
Alfonso Acosta
2016-11-29 07:42:59 -08:00
committed by GitHub
15 changed files with 667 additions and 160 deletions
+159 -146
View File
@@ -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...)
+20 -11
View File
@@ -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 {
+3 -2
View File
@@ -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
}
@@ -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()
+150
View File
@@ -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
}
+168
View File
@@ -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"
}
+2
View File
@@ -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)}
+2
View File
@@ -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),
} {
+5
View File
@@ -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")
+10
View File
@@ -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)
+14
View File
@@ -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)
+86
View File
@@ -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
}
+2
View File
@@ -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)
)
+12
View File
@@ -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
+30 -1
View File
@@ -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
}