Files
sless/sqs-operator/internal/controller/queueservice_controller.go
T
Naeel d1d1bffd7c v0.1.13: Strategy Recreate + preStop + liveness fix (ERR-SQS-06 root cause)
- ensureDeployment: Strategy Recreate (not RollingUpdate) — prevents H2 file lock
  when two pods mount same PVC simultaneously during rollout restart
- preStop: sleep 3 — graceful H2 shutdown before SIGTERM
- livenessProbe timeoutSeconds: 3 — prevents false positive on GC pause
- terminationGracePeriodSeconds: 15
- doc: thinking log, ERR-SQS-06, progress.md, architecture SVG schema
2026-04-09 07:49:11 +03:00

1047 lines
39 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.
// Изменён: 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)
}