Files
sless/sqs-operator/internal/controller/queueservice_controller.go
T
Naeel 336ee7b869 fix(sqs-operator): persistence, OOM, Ingress routing (v0.1.2→v0.1.4)
Problem 1: elasticmq-native does not support H2 JDBC persistence
- GraalVM native build excludes H2 driver
- Fix: switch to softwaremill/elasticmq:1.7.1 (JVM image)

Problem 2: OOMKilled — JVM requires >150MB, spec.MemoryMB=64 too low
- Fix: enforce minimum 256Mi for JVM; add -Xmx (75% of limit)

Problem 3: AccessDeniedException on /data — PVC mounted as root:root
- JVM image runs as uid=999 (elasticmq)
- Fix: add fsGroup=999 to PodSecurityContext

Problem 4: 404 through HTTPS — nginx rewrite stripped context-path
- ElasticMQ JVM listens at context-path (/sqs/{tenantId})
- Native image tolerated /, JVM does not
- Fix: remove rewrite-target annotation, use plain Prefix path

Result: S6 PASS — 8 queues + messages survived graceful pod kill
2026-04-07 15:34:15 +03:00

725 lines
26 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-07
// 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"
// elasticMQPort — порт на котором ElasticMQ слушает SQS HTTP запросы.
elasticMQPort = 9324
)
// 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
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))
}
// Все ресурсы созданы — переходим в 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).
// Если pod упал — переходим в Failed для последующего восстановления.
func (r *QueueServiceReconciler) ensureHealthy(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) {
qs.Status.Phase = sqsv1alpha1.QueueServicePhasePending
qs.Status.Message = "deployment disappeared, reprovisioning"
_ = r.Status().Update(ctx, qs)
return ctrl.Result{Requeue: true}, nil
}
return ctrl.Result{}, fmt.Errorf("health check get deployment: %w", err)
}
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 через 30 секунд
return ctrl.Result{RequeueAfter: 30 * 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
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)
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,
},
},
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 }(),
},
Containers: []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,
FailureThreshold: 5,
},
},
},
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: []corev1.ServicePort{
{
Name: "sqs-http",
Port: int32(elasticMQPort),
TargetPort: intstr.FromInt(elasticMQPort),
Protocol: corev1.ProtocolTCP,
},
},
},
}
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)
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}/...
// (context-path в конфиге определяет listen path у JVM образа, в отличие от native).
"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
}