mirror of
https://github.com/weaveworks/scope.git
synced 2026-07-22 06:46:50 +00:00
169 lines
5.1 KiB
Go
169 lines
5.1 KiB
Go
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"
|
|
}
|