MVP Single Node with Manual Config in Docker

This commit is contained in:
ayurishchev committed 2026-08-16 21:43:02 +03:00
commit 8212699f23
27 files changed
+3818

No files matched your search

+3
View File
@@ -0,0 +1,3 @@
module hpnn/hc
go 1.22
+168
View File
@@ -0,0 +1,168 @@
package main
import (
"crypto/sha256"
"encoding/binary"
"encoding/hex"
"fmt"
"hash/fnv"
"sort"
)
// Раскладка слотов по членам пула — §5 дизайн-концепции hpnn_v1.
//
// Свойства, ради которых взят именно Maglev-подобный алгоритм:
// - детерминированность: одинаковый вход даёт одинаковую раскладку на любом
// узле и между перезапусками;
// - равномерность в пределах ±1 слота от идеальной доли;
// - минимальное возмущение: при выбытии члена переезжают только его слоты,
// чужие остаются на месте.
//
// Ключ члена — "адрес:порт", а не идентификатор объекта: удаление и повторное
// создание члена с теми же адресом и портом обязано давать ту же раскладку
// (§5.1).
// hashKey — детерминированный хэш ключа члена с солью seed.
func hashKey(key string, seed uint32, salt uint32) uint32 {
h := fnv.New32a()
var buf [4]byte
binary.BigEndian.PutUint32(buf[:], seed)
_, _ = h.Write(buf[:])
binary.BigEndian.PutUint32(buf[:], salt)
_, _ = h.Write(buf[:])
_, _ = h.Write([]byte(key))
return h.Sum32()
}
// BuildSlotTable возвращает срез длиной slots, где каждый элемент — ID члена
// пула, обслуживающего этот слот. Пустой список членов даёт nil: вызывающая
// сторона трактует это как fail-close.
func BuildSlotTable(poolID uint32, members []*Member, slots int) []int {
live := make([]*Member, 0, len(members))
for _, m := range members {
if m.Active() {
live = append(live, m)
}
}
if len(live) == 0 || slots <= 0 {
return nil
}
sort.Slice(live, func(i, j int) bool { return live[i].Key() < live[j].Key() })
// Перестановка каждого члена: последовательность (offset + j*skip) mod M.
// Она обойдёт все M слотов, только если skip взаимно прост с M. При M —
// степени двойки (по умолчанию 1024) это означает «skip нечётный»: чётный
// шаг покрыл бы лишь половину слотов и заполнение зациклилось бы.
perm := make([][]int, len(live))
for i, m := range live {
offset := int(hashKey(m.Key(), poolID, 1) % uint32(slots))
skip := int(hashKey(m.Key(), poolID, 2)%uint32(slots-1)) + 1
if skip%2 == 0 {
skip++
}
p := make([]int, slots)
for j := 0; j < slots; j++ {
p[j] = (offset + j*skip) % slots
}
perm[i] = p
}
// Квоты по весам: минимальный вес получает не менее одного слота.
total := 0
for _, m := range live {
total += m.WeightOrDefault()
}
quota := make([]int, len(live))
assigned := 0
for i, m := range live {
quota[i] = slots * m.WeightOrDefault() / total
if quota[i] == 0 {
quota[i] = 1
}
assigned += quota[i]
}
// Остаток раздаём по кругу, чтобы сумма квот совпала с числом слотов.
for i := 0; assigned < slots; i = (i + 1) % len(live) {
quota[i]++
assigned++
}
table := make([]int, slots)
for i := range table {
table[i] = -1
}
next := make([]int, len(live))
filled := 0
for filled < slots {
progress := false
for i := range live {
if quota[i] == 0 {
continue
}
for next[i] < slots {
c := perm[i][next[i]]
next[i]++
if table[c] == -1 {
table[c] = live[i].ID
quota[i]--
filled++
progress = true
break
}
}
if filled == slots {
break
}
}
if !progress {
// Недостижимо при нечётном skip, но лучше выйти, чем зациклиться.
break
}
}
// Страховка: не покрытые слоты отдаём первому живому члену.
for i := range table {
if table[i] == -1 {
table[i] = live[0].ID
}
}
return table
}
// Digest — SHA-256 от сериализованной раскладки (§5.3 дизайна). Служит для
// сравнения раскладок между узлами и между перезапусками.
func Digest(table []int) string {
h := sha256.New()
var buf [4]byte
for _, id := range table {
binary.BigEndian.PutUint32(buf[:], uint32(id))
_, _ = h.Write(buf[:])
}
return hex.EncodeToString(h.Sum(nil))[:16]
}
// SlotCounts возвращает распределение слотов по членам пула.
func SlotCounts(table []int) map[int]int {
out := map[int]int{}
for _, id := range table {
out[id]++
}
return out
}
// FlowBundle рендерит директивы для ovs-ofctl bundle: замена таблицы слотов
// целиком одной атомарной транзакцией. Первая строка сносит прежнее
// содержимое таблицы, остальные заливают новое — датапас не проходит через
// состояние с полупустой раскладкой.
//
// Формат файла бандла: каждая строка начинается с типа сообщения ("flow") и
// команды ("add" / "delete" / "modify").
func FlowBundle(table []int, slotTable int, nextTable int) string {
out := fmt.Sprintf("flow delete table=%d\n", slotTable)
// Основание fail-close: пока слотов нет, трафик на VIP отбрасывается.
out += fmt.Sprintf("flow add table=%d,priority=0 actions=drop\n", slotTable)
for slot, id := range table {
out += fmt.Sprintf("flow add table=%d,priority=100,ip,reg1=0x%x actions=load:0x%x->NXM_NX_REG2[],goto_table:%d\n",
slotTable, slot, id, nextTable)
}
return out
}
+505
View File
@@ -0,0 +1,505 @@
// 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"`
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"`
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
}
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"`
}
// --- пробы -------------------------------------------------------------------
// 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,
}
if !m.lastChange.IsZero() {
v.Since = m.lastChange.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)
}
}
})
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: перезаливка всего пайплайна стирает таблицу
// слотов, поэтому её нужно восстановить в актуальном составе.
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.mu.Lock()
digest := a.digest
a.mu.Unlock()
fmt.Fprintf(w, "таблица слотов перезалита, дайджест %s\n", digest)
})
log.Printf("API на http://%s (/status, /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()
for _, c := range changed {
if c {
if err := agent.apply(false); err != nil {
log.Printf("ошибка применения: %v", err)
}
break
}
}
}
}
}