package metrics import ( "context" "fmt" "log" "os" "strings" "sync" "time" "github.com/aws/aws-sdk-go-v2/aws" "github.com/aws/aws-sdk-go-v2/config" "github.com/aws/aws-sdk-go-v2/credentials" "github.com/aws/aws-sdk-go-v2/service/s3" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/runtime/schema" "k8s.io/client-go/dynamic" "k8s.io/client-go/kubernetes" "k8s.io/client-go/rest" "k8s.io/client-go/tools/clientcmd" "admin-console/internal/model" ) // fission CRD group-version-resources. var ( gvrFunctions = schema.GroupVersionResource{Group: "fission.io", Version: "v1", Resource: "functions"} gvrPackages = schema.GroupVersionResource{Group: "fission.io", Version: "v1", Resource: "packages"} gvrEnvironments = schema.GroupVersionResource{Group: "fission.io", Version: "v1", Resource: "environments"} gvrHTTPTriggers = schema.GroupVersionResource{Group: "fission.io", Version: "v1", Resource: "httptriggers"} gvrTimeTriggers = schema.GroupVersionResource{Group: "fission.io", Version: "v1", Resource: "timetriggers"} ) // K8sClient — обёртка над kubernetes.Clientset для нужд admin-console. type K8sClient struct { cs *kubernetes.Clientset dc dynamic.Interface s3 *s3.Client s3Bucket string } // NewK8sClient создаёт клиента из in-cluster конфига, с fallback на ~/.kube/config. func NewK8sClient() (*K8sClient, error) { cfg, err := rest.InClusterConfig() if err != nil { cfg, err = clientcmd.BuildConfigFromFlags("", clientcmd.RecommendedHomeFile) if err != nil { return nil, fmt.Errorf("no k8s config: %w", err) } } cs, err := kubernetes.NewForConfig(cfg) if err != nil { return nil, err } dc, err := dynamic.NewForConfig(cfg) if err != nil { return nil, fmt.Errorf("dynamic client: %w", err) } k := &K8sClient{cs: cs, dc: dc} // S3 — опционально (не падаем если нет кредов) k.initS3() return k, nil } // initS3 создаёт S3-клиент из переменных окружения. func (k *K8sClient) initS3() { endpoint := os.Getenv("S3_ENDPOINT") bucket := os.Getenv("S3_BUCKET") accessKey := os.Getenv("S3_ACCESS_KEY") secretKey := os.Getenv("S3_SECRET_KEY") if endpoint == "" || bucket == "" || accessKey == "" || secretKey == "" { log.Println("[admin-console] S3 env vars not set — storage metrics disabled") return } cfg, err := config.LoadDefaultConfig(context.Background(), config.WithRegion("ru-msk-1"), config.WithCredentialsProvider(credentials.NewStaticCredentialsProvider(accessKey, secretKey, "")), ) if err != nil { log.Printf("[admin-console] S3 config error: %v", err) return } k.s3 = s3.NewFromConfig(cfg, func(o *s3.Options) { o.BaseEndpoint = aws.String(endpoint) o.UsePathStyle = true }) k.s3Bucket = bucket log.Printf("[admin-console] S3 client ready, bucket=%s", bucket) } // ListClientNamespaces возвращает namespace-ы с лейблом managed-by=fission-console. // Для каждого считает количество объектов Fission CRD через Dynamic клиент. func (k *K8sClient) ListClientNamespaces(ctx context.Context) ([]model.NamespaceInfo, error) { nsList, err := k.cs.CoreV1().Namespaces().List(ctx, metav1.ListOptions{ LabelSelector: "managed-by=fission-console", }) if err != nil { return nil, err } result := make([]model.NamespaceInfo, 0, len(nsList.Items)) for _, ns := range nsList.Items { info := model.NamespaceInfo{ Name: ns.Name, UserEmail: ns.Labels["user"], CreatedAt: ns.CreationTimestamp.Format(time.RFC3339), } k.fillCRDCounts(ctx, ns.Name, &info) result = append(result, info) } return result, nil } // fillCRDCounts заполняет счётчики Fission-ресурсов для одного namespace. func (k *K8sClient) fillCRDCounts(ctx context.Context, ns string, info *model.NamespaceInfo) { count := func(gvr schema.GroupVersionResource) int { list, err := k.dc.Resource(gvr).Namespace(ns).List(ctx, metav1.ListOptions{Limit: 1}) if err != nil { return 0 } // List возвращает items и list-metadata с remaining count. // Используем оставшийся счёт для точности. if rem := list.GetRemainingItemCount(); rem != nil { return int(*rem) + len(list.Items) } return len(list.Items) } var wg sync.WaitGroup wg.Add(5) go func() { defer wg.Done(); info.Functions = count(gvrFunctions) }() go func() { defer wg.Done(); info.Packages = count(gvrPackages) }() go func() { defer wg.Done(); info.Environments = count(gvrEnvironments) }() go func() { defer wg.Done(); info.HTTPTriggers = count(gvrHTTPTriggers) }() go func() { defer wg.Done(); info.TimeTriggers = count(gvrTimeTriggers) }() wg.Wait() } // CollectUsage возвращает агрегированную статистику по всем клиентским namespace-ам. func (k *K8sClient) CollectUsage(ctx context.Context) ([]model.UsageEntry, error) { nsList, err := k.cs.CoreV1().Namespaces().List(ctx, metav1.ListOptions{ LabelSelector: "managed-by=fission-console", }) if err != nil { return nil, err } now := time.Now().UTC() startOfMonth := time.Date(now.Year(), now.Month(), 1, 0, 0, 0, 0, time.UTC) result := make([]model.UsageEntry, 0, len(nsList.Items)) for _, ns := range nsList.Items { entry := model.UsageEntry{ Namespace: ns.Name, UserEmail: ns.Labels["user"], PeriodStart: startOfMonth.Format("2006-01-02"), PeriodEnd: now.Format("2006-01-02"), } // Считаем функции через Dynamic client. if list, err := k.dc.Resource(gvrFunctions).Namespace(ns.Name).List(ctx, metav1.ListOptions{Limit: 1}); err == nil { entry.Functions = len(list.Items) if rem := list.GetRemainingItemCount(); rem != nil { entry.Functions += int(*rem) } } // Считаем S3-хранилище. entry.StorageGB = k.namespaceStorage(ctx, ns.Name) // Invocations / CPU / Memory — TODO. result = append(result, entry) } return result, nil } // namespaceStorage возвращает суммарный размер объектов S3 с префиксом "{ns}/". func (k *K8sClient) namespaceStorage(ctx context.Context, ns string) float64 { if k.s3 == nil { return 0 } var totalBytes int64 paginator := s3.NewListObjectsV2Paginator(k.s3, &s3.ListObjectsV2Input{ Bucket: aws.String(k.s3Bucket), Prefix: aws.String(ns + "/"), }) for paginator.HasMorePages() { page, err := paginator.NextPage(ctx) if err != nil { log.Printf("[admin-console] S3 list error ns=%s: %v", ns, err) return 0 } for _, obj := range page.Contents { if obj.Size != nil { totalBytes += *obj.Size } } } return float64(totalBytes) / (1024 * 1024 * 1024) } // ListUsers собирает пользователей из secret fission-console-users в каждом // клиентском namespace. func (k *K8sClient) ListUsers(ctx context.Context) ([]model.UserInfo, error) { nsList, err := k.cs.CoreV1().Namespaces().List(ctx, metav1.ListOptions{ LabelSelector: "managed-by=fission-console", }) if err != nil { return nil, err } var users []model.UserInfo for _, ns := range nsList.Items { // label user содержит email владельца namespace email := ns.Labels["user"] if email == "" { // fallback: попробовать прочитать из secret secret, err := k.cs.CoreV1().Secrets(ns.Name).Get(ctx, "fission-console-users", metav1.GetOptions{}) if err == nil { for key := range secret.Data { if strings.Contains(key, "@") { email = key break } } } } users = append(users, model.UserInfo{ Email: email, Namespace: ns.Name, CreatedAt: ns.CreationTimestamp.Format(time.RFC3339), }) } return users, nil }