From cca4da935f9cdfbc0abb1651cd55720a21d4e668 Mon Sep 17 00:00:00 2001 From: xuziyi Date: Tue, 20 Aug 2024 18:04:22 +0800 Subject: [PATCH] feat: get os traffic from minio (#4968) * get os traffic from minio * optimize get obj traffic within the last hour --------- Co-authored-by: jiahui --- controllers/pkg/database/interface.go | 6 + controllers/pkg/database/mongo/account.go | 147 ++++++++++++++++++ .../pkg/database/mongo/account_test.go | 119 ++++++++++++++ .../pkg/objectstorage/metric_parser.go | 34 ++-- .../pkg/objectstorage/objectstorage.go | 118 +++++++++++--- .../pkg/objectstorage/objectstorage_test.go | 18 ++- controllers/pkg/types/traffic.go | 28 ++++ .../controllers/monitor_controller.go | 91 +++++++++-- controllers/resources/main.go | 4 + 9 files changed, 514 insertions(+), 51 deletions(-) create mode 100644 controllers/pkg/types/traffic.go diff --git a/controllers/pkg/database/interface.go b/controllers/pkg/database/interface.go index 6f364ac4b..0d4b110ca 100644 --- a/controllers/pkg/database/interface.go +++ b/controllers/pkg/database/interface.go @@ -36,6 +36,7 @@ type Interface interface { Account Traffic CVM + Creator } type CVM interface { @@ -49,6 +50,10 @@ type Account interface { GetBillingHistoryNamespaceList(ns *accountv1.NamespaceBillingHistorySpec, owner string) ([]string, error) GetBillingHistoryNamespaces(startTime, endTime *time.Time, billType int, owner string) ([]string, error) SaveBillings(billing ...*resources.Billing) error + SaveObjTraffic(obs ...*types.ObjectStorageTraffic) error + GetAllLatestObjTraffic(startTime, endTime time.Time) ([]types.ObjectStorageTraffic, error) + HandlerTimeObjBucketSentTraffic(startTime, endTime time.Time, bucket string) (int64, error) + GetTimeObjBucketBucket(startTime, endTime time.Time) ([]string, error) QueryBillingRecords(billingRecordQuery *accountv1.BillingRecordQuery, owner string) error GetUnsettingBillingHandler(owner string) ([]resources.BillingHandler, error) UpdateBillingStatus(orderID string, status resources.BillingStatus) error @@ -117,6 +122,7 @@ type Creator interface { CreateBillingIfNotExist() error //suffix by day, eg: monitor_20200101 CreateMonitorTimeSeriesIfNotExist(collTime time.Time) error + CreateTTLTrafficTimeSeries() error } type MeteringOwnerTimeResult struct { diff --git a/controllers/pkg/database/mongo/account.go b/controllers/pkg/database/mongo/account.go index 2ca0d2a8d..53c8adb3a 100644 --- a/controllers/pkg/database/mongo/account.go +++ b/controllers/pkg/database/mongo/account.go @@ -22,6 +22,8 @@ import ( "strings" "time" + "github.com/labring/sealos/controllers/pkg/types" + "github.com/labring/sealos/controllers/pkg/utils/env" "github.com/labring/sealos/controllers/pkg/common" @@ -57,6 +59,7 @@ const ( DefaultMeteringConn = "metering" DefaultMonitorConn = "monitor" DefaultBillingConn = "billing" + DefaultObjTrafficConn = "objectstorage-traffic" DefaultUserConn = "user" DefaultPricesConn = "prices" DefaultPropertiesConn = "properties" @@ -75,6 +78,7 @@ type mongoDB struct { MonitorConnPrefix string MeteringConn string BillingConn string + ObjTrafficConn string PropertiesConn string TrafficConn string } @@ -246,6 +250,130 @@ func (m *mongoDB) SaveBillings(billing ...*resources.Billing) error { return err } +func (m *mongoDB) SaveObjTraffic(obs ...*types.ObjectStorageTraffic) error { + traffic := make([]interface{}, len(obs)) + for i, ob := range obs { + traffic[i] = ob + } + _, err := m.getObjTrafficCollection().InsertMany(context.Background(), traffic) + return err +} + +func (m *mongoDB) GetAllLatestObjTraffic(startTime, endTime time.Time) ([]types.ObjectStorageTraffic, error) { + pipeline := []bson.M{ + { + "$match": bson.M{ + "time": bson.M{ + "$gt": startTime, + "$lte": endTime, + }, + }, + }, + { + "$sort": bson.M{"time": -1}, + }, + { + "$group": bson.M{ + "_id": bson.M{"user": "$user", "bucket": "$bucket"}, + "latestDoc": bson.M{"$first": "$$ROOT"}, + }, + }, + { + "$replaceRoot": bson.M{"newRoot": "$latestDoc"}, + }, + } + + cursor, err := m.getObjTrafficCollection().Aggregate(context.Background(), pipeline) + if err != nil { + return nil, err + } + defer cursor.Close(context.Background()) + + var results []types.ObjectStorageTraffic + if err = cursor.All(context.Background(), &results); err != nil { + return nil, err + } + + return results, nil +} + +func (m *mongoDB) HandlerTimeObjBucketSentTraffic(startTime, endTime time.Time, bucket string) (int64, error) { + pipeline := []bson.M{ + { + "$match": bson.M{ + "time": bson.M{ + "$gt": startTime, + "$lte": endTime, + }, + "bucket": bucket, + }, + }, + { + "$group": bson.M{ + "_id": nil, + "totalSent": bson.M{"$sum": "$sent"}, + }, + }, + } + cursor, err := m.getObjTrafficCollection().Aggregate(context.Background(), pipeline) + if err != nil { + return 0, err + } + defer cursor.Close(context.Background()) + + var result struct { + TotalSent int64 `bson:"totalSent"` + } + if cursor.Next(context.Background()) { + if err := cursor.Decode(&result); err != nil { + return 0, err + } + return result.TotalSent, nil + } + if err := cursor.Err(); err != nil { + return 0, err + } + return 0, nil +} + +func (m *mongoDB) GetTimeObjBucketBucket(startTime, endTime time.Time) ([]string, error) { + pipeline := []bson.M{ + { + "$match": bson.M{ + "time": bson.M{ + "$gt": startTime, + "$lte": endTime, + }, + }, + }, + { + "$group": bson.M{ + "_id": nil, + "buckets": bson.M{"$addToSet": "$bucket"}, + }, + }, + } + cursor, err := m.getObjTrafficCollection().Aggregate(context.Background(), pipeline) + if err != nil { + return nil, err + } + defer cursor.Close(context.Background()) + + var result struct { + Buckets []string `bson:"buckets"` + } + if cursor.Next(context.Background()) { + if err := cursor.Decode(&result); err != nil { + return nil, err + } + return result.Buckets, nil + } + if err := cursor.Err(); err != nil { + return nil, err + } + return nil, nil +} + // InsertMonitor insert monitor data to mongodb collection monitor + time (eg: monitor_20200101) // The monitor data is saved daily 2020-12-01 00:00:00 - 2020-12-01 23:59:59 => monitor_20201201 func (m *mongoDB) InsertMonitor(ctx context.Context, monitors ...*resources.Monitor) error { @@ -885,6 +1013,9 @@ func (m *mongoDB) getMonitorCollectionName(collTime time.Time) string { func (m *mongoDB) getBillingCollection() *mongo.Collection { return m.Client.Database(m.AccountDB).Collection(m.BillingConn) } +func (m *mongoDB) getObjTrafficCollection() *mongo.Collection { + return m.Client.Database(m.AccountDB).Collection(m.ObjTrafficConn) +} func (m *mongoDB) getPropertiesCollection() *mongo.Collection { return m.Client.Database(m.AccountDB).Collection(m.PropertiesConn) @@ -941,6 +1072,21 @@ func (m *mongoDB) CreateTimeSeriesIfNotExist(dbName, collectionName string) erro return m.Client.Database(dbName).RunCommand(context.TODO(), cmd).Err() } +func (m *mongoDB) CreateTTLTrafficTimeSeries() error { + // Check if the collection already exists + if exist, err := m.collectionExist(m.AccountDB, m.ObjTrafficConn); exist || err != nil { + return err + } + // If the collection does not exist, create it + cmd := bson.D{ + primitive.E{Key: "create", Value: m.ObjTrafficConn}, + primitive.E{Key: "timeseries", Value: bson.D{{Key: "timeField", Value: "time"}}}, + //default ttl set 30 days + primitive.E{Key: "expireAfterSeconds", Value: 30 * 24 * 60 * 60}, + } + return m.Client.Database(m.AccountDB).RunCommand(context.TODO(), cmd).Err() +} + func (m *mongoDB) DropMonitorCollectionsOlderThan(days int) error { db := m.Client.Database(m.AccountDB) // Get the current time minus the number of days @@ -985,6 +1131,7 @@ func NewMongoInterface(ctx context.Context, URL string) (database.Interface, err MeteringConn: DefaultMeteringConn, MonitorConnPrefix: DefaultMonitorConn, BillingConn: DefaultBillingConn, + ObjTrafficConn: DefaultObjTrafficConn, PropertiesConn: DefaultPropertiesConn, TrafficConn: env.GetEnvWithDefault(EnvTrafficConn, DefaultTrafficConn), CvmConn: env.GetEnvWithDefault(EnvCVMConn, DefaultCVMConn), diff --git a/controllers/pkg/database/mongo/account_test.go b/controllers/pkg/database/mongo/account_test.go index d8cbfa0fe..f8480747a 100644 --- a/controllers/pkg/database/mongo/account_test.go +++ b/controllers/pkg/database/mongo/account_test.go @@ -22,6 +22,8 @@ import ( "testing" "time" + "github.com/labring/sealos/controllers/pkg/types" + "github.com/labring/sealos/controllers/pkg/resources" "github.com/dustin/go-humanize" @@ -595,3 +597,120 @@ func Test_mongoDB_GetDistinctMonitorCombinations(t *testing.T) { } t.Logf("monitorCombinations: %v", monitorCombinations) } + +func Test_mongoDB_CreateTTLTrafficTimeSeries(t *testing.T) { + dbCTX := context.Background() + + m, err := NewMongoInterface(dbCTX, os.Getenv("MONGODB_URI")) + if err != nil { + t.Errorf("failed to connect mongo: error = %v", err) + } + defer func() { + if err = m.Disconnect(dbCTX); err != nil { + t.Errorf("failed to disconnect mongo: error = %v", err) + } + }() + + if err = m.CreateTTLTrafficTimeSeries(); err != nil { + t.Fatalf("failed to create TTL traffic time series: %v", err) + } + t.Logf("create TTL traffic time series success") +} + +func Test_mongoDB_SaveObjTraffic(t *testing.T) { + dbCTX := context.Background() + + m, err := NewMongoInterface(dbCTX, os.Getenv("MONGODB_URI")) + if err != nil { + t.Errorf("failed to connect mongo: error = %v", err) + } + defer func() { + if err = m.Disconnect(dbCTX); err != nil { + t.Errorf("failed to disconnect mongo: error = %v", err) + } + }() + var traffic []*types.ObjectStorageTraffic + for i := 0; i < 10; i++ { + traffic = append(traffic, &types.ObjectStorageTraffic{ + Time: time.Now().UTC(), + User: "user-" + fmt.Sprint(i), + Bucket: "bucket-" + fmt.Sprint(i), + TotalSent: int64(1000 + i), + Sent: int64(100 + i), + }) + } + if err = m.SaveObjTraffic(traffic...); err != nil { + t.Fatalf("failed to save object storage traffic: %v", err) + } + t.Logf("save object storage traffic success") +} + +func Test_mongoDB_GetAllLatestObjTraffic(t *testing.T) { + dbCTX := context.Background() + + m, err := NewMongoInterface(dbCTX, os.Getenv("MONGODB_URI")) + if err != nil { + t.Errorf("failed to connect mongo: error = %v", err) + } + defer func() { + if err = m.Disconnect(dbCTX); err != nil { + t.Errorf("failed to disconnect mongo: error = %v", err) + } + }() + + traffic, err := m.GetAllLatestObjTraffic(time.Now().Add(-time.Hour), time.Now()) + if err != nil { + t.Fatalf("failed to save object storage traffic: %v", err) + } + t.Logf("save object storage traffic success") + for _, tf := range traffic { + t.Logf("traffic: %#+v", tf) + } +} + +func Test_mongoDB_HandlerTimeObjBucketSentTraffic(t *testing.T) { + dbCTX := context.Background() + + m, err := NewMongoInterface(dbCTX, os.Getenv("MONGODB_URI")) + if err != nil { + t.Errorf("failed to connect mongo: error = %v", err) + } + defer func() { + if err = m.Disconnect(dbCTX); err != nil { + t.Errorf("failed to disconnect mongo: error = %v", err) + } + }() + + bytes, err := m.HandlerTimeObjBucketSentTraffic(time.Now().UTC().Add(-time.Hour), time.Now().UTC(), "bucket-6") + if err != nil { + t.Fatalf("failed to handle time object bucket usage: %v", err) + } + t.Logf("handle time object bucket usage success: %v", bytes) +} + +func init() { + os.Setenv("MONGODB_URI", "") +} + +func Test_mongoDB_GetTimeObjBucketBucket(t *testing.T) { + dbCTX := context.Background() + + m, err := NewMongoInterface(dbCTX, os.Getenv("MONGODB_URI")) + if err != nil { + t.Errorf("failed to connect mongo: error = %v", err) + } + defer func() { + if err = m.Disconnect(dbCTX); err != nil { + t.Errorf("failed to disconnect mongo: error = %v", err) + } + }() + + buckets, err := m.GetTimeObjBucketBucket(time.Now().UTC().Add(-10*time.Hour), time.Now().UTC()) + if err != nil { + t.Fatalf("failed to get time object bucket bucket: %v", err) + } + t.Logf("get time object bucket bucket success: len: %v", len(buckets)) + for _, bucket := range buckets { + t.Logf("bucket: %#+v", bucket) + } +} diff --git a/controllers/pkg/objectstorage/metric_parser.go b/controllers/pkg/objectstorage/metric_parser.go index 1d4b0ee60..47ec33454 100644 --- a/controllers/pkg/objectstorage/metric_parser.go +++ b/controllers/pkg/objectstorage/metric_parser.go @@ -75,13 +75,18 @@ func NewMetricsClient(endpoint string, accessKeyID, secretAccessKey string, secu return privateNewMetricsClient(endpointURL, jwtToken, secure) } -// BucketUsageTotalBytesMetrics - returns Bucket Metrics in Prometheus format +// BucketUsageTotalBytesMetrics - returns Bucket Usage Total Metrics in Prometheus format func (client *MetricsClient) BucketUsageTotalBytesMetrics(ctx context.Context) ([]*prom2json.Family, error) { - return client.fetchMetrics(ctx, "bucket", "minio_bucket_usage_total_bytes") + return client.fetchMetrics(ctx, "bucket", []string{"minio_bucket_usage_total_bytes"}) +} + +// BucketUsageAndTrafficBytesMetrics - returns Bucket Usage And Traffic Metrics in Prometheus format +func (client *MetricsClient) BucketUsageAndTrafficBytesMetrics(ctx context.Context) ([]*prom2json.Family, error) { + return client.fetchMetrics(ctx, "bucket", []string{"minio_bucket_usage_total_bytes", "minio_bucket_traffic_sent_bytes", "minio_bucket_traffic_received_bytes"}) } // fetchMetrics - returns Metrics of given subsystem in Prometheus format -func (client *MetricsClient) fetchMetrics(ctx context.Context, subSystem string, metricsName string) ([]*prom2json.Family, error) { +func (client *MetricsClient) fetchMetrics(ctx context.Context, subSystem string, metricsName []string) ([]*prom2json.Family, error) { reqData := metricsRequestData{ relativePath: "/v2/metrics/" + subSystem, } @@ -119,7 +124,7 @@ func closeResponse(resp *http.Response) { } } -func parsePrometheusResults(reader io.Reader, prefix string) (results []*prom2json.Family, err error) { +func parsePrometheusResults(reader io.Reader, prefix []string) (results []*prom2json.Family, err error) { filteredReader, err := filterMetricsByPrefix(reader, prefix) if err != nil { return nil, err @@ -136,10 +141,11 @@ func parsePrometheusResults(reader io.Reader, prefix string) (results []*prom2js }() for mf := range mfChan { - if !strings.Contains(mf.GetName(), prefix) { - continue + for i := range prefix { + if strings.Contains(mf.GetName(), prefix[i]) { + results = append(results, prom2json.NewFamily(mf)) + } } - results = append(results, prom2json.NewFamily(mf)) } if err := <-errChan; err != nil { return nil, err @@ -147,7 +153,7 @@ func parsePrometheusResults(reader io.Reader, prefix string) (results []*prom2js return results, nil } -func filterMetricsByPrefix(reader io.Reader, prefix string) (io.Reader, error) { +func filterMetricsByPrefix(reader io.Reader, prefix []string) (io.Reader, error) { var buf bytes.Buffer for { line, err := readLine(reader) @@ -156,11 +162,17 @@ func filterMetricsByPrefix(reader io.Reader, prefix string) (io.Reader, error) { } else if err != nil { return nil, err } - if bytes.HasPrefix(line, []byte("#")) || !bytes.HasPrefix(line, []byte(prefix)) { + + if bytes.HasPrefix(line, []byte("#")) { continue } - if _, err := buf.Write(line); err != nil { - return nil, err + + for i := range prefix { + if bytes.HasPrefix(line, []byte(prefix[i])) { + if _, err := buf.Write(line); err != nil { + return nil, err + } + } } } return &buf, nil diff --git a/controllers/pkg/objectstorage/objectstorage.go b/controllers/pkg/objectstorage/objectstorage.go index f9d765370..ccc2985e3 100644 --- a/controllers/pkg/objectstorage/objectstorage.go +++ b/controllers/pkg/objectstorage/objectstorage.go @@ -199,6 +199,10 @@ func extractValues(result1, result2 model.Value) (int64, int64) { type MetricData struct { // key: bucket name, value: usage Usage map[string]int64 + // key: bucket name, value: traffic sent + Sent map[string]int64 + // key: bucket name, value: traffic received + Received map[string]int64 } type Metrics map[string]MetricData @@ -237,6 +241,57 @@ func QueryUserUsage(client *MetricsClient) (Metrics, error) { return obMetrics, err } +func QueryUserUsageAndTraffic(client *MetricsClient) (Metrics, error) { + obMetrics := make(Metrics) + bucketMetrics, err := client.BucketUsageAndTrafficBytesMetrics(context.TODO()) + if err != nil { + return nil, fmt.Errorf("failed to get bucket traffic metrics: %w", err) + } + for _, bucketMetric := range bucketMetrics { + if !isUsageAndTrafficBytesTargetMetric(bucketMetric.Name) { + continue + } + + for _, metrics := range bucketMetric.Metrics { + promMetrics := metrics.(prom2json.Metric) + floatValue, err := strconv.ParseFloat(promMetrics.Value, 64) + if err != nil { + return nil, fmt.Errorf("failed to parse %s to float value", promMetrics.Value) + } + intValue := int64(floatValue) + if bucket := promMetrics.Labels["bucket"]; bucket != "" { + //fmt.Println("debug info", "type:", bucketMetric.Name, "promMetrics:", promMetrics) + user := getUserWithBucket(bucket) + if user == "" { + //fmt.Println("debug info", "false bucket:", bucket) + continue + } + //fmt.Println("debug info", "true bucket:", bucket, "user:", user) + metricData, exists := obMetrics[user] + if !exists { + metricData = MetricData{ + Usage: make(map[string]int64), + Sent: make(map[string]int64), + Received: make(map[string]int64), + } + } + if bucketMetric.Name == "minio_bucket_usage_total_bytes" { + metricData.Usage[bucket] += intValue + } + if bucketMetric.Name == "minio_bucket_traffic_sent_bytes" { + metricData.Sent[bucket] += intValue + } + if bucketMetric.Name == "minio_bucket_traffic_received_bytes" { + metricData.Received[bucket] += intValue + } + obMetrics[user] = metricData + } + } + } + + return obMetrics, err +} + func isUsageBytesTargetMetric(name string) bool { targetMetrics := []string{ "minio_bucket_usage_total_bytes", @@ -249,8 +304,27 @@ func isUsageBytesTargetMetric(name string) bool { return false } +func isUsageAndTrafficBytesTargetMetric(name string) bool { + targetMetrics := []string{ + "minio_bucket_usage_total_bytes", + "minio_bucket_traffic_sent_bytes", + "minio_bucket_traffic_received_bytes", + } + for _, target := range targetMetrics { + if name == target { + return true + } + } + return false +} + func getUserWithBucket(bucket string) string { - return strings.Split(bucket, "-")[0] + re := regexp.MustCompile(`^([a-zA-Z0-9]{8})-(.*)$`) + matches := re.FindStringSubmatch(bucket) + if len(matches) == 3 { + return matches[1] + } + return "" } /* @@ -268,25 +342,25 @@ func getUserWithBucket(bucket string) string { 6/halo-faxdridb-pg-yhxnjm.tar.gz */ -func GetUserBakFileSize(client *minio.Client) map[string]int64 { - bucket := "file-backup" - userUsageMap := make(map[string]int64) - objectsCh := client.ListObjects(context.Background(), bucket, minio.ListObjectsOptions{Recursive: true}) - for object := range objectsCh { - user := extractNamespace(object.Key) - if user != "" { - userUsageMap[user] += object.Size - } - } +//func GetUserBakFileSize(client *minio.Client) map[string]int64 { +// bucket := "file-backup" +// userUsageMap := make(map[string]int64) +// objectsCh := client.ListObjects(context.Background(), bucket, minio.ListObjectsOptions{Recursive: true}) +// for object := range objectsCh { +// user := extractNamespace(object.Key) +// if user != "" { +// userUsageMap[user] += object.Size +// } +// } +// +// return userUsageMap +//} - return userUsageMap -} - -func extractNamespace(input string) string { - re := regexp.MustCompile(`ns-(\w+)`) - matches := re.FindStringSubmatch(input) - if len(matches) < 2 { - return "" - } - return matches[1] -} +//func extractNamespace(input string) string { +// re := regexp.MustCompile(`ns-(\w+)`) +// matches := re.FindStringSubmatch(input) +// if len(matches) < 2 { +// return "" +// } +// return matches[1] +//} diff --git a/controllers/pkg/objectstorage/objectstorage_test.go b/controllers/pkg/objectstorage/objectstorage_test.go index 7cd466b0c..e10c48d24 100644 --- a/controllers/pkg/objectstorage/objectstorage_test.go +++ b/controllers/pkg/objectstorage/objectstorage_test.go @@ -62,14 +62,20 @@ func TestQueryUserUsage(t *testing.T) { } } -func TestGetUserBakFileSize(t *testing.T) { - objStorageClient, err := objectstoragev1.NewOSClient("objectstorageapi.192.168.0.55.nip.io", "username", "passw0rd") +func TestQueryUserTraffic(t *testing.T) { + obClient, err := NewMetricsClient("objectstorageapi.192.168.0.55.nip.io", "username", "passw0rd", false) if err != nil { - t.Fatalf("NewOSClient error: %v", err) + t.Error(err) } - bakFileSize := GetUserBakFileSize(objStorageClient) + metrics, err := QueryUserUsageAndTraffic(obClient) if err != nil { - t.Fatalf("GetUserBakFileSize error: %v", err) + t.Error(err) + } + + for user, metric := range metrics { + fmt.Println("user:", user) + fmt.Println("usage:", metric.Usage) + fmt.Println("sent:", metric.Sent) + fmt.Println("received:", metric.Received) } - fmt.Println(bakFileSize) } diff --git a/controllers/pkg/types/traffic.go b/controllers/pkg/types/traffic.go new file mode 100644 index 000000000..44adde50d --- /dev/null +++ b/controllers/pkg/types/traffic.go @@ -0,0 +1,28 @@ +// Copyright © 2024 sealos. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package types + +import "time" + +type ObjectStorageTraffic struct { + Time time.Time `json:"time" bson:"time"` + User string `json:"user" bson:"user"` + Bucket string `json:"bucket" bson:"bucket"` + //bytes + TotalSent int64 `json:"totalSent" bson:"totalSent"` + + //The sent traffic since the last time + Sent int64 `json:"sent" bson:"sent"` +} diff --git a/controllers/resources/controllers/monitor_controller.go b/controllers/resources/controllers/monitor_controller.go index dfc55a7b9..bc1995c37 100644 --- a/controllers/resources/controllers/monitor_controller.go +++ b/controllers/resources/controllers/monitor_controller.go @@ -25,6 +25,8 @@ import ( "sync" "time" + "github.com/labring/sealos/controllers/pkg/types" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" appv1 "github.com/labring/sealos/controllers/app/api/v1" @@ -80,7 +82,8 @@ type MonitorReconciler struct { TrafficClient database.Interface Properties *resources.PropertyTypeLS PromURL string - currentObjectMetrics map[string]objstorage.MetricData + lastObjectMetrics objstorage.Metrics + currentObjectMetrics objstorage.Metrics ObjStorageClient *minio.Client ObjStorageMetricsClient *objstorage.MetricsClient ObjStorageUserBackupSize map[string]int64 @@ -288,22 +291,52 @@ func (r *MonitorReconciler) processNamespaceList(namespaceList *corev1.Namespace }(&namespaceList.Items[i]) } wg.Wait() + if err := r.monitorObjectStorageTraffic(); err != nil { + r.Logger.Error(err, "failed to monitor object storage traffic") + } logger.Info("end processNamespaceList", "time", time.Now().Format("2006-01-02 15:04:05")) return nil } func (r *MonitorReconciler) preMonitorResourceUsage() error { if r.ObjStorageMetricsClient != nil { - metrics, err := objstorage.QueryUserUsage(r.ObjStorageMetricsClient) + metrics, err := objstorage.QueryUserUsageAndTraffic(r.ObjStorageMetricsClient) if err != nil { + r.lastObjectMetrics = r.currentObjectMetrics return fmt.Errorf("failed to query object storage metrics: %w", err) } + if r.currentObjectMetrics != nil { + r.lastObjectMetrics = r.currentObjectMetrics + } else { + latestObjTrafficSentMetrics := make(objstorage.Metrics) + startTime, endTime := time.Now().UTC().Add(-time.Hour), time.Now().UTC() + traffic, err := r.DBClient.GetAllLatestObjTraffic(startTime, endTime) + if err != nil { + return fmt.Errorf("failed to get all latest object storage traffic: %w", err) + } + for i := range traffic { + user := traffic[i].User + bucket := traffic[i].Bucket + if _, ok := metrics[user]; !ok { + continue + } + + if traffic[i].Time.Before(time.Now().Add(-time.Hour)) { + continue + } + + if _, ok := latestObjTrafficSentMetrics[user]; !ok { + latestObjTrafficSentMetrics[user] = objstorage.MetricData{ + Sent: make(map[string]int64), + } + } + + latestObjTrafficSentMetrics[user].Sent[bucket] = traffic[i].TotalSent + } + r.lastObjectMetrics = latestObjTrafficSentMetrics + } r.currentObjectMetrics = metrics - logger.Info("success query object storage resource usage", "time", time.Now().Format("2006-01-02 15:04:05")) - } - if r.ObjStorageClient != nil { - r.ObjStorageUserBackupSize = objstorage.GetUserBakFileSize(r.ObjStorageClient) - logger.Info("success query object storage backup size", "time", time.Now().Format("2006-01-02 15:04:05")) + logger.Info("success query object storage usage and traffic metrics", "time", time.Now().Format("2006-01-02 15:04:05")) } return nil } @@ -457,7 +490,7 @@ func (r *MonitorReconciler) monitorDatabaseBackupUsage(namespace string, resUsed for i := range backupList.Items { backup := &backupList.Items[i] backupRes := resources.NewResourceNamed(backup) - fmt.Printf("backup name: %v, backup size: %v, backupRes: %s \n", backupList.Items[i].Name, backupList.Items[i].Status.TotalSize, backupRes.String()) + //fmt.Printf("backup name: %v, backup size: %v, backupRes: %s \n", backupList.Items[i].Name, backupList.Items[i].Status.TotalSize, backupRes.String()) if resUsed[backupRes.String()] == nil { resNamed[backupRes.String()] = backupRes resUsed[backupRes.String()] = initResources() @@ -533,6 +566,41 @@ func (r *MonitorReconciler) monitorObjectStorageUsage(namespace string, resMap m return nil } +func (r *MonitorReconciler) monitorObjectStorageTraffic() error { + if r.currentObjectMetrics == nil { + return nil + } + var objTraffic []*types.ObjectStorageTraffic + now := time.Now().UTC() + for user, metric := range r.currentObjectMetrics { + if len(metric.Sent) == 0 { + continue + } + for bucket, m := range metric.Sent { + sent := int64(0) + if r.lastObjectMetrics != nil && r.lastObjectMetrics[user].Sent != nil { + if _, ok := r.lastObjectMetrics[user].Sent[bucket]; ok { + ss := m - r.lastObjectMetrics[user].Sent[bucket] + if ss > 0 { + sent = ss + } + } + } + objTraffic = append(objTraffic, &types.ObjectStorageTraffic{ + Time: now, + User: user, + Bucket: bucket, + TotalSent: m, + Sent: sent, + }) + } + } + if err := r.DBClient.SaveObjTraffic(objTraffic...); err != nil { + return fmt.Errorf("failed to save object storage traffic: %w", err) + } + return nil +} + func (r *MonitorReconciler) MonitorTrafficUsed(startTime, endTime time.Time) error { logger.Info("start getTrafficUsed", "startTime", startTime.Format(time.RFC3339), "endTime", endTime.Format(time.RFC3339)) execTime := time.Now().UTC() @@ -551,9 +619,9 @@ func (r *MonitorReconciler) MonitorTrafficUsed(startTime, endTime time.Time) err } func (r *MonitorReconciler) monitorObjectStorageTrafficUsed(startTime, endTime time.Time) error { - buckets, err := objstorage.ListAllObjectStorageBucket(r.ObjStorageClient) + buckets, err := r.DBClient.GetTimeObjBucketBucket(startTime, endTime) if err != nil { - return fmt.Errorf("failed to list object storage buckets: %w", err) + return fmt.Errorf("failed to get object storage buckets: %w", err) } r.Logger.Info("object storage buckets", "buckets len", len(buckets)) wg, _ := errgroup.WithContext(context.Background()) @@ -571,7 +639,7 @@ func (r *MonitorReconciler) monitorObjectStorageTrafficUsed(startTime, endTime t } func (r *MonitorReconciler) handlerObjectStorageTrafficUsed(startTime, endTime time.Time, bucket string) error { - bytes, err := objstorage.GetObjectStorageFlow(r.PromURL, bucket, r.ObjectStorageInstance, startTime, endTime) + bytes, err := r.DBClient.HandlerTimeObjBucketSentTraffic(startTime, endTime, bucket) if err != nil { return fmt.Errorf("failed to get object storage flow: %w", err) } @@ -633,7 +701,6 @@ func (r *MonitorReconciler) handlerTrafficUsed(startTime, endTime time.Time, mon Time: endTime.Add(-1 * time.Minute), Type: monitor.Type, } - r.Logger.Info("monitor traffic used", "monitor", ro) err = r.DBClient.InsertMonitor(context.Background(), &ro) if err != nil { return fmt.Errorf("failed to insert monitor: %w", err) diff --git a/controllers/resources/main.go b/controllers/resources/main.go index 5677a7707..61f70d7f7 100644 --- a/controllers/resources/main.go +++ b/controllers/resources/main.go @@ -185,6 +185,10 @@ func main() { } else { reconciler.Logger.Info("minio info not found, please check env: MINIO_ENDPOINT, MINIO_AK, MINIO_SK, MINIO_METRICS_ADDR") } + err = reconciler.DBClient.CreateTTLTrafficTimeSeries() + if err != nil { + reconciler.Logger.Error(err, "failed to create ttl traffic time series") + } // timer creates tomorrow's timing table in advance to ensure that tomorrow's table exists // Execute immediately and then every 24 hours. time.AfterFunc(time.Until(getNextMidnight()), func() {