Compare commits
4 commits
a1df88a1f2
...
cdbb5012b1
| Author | SHA1 | Date | |
|---|---|---|---|
|
cdbb5012b1 |
|||
|
1ce84ffd58 |
|||
|
d8ece59d1c |
|||
|
2ce7ebafe7 |
22 changed files with 483 additions and 91 deletions
|
|
@ -1,15 +1,22 @@
|
|||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"flag"
|
||||
"fmt"
|
||||
"log/slog"
|
||||
"os"
|
||||
"os/signal"
|
||||
"syscall"
|
||||
"time"
|
||||
|
||||
agentapi "git.g3e.fr/syonad/two/internal/api/agent"
|
||||
configuration "git.g3e.fr/syonad/two/internal/config/agent"
|
||||
dispatcher "git.g3e.fr/syonad/two/internal/dispatcher/agent"
|
||||
"git.g3e.fr/syonad/two/internal/migration"
|
||||
agentmetrics "git.g3e.fr/syonad/two/internal/prometheus/agent"
|
||||
"git.g3e.fr/syonad/two/internal/watchdog"
|
||||
"git.g3e.fr/syonad/two/internal/watchdog/notify"
|
||||
"git.g3e.fr/syonad/two/pkg/db/kv"
|
||||
"git.g3e.fr/syonad/two/pkg/logger"
|
||||
promserver "git.g3e.fr/syonad/two/pkg/prometheus"
|
||||
|
|
@ -17,6 +24,8 @@ import (
|
|||
"github.com/prometheus/client_golang/prometheus"
|
||||
)
|
||||
|
||||
const shutdownTimeout = 20 * time.Second
|
||||
|
||||
func main() {
|
||||
confFile := flag.String("config", "/etc/two/agent.yml", "config file path")
|
||||
flag.Parse()
|
||||
|
|
@ -30,7 +39,12 @@ func main() {
|
|||
log := logger.New(cfg.Logger.Level, cfg.Logger.Debug)
|
||||
|
||||
db := kv.InitDB(kv.Config{Path: cfg.Database.Path}, false)
|
||||
defer db.Close()
|
||||
closeDB := true
|
||||
defer func() {
|
||||
if closeDB {
|
||||
db.Close()
|
||||
}
|
||||
}()
|
||||
|
||||
// Avant tout démarrage de service : la DB peut porter l'ancien vocabulaire
|
||||
// d'états, et des ressources transitoires orphelines d'un arrêt précédent.
|
||||
|
|
@ -56,13 +70,69 @@ func main() {
|
|||
"debug", cfg.Logger.Debug,
|
||||
)
|
||||
|
||||
ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
|
||||
defer stop()
|
||||
|
||||
d := dispatcher.New(q, db, cfg, log.With(slog.String("component", "dispatcher")))
|
||||
go agentapi.New(d, db, log.With(slog.String("component", "api"))).Start(apiAddr)
|
||||
go promserver.Start(promAddr, registry)
|
||||
|
||||
apiSrv := agentapi.New(d, db, log.With(slog.String("component", "api")), apiAddr)
|
||||
go apiSrv.Start()
|
||||
|
||||
promSrv := promserver.New(promAddr, registry)
|
||||
go promSrv.Start()
|
||||
|
||||
var adminSrv *kv.AdminServer
|
||||
if cfg.Admin.Enabled {
|
||||
adminAddr := fmt.Sprintf("%s:%d", cfg.Admin.Address, cfg.Admin.Port)
|
||||
go kv.NewAdminServer(db, log.With(slog.String("component", "admin"))).Start(adminAddr)
|
||||
adminSrv = kv.NewAdminServer(db, log.With(slog.String("component", "admin")), adminAddr)
|
||||
go adminSrv.Start()
|
||||
}
|
||||
|
||||
select {}
|
||||
if cfg.Watchdog.Enabled {
|
||||
wlog := log.With(slog.String("component", "watchdog"))
|
||||
go watchdog.New(db, cfg, notify.NewStderr(wlog), wlog,
|
||||
time.Duration(cfg.Watchdog.IntervalSeconds)*time.Second,
|
||||
).Run(ctx)
|
||||
}
|
||||
|
||||
<-ctx.Done()
|
||||
stop()
|
||||
|
||||
servers := map[string]httpShutdowner{"api": apiSrv, "prometheus": promSrv}
|
||||
if adminSrv != nil {
|
||||
servers["admin"] = adminSrv
|
||||
}
|
||||
closeDB = shutdown(log, q, servers, shutdownTimeout)
|
||||
}
|
||||
|
||||
type httpShutdowner interface {
|
||||
Shutdown(context.Context) error
|
||||
}
|
||||
|
||||
func shutdown(log *slog.Logger, q *worker.Queue, servers map[string]httpShutdowner, timeout time.Duration) bool {
|
||||
log.Info("shutting down", "timeout", timeout)
|
||||
|
||||
ctx, cancel := context.WithTimeout(context.Background(), timeout)
|
||||
defer cancel()
|
||||
|
||||
for name, srv := range servers {
|
||||
if err := srv.Shutdown(ctx); err != nil {
|
||||
log.Error("http server shutdown", "server", name, "error", err)
|
||||
}
|
||||
}
|
||||
|
||||
drained := make(chan struct{})
|
||||
go func() {
|
||||
q.Stop()
|
||||
close(drained)
|
||||
}()
|
||||
|
||||
select {
|
||||
case <-drained:
|
||||
log.Info("workers drained")
|
||||
return true
|
||||
case <-ctx.Done():
|
||||
log.Error("workers still running after timeout, leaving database untouched")
|
||||
return false
|
||||
}
|
||||
}
|
||||
|
|
|
|||
98
cmd/agent/main_test.go
Normal file
98
cmd/agent/main_test.go
Normal file
|
|
@ -0,0 +1,98 @@
|
|||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"io"
|
||||
"log/slog"
|
||||
"strings"
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"git.g3e.fr/syonad/two/pkg/worker"
|
||||
)
|
||||
|
||||
type fakeServer struct {
|
||||
called atomic.Bool
|
||||
err error
|
||||
}
|
||||
|
||||
func (f *fakeServer) Shutdown(context.Context) error {
|
||||
f.called.Store(true)
|
||||
return f.err
|
||||
}
|
||||
|
||||
func discardLogger() *slog.Logger {
|
||||
return slog.New(slog.NewTextHandler(io.Discard, nil))
|
||||
}
|
||||
|
||||
func TestShutdown_DrainReussi(t *testing.T) {
|
||||
q := worker.New(10)
|
||||
q.Start(2)
|
||||
|
||||
var done atomic.Int32
|
||||
for range 3 {
|
||||
q.Submit(func() {
|
||||
time.Sleep(20 * time.Millisecond)
|
||||
done.Add(1)
|
||||
})
|
||||
}
|
||||
|
||||
api, prom := &fakeServer{}, &fakeServer{}
|
||||
servers := map[string]httpShutdowner{"api": api, "prometheus": prom}
|
||||
|
||||
if !shutdown(discardLogger(), q, servers, 5*time.Second) {
|
||||
t.Fatal("un drainage réussi doit autoriser la fermeture de la base")
|
||||
}
|
||||
if !api.called.Load() || !prom.called.Load() {
|
||||
t.Error("tous les serveurs HTTP doivent être arrêtés")
|
||||
}
|
||||
if got := done.Load(); got != 3 {
|
||||
t.Errorf("les 3 tâches devaient se terminer, %d terminées", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestShutdown_TimeoutLaisseLaBaseIntacte(t *testing.T) {
|
||||
q := worker.New(10)
|
||||
q.Start(1)
|
||||
q.Submit(func() { time.Sleep(2 * time.Second) })
|
||||
|
||||
var buf strings.Builder
|
||||
log := slog.New(slog.NewTextHandler(&buf, nil))
|
||||
|
||||
if shutdown(log, q, map[string]httpShutdowner{}, 50*time.Millisecond) {
|
||||
t.Fatal("un drainage incomplet ne doit pas autoriser la fermeture de la base")
|
||||
}
|
||||
if !strings.Contains(buf.String(), "leaving database untouched") {
|
||||
t.Errorf("le dépassement devrait être logué, obtenu %q", buf.String())
|
||||
}
|
||||
}
|
||||
|
||||
func TestShutdown_ErreurServeurNEmpechePasLeDrainage(t *testing.T) {
|
||||
q := worker.New(10)
|
||||
q.Start(1)
|
||||
|
||||
var buf strings.Builder
|
||||
log := slog.New(slog.NewTextHandler(&buf, nil))
|
||||
servers := map[string]httpShutdowner{
|
||||
"api": &fakeServer{err: errors.New("boom")},
|
||||
"prometheus": &fakeServer{},
|
||||
}
|
||||
|
||||
if !shutdown(log, q, servers, 5*time.Second) {
|
||||
t.Fatal("une erreur d'arrêt HTTP ne doit pas empêcher le drainage")
|
||||
}
|
||||
if !strings.Contains(buf.String(), "http server shutdown") {
|
||||
t.Errorf("l'erreur devrait être loguée, obtenu %q", buf.String())
|
||||
}
|
||||
}
|
||||
|
||||
func TestShutdown_SansServeur(t *testing.T) {
|
||||
q := worker.New(10)
|
||||
q.Start(1)
|
||||
|
||||
if !shutdown(discardLogger(), q, map[string]httpShutdowner{}, 5*time.Second) {
|
||||
t.Fatal("l'absence de serveur ne doit pas empêcher un arrêt propre")
|
||||
}
|
||||
}
|
||||
|
|
@ -51,6 +51,13 @@ qemu:
|
|||
monitor_dir: "/run/two/vms/monitor"
|
||||
qmp_dir: "/run/two/vms/qmp"
|
||||
|
||||
# Consistency watchdog: periodically checks that resources marked running in the
|
||||
# database still exist on the system, and reports the gaps. Read-only, never repairs.
|
||||
# A persisting gap is reported at every tick (no deduplication).
|
||||
watchdog:
|
||||
enabled: true
|
||||
interval_seconds: 60
|
||||
|
||||
# Admin API (read-only DB inspection, loopback only)
|
||||
admin:
|
||||
enabled: false
|
||||
|
|
|
|||
|
|
@ -23,5 +23,5 @@ func newTestServer(t *testing.T) (*Server, *badger.DB) {
|
|||
cfg := &configuration.Config{DefaultInterface: "br-test"}
|
||||
logger := slog.New(slog.NewTextHandler(io.Discard, nil))
|
||||
d := dispatcher.New(q, db, cfg, logger)
|
||||
return New(d, db, logger), db
|
||||
return New(d, db, logger, "127.0.0.1:0"), db
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,8 +1,10 @@
|
|||
package agentapi
|
||||
|
||||
import (
|
||||
"context"
|
||||
"crypto/rand"
|
||||
"encoding/hex"
|
||||
"errors"
|
||||
"log/slog"
|
||||
"net/http"
|
||||
"time"
|
||||
|
|
@ -15,13 +17,11 @@ type Server struct {
|
|||
dispatcher *dispatcher.Dispatcher
|
||||
db *badger.DB
|
||||
logger *slog.Logger
|
||||
srv *http.Server
|
||||
}
|
||||
|
||||
func New(d *dispatcher.Dispatcher, db *badger.DB, logger *slog.Logger) *Server {
|
||||
return &Server{dispatcher: d, db: db, logger: logger}
|
||||
}
|
||||
|
||||
func (s *Server) Start(address string) {
|
||||
func New(d *dispatcher.Dispatcher, db *badger.DB, logger *slog.Logger, address string) *Server {
|
||||
s := &Server{dispatcher: d, db: db, logger: logger}
|
||||
mux := http.NewServeMux()
|
||||
mux.HandleFunc("/vpcs", s.VpcsHandler)
|
||||
mux.HandleFunc("/vpcs/", s.VpcByNameHandler)
|
||||
|
|
@ -29,12 +29,21 @@ func (s *Server) Start(address string) {
|
|||
mux.HandleFunc("/subnets/", s.SubnetByNameHandler)
|
||||
mux.HandleFunc("/vms", s.VmsHandler)
|
||||
mux.HandleFunc("/vms/", s.VmByNameHandler)
|
||||
s.logger.Info("API server listening", "address", address)
|
||||
if err := http.ListenAndServe(address, s.logMiddleware(mux)); err != nil {
|
||||
s.srv = &http.Server{Addr: address, Handler: s.logMiddleware(mux)}
|
||||
return s
|
||||
}
|
||||
|
||||
func (s *Server) Start() {
|
||||
s.logger.Info("API server listening", "address", s.srv.Addr)
|
||||
if err := s.srv.ListenAndServe(); err != nil && !errors.Is(err, http.ErrServerClosed) {
|
||||
s.logger.Error("API server stopped", "error", err)
|
||||
}
|
||||
}
|
||||
|
||||
func (s *Server) Shutdown(ctx context.Context) error {
|
||||
return s.srv.Shutdown(ctx)
|
||||
}
|
||||
|
||||
type statusWriter struct {
|
||||
http.ResponseWriter
|
||||
status int
|
||||
|
|
|
|||
|
|
@ -81,3 +81,68 @@ database:
|
|||
t.Errorf("attendu %q, obtenu %q", "/opt/two/data", cfg.Database.Path)
|
||||
}
|
||||
}
|
||||
|
||||
func TestLoadConfig_WatchdogDefauts(t *testing.T) {
|
||||
path := writeYAML(t, "")
|
||||
cfg, err := LoadConfig(path)
|
||||
if err != nil {
|
||||
t.Fatalf("LoadConfig a échoué : %v", err)
|
||||
}
|
||||
if cfg.Watchdog.Enabled {
|
||||
t.Error("watchdog.enabled devrait être false par défaut")
|
||||
}
|
||||
if cfg.Watchdog.IntervalSeconds != 60 {
|
||||
t.Errorf("watchdog.interval_seconds attendu 60, obtenu %d", cfg.Watchdog.IntervalSeconds)
|
||||
}
|
||||
}
|
||||
|
||||
func TestLoadConfig_WatchdogValeursExplicites(t *testing.T) {
|
||||
path := writeYAML(t, `
|
||||
watchdog:
|
||||
enabled: true
|
||||
interval_seconds: 30
|
||||
`)
|
||||
cfg, err := LoadConfig(path)
|
||||
if err != nil {
|
||||
t.Fatalf("LoadConfig a échoué : %v", err)
|
||||
}
|
||||
if !cfg.Watchdog.Enabled {
|
||||
t.Error("watchdog.enabled attendu true")
|
||||
}
|
||||
if cfg.Watchdog.IntervalSeconds != 30 {
|
||||
t.Errorf("watchdog.interval_seconds attendu 30, obtenu %d", cfg.Watchdog.IntervalSeconds)
|
||||
}
|
||||
}
|
||||
|
||||
func TestLoadConfig_WatchdogActiveSansIntervalle(t *testing.T) {
|
||||
path := writeYAML(t, `
|
||||
watchdog:
|
||||
enabled: true
|
||||
`)
|
||||
cfg, err := LoadConfig(path)
|
||||
if err != nil {
|
||||
t.Fatalf("LoadConfig a échoué : %v", err)
|
||||
}
|
||||
if !cfg.Watchdog.Enabled {
|
||||
t.Error("watchdog.enabled attendu true")
|
||||
}
|
||||
if cfg.Watchdog.IntervalSeconds != 60 {
|
||||
t.Errorf("watchdog.interval_seconds attendu 60 (défaut viper), obtenu %d", cfg.Watchdog.IntervalSeconds)
|
||||
}
|
||||
}
|
||||
|
||||
func TestLoadConfig_ExempleFourniEstValide(t *testing.T) {
|
||||
cfg, err := LoadConfig("../../../conf/agent/config.exemple.yml")
|
||||
if err != nil {
|
||||
t.Fatalf("config.exemple.yml illisible : %v", err)
|
||||
}
|
||||
if !cfg.Watchdog.Enabled {
|
||||
t.Error("config.exemple.yml devrait activer le watchdog")
|
||||
}
|
||||
if cfg.Watchdog.IntervalSeconds != 60 {
|
||||
t.Errorf("config.exemple.yml : interval_seconds attendu 60, obtenu %d", cfg.Watchdog.IntervalSeconds)
|
||||
}
|
||||
if cfg.QEMU.QMPDir == "" {
|
||||
t.Error("config.exemple.yml devrait définir qemu.qmp_dir")
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -44,6 +44,10 @@ type Config struct {
|
|||
MonitorDir string `mapstructure:"monitor_dir"`
|
||||
QMPDir string `mapstructure:"qmp_dir"`
|
||||
} `mapstructure:"qemu"`
|
||||
Watchdog struct {
|
||||
Enabled bool `mapstructure:"enabled"`
|
||||
IntervalSeconds int `mapstructure:"interval_seconds"`
|
||||
} `mapstructure:"watchdog"`
|
||||
DefaultInterface string `mapstructure:"default_interface"`
|
||||
Interfaces map[string]string `mapstructure:"interfaces"`
|
||||
}
|
||||
|
|
@ -69,6 +73,8 @@ func LoadConfig(path string) (*Config, error) {
|
|||
v.SetDefault("qemu.serial_dir", "/run/two/vms/serial")
|
||||
v.SetDefault("qemu.monitor_dir", "/run/two/vms/monitor")
|
||||
v.SetDefault("qemu.qmp_dir", "/run/two/vms/qmp")
|
||||
v.SetDefault("watchdog.enabled", false)
|
||||
v.SetDefault("watchdog.interval_seconds", 60)
|
||||
v.SetDefault("admin.enabled", false)
|
||||
v.SetDefault("admin.address", "127.0.0.1")
|
||||
v.SetDefault("admin.port", 9091)
|
||||
|
|
|
|||
|
|
@ -24,7 +24,7 @@ const (
|
|||
func subnetIfaceNames(subnetName string) (hostVeth, nsVeth, bridge string, err error) {
|
||||
parts := strings.SplitN(subnetName, "-", 2)
|
||||
if len(parts) < 2 || parts[1] == "" {
|
||||
return "", "", "", fmt.Errorf("nom de subnet %q sans identifiant après le tiret, interfaces indéductibles", subnetName)
|
||||
return "", "", "", fmt.Errorf("subnet name %q has no identifier after the dash, interface names cannot be derived", subnetName)
|
||||
}
|
||||
id := parts[1]
|
||||
return "v-" + id + "-e", "v-" + id + "-i", "br-" + id, nil
|
||||
|
|
@ -37,13 +37,13 @@ func dnsmasqName(vpc, bridge string) string {
|
|||
func CheckSubnets(db *badger.DB, u unitChecker, n notify.Notifier) error {
|
||||
pairs, err := kv.ListByPrefix(db, prefixSubnet)
|
||||
if err != nil {
|
||||
return fmt.Errorf("watchdog: lecture des subnets: %w", err)
|
||||
return fmt.Errorf("watchdog: listing subnets: %w", err)
|
||||
}
|
||||
|
||||
for _, name := range resourceNames(pairs, prefixSubnet) {
|
||||
st, err := state.Get(db, prefixSubnet+name)
|
||||
if err != nil {
|
||||
n.Notify(kindSubnet, name, fmt.Sprintf("état illisible en base: %v", err))
|
||||
n.Notify(kindSubnet, name, fmt.Sprintf("state unreadable in database: %v", err))
|
||||
continue
|
||||
}
|
||||
if st != state.Running {
|
||||
|
|
@ -63,13 +63,13 @@ func checkSubnet(db *badger.DB, name string, u unitChecker, n notify.Notifier) {
|
|||
|
||||
vpc, err := kv.GetFromDB(db, prefixSubnet+name+"/vpc")
|
||||
if err != nil {
|
||||
n.Notify(kindSubnet, name, fmt.Sprintf("vpc illisible en base: %v", err))
|
||||
n.Notify(kindSubnet, name, fmt.Sprintf("vpc unreadable in database: %v", err))
|
||||
return
|
||||
}
|
||||
|
||||
mode, err := kv.GetFromDB(db, prefixSubnet+name+"/mode")
|
||||
if err != nil {
|
||||
n.Notify(kindSubnet, name, fmt.Sprintf("mode illisible en base: %v", err))
|
||||
n.Notify(kindSubnet, name, fmt.Sprintf("mode unreadable in database: %v", err))
|
||||
return
|
||||
}
|
||||
|
||||
|
|
@ -85,7 +85,7 @@ func checkSubnet(db *badger.DB, name string, u unitChecker, n notify.Notifier) {
|
|||
checkVxlanIface(db, name, n)
|
||||
case modeBridge:
|
||||
default:
|
||||
n.Notify(kindSubnet, name, fmt.Sprintf("mode inconnu %q", mode))
|
||||
n.Notify(kindSubnet, name, fmt.Sprintf("unknown mode %q", mode))
|
||||
}
|
||||
|
||||
checkSubnetNetns(name, vpc, nsVeth, bridge, n)
|
||||
|
|
@ -93,7 +93,7 @@ func checkSubnet(db *badger.DB, name string, u unitChecker, n notify.Notifier) {
|
|||
dnsName := dnsmasqName(vpc, bridge)
|
||||
conf := filepath.Join(dhcp.DefaultConfDir, dnsName+".conf")
|
||||
if _, err := os.Stat(conf); err != nil {
|
||||
n.Notify(kindSubnet, name, fmt.Sprintf("config dnsmasq absente (%s): %v", conf, err))
|
||||
n.Notify(kindSubnet, name, fmt.Sprintf("dnsmasq config missing (%s): %v", conf, err))
|
||||
}
|
||||
|
||||
checkUnit(kindSubnet, name, "dnsmasq@"+dnsName+".service", u, n)
|
||||
|
|
@ -102,12 +102,12 @@ func checkSubnet(db *badger.DB, name string, u unitChecker, n notify.Notifier) {
|
|||
func checkVxlanIface(db *badger.DB, name string, n notify.Notifier) {
|
||||
raw, err := kv.GetFromDB(db, prefixSubnet+name+"/vxlan_id")
|
||||
if err != nil {
|
||||
n.Notify(kindSubnet, name, fmt.Sprintf("vxlan_id illisible en base: %v", err))
|
||||
n.Notify(kindSubnet, name, fmt.Sprintf("vxlan_id unreadable in database: %v", err))
|
||||
return
|
||||
}
|
||||
id, err := strconv.Atoi(raw)
|
||||
if err != nil {
|
||||
n.Notify(kindSubnet, name, fmt.Sprintf("vxlan_id invalide %q: %v", raw, err))
|
||||
n.Notify(kindSubnet, name, fmt.Sprintf("invalid vxlan_id %q: %v", raw, err))
|
||||
return
|
||||
}
|
||||
if p := linkProblem(fmt.Sprintf("vxlan-%d", id)); p != "" {
|
||||
|
|
@ -117,18 +117,18 @@ func checkVxlanIface(db *badger.DB, name string, n notify.Notifier) {
|
|||
|
||||
func checkSubnetNetns(name, vpc, nsVeth, bridge string, n notify.Notifier) {
|
||||
if !netns.Exist(vpc) {
|
||||
n.Notify(kindSubnet, name, "netns "+vpc+" absent (/var/run/netns/"+vpc+")")
|
||||
n.Notify(kindSubnet, name, "netns "+vpc+" missing (/var/run/netns/"+vpc+")")
|
||||
return
|
||||
}
|
||||
|
||||
if err := netns.Call(vpc, func() error {
|
||||
for _, iface := range []string{nsVeth, bridge} {
|
||||
if p := linkProblem(iface); p != "" {
|
||||
n.Notify(kindSubnet, name, p+" (dans le netns "+vpc+")")
|
||||
n.Notify(kindSubnet, name, p+" (in netns "+vpc+")")
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}); err != nil {
|
||||
n.Notify(kindSubnet, name, fmt.Sprintf("entrée dans le netns %s impossible: %v", vpc, err))
|
||||
n.Notify(kindSubnet, name, fmt.Sprintf("cannot enter netns %s: %v", vpc, err))
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -85,7 +85,7 @@ func TestCheckSubnets_VPCManquantEnBase(t *testing.T) {
|
|||
if len(got) != 1 {
|
||||
t.Fatalf("attendu 1 notification, obtenu %d : %v", len(got), got)
|
||||
}
|
||||
if !strings.Contains(got[0].problem, "vpc illisible") {
|
||||
if !strings.Contains(got[0].problem, "vpc unreadable") {
|
||||
t.Errorf("problem = %q, devrait porter sur le vpc", got[0].problem)
|
||||
}
|
||||
}
|
||||
|
|
@ -100,7 +100,7 @@ func TestCheckSubnets_ModeManquantEnBase(t *testing.T) {
|
|||
t.Fatalf("erreur inattendue: %v", err)
|
||||
}
|
||||
|
||||
if !r.hasProblemContaining("mode illisible") {
|
||||
if !r.hasProblemContaining("mode unreadable") {
|
||||
t.Errorf("devrait signaler un mode illisible, obtenu %v", r.calls)
|
||||
}
|
||||
}
|
||||
|
|
@ -114,7 +114,7 @@ func TestCheckSubnets_ModeInconnu(t *testing.T) {
|
|||
t.Fatalf("erreur inattendue: %v", err)
|
||||
}
|
||||
|
||||
if !r.hasProblemContaining(`mode inconnu "macvlan"`) {
|
||||
if !r.hasProblemContaining(`unknown mode "macvlan"`) {
|
||||
t.Errorf("devrait signaler un mode inconnu, obtenu %v", r.calls)
|
||||
}
|
||||
}
|
||||
|
|
@ -147,7 +147,7 @@ func TestCheckSubnets_ModeVxlanSansVxlanID(t *testing.T) {
|
|||
t.Fatalf("erreur inattendue: %v", err)
|
||||
}
|
||||
|
||||
if !r.hasProblemContaining("vxlan_id illisible") {
|
||||
if !r.hasProblemContaining("vxlan_id unreadable") {
|
||||
t.Errorf("devrait signaler un vxlan_id illisible, obtenu %v", r.calls)
|
||||
}
|
||||
}
|
||||
|
|
@ -162,7 +162,7 @@ func TestCheckSubnets_ModeVxlanVxlanIDInvalide(t *testing.T) {
|
|||
t.Fatalf("erreur inattendue: %v", err)
|
||||
}
|
||||
|
||||
if !r.hasProblemContaining("vxlan_id invalide") {
|
||||
if !r.hasProblemContaining("invalid vxlan_id") {
|
||||
t.Errorf("devrait signaler un vxlan_id invalide, obtenu %v", r.calls)
|
||||
}
|
||||
}
|
||||
|
|
@ -191,7 +191,7 @@ func TestCheckSubnets_ConfigDnsmasqAbsente(t *testing.T) {
|
|||
t.Fatalf("erreur inattendue: %v", err)
|
||||
}
|
||||
|
||||
if !r.hasProblemContaining("config dnsmasq absente") {
|
||||
if !r.hasProblemContaining("dnsmasq config missing") {
|
||||
t.Errorf("devrait signaler la config dnsmasq absente, obtenu %v", r.calls)
|
||||
}
|
||||
if !r.hasProblemContaining("vp-admin_br-000042.conf") {
|
||||
|
|
@ -242,7 +242,7 @@ func TestCheckSubnets_UnitIllisible(t *testing.T) {
|
|||
t.Fatalf("erreur inattendue: %v", err)
|
||||
}
|
||||
|
||||
if !r.hasProblemContaining("unit dnsmasq@vp-admin_br-000042.service illisible") {
|
||||
if !r.hasProblemContaining("unit dnsmasq@vp-admin_br-000042.service unreadable") {
|
||||
t.Errorf("devrait signaler l'unit illisible, obtenu %v", r.calls)
|
||||
}
|
||||
}
|
||||
|
|
@ -271,7 +271,7 @@ func TestCheckSubnets_EtatCorrompuNInterrompPasLaBoucle(t *testing.T) {
|
|||
t.Fatalf("un état corrompu ne doit pas faire échouer CheckSubnets: %v", err)
|
||||
}
|
||||
|
||||
if !strings.Contains(strings.Join(problems(r.forName("br-corrompu")), " "), "état illisible") {
|
||||
if !strings.Contains(strings.Join(problems(r.forName("br-corrompu")), " "), "state unreadable") {
|
||||
t.Errorf("devrait signaler l'état corrompu, obtenu %v", r.calls)
|
||||
}
|
||||
if len(r.forName("br-000042")) == 0 {
|
||||
|
|
|
|||
|
|
@ -23,18 +23,18 @@ func tapName(tapID int) string {
|
|||
|
||||
func CheckVMs(db *badger.DB, cfg *configuration.Config, u unitChecker, n notify.Notifier) error {
|
||||
if cfg == nil {
|
||||
return errors.New("watchdog: configuration requise pour vérifier les VMs")
|
||||
return errors.New("watchdog: configuration required to check vms")
|
||||
}
|
||||
|
||||
pairs, err := kv.ListByPrefix(db, prefixVM)
|
||||
if err != nil {
|
||||
return fmt.Errorf("watchdog: lecture des vm: %w", err)
|
||||
return fmt.Errorf("watchdog: listing vms: %w", err)
|
||||
}
|
||||
|
||||
for _, name := range resourceNames(pairs, prefixVM) {
|
||||
st, err := state.Get(db, prefixVM+name)
|
||||
if err != nil {
|
||||
n.Notify(kindVM, name, fmt.Sprintf("état illisible en base: %v", err))
|
||||
n.Notify(kindVM, name, fmt.Sprintf("state unreadable in database: %v", err))
|
||||
continue
|
||||
}
|
||||
if st != state.Running {
|
||||
|
|
@ -48,13 +48,13 @@ func CheckVMs(db *badger.DB, cfg *configuration.Config, u unitChecker, n notify.
|
|||
func checkVM(db *badger.DB, cfg *configuration.Config, name string, u unitChecker, n notify.Notifier) {
|
||||
subnetName, err := kv.GetFromDB(db, prefixVM+name+"/subnet")
|
||||
if err != nil {
|
||||
n.Notify(kindVM, name, fmt.Sprintf("subnet illisible en base: %v", err))
|
||||
n.Notify(kindVM, name, fmt.Sprintf("subnet unreadable in database: %v", err))
|
||||
return
|
||||
}
|
||||
|
||||
vpc, err := kv.GetFromDB(db, prefixSubnet+subnetName+"/vpc")
|
||||
if err != nil {
|
||||
n.Notify(kindVM, name, fmt.Sprintf("vpc du subnet %s illisible en base: %v", subnetName, err))
|
||||
n.Notify(kindVM, name, fmt.Sprintf("vpc of subnet %s unreadable in database: %v", subnetName, err))
|
||||
return
|
||||
}
|
||||
|
||||
|
|
@ -67,33 +67,33 @@ func checkVM(db *badger.DB, cfg *configuration.Config, name string, u unitChecke
|
|||
func checkVMTap(db *badger.DB, name, vpc string, n notify.Notifier) {
|
||||
raw, err := kv.GetFromDB(db, prefixVM+name+"/tap_id")
|
||||
if err != nil {
|
||||
n.Notify(kindVM, name, fmt.Sprintf("tap_id illisible en base: %v", err))
|
||||
n.Notify(kindVM, name, fmt.Sprintf("tap_id unreadable in database: %v", err))
|
||||
return
|
||||
}
|
||||
tapID, err := strconv.Atoi(raw)
|
||||
if err != nil {
|
||||
n.Notify(kindVM, name, fmt.Sprintf("tap_id invalide %q: %v", raw, err))
|
||||
n.Notify(kindVM, name, fmt.Sprintf("invalid tap_id %q: %v", raw, err))
|
||||
return
|
||||
}
|
||||
|
||||
if !netns.Exist(vpc) {
|
||||
n.Notify(kindVM, name, "netns "+vpc+" absent (/var/run/netns/"+vpc+")")
|
||||
n.Notify(kindVM, name, "netns "+vpc+" missing (/var/run/netns/"+vpc+")")
|
||||
return
|
||||
}
|
||||
|
||||
if err := netns.Call(vpc, func() error {
|
||||
if p := linkProblem(tapName(tapID)); p != "" {
|
||||
n.Notify(kindVM, name, p+" (dans le netns "+vpc+")")
|
||||
n.Notify(kindVM, name, p+" (in netns "+vpc+")")
|
||||
}
|
||||
return nil
|
||||
}); err != nil {
|
||||
n.Notify(kindVM, name, fmt.Sprintf("entrée dans le netns %s impossible: %v", vpc, err))
|
||||
n.Notify(kindVM, name, fmt.Sprintf("cannot enter netns %s: %v", vpc, err))
|
||||
}
|
||||
}
|
||||
|
||||
func checkVMQemu(cfg *configuration.Config, name string, n notify.Notifier) {
|
||||
sock := filepath.Join(cfg.QEMU.QMPDir, name+".sock")
|
||||
if _, err := qmp.Send(sock, nil); err != nil {
|
||||
n.Notify(kindVM, name, fmt.Sprintf("qemu ne répond pas sur %s: %v", sock, err))
|
||||
n.Notify(kindVM, name, fmt.Sprintf("qemu not responding on %s: %v", sock, err))
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -86,7 +86,7 @@ func TestCheckVMs_SubnetManquantEnBase(t *testing.T) {
|
|||
if len(got) != 1 {
|
||||
t.Fatalf("attendu 1 notification, obtenu %d : %v", len(got), got)
|
||||
}
|
||||
if !strings.Contains(got[0].problem, "subnet illisible") {
|
||||
if !strings.Contains(got[0].problem, "subnet unreadable") {
|
||||
t.Errorf("problem = %q, devrait porter sur le subnet", got[0].problem)
|
||||
}
|
||||
}
|
||||
|
|
@ -101,7 +101,7 @@ func TestCheckVMs_VPCDuSubnetManquant(t *testing.T) {
|
|||
t.Fatalf("erreur inattendue: %v", err)
|
||||
}
|
||||
|
||||
if !r.hasProblemContaining("vpc du subnet br-000042 illisible") {
|
||||
if !r.hasProblemContaining("vpc of subnet br-000042 unreadable") {
|
||||
t.Errorf("devrait signaler le vpc introuvable, obtenu %v", r.calls)
|
||||
}
|
||||
}
|
||||
|
|
@ -115,7 +115,7 @@ func TestCheckVMs_TapIDManquant(t *testing.T) {
|
|||
t.Fatalf("erreur inattendue: %v", err)
|
||||
}
|
||||
|
||||
if !r.hasProblemContaining("tap_id illisible") {
|
||||
if !r.hasProblemContaining("tap_id unreadable") {
|
||||
t.Errorf("devrait signaler un tap_id illisible, obtenu %v", r.calls)
|
||||
}
|
||||
}
|
||||
|
|
@ -129,7 +129,7 @@ func TestCheckVMs_TapIDInvalide(t *testing.T) {
|
|||
t.Fatalf("erreur inattendue: %v", err)
|
||||
}
|
||||
|
||||
if !r.hasProblemContaining("tap_id invalide") {
|
||||
if !r.hasProblemContaining("invalid tap_id") {
|
||||
t.Errorf("devrait signaler un tap_id invalide, obtenu %v", r.calls)
|
||||
}
|
||||
}
|
||||
|
|
@ -144,7 +144,7 @@ func TestCheckVMs_QemuNeRepondPas(t *testing.T) {
|
|||
t.Fatalf("erreur inattendue: %v", err)
|
||||
}
|
||||
|
||||
if !r.hasProblemContaining("qemu ne répond pas") {
|
||||
if !r.hasProblemContaining("qemu not responding") {
|
||||
t.Errorf("devrait signaler que qemu ne répond pas, obtenu %v", r.calls)
|
||||
}
|
||||
if !r.hasProblemContaining("i-test1.sock") {
|
||||
|
|
@ -220,7 +220,7 @@ func TestCheckVMs_EtatCorrompuNInterrompPasLaBoucle(t *testing.T) {
|
|||
t.Fatalf("un état corrompu ne doit pas faire échouer CheckVMs: %v", err)
|
||||
}
|
||||
|
||||
if !strings.Contains(strings.Join(problems(r.forName("i-corrompu")), " "), "état illisible") {
|
||||
if !strings.Contains(strings.Join(problems(r.forName("i-corrompu")), " "), "state unreadable") {
|
||||
t.Errorf("devrait signaler l'état corrompu, obtenu %v", r.calls)
|
||||
}
|
||||
if len(r.forName("i-test1")) == 0 {
|
||||
|
|
|
|||
|
|
@ -17,7 +17,7 @@ const vpcBridge = "br-public"
|
|||
func vpcIfaceNames(vpcName string) (host, ns string, err error) {
|
||||
parts := strings.SplitN(vpcName, "-", 2)
|
||||
if len(parts) < 2 || parts[1] == "" {
|
||||
return "", "", fmt.Errorf("nom de VPC %q sans identifiant après le tiret, interfaces indéductibles", vpcName)
|
||||
return "", "", fmt.Errorf("vpc name %q has no identifier after the dash, interface names cannot be derived", vpcName)
|
||||
}
|
||||
return "vp-" + parts[1] + "-e", "vp-" + parts[1] + "-i", nil
|
||||
}
|
||||
|
|
@ -25,13 +25,13 @@ func vpcIfaceNames(vpcName string) (host, ns string, err error) {
|
|||
func CheckVPCs(db *badger.DB, n notify.Notifier) error {
|
||||
pairs, err := kv.ListByPrefix(db, prefixVPC)
|
||||
if err != nil {
|
||||
return fmt.Errorf("watchdog: lecture des vpc: %w", err)
|
||||
return fmt.Errorf("watchdog: listing vpcs: %w", err)
|
||||
}
|
||||
|
||||
for _, name := range resourceNames(pairs, prefixVPC) {
|
||||
st, err := state.Get(db, prefixVPC+name)
|
||||
if err != nil {
|
||||
n.Notify(kindVPC, name, fmt.Sprintf("état illisible en base: %v", err))
|
||||
n.Notify(kindVPC, name, fmt.Sprintf("state unreadable in database: %v", err))
|
||||
continue
|
||||
}
|
||||
if st != state.Running {
|
||||
|
|
@ -50,7 +50,7 @@ func checkVPC(name string, n notify.Notifier) {
|
|||
}
|
||||
|
||||
if !netns.Exist(name) {
|
||||
n.Notify(kindVPC, name, "netns absent (/var/run/netns/"+name+")")
|
||||
n.Notify(kindVPC, name, "netns missing (/var/run/netns/"+name+")")
|
||||
return
|
||||
}
|
||||
|
||||
|
|
@ -61,11 +61,11 @@ func checkVPC(name string, n notify.Notifier) {
|
|||
if err := netns.Call(name, func() error {
|
||||
for _, iface := range []string{nsVeth, vpcBridge} {
|
||||
if p := linkProblem(iface); p != "" {
|
||||
n.Notify(kindVPC, name, p+" (dans le netns)")
|
||||
n.Notify(kindVPC, name, p+" (in netns)")
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}); err != nil {
|
||||
n.Notify(kindVPC, name, fmt.Sprintf("entrée dans le netns impossible: %v", err))
|
||||
n.Notify(kindVPC, name, fmt.Sprintf("cannot enter netns: %v", err))
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -122,7 +122,7 @@ func TestCheckVPCs_EtatCorrompu(t *testing.T) {
|
|||
if len(got) != 1 {
|
||||
t.Fatalf("attendu 1 notification pour l'état corrompu, obtenu %d", len(got))
|
||||
}
|
||||
if !strings.Contains(got[0].problem, "état illisible") {
|
||||
if !strings.Contains(got[0].problem, "state unreadable") {
|
||||
t.Errorf("problem = %q, devrait mentionner un état illisible", got[0].problem)
|
||||
}
|
||||
if len(r.forName("vp-suivant")) == 0 {
|
||||
|
|
@ -143,7 +143,7 @@ func TestCheckVPCs_NomIndeductibleNePaniquePas(t *testing.T) {
|
|||
if len(got) != 1 {
|
||||
t.Fatalf("attendu 1 notification, obtenu %d : %v", len(got), got)
|
||||
}
|
||||
if !strings.Contains(got[0].problem, "indéductibles") {
|
||||
if !strings.Contains(got[0].problem, "cannot be derived") {
|
||||
t.Errorf("problem = %q, devrait porter sur les interfaces indéductibles", got[0].problem)
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -16,7 +16,7 @@ func NewStderr(l *slog.Logger) *StderrNotifier {
|
|||
}
|
||||
|
||||
func (n *StderrNotifier) Notify(kind, name, problem string) {
|
||||
n.logger.Error("watchdog: incohérence détectée",
|
||||
n.logger.Error("watchdog: inconsistency detected",
|
||||
"kind", kind,
|
||||
"name", name,
|
||||
"problem", problem,
|
||||
|
|
|
|||
|
|
@ -52,7 +52,7 @@ func checkUnit(kind, name, unit string, u unitChecker, n notify.Notifier) {
|
|||
}
|
||||
st, err := u.Status(unit)
|
||||
if err != nil {
|
||||
n.Notify(kind, name, fmt.Sprintf("unit %s illisible: %v", unit, err))
|
||||
n.Notify(kind, name, fmt.Sprintf("unit %s unreadable: %v", unit, err))
|
||||
return
|
||||
}
|
||||
if st.ActiveState != "active" {
|
||||
|
|
@ -64,7 +64,7 @@ func linkProblem(iface string) string {
|
|||
up, err := netif.LinkIsUp(iface)
|
||||
switch {
|
||||
case err != nil:
|
||||
return fmt.Sprintf("interface %s introuvable: %v", iface, err)
|
||||
return fmt.Sprintf("interface %s not found: %v", iface, err)
|
||||
case !up:
|
||||
return fmt.Sprintf("interface %s down", iface)
|
||||
}
|
||||
|
|
|
|||
|
|
@ -28,7 +28,7 @@ func New(db *badger.DB, cfg *configuration.Config, n notify.Notifier, logger *sl
|
|||
logger = slog.Default()
|
||||
}
|
||||
if interval <= 0 {
|
||||
logger.Warn("watchdog: intervalle invalide, valeur par défaut appliquée",
|
||||
logger.Warn("watchdog: invalid interval, default applied",
|
||||
"interval", interval, "default", defaultInterval)
|
||||
interval = defaultInterval
|
||||
}
|
||||
|
|
@ -45,12 +45,12 @@ func (w *Watchdog) Run(ctx context.Context) {
|
|||
ticker := time.NewTicker(w.interval)
|
||||
defer ticker.Stop()
|
||||
|
||||
w.logger.Info("watchdog: démarrage", "interval", w.interval)
|
||||
w.logger.Info("watchdog: starting", "interval", w.interval)
|
||||
|
||||
for {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
w.logger.Info("watchdog: arrêt")
|
||||
w.logger.Info("watchdog: stopping")
|
||||
return
|
||||
case <-ticker.C:
|
||||
w.tick()
|
||||
|
|
@ -63,13 +63,13 @@ func (w *Watchdog) tick() {
|
|||
defer closeUnits()
|
||||
|
||||
if err := CheckVPCs(w.db, w.notifier); err != nil {
|
||||
w.logger.Error("watchdog: vérification des vpc", "err", err)
|
||||
w.logger.Error("watchdog: vpc check failed", "err", err)
|
||||
}
|
||||
if err := CheckSubnets(w.db, u, w.notifier); err != nil {
|
||||
w.logger.Error("watchdog: vérification des subnets", "err", err)
|
||||
w.logger.Error("watchdog: subnet check failed", "err", err)
|
||||
}
|
||||
if err := CheckVMs(w.db, w.cfg, u, w.notifier); err != nil {
|
||||
w.logger.Error("watchdog: vérification des vm", "err", err)
|
||||
w.logger.Error("watchdog: vm check failed", "err", err)
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -77,13 +77,13 @@ func (w *Watchdog) units() (unitChecker, func()) {
|
|||
m, err := systemd.New()
|
||||
if err != nil {
|
||||
if !w.dbusDown {
|
||||
w.logger.Warn("watchdog: connexion systemd impossible, vérification des units désactivée", "err", err)
|
||||
w.logger.Warn("watchdog: systemd unreachable, unit checks disabled", "err", err)
|
||||
w.dbusDown = true
|
||||
}
|
||||
return nil, func() {}
|
||||
}
|
||||
if w.dbusDown {
|
||||
w.logger.Info("watchdog: connexion systemd rétablie")
|
||||
w.logger.Info("watchdog: systemd connection restored")
|
||||
w.dbusDown = false
|
||||
}
|
||||
return m, m.Close
|
||||
|
|
|
|||
|
|
@ -61,7 +61,7 @@ func TestNew_IntervalleInvalideUtiliseLeDefaut(t *testing.T) {
|
|||
}
|
||||
}
|
||||
|
||||
if !strings.Contains(buf.String(), "intervalle invalide") {
|
||||
if !strings.Contains(buf.String(), "invalid interval") {
|
||||
t.Error("un intervalle invalide devrait être signalé dans les logs")
|
||||
}
|
||||
}
|
||||
|
|
@ -143,11 +143,11 @@ func TestRun_ErreurDeBaseLogueeEtBoucleContinue(t *testing.T) {
|
|||
for _, line := range logLines(t, buf) {
|
||||
msg, _ := line["msg"].(string)
|
||||
switch {
|
||||
case strings.Contains(msg, "vérification des vpc"):
|
||||
case strings.Contains(msg, "vpc check failed"):
|
||||
vpc = true
|
||||
case strings.Contains(msg, "vérification des subnets"):
|
||||
case strings.Contains(msg, "subnet check failed"):
|
||||
subnet = true
|
||||
case strings.Contains(msg, "vérification des vm"):
|
||||
case strings.Contains(msg, "vm check failed"):
|
||||
vm = true
|
||||
}
|
||||
}
|
||||
|
|
@ -168,10 +168,10 @@ func TestRun_LogueDemarrageEtArret(t *testing.T) {
|
|||
w.Run(ctx)
|
||||
|
||||
out := buf.String()
|
||||
if !strings.Contains(out, "watchdog: démarrage") {
|
||||
if !strings.Contains(out, "watchdog: starting") {
|
||||
t.Error("le démarrage devrait être logué")
|
||||
}
|
||||
if !strings.Contains(out, "watchdog: arrêt") {
|
||||
if !strings.Contains(out, "watchdog: stopping") {
|
||||
t.Error("l'arrêt devrait être logué")
|
||||
}
|
||||
}
|
||||
|
|
@ -190,7 +190,7 @@ func TestUnits_ConnexionSystemdIndisponibleSignaleeUneSeuleFois(t *testing.T) {
|
|||
|
||||
var warnings int
|
||||
for _, line := range logLines(t, buf) {
|
||||
if msg, _ := line["msg"].(string); strings.Contains(msg, "connexion systemd impossible") {
|
||||
if msg, _ := line["msg"].(string); strings.Contains(msg, "systemd unreachable") {
|
||||
warnings++
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,6 +1,8 @@
|
|||
package kv
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"log/slog"
|
||||
"net/http"
|
||||
|
|
@ -12,21 +14,28 @@ import (
|
|||
type AdminServer struct {
|
||||
db *badger.DB
|
||||
logger *slog.Logger
|
||||
srv *http.Server
|
||||
}
|
||||
|
||||
func NewAdminServer(db *badger.DB, logger *slog.Logger) *AdminServer {
|
||||
return &AdminServer{db: db, logger: logger}
|
||||
}
|
||||
|
||||
func (s *AdminServer) Start(address string) {
|
||||
func NewAdminServer(db *badger.DB, logger *slog.Logger, address string) *AdminServer {
|
||||
s := &AdminServer{db: db, logger: logger}
|
||||
mux := http.NewServeMux()
|
||||
mux.HandleFunc("/db", s.dbHandler)
|
||||
s.logger.Info("admin server listening", "address", address)
|
||||
if err := http.ListenAndServe(address, mux); err != nil {
|
||||
s.srv = &http.Server{Addr: address, Handler: mux}
|
||||
return s
|
||||
}
|
||||
|
||||
func (s *AdminServer) Start() {
|
||||
s.logger.Info("admin server listening", "address", s.srv.Addr)
|
||||
if err := s.srv.ListenAndServe(); err != nil && !errors.Is(err, http.ErrServerClosed) {
|
||||
s.logger.Error("admin server stopped", "error", err)
|
||||
}
|
||||
}
|
||||
|
||||
func (s *AdminServer) Shutdown(ctx context.Context) error {
|
||||
return s.srv.Shutdown(ctx)
|
||||
}
|
||||
|
||||
func (s *AdminServer) dbHandler(w http.ResponseWriter, r *http.Request) {
|
||||
if r.Method != http.MethodGet {
|
||||
http.Error(w, "method not allowed", http.StatusMethodNotAllowed)
|
||||
|
|
|
|||
|
|
@ -1,6 +1,8 @@
|
|||
package promserver
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"log"
|
||||
"net/http"
|
||||
|
||||
|
|
@ -8,13 +10,30 @@ import (
|
|||
"github.com/prometheus/client_golang/prometheus/promhttp"
|
||||
)
|
||||
|
||||
// Start launches the Prometheus metrics HTTP server on the given address.
|
||||
// The provided registry is used to expose metrics at /metrics.
|
||||
func Start(address string, registry *prometheus.Registry) {
|
||||
// Server exposes the Prometheus metrics endpoint.
|
||||
type Server struct {
|
||||
srv *http.Server
|
||||
}
|
||||
|
||||
// New builds the metrics server for the given address and registry.
|
||||
func New(address string, registry *prometheus.Registry) *Server {
|
||||
mux := http.NewServeMux()
|
||||
mux.Handle("/metrics", promhttp.HandlerFor(registry, promhttp.HandlerOpts{
|
||||
EnableOpenMetrics: true,
|
||||
}))
|
||||
log.Printf("Prometheus server listening on %s", address)
|
||||
log.Fatal(http.ListenAndServe(address, mux))
|
||||
return &Server{srv: &http.Server{Addr: address, Handler: mux}}
|
||||
}
|
||||
|
||||
// Start blocks until the server stops. A failure is logged, never fatal: an
|
||||
// unavailable metrics endpoint must not bring the whole agent down.
|
||||
func (s *Server) Start() {
|
||||
log.Printf("Prometheus server listening on %s", s.srv.Addr)
|
||||
if err := s.srv.ListenAndServe(); err != nil && !errors.Is(err, http.ErrServerClosed) {
|
||||
log.Printf("Prometheus server stopped: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
// Shutdown stops the server, waiting for in-flight requests until ctx expires.
|
||||
func (s *Server) Shutdown(ctx context.Context) error {
|
||||
return s.srv.Shutdown(ctx)
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,6 +1,9 @@
|
|||
package worker
|
||||
|
||||
import "log"
|
||||
import (
|
||||
"log"
|
||||
"sync"
|
||||
)
|
||||
|
||||
// Task is a function to be executed asynchronously by a worker.
|
||||
type Task func()
|
||||
|
|
@ -8,6 +11,9 @@ type Task func()
|
|||
// Queue is a FIFO channel-backed task queue consumed by worker goroutines.
|
||||
type Queue struct {
|
||||
tasks chan Task
|
||||
wg sync.WaitGroup
|
||||
mu sync.RWMutex
|
||||
stopped bool
|
||||
}
|
||||
|
||||
// New creates a Queue with the given channel buffer size.
|
||||
|
|
@ -15,8 +21,16 @@ func New(bufferSize int) *Queue {
|
|||
return &Queue{tasks: make(chan Task, bufferSize)}
|
||||
}
|
||||
|
||||
// Submit enqueues a task. Blocks if the queue is full.
|
||||
// Submit enqueues a task. Blocks if the queue is full. Tasks submitted after
|
||||
// Stop are rejected and logged rather than enqueued.
|
||||
func (q *Queue) Submit(t Task) {
|
||||
q.mu.RLock()
|
||||
defer q.mu.RUnlock()
|
||||
|
||||
if q.stopped {
|
||||
log.Print("worker: queue stopped, task rejected")
|
||||
return
|
||||
}
|
||||
q.tasks <- t
|
||||
}
|
||||
|
||||
|
|
@ -24,10 +38,28 @@ func (q *Queue) Submit(t Task) {
|
|||
func (q *Queue) Start(n int) {
|
||||
log.Printf("worker: starting %d workers", n)
|
||||
for i := range n {
|
||||
q.wg.Add(1)
|
||||
go func(id int) {
|
||||
defer q.wg.Done()
|
||||
for task := range q.tasks {
|
||||
task()
|
||||
}
|
||||
}(i)
|
||||
}
|
||||
}
|
||||
|
||||
// Stop rejects new tasks, then waits for the queued and in-flight ones to
|
||||
// finish. The write lock is what makes closing the channel safe: it is only
|
||||
// taken once every in-flight Submit has released its read lock.
|
||||
func (q *Queue) Stop() {
|
||||
q.mu.Lock()
|
||||
if q.stopped {
|
||||
q.mu.Unlock()
|
||||
return
|
||||
}
|
||||
q.stopped = true
|
||||
close(q.tasks)
|
||||
q.mu.Unlock()
|
||||
|
||||
q.wg.Wait()
|
||||
}
|
||||
|
|
|
|||
|
|
@ -101,3 +101,79 @@ func TestQueue_SubmitBlocksWhenFull(t *testing.T) {
|
|||
t.Error("Submit aurait dû se débloquer après démarrage d'un worker")
|
||||
}
|
||||
}
|
||||
|
||||
func TestStop_AttendLesTachesEnCours(t *testing.T) {
|
||||
q := New(10)
|
||||
q.Start(2)
|
||||
|
||||
var mu sync.Mutex
|
||||
done := 0
|
||||
for range 5 {
|
||||
q.Submit(func() {
|
||||
time.Sleep(20 * time.Millisecond)
|
||||
mu.Lock()
|
||||
done++
|
||||
mu.Unlock()
|
||||
})
|
||||
}
|
||||
|
||||
q.Stop()
|
||||
|
||||
mu.Lock()
|
||||
defer mu.Unlock()
|
||||
if done != 5 {
|
||||
t.Errorf("Stop devrait attendre les 5 tâches, %d terminées", done)
|
||||
}
|
||||
}
|
||||
|
||||
func TestStop_RejetteLesTachesSuivantes(t *testing.T) {
|
||||
q := New(10)
|
||||
q.Start(1)
|
||||
q.Stop()
|
||||
|
||||
executed := make(chan struct{}, 1)
|
||||
q.Submit(func() { executed <- struct{}{} })
|
||||
|
||||
select {
|
||||
case <-executed:
|
||||
t.Error("une tâche soumise après Stop ne doit pas être exécutée")
|
||||
case <-time.After(100 * time.Millisecond):
|
||||
}
|
||||
}
|
||||
|
||||
func TestStop_Idempotent(t *testing.T) {
|
||||
q := New(10)
|
||||
q.Start(1)
|
||||
q.Stop()
|
||||
q.Stop()
|
||||
}
|
||||
|
||||
func TestStop_SansTacheEnCours(t *testing.T) {
|
||||
q := New(10)
|
||||
q.Start(3)
|
||||
|
||||
done := make(chan struct{})
|
||||
go func() { q.Stop(); close(done) }()
|
||||
|
||||
select {
|
||||
case <-done:
|
||||
case <-time.After(2 * time.Second):
|
||||
t.Fatal("Stop n'a pas rendu la main sur une file vide")
|
||||
}
|
||||
}
|
||||
|
||||
func TestSubmit_ConcurrentAvecStop(t *testing.T) {
|
||||
q := New(100)
|
||||
q.Start(4)
|
||||
|
||||
var wg sync.WaitGroup
|
||||
for range 20 {
|
||||
wg.Add(1)
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
q.Submit(func() {})
|
||||
}()
|
||||
}
|
||||
q.Stop()
|
||||
wg.Wait()
|
||||
}
|
||||
|
|
|
|||
|
|
@ -5,6 +5,7 @@ After=network.target
|
|||
[Service]
|
||||
Type=simple
|
||||
ExecStart=/opt/two/bin/agent
|
||||
TimeoutStopSec=30
|
||||
|
||||
[Install]
|
||||
WantedBy=multi-user.target
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue