60 lines
1.6 KiB
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")
|
|
}
|