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