// Package dockerevents подписывается на события контейнеров Docker (raw HTTP, // без зависимости от docker SDK) и вызывает trigger как сигнал к немедленному reconcile. // События — только триггер; источник истины — Traefik API. package dockerevents import ( "bufio" "context" "encoding/json" "fmt" "log/slog" "net" "net/http" "net/url" "strings" "time" ) const eventsFilters = `{"type":["container"],"event":["start","die","destroy"]}` // Watcher читает поток /events. type Watcher struct { host string log *slog.Logger trigger func() // для тестов minBackoff, maxBackoff time.Duration } // New создаёт Watcher. host: unix:///path/to.sock или tcp://host:port. func New(host string, log *slog.Logger, trigger func()) *Watcher { return &Watcher{host: host, log: log, trigger: trigger, minBackoff: time.Second, maxBackoff: time.Minute} } func (w *Watcher) client() (*http.Client, string, error) { u, err := url.Parse(w.host) if err != nil { return nil, "", fmt.Errorf("docker.host: %w", err) } dialer := &net.Dialer{Timeout: 5 * time.Second} tr := &http.Transport{ResponseHeaderTimeout: 10 * time.Second} switch u.Scheme { case "unix": path := u.Path if path == "" { return nil, "", fmt.Errorf("docker.host: пустой путь сокета") } tr.DialContext = func(ctx context.Context, _, _ string) (net.Conn, error) { return dialer.DialContext(ctx, "unix", path) } return &http.Client{Transport: tr}, "http://docker", nil case "tcp", "http": tr.DialContext = dialer.DialContext return &http.Client{Transport: tr}, "http://" + u.Host, nil default: return nil, "", fmt.Errorf("docker.host: неподдерживаемая схема %q (unix|tcp)", u.Scheme) } } // Run блокируется до отмены ctx, переподключаясь при обрывах. Ошибки не фатальны. func (w *Watcher) Run(ctx context.Context) { hc, base, err := w.client() if err != nil { w.log.Error("Docker events отключены", "error", err) return } backoff := w.minBackoff failing := false for ctx.Err() == nil { connected, err := w.stream(ctx, hc, base) if ctx.Err() != nil { return } if connected { backoff = w.minBackoff failing = false } if !failing { w.log.Warn("Docker events недоступны, работаем по опросу; повторные попытки в фоне", "error", err) failing = true } else { w.log.Debug("Docker events: повторная попытка не удалась", "error", err) } select { case <-ctx.Done(): return case <-time.After(backoff): } if backoff *= 2; backoff > w.maxBackoff { backoff = w.maxBackoff } } } // stream возвращает connected=true, если соединение было установлено (ответ 200). func (w *Watcher) stream(ctx context.Context, hc *http.Client, base string) (bool, error) { q := url.Values{"filters": {eventsFilters}} req, err := http.NewRequestWithContext(ctx, http.MethodGet, base+"/events?"+q.Encode(), nil) if err != nil { return false, err } resp, err := hc.Do(req) if err != nil { return false, err } defer resp.Body.Close() if resp.StatusCode != http.StatusOK { return false, fmt.Errorf("docker events: статус %d", resp.StatusCode) } w.log.Info("подписка на Docker events установлена") sc := bufio.NewScanner(resp.Body) sc.Buffer(make([]byte, 64*1024), 1<<20) for sc.Scan() { line := strings.TrimSpace(sc.Text()) if line == "" { continue } var ev struct { Type string `json:"Type"` Action string `json:"Action"` } if err := json.Unmarshal([]byte(line), &ev); err != nil { w.log.Debug("Docker events: не удалось разобрать событие", "error", err) continue } w.log.Debug("событие Docker", "type", ev.Type, "action", ev.Action) w.trigger() } if err := sc.Err(); err != nil { return true, err } return true, fmt.Errorf("docker events: поток закрыт") }