Files
sless/services/event-dispatcher/watcher.go
T

130 lines
4.1 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
// Изменено: 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)
}