package demos import ( "bytes" "context" "fmt" "io" "log" "os" "os/exec" "path/filepath" "strings" "time" "github.com/omnex/control-plane/api/internal/config" "go.yaml.in/yaml/v2" k8sCoreV1 "k8s.io/api/core/v1" k8sErrors "k8s.io/apimachinery/pkg/api/errors" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/client-go/kubernetes" "k8s.io/client-go/kubernetes/scheme" "k8s.io/client-go/rest" "k8s.io/client-go/tools/remotecommand" metricsclient "k8s.io/metrics/pkg/client/clientset/versioned" ) // HelmProvisioner : implémente Provisioner en appelant l'exécutable Helm. // Nécessite que helm soit installé dans le container (ex: dans /usr/local/bin/helm). type HelmProvisioner struct { cfg *config.Config k8sClient *kubernetes.Clientset restConfig *rest.Config // pour l'exec dans les pods (pg_dump/restore) metricsClient *metricsclient.Clientset // optionnel : nil si metrics-server indisponible chartsDir string // chemin vers le dossier des charts (ex: /charts) frontendImage string backendImage string baseDomain string helmPath string // chemin vers l'exécutable helm (default: "helm") } // NewHelmProvisioner crée un nouveau provisioner Helm. // chartsDir : chemin absolu vers le dossier contenant les charts (backend/, frontend/) // helmPath : chemin vers l'exécutable helm (optionnel, default: "helm") // metricsClient : client metrics.k8s.io pour l'usage CPU/mémoire live (optionnel, peut être nil). func NewHelmProvisioner( cfg *config.Config, k8sClient *kubernetes.Clientset, restConfig *rest.Config, metricsClient *metricsclient.Clientset, chartsDir string, helmPath string, ) (*HelmProvisioner, error) { if chartsDir == "" { chartsDir = "/charts" } if helmPath == "" { helmPath = "helm" } // Vérifier que helm est disponible if _, err := exec.LookPath(helmPath); err != nil { return nil, fmt.Errorf("helm exécutable non trouvé à %s: %w", helmPath, err) } return &HelmProvisioner{ cfg: cfg, k8sClient: k8sClient, restConfig: restConfig, metricsClient: metricsClient, chartsDir: chartsDir, frontendImage: cfg.FrontendImage, backendImage: cfg.BackendImage, baseDomain: cfg.DemoDomain, helmPath: helmPath, }, nil } // sharedTraefikNamespace / sharedTraefikRelease : le routeur Traefik + WAF // Coraza est installé UNE SEULE FOIS pour toutes les démos (voir // EnsureSharedInfra), jamais par démo — chaque démo installe seulement le // chart "ingressroute" qui la raccorde à ce Traefik partagé. const ( sharedTraefikNamespace = "traefik" sharedTraefikRelease = "gestion-traefik" ) // Provision déploie une démo avec Helm : // 1. Crée le namespace // 2. Installe PostgreSQL pour la démo // 3. Installe Redis pour la démo // 4. Installe le chart backend // 5. Installe le chart frontend // 6. Raccorde la démo au Traefik partagé (WAF + routage HTTPS du sous-domaine) func (h *HelmProvisioner) Provision(d Demo, resources []ExternalResource, cfg ProvisionConfig) error { // Créer le namespace if err := h.createNamespace(d.Namespace); err != nil { return fmt.Errorf("création namespace %s: %w", d.Namespace, err) } // Préparer les valeurs. Pas de persistence à la création : les démos // d'essai sont éphémères (TTL 30j) — voir UpgradeToPaid pour le passage // en stockage persistant lorsqu'un client devient payant. postgresValues := h.buildPostgresValues(d, false) redisValues := h.buildRedisValues(d, false) backendValues := h.buildBackendValues(d, resources, cfg) frontendValues := h.buildFrontendValues(d) ingressValues := h.buildIngressRouteValues(d) // Installer PostgreSQL (nécessaire pour le backend) if err := h.installChart(d.Namespace, "postgresql", postgresValues); err != nil { h.deleteNamespace(d.Namespace) return fmt.Errorf("déploiement postgresql: %w", err) } // Installer Redis (nécessaire pour le backend) if err := h.installChart(d.Namespace, "redis", redisValues); err != nil { h.deleteNamespace(d.Namespace) return fmt.Errorf("déploiement redis: %w", err) } // Installer le backend if err := h.installChart(d.Namespace, "backend", backendValues); err != nil { h.deleteNamespace(d.Namespace) return fmt.Errorf("déploiement backend: %w", err) } // Installer le frontend if err := h.installChart(d.Namespace, "frontend", frontendValues); err != nil { h.deleteNamespace(d.Namespace) return fmt.Errorf("déploiement frontend: %w", err) } // Raccorder au Traefik partagé : sans ça la démo n'est ni routée ni // protégée par le WAF, même si backend/frontend tournent. if err := h.installChart(d.Namespace, "ingressroute", ingressValues); err != nil { h.deleteNamespace(d.Namespace) return fmt.Errorf("déploiement ingressroute: %w", err) } // Attendre que les pods soient prêts if err := h.waitForRollout(d.Namespace); err != nil { log.Printf("Warning: rollout check échoué pour %s: %v", d.Namespace, err) } return nil } // EnsureSharedInfra installe (ou met à jour, de façon idempotente) le // Traefik partagé + WAF Coraza. À appeler une fois au démarrage de l'API, // avant toute démo — voir cmd/api/main.go. Contrairement à Provision, ceci // n'est PAS répété par démo. func (h *HelmProvisioner) EnsureSharedInfra() error { if err := h.createNamespace(sharedTraefikNamespace); err != nil { return fmt.Errorf("création namespace %s: %w", sharedTraefikNamespace, err) } // Le subchart officiel traefik/traefik est vendorisé dans le repo // (deploy/chart-gestion/traefik/charts/traefik-*.tgz + Chart.lock) : // pas de "helm dependency update" au runtime, qui échouerait de toute // façon puisque /charts est monté en lecture seule. if err := h.upgradeInstallChart(sharedTraefikNamespace, sharedTraefikRelease, "traefik", nil); err != nil { return fmt.Errorf("déploiement traefik partagé: %w", err) } log.Printf("Traefik partagé (WAF Coraza + rate-limit) prêt dans le namespace %s", sharedTraefikNamespace) return nil } // Teardown détruit une démo en supprimant son namespace (cascading). func (h *HelmProvisioner) Teardown(d Demo) error { return h.deleteNamespace(d.Namespace) } // Identifiants postgres des démos — mêmes valeurs codées en dur que // buildPostgresValues/buildBackendValues (voir ces fonctions). const ( demoDBUser = "postgres" demoDBName = "demo_db" demoDBPass = "demo-postgres-pass" ) // UpgradeToPaid bascule postgresql/redis d'une démo en stockage persistant // lorsqu'un client passe en abonnement payant, en préservant les données // postgres déjà créées pendant l'essai : pg_dump avant la bascule (le volume // actuel est un emptyDir éphémère), puis restore une fois le nouveau volume // persistant en place. Redis n'est pas migré (cache/sessions — reconstruit // naturellement), seule sa persistence est activée pour la suite. func (h *HelmProvisioner) UpgradeToPaid(d Demo) error { ctx, cancel := context.WithTimeout(context.Background(), 10*time.Minute) defer cancel() pgPod, err := h.findPod(ctx, d.Namespace, "postgresql") if err != nil { return fmt.Errorf("pod postgresql introuvable: %w", err) } dumpCmd := fmt.Sprintf("PGPASSWORD=%s pg_dump -h localhost -U %s %s", demoDBPass, demoDBUser, demoDBName) dump, stderr, err := h.execInPod(ctx, d.Namespace, pgPod, "postgresql", []string{"sh", "-c", dumpCmd}, nil) if err != nil { return fmt.Errorf("pg_dump: %w (%s)", err, stderr) } // Bascule postgresql en stockage persistant : le pod redémarre (nouveau // PVC, initialement vide) — c'est pour ça que le dump ci-dessus doit // être pris AVANT cet appel. helm upgrade --wait attend que le nouveau // pod soit prêt (readinessProbe = pg_isready) avant de retourner. if err := h.upgradeInstallChart(d.Namespace, d.Namespace+"-postgresql", "postgresql", h.buildPostgresValues(d, true)); err != nil { return fmt.Errorf("upgrade postgresql (persistence): %w", err) } newPgPod, err := h.findPod(ctx, d.Namespace, "postgresql") if err != nil { return fmt.Errorf("pod postgresql introuvable après upgrade: %w", err) } if strings.TrimSpace(dump) != "" { restoreCmd := fmt.Sprintf("PGPASSWORD=%s psql -h localhost -U %s %s", demoDBPass, demoDBUser, demoDBName) _, stderr, err := h.execInPod(ctx, d.Namespace, newPgPod, "postgresql", []string{"sh", "-c", restoreCmd}, strings.NewReader(dump)) if err != nil { return fmt.Errorf("restore pg_dump: %w (%s)", err, stderr) } } if err := h.upgradeInstallChart(d.Namespace, d.Namespace+"-redis", "redis", h.buildRedisValues(d, true)); err != nil { return fmt.Errorf("upgrade redis (persistence): %w", err) } log.Printf("Démo %s basculée en stockage persistant (données postgres migrées)", d.Namespace) return nil } // findPod retourne le nom du premier pod du namespace dont le nom contient nameContains. func (h *HelmProvisioner) findPod(ctx context.Context, namespace, nameContains string) (string, error) { pods, err := h.k8sClient.CoreV1().Pods(namespace).List(ctx, metav1.ListOptions{}) if err != nil { return "", err } for _, p := range pods.Items { if strings.Contains(p.Name, nameContains) { return p.Name, nil } } return "", fmt.Errorf("aucun pod contenant %q dans %s", nameContains, namespace) } // execInPod exécute une commande dans un conteneur et capture stdout/stderr. // stdin peut être nil (aucune entrée à fournir). func (h *HelmProvisioner) execInPod(ctx context.Context, namespace, pod, container string, command []string, stdin io.Reader) (string, string, error) { req := h.k8sClient.CoreV1().RESTClient().Post(). Resource("pods"). Name(pod). Namespace(namespace). SubResource("exec"). VersionedParams(&k8sCoreV1.PodExecOptions{ Container: container, Command: command, Stdin: stdin != nil, Stdout: true, Stderr: true, }, scheme.ParameterCodec) executor, err := remotecommand.NewSPDYExecutor(h.restConfig, "POST", req.URL()) if err != nil { return "", "", fmt.Errorf("exec setup: %w", err) } var stdout, stderrBuf bytes.Buffer err = executor.StreamWithContext(ctx, remotecommand.StreamOptions{ Stdin: stdin, Stdout: &stdout, Stderr: &stderrBuf, }) if err != nil { return stdout.String(), stderrBuf.String(), err } return stdout.String(), stderrBuf.String(), nil } // createNamespace crée un namespace Kubernetes. func (h *HelmProvisioner) createNamespace(name string) error { ns := &k8sCoreV1.Namespace{ ObjectMeta: metav1.ObjectMeta{ Name: name, Labels: map[string]string{ "omnex.app/demo": "true", }, }, } _, err := h.k8sClient.CoreV1().Namespaces().Create(context.Background(), ns, metav1.CreateOptions{}) if err != nil && !k8sErrors.IsAlreadyExists(err) { return err } return nil } // deleteNamespace supprime un namespace (avec cascading). func (h *HelmProvisioner) deleteNamespace(name string) error { propagationPolicy := metav1.DeletePropagationForeground return h.k8sClient.CoreV1().Namespaces().Delete(context.Background(), name, metav1.DeleteOptions{ PropagationPolicy: &propagationPolicy, }) } // installChart installe (première install, jamais réappliqué) un chart Helm // avec des valeurs personnalisées. Release name : "-". // Adapté aux charts par démo (namespace toujours neuf). func (h *HelmProvisioner) installChart(namespace, chartName string, values map[string]interface{}) error { release := fmt.Sprintf("%s-%s", namespace, chartName) valuesFile, cleanup, err := h.writeValuesFile(namespace, chartName, values) if err != nil { return err } defer cleanup() args := []string{ "install", release, filepath.Join(h.chartsDir, chartName), "--namespace", namespace, "--values", valuesFile, "--wait", "--timeout", "5m", } if err := h.runHelm(args...); err != nil { return fmt.Errorf("helm install %s: %w", chartName, err) } log.Printf("Chart %s installé dans %s", chartName, namespace) return nil } // upgradeInstallChart installe ou met à jour un chart de façon idempotente // (helm upgrade --install) — pour l'infra partagée (Traefik), appelée à // chaque démarrage de l'API et qui ne doit pas échouer si déjà en place. func (h *HelmProvisioner) upgradeInstallChart(namespace, release, chartName string, values map[string]interface{}) error { args := []string{ "upgrade", release, filepath.Join(h.chartsDir, chartName), "--install", "--namespace", namespace, "--wait", "--timeout", "5m", } if values != nil { valuesFile, cleanup, err := h.writeValuesFile(namespace, chartName, values) if err != nil { return err } defer cleanup() args = append(args, "--values", valuesFile) } if err := h.runHelm(args...); err != nil { return fmt.Errorf("helm upgrade --install %s: %w", chartName, err) } log.Printf("Chart %s (release %s) à jour dans %s", chartName, release, namespace) return nil } // writeValuesFile sérialise des valeurs Helm dans un fichier temporaire et // retourne son chemin + une fonction de nettoyage à appeler en defer. func (h *HelmProvisioner) writeValuesFile(namespace, chartName string, values map[string]interface{}) (string, func(), error) { valuesFile := filepath.Join("/tmp", fmt.Sprintf("values-%s-%s.yaml", namespace, chartName)) valuesYAML, err := yaml.Marshal(values) if err != nil { return "", nil, fmt.Errorf("marshal YAML: %w", err) } if err := os.WriteFile(valuesFile, valuesYAML, 0644); err != nil { return "", nil, fmt.Errorf("write values file: %w", err) } return valuesFile, func() { os.Remove(valuesFile) }, nil } // runHelm exécute une commande helm et remonte stdout/stderr en cas d'erreur. func (h *HelmProvisioner) runHelm(args ...string) error { cmd := exec.Command(h.helmPath, args...) var stdout, stderr bytes.Buffer cmd.Stdout = &stdout cmd.Stderr = &stderr if err := cmd.Run(); err != nil { return fmt.Errorf("%v (stdout: %s, stderr: %s)", err, stdout.String(), stderr.String()) } return nil } func parseImage(image string) (repo string, tag string) { if image == "" { return "", "helm" } // Chercher le dernier ":" (pour gérer les images avec port dans le repo) idx := strings.LastIndex(image, ":") if idx == -1 { return image, "helm" } return image[:idx], image[idx+1:] } // buildBackendValues construit les valeurs pour le chart backend. func (h *HelmProvisioner) buildBackendValues(d Demo, resources []ExternalResource, cfg ProvisionConfig) map[string]interface{} { // Parser l'image backend pour séparer repository et tag backendRepo, backendTag := parseImage(h.backendImage) storageDriver := cfg.StorageDriver if storageDriver == "" { storageDriver = StorageLocal } // PVC uploads inutile si le stockage est S3 (rien n'est écrit sur le // disque local dans ce cas). uploadsEnabled := storageDriver == StorageLocal env := map[string]string{ "DB_HOST": fmt.Sprintf("%s-postgresql-postgresql", d.Namespace), "DB_PORT": "5432", "DB_PASSWORD": "demo-postgres-pass", "DB_USER": "postgres", "DB_NAME": "demo_db", "DB_SSLMODE": "disable", "REDIS_HOST": fmt.Sprintf("%s-redis-redis", d.Namespace), "REDIS_PORT": "6379", "REDIS_PASSWORD": "demo-redis-pass", // Mot de passe Redis "API_PORT": "8080", "NODE_ENV": "production", "STORAGE_DRIVER": storageDriver, } if storageDriver == StorageS3 { env["S3_BUCKET"] = cfg.S3Bucket env["S3_ENDPOINT"] = cfg.S3Endpoint } secrets := h.buildSecrets(resources) if cfg.TelegramBotToken != "" { secrets["TELEGRAM_BOT_TOKEN"] = cfg.TelegramBotToken } if cfg.NowPaymentsAPIKey != "" { secrets["NOWPAYMENTS_API_KEY"] = cfg.NowPaymentsAPIKey } if cfg.NowPaymentsIPNSecret != "" { // Écrase le secret pré-provisionné par le pool (ServiceNowPayments) : // contrairement au webhook Telegram, l'IPN NowPayments doit // correspondre à ce qui est réellement configuré sur le compte // marchand NowPayments, pas à une valeur générée par le pool. secrets["NOWPAYMENTS_IPN_SECRET"] = cfg.NowPaymentsIPNSecret } values := map[string]interface{}{ "replicaCount": 1, "image": map[string]interface{}{ "repository": backendRepo, "tag": backendTag, "pullPolicy": "Always", }, "service": map[string]interface{}{ "type": "ClusterIP", "port": 8080, }, "autoscaling": map[string]interface{}{ "enabled": false, }, "persistence": map[string]interface{}{ "uploads": map[string]interface{}{ "enabled": uploadsEnabled, }, }, "env": env, "secrets": secrets, } return values } // buildFrontendValues construit les valeurs pour le chart frontend. func (h *HelmProvisioner) buildFrontendValues(d Demo) map[string]interface{} { // Nom réel du Service backend : "-backend" (release helm) + // "-gestion-backend" (nom du chart, ajouté par le helper fullname). backendURL := fmt.Sprintf("http://%s-backend-gestion-backend.%s.svc.cluster.local:8080", d.Namespace, d.Namespace) // Parser l'image frontend pour séparer repository et tag frontendRepo, frontendTag := parseImage(h.frontendImage) return map[string]interface{}{ "replicaCount": 1, "image": map[string]interface{}{ "repository": frontendRepo, "tag": frontendTag, "pullPolicy": "Always", }, "apiUrl": backendURL, "service": map[string]interface{}{ "type": "ClusterIP", "port": 80, }, "autoscaling": map[string]interface{}{ "enabled": false, }, } } // buildIngressRouteValues construit les valeurs pour le chart ingressroute, // qui raccorde cette démo au Traefik partagé (WAF + routage HTTPS). func (h *HelmProvisioner) buildIngressRouteValues(d Demo) map[string]interface{} { return map[string]interface{}{ "host": fmt.Sprintf("%s.%s", d.Namespace, h.baseDomain), "traefikNamespace": sharedTraefikNamespace, } } // buildSecrets extrait les références des secrets des ressources externes. func (h *HelmProvisioner) buildSecrets(resources []ExternalResource) map[string]string { secrets := map[string]string{} for _, r := range resources { switch r.Service { case ServiceTelegram: secrets["TELEGRAM_WEBHOOK_SECRET"] = r.SecretRef case ServiceNowPayments: secrets["NOWPAYMENTS_IPN_SECRET"] = r.SecretRef } } return secrets } // buildPostgresValues construit les valeurs pour le chart postgresql. func (h *HelmProvisioner) buildPostgresValues(d Demo, persistent bool) map[string]interface{} { return map[string]interface{}{ "auth": map[string]interface{}{ "password": "demo-postgres-pass", // Mot de passe OBLIGATOIRE (champ correct pour le chart) "username": "postgres", "database": "demo_db", }, "service": map[string]interface{}{ "type": "ClusterIP", "port": 5432, }, "persistence": map[string]interface{}{ "enabled": persistent, }, } } // buildRedisValues construit les valeurs pour le chart redis. func (h *HelmProvisioner) buildRedisValues(d Demo, persistent bool) map[string]interface{} { return map[string]interface{}{ "service": map[string]interface{}{ "type": "ClusterIP", "port": 6379, }, "auth": map[string]interface{}{ "password": "demo-redis-pass", // Mot de passe simple pour les démos }, "persistence": map[string]interface{}{ "enabled": persistent, }, } } // waitForRollout attend que tous les déploiements dans le namespace soient prêts. func (h *HelmProvisioner) waitForRollout(namespace string) error { ctx, cancel := context.WithTimeout(context.Background(), 5*time.Minute) defer cancel() for { select { case <-ctx.Done(): return ctx.Err() case <-time.After(5 * time.Second): // Utiliser kubectl pour vérifier le rollout cmd := exec.Command("kubectl", "rollout", "status", "deployment", "-n", namespace, "--timeout=30s") var stderr bytes.Buffer cmd.Stderr = &stderr if err := cmd.Run(); err != nil { // Si c'est juste un timeout, on continue d'attendre if strings.Contains(stderr.String(), "timed out") { continue } return fmt.Errorf("rollout check: %v: %s", err, stderr.String()) } // Si la commande réussit, tous les déploiements sont prêts return nil } } } // componentResourceLimits : limites CPU/mémoire configurées dans chaque chart // (deploy/chart-gestion/*/values.yaml), utilisées pour calculer un % d'usage. var componentResourceLimits = map[string]struct { cpuMilli int64 memMi int64 }{ "backend": {cpuMilli: 500, memMi: 256}, "frontend": {cpuMilli: 200, memMi: 128}, "postgresql": {cpuMilli: 1000, memMi: 1024}, "redis": {cpuMilli: 500, memMi: 512}, } // componentKey identifie le composant (backend/frontend/postgresql/redis) à // partir du nom de pod généré par Helm (ex: "demo-xxx-backend-...-6d59f4-abcde"). func componentKey(podName string) string { switch { case strings.Contains(podName, "backend"): return "backend" case strings.Contains(podName, "frontend"): return "frontend" case strings.Contains(podName, "postgresql"): return "postgresql" case strings.Contains(podName, "redis"): return "redis" default: return "" } } // GetResourceState : phase + usage CPU/mémoire live de chaque pod d'une démo. // L'usage (metrics-server) est best-effort : s'il est indisponible, seule la // phase du pod est renseignée plutôt que de faire échouer tout l'appel. func (h *HelmProvisioner) GetResourceState( ctx context.Context, namespace string, ) (ResourceState, error) { pods, err := h.k8sClient. CoreV1(). Pods(namespace). List(ctx, metav1.ListOptions{}) if err != nil { return ResourceState{}, err } usage := map[string][2]int64{} // pod name -> [cpuMilli, memMi] if h.metricsClient != nil { if list, err := h.metricsClient.MetricsV1beta1().PodMetricses(namespace).List(ctx, metav1.ListOptions{}); err == nil { for _, m := range list.Items { var cpu, mem int64 for _, c := range m.Containers { cpu += c.Usage.Cpu().MilliValue() mem += c.Usage.Memory().Value() / (1024 * 1024) } usage[m.Name] = [2]int64{cpu, mem} } } } var state ResourceState for _, pod := range pods.Items { key := componentKey(pod.Name) if key == "" { continue } cs := ComponentState{Phase: string(pod.Status.Phase)} if limits, ok := componentResourceLimits[key]; ok { cs.CPULimitMilli = limits.cpuMilli cs.MemoryLimitMi = limits.memMi } if u, ok := usage[pod.Name]; ok { cs.CPUMilli = u[0] cs.MemoryMi = u[1] } switch key { case "backend": state.API = cs case "frontend": state.Web = cs case "postgresql": state.DB = cs case "redis": state.DBM = cs } } return state, nil }