Files
hpnn-proto/hc/main.go
T

547 lines
18 KiB
Go
Raw Normal View History

// hcd — подсистема health-check узла балансировщика (§7 дизайн-концепции
// hpnn_v1, адаптированная к одноузловому прототипу).
//
// Демон проверяет живость членов пула, ведёт их состояние с гистерезисом и
// отражает состав пула в датапасе: пересчитывает Maglev-раскладку слотов и
// заливает таблицу слотов OVS одной атомарной транзакцией.
//
// Пробы уходят с уникального адреса узла в сегменте бэкенда, а не с VIP —
// иначе ответ вернулся бы не тому, кто проверял (§7.1).
package main
import (
"context"
"encoding/json"
"flag"
"fmt"
"log"
"net"
"net/http"
"os"
"os/exec"
"os/signal"
"sort"
"sync"
"syscall"
"time"
)
// --- конфигурация ------------------------------------------------------------
type Config struct {
Bridge string `json:"bridge"`
OFVersion string `json:"of_version"`
PoolID uint32 `json:"pool_id"`
Slots int `json:"slots"`
SlotTable int `json:"slot_table"`
DNATTable int `json:"dnat_table"`
AdjTable int `json:"adj_table"` // таблица 21 — next-hop MAC членов (шаг 5)
Probe string `json:"probe"`
HTTPPath string `json:"http_path"`
Interval string `json:"interval"`
Timeout string `json:"timeout"`
Rise int `json:"rise"`
Fall int `json:"fall"`
API string `json:"api"`
ApplyScript string `json:"apply_script"`
Members []*Member `json:"members"`
interval time.Duration
timeout time.Duration
}
type Member struct {
ID int `json:"id"`
Name string `json:"name"`
Address string `json:"address"`
Port int `json:"port"`
Weight int `json:"weight"`
Source string `json:"source"`
Iface string `json:"iface"` // hcif-порт, на котором резолвится MAC (шаг 5)
OFPort int `json:"ofport"` // выходной порт члена в OpenFlow
mu sync.Mutex
state string // up | down
admin string // enabled | drain
okStreak int
failStreak int
lastChange time.Time
lastErr string
lastLatency time.Duration
probes uint64
failures uint64
transitions uint64
mac string // резолвлен hcd через neigh-таблицу ядра, пусто = не резолвлен
macSince time.Time
}
func (m *Member) Key() string { return fmt.Sprintf("%s:%d", m.Address, m.Port) }
func (m *Member) WeightOrDefault() int {
if m.Weight <= 0 {
return 1
}
return m.Weight
}
// Active — член участвует в раскладке слотов: жив и не выведен на дренаж.
func (m *Member) Active() bool {
m.mu.Lock()
defer m.mu.Unlock()
return m.state == "up" && m.admin == "enabled"
}
type memberView struct {
Name string `json:"name"`
Address string `json:"address"`
Port int `json:"port"`
State string `json:"state"`
Admin string `json:"admin"`
Active bool `json:"active"`
Since string `json:"since"`
LastError string `json:"last_error,omitempty"`
LatencyMS int64 `json:"latency_ms"`
Probes uint64 `json:"probes"`
Failures uint64 `json:"failures"`
Transitions uint64 `json:"transitions"`
Slots int `json:"slots"`
Source string `json:"probe_source"`
Iface string `json:"iface"`
MAC string `json:"mac"`
MACSince string `json:"mac_since,omitempty"`
}
// --- пробы -------------------------------------------------------------------
// probe возвращает nil, если член ответил. Источник соединения жёстко
// привязан к адресу узла в сегменте этого члена.
func probe(ctx context.Context, cfg *Config, m *Member) error {
dialer := &net.Dialer{
Timeout: cfg.timeout,
LocalAddr: &net.TCPAddr{IP: net.ParseIP(m.Source)},
}
target := fmt.Sprintf("%s:%d", m.Address, m.Port)
if cfg.Probe == "tcp" {
conn, err := dialer.DialContext(ctx, "tcp", target)
if err != nil {
return err
}
return conn.Close()
}
client := &http.Client{
Timeout: cfg.timeout,
Transport: &http.Transport{
DialContext: dialer.DialContext,
DisableKeepAlives: true, // каждая проба — новое соединение
},
}
req, err := http.NewRequestWithContext(ctx, http.MethodGet,
fmt.Sprintf("http://%s%s", target, cfg.HTTPPath), nil)
if err != nil {
return err
}
resp, err := client.Do(req)
if err != nil {
return err
}
defer resp.Body.Close()
if resp.StatusCode < 200 || resp.StatusCode > 299 {
return fmt.Errorf("код ответа %d", resp.StatusCode)
}
return nil
}
// observe применяет результат пробы с гистерезисом и возвращает true, если
// состояние члена изменилось.
func observe(m *Member, cfg *Config, err error, latency time.Duration) bool {
m.mu.Lock()
defer m.mu.Unlock()
m.probes++
m.lastLatency = latency
changed := false
if err == nil {
m.lastErr = ""
m.failStreak = 0
m.okStreak++
if m.state != "up" && m.okStreak >= cfg.Rise {
m.state = "up"
m.lastChange = time.Now()
m.transitions++
changed = true
log.Printf("член %s (%s) -> up после %d успешных проб", m.Name, m.Key(), m.okStreak)
}
} else {
m.failures++
m.lastErr = err.Error()
m.okStreak = 0
m.failStreak++
if m.state != "down" && m.failStreak >= cfg.Fall {
m.state = "down"
m.lastChange = time.Now()
m.transitions++
changed = true
log.Printf("член %s (%s) -> down после %d неудач: %v", m.Name, m.Key(), m.failStreak, err)
}
}
return changed
}
// --- применение раскладки ----------------------------------------------------
type Agent struct {
cfg *Config
mu sync.Mutex
table []int
digest string
applies uint64
errors uint64
}
// apply пересчитывает раскладку и, если она изменилась, заливает таблицу
// слотов атомарным бандлом. force заставляет залить даже неизменившуюся —
// нужно после перезаливки всего пайплайна, которая стирает слоты.
func (a *Agent) apply(force bool) error {
table := BuildSlotTable(a.cfg.PoolID, a.cfg.Members, a.cfg.Slots)
digest := Digest(table)
a.mu.Lock()
unchanged := digest == a.digest && !force
a.mu.Unlock()
if unchanged {
return nil
}
bundle := FlowBundle(table, a.cfg.SlotTable, a.cfg.DNATTable)
path := "/var/run/openvswitch/slots.bundle"
if err := os.WriteFile(path, []byte(bundle), 0o644); err != nil {
return err
}
// ovs-ofctl bundle выполняет delete+add как одну транзакцию: датапас не
// проходит через состояние с полупустой таблицей слотов.
out, err := exec.Command("ovs-ofctl", "-O", a.cfg.OFVersion, "bundle", a.cfg.Bridge, path).CombinedOutput()
a.mu.Lock()
defer a.mu.Unlock()
if err != nil {
a.errors++
return fmt.Errorf("ovs-ofctl bundle: %v: %s", err, out)
}
a.table = table
a.digest = digest
a.applies++
counts := SlotCounts(table)
parts := make([]string, 0, len(counts))
for _, m := range a.cfg.Members {
if n, ok := counts[m.ID]; ok {
parts = append(parts, fmt.Sprintf("%s=%d", m.Name, n))
}
}
sort.Strings(parts)
if len(table) == 0 {
log.Printf("раскладка применена: живых членов нет, fail-close (трафик на VIP отбрасывается)")
} else {
log.Printf("раскладка применена: слоты %v, дайджест %s", parts, digest)
}
return nil
}
func (a *Agent) slotsOf(id int) int {
a.mu.Lock()
defer a.mu.Unlock()
return SlotCounts(a.table)[id]
}
// --- HTTP API ----------------------------------------------------------------
func (a *Agent) memberByName(name string) *Member {
for _, m := range a.cfg.Members {
if m.Name == name {
return m
}
}
return nil
}
func (a *Agent) setAdmin(w http.ResponseWriter, r *http.Request, admin string) {
name := r.URL.Query().Get("member")
m := a.memberByName(name)
if m == nil {
http.Error(w, fmt.Sprintf("член %q не найден\n", name), http.StatusNotFound)
return
}
m.mu.Lock()
prev := m.admin
m.admin = admin
m.mu.Unlock()
if prev != admin {
log.Printf("член %s: admin %s -> %s", m.Name, prev, admin)
}
if err := a.apply(false); err != nil {
log.Printf("ошибка применения: %v", err)
}
fmt.Fprintf(w, "член %s: admin=%s\n", m.Name, admin)
}
func (a *Agent) status() []memberView {
out := make([]memberView, 0, len(a.cfg.Members))
for _, m := range a.cfg.Members {
m.mu.Lock()
v := memberView{
Name: m.Name, Address: m.Address, Port: m.Port,
State: m.state, Admin: m.admin,
Active: m.state == "up" && m.admin == "enabled",
LastError: m.lastErr, LatencyMS: m.lastLatency.Milliseconds(),
Probes: m.probes, Failures: m.failures, Transitions: m.transitions,
Source: m.Source, Iface: m.Iface, MAC: m.mac,
}
if !m.lastChange.IsZero() {
v.Since = m.lastChange.Format(time.RFC3339)
}
if !m.macSince.IsZero() {
v.MACSince = m.macSince.Format(time.RFC3339)
}
m.mu.Unlock()
v.Slots = a.slotsOf(m.ID)
out = append(out, v)
}
return out
}
func (a *Agent) serve() {
mux := http.NewServeMux()
mux.HandleFunc("/status", func(w http.ResponseWriter, r *http.Request) {
a.mu.Lock()
digest := a.digest
a.mu.Unlock()
w.Header().Set("Content-Type", "application/json")
_ = json.NewEncoder(w).Encode(map[string]any{
"pool_id": a.cfg.PoolID,
"slots": a.cfg.Slots,
"slot_table_digest": digest,
"probe": a.cfg.Probe,
"interval": a.cfg.Interval,
"rise": a.cfg.Rise,
"fall": a.cfg.Fall,
"members": a.status(),
})
})
// Человекочитаемый вид для lbctl health: форматирование здесь избавляет
// образ балансировщика от зависимости на jq или python.
mux.HandleFunc("/status.txt", func(w http.ResponseWriter, r *http.Request) {
a.mu.Lock()
digest := a.digest
a.mu.Unlock()
if digest == "" {
digest = "-"
}
w.Header().Set("Content-Type", "text/plain; charset=utf-8")
fmt.Fprintf(w, "пул %d: слотов %d, проба %s каждые %s (rise=%d fall=%d)\n",
a.cfg.PoolID, a.cfg.Slots, a.cfg.Probe, a.cfg.Interval, a.cfg.Rise, a.cfg.Fall)
fmt.Fprintf(w, "дайджест раскладки: %s\n\n", digest)
fmt.Fprintf(w, "%-6s %-17s %-10s %-9s %-7s %7s %7s %7s %5s %s\n",
"ЧЛЕН", "АДРЕС", "СОСТОЯНИЕ", "ADMIN", "В ПУЛЕ", "СЛОТОВ", "ПРОБ", "НЕУДАЧ", "МС", "ИСТОЧНИК ПРОБ")
for _, v := range a.status() {
inPool := "нет"
if v.Active {
inPool = "да"
}
fmt.Fprintf(w, "%-6s %-17s %-10s %-9s %-7s %7d %7d %7d %5d %s\n",
v.Name, fmt.Sprintf("%s:%d", v.Address, v.Port), v.State, v.Admin, inPool,
v.Slots, v.Probes, v.Failures, v.LatencyMS, v.Source)
if v.LastError != "" {
fmt.Fprintf(w, " последняя ошибка: %s\n", v.LastError)
}
}
})
// Резолв MAC членов пула (шаг 5): человекочитаемая таблица для lbctl neigh.
mux.HandleFunc("/neigh.txt", func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "text/plain; charset=utf-8")
fmt.Fprintf(w, "%-6s %-17s %-10s %-17s %-8s %s\n",
"ЧЛЕН", "АДРЕС", "ИНТЕРФЕЙС", "MAC", "РЕЗОЛВЛЕН", "С МОМЕНТА")
for _, v := range a.status() {
mac := v.MAC
resolved := "нет"
if mac == "" {
mac = "-"
} else {
resolved = "да"
}
since := v.MACSince
if since == "" {
since = "-"
}
fmt.Fprintf(w, "%-6s %-17s %-10s %-17s %-8s %s\n",
v.Name, v.Address, v.Iface, mac, resolved, since)
}
})
mux.HandleFunc("/metrics", func(w http.ResponseWriter, r *http.Request) {
a.mu.Lock()
applies, errs := a.applies, a.errors
a.mu.Unlock()
w.Header().Set("Content-Type", "text/plain; version=0.0.4")
fmt.Fprintf(w, "# HELP hpnn_member_up Член пула жив по данным проб\n# TYPE hpnn_member_up gauge\n")
for _, v := range a.status() {
up := 0
if v.State == "up" {
up = 1
}
fmt.Fprintf(w, "hpnn_member_up{member=%q,address=%q} %d\n", v.Name, v.Address, up)
}
fmt.Fprintf(w, "# HELP hpnn_member_active Член участвует в раскладке слотов\n# TYPE hpnn_member_active gauge\n")
for _, v := range a.status() {
act := 0
if v.Active {
act = 1
}
fmt.Fprintf(w, "hpnn_member_active{member=%q} %d\n", v.Name, act)
}
fmt.Fprintf(w, "# HELP hpnn_member_slots Слотов у члена пула\n# TYPE hpnn_member_slots gauge\n")
for _, v := range a.status() {
fmt.Fprintf(w, "hpnn_member_slots{member=%q} %d\n", v.Name, v.Slots)
}
fmt.Fprintf(w, "# HELP hpnn_probes_total Всего проб\n# TYPE hpnn_probes_total counter\n")
for _, v := range a.status() {
fmt.Fprintf(w, "hpnn_probes_total{member=%q} %d\n", v.Name, v.Probes)
}
fmt.Fprintf(w, "# HELP hpnn_probe_failures_total Неудачных проб\n# TYPE hpnn_probe_failures_total counter\n")
for _, v := range a.status() {
fmt.Fprintf(w, "hpnn_probe_failures_total{member=%q} %d\n", v.Name, v.Failures)
}
fmt.Fprintf(w, "# HELP hpnn_probe_latency_ms Задержка последней пробы\n# TYPE hpnn_probe_latency_ms gauge\n")
for _, v := range a.status() {
fmt.Fprintf(w, "hpnn_probe_latency_ms{member=%q} %d\n", v.Name, v.LatencyMS)
}
fmt.Fprintf(w, "# HELP hpnn_slot_table_applies_total Заливок таблицы слотов\n# TYPE hpnn_slot_table_applies_total counter\nhpnn_slot_table_applies_total %d\n", applies)
fmt.Fprintf(w, "# HELP hpnn_slot_table_errors_total Ошибок заливки\n# TYPE hpnn_slot_table_errors_total counter\nhpnn_slot_table_errors_total %d\n", errs)
})
mux.HandleFunc("/drain", func(w http.ResponseWriter, r *http.Request) { a.setAdmin(w, r, "drain") })
mux.HandleFunc("/enable", func(w http.ResponseWriter, r *http.Request) { a.setAdmin(w, r, "enabled") })
// Вызывается из apply.sh: перезаливка всего пайплайна (replace-flows)
// стирает и таблицу слотов, и правила adjacency (таблица 21) — обе нужно
// восстановить в актуальном составе. force=true у обеих: MAC/раскладка в
// памяти демона могли не измениться, но датапас сброшен и нуждается в
// повторной заливке уже известных значений.
mux.HandleFunc("/reapply", func(w http.ResponseWriter, r *http.Request) {
if err := a.apply(true); err != nil {
http.Error(w, err.Error()+"\n", http.StatusInternalServerError)
return
}
a.syncNeighbors(true)
a.mu.Lock()
digest := a.digest
a.mu.Unlock()
fmt.Fprintf(w, "таблица слотов перезалита, дайджест %s\n", digest)
})
log.Printf("API на http://%s (/status, /neigh.txt, /metrics, /drain, /enable, /reapply)", a.cfg.API)
if err := http.ListenAndServe(a.cfg.API, mux); err != nil {
log.Fatalf("API: %v", err)
}
}
// --- запуск ------------------------------------------------------------------
func loadConfig(path string) (*Config, error) {
raw, err := os.ReadFile(path)
if err != nil {
return nil, err
}
cfg := &Config{}
if err := json.Unmarshal(raw, cfg); err != nil {
return nil, err
}
if cfg.interval, err = time.ParseDuration(cfg.Interval); err != nil {
return nil, fmt.Errorf("interval: %w", err)
}
if cfg.timeout, err = time.ParseDuration(cfg.Timeout); err != nil {
return nil, fmt.Errorf("timeout: %w", err)
}
if len(cfg.Members) == 0 {
return nil, fmt.Errorf("пустой список членов пула")
}
for _, m := range cfg.Members {
// Стартуем с down: член войдёт в раскладку только после Rise
// успешных проб. Балансировать на непроверенный бэкенд нельзя.
m.state = "down"
m.admin = "enabled"
}
return cfg, nil
}
func main() {
path := flag.String("config", "/var/run/openvswitch/hc.json", "путь к конфигурации")
flag.Parse()
log.SetFlags(log.LstdFlags)
log.SetPrefix("[hcd] ")
cfg, err := loadConfig(*path)
if err != nil {
log.Fatalf("конфигурация: %v", err)
}
log.Printf("пул %d: %d членов, %d слотов, проба %s каждые %s (rise=%d fall=%d)",
cfg.PoolID, len(cfg.Members), cfg.Slots, cfg.Probe, cfg.Interval, cfg.Rise, cfg.Fall)
agent := &Agent{cfg: cfg}
// Стартовое состояние — пустой пул: до первых успешных проб датапас
// работает в fail-close.
if err := agent.apply(true); err != nil {
log.Printf("ошибка стартового применения: %v", err)
}
go agent.serve()
stop := make(chan os.Signal, 1)
signal.Notify(stop, syscall.SIGTERM, syscall.SIGINT)
ticker := time.NewTicker(cfg.interval)
defer ticker.Stop()
for {
select {
case <-stop:
log.Printf("остановка")
return
case <-ticker.C:
var wg sync.WaitGroup
changed := make([]bool, len(cfg.Members))
for i, m := range cfg.Members {
wg.Add(1)
go func(i int, m *Member) {
defer wg.Done()
ctx, cancel := context.WithTimeout(context.Background(), cfg.timeout)
defer cancel()
start := time.Now()
err := probe(ctx, cfg, m)
changed[i] = observe(m, cfg, err, time.Since(start))
}(i, m)
}
wg.Wait()
// Резолв next-hop MAC членов пула — на том же тикере, что и пробы:
// именно проба вызывает первый исходящий пакет с hcif-порта и тем
// самым запускает ARP ядра для ещё не резолвленных адресов.
agent.syncNeighbors(false)
for _, c := range changed {
if c {
if err := agent.apply(false); err != nil {
log.Printf("ошибка применения: %v", err)
}
break
}
}
}
}
}