Files
sless/services/event-dispatcher/main.go

60 lines
1.6 KiB
Go

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