diff --git a/internal/netif/link.go b/internal/netif/link.go new file mode 100644 index 0000000..4e1a2b4 --- /dev/null +++ b/internal/netif/link.go @@ -0,0 +1,15 @@ +package netif + +import ( + "net" + + "github.com/vishvananda/netlink" +) + +func LinkIsUp(name string) (bool, error) { + link, err := netlink.LinkByName(name) + if err != nil { + return false, err + } + return link.Attrs().Flags&net.FlagUp != 0, nil +} diff --git a/internal/watchdog/check_vpc.go b/internal/watchdog/check_vpc.go new file mode 100644 index 0000000..5cf252d --- /dev/null +++ b/internal/watchdog/check_vpc.go @@ -0,0 +1,71 @@ +package watchdog + +import ( + "fmt" + "strings" + + "git.g3e.fr/syonad/two/internal/netns" + "git.g3e.fr/syonad/two/internal/state" + "git.g3e.fr/syonad/two/internal/watchdog/notify" + "git.g3e.fr/syonad/two/pkg/db/kv" + + "github.com/dgraph-io/badger/v4" +) + +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 "vp-" + parts[1] + "-e", "vp-" + parts[1] + "-i", nil +} + +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) + } + + 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)) + continue + } + if st != state.Running { + continue + } + checkVPC(name, n) + } + return nil +} + +func checkVPC(name string, n notify.Notifier) { + hostVeth, nsVeth, err := vpcIfaceNames(name) + if err != nil { + n.Notify(kindVPC, name, err.Error()) + return + } + + if !netns.Exist(name) { + n.Notify(kindVPC, name, "netns absent (/var/run/netns/"+name+")") + return + } + + if p := linkProblem(hostVeth); p != "" { + n.Notify(kindVPC, name, p) + } + + 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)") + } + } + return nil + }); err != nil { + n.Notify(kindVPC, name, fmt.Sprintf("entrée dans le netns impossible: %v", err)) + } +} diff --git a/internal/watchdog/check_vpc_test.go b/internal/watchdog/check_vpc_test.go new file mode 100644 index 0000000..1b37248 --- /dev/null +++ b/internal/watchdog/check_vpc_test.go @@ -0,0 +1,177 @@ +package watchdog + +import ( + "strings" + "testing" + + "git.g3e.fr/syonad/two/internal/state" + "git.g3e.fr/syonad/two/pkg/db/kv" +) + +func TestVPCIfaceNames_NomStandard(t *testing.T) { + host, ns, err := vpcIfaceNames("vp-admin") + if err != nil { + t.Fatalf("erreur inattendue: %v", err) + } + if host != "vp-admin-e" || ns != "vp-admin-i" { + t.Errorf("obtenu (%q, %q), attendu (vp-admin-e, vp-admin-i)", host, ns) + } +} + +func TestVPCIfaceNames_IdentifiantNumerique(t *testing.T) { + host, ns, err := vpcIfaceNames("vpc-000003") + if err != nil { + t.Fatalf("erreur inattendue: %v", err) + } + if host != "vp-000003-e" || ns != "vp-000003-i" { + t.Errorf("obtenu (%q, %q), attendu (vp-000003-e, vp-000003-i)", host, ns) + } +} + +func TestVPCIfaceNames_PlusieursTirets(t *testing.T) { + host, _, err := vpcIfaceNames("vp-admin-prod") + if err != nil { + t.Fatalf("erreur inattendue: %v", err) + } + if host != "vp-admin-prod-e" { + t.Errorf("obtenu %q, attendu vp-admin-prod-e", host) + } +} + +func TestVPCIfaceNames_NomSansTiret(t *testing.T) { + if _, _, err := vpcIfaceNames("admin"); err == nil { + t.Fatal("un nom sans tiret devrait produire une erreur, pas une panique") + } +} + +func TestVPCIfaceNames_TiretFinal(t *testing.T) { + if _, _, err := vpcIfaceNames("vp-"); err == nil { + t.Fatal("un identifiant vide devrait produire une erreur") + } +} + +func TestVPCIfaceNames_NomVide(t *testing.T) { + if _, _, err := vpcIfaceNames(""); err == nil { + t.Fatal("un nom vide devrait produire une erreur") + } +} + +func TestCheckVPCs_BaseVide(t *testing.T) { + db := newTestDB(t) + r := &recorder{} + + if err := CheckVPCs(db, r); err != nil { + t.Fatalf("erreur inattendue: %v", err) + } + if len(r.calls) != 0 { + t.Errorf("aucune notification attendue, obtenu %v", r.calls) + } +} + +func TestCheckVPCs_IgnoreLesEtatsNonRunning(t *testing.T) { + db := newTestDB(t) + for _, s := range []state.State{state.Creating, state.Deleting, state.Error, state.Deleted} { + seedResource(t, db, prefixVPC, "vp-"+string(s), s) + } + r := &recorder{} + + if err := CheckVPCs(db, r); err != nil { + t.Fatalf("erreur inattendue: %v", err) + } + if len(r.calls) != 0 { + t.Errorf("aucune notification attendue pour des états non-running, obtenu %v", r.calls) + } +} + +func TestCheckVPCs_SignaleUnVPCRunningAbsent(t *testing.T) { + db := newTestDB(t) + seedResource(t, db, prefixVPC, "vp-fantome", state.Running) + r := &recorder{} + + if err := CheckVPCs(db, r); err != nil { + t.Fatalf("erreur inattendue: %v", err) + } + + got := r.forName("vp-fantome") + if len(got) == 0 { + t.Fatal("un VPC running inexistant sur le système doit être signalé") + } + for _, c := range got { + if c.kind != kindVPC { + t.Errorf("kind = %q, attendu %q", c.kind, kindVPC) + } + if c.problem == "" { + t.Error("problem ne doit pas être vide") + } + } +} + +func TestCheckVPCs_EtatCorrompu(t *testing.T) { + db := newTestDB(t) + if err := kv.AddInDB(db, prefixVPC+"vp-corrompu"+stateSuffix, "n_importe_quoi"); err != nil { + t.Fatalf("préparation: %v", err) + } + seedResource(t, db, prefixVPC, "vp-suivant", state.Running) + r := &recorder{} + + if err := CheckVPCs(db, r); err != nil { + t.Fatalf("un état corrompu ne doit pas faire échouer CheckVPCs: %v", err) + } + + got := r.forName("vp-corrompu") + 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") { + t.Errorf("problem = %q, devrait mentionner un état illisible", got[0].problem) + } + if len(r.forName("vp-suivant")) == 0 { + t.Error("une clé corrompue ne doit pas empêcher l'examen des VPC suivants") + } +} + +func TestCheckVPCs_NomIndeductibleNePaniquePas(t *testing.T) { + db := newTestDB(t) + seedResource(t, db, prefixVPC, "admin", state.Running) + r := &recorder{} + + if err := CheckVPCs(db, r); err != nil { + t.Fatalf("erreur inattendue: %v", err) + } + + got := r.forName("admin") + if len(got) != 1 { + t.Fatalf("attendu 1 notification, obtenu %d : %v", len(got), got) + } + if !strings.Contains(got[0].problem, "indéductibles") { + t.Errorf("problem = %q, devrait porter sur les interfaces indéductibles", got[0].problem) + } +} + +func TestCheckVPCs_OrdreDeterministe(t *testing.T) { + db := newTestDB(t) + for _, name := range []string{"vp-c", "vp-a", "vp-b"} { + seedResource(t, db, prefixVPC, name, state.Running) + } + + var first []string + for i := range 5 { + r := &recorder{} + if err := CheckVPCs(db, r); err != nil { + t.Fatalf("erreur inattendue: %v", err) + } + var order []string + for _, c := range r.calls { + if len(order) == 0 || order[len(order)-1] != c.name { + order = append(order, c.name) + } + } + if i == 0 { + first = order + continue + } + if strings.Join(order, ",") != strings.Join(first, ",") { + t.Fatalf("itération %d : ordre %v, attendu %v", i, order, first) + } + } +} diff --git a/internal/watchdog/helpers_test.go b/internal/watchdog/helpers_test.go new file mode 100644 index 0000000..41b9dcc --- /dev/null +++ b/internal/watchdog/helpers_test.go @@ -0,0 +1,61 @@ +package watchdog + +import ( + "strings" + "testing" + + "git.g3e.fr/syonad/two/internal/state" + "git.g3e.fr/syonad/two/internal/watchdog/notify" + "git.g3e.fr/syonad/two/pkg/db/kv" + + "github.com/dgraph-io/badger/v4" +) + +type notification struct { + kind string + name string + problem string +} + +type recorder struct { + calls []notification +} + +var _ notify.Notifier = (*recorder)(nil) + +func (r *recorder) Notify(kind, name, problem string) { + r.calls = append(r.calls, notification{kind: kind, name: name, problem: problem}) +} + +func (r *recorder) forName(name string) []notification { + var out []notification + for _, c := range r.calls { + if c.name == name { + out = append(out, c) + } + } + return out +} + +func (r *recorder) hasProblemContaining(substr string) bool { + for _, c := range r.calls { + if strings.Contains(c.problem, substr) { + return true + } + } + return false +} + +func newTestDB(t *testing.T) *badger.DB { + t.Helper() + db := kv.InitDB(kv.Config{Path: t.TempDir()}, false) + t.Cleanup(func() { db.Close() }) + return db +} + +func seedResource(t *testing.T, db *badger.DB, prefix, name string, s state.State) { + t.Helper() + if err := state.Set(db, prefix+name, s); err != nil { + t.Fatalf("seedResource %s%s: %v", prefix, name, err) + } +} diff --git a/internal/watchdog/resources.go b/internal/watchdog/resources.go new file mode 100644 index 0000000..76fc452 --- /dev/null +++ b/internal/watchdog/resources.go @@ -0,0 +1,52 @@ +package watchdog + +import ( + "fmt" + "sort" + "strings" + + "git.g3e.fr/syonad/two/internal/netif" +) + +const ( + prefixVPC = "vpc/" + prefixSubnet = "subnet/" + prefixVM = "vm/" + + kindVPC = "vpc" + kindSubnet = "subnet" + kindVM = "vm" + + stateSuffix = "/state" +) + +func resourceNames(pairs map[string]string, prefix string) []string { + names := make([]string, 0, len(pairs)) + for key := range pairs { + name, ok := strings.CutPrefix(key, prefix) + if !ok { + continue + } + name, ok = strings.CutSuffix(name, stateSuffix) + if !ok { + continue + } + if name == "" || strings.Contains(name, "/") { + continue + } + names = append(names, name) + } + sort.Strings(names) + return names +} + +func linkProblem(iface string) string { + up, err := netif.LinkIsUp(iface) + switch { + case err != nil: + return fmt.Sprintf("interface %s introuvable: %v", iface, err) + case !up: + return fmt.Sprintf("interface %s down", iface) + } + return "" +} diff --git a/internal/watchdog/resources_test.go b/internal/watchdog/resources_test.go new file mode 100644 index 0000000..6d0ed5c --- /dev/null +++ b/internal/watchdog/resources_test.go @@ -0,0 +1,81 @@ +package watchdog + +import ( + "reflect" + "strings" + "testing" +) + +func TestResourceNames_ExtraitDepuisLesClesState(t *testing.T) { + pairs := map[string]string{ + "vpc/vp-admin/state": "running", + "vpc/vp-lab/state": "creating", + } + want := []string{"vp-admin", "vp-lab"} + if got := resourceNames(pairs, prefixVPC); !reflect.DeepEqual(got, want) { + t.Errorf("resourceNames = %v, attendu %v", got, want) + } +} + +func TestResourceNames_IgnoreLesAutresCles(t *testing.T) { + pairs := map[string]string{ + "subnet/br-000042/state": "running", + "subnet/br-000042/vpc": "vp-admin", + "subnet/br-000042/cidr": "10.0.0.0/24", + "subnet/br-000042/local_iface": "br-000042", + } + want := []string{"br-000042"} + if got := resourceNames(pairs, prefixSubnet); !reflect.DeepEqual(got, want) { + t.Errorf("resourceNames = %v, attendu %v", got, want) + } +} + +func TestResourceNames_IgnoreLesClesImbriquees(t *testing.T) { + pairs := map[string]string{ + "vm/i-test1/state": "running", + "vm/i-test1/disk/vda/state": "attached", + "vm/i-test1/disk/vdb/state": "attached", + "vm/i-test1/dhcp/10.0.0.5": "00:22:33:44:55:66", + } + want := []string{"i-test1"} + if got := resourceNames(pairs, prefixVM); !reflect.DeepEqual(got, want) { + t.Errorf("resourceNames = %v, attendu %v", got, want) + } +} + +func TestResourceNames_Trie(t *testing.T) { + pairs := map[string]string{ + "vpc/vp-c/state": "running", + "vpc/vp-a/state": "running", + "vpc/vp-b/state": "running", + } + want := []string{"vp-a", "vp-b", "vp-c"} + for i := range 20 { + if got := resourceNames(pairs, prefixVPC); !reflect.DeepEqual(got, want) { + t.Fatalf("itération %d : resourceNames = %v, attendu %v", i, got, want) + } + } +} + +func TestResourceNames_MapVide(t *testing.T) { + if got := resourceNames(map[string]string{}, prefixVPC); len(got) != 0 { + t.Errorf("resourceNames = %v, attendu vide", got) + } +} + +func TestResourceNames_PrefixeNonCorrespondant(t *testing.T) { + pairs := map[string]string{"vm/i-test1/state": "running"} + if got := resourceNames(pairs, prefixVPC); len(got) != 0 { + t.Errorf("resourceNames = %v, attendu vide", got) + } +} + +func TestLinkProblem_InterfaceInexistante(t *testing.T) { + p := linkProblem("interface-qui-nexiste-pas-42") + if p == "" { + t.Fatal("linkProblem devrait signaler un problème pour une interface inexistante") + } + if !strings.Contains(p, "interface-qui-nexiste-pas-42") { + t.Errorf("le message %q devrait nommer l'interface", p) + } +}