Skip to content

Commit c35c4b5

Browse files
authored
feat: Include cosmos-db cleanup (#209)
Signed-off-by: Jorge Turrado <jorge_turrado@hotmail.es>
1 parent b8adae4 commit c35c4b5

5 files changed

Lines changed: 330 additions & 41 deletions

File tree

garbage-collector/go.mod

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -6,6 +6,7 @@ require (
66
cloud.google.com/go/spanner v1.94.0
77
github.com/Azure/azure-sdk-for-go/sdk/azcore v1.22.0
88
github.com/Azure/azure-sdk-for-go/sdk/azidentity v1.14.0
9+
github.com/Azure/azure-sdk-for-go/sdk/resourcemanager/cosmos/armcosmos/v3 v3.4.0
910
github.com/Azure/azure-sdk-for-go/sdk/resourcemanager/eventhub/armeventhub v1.3.0
1011
github.com/Azure/azure-sdk-for-go/sdk/resourcemanager/resources/armresources v1.2.0
1112
github.com/Azure/azure-sdk-for-go/sdk/resourcemanager/servicebus/armservicebus v1.2.0

garbage-collector/go.sum

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -20,6 +20,8 @@ github.com/Azure/azure-sdk-for-go/sdk/azidentity/cache v0.4.0 h1:xFaZZ+IubdftrDH
2020
github.com/Azure/azure-sdk-for-go/sdk/azidentity/cache v0.4.0/go.mod h1:mCBhUhlMjLLJKr5aqw2TNS/VqJOie8MzWq3DAMJeKso=
2121
github.com/Azure/azure-sdk-for-go/sdk/internal v1.12.0 h1:fhqpLE3UEXi9lPaBRpQ6XuRW0nU7hgg4zlmZZa+a9q4=
2222
github.com/Azure/azure-sdk-for-go/sdk/internal v1.12.0/go.mod h1:7dCRMLwisfRH3dBupKeNCioWYUZ4SS09Z14H+7i8ZoY=
23+
github.com/Azure/azure-sdk-for-go/sdk/resourcemanager/cosmos/armcosmos/v3 v3.4.0 h1:+EhRnIOLvffCvUMUfP+MgOp6PrtN1d6xt94DZtrC3lA=
24+
github.com/Azure/azure-sdk-for-go/sdk/resourcemanager/cosmos/armcosmos/v3 v3.4.0/go.mod h1:Bb7kqorvA2acMCNFac+2ldoQWi7QrcMdH+9Gg9C7fSM=
2325
github.com/Azure/azure-sdk-for-go/sdk/resourcemanager/eventhub/armeventhub v1.3.0 h1:4hGvxD72TluuFIXVr8f4XkKZfqAa7Pj61t0jmQ7+kes=
2426
github.com/Azure/azure-sdk-for-go/sdk/resourcemanager/eventhub/armeventhub v1.3.0/go.mod h1:TSH7DcFItwAufy0Lz+Ft2cyopExCpxbOxI5SkH4dRNo=
2527
github.com/Azure/azure-sdk-for-go/sdk/resourcemanager/internal/v2 v2.0.0 h1:PTFGRSlMKCQelWwxUyYVEUqseBJVemLyqWJjvMyt0do=
Lines changed: 152 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,152 @@
1+
package cosmosdb
2+
3+
import (
4+
"context"
5+
"fmt"
6+
"time"
7+
8+
"github.com/Azure/azure-sdk-for-go/sdk/azcore"
9+
"github.com/Azure/azure-sdk-for-go/sdk/resourcemanager/cosmos/armcosmos/v3"
10+
"github.com/Azure/azure-sdk-for-go/sdk/resourcemanager/resources/armresources"
11+
12+
"github.com/kedacore/testing-infrastructure/garbage-colletor/internal/config"
13+
"github.com/kedacore/testing-infrastructure/garbage-colletor/internal/core"
14+
)
15+
16+
type Cleaner struct {
17+
rgClient *armresources.ResourceGroupsClient
18+
accountsClient *armcosmos.DatabaseAccountsClient
19+
sqlClient *armcosmos.SQLResourcesClient
20+
cfg config.Config
21+
}
22+
23+
func New(cred azcore.TokenCredential, cfg config.Config) (*Cleaner, error) {
24+
rgClient, err := armresources.NewResourceGroupsClient(cfg.AzureSubscriptionID, cred, nil)
25+
if err != nil {
26+
return nil, fmt.Errorf("creating azure resource groups client: %w", err)
27+
}
28+
accountsClient, err := armcosmos.NewDatabaseAccountsClient(cfg.AzureSubscriptionID, cred, nil)
29+
if err != nil {
30+
return nil, fmt.Errorf("creating azure cosmosdb accounts client: %w", err)
31+
}
32+
sqlClient, err := armcosmos.NewSQLResourcesClient(cfg.AzureSubscriptionID, cred, nil)
33+
if err != nil {
34+
return nil, fmt.Errorf("creating azure cosmosdb sql resources client: %w", err)
35+
}
36+
37+
return &Cleaner{
38+
rgClient: rgClient,
39+
accountsClient: accountsClient,
40+
sqlClient: sqlClient,
41+
cfg: cfg,
42+
}, nil
43+
}
44+
45+
func (c *Cleaner) Name() string {
46+
return "azure-cosmosdb"
47+
}
48+
49+
func (c *Cleaner) Run(ctx context.Context) core.Result {
50+
result := core.Result{Name: c.Name(), DryRun: c.cfg.DryRun}
51+
cutoff := time.Now().Add(-c.cfg.MaxAge)
52+
53+
rgPager := c.rgClient.NewListPager(nil)
54+
for rgPager.More() {
55+
rgPage, err := rgPager.NextPage(ctx)
56+
if err != nil {
57+
fmt.Printf("[%s] listing resource groups failed: %v\n", c.Name(), err)
58+
result.Errors++
59+
break
60+
}
61+
62+
for _, rg := range rgPage.Value {
63+
rgName := str(rg.Name)
64+
if rgName == "" {
65+
continue
66+
}
67+
68+
accountPager := c.accountsClient.NewListByResourceGroupPager(rgName, nil)
69+
for accountPager.More() {
70+
accountPage, err := accountPager.NextPage(ctx)
71+
if err != nil {
72+
fmt.Printf("[%s] listing cosmosdb accounts failed in resource group %s: %v\n", c.Name(), rgName, err)
73+
result.Errors++
74+
break
75+
}
76+
77+
for _, account := range accountPage.Value {
78+
accountName := str(account.Name)
79+
if accountName == "" {
80+
continue
81+
}
82+
83+
c.cleanupSQLDatabases(ctx, rgName, accountName, cutoff, &result)
84+
}
85+
}
86+
}
87+
}
88+
89+
return result
90+
}
91+
92+
func (c *Cleaner) cleanupSQLDatabases(ctx context.Context, resourceGroup, account string, cutoff time.Time, result *core.Result) {
93+
pager := c.sqlClient.NewListSQLDatabasesPager(resourceGroup, account, nil)
94+
for pager.More() {
95+
page, err := pager.NextPage(ctx)
96+
if err != nil {
97+
fmt.Printf("[%s] listing sql databases failed in %s/%s: %v\n", c.Name(), resourceGroup, account, err)
98+
result.Errors++
99+
break
100+
}
101+
102+
for _, db := range page.Value {
103+
result.Found++
104+
name := str(db.Name)
105+
if name == "" {
106+
continue
107+
}
108+
109+
createdAt := databaseCreatedAt(db)
110+
if createdAt == nil || createdAt.After(cutoff) {
111+
continue
112+
}
113+
114+
if c.cfg.DryRun {
115+
fmt.Printf("[%s][dry-run] would delete sql database %s/%s/%s (created %s)\n", c.Name(), resourceGroup, account, name, createdAt.UTC().Format(time.RFC3339))
116+
result.Deleted++
117+
continue
118+
}
119+
120+
poller, err := c.sqlClient.BeginDeleteSQLDatabase(ctx, resourceGroup, account, name, nil)
121+
if err != nil {
122+
fmt.Printf("[%s] delete failed for sql database %s/%s/%s: %v\n", c.Name(), resourceGroup, account, name, err)
123+
result.Errors++
124+
continue
125+
}
126+
if _, err := poller.PollUntilDone(ctx, nil); err != nil {
127+
fmt.Printf("[%s] delete failed for sql database %s/%s/%s: %v\n", c.Name(), resourceGroup, account, name, err)
128+
result.Errors++
129+
continue
130+
}
131+
132+
fmt.Printf("[%s] deleted sql database %s/%s/%s\n", c.Name(), resourceGroup, account, name)
133+
result.Deleted++
134+
}
135+
}
136+
}
137+
138+
// Cosmos DB exposes the last update time as `_ts`, an epoch seconds timestamp.
139+
func databaseCreatedAt(db *armcosmos.SQLDatabaseGetResults) *time.Time {
140+
if db == nil || db.Properties == nil || db.Properties.Resource == nil || db.Properties.Resource.Ts == nil {
141+
return nil
142+
}
143+
createdAt := time.Unix(int64(*db.Properties.Resource.Ts), 0)
144+
return &createdAt
145+
}
146+
147+
func str(v *string) string {
148+
if v == nil {
149+
return ""
150+
}
151+
return *v
152+
}

garbage-collector/internal/config/config.go

Lines changed: 4 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,6 @@
11
package config
22

33
import (
4-
"errors"
54
"fmt"
65
"os"
76
"strconv"
@@ -24,26 +23,18 @@ type Config struct {
2423
}
2524

2625
func Load(dryRun bool) (Config, error) {
27-
subscriptionID := os.Getenv("AZURE_SUBSCRIPTION_ID")
28-
if subscriptionID == "" {
29-
return Config{}, errors.New("AZURE_SUBSCRIPTION_ID is required")
30-
}
31-
32-
gcpProjectID := os.Getenv("GCP_PROJECT_ID")
33-
if gcpProjectID == "" {
34-
return Config{}, errors.New("GCP_PROJECT_ID is required")
35-
}
36-
3726
maxAgeHours, err := readMaxAgeHours()
3827
if err != nil {
3928
return Config{}, err
4029
}
4130

31+
// Provider credentials are validated by each cleaner, so a partial
32+
// environment is enough to run a subset of them.
4233
return Config{
43-
AzureSubscriptionID: subscriptionID,
34+
AzureSubscriptionID: os.Getenv("AZURE_SUBSCRIPTION_ID"),
4435
ExcludedSBTopicSuffixes: []string{PermanentServiceBusTopicSuffix},
4536
AWSRegion: DefaultAWSRegion,
46-
GCPProjectID: gcpProjectID,
37+
GCPProjectID: os.Getenv("GCP_PROJECT_ID"),
4738
MaxAge: time.Duration(maxAgeHours) * time.Hour,
4839
DryRun: dryRun,
4940
}, nil

0 commit comments

Comments
 (0)