Files
fission-console/admin-console/internal/metrics/k8s.go
T

255 lines
7.8 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
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-ов БЕЗ CRD-счётчиков (быстро).
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 {
result = append(result, model.NamespaceInfo{
Name: ns.Name,
UserEmail: ns.Labels["user"],
CreatedAt: ns.CreationTimestamp.Format(time.RFC3339),
})
}
return result, nil
}
// GetNamespaceDetail возвращает информацию об одном namespace с CRD-счётчиками.
func (k *K8sClient) GetNamespaceDetail(ctx context.Context, nsName string) (*model.NamespaceInfo, error) {
ns, err := k.cs.CoreV1().Namespaces().Get(ctx, nsName, metav1.GetOptions{})
if err != nil {
return nil, err
}
info := &model.NamespaceInfo{
Name: ns.Name,
UserEmail: ns.Labels["user"],
CreatedAt: ns.CreationTimestamp.Format(time.RFC3339),
}
k.fillCRDCounts(ctx, nsName, info)
return info, 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
}
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 {
email := ns.Labels["user"]
if email == "" {
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
}