// 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") }