// Изменено: 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) } }