// Изменён: 2026-04-08 // queueservice_controller.go — Reconciler для QueueService CRD. // Управляет жизненным циклом ElasticMQ инстанса: создание/удаление всех k8s ресурсов. // // Фазы: Pending → Provisioning → Ready / Failed // // Ресурсы на тенанта (в namespace sless-fn-{tenantId}): // - Secret sqs-creds-{tenantId} — accessKey/secretKey // - ConfigMap sqs-cfg-{tenantId} — elasticmq.conf // - PVC sqs-data-{tenantId} — H2 persistence (не удаляется при удалении QS) // - Deployment sqs-{tenantId} — ElasticMQ Native pod // - Service sqs-svc-{tenantId} — ClusterIP :9324 // - Ingress sqs-ing-{tenantId} — sqs.kube5s.ru/sqs/{tenantId} → Service package controller import ( "context" "fmt" "time" appsv1 "k8s.io/api/apps/v1" corev1 "k8s.io/api/core/v1" netv1 "k8s.io/api/networking/v1" "k8s.io/apimachinery/pkg/api/errors" "k8s.io/apimachinery/pkg/api/resource" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/runtime" "k8s.io/apimachinery/pkg/util/intstr" ctrl "sigs.k8s.io/controller-runtime" "sigs.k8s.io/controller-runtime/pkg/client" "sigs.k8s.io/controller-runtime/pkg/log" sqsv1alpha1 "gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sqs-operator/api/v1alpha1" "gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sqs-operator/internal/elasticmq" ) const ( // sqsFinalizer — добавляется к QueueService CR для управления cleanup при удалении. sqsFinalizer = "sqs.kube5s.ru/sqs-finalizer" // elasticMQImage — образ ElasticMQ JVM (полный, с поддержкой H2 JDBC persistence). // Native-образ (elasticmq-native) не включает H2 в GraalVM native image — persistence не работает. elasticMQImage = "softwaremill/elasticmq:1.7.1" // elasticMQUIImage — Web UI для мониторинга очередей (появился в v1.7.0). // Запускается как sidecar, обращается к ElasticMQ через localhost:9324. elasticMQUIImage = "softwaremill/elasticmq-ui:latest" // elasticMQPort — порт на котором ElasticMQ слушает SQS HTTP запросы. elasticMQPort = 9324 // elasticMQUIPort — порт UI контейнера. elasticMQUIPort = 3000 ) // QueueServiceReconciler управляет QueueService CRD. type QueueServiceReconciler struct { client.Client Scheme *runtime.Scheme SQSExternalHost string // из env SQS_EXTERNAL_HOST, например sqs.kube5s.ru ElasticMQImage string // переопределяемый образ (для тестов / обновлений) } // +kubebuilder:rbac:groups=sqs.kube5s.ru,resources=queueservices,verbs=get;list;watch;create;update;patch;delete // +kubebuilder:rbac:groups=sqs.kube5s.ru,resources=queueservices/status,verbs=get;update;patch // +kubebuilder:rbac:groups=sqs.kube5s.ru,resources=queueservices/finalizers,verbs=update // +kubebuilder:rbac:groups=apps,resources=deployments,verbs=get;list;watch;create;update;patch;delete // +kubebuilder:rbac:groups="",resources=services,verbs=get;list;watch;create;update;patch;delete // +kubebuilder:rbac:groups="",resources=configmaps,verbs=get;list;watch;create;update;patch;delete // +kubebuilder:rbac:groups="",resources=secrets,verbs=get;list;watch;create;update;patch;delete // +kubebuilder:rbac:groups="",resources=persistentvolumeclaims,verbs=get;list;watch;create;update;patch;delete // +kubebuilder:rbac:groups="",resources=namespaces,verbs=get;list;watch;create // +kubebuilder:rbac:groups=networking.k8s.io,resources=ingresses,verbs=get;list;watch;create;update;patch;delete // Reconcile — главный цикл управления QueueService. func (r *QueueServiceReconciler) Reconcile(ctx context.Context, req ctrl.Request) (ctrl.Result, error) { logger := log.FromContext(ctx) qs := &sqsv1alpha1.QueueService{} if err := r.Get(ctx, req.NamespacedName, qs); err != nil { if errors.IsNotFound(err) { return ctrl.Result{}, nil } return ctrl.Result{}, fmt.Errorf("get QueueService: %w", err) } // Если объект удаляется — очистить ресурсы и снять finalizer if !qs.DeletionTimestamp.IsZero() { return r.handleDeletion(ctx, qs) } // Добавляем finalizer при первом создании if !containsString(qs.Finalizers, sqsFinalizer) { qs.Finalizers = append(qs.Finalizers, sqsFinalizer) if err := r.Update(ctx, qs); err != nil { return ctrl.Result{}, fmt.Errorf("add finalizer: %w", err) } return ctrl.Result{Requeue: true}, nil } tenantNS := "sless-fn-" + qs.Spec.TenantID switch qs.Status.Phase { case "", sqsv1alpha1.QueueServicePhasePending: return r.provision(ctx, qs, tenantNS) case sqsv1alpha1.QueueServicePhaseProvisioning: return r.checkReady(ctx, qs, tenantNS) case sqsv1alpha1.QueueServicePhaseReady: return r.ensureHealthy(ctx, qs, tenantNS) case sqsv1alpha1.QueueServicePhaseFailed: return r.recoverFromFailed(ctx, qs, tenantNS) } logger.Info("unknown phase, skipping", "phase", qs.Status.Phase) return ctrl.Result{}, nil } // provision создаёт все k8s ресурсы для тенанта (фаза Pending → Provisioning). func (r *QueueServiceReconciler) provision(ctx context.Context, qs *sqsv1alpha1.QueueService, tenantNS string) (ctrl.Result, error) { logger := log.FromContext(ctx) tenantID := qs.Spec.TenantID logger.Info("provisioning QueueService", "tenant", tenantID, "namespace", tenantNS) // 1. Убедиться что namespace существует if err := r.ensureNamespace(ctx, tenantNS, qs); err != nil { return r.setFailed(ctx, qs, fmt.Sprintf("ensure namespace: %v", err)) } // 2. Secret с credentials (если ещё нет) secretName := "sqs-creds-" + tenantID if err := r.ensureSecret(ctx, qs, tenantNS, secretName); err != nil { return r.setFailed(ctx, qs, fmt.Sprintf("ensure secret: %v", err)) } // 3. ConfigMap с elasticmq.conf cfgMapName := "sqs-cfg-" + tenantID if err := r.ensureConfigMap(ctx, qs, tenantNS, cfgMapName); err != nil { return r.setFailed(ctx, qs, fmt.Sprintf("ensure configmap: %v", err)) } // 4. PVC для H2 persistence pvcName := "sqs-data-" + tenantID if err := r.ensurePVC(ctx, qs, tenantNS, pvcName); err != nil { return r.setFailed(ctx, qs, fmt.Sprintf("ensure pvc: %v", err)) } // 5. Deployment ElasticMQ deployName := "sqs-" + tenantID if err := r.ensureDeployment(ctx, qs, tenantNS, deployName, cfgMapName, pvcName); err != nil { return r.setFailed(ctx, qs, fmt.Sprintf("ensure deployment: %v", err)) } // 6. Service (ClusterIP) svcName := "sqs-svc-" + tenantID if err := r.ensureService(ctx, qs, tenantNS, svcName, deployName); err != nil { return r.setFailed(ctx, qs, fmt.Sprintf("ensure service: %v", err)) } // 7. Ingress (SQS API) ingName := "sqs-ing-" + tenantID if err := r.ensureIngress(ctx, qs, tenantNS, ingName, svcName); err != nil { return r.setFailed(ctx, qs, fmt.Sprintf("ensure ingress: %v", err)) } // 8. Ingress для Web UI (если EnableUI) if qs.Spec.EnableUI { uiIngName := "sqs-ing-ui-" + tenantID if err := r.ensureIngressUI(ctx, qs, tenantNS, uiIngName, svcName); err != nil { return r.setFailed(ctx, qs, fmt.Sprintf("ensure ingress-ui: %v", err)) } // Отдельный ingress для /_next/ ресурсов Next.js (без rewrite — basePath="" в образе). assetsIngName := "sqs-ing-ui-assets-" + tenantID if err := r.ensureIngressUIAssets(ctx, qs, tenantNS, assetsIngName, svcName); err != nil { return r.setFailed(ctx, qs, fmt.Sprintf("ensure ingress-ui-assets: %v", err)) } // Ingress для внутренних маршрутов Next.js (/queues/[name]). // Next.js собран с basePath="" — browser navigates /queues/xxx напрямую без /sqs-ui/tenantID/ префикса. queuesIngName := "sqs-ing-ui-queues-" + tenantID if err := r.ensureIngressUIQueues(ctx, qs, tenantNS, queuesIngName, svcName); err != nil { return r.setFailed(ctx, qs, fmt.Sprintf("ensure ingress-ui-queues: %v", err)) } } // Все ресурсы созданы — переходим в Provisioning, ждём готовности pod qs.Status.Phase = sqsv1alpha1.QueueServicePhaseProvisioning qs.Status.Message = "resources created, waiting for pod ready" qs.Status.SecretName = secretName if err := r.Status().Update(ctx, qs); err != nil { return ctrl.Result{}, fmt.Errorf("update status provisioning: %w", err) } // Проверяем через 3 секунды return ctrl.Result{RequeueAfter: 3 * time.Second}, nil } // checkReady проверяет готовность Deployment (фаза Provisioning → Ready). func (r *QueueServiceReconciler) checkReady(ctx context.Context, qs *sqsv1alpha1.QueueService, tenantNS string) (ctrl.Result, error) { deployName := "sqs-" + qs.Spec.TenantID deploy := &appsv1.Deployment{} if err := r.Get(ctx, client.ObjectKey{Namespace: tenantNS, Name: deployName}, deploy); err != nil { if errors.IsNotFound(err) { // Deployment пропал — возвращаемся в Pending для пересоздания qs.Status.Phase = sqsv1alpha1.QueueServicePhasePending qs.Status.Message = "deployment not found, reprovisioning" _ = r.Status().Update(ctx, qs) return ctrl.Result{Requeue: true}, nil } return ctrl.Result{}, fmt.Errorf("get deployment: %w", err) } if deploy.Status.AvailableReplicas < 1 { // Ещё не готов — ждём return ctrl.Result{RequeueAfter: 3 * time.Second}, nil } // Pod готов — переходим в Ready now := metav1.Now() endpoint := fmt.Sprintf("https://%s/sqs/%s", r.SQSExternalHost, qs.Spec.TenantID) qs.Status.Phase = sqsv1alpha1.QueueServicePhaseReady qs.Status.Endpoint = endpoint qs.Status.Message = "" qs.Status.ReadyAt = &now if err := r.Status().Update(ctx, qs); err != nil { return ctrl.Result{}, fmt.Errorf("update status ready: %w", err) } log.FromContext(ctx).Info("QueueService ready", "tenant", qs.Spec.TenantID, "endpoint", endpoint) return ctrl.Result{}, nil } // ensureHealthy мониторит состояние готового инстанса (фаза Ready). // Проверяет наличие всех критических ресурсов (Deployment, Service, ConfigMap, Ingress). // Если любой ресурс пропал — переход в Pending для пересоздания (OwnerReference не работает cross-namespace). // Если pod упал — переходим в Failed для последующего восстановления. func (r *QueueServiceReconciler) ensureHealthy(ctx context.Context, qs *sqsv1alpha1.QueueService, tenantNS string) (ctrl.Result, error) { tenantID := qs.Spec.TenantID // Проверяем критические ресурсы — если удалены вручную, пересоздаём через provision. checkResources := []struct { name string obj client.Object }{ {"sqs-" + tenantID, &appsv1.Deployment{}}, {"sqs-svc-" + tenantID, &corev1.Service{}}, {"sqs-cfg-" + tenantID, &corev1.ConfigMap{}}, {"sqs-ing-" + tenantID, &netv1.Ingress{}}, } if qs.Spec.EnableUI { checkResources = append(checkResources, struct { name string obj client.Object }{"sqs-ing-ui-" + tenantID, &netv1.Ingress{}}) checkResources = append(checkResources, struct { name string obj client.Object }{"sqs-ing-ui-assets-" + tenantID, &netv1.Ingress{}}) checkResources = append(checkResources, struct { name string obj client.Object }{"sqs-ing-ui-queues-" + tenantID, &netv1.Ingress{}}) } for _, res := range checkResources { if err := r.Get(ctx, client.ObjectKey{Namespace: tenantNS, Name: res.name}, res.obj); err != nil { if errors.IsNotFound(err) { log.FromContext(ctx).Info("resource disappeared, reprovisioning", "resource", res.name) qs.Status.Phase = sqsv1alpha1.QueueServicePhasePending qs.Status.Message = fmt.Sprintf("%s disappeared, reprovisioning", res.name) _ = r.Status().Update(ctx, qs) return ctrl.Result{Requeue: true}, nil } return ctrl.Result{}, fmt.Errorf("health check get %s: %w", res.name, err) } } // Deployment существует (проверен в цикле выше) — проверяем готовность pod deploy := &appsv1.Deployment{} _ = r.Get(ctx, client.ObjectKey{Namespace: tenantNS, Name: "sqs-" + tenantID}, deploy) if deploy.Status.AvailableReplicas < 1 { qs.Status.Phase = sqsv1alpha1.QueueServicePhaseFailed qs.Status.Message = "pod unavailable" _ = r.Status().Update(ctx, qs) return ctrl.Result{RequeueAfter: 10 * time.Second}, nil } // Всё хорошо — следующий check через 15 секунд return ctrl.Result{RequeueAfter: 15 * time.Second}, nil } // recoverFromFailed пробует восстановиться из Failed состояния. func (r *QueueServiceReconciler) recoverFromFailed(ctx context.Context, qs *sqsv1alpha1.QueueService, tenantNS string) (ctrl.Result, error) { deployName := "sqs-" + qs.Spec.TenantID deploy := &appsv1.Deployment{} if err := r.Get(ctx, client.ObjectKey{Namespace: tenantNS, Name: deployName}, deploy); err != nil { if errors.IsNotFound(err) { // Возвращаемся в Pending — пересоздадим всё qs.Status.Phase = sqsv1alpha1.QueueServicePhasePending qs.Status.Message = "recovering: reprovisioning" _ = r.Status().Update(ctx, qs) return ctrl.Result{Requeue: true}, nil } return ctrl.Result{RequeueAfter: 30 * time.Second}, nil } if deploy.Status.AvailableReplicas >= 1 { // Pod восстановился now := metav1.Now() qs.Status.Phase = sqsv1alpha1.QueueServicePhaseReady qs.Status.Message = "" qs.Status.ReadyAt = &now _ = r.Status().Update(ctx, qs) return ctrl.Result{}, nil } return ctrl.Result{RequeueAfter: 30 * time.Second}, nil } // handleDeletion удаляет все ресурсы тенанта кроме PVC (данные сохраняются). func (r *QueueServiceReconciler) handleDeletion(ctx context.Context, qs *sqsv1alpha1.QueueService) (ctrl.Result, error) { logger := log.FromContext(ctx) tenantID := qs.Spec.TenantID tenantNS := "sless-fn-" + tenantID logger.Info("deleting QueueService resources", "tenant", tenantID) deleteOpts := []client.DeleteOption{client.PropagationPolicy(metav1.DeletePropagationForeground)} // Удаляем Ingress UI queues (/queues/*) если был создан uiQueuesIng := &netv1.Ingress{} if err := r.Get(ctx, client.ObjectKey{Namespace: tenantNS, Name: "sqs-ing-ui-queues-" + tenantID}, uiQueuesIng); err == nil { _ = r.Delete(ctx, uiQueuesIng, deleteOpts...) } // Удаляем Ingress UI assets (/_next/) если был создан uiAssetsIng := &netv1.Ingress{} if err := r.Get(ctx, client.ObjectKey{Namespace: tenantNS, Name: "sqs-ing-ui-assets-" + tenantID}, uiAssetsIng); err == nil { _ = r.Delete(ctx, uiAssetsIng, deleteOpts...) } // Удаляем Ingress UI (если был создан) uiIng := &netv1.Ingress{} if err := r.Get(ctx, client.ObjectKey{Namespace: tenantNS, Name: "sqs-ing-ui-" + tenantID}, uiIng); err == nil { _ = r.Delete(ctx, uiIng, deleteOpts...) } // Удаляем Ingress SQS API ing := &netv1.Ingress{} if err := r.Get(ctx, client.ObjectKey{Namespace: tenantNS, Name: "sqs-ing-" + tenantID}, ing); err == nil { _ = r.Delete(ctx, ing, deleteOpts...) } // Удаляем Service svc := &corev1.Service{} if err := r.Get(ctx, client.ObjectKey{Namespace: tenantNS, Name: "sqs-svc-" + tenantID}, svc); err == nil { _ = r.Delete(ctx, svc, deleteOpts...) } // Удаляем Deployment deploy := &appsv1.Deployment{} if err := r.Get(ctx, client.ObjectKey{Namespace: tenantNS, Name: "sqs-" + tenantID}, deploy); err == nil { _ = r.Delete(ctx, deploy, deleteOpts...) } // Удаляем ConfigMap cm := &corev1.ConfigMap{} if err := r.Get(ctx, client.ObjectKey{Namespace: tenantNS, Name: "sqs-cfg-" + tenantID}, cm); err == nil { _ = r.Delete(ctx, cm, deleteOpts...) } // Удаляем Secret с credentials secret := &corev1.Secret{} if err := r.Get(ctx, client.ObjectKey{Namespace: tenantNS, Name: "sqs-creds-" + tenantID}, secret); err == nil { _ = r.Delete(ctx, secret, deleteOpts...) } // PVC НЕ удаляем — данные тенанта сохраняются для возможного восстановления. // Удаление PVC — отдельная административная операция. // Снимаем finalizer qs.Finalizers = removeString(qs.Finalizers, sqsFinalizer) if err := r.Update(ctx, qs); err != nil { return ctrl.Result{}, fmt.Errorf("remove finalizer: %w", err) } logger.Info("QueueService deleted", "tenant", tenantID) return ctrl.Result{}, nil } // ensureNamespace создаёт namespace если не существует. func (r *QueueServiceReconciler) ensureNamespace(ctx context.Context, ns string, qs *sqsv1alpha1.QueueService) error { namespace := &corev1.Namespace{} err := r.Get(ctx, client.ObjectKey{Name: ns}, namespace) if err == nil { return nil // уже существует } if !errors.IsNotFound(err) { return fmt.Errorf("get namespace: %w", err) } namespace = &corev1.Namespace{ ObjectMeta: metav1.ObjectMeta{ Name: ns, Labels: map[string]string{ "sqs.kube5s.ru/managed-by": "sqs-operator", "sqs.kube5s.ru/tenant": qs.Spec.TenantID, }, }, } return r.Create(ctx, namespace) } // ensureSecret создаёт Secret с credentials если не существует. func (r *QueueServiceReconciler) ensureSecret(ctx context.Context, qs *sqsv1alpha1.QueueService, ns, name string) error { secret := &corev1.Secret{} if err := r.Get(ctx, client.ObjectKey{Namespace: ns, Name: name}, secret); err == nil { return nil // уже существует, не перезаписываем } else if !errors.IsNotFound(err) { return fmt.Errorf("get secret: %w", err) } // Генерируем новые credentials creds, err := elasticmq.GenerateCredentials(qs.Spec.TenantID) if err != nil { return fmt.Errorf("generate credentials: %w", err) } secret = &corev1.Secret{ ObjectMeta: metav1.ObjectMeta{ Name: name, Namespace: ns, Labels: sqsLabels(qs.Spec.TenantID), }, Type: corev1.SecretTypeOpaque, StringData: map[string]string{ "accessKey": creds.AccessKey, "secretKey": creds.SecretKey, }, } return r.Create(ctx, secret) } // ensureConfigMap создаёт ConfigMap с elasticmq.conf если не существует. func (r *QueueServiceReconciler) ensureConfigMap(ctx context.Context, qs *sqsv1alpha1.QueueService, ns, name string) error { cm := &corev1.ConfigMap{} if err := r.Get(ctx, client.ObjectKey{Namespace: ns, Name: name}, cm); err == nil { return nil } else if !errors.IsNotFound(err) { return fmt.Errorf("get configmap: %w", err) } conf := elasticmq.GenerateConfig(qs.Spec.TenantID, r.SQSExternalHost, qs.Spec.Persistence, qs.Spec.AutoCreateQueues) cm = &corev1.ConfigMap{ ObjectMeta: metav1.ObjectMeta{ Name: name, Namespace: ns, Labels: sqsLabels(qs.Spec.TenantID), }, Data: map[string]string{ // Ключ соответствует subPath при монтировании в pod. "elasticmq.conf": conf, }, } return r.Create(ctx, cm) } // ensurePVC создаёт PVC для H2 persistence если не существует. func (r *QueueServiceReconciler) ensurePVC(ctx context.Context, qs *sqsv1alpha1.QueueService, ns, name string) error { pvc := &corev1.PersistentVolumeClaim{} if err := r.Get(ctx, client.ObjectKey{Namespace: ns, Name: name}, pvc); err == nil { return nil } else if !errors.IsNotFound(err) { return fmt.Errorf("get pvc: %w", err) } storageSize := resource.MustParse(fmt.Sprintf("%dMi", qs.Spec.StorageMB)) pvc = &corev1.PersistentVolumeClaim{ ObjectMeta: metav1.ObjectMeta{ Name: name, Namespace: ns, Labels: sqsLabels(qs.Spec.TenantID), }, Spec: corev1.PersistentVolumeClaimSpec{ AccessModes: []corev1.PersistentVolumeAccessMode{corev1.ReadWriteOnce}, Resources: corev1.VolumeResourceRequirements{ Requests: corev1.ResourceList{ corev1.ResourceStorage: storageSize, }, }, }, } return r.Create(ctx, pvc) } // ensureDeployment создаёт Deployment с ElasticMQ Native если не существует. func (r *QueueServiceReconciler) ensureDeployment(ctx context.Context, qs *sqsv1alpha1.QueueService, ns, name, cfgMapName, pvcName string) error { deploy := &appsv1.Deployment{} if err := r.Get(ctx, client.ObjectKey{Namespace: ns, Name: name}, deploy); err == nil { return nil } else if !errors.IsNotFound(err) { return fmt.Errorf("get deployment: %w", err) } image := r.elasticMQImageName() replicas := int32(1) port := int32(elasticMQPort) // JVM-образ ElasticMQ требует минимум ~150MB. Форсируем нижнюю границу 256Mi. // MemoryMB из spec — лимит; request = половина лимита, но не меньше 128Mi. memMB := qs.Spec.MemoryMB if memMB < 256 { memMB = 256 } memLimit := resource.MustParse(fmt.Sprintf("%dMi", memMB)) memReqMB := memMB / 2 if memReqMB < 128 { memReqMB = 128 } memRequest := resource.MustParse(fmt.Sprintf("%dMi", memReqMB)) // -Xmx = 75% от лимита, чтобы JVM не выбил OOMKill при GC overhead. xmxMB := memMB * 3 / 4 jvmOpts := fmt.Sprintf("-Xmx%dm -Xms64m", xmxMB) cpuRequest := resource.MustParse("10m") cpuLimit := resource.MustParse("500m") deploy = &appsv1.Deployment{ ObjectMeta: metav1.ObjectMeta{ Name: name, Namespace: ns, Labels: sqsLabels(qs.Spec.TenantID), }, Spec: appsv1.DeploymentSpec{ Replicas: &replicas, Selector: &metav1.LabelSelector{ MatchLabels: map[string]string{ "sqs.kube5s.ru/tenant": qs.Spec.TenantID, }, }, // Recreate: старый pod убивается ДО создания нового. // H2 MVStore держит FileChannel.lock() на /data/elasticmq.mv.db — // при RollingUpdate два pod лезут в один PVC одновременно = ERR-SQS-06. Strategy: appsv1.DeploymentStrategy{ Type: appsv1.RecreateDeploymentStrategyType, }, Template: corev1.PodTemplateSpec{ ObjectMeta: metav1.ObjectMeta{ Labels: sqsLabels(qs.Spec.TenantID), }, Spec: corev1.PodSpec{ // fsGroup=999 — uid elasticmq в JVM-образе. // Kubernetes при монтировании PVC сделает chown :999 на /data. SecurityContext: &corev1.PodSecurityContext{ FSGroup: func() *int64 { v := int64(999); return &v }(), }, // 15 сек — достаточно для graceful shutdown JVM + H2 fsync. TerminationGracePeriodSeconds: func() *int64 { v := int64(15); return &v }(), Containers: func() []corev1.Container { ctrs := []corev1.Container{ { Name: "elasticmq", Image: image, Ports: []corev1.ContainerPort{ {ContainerPort: port, Protocol: corev1.ProtocolTCP}, }, // JVM-образ имеет ENTRYPOINT: java -Dconfig.file=/opt/elasticmq.conf -jar ... // JAVA_TOOL_OPTIONS задаёт Xmx чтобы JVM не выбил OOMKill ниже k8s лимита. Env: []corev1.EnvVar{ { Name: "JAVA_TOOL_OPTIONS", Value: jvmOpts, }, }, Resources: corev1.ResourceRequirements{ Requests: corev1.ResourceList{ corev1.ResourceCPU: cpuRequest, corev1.ResourceMemory: memRequest, }, Limits: corev1.ResourceList{ corev1.ResourceCPU: cpuLimit, corev1.ResourceMemory: memLimit, }, }, VolumeMounts: []corev1.VolumeMount{ { // Монтируем конфиг на путь который JVM-образ читает по дефолту. Name: "elasticmq-config", MountPath: "/opt/elasticmq.conf", SubPath: "elasticmq.conf", ReadOnly: true, }, { // PVC для H2 persistence Name: "elasticmq-data", MountPath: "/data", }, }, ReadinessProbe: &corev1.Probe{ ProbeHandler: corev1.ProbeHandler{ HTTPGet: &corev1.HTTPGetAction{ Path: "/health", Port: intstr.FromInt(int(port)), Scheme: corev1.URISchemeHTTP, }, }, InitialDelaySeconds: 2, PeriodSeconds: 5, FailureThreshold: 3, }, LivenessProbe: &corev1.Probe{ ProbeHandler: corev1.ProbeHandler{ HTTPGet: &corev1.HTTPGetAction{ Path: "/health", Port: intstr.FromInt(int(port)), Scheme: corev1.URISchemeHTTP, }, }, InitialDelaySeconds: 5, PeriodSeconds: 10, TimeoutSeconds: 3, FailureThreshold: 5, }, // preStop: 3 сек на graceful shutdown H2 перед SIGTERM. Lifecycle: &corev1.Lifecycle{ PreStop: &corev1.LifecycleHandler{ Exec: &corev1.ExecAction{ Command: []string{"sh", "-c", "sleep 3"}, }, }, }, }, } if qs.Spec.EnableUI { // SQS_ENDPOINT включает context-path: ElasticMQ принимает запросы // только на /sqs/{tenantID}, root / возвращает текстовую ошибку "The request...". ctrs = append(ctrs, corev1.Container{ Name: "elasticmq-ui", Image: elasticMQUIImage, Ports: []corev1.ContainerPort{ {ContainerPort: int32(elasticMQUIPort), Protocol: corev1.ProtocolTCP}, }, Env: []corev1.EnvVar{ { Name: "SQS_ENDPOINT", Value: fmt.Sprintf("http://localhost:%d/sqs/%s", elasticMQPort, qs.Spec.TenantID), }, }, Resources: corev1.ResourceRequirements{ Requests: corev1.ResourceList{ corev1.ResourceCPU: resource.MustParse("10m"), corev1.ResourceMemory: resource.MustParse("64Mi"), }, Limits: corev1.ResourceList{ corev1.ResourceCPU: resource.MustParse("200m"), corev1.ResourceMemory: resource.MustParse("128Mi"), }, }, }) } return ctrs }(), Volumes: []corev1.Volume{ { Name: "elasticmq-config", VolumeSource: corev1.VolumeSource{ ConfigMap: &corev1.ConfigMapVolumeSource{ LocalObjectReference: corev1.LocalObjectReference{Name: cfgMapName}, }, }, }, { Name: "elasticmq-data", VolumeSource: corev1.VolumeSource{ PersistentVolumeClaim: &corev1.PersistentVolumeClaimVolumeSource{ ClaimName: pvcName, }, }, }, }, }, }, }, } return r.Create(ctx, deploy) } // ensureService создаёт ClusterIP Service для ElasticMQ если не существует. func (r *QueueServiceReconciler) ensureService(ctx context.Context, qs *sqsv1alpha1.QueueService, ns, name, deployName string) error { svc := &corev1.Service{} if err := r.Get(ctx, client.ObjectKey{Namespace: ns, Name: name}, svc); err == nil { return nil } else if !errors.IsNotFound(err) { return fmt.Errorf("get service: %w", err) } svc = &corev1.Service{ ObjectMeta: metav1.ObjectMeta{ Name: name, Namespace: ns, Labels: sqsLabels(qs.Spec.TenantID), }, Spec: corev1.ServiceSpec{ Type: corev1.ServiceTypeClusterIP, Selector: map[string]string{ "sqs.kube5s.ru/tenant": qs.Spec.TenantID, }, Ports: func() []corev1.ServicePort { // Базовый порт ElasticMQ — всегда ports := []corev1.ServicePort{ { Name: "sqs-http", Port: int32(elasticMQPort), TargetPort: intstr.FromInt(elasticMQPort), Protocol: corev1.ProtocolTCP, }, } // Порт UI — только если enableUI=true; при self-healing сервис пересоздаётся с нужными портами if qs.Spec.EnableUI { ports = append(ports, corev1.ServicePort{ Name: "ui-http", Port: int32(elasticMQUIPort), TargetPort: intstr.FromInt(elasticMQUIPort), Protocol: corev1.ProtocolTCP, }) } return ports }(), }, } return r.Create(ctx, svc) } // ensureIngress создаёт Ingress с path-based routing для тенанта. // Путь: sqs.kube5s.ru/sqs/{tenantId} → Service :9324 (без rewrite — ElasticMQ JVM слушает с context-path). // Nginx Ingress Controller автоматически мерджит правила от разных Ingress объектов. func (r *QueueServiceReconciler) ensureIngress(ctx context.Context, qs *sqsv1alpha1.QueueService, ns, name, svcName string) error { ing := &netv1.Ingress{} if err := r.Get(ctx, client.ObjectKey{Namespace: ns, Name: name}, ing); err == nil { return nil } else if !errors.IsNotFound(err) { return fmt.Errorf("get ingress: %w", err) } // ImplementationSpecific нужен для regex, Prefix — для простого prefix-matching без rewrite. pathType := netv1.PathTypePrefix svcPort := int32(elasticMQPort) // MT03 (auth/isolation): configuration-snippet НЕ используется — nginx Ingress Controller // блокирует его как "risky annotation" (CVE-2021-25742 mitigation, включён по умолчанию с v1.9+). // Изоляция тенантов будет обеспечена через Keycloak JWT в проде (не через SigV4/snippet). ing = &netv1.Ingress{ ObjectMeta: metav1.ObjectMeta{ Name: name, Namespace: ns, Labels: sqsLabels(qs.Spec.TenantID), Annotations: map[string]string{ // Без rewrite-target: ElasticMQ JVM слушает по полному пути /sqs/{tenantId}/... "nginx.ingress.kubernetes.io/proxy-read-timeout": "60", "nginx.ingress.kubernetes.io/proxy-send-timeout": "60", "nginx.ingress.kubernetes.io/proxy-body-size": "10m", "cert-manager.io/cluster-issuer": "letsencrypt-prod", }, }, Spec: netv1.IngressSpec{ IngressClassName: strPtr("nginx"), TLS: []netv1.IngressTLS{ { Hosts: []string{r.SQSExternalHost}, SecretName: "sqs-tls-" + qs.Spec.TenantID, }, }, Rules: []netv1.IngressRule{ { Host: r.SQSExternalHost, IngressRuleValue: netv1.IngressRuleValue{ HTTP: &netv1.HTTPIngressRuleValue{ Paths: []netv1.HTTPIngressPath{ { // Простой prefix: /sqs/{tenantId} форвардится как есть в ElasticMQ. // ElasticMQ JVM обрабатывает полный путь включая context-path. Path: "/sqs/" + qs.Spec.TenantID, PathType: &pathType, Backend: netv1.IngressBackend{ Service: &netv1.IngressServiceBackend{ Name: svcName, Port: netv1.ServiceBackendPort{ Number: svcPort, }, }, }, }, }, }, }, }, }, }, } return r.Create(ctx, ing) } // setFailed переводит QueueService в фазу Failed с сообщением об ошибке. func (r *QueueServiceReconciler) setFailed(ctx context.Context, qs *sqsv1alpha1.QueueService, message string) (ctrl.Result, error) { log.FromContext(ctx).Error(fmt.Errorf(message), "QueueService failed", "tenant", qs.Spec.TenantID) qs.Status.Phase = sqsv1alpha1.QueueServicePhaseFailed qs.Status.Message = message _ = r.Status().Update(ctx, qs) return ctrl.Result{RequeueAfter: 30 * time.Second}, nil } // elasticMQImageName возвращает имя образа — из поля или константы. func (r *QueueServiceReconciler) elasticMQImageName() string { if r.ElasticMQImage != "" { return r.ElasticMQImage } return elasticMQImage } // sqsLabels возвращает стандартные labels для всех ресурсов тенанта. func sqsLabels(tenantID string) map[string]string { return map[string]string{ "app.kubernetes.io/managed-by": "sqs-operator", "app.kubernetes.io/component": "elasticmq", "sqs.kube5s.ru/tenant": tenantID, } } // SetupWithManager регистрирует контроллер в operator manager. func (r *QueueServiceReconciler) SetupWithManager(mgr ctrl.Manager) error { return ctrl.NewControllerManagedBy(mgr). For(&sqsv1alpha1.QueueService{}). Owns(&appsv1.Deployment{}). Owns(&corev1.Service{}). Owns(&corev1.ConfigMap{}). Owns(&netv1.Ingress{}). Complete(r) } // containsString проверяет наличие строки в слайсе. func containsString(slice []string, s string) bool { for _, item := range slice { if item == s { return true } } return false } // removeString удаляет строку из слайса. func removeString(slice []string, s string) []string { result := make([]string, 0, len(slice)) for _, item := range slice { if item != s { result = append(result, item) } } return result } // strPtr возвращает указатель на строку. func strPtr(s string) *string { return &s } // ensureIngressUI создаёт отдельный Ingress для ElasticMQ Web UI. // Путь: sqs.kube5s.ru/sqs-ui/{tenantId}/ → Service :3000 // Используется rewrite-target для корректной работы SPA: /sqs-ui/{tenantId}(/|$)(.*) → /$2 func (r *QueueServiceReconciler) ensureIngressUI(ctx context.Context, qs *sqsv1alpha1.QueueService, ns, name, svcName string) error { ing := &netv1.Ingress{} if err := r.Get(ctx, client.ObjectKey{Namespace: ns, Name: name}, ing); err == nil { return nil } else if !errors.IsNotFound(err) { return fmt.Errorf("get ingress-ui: %w", err) } // ImplementationSpecific + regex path для корректного rewrite в nginx. // Паттерн /sqs-ui/{tenantId}(/|$)(.*) захватывает suffix в группу $2. // rewrite-target: /$2 — убирает prefix, передаёт только suffix в UI контейнер. pathType := netv1.PathTypeImplementationSpecific uiSvcPort := int32(elasticMQUIPort) uiPath := "/sqs-ui/" + qs.Spec.TenantID + "(/|$)(.*)" ing = &netv1.Ingress{ ObjectMeta: metav1.ObjectMeta{ Name: name, Namespace: ns, Labels: sqsLabels(qs.Spec.TenantID), Annotations: map[string]string{ // rewrite-target: убираем /sqs-ui/{tenantId} prefix, UI видит путь от корня. "nginx.ingress.kubernetes.io/rewrite-target": "/$2", "nginx.ingress.kubernetes.io/use-regex": "true", "nginx.ingress.kubernetes.io/proxy-read-timeout": "30", "cert-manager.io/cluster-issuer": "letsencrypt-prod", }, }, Spec: netv1.IngressSpec{ IngressClassName: strPtr("nginx"), TLS: []netv1.IngressTLS{ { Hosts: []string{r.SQSExternalHost}, SecretName: "sqs-tls-" + qs.Spec.TenantID, }, }, Rules: []netv1.IngressRule{ { Host: r.SQSExternalHost, IngressRuleValue: netv1.IngressRuleValue{ HTTP: &netv1.HTTPIngressRuleValue{ Paths: []netv1.HTTPIngressPath{ { Path: uiPath, PathType: &pathType, Backend: netv1.IngressBackend{ Service: &netv1.IngressServiceBackend{ Name: svcName, Port: netv1.ServiceBackendPort{ Number: uiSvcPort, }, }, }, }, }, }, }, }, }, }, } return r.Create(ctx, ing) } // ensureIngressUIAssets создаёт Ingress для статических ресурсов Next.js (/_next/). // Next.js собран с basePath="" — ресурсы запрашиваются по абсолютному пути /_next/... // Без rewrite: nginx передаёт /_next/... напрямую в UI-контейнер (port 3000). // Все UI-поды одного образа отдают идентичные /_next/ файлы — конфликтов нет. func (r *QueueServiceReconciler) ensureIngressUIAssets(ctx context.Context, qs *sqsv1alpha1.QueueService, ns, name, svcName string) error { ing := &netv1.Ingress{} if err := r.Get(ctx, client.ObjectKey{Namespace: ns, Name: name}, ing); err == nil { return nil } else if !errors.IsNotFound(err) { return fmt.Errorf("get ui-assets ingress: %w", err) } // Prefix без rewrite-target — nginx пробрасывает /_next/... как есть в Next.js сервер. pathType := netv1.PathTypePrefix ing = &netv1.Ingress{ ObjectMeta: metav1.ObjectMeta{ Name: name, Namespace: ns, Labels: sqsLabels(qs.Spec.TenantID), Annotations: map[string]string{ "nginx.ingress.kubernetes.io/proxy-read-timeout": "30", "cert-manager.io/cluster-issuer": "letsencrypt-prod", }, }, Spec: netv1.IngressSpec{ IngressClassName: strPtr("nginx"), TLS: []netv1.IngressTLS{ { Hosts: []string{r.SQSExternalHost}, SecretName: "sqs-tls-" + qs.Spec.TenantID, }, }, Rules: []netv1.IngressRule{ { Host: r.SQSExternalHost, IngressRuleValue: netv1.IngressRuleValue{ HTTP: &netv1.HTTPIngressRuleValue{ Paths: []netv1.HTTPIngressPath{ { Path: "/_next/", PathType: &pathType, Backend: netv1.IngressBackend{ Service: &netv1.IngressServiceBackend{ Name: svcName, Port: netv1.ServiceBackendPort{ Number: int32(elasticMQUIPort), }, }, }, }, }, }, }, }, }, }, } return r.Create(ctx, ing) } // ensureIngressUIQueues создаёт Ingress для внутренних маршрутов Next.js (/queues/*). // Next.js собран с basePath="" — при навигации браузер обращается к /queues/[name] напрямую. // Без rewrite-target: nginx передаёт /queues/... как есть в UI-контейнер (port 3000). func (r *QueueServiceReconciler) ensureIngressUIQueues(ctx context.Context, qs *sqsv1alpha1.QueueService, ns, name, svcName string) error { ing := &netv1.Ingress{} if err := r.Get(ctx, client.ObjectKey{Namespace: ns, Name: name}, ing); err == nil { return nil } else if !errors.IsNotFound(err) { return fmt.Errorf("get ui-queues ingress: %w", err) } pathType := netv1.PathTypePrefix ing = &netv1.Ingress{ ObjectMeta: metav1.ObjectMeta{ Name: name, Namespace: ns, Labels: sqsLabels(qs.Spec.TenantID), Annotations: map[string]string{ "nginx.ingress.kubernetes.io/proxy-read-timeout": "30", "cert-manager.io/cluster-issuer": "letsencrypt-prod", }, }, Spec: netv1.IngressSpec{ IngressClassName: strPtr("nginx"), TLS: []netv1.IngressTLS{ { Hosts: []string{r.SQSExternalHost}, SecretName: "sqs-tls-" + qs.Spec.TenantID, }, }, Rules: []netv1.IngressRule{ { Host: r.SQSExternalHost, IngressRuleValue: netv1.IngressRuleValue{ HTTP: &netv1.HTTPIngressRuleValue{ Paths: []netv1.HTTPIngressPath{ { Path: "/queues", PathType: &pathType, Backend: netv1.IngressBackend{ Service: &netv1.IngressServiceBackend{ Name: svcName, Port: netv1.ServiceBackendPort{ Number: int32(elasticMQUIPort), }, }, }, }, }, }, }, }, }, }, } return r.Create(ctx, ing) }