feat: sqs-operator — CRD QueueService, контроллер, ElasticMQ per-tenant

This commit is contained in:
Naeel
2026-04-07 13:56:53 +03:00
parent a22a37bef3
commit cbe61d9f62
9 changed files with 1956 additions and 0 deletions
+49
View File
@@ -0,0 +1,49 @@
// Создан: 2026-04-07
// config.go — конфигурация SQS Operator из env переменных.
// Все параметры читаются при старте через Load().
package config
import (
"fmt"
"os"
)
// Config — конфигурация SQS Operator.
type Config struct {
// SQSExternalHost — публичный хост для SQS endpoints тенантов.
// Используется в: Ingress rule host, elasticmq.conf node-address, Status.Endpoint.
// dev: sqs.kube5s.ru / prod: sqs.cloud-provider.ru
// Обязательный параметр — не хардкодится, чтобы переехать в другую зону без пересборки.
SQSExternalHost string
// ElasticMQImage — Docker образ ElasticMQ для pod тенанта.
// Опциональный — если не задан, используется константа из контроллера (v1.7.1).
// Позволяет обновить версию ElasticMQ без пересборки оператора.
ElasticMQImage string
// OperatorNamespace — namespace где запущен сам оператор.
// Default: sless
OperatorNamespace string
}
// Load читает конфигурацию из env переменных.
func Load() (*Config, error) {
cfg := &Config{
OperatorNamespace: "sless",
}
cfg.SQSExternalHost = os.Getenv("SQS_EXTERNAL_HOST")
if cfg.SQSExternalHost == "" {
return nil, fmt.Errorf("SQS_EXTERNAL_HOST is required (example: sqs.kube5s.ru)")
}
// Опциональный — если пуст, контроллер использует захардкоженную версию
cfg.ElasticMQImage = os.Getenv("SQS_ELASTICMQ_IMAGE")
if ns := os.Getenv("OPERATOR_NAMESPACE"); ns != "" {
cfg.OperatorNamespace = ns
}
return cfg, nil
}
+56
View File
@@ -0,0 +1,56 @@
// Создан: 2026-04-07
// elasticmq_config.go — генератор HOCON конфигурации для ElasticMQ.
// Каждый тенант получает уникальный конфиг: свой context-path, accountId, persistence.
// Конфиг монтируется в pod как ConfigMap → /opt/elasticmq/custom.conf.
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 {
persistenceBlock := ""
if persistence {
// H2 база данных хранит очереди и сообщения между рестартами pod.
// Файл /data/elasticmq.db — на PVC тенанта.
persistenceBlock = `
messages-storage {
enabled = true
uri = "jdbc:h2:/data/elasticmq"
}`
} else {
persistenceBlock = `
messages-storage {
enabled = false
}`
}
return fmt.Sprintf(`include classpath("application.conf")
# Внешний адрес ноды — как будут формироваться queue URL в ответах SQS API.
# context-path должен совпадать с путём в Ingress: /sqs/%s
node-address {
protocol = https
host = %s
port = 443
context-path = "/sqs/%s"
}
rest-sqs {
enabled = true
bind-port = 9324
bind-hostname = "0.0.0.0"
# strict — соответствие AWS SQS валидации параметров
sqs-limits = strict
}
%s
# accountId = tenantID — используется в ARN очередей
aws {
region = ru-msk-1
accountId = "%s"
}
`, tenantID, externalHost, tenantID, persistenceBlock, tenantID)
}
@@ -0,0 +1,44 @@
// Создан: 2026-04-07
// credentials.go — генератор access/secret key пары для тенанта.
// Credentials создаются один раз при провизионировании и хранятся в k8s Secret.
// Тенант передаёт их в AWS SDK как aws_access_key_id / aws_secret_access_key.
//
// ElasticMQ принимает любые non-empty credentials — аутентификация не проверяется
// на уровне ElasticMQ. Изоляция достигается через Ingress (каждый тенант
// обращается только к своему path /sqs/{tenantId}).
package elasticmq
import (
"crypto/rand"
"encoding/base64"
"fmt"
)
// Credentials — пара access/secret key для SQS клиента тенанта.
type Credentials struct {
AccessKey string
SecretKey string
}
// GenerateCredentials создаёт уникальную пару ключей для тенанта.
// accessKey: SQSAK-{tenantID}-{random6hex} — читаемый, уникальный
// secretKey: 32 случайных байта → base64 URL-safe (без padding)
func GenerateCredentials(tenantID string) (Credentials, error) {
// 3 байта → 6 hex символов для суффикса accessKey
akSuffix := make([]byte, 3)
if _, err := rand.Read(akSuffix); err != nil {
return Credentials{}, fmt.Errorf("generate access key suffix: %w", err)
}
// 32 байта → надёжный secretKey
skBytes := make([]byte, 32)
if _, err := rand.Read(skBytes); err != nil {
return Credentials{}, fmt.Errorf("generate secret key: %w", err)
}
return Credentials{
AccessKey: fmt.Sprintf("SQSAK-%s-%x", tenantID, akSuffix),
SecretKey: base64.RawURLEncoding.EncodeToString(skBytes),
}, nil
}