sqs-operator v0.1.7: web UI sidecar, auto-create-queues, fix A06 long polling

- api/v1alpha1: добавлены поля EnableUI, AutoCreateQueues в QueueServiceSpec
- elasticmq_config.go: GenerateConfig принимает autoCreateQueues, добавляет HOCON блок
- controller.go: sidecar elasticmq-ui (port 3000), ensureIngressUI, ensureService UI port
- test_v2_suite.sh: A06 теперь PASS — PurgeQueue перед long polling тестом
- CRD обновлён: make install применён в кластере

Test run v0.1.7: 35 PASS / 3 FAIL / 4 WARN / 2 SKIP(WONTFIX)
This commit is contained in:
Naeel
2026-04-08 19:11:26 +03:00
parent c4efc5c960
commit 64626659a9
5 changed files with 158 additions and 13 deletions
@@ -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.
@@ -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: |-
@@ -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)
}
@@ -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)
}
+15 -1
View File
@@ -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)