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:
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
|
||||
Reference in New Issue
Block a user