@@ -0,0 +1,96 @@
|
||||
// Package demos : cœur métier — cycle de vie des démos client déployées
|
||||
// dans le cluster (une démo = un namespace isolé, TTL 30 jours).
|
||||
package demos
|
||||
|
||||
import "time"
|
||||
|
||||
// Paramètres capacité / durée de vie.
|
||||
const (
|
||||
TTL = 30 * 24 * time.Hour // durée de vie d'une démo
|
||||
MaxConcurrentDemos = 5 // capacité cible
|
||||
PoolSizePerService = 5 // slots par service externe poolé
|
||||
)
|
||||
|
||||
// Services externes gérés par pool pré-provisionné (non auto-créables par API).
|
||||
// TomTom n'y figure pas : ses clés sont créées à la volée par le worker.
|
||||
const (
|
||||
ServiceTelegram = "telegram"
|
||||
ServiceNowPayments = "nowpayments"
|
||||
)
|
||||
|
||||
// PooledServices : services nécessitant un slot de pool par démo.
|
||||
var PooledServices = []string{ServiceTelegram, ServiceNowPayments}
|
||||
|
||||
// Status : état d'une démo.
|
||||
type Status string
|
||||
|
||||
const (
|
||||
StatusPending Status = "pending" // créée, pas encore prise par le worker
|
||||
StatusProvisioning Status = "provisioning" // worker en train de déployer
|
||||
StatusReady Status = "ready" // démo accessible
|
||||
StatusExpiring Status = "expiring" // teardown en cours
|
||||
StatusExpired Status = "expired" // détruite / TTL atteint
|
||||
StatusFailed Status = "failed" // échec de provisioning
|
||||
)
|
||||
|
||||
func (s Status) Valid() bool {
|
||||
switch s {
|
||||
case StatusPending, StatusProvisioning, StatusReady, StatusExpiring, StatusExpired, StatusFailed:
|
||||
return true
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
// Active indique si la démo occupe des ressources (compte dans la capacité).
|
||||
func (s Status) Active() bool {
|
||||
switch s {
|
||||
case StatusPending, StatusProvisioning, StatusReady, StatusExpiring:
|
||||
return true
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
// Extendable : on ne prolonge que des démos encore vivantes.
|
||||
func (s Status) Extendable() bool {
|
||||
switch s {
|
||||
case StatusProvisioning, StatusReady:
|
||||
return true
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
// Demo : modèle GORM d'une démo.
|
||||
type Demo struct {
|
||||
ID string `gorm:"type:uuid;primaryKey" json:"id"`
|
||||
Username string `gorm:"type:varchar(64);index" json:"username,omitempty"`
|
||||
LeadID string `gorm:"type:varchar(36);index" json:"lead_id,omitempty"`
|
||||
Status Status `gorm:"size:20;not null;index" json:"status"`
|
||||
Namespace string `gorm:"size:63;uniqueIndex" json:"namespace"`
|
||||
URL string `gorm:"size:255" json:"url"`
|
||||
TypeAbo string `gorm:"type:varchar(35)" json:"type_abonnement"`
|
||||
CreatedAt time.Time `json:"created_at"`
|
||||
ExpiresAt time.Time `json:"expires_at"`
|
||||
}
|
||||
|
||||
func (Demo) TableName() string { return "demos" }
|
||||
|
||||
// ResourceStatus : état d'un slot de pool.
|
||||
type ResourceStatus string
|
||||
|
||||
const (
|
||||
ResourceFree ResourceStatus = "free"
|
||||
ResourceBorrowed ResourceStatus = "borrowed"
|
||||
)
|
||||
|
||||
// ExternalResource : slot pré-provisionné (table external_pool).
|
||||
// On stocke une *référence* au secret (SecretRef), jamais le secret en clair.
|
||||
type ExternalResource struct {
|
||||
ID uint `gorm:"primaryKey" json:"id"`
|
||||
Service string `gorm:"size:20;not null;index:idx_service_status" json:"service"`
|
||||
Label string `gorm:"size:120" json:"label"` // ex. "@omnex_demo_bot_1"
|
||||
SecretRef string `gorm:"size:200;not null" json:"-"` // clé dans le secret manager
|
||||
Status ResourceStatus `gorm:"size:20;not null;index:idx_service_status" json:"status"`
|
||||
DemoID string `gorm:"type:varchar(36);index" json:"demo_id,omitempty"` // FK optionnelle (vide si libre)
|
||||
}
|
||||
|
||||
func (ExternalResource) TableName() string { return "external_pool" }
|
||||
@@ -0,0 +1,104 @@
|
||||
package demos
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"net/http"
|
||||
|
||||
"github.com/gin-gonic/gin"
|
||||
|
||||
"github.com/omnex/control-plane/api/internal/auth"
|
||||
)
|
||||
|
||||
type Handler struct{ svc *Service }
|
||||
|
||||
func NewHandler(svc *Service) *Handler { return &Handler{svc: svc} }
|
||||
|
||||
type createRequest struct {
|
||||
LeadID string `json:"lead_id" binding:"omitempty,uuid4"`
|
||||
Username string `json:"username" binding:"omitempty,min=3,max=64,alphanum"`
|
||||
}
|
||||
|
||||
// Create : POST /demos — provisionne une démo (202 Accepted, worker asynchrone).
|
||||
func (h *Handler) Create(c *gin.Context) {
|
||||
var req createRequest
|
||||
if err := c.ShouldBindJSON(&req); err != nil {
|
||||
c.JSON(http.StatusBadRequest, gin.H{"error": "requête invalide"})
|
||||
return
|
||||
}
|
||||
d, err := h.svc.Create(req.LeadID, req.Username)
|
||||
if err != nil {
|
||||
if errors.Is(err, ErrCapacityReached) {
|
||||
c.JSON(http.StatusConflict, gin.H{"error": "capacité maximale de démos atteinte"})
|
||||
return
|
||||
}
|
||||
c.JSON(http.StatusInternalServerError, gin.H{"error": "erreur serveur"})
|
||||
return
|
||||
}
|
||||
c.JSON(http.StatusAccepted, d)
|
||||
}
|
||||
|
||||
// List : GET /demos.
|
||||
func (h *Handler) List(c *gin.Context) {
|
||||
items, err := h.svc.List()
|
||||
if err != nil {
|
||||
c.JSON(http.StatusInternalServerError, gin.H{"error": "erreur serveur"})
|
||||
return
|
||||
}
|
||||
c.JSON(http.StatusOK, gin.H{"items": items})
|
||||
}
|
||||
|
||||
// ListMine : GET /demos/mine — démos rattachées au client authentifié.
|
||||
func (h *Handler) ListMine(c *gin.Context) {
|
||||
p := auth.PrincipalFrom(c)
|
||||
if p == nil {
|
||||
c.JSON(http.StatusUnauthorized, gin.H{"error": "non authentifié"})
|
||||
return
|
||||
}
|
||||
items, err := h.svc.ListForUser(p.Username)
|
||||
if err != nil {
|
||||
c.JSON(http.StatusInternalServerError, gin.H{"error": "erreur serveur"})
|
||||
return
|
||||
}
|
||||
c.JSON(http.StatusOK, gin.H{"items": items})
|
||||
}
|
||||
|
||||
// Get : GET /demos/:id.
|
||||
func (h *Handler) Get(c *gin.Context) {
|
||||
d, err := h.svc.Get(c.Param("id"))
|
||||
if err != nil {
|
||||
c.JSON(http.StatusNotFound, gin.H{"error": "démo introuvable"})
|
||||
return
|
||||
}
|
||||
c.JSON(http.StatusOK, d)
|
||||
}
|
||||
|
||||
// Delete : DELETE /demos/:id — teardown + libération du pool.
|
||||
func (h *Handler) Delete(c *gin.Context) {
|
||||
d, err := h.svc.Delete(c.Param("id"))
|
||||
if err != nil {
|
||||
if errors.Is(err, ErrNotFound) {
|
||||
c.JSON(http.StatusNotFound, gin.H{"error": "démo introuvable"})
|
||||
return
|
||||
}
|
||||
c.JSON(http.StatusInternalServerError, gin.H{"error": "erreur serveur"})
|
||||
return
|
||||
}
|
||||
c.JSON(http.StatusOK, d)
|
||||
}
|
||||
|
||||
// Extend : POST /demos/:id/extend — prolonge de 30 jours.
|
||||
func (h *Handler) Extend(c *gin.Context) {
|
||||
d, err := h.svc.Extend(c.Param("id"))
|
||||
if err != nil {
|
||||
switch {
|
||||
case errors.Is(err, ErrNotFound):
|
||||
c.JSON(http.StatusNotFound, gin.H{"error": "démo introuvable"})
|
||||
case errors.Is(err, ErrNotExtendable):
|
||||
c.JSON(http.StatusConflict, gin.H{"error": "démo non prolongeable"})
|
||||
default:
|
||||
c.JSON(http.StatusInternalServerError, gin.H{"error": "erreur serveur"})
|
||||
}
|
||||
return
|
||||
}
|
||||
c.JSON(http.StatusOK, d)
|
||||
}
|
||||
@@ -0,0 +1,337 @@
|
||||
package demos
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"fmt"
|
||||
"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"
|
||||
)
|
||||
|
||||
// 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
|
||||
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")
|
||||
func NewHelmProvisioner(
|
||||
cfg *config.Config,
|
||||
k8sClient *kubernetes.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,
|
||||
chartsDir: chartsDir,
|
||||
frontendImage: cfg.FrontendImage,
|
||||
backendImage: cfg.BackendImage,
|
||||
baseDomain: cfg.DemoDomain,
|
||||
helmPath: helmPath,
|
||||
}, nil
|
||||
}
|
||||
|
||||
// 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
|
||||
func (h *HelmProvisioner) Provision(d Demo, resources []ExternalResource) 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
|
||||
postgresValues := h.buildPostgresValues(d)
|
||||
redisValues := h.buildRedisValues(d)
|
||||
backendValues := h.buildBackendValues(d, resources)
|
||||
frontendValues := h.buildFrontendValues(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)
|
||||
}
|
||||
|
||||
// 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
|
||||
}
|
||||
|
||||
// Teardown détruit une démo en supprimant son namespace (cascading).
|
||||
func (h *HelmProvisioner) Teardown(d Demo) error {
|
||||
return h.deleteNamespace(d.Namespace)
|
||||
}
|
||||
|
||||
// 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 un chart Helm avec des valeurs personnalisées.
|
||||
func (h *HelmProvisioner) installChart(namespace, chartName string, values map[string]interface{}) error {
|
||||
chartPath := filepath.Join(h.chartsDir, chartName)
|
||||
|
||||
// Créer un fichier values.yaml temporaire
|
||||
valuesFile := filepath.Join("/tmp", fmt.Sprintf("values-%s-%s.yaml", namespace, chartName))
|
||||
valuesYAML, err := yaml.Marshal(values)
|
||||
if err != nil {
|
||||
return fmt.Errorf("marshal YAML: %w", err)
|
||||
}
|
||||
if err := os.WriteFile(valuesFile, valuesYAML, 0644); err != nil {
|
||||
return fmt.Errorf("write values file: %w", err)
|
||||
}
|
||||
defer os.Remove(valuesFile)
|
||||
|
||||
// Construire la commande helm install
|
||||
args := []string{
|
||||
"install",
|
||||
fmt.Sprintf("%s-%s", namespace, chartName),
|
||||
chartPath,
|
||||
"--namespace", namespace,
|
||||
"--values", valuesFile,
|
||||
"--wait",
|
||||
"--timeout", "5m",
|
||||
}
|
||||
|
||||
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("helm install %s: %v (stdout: %s, stderr: %s)", chartName, err, stdout.String(), stderr.String())
|
||||
}
|
||||
|
||||
log.Printf("Chart %s installé dans %s", chartName, namespace)
|
||||
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) map[string]interface{} {
|
||||
// Parser l'image backend pour séparer repository et tag
|
||||
backendRepo, backendTag := parseImage(h.backendImage)
|
||||
|
||||
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": false,
|
||||
},
|
||||
},
|
||||
"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",
|
||||
},
|
||||
"secrets": h.buildSecrets(resources),
|
||||
}
|
||||
|
||||
return values
|
||||
}
|
||||
|
||||
// buildFrontendValues construit les valeurs pour le chart frontend.
|
||||
func (h *HelmProvisioner) buildFrontendValues(d Demo) map[string]interface{} {
|
||||
backendURL := fmt.Sprintf("http://%s-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,
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
// 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) 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": false,
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
// buildRedisValues construit les valeurs pour le chart redis.
|
||||
func (h *HelmProvisioner) buildRedisValues(d Demo) 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": false, // Pas de persistence pour les démos
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
// 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
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,145 @@
|
||||
package demos
|
||||
|
||||
import "sync"
|
||||
|
||||
// MemStore : Store en mémoire (dev/tests).
|
||||
type MemStore struct {
|
||||
mu sync.RWMutex
|
||||
items map[string]Demo
|
||||
}
|
||||
|
||||
func NewMemStore() *MemStore { return &MemStore{items: make(map[string]Demo)} }
|
||||
|
||||
func (m *MemStore) Create(d Demo) (Demo, error) {
|
||||
m.mu.Lock()
|
||||
defer m.mu.Unlock()
|
||||
m.items[d.ID] = d
|
||||
return d, nil
|
||||
}
|
||||
|
||||
func (m *MemStore) Get(id string) (Demo, bool) {
|
||||
m.mu.RLock()
|
||||
defer m.mu.RUnlock()
|
||||
d, ok := m.items[id]
|
||||
return d, ok
|
||||
}
|
||||
|
||||
func (m *MemStore) List() ([]Demo, error) {
|
||||
m.mu.RLock()
|
||||
defer m.mu.RUnlock()
|
||||
out := make([]Demo, 0, len(m.items))
|
||||
for _, d := range m.items {
|
||||
out = append(out, d)
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
func (m *MemStore) ListByUsername(username string) ([]Demo, error) {
|
||||
m.mu.RLock()
|
||||
defer m.mu.RUnlock()
|
||||
out := make([]Demo, 0)
|
||||
for _, d := range m.items {
|
||||
if d.Username == username {
|
||||
out = append(out, d)
|
||||
}
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
func (m *MemStore) Update(d Demo) (Demo, error) {
|
||||
m.mu.Lock()
|
||||
defer m.mu.Unlock()
|
||||
m.items[d.ID] = d
|
||||
return d, nil
|
||||
}
|
||||
|
||||
func (m *MemStore) CountActive() (int, error) {
|
||||
m.mu.RLock()
|
||||
defer m.mu.RUnlock()
|
||||
n := 0
|
||||
for _, d := range m.items {
|
||||
if d.Status.Active() {
|
||||
n++
|
||||
}
|
||||
}
|
||||
return n, nil
|
||||
}
|
||||
|
||||
// MemPool : Pool en mémoire (dev/tests), PoolSizePerService slots par service.
|
||||
type MemPool struct {
|
||||
mu sync.Mutex
|
||||
res []*ExternalResource
|
||||
next uint
|
||||
}
|
||||
|
||||
// NewMemPool construit un pool rempli (PoolSizePerService slots libres par service).
|
||||
func NewMemPool() *MemPool {
|
||||
p := &MemPool{}
|
||||
for _, svc := range PooledServices {
|
||||
for i := 1; i <= PoolSizePerService; i++ {
|
||||
p.next++
|
||||
p.res = append(p.res, &ExternalResource{
|
||||
ID: p.next,
|
||||
Service: svc,
|
||||
Label: placeholderLabel(svc, i),
|
||||
SecretRef: "secretref://" + svc + "/slot-" + itoa(i),
|
||||
Status: ResourceFree,
|
||||
})
|
||||
}
|
||||
}
|
||||
return p
|
||||
}
|
||||
|
||||
func (p *MemPool) Borrow(demoID string) ([]ExternalResource, error) {
|
||||
p.mu.Lock()
|
||||
defer p.mu.Unlock()
|
||||
// Vérifie d'abord la dispo de chaque service (atomicité).
|
||||
pick := make(map[string]*ExternalResource, len(PooledServices))
|
||||
for _, svc := range PooledServices {
|
||||
var found *ExternalResource
|
||||
for _, r := range p.res {
|
||||
if r.Service == svc && r.Status == ResourceFree {
|
||||
found = r
|
||||
break
|
||||
}
|
||||
}
|
||||
if found == nil {
|
||||
return nil, ErrPoolExhausted
|
||||
}
|
||||
pick[svc] = found
|
||||
}
|
||||
var out []ExternalResource
|
||||
for _, r := range pick {
|
||||
r.Status = ResourceBorrowed
|
||||
r.DemoID = demoID
|
||||
out = append(out, *r)
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
func (p *MemPool) Return(demoID string) error {
|
||||
p.mu.Lock()
|
||||
defer p.mu.Unlock()
|
||||
for _, r := range p.res {
|
||||
if r.DemoID == demoID {
|
||||
r.Status = ResourceFree
|
||||
r.DemoID = ""
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (p *MemPool) FreeCount() (map[string]int, error) {
|
||||
p.mu.Lock()
|
||||
defer p.mu.Unlock()
|
||||
out := map[string]int{}
|
||||
for _, svc := range PooledServices {
|
||||
out[svc] = 0
|
||||
}
|
||||
for _, r := range p.res {
|
||||
if r.Status == ResourceFree {
|
||||
out[r.Service]++
|
||||
}
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
@@ -0,0 +1,127 @@
|
||||
package demos
|
||||
|
||||
import (
|
||||
"errors"
|
||||
|
||||
"gorm.io/gorm"
|
||||
"gorm.io/gorm/clause"
|
||||
)
|
||||
|
||||
// ErrPoolExhausted : plus de slot libre pour au moins un service requis.
|
||||
var ErrPoolExhausted = errors.New("pool de ressources externes épuisé")
|
||||
|
||||
// Pool : allocation des ressources externes pré-provisionnées.
|
||||
type Pool interface {
|
||||
// Borrow réserve un slot par service requis, de façon atomique.
|
||||
// Rend ErrPoolExhausted si un service n'a plus de slot libre.
|
||||
Borrow(demoID string) ([]ExternalResource, error)
|
||||
// Return libère tous les slots d'une démo.
|
||||
Return(demoID string) error
|
||||
// FreeCount renvoie le nombre de slots libres par service.
|
||||
FreeCount() (map[string]int, error)
|
||||
}
|
||||
|
||||
// GormPool : Pool adossé à PostgreSQL via GORM (verrou transactionnel).
|
||||
type GormPool struct{ db *gorm.DB }
|
||||
|
||||
func NewGormPool(db *gorm.DB) *GormPool { return &GormPool{db: db} }
|
||||
|
||||
func (p *GormPool) Borrow(demoID string) ([]ExternalResource, error) {
|
||||
var borrowed []ExternalResource
|
||||
err := p.db.Transaction(func(tx *gorm.DB) error {
|
||||
for _, svc := range PooledServices {
|
||||
var r ExternalResource
|
||||
// SELECT ... FOR UPDATE SKIP LOCKED : évite la course entre deux démos.
|
||||
err := tx.Clauses(clause.Locking{Strength: "UPDATE", Options: "SKIP LOCKED"}).
|
||||
Where("service = ? AND status = ?", svc, ResourceFree).
|
||||
First(&r).Error
|
||||
if errors.Is(err, gorm.ErrRecordNotFound) {
|
||||
return ErrPoolExhausted
|
||||
}
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
r.Status = ResourceBorrowed
|
||||
r.DemoID = demoID
|
||||
if err := tx.Save(&r).Error; err != nil {
|
||||
return err
|
||||
}
|
||||
borrowed = append(borrowed, r)
|
||||
}
|
||||
return nil
|
||||
})
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return borrowed, nil
|
||||
}
|
||||
|
||||
func (p *GormPool) Return(demoID string) error {
|
||||
return p.db.Model(&ExternalResource{}).
|
||||
Where("demo_id = ?", demoID).
|
||||
Updates(map[string]any{"status": ResourceFree, "demo_id": ""}).Error
|
||||
}
|
||||
|
||||
func (p *GormPool) FreeCount() (map[string]int, error) {
|
||||
type row struct {
|
||||
Service string
|
||||
N int
|
||||
}
|
||||
var rows []row
|
||||
err := p.db.Model(&ExternalResource{}).
|
||||
Select("service, count(*) as n").
|
||||
Where("status = ?", ResourceFree).
|
||||
Group("service").Scan(&rows).Error
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
out := map[string]int{}
|
||||
for _, svc := range PooledServices {
|
||||
out[svc] = 0
|
||||
}
|
||||
for _, r := range rows {
|
||||
out[r.Service] = r.N
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
// SeedPool crée les slots manquants pour atteindre PoolSizePerService par service.
|
||||
// Idempotent : n'ajoute que ce qui manque. SecretRef pointe vers le secret manager.
|
||||
func SeedPool(db *gorm.DB) error {
|
||||
for _, svc := range PooledServices {
|
||||
var n int64
|
||||
if err := db.Model(&ExternalResource{}).Where("service = ?", svc).Count(&n).Error; err != nil {
|
||||
return err
|
||||
}
|
||||
for i := int(n) + 1; i <= PoolSizePerService; i++ {
|
||||
res := ExternalResource{
|
||||
Service: svc,
|
||||
Label: placeholderLabel(svc, i),
|
||||
SecretRef: "secretref://" + svc + "/slot-" + itoa(i),
|
||||
Status: ResourceFree,
|
||||
}
|
||||
if err := db.Create(&res).Error; err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func placeholderLabel(svc string, i int) string {
|
||||
return svc + "-slot-" + itoa(i)
|
||||
}
|
||||
|
||||
func itoa(i int) string {
|
||||
if i == 0 {
|
||||
return "0"
|
||||
}
|
||||
var b [20]byte
|
||||
pos := len(b)
|
||||
for i > 0 {
|
||||
pos--
|
||||
b[pos] = byte('0' + i%10)
|
||||
i /= 10
|
||||
}
|
||||
return string(b[pos:])
|
||||
}
|
||||
@@ -0,0 +1,18 @@
|
||||
package demos
|
||||
|
||||
// Provisioner : abstraction du worker qui déploie/détruit la stack dans le cluster.
|
||||
// L'implémentation réelle (Helm + client-go) sera branchée plus tard ; l'API se
|
||||
// contente d'enfiler les intentions.
|
||||
type Provisioner interface {
|
||||
// Provision demande le déploiement d'une démo (asynchrone).
|
||||
Provision(d Demo, resources []ExternalResource) error
|
||||
// Teardown demande la destruction d'une démo.
|
||||
Teardown(d Demo) error
|
||||
}
|
||||
|
||||
// NoopProvisioner : implémentation neutre (aucun déploiement réel).
|
||||
// Sert au scaffolding et aux tests tant que le worker n'existe pas.
|
||||
type NoopProvisioner struct{}
|
||||
|
||||
func (NoopProvisioner) Provision(Demo, []ExternalResource) error { return nil }
|
||||
func (NoopProvisioner) Teardown(Demo) error { return nil }
|
||||
@@ -0,0 +1,172 @@
|
||||
package demos
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"log"
|
||||
"time"
|
||||
|
||||
"github.com/google/uuid"
|
||||
|
||||
"github.com/omnex/control-plane/api/internal/auth"
|
||||
)
|
||||
|
||||
var (
|
||||
ErrCapacityReached = errors.New("capacité maximale de démos atteinte")
|
||||
ErrNotFound = errors.New("démo introuvable")
|
||||
ErrNotExtendable = errors.New("démo non prolongeable dans son état actuel")
|
||||
)
|
||||
|
||||
// Config du service (injectable pour les tests).
|
||||
type Config struct {
|
||||
BaseDomain string // ex. "demo.omnex.app"
|
||||
TTL time.Duration // durée de vie
|
||||
MaxDemos int // capacité
|
||||
}
|
||||
|
||||
// Service : logique métier des démos (indépendante du transport HTTP).
|
||||
type Service struct {
|
||||
store Store
|
||||
pool Pool
|
||||
prov Provisioner
|
||||
cfg Config
|
||||
now func() time.Time // horloge injectable
|
||||
}
|
||||
|
||||
func NewService(store Store, pool Pool, prov Provisioner, cfg Config) *Service {
|
||||
if cfg.TTL == 0 {
|
||||
cfg.TTL = TTL
|
||||
}
|
||||
if cfg.MaxDemos == 0 {
|
||||
cfg.MaxDemos = MaxConcurrentDemos
|
||||
}
|
||||
if cfg.BaseDomain == "" {
|
||||
cfg.BaseDomain = "demo.omnex.app"
|
||||
}
|
||||
return &Service{store: store, pool: pool, prov: prov, cfg: cfg, now: time.Now}
|
||||
}
|
||||
|
||||
// Create réserve la capacité + le pool, persiste la démo et déclenche le worker.
|
||||
// username (optionnel) rattache la démo à un client existant.
|
||||
func (s *Service) Create(leadID, username string) (Demo, error) {
|
||||
active, err := s.store.CountActive()
|
||||
if err != nil {
|
||||
return Demo{}, err
|
||||
}
|
||||
if active >= s.cfg.MaxDemos {
|
||||
return Demo{}, ErrCapacityReached
|
||||
}
|
||||
|
||||
id := uuid.NewString()
|
||||
demo := Demo{
|
||||
ID: id,
|
||||
Username: auth.NormalizeUsername(username),
|
||||
LeadID: leadID,
|
||||
Status: StatusPending,
|
||||
Namespace: "demo-" + shortID(id),
|
||||
CreatedAt: s.now().UTC(),
|
||||
ExpiresAt: s.now().UTC().Add(s.cfg.TTL),
|
||||
}
|
||||
demo.URL = "https://" + demo.Namespace + "." + s.cfg.BaseDomain
|
||||
|
||||
// Réserve les ressources externes ; ErrPoolExhausted => capacité atteinte.
|
||||
resources, err := s.pool.Borrow(id)
|
||||
if err != nil {
|
||||
if errors.Is(err, ErrPoolExhausted) {
|
||||
return Demo{}, ErrCapacityReached
|
||||
}
|
||||
return Demo{}, err
|
||||
}
|
||||
|
||||
demo.Status = StatusProvisioning
|
||||
created, err := s.store.Create(demo)
|
||||
if err != nil {
|
||||
_ = s.pool.Return(id) // fail-secure : on rend ce qu'on a emprunté
|
||||
return Demo{}, err
|
||||
}
|
||||
|
||||
// Lance le déploiement de manière asynchrone (worker inline via goroutine).
|
||||
go func() {
|
||||
if err := s.prov.Provision(created, resources); err != nil {
|
||||
log.Printf("Provisioning échoué pour %s: %v", created.ID, err)
|
||||
created.Status = StatusFailed
|
||||
_, _ = s.store.Update(created)
|
||||
_ = s.pool.Return(id)
|
||||
return
|
||||
}
|
||||
created.Status = StatusReady
|
||||
_, _ = s.store.Update(created)
|
||||
}()
|
||||
|
||||
return created, nil
|
||||
}
|
||||
|
||||
func (s *Service) Get(id string) (Demo, error) {
|
||||
d, ok := s.store.Get(id)
|
||||
if !ok {
|
||||
return Demo{}, ErrNotFound
|
||||
}
|
||||
return d, nil
|
||||
}
|
||||
|
||||
func (s *Service) List() ([]Demo, error) {
|
||||
return s.store.List()
|
||||
}
|
||||
|
||||
// ListForUser retourne uniquement les démos rattachées à ce client.
|
||||
func (s *Service) ListForUser(username string) ([]Demo, error) {
|
||||
return s.store.ListByUsername(auth.NormalizeUsername(username))
|
||||
}
|
||||
|
||||
// Delete déclenche le teardown et libère le pool.
|
||||
func (s *Service) Delete(id string) (Demo, error) {
|
||||
d, ok := s.store.Get(id)
|
||||
if !ok {
|
||||
return Demo{}, ErrNotFound
|
||||
}
|
||||
if d.Status == StatusExpired {
|
||||
return d, nil // déjà détruite (idempotent)
|
||||
}
|
||||
d.Status = StatusExpiring
|
||||
if _, err := s.store.Update(d); err != nil {
|
||||
return Demo{}, err
|
||||
}
|
||||
if err := s.prov.Teardown(d); err != nil {
|
||||
return Demo{}, err
|
||||
}
|
||||
if err := s.pool.Return(id); err != nil {
|
||||
return Demo{}, err
|
||||
}
|
||||
d.Status = StatusExpired
|
||||
return s.store.Update(d)
|
||||
}
|
||||
|
||||
// Extend prolonge la démo de TTL à partir de son échéance courante.
|
||||
func (s *Service) Extend(id string) (Demo, error) {
|
||||
d, ok := s.store.Get(id)
|
||||
if !ok {
|
||||
return Demo{}, ErrNotFound
|
||||
}
|
||||
if !d.Status.Extendable() {
|
||||
return Demo{}, ErrNotExtendable
|
||||
}
|
||||
base := d.ExpiresAt
|
||||
if now := s.now().UTC(); base.Before(now) {
|
||||
base = now
|
||||
}
|
||||
d.ExpiresAt = base.Add(s.cfg.TTL)
|
||||
return s.store.Update(d)
|
||||
}
|
||||
|
||||
// shortID : 8 premiers caractères hex d'un uuid pour un nom de namespace court.
|
||||
func shortID(id string) string {
|
||||
clean := ""
|
||||
for _, c := range id {
|
||||
if c != '-' {
|
||||
clean += string(c)
|
||||
}
|
||||
if len(clean) == 8 {
|
||||
break
|
||||
}
|
||||
}
|
||||
return clean
|
||||
}
|
||||
@@ -0,0 +1,70 @@
|
||||
package demos
|
||||
|
||||
import (
|
||||
"errors"
|
||||
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
// Store : persistance des démos.
|
||||
type Store interface {
|
||||
Create(d Demo) (Demo, error)
|
||||
Get(id string) (Demo, bool)
|
||||
List() ([]Demo, error)
|
||||
ListByUsername(username string) ([]Demo, error)
|
||||
Update(d Demo) (Demo, error)
|
||||
CountActive() (int, error)
|
||||
}
|
||||
|
||||
// GormStore : Store adossé à PostgreSQL via GORM.
|
||||
type GormStore struct{ db *gorm.DB }
|
||||
|
||||
func NewGormStore(db *gorm.DB) *GormStore { return &GormStore{db: db} }
|
||||
|
||||
func (s *GormStore) Create(d Demo) (Demo, error) {
|
||||
if err := s.db.Create(&d).Error; err != nil {
|
||||
return Demo{}, err
|
||||
}
|
||||
return d, nil
|
||||
}
|
||||
|
||||
func (s *GormStore) Get(id string) (Demo, bool) {
|
||||
var d Demo
|
||||
err := s.db.Where("id = ?", id).First(&d).Error
|
||||
if errors.Is(err, gorm.ErrRecordNotFound) || err != nil {
|
||||
return Demo{}, false
|
||||
}
|
||||
return d, true
|
||||
}
|
||||
|
||||
func (s *GormStore) List() ([]Demo, error) {
|
||||
var out []Demo
|
||||
if err := s.db.Order("created_at DESC").Find(&out).Error; err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
func (s *GormStore) ListByUsername(username string) ([]Demo, error) {
|
||||
var out []Demo
|
||||
if err := s.db.Where("username = ?", username).Order("created_at DESC").Find(&out).Error; err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
func (s *GormStore) Update(d Demo) (Demo, error) {
|
||||
if err := s.db.Save(&d).Error; err != nil {
|
||||
return Demo{}, err
|
||||
}
|
||||
return d, nil
|
||||
}
|
||||
|
||||
// CountActive compte les démos occupant de la capacité.
|
||||
func (s *GormStore) CountActive() (int, error) {
|
||||
var n int64
|
||||
err := s.db.Model(&Demo{}).
|
||||
Where("status IN ?", []Status{StatusPending, StatusProvisioning, StatusReady, StatusExpiring}).
|
||||
Count(&n).Error
|
||||
return int(n), err
|
||||
}
|
||||
Reference in New Issue
Block a user