diff --git a/sqs-operator/api/v1alpha1/queueservice_types.go b/sqs-operator/api/v1alpha1/queueservice_types.go index 8f2abf6..ee21565 100644 --- a/sqs-operator/api/v1alpha1/queueservice_types.go +++ b/sqs-operator/api/v1alpha1/queueservice_types.go @@ -1,4 +1,4 @@ -// Изменён: 2026-04-07 +// Изменён: 2026-04-08 // queueservice_types.go — CRD типы для QueueService. // Описывает Spec (желаемое состояние) и Status (наблюдаемое состояние) ElasticMQ инстанса тенанта. @@ -54,6 +54,20 @@ type QueueServiceSpec struct { // +kubebuilder:default=true // +optional Persistence bool `json:"persistence,omitempty"` + + // EnableUI — добавить sidecar контейнер ElasticMQ Web UI (softwaremill/elasticmq-ui). + // UI доступен по https://{SQS_EXTERNAL_HOST}/sqs-ui/{tenantId}/ + // UI контейнер обращается к ElasticMQ через localhost:9324 (sidecar в одном поде). + // +kubebuilder:default=false + // +optional + EnableUI bool `json:"enableUI,omitempty"` + + // AutoCreateQueues — автоматически создавать очередь при первом обращении. + // Если true: SendMessage/ReceiveMessage в несуществующую очередь создаёт её автоматически. + // Если false: возвращает NonExistentQueue ошибку (поведение AWS SQS по умолчанию). + // +kubebuilder:default=false + // +optional + AutoCreateQueues bool `json:"autoCreateQueues,omitempty"` } // QueueServiceStatus — наблюдаемое состояние QueueService. diff --git a/sqs-operator/config/crd/bases/sqs.kube5s.ru_queueservices.yaml b/sqs-operator/config/crd/bases/sqs.kube5s.ru_queueservices.yaml index 375afad..5cd23e9 100644 --- a/sqs-operator/config/crd/bases/sqs.kube5s.ru_queueservices.yaml +++ b/sqs-operator/config/crd/bases/sqs.kube5s.ru_queueservices.yaml @@ -55,6 +55,20 @@ spec: description: QueueServiceSpec — желаемое состояние QueueService (параметры тенанта). properties: + autoCreateQueues: + default: false + description: |- + AutoCreateQueues — автоматически создавать очередь при первом обращении. + Если true: SendMessage/ReceiveMessage в несуществующую очередь создаёт её автоматически. + Если false: возвращает NonExistentQueue ошибку (поведение AWS SQS по умолчанию). + type: boolean + enableUI: + default: false + description: |- + EnableUI — добавить sidecar контейнер ElasticMQ Web UI (softwaremill/elasticmq-ui). + UI доступен по https://{SQS_EXTERNAL_HOST}/sqs-ui/{tenantId}/ + UI контейнер обращается к ElasticMQ через localhost:9324 (sidecar в одном поде). + type: boolean memoryMB: default: 64 description: |- diff --git a/sqs-operator/internal/controller/queueservice_controller.go b/sqs-operator/internal/controller/queueservice_controller.go index 3bf8a39..12a002d 100644 --- a/sqs-operator/internal/controller/queueservice_controller.go +++ b/sqs-operator/internal/controller/queueservice_controller.go @@ -1,4 +1,4 @@ -// Изменён: 2026-04-07 +// Изменён: 2026-04-08 // queueservice_controller.go — Reconciler для QueueService CRD. // Управляет жизненным циклом ElasticMQ инстанса: создание/удаление всех k8s ресурсов. // @@ -41,8 +41,13 @@ const ( // 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. @@ -148,12 +153,20 @@ func (r *QueueServiceReconciler) provision(ctx context.Context, qs *sqsv1alpha1. return r.setFailed(ctx, qs, fmt.Sprintf("ensure service: %v", err)) } - // 7. Ingress + // 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)) + } + } + // Все ресурсы созданы — переходим в Provisioning, ждём готовности pod qs.Status.Phase = sqsv1alpha1.QueueServicePhaseProvisioning qs.Status.Message = "resources created, waiting for pod ready" @@ -218,6 +231,12 @@ func (r *QueueServiceReconciler) ensureHealthy(ctx context.Context, qs *sqsv1alp {"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{}}) + } for _, res := range checkResources { if err := r.Get(ctx, client.ObjectKey{Namespace: tenantNS, Name: res.name}, res.obj); err != nil { if errors.IsNotFound(err) { @@ -283,7 +302,13 @@ func (r *QueueServiceReconciler) handleDeletion(ctx context.Context, qs *sqsv1al deleteOpts := []client.DeleteOption{client.PropagationPolicy(metav1.DeletePropagationForeground)} - // Удаляем Ingress + // Удаляем 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...) @@ -388,7 +413,7 @@ func (r *QueueServiceReconciler) ensureConfigMap(ctx context.Context, qs *sqsv1a return fmt.Errorf("get configmap: %w", err) } - conf := elasticmq.GenerateConfig(qs.Spec.TenantID, r.SQSExternalHost, qs.Spec.Persistence) + conf := elasticmq.GenerateConfig(qs.Spec.TenantID, r.SQSExternalHost, qs.Spec.Persistence, qs.Spec.AutoCreateQueues) cm = &corev1.ConfigMap{ ObjectMeta: metav1.ObjectMeta{ @@ -741,3 +766,70 @@ func removeString(slice []string, s string) []string { 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) +} diff --git a/sqs-operator/internal/elasticmq/elasticmq_config.go b/sqs-operator/internal/elasticmq/elasticmq_config.go index bf27e2d..8371a65 100644 --- a/sqs-operator/internal/elasticmq/elasticmq_config.go +++ b/sqs-operator/internal/elasticmq/elasticmq_config.go @@ -1,4 +1,4 @@ -// Изменён: 2026-04-07 +// Изменён: 2026-04-08 // elasticmq_config.go — генератор HOCON конфигурации для ElasticMQ. // Каждый тенант получает уникальный конфиг: свой context-path, accountId, persistence. // Конфиг монтируется в pod как ConfigMap → /opt/elasticmq.conf (путь из ENTRYPOINT JVM-образа). @@ -8,10 +8,11 @@ package elasticmq import "fmt" // GenerateConfig генерирует HOCON конфигурацию для ElasticMQ инстанса тенанта. -// tenantID — ID тенанта (используется в context-path и accountId). -// externalHost — публичный хост (SQS_EXTERNAL_HOST env var, например sqs.kube5s.ru). -// persistence — если true, сообщения сохраняются в H2 на PVC /data. -func GenerateConfig(tenantID, externalHost string, persistence bool) string { +// tenantID — ID тенанта (используется в context-path и accountId). +// externalHost — публичный хост (SQS_EXTERNAL_HOST env var, например sqs.kube5s.ru). +// persistence — если true, сообщения сохраняются в H2 на PVC /data. +// autoCreateQueues — если true, очередь создаётся автоматически при первом обращении. +func GenerateConfig(tenantID, externalHost string, persistence bool, autoCreateQueues bool) string { persistenceBlock := "" if persistence { // H2 база данных хранит очереди и сообщения между рестартами pod. @@ -28,6 +29,16 @@ messages-storage { }` } + autoCreateBlock := "" + if autoCreateQueues { + // При первом обращении к несуществующей очереди — создать её автоматически. + // Полезно для dev/mvp: не нужен явный CreateQueue перед SendMessage. + autoCreateBlock = ` +auto-create-queues { + enabled = true +}` + } + return fmt.Sprintf(`include classpath("application.conf") # Внешний адрес ноды — как будут формироваться queue URL в ответах SQS API. @@ -46,11 +57,11 @@ rest-sqs { # strict — соответствие AWS SQS валидации параметров sqs-limits = strict } -%s +%s%s # accountId = tenantID — используется в ARN очередей aws { region = ru-msk-1 accountId = "%s" } -`, tenantID, externalHost, tenantID, persistenceBlock, tenantID) +`, tenantID, externalHost, tenantID, persistenceBlock, autoCreateBlock, tenantID) } diff --git a/sqs-operator/test_v2_suite.sh b/sqs-operator/test_v2_suite.sh index f4a7c2e..0da4534 100755 --- a/sqs-operator/test_v2_suite.sh +++ b/sqs-operator/test_v2_suite.sh @@ -155,7 +155,18 @@ kubectl delete queueservice "$TIMING_CR" -n "$CR_NAMESPACE" --ignore-not-found=t --wait=true --timeout=30s 2>&1 | grep -v "^$" | head -3 kubectl delete ns "sless-fn-${TIMING_TENANT}" --ignore-not-found=true \ --wait=false 2>/dev/null -sleep 5 +# Ждём полного удаления NS — если завис (Terminating), принудительно убираем finalizer +for i in $(seq 1 20); do + STATUS=$(kubectl get ns "sless-fn-${TIMING_TENANT}" -o jsonpath='{.status.phase}' 2>/dev/null) + [ -z "$STATUS" ] && break + if [ "$STATUS" = "Terminating" ] && [ "$i" -ge 10 ]; then + # NS завис в Terminating — принудительно убираем finalizer + kubectl get ns "sless-fn-${TIMING_TENANT}" -o json \ + | python3 -c "import sys,json; d=json.load(sys.stdin); d['spec']['finalizers']=[]; print(json.dumps(d))" \ + | kubectl replace --raw "/api/v1/namespaces/sless-fn-${TIMING_TENANT}/finalize" -f - >/dev/null 2>&1 + fi + sleep 2 +done echo " Создаём QueueService..." T0_APPLY=$(date +%s) @@ -454,6 +465,9 @@ echo "$R" | grep -q "DeleteMessageBatchResponse\|ResultCode\|ResponseMetadata" \ # A06: Long polling (WaitTimeSeconds=3, пустая очередь) echo "--- A06 Long polling 3s" sleep 6 # Ждём VisibilityTimeout для сообщения из A02 +# Чистим очередь — для long polling нужна гарантированно пустая очередь +sqs "Action=PurgeQueue&QueueUrl=${ADV_URL}&Version=2012-11-05" > /dev/null 2>&1 +sleep 1 START=$(date +%s) R=$(sqs "Action=ReceiveMessage&QueueUrl=${ADV_URL}&MaxNumberOfMessages=1&WaitTimeSeconds=3&Version=2012-11-05") WAITED=$(elapsed $START)