Repository navigation
feat(aws): add elasticache endpoints #786
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
8175d5e
7e4e49f
519d190
0ba6ab4
e8fa82a
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,226 @@ | ||
| package aws | ||
|
|
||
| import ( | ||
| "context" | ||
| "fmt" | ||
| "strings" | ||
| "sync" | ||
| "time" | ||
|
|
||
| "github.com/aws/aws-sdk-go/aws" | ||
| "github.com/aws/aws-sdk-go/aws/credentials/stscreds" | ||
| "github.com/aws/aws-sdk-go/aws/session" | ||
| "github.com/aws/aws-sdk-go/service/ec2" | ||
| "github.com/aws/aws-sdk-go/service/elasticache" | ||
| "github.com/pkg/errors" | ||
| "github.com/projectdiscovery/cloudlist/pkg/schema" | ||
| "github.com/projectdiscovery/gologger" | ||
| ) | ||
|
|
||
| // elastiCacheProvider is a provider for AWS ElastiCache API. | ||
| type elastiCacheProvider struct { | ||
| options ProviderOptions | ||
| elastiCacheClient *elasticache.ElastiCache | ||
| session *session.Session | ||
| regions *ec2.DescribeRegionsOutput | ||
| } | ||
|
|
||
| func (ep *elastiCacheProvider) name() string { | ||
| return "elasticache" | ||
| } | ||
|
|
||
| // GetResource returns all the resources in the store for a provider. | ||
| func (ep *elastiCacheProvider) GetResource(ctx context.Context) (*schema.Resources, error) { | ||
| list := schema.NewResources() | ||
| var wg sync.WaitGroup | ||
| var mu sync.Mutex | ||
| var errs []error | ||
|
|
||
| for _, region := range ep.regions.Regions { | ||
| for _, client := range ep.getElastiCacheClients(region.RegionName) { | ||
| wg.Add(1) | ||
|
|
||
| go func(client *elasticache.ElastiCache) { | ||
| defer wg.Done() | ||
| defer func() { | ||
| if r := recover(); r != nil { | ||
| mu.Lock() | ||
| errs = append(errs, fmt.Errorf("panic in elasticache provider: %v", r)) | ||
| mu.Unlock() | ||
| } | ||
| }() | ||
|
|
||
| resources, err := ep.listElastiCacheResources(client) | ||
| mu.Lock() | ||
| defer mu.Unlock() | ||
| if resources != nil { | ||
| list.Merge(resources) | ||
| } | ||
| if err != nil { | ||
| errs = append(errs, err) | ||
| } | ||
| }(client) | ||
| } | ||
| } | ||
| wg.Wait() | ||
| if len(errs) > 0 && len(list.Items) == 0 { | ||
| return nil, fmt.Errorf("elasticache: all workers failed: %v", errs) | ||
| } | ||
| if len(errs) > 0 { | ||
| gologger.Warning().Msgf("elasticache: some listings failed: %v", errs) | ||
| } | ||
| return list, nil | ||
| } | ||
|
|
||
| func (ep *elastiCacheProvider) listElastiCacheResources(client *elasticache.ElastiCache) (*schema.Resources, error) { | ||
| list := schema.NewResources() | ||
| appendEndpoint := func(endpoint *elasticache.Endpoint, metadata map[string]string) { | ||
| if endpoint == nil || aws.StringValue(endpoint.Address) == "" { | ||
| return | ||
| } | ||
| list.Append(&schema.Resource{ | ||
| ID: ep.options.Id, | ||
| Provider: providerName, | ||
| DNSName: aws.StringValue(endpoint.Address), | ||
| Public: true, | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. 🎯 Functional Correctness | 🟡 Minor | ⚡ Quick win 🔎 Supported by static analysis🏁 Script executed: #!/bin/bash
rg -n -C2 'Public:' pkg/providers/aws/rds.go pkg/providers/aws/elb.goRepository: projectdiscovery/cloudlist Length of output: 1339 🏁 Script executed: set -e
printf '%s\n' '--- elasticache.go ---'
nl -ba pkg/providers/aws/elasticache.go | sed -n '1,180p'
printf '%s\n' '--- Public field declarations/usages ---'
rg -n -C3 -F -- 'Public bool' .
rg -n -C2 -F -- '.Public' pkg | head -200
printf '%s\n' '--- schema resource definitions ---'
rg -n -C3 -F -- 'type Resource' pkg internal . 2>/dev/null | head -240Repository: projectdiscovery/cloudlist Length of output: 24457 🏁 Script executed: set -e
printf '%s\n' '--- RDS endpoint construction and metadata ---'
nl -ba pkg/providers/aws/rds.go | sed -n '70,140p;150,195p'
printf '%s\n' '--- Resources.Append and Resource definition ---'
nl -ba pkg/schema/schema.go | sed -n '60,120p;170,205p'Repository: projectdiscovery/cloudlist Length of output: 9337 🌐 Web query:
💡 Result: 🏁 Script executed: set -e
printf '%s\n' '--- dependency version ---'
rg -n -F -- 'github.com/aws/aws-sdk-go' go.mod go.sum
printf '%s\n' '--- repository Public semantics in docs and filters ---'
rg -n -i -C3 -- 'public facing|publicly accessible|public endpoint|private endpoint|Public:' README.md DESIGN.md docs pkg cmd .github 2>/dev/null | head -300
printf '%s\n' '--- serverless endpoint-related fields in repository ---'
rg -n -i -C2 -- 'ServerlessCache|NetworkType|public endpoint|publicly accessible' . --glob '!go.sum' --glob '!vendor/**' | head -300Repository: projectdiscovery/cloudlist Length of output: 20059 🌐 Web query:
💡 Result: Classify ElastiCache endpoints by connection type.
Mark node-based endpoints as private. For serverless caches, set The RDS implementation is not an aligned precedent because it also sets 🤖 Prompt for AI Agents |
||
| Service: ep.name(), | ||
| Metadata: metadata, | ||
| }) | ||
| } | ||
|
|
||
| err := client.DescribeReplicationGroupsPages(&elasticache.DescribeReplicationGroupsInput{}, func(page *elasticache.DescribeReplicationGroupsOutput, _ bool) bool { | ||
| for _, group := range page.ReplicationGroups { | ||
| var metadata map[string]string | ||
| if ep.options.ExtendedMetadata { | ||
| metadata = getReplicationGroupMetadata(group) | ||
| } | ||
| appendEndpoint(group.ConfigurationEndpoint, metadata) | ||
| for _, nodeGroup := range group.NodeGroups { | ||
| appendEndpoint(nodeGroup.PrimaryEndpoint, metadata) | ||
| appendEndpoint(nodeGroup.ReaderEndpoint, metadata) | ||
| for _, member := range nodeGroup.NodeGroupMembers { | ||
| appendEndpoint(member.ReadEndpoint, metadata) | ||
| } | ||
| } | ||
| } | ||
| return true | ||
| }) | ||
| if err != nil { | ||
| return list, errors.Wrap(err, "could not describe elasticache replication groups") | ||
| } | ||
|
coderabbitai[bot] marked this conversation as resolved.
|
||
|
|
||
| // Node endpoints are only returned when ShowCacheNodeInfo is set; Memcached | ||
| // clients connect to them directly, so they are reachable endpoints too. | ||
| err = client.DescribeCacheClustersPages(&elasticache.DescribeCacheClustersInput{ShowCacheNodeInfo: aws.Bool(true)}, func(page *elasticache.DescribeCacheClustersOutput, _ bool) bool { | ||
| for _, cluster := range page.CacheClusters { | ||
| var metadata map[string]string | ||
| if ep.options.ExtendedMetadata { | ||
| metadata = getCacheClusterMetadata(cluster) | ||
| } | ||
| appendEndpoint(cluster.ConfigurationEndpoint, metadata) | ||
| for _, node := range cluster.CacheNodes { | ||
| appendEndpoint(node.Endpoint, metadata) | ||
| } | ||
| } | ||
| return true | ||
| }) | ||
| if err != nil { | ||
| return list, errors.Wrap(err, "could not describe elasticache cache clusters") | ||
| } | ||
|
|
||
| // Serverless caches are not available in every region, so a failure here | ||
| // must not discard the clusters already found. | ||
| _ = client.DescribeServerlessCachesPages(&elasticache.DescribeServerlessCachesInput{}, func(page *elasticache.DescribeServerlessCachesOutput, _ bool) bool { | ||
| for _, cache := range page.ServerlessCaches { | ||
| var metadata map[string]string | ||
| if ep.options.ExtendedMetadata { | ||
| metadata = getServerlessCacheMetadata(cache) | ||
| } | ||
| appendEndpoint(cache.Endpoint, metadata) | ||
| appendEndpoint(cache.ReaderEndpoint, metadata) | ||
| } | ||
| return true | ||
| }) | ||
|
|
||
| return list, nil | ||
| } | ||
|
|
||
| func getReplicationGroupMetadata(group *elasticache.ReplicationGroup) map[string]string { | ||
| metadata := make(map[string]string) | ||
| schema.AddMetadata(metadata, "replication_group_id", group.ReplicationGroupId) | ||
| schema.AddMetadata(metadata, "arn", group.ARN) | ||
| schema.AddMetadata(metadata, "status", group.Status) | ||
| schema.AddMetadata(metadata, "cache_node_type", group.CacheNodeType) | ||
| schema.AddMetadata(metadata, "cluster_mode", group.ClusterMode) | ||
| if group.TransitEncryptionEnabled != nil { | ||
| metadata["transit_encryption_enabled"] = fmt.Sprintf("%t", *group.TransitEncryptionEnabled) | ||
| } | ||
| if group.AuthTokenEnabled != nil { | ||
| metadata["auth_token_enabled"] = fmt.Sprintf("%t", *group.AuthTokenEnabled) | ||
| } | ||
| return metadata | ||
| } | ||
|
|
||
| func getCacheClusterMetadata(cluster *elasticache.CacheCluster) map[string]string { | ||
| metadata := make(map[string]string) | ||
| schema.AddMetadata(metadata, "cache_cluster_id", cluster.CacheClusterId) | ||
| schema.AddMetadata(metadata, "arn", cluster.ARN) | ||
| schema.AddMetadata(metadata, "replication_group_id", cluster.ReplicationGroupId) | ||
| schema.AddMetadata(metadata, "engine", cluster.Engine) | ||
| schema.AddMetadata(metadata, "engine_version", cluster.EngineVersion) | ||
| schema.AddMetadata(metadata, "status", cluster.CacheClusterStatus) | ||
| schema.AddMetadata(metadata, "cache_node_type", cluster.CacheNodeType) | ||
| schema.AddMetadata(metadata, "availability_zone", cluster.PreferredAvailabilityZone) | ||
| if cluster.TransitEncryptionEnabled != nil { | ||
| metadata["transit_encryption_enabled"] = fmt.Sprintf("%t", *cluster.TransitEncryptionEnabled) | ||
| } | ||
| if cluster.CacheClusterCreateTime != nil { | ||
| metadata["created_at"] = cluster.CacheClusterCreateTime.Format(time.RFC3339) | ||
| } | ||
| return metadata | ||
| } | ||
|
|
||
| func getServerlessCacheMetadata(cache *elasticache.ServerlessCache) map[string]string { | ||
| metadata := make(map[string]string) | ||
| schema.AddMetadata(metadata, "serverless_cache_name", cache.ServerlessCacheName) | ||
| schema.AddMetadata(metadata, "arn", cache.ARN) | ||
| schema.AddMetadata(metadata, "engine", cache.Engine) | ||
| schema.AddMetadata(metadata, "engine_version", cache.FullEngineVersion) | ||
| schema.AddMetadata(metadata, "status", cache.Status) | ||
| if len(cache.SecurityGroupIds) > 0 { | ||
| metadata["security_group_ids"] = strings.Join(aws.StringValueSlice(cache.SecurityGroupIds), ",") | ||
| } | ||
| if cache.CreateTime != nil { | ||
| metadata["created_at"] = cache.CreateTime.Format(time.RFC3339) | ||
| } | ||
| return metadata | ||
| } | ||
|
|
||
| func (ep *elastiCacheProvider) getElastiCacheClients(region *string) []*elasticache.ElastiCache { | ||
| clients := make([]*elasticache.ElastiCache, 0) | ||
|
|
||
| clients = append(clients, elasticache.New( | ||
| ep.session, | ||
| aws.NewConfig().WithRegion(aws.StringValue(region)), | ||
| )) | ||
|
|
||
| if ep.options.AssumeRoleName == "" || len(ep.options.AccountIds) < 1 { | ||
| return clients | ||
| } | ||
|
|
||
| for _, accountId := range ep.options.AccountIds { | ||
| roleARN := fmt.Sprintf("arn:aws:iam::%s:role/%s", accountId, ep.options.AssumeRoleName) | ||
| creds := stscreds.NewCredentials(ep.session, roleARN) | ||
|
|
||
| assumeSession, err := session.NewSession(&aws.Config{ | ||
| Region: region, | ||
| Credentials: creds, | ||
| }) | ||
| if err != nil { | ||
| continue | ||
| } | ||
|
|
||
| clients = append(clients, elasticache.New(assumeSession)) | ||
| } | ||
| return clients | ||
| } | ||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
🩺 Stability & Availability | 🟡 Minor | ⚡ Quick win
🔎 Supported by static analysis
🏁 Script executed:
Repository: projectdiscovery/cloudlist
Length of output: 6943
🏁 Script executed:
Repository: projectdiscovery/cloudlist
Length of output: 17299
🏁 Script executed:
Repository: projectdiscovery/cloudlist
Length of output: 14395
Pass
ctxto the AWS pagination calls.GetResourcereceivesctx, but each worker callslistElastiCacheResourceswithout it and then waits inwg.Wait(). The helper uses non-context pagination methods, so cancellation cannot stop an active request or prevent later listing calls in that worker.Use the context-aware methods and update the test callers.
Suggested fix
Update the three
listElastiCacheResourcestest callers to passcontext.Background().📝 Committable suggestion
🤖 Prompt for AI Agents