feat: event-trigger (Вариант A) — TriggerTypeEvent, event-dispatcher, reconcileEvent
This commit is contained in:
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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")
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
Reference in New Issue
Block a user