diff --git a/api/v1alpha1/trigger_types.go b/api/v1alpha1/trigger_types.go index f550b4b..6bb1e45 100644 --- a/api/v1alpha1/trigger_types.go +++ b/api/v1alpha1/trigger_types.go @@ -1,6 +1,8 @@ -// Изменено: 2026-03-08 -// Описание CRD Trigger — триггер для функции (HTTP или Cron). +// Изменено: 2026-03-19 +// Описание CRD Trigger — триггер для функции (HTTP, Cron или Event). // Один Trigger ссылается на одну Function и определяет способ вызова. +// Event-тип: event-dispatcher подписывается на AMQP очередь и при сообщении +// вызывает функцию по внутреннему HTTP. package v1alpha1 @@ -16,6 +18,9 @@ const ( TriggerTypeHTTP TriggerType = "http" // TriggerTypeCron — функция вызывается по расписанию (k8s CronJob) TriggerTypeCron TriggerType = "cron" + // TriggerTypeEvent — функция вызывается при получении сообщения из AMQP очереди. + // event-dispatcher подписывается на spec.queue в RabbitMQ и делает POST на HTTP endpoint функции. + TriggerTypeEvent TriggerType = "event" ) // TriggerSpec — желаемое состояние триггера. @@ -30,8 +35,8 @@ type TriggerSpec struct { // +kubebuilder:validation:Required FunctionRef string `json:"functionRef"` - // Type — тип триггера: http или cron - // +kubebuilder:validation:Enum=http;cron + // Type — тип триггера: http, cron или event + // +kubebuilder:validation:Enum=http;cron;event // +kubebuilder:validation:Required Type TriggerType `json:"type"` @@ -42,6 +47,10 @@ type TriggerSpec struct { // Актуально для cron: запускаем pod заранее чтобы избежать cold start. // +kubebuilder:default=300 PreWarmSeconds int32 `json:"preWarmSeconds,omitempty"` + + // Queue — имя AMQP очереди в RabbitMQ (только для type=event). + // event-dispatcher подпишется на эту очередь и вызовет функцию при каждом сообщении. + Queue string `json:"queue,omitempty"` } // TriggerStatus — наблюдаемое состояние триггера (заполняет контроллер). @@ -65,6 +74,7 @@ type TriggerStatus struct { //+kubebuilder:printcolumn:name="Function",type=string,JSONPath=`.spec.functionRef` //+kubebuilder:printcolumn:name="Active",type=boolean,JSONPath=`.status.active` //+kubebuilder:printcolumn:name="URL",type=string,JSONPath=`.status.url` +//+kubebuilder:printcolumn:name="Queue",type=string,JSONPath=`.spec.queue` //+kubebuilder:printcolumn:name="Age",type=date,JSONPath=`.metadata.creationTimestamp` // Trigger — ресурс для управления способом вызова Function. diff --git a/controllers/trigger_controller.go b/controllers/trigger_controller.go index f637c0a..6d613b6 100644 --- a/controllers/trigger_controller.go +++ b/controllers/trigger_controller.go @@ -124,6 +124,9 @@ func (r *TriggerReconciler) Reconcile(ctx context.Context, req ctrl.Request) (ct case slessv1alpha1.TriggerTypeCron: logger.Info("reconcile cron trigger", "trigger", tr.Name) return r.reconcileCron(ctx, tr, fn) + case slessv1alpha1.TriggerTypeEvent: + logger.Info("reconcile event trigger", "trigger", tr.Name) + return r.reconcileEvent(ctx, tr, fn) } return ctrl.Result{}, nil @@ -325,3 +328,44 @@ func (r *TriggerReconciler) SetupWithManager(mgr ctrl.Manager) error { For(&slessv1alpha1.Trigger{}). Complete(r) } + +// reconcileEvent обрабатывает Trigger{type:event}. +// Оператор не управляет AMQP напрямую — это задача event-dispatcher. +// Здесь: убеждаемся что Service функции существует (dispatcher использует его для POST), +// обновляем статус триггера. +func (r *TriggerReconciler) reconcileEvent(ctx context.Context, tr *slessv1alpha1.Trigger, fn *slessv1alpha1.Function) (ctrl.Result, error) { + if tr.Spec.Queue == "" { + tr.Status.Active = false + tr.Status.Message = "queue is required for type=event" + _ = r.Status().Update(ctx, tr) + return ctrl.Result{}, nil + } + + deployNS := "sless-fn-" + tr.Namespace + + // Service нужен event-dispatcher для доставки сообщений в функцию по HTTP. + // Имя Service совпадает с именем Function — dispatcher строит URL как + // http://{functionRef}.{deployNS}.svc.cluster.local:8080/ + wantSvc := &corev1.Service{ + ObjectMeta: metav1.ObjectMeta{Name: fn.Name, Namespace: deployNS}, + Spec: corev1.ServiceSpec{ + Selector: map[string]string{"app": fn.Name}, + Ports: []corev1.ServicePort{{Port: 8080, Protocol: corev1.ProtocolTCP}}, + }, + } + existingSvc := &corev1.Service{} + if err := r.Get(ctx, client.ObjectKey{Name: fn.Name, Namespace: deployNS}, existingSvc); err != nil { + if errors.IsNotFound(err) { + if err := r.Create(ctx, wantSvc); err != nil { + return ctrl.Result{}, fmt.Errorf("create service for event trigger: %w", err) + } + } else { + return ctrl.Result{}, fmt.Errorf("get service: %w", err) + } + } + + tr.Status.Active = true + tr.Status.Message = fmt.Sprintf("listening on queue %q via event-dispatcher", tr.Spec.Queue) + _ = r.Status().Update(ctx, tr) + return ctrl.Result{}, nil +} diff --git a/deployments/k8s/event-dispatcher.yaml b/deployments/k8s/event-dispatcher.yaml new file mode 100644 index 0000000..75dfac1 --- /dev/null +++ b/deployments/k8s/event-dispatcher.yaml @@ -0,0 +1,77 @@ +# Изменено: 2026-03-19 +# event-dispatcher — отдельный сервис для обработки event-триггеров. +# Следит за Trigger CRD{type:event}, подписывается на AMQP очереди, +# при сообщении делает POST на внутренний HTTP endpoint функции. +# +# Требует: +# - SecretRef: sless-operator-secret (RABBITMQ_URL) +# - ClusterRole: event-dispatcher-role (чтение Trigger CRD, Namespace) +--- +apiVersion: v1 +kind: ServiceAccount +metadata: + name: event-dispatcher + namespace: sless +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: ClusterRole +metadata: + name: event-dispatcher-role +rules: + # Нужно читать Trigger CRD по всем namespace (event-dispatcher глобальный) + - apiGroups: ["sless.kube5s.ru"] + resources: ["triggers"] + verbs: ["get", "list", "watch"] + # Нужно читать namespace для построения URLs функций + - apiGroups: [""] + resources: ["namespaces"] + verbs: ["get", "list", "watch"] +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: ClusterRoleBinding +metadata: + name: event-dispatcher-rolebinding +subjects: + - kind: ServiceAccount + name: event-dispatcher + namespace: sless +roleRef: + kind: ClusterRole + name: event-dispatcher-role + apiGroup: rbac.authorization.k8s.io +--- +apiVersion: apps/v1 +kind: Deployment +metadata: + name: event-dispatcher + namespace: sless + labels: + app: event-dispatcher +spec: + replicas: 1 + selector: + matchLabels: + app: event-dispatcher + template: + metadata: + labels: + app: event-dispatcher + spec: + serviceAccountName: event-dispatcher + containers: + - name: event-dispatcher + image: naeel/sless-event-dispatcher:v0.1.0 + imagePullPolicy: Always + env: + - name: RABBITMQ_URL + valueFrom: + secretKeyRef: + name: sless-operator-secret + key: RABBITMQ_URL + resources: + requests: + cpu: 50m + memory: 64Mi + limits: + cpu: 200m + memory: 128Mi diff --git a/go.mod b/go.mod index 073603b..425ba09 100644 --- a/go.mod +++ b/go.mod @@ -54,6 +54,7 @@ require ( github.com/prometheus/client_model v0.3.0 // indirect github.com/prometheus/common v0.37.0 // indirect github.com/prometheus/procfs v0.8.0 // indirect + github.com/rabbitmq/amqp091-go v1.10.0 // indirect github.com/rs/xid v1.6.0 // indirect github.com/spf13/pflag v1.0.5 // indirect github.com/tinylib/msgp v1.6.1 // indirect diff --git a/go.sum b/go.sum index ce95ff4..5e3352e 100644 --- a/go.sum +++ b/go.sum @@ -270,6 +270,8 @@ github.com/prometheus/procfs v0.6.0/go.mod h1:cz+aTbrPOrUb4q7XlbU9ygM+/jj0fzG6c1 github.com/prometheus/procfs v0.7.3/go.mod h1:cz+aTbrPOrUb4q7XlbU9ygM+/jj0fzG6c1xBZuNvfVA= github.com/prometheus/procfs v0.8.0 h1:ODq8ZFEaYeCaZOJlZZdJA2AbQR98dSHSM1KW/You5mo= github.com/prometheus/procfs v0.8.0/go.mod h1:z7EfXMXOkbkqb9IINtpCn86r/to3BnA0uaxHdg830/4= +github.com/rabbitmq/amqp091-go v1.10.0 h1:STpn5XsHlHGcecLmMFCtg7mqq0RnD+zFr4uzukfVhBw= +github.com/rabbitmq/amqp091-go v1.10.0/go.mod h1:Hy4jKW5kQART1u+JkDTF9YYOQUHXqMuhrgxOEeS7G4o= github.com/rogpeppe/go-internal v1.3.0/go.mod h1:M8bDsm7K2OlrFYOpmOWEs/qY81heoFRclV5y23lUDJ4= github.com/rs/xid v1.6.0 h1:fV591PaemRlL6JfRxGDEPl69wICngIQ3shQtzfy2gxU= github.com/rs/xid v1.6.0/go.mod h1:7XoLgs4eV+QndskICGsho+ADou8ySMSjJKDIan90Nz0= @@ -305,6 +307,7 @@ go.uber.org/atomic v1.7.0/go.mod h1:fEN4uk6kAWBTFdckzkM89CLk9XfWZrxpCo0nPH17wJc= go.uber.org/goleak v1.1.10/go.mod h1:8a7PlsEVH3e/a/GLqe5IIrQx6GzcnRmZEufDUTk4A7A= go.uber.org/goleak v1.2.0 h1:xqgm/S+aQvhWFTtR0XK3Jvg7z8kGV8P4X14IzwN3Eqk= go.uber.org/goleak v1.2.0/go.mod h1:XJYK+MuIchqpmGmUSAzotztawfKvYLUIgg7guXrwVUo= +go.uber.org/goleak v1.3.0 h1:2K3zAYmnTNqV73imy9J1T3WC+gmCePx2hEGkimedGto= go.uber.org/multierr v1.6.0 h1:y6IPFStTAIT5Ytl7/XYmHvzXQ7S3g/IeZW9hyZ5thw4= go.uber.org/multierr v1.6.0/go.mod h1:cdWPpRnG4AhwMwsgIHip0KRBQjJy5kYEpYjJxpXp9iU= go.uber.org/zap v1.19.0/go.mod h1:xg/QME4nWcxGxrpdeYfq7UvYrLh66cuVKdrbD1XF/NI= diff --git a/internal/config/config.go b/internal/config/config.go index 6c7855c..ebc243e 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -1,4 +1,4 @@ -// Изменено: 2026-03-11 +// Изменено: 2026-03-19 // Конфигурация сервиса — читается из env переменных при старте. // Все компоненты (API, builder, runner) получают конфиг через эту структуру. // Используем env а не файлы конфигурации — стандарт для k8s (ConfigMap/Secret → env). @@ -56,6 +56,11 @@ type Config struct { // Если не заданы, EnsureProject пропускается (только DockerHub/без проектов). HarborUser string HarborPass string + + // RabbitMQURL — AMQP URL брокера для event-dispatcher. + // Формат: amqp://user:pass@host:5672/ + // Опционально — если не задан, event-триггеры не будут обрабатываться. + RabbitMQURL string } func Load() (*Config, error) { @@ -133,9 +138,12 @@ func Load() (*Config, error) { cfg.APIToken = os.Getenv("SLESS_API_TOKEN") // Harbor API креды — опциональны. - // Если не заданы, EnsureProject пропускается (push работает через docker-кред в kaniko). + // Если не заданы, EnsureProject пропускается (push работает через docker-кред в канiko). cfg.HarborUser = os.Getenv("HARBOR_USER") cfg.HarborPass = os.Getenv("HARBOR_PASS") + // RabbitMQ — опционально. Нужен только если используются event-триггеры. + cfg.RabbitMQURL = os.Getenv("RABBITMQ_URL") + return cfg, nil } diff --git a/services/event-dispatcher/dispatcher.go b/services/event-dispatcher/dispatcher.go new file mode 100644 index 0000000..63ddecd --- /dev/null +++ b/services/event-dispatcher/dispatcher.go @@ -0,0 +1,197 @@ +// Изменено: 2026-03-19 +// Dispatcher — управляет AMQP consumer-ами. +// Для каждого Trigger{type:event} держит один goroutine-consumer: +// - Subscribe(key, queue, targetURL) — открывает consumer на очередь +// - Unsubscribe(key) — закрывает consumer +// +// При получении сообщения из очереди: POST на targetURL с телом сообщения. +// 2xx → ack, остальное → nack (requeue=true). +// При потере AMQP соединения — переподключение с экспоненциальной задержкой. + +package main + +import ( + "bytes" + "context" + "fmt" + "net/http" + "sync" + "time" + + "github.com/go-logr/logr" + amqp "github.com/rabbitmq/amqp091-go" +) + +// consumerEntry — запись об активном consumer-е. +type consumerEntry struct { + queue string + targetURL string + cancel context.CancelFunc +} + +// Dispatcher управляет пулом AMQP consumer-ов (один на каждый event Trigger). +type Dispatcher struct { + rabbitURL string + logger logr.Logger + httpClient *http.Client + + mu sync.Mutex + consumers map[string]*consumerEntry // key = namespace/name триггера +} + +// NewDispatcher создаёт Dispatcher. +func NewDispatcher(rabbitURL string, logger logr.Logger) *Dispatcher { + return &Dispatcher{ + rabbitURL: rabbitURL, + logger: logger.WithName("dispatcher"), + httpClient: &http.Client{Timeout: 30 * time.Second}, + consumers: make(map[string]*consumerEntry), + } +} + +// Subscribe регистрирует consumer на очередь queue. +// key — уникальный идентификатор триггера (namespace/name). +// targetURL — внутренний HTTP URL функции (http://svc.ns.svc.cluster.local:8080/). +// Если consumer для этого key уже есть — он переподписывается с новыми параметрами. +func (d *Dispatcher) Subscribe(key, queue, targetURL string) { + d.mu.Lock() + defer d.mu.Unlock() + + // Если уже подписан — отменяем старый и создаём новый + if entry, ok := d.consumers[key]; ok { + if entry.queue == queue && entry.targetURL == targetURL { + return // ничего не изменилось + } + entry.cancel() + } + + ctx, cancel := context.WithCancel(context.Background()) + entry := &consumerEntry{queue: queue, targetURL: targetURL, cancel: cancel} + d.consumers[key] = entry + + go d.runConsumer(ctx, key, queue, targetURL) + d.logger.Info("subscribed", "key", key, "queue", queue, "target", targetURL) +} + +// Unsubscribe закрывает consumer для триггера key. +func (d *Dispatcher) Unsubscribe(key string) { + d.mu.Lock() + defer d.mu.Unlock() + + if entry, ok := d.consumers[key]; ok { + entry.cancel() + delete(d.consumers, key) + d.logger.Info("unsubscribed", "key", key) + } +} + +// runConsumer — горутина одного consumer-а. +// Держит AMQP соединение и channel. При разрыве переподключается. +// Завершается когда ctx отменён. +func (d *Dispatcher) runConsumer(ctx context.Context, key, queue, targetURL string) { + backoff := time.Second + for { + if ctx.Err() != nil { + return + } + err := d.consumeLoop(ctx, queue, targetURL) + if ctx.Err() != nil { + return // нормальное завершение + } + d.logger.Error(err, "consumer loop error, reconnecting", "key", key, "backoff", backoff) + select { + case <-ctx.Done(): + return + case <-time.After(backoff): + } + // Экспоненциальная задержка, максимум 30 секунд + if backoff < 30*time.Second { + backoff *= 2 + } + } +} + +// consumeLoop устанавливает соединение, объявляет очередь и читает сообщения. +// Возвращает ошибку при потере соединения (вызывающий перезапустит). +func (d *Dispatcher) consumeLoop(ctx context.Context, queue, targetURL string) error { + conn, err := amqp.Dial(d.rabbitURL) + if err != nil { + return fmt.Errorf("amqp dial: %w", err) + } + defer conn.Close() + + ch, err := conn.Channel() + if err != nil { + return fmt.Errorf("amqp channel: %w", err) + } + defer ch.Close() + + // Объявляем очередь как durable — она переживёт рестарт брокера + _, err = ch.QueueDeclare(queue, true, false, false, false, nil) + if err != nil { + return fmt.Errorf("queue declare %q: %w", queue, err) + } + + // prefetch=1: не берём следующее сообщение пока не обработали текущее + if err = ch.Qos(1, 0, false); err != nil { + return fmt.Errorf("qos: %w", err) + } + + msgs, err := ch.Consume(queue, "" /*auto-tag*/, false /*autoAck*/, false, false, false, nil) + if err != nil { + return fmt.Errorf("consume: %w", err) + } + + connClose := conn.NotifyClose(make(chan *amqp.Error, 1)) + d.logger.Info("consuming", "queue", queue, "target", targetURL) + + for { + select { + case <-ctx.Done(): + return nil + case amqpErr := <-connClose: + return fmt.Errorf("connection closed: %v", amqpErr) + case msg, ok := <-msgs: + if !ok { + return fmt.Errorf("messages channel closed") + } + d.handleMessage(ctx, msg, targetURL) + } + } +} + +// handleMessage выполняет HTTP POST на targetURL, ack/nack по результату. +func (d *Dispatcher) handleMessage(ctx context.Context, msg amqp.Delivery, targetURL string) { + logger := d.logger.WithValues("target", targetURL, "contentType", msg.ContentType) + + req, err := http.NewRequestWithContext(ctx, http.MethodPost, targetURL, bytes.NewReader(msg.Body)) + if err != nil { + logger.Error(err, "failed to build request, nacking") + _ = msg.Nack(false, true) + return + } + + contentType := msg.ContentType + if contentType == "" { + contentType = "application/octet-stream" + } + req.Header.Set("Content-Type", contentType) + req.ContentLength = int64(len(msg.Body)) + + resp, err := d.httpClient.Do(req) + if err != nil { + logger.Error(err, "http call failed, nacking") + _ = msg.Nack(false, true) + return + } + resp.Body.Close() + + if resp.StatusCode >= 200 && resp.StatusCode < 300 { + _ = msg.Ack(false) + logger.V(1).Info("message delivered", "status", resp.StatusCode) + } else { + // Функция вернула ошибку — requeue для повтора + logger.Info("function returned non-2xx, nacking", "status", resp.StatusCode) + _ = msg.Nack(false, true) + } +} diff --git a/services/event-dispatcher/main.go b/services/event-dispatcher/main.go new file mode 100644 index 0000000..42de5fc --- /dev/null +++ b/services/event-dispatcher/main.go @@ -0,0 +1,59 @@ +// Package main — точка входа event-dispatcher. +// Изменено: 2026-03-19 +// Читает RABBITMQ_URL из env, запускает Watcher (следит за Trigger CRD) +// и Dispatcher (держит AMQP consumer-ы, роутит сообщения в функции). +// Graceful shutdown по SIGTERM/SIGINT. + +package main + +import ( + "context" + "os" + "os/signal" + "syscall" + + "k8s.io/client-go/rest" + "k8s.io/client-go/tools/clientcmd" + ctrl "sigs.k8s.io/controller-runtime" + "sigs.k8s.io/controller-runtime/pkg/log/zap" +) + +func main() { + ctrl.SetLogger(zap.New(zap.UseDevMode(false))) + logger := ctrl.Log.WithName("event-dispatcher") + + rabbitURL := os.Getenv("RABBITMQ_URL") + if rabbitURL == "" { + logger.Error(nil, "RABBITMQ_URL is required") + os.Exit(1) + } + + // k8s config: в кластере — InClusterConfig, снаружи — KUBECONFIG + cfg, err := rest.InClusterConfig() + if err != nil { + // fallback для локальной разработки/тестов + kubeconfig := os.Getenv("KUBECONFIG") + cfg, err = clientcmd.BuildConfigFromFlags("", kubeconfig) + if err != nil { + logger.Error(err, "failed to build k8s config") + os.Exit(1) + } + } + + ctx, cancel := signal.NotifyContext(context.Background(), syscall.SIGTERM, syscall.SIGINT) + defer cancel() + + disp := NewDispatcher(rabbitURL, logger) + watcher, err := NewWatcher(cfg, disp, logger) + if err != nil { + logger.Error(err, "failed to create watcher") + os.Exit(1) + } + + logger.Info("starting event-dispatcher") + if err := watcher.Run(ctx); err != nil { + logger.Error(err, "watcher stopped with error") + os.Exit(1) + } + logger.Info("event-dispatcher stopped") +} diff --git a/services/event-dispatcher/watcher.go b/services/event-dispatcher/watcher.go new file mode 100644 index 0000000..185c86d --- /dev/null +++ b/services/event-dispatcher/watcher.go @@ -0,0 +1,129 @@ +// Изменено: 2026-03-19 +// Watcher — следит за Trigger CRD через k8s informer. +// При появлении Trigger{type:event} вызывает Dispatcher.Subscribe. +// При удалении — Dispatcher.Unsubscribe. +// При изменении (другая очередь или функция) — Subscribe автоматически переподключает. +// +// targetURL строится как: +// http://{functionRef}.{namespace}.svc.cluster.local:8080/ +// Это внутренний адрес Service функции — оператор создаёт его при reconcileEvent. + +package main + +import ( + "context" + "fmt" + + "github.com/go-logr/logr" + "k8s.io/apimachinery/pkg/runtime" + utilruntime "k8s.io/apimachinery/pkg/util/runtime" + clientgoscheme "k8s.io/client-go/kubernetes/scheme" + "k8s.io/client-go/rest" + ctrl "sigs.k8s.io/controller-runtime" + "sigs.k8s.io/controller-runtime/pkg/client" + + slessv1alpha1 "gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sless/api/v1alpha1" +) + +var watcherScheme = runtime.NewScheme() + +func init() { + utilruntime.Must(clientgoscheme.AddToScheme(watcherScheme)) + utilruntime.Must(slessv1alpha1.AddToScheme(watcherScheme)) +} + +// Watcher запускает controller-runtime manager который следит за Trigger CRD. +type Watcher struct { + manager ctrl.Manager + dispatcher *Dispatcher + logger logr.Logger +} + +// NewWatcher создаёт Watcher. +func NewWatcher(cfg *rest.Config, disp *Dispatcher, logger logr.Logger) (*Watcher, error) { + mgr, err := ctrl.NewManager(cfg, ctrl.Options{ + Scheme: watcherScheme, + MetricsBindAddress: "0", // метрики не нужны — это не оператор + LeaderElection: false, + }) + if err != nil { + return nil, fmt.Errorf("create manager: %w", err) + } + + w := &Watcher{ + manager: mgr, + dispatcher: disp, + logger: logger.WithName("watcher"), + } + + // Регистрируем reconciler для Trigger + if err := (&triggerEventReconciler{ + Client: mgr.GetClient(), + dispatcher: disp, + logger: w.logger, + }).SetupWithManager(mgr); err != nil { + return nil, fmt.Errorf("setup reconciler: %w", err) + } + + return w, nil +} + +// Run запускает manager (блокирует до ctx.Done()). +func (w *Watcher) Run(ctx context.Context) error { + return w.manager.Start(ctx) +} + +// triggerEventReconciler — controller-runtime reconciler только для Trigger{type:event}. +type triggerEventReconciler struct { + client.Client + dispatcher *Dispatcher + logger logr.Logger +} + +func (r *triggerEventReconciler) Reconcile(ctx context.Context, req ctrl.Request) (ctrl.Result, error) { + tr := &slessv1alpha1.Trigger{} + if err := r.Get(ctx, req.NamespacedName, tr); err != nil { + if client.IgnoreNotFound(err) == nil { + // Триггер удалён — убираем consumer + r.dispatcher.Unsubscribe(req.String()) + return ctrl.Result{}, nil + } + return ctrl.Result{}, err + } + + // Нас интересуют только event-триггеры + if tr.Spec.Type != slessv1alpha1.TriggerTypeEvent { + return ctrl.Result{}, nil + } + + // При удалении — убираем consumer + if !tr.DeletionTimestamp.IsZero() { + r.dispatcher.Unsubscribe(req.String()) + return ctrl.Result{}, nil + } + + // Триггер отключён — не подписываемся + if !tr.Spec.Enabled { + r.dispatcher.Unsubscribe(req.String()) + return ctrl.Result{}, nil + } + + if tr.Spec.Queue == "" { + r.logger.Info("event trigger has no queue, skipping", "trigger", req.String()) + return ctrl.Result{}, nil + } + + // Строим внутренний URL Service функции. + // Оператор при reconcileEvent создаёт Service с именем functionRef в namespace триггера. + targetURL := fmt.Sprintf("http://%s.%s.svc.cluster.local:8080/", + tr.Spec.FunctionRef, tr.Namespace) + + r.dispatcher.Subscribe(req.String(), tr.Spec.Queue, targetURL) + return ctrl.Result{}, nil +} + +func (r *triggerEventReconciler) SetupWithManager(mgr ctrl.Manager) error { + return ctrl.NewControllerManagedBy(mgr). + For(&slessv1alpha1.Trigger{}). + Complete(r) +}