refactor(sqs-operator): переделка через Operator SDK v1.37.0
- operator-sdk init + create api (QueueService kind, group sqs, v1alpha1) - Стандартная kubebuilder структура: cmd/, internal/controller/, config/ - CRD types с kubebuilder маркерами (validation, defaults, printcolumns) - Reconciler перенесён из ручного кода в internal/controller/ - controller-gen v0.17.0 (совместимость с Go 1.26.1) - Автогенерация: deepcopy, CRD YAML, RBAC ClusterRole - Удалены все ручные файлы (controllers/, main.go, deployments/)
This commit is contained in:
@@ -0,0 +1,705 @@
|
||||
// Изменён: 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 Native (GraalVM, ~70MB, старт <0.5с).
|
||||
elasticMQImage = "softwaremill/elasticmq-native: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{
|
||||
"custom.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)
|
||||
|
||||
memLimit := resource.MustParse(fmt.Sprintf("%dMi", qs.Spec.MemoryMB))
|
||||
memRequest := resource.MustParse("32Mi")
|
||||
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{
|
||||
Containers: []corev1.Container{
|
||||
{
|
||||
Name: "elasticmq",
|
||||
Image: image,
|
||||
Ports: []corev1.ContainerPort{
|
||||
{ContainerPort: port, Protocol: corev1.ProtocolTCP},
|
||||
},
|
||||
// ElasticMQ Native читает конфиг через системное свойство JVM.
|
||||
// Путь должен совпадать с mountPath в VolumeMounts.
|
||||
Env: []corev1.EnvVar{
|
||||
{
|
||||
Name: "JAVA_TOOL_OPTIONS",
|
||||
Value: "-Dconfig.file=/opt/elasticmq/custom.conf",
|
||||
},
|
||||
},
|
||||
Resources: corev1.ResourceRequirements{
|
||||
Requests: corev1.ResourceList{
|
||||
corev1.ResourceCPU: cpuRequest,
|
||||
corev1.ResourceMemory: memRequest,
|
||||
},
|
||||
Limits: corev1.ResourceList{
|
||||
corev1.ResourceCPU: cpuLimit,
|
||||
corev1.ResourceMemory: memLimit,
|
||||
},
|
||||
},
|
||||
VolumeMounts: []corev1.VolumeMount{
|
||||
{
|
||||
// ConfigMap монтируем как единственный файл custom.conf
|
||||
Name: "elasticmq-config",
|
||||
MountPath: "/opt/elasticmq/custom.conf",
|
||||
SubPath: "custom.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
|
||||
// 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)
|
||||
}
|
||||
|
||||
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{
|
||||
// Nginx rewrite: убираем prefix /sqs/{tenantId} перед проксированием в ElasticMQ.
|
||||
// ElasticMQ получает чистый SQS path (например /?Action=CreateQueue).
|
||||
"nginx.ingress.kubernetes.io/rewrite-target": "/$2",
|
||||
"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{
|
||||
{
|
||||
// Regex capture group ($2) передаётся в rewrite-target.
|
||||
// /sqs/{tenantId}(/|$)(.*) → /$2
|
||||
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
|
||||
}
|
||||
@@ -0,0 +1,84 @@
|
||||
/*
|
||||
Copyright 2026.
|
||||
|
||||
Licensed under the Apache License, Version 2.0 (the "License");
|
||||
you may not use this file except in compliance with the License.
|
||||
You may obtain a copy of the License at
|
||||
|
||||
http://www.apache.org/licenses/LICENSE-2.0
|
||||
|
||||
Unless required by applicable law or agreed to in writing, software
|
||||
distributed under the License is distributed on an "AS IS" BASIS,
|
||||
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
See the License for the specific language governing permissions and
|
||||
limitations under the License.
|
||||
*/
|
||||
|
||||
package controller
|
||||
|
||||
import (
|
||||
"context"
|
||||
|
||||
. "github.com/onsi/ginkgo/v2"
|
||||
. "github.com/onsi/gomega"
|
||||
"k8s.io/apimachinery/pkg/api/errors"
|
||||
"k8s.io/apimachinery/pkg/types"
|
||||
"sigs.k8s.io/controller-runtime/pkg/reconcile"
|
||||
|
||||
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||||
|
||||
sqsv1alpha1 "gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sqs-operator/api/v1alpha1"
|
||||
)
|
||||
|
||||
var _ = Describe("QueueService Controller", func() {
|
||||
Context("When reconciling a resource", func() {
|
||||
const resourceName = "test-resource"
|
||||
|
||||
ctx := context.Background()
|
||||
|
||||
typeNamespacedName := types.NamespacedName{
|
||||
Name: resourceName,
|
||||
Namespace: "default", // TODO(user):Modify as needed
|
||||
}
|
||||
queueservice := &sqsv1alpha1.QueueService{}
|
||||
|
||||
BeforeEach(func() {
|
||||
By("creating the custom resource for the Kind QueueService")
|
||||
err := k8sClient.Get(ctx, typeNamespacedName, queueservice)
|
||||
if err != nil && errors.IsNotFound(err) {
|
||||
resource := &sqsv1alpha1.QueueService{
|
||||
ObjectMeta: metav1.ObjectMeta{
|
||||
Name: resourceName,
|
||||
Namespace: "default",
|
||||
},
|
||||
// TODO(user): Specify other spec details if needed.
|
||||
}
|
||||
Expect(k8sClient.Create(ctx, resource)).To(Succeed())
|
||||
}
|
||||
})
|
||||
|
||||
AfterEach(func() {
|
||||
// TODO(user): Cleanup logic after each test, like removing the resource instance.
|
||||
resource := &sqsv1alpha1.QueueService{}
|
||||
err := k8sClient.Get(ctx, typeNamespacedName, resource)
|
||||
Expect(err).NotTo(HaveOccurred())
|
||||
|
||||
By("Cleanup the specific resource instance QueueService")
|
||||
Expect(k8sClient.Delete(ctx, resource)).To(Succeed())
|
||||
})
|
||||
It("should successfully reconcile the resource", func() {
|
||||
By("Reconciling the created resource")
|
||||
controllerReconciler := &QueueServiceReconciler{
|
||||
Client: k8sClient,
|
||||
Scheme: k8sClient.Scheme(),
|
||||
}
|
||||
|
||||
_, err := controllerReconciler.Reconcile(ctx, reconcile.Request{
|
||||
NamespacedName: typeNamespacedName,
|
||||
})
|
||||
Expect(err).NotTo(HaveOccurred())
|
||||
// TODO(user): Add more specific assertions depending on your controller's reconciliation logic.
|
||||
// Example: If you expect a certain status condition after reconciliation, verify it here.
|
||||
})
|
||||
})
|
||||
})
|
||||
@@ -0,0 +1,90 @@
|
||||
/*
|
||||
Copyright 2026.
|
||||
|
||||
Licensed under the Apache License, Version 2.0 (the "License");
|
||||
you may not use this file except in compliance with the License.
|
||||
You may obtain a copy of the License at
|
||||
|
||||
http://www.apache.org/licenses/LICENSE-2.0
|
||||
|
||||
Unless required by applicable law or agreed to in writing, software
|
||||
distributed under the License is distributed on an "AS IS" BASIS,
|
||||
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
See the License for the specific language governing permissions and
|
||||
limitations under the License.
|
||||
*/
|
||||
|
||||
package controller
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"path/filepath"
|
||||
"runtime"
|
||||
"testing"
|
||||
|
||||
. "github.com/onsi/ginkgo/v2"
|
||||
. "github.com/onsi/gomega"
|
||||
|
||||
"k8s.io/client-go/kubernetes/scheme"
|
||||
"k8s.io/client-go/rest"
|
||||
"sigs.k8s.io/controller-runtime/pkg/client"
|
||||
"sigs.k8s.io/controller-runtime/pkg/envtest"
|
||||
logf "sigs.k8s.io/controller-runtime/pkg/log"
|
||||
"sigs.k8s.io/controller-runtime/pkg/log/zap"
|
||||
|
||||
sqsv1alpha1 "gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sqs-operator/api/v1alpha1"
|
||||
//+kubebuilder:scaffold:imports
|
||||
)
|
||||
|
||||
// These tests use Ginkgo (BDD-style Go testing framework). Refer to
|
||||
// http://onsi.github.io/ginkgo/ to learn more about Ginkgo.
|
||||
|
||||
var cfg *rest.Config
|
||||
var k8sClient client.Client
|
||||
var testEnv *envtest.Environment
|
||||
|
||||
func TestControllers(t *testing.T) {
|
||||
RegisterFailHandler(Fail)
|
||||
|
||||
RunSpecs(t, "Controller Suite")
|
||||
}
|
||||
|
||||
var _ = BeforeSuite(func() {
|
||||
logf.SetLogger(zap.New(zap.WriteTo(GinkgoWriter), zap.UseDevMode(true)))
|
||||
|
||||
By("bootstrapping test environment")
|
||||
testEnv = &envtest.Environment{
|
||||
CRDDirectoryPaths: []string{filepath.Join("..", "..", "config", "crd", "bases")},
|
||||
ErrorIfCRDPathMissing: true,
|
||||
|
||||
// The BinaryAssetsDirectory is only required if you want to run the tests directly
|
||||
// without call the makefile target test. If not informed it will look for the
|
||||
// default path defined in controller-runtime which is /usr/local/kubebuilder/.
|
||||
// Note that you must have the required binaries setup under the bin directory to perform
|
||||
// the tests directly. When we run make test it will be setup and used automatically.
|
||||
BinaryAssetsDirectory: filepath.Join("..", "..", "bin", "k8s",
|
||||
fmt.Sprintf("1.29.0-%s-%s", runtime.GOOS, runtime.GOARCH)),
|
||||
}
|
||||
|
||||
var err error
|
||||
// cfg is defined in this file globally.
|
||||
cfg, err = testEnv.Start()
|
||||
Expect(err).NotTo(HaveOccurred())
|
||||
Expect(cfg).NotTo(BeNil())
|
||||
|
||||
err = sqsv1alpha1.AddToScheme(scheme.Scheme)
|
||||
Expect(err).NotTo(HaveOccurred())
|
||||
|
||||
//+kubebuilder:scaffold:scheme
|
||||
|
||||
k8sClient, err = client.New(cfg, client.Options{Scheme: scheme.Scheme})
|
||||
Expect(err).NotTo(HaveOccurred())
|
||||
Expect(k8sClient).NotTo(BeNil())
|
||||
|
||||
})
|
||||
|
||||
var _ = AfterSuite(func() {
|
||||
By("tearing down the test environment")
|
||||
err := testEnv.Stop()
|
||||
Expect(err).NotTo(HaveOccurred())
|
||||
})
|
||||
Reference in New Issue
Block a user