From fef04b10fc483dda961576b261d7f2bf01dd3a0c Mon Sep 17 00:00:00 2001 From: GnomeZworc Date: Thu, 13 Aug 2026 23:43:16 +0200 Subject: [PATCH 1/2] f-37: migration: add state change at boot #37 Signed-off-by: GnomeZworc --- cmd/agent/main.go | 8 ++ internal/migration/state.go | 84 +++++++++++++++ internal/migration/state_test.go | 175 +++++++++++++++++++++++++++++++ 3 files changed, 267 insertions(+) create mode 100644 internal/migration/state.go create mode 100644 internal/migration/state_test.go diff --git a/cmd/agent/main.go b/cmd/agent/main.go index 7083076..0425d4e 100644 --- a/cmd/agent/main.go +++ b/cmd/agent/main.go @@ -8,6 +8,7 @@ import ( 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/pkg/db/kv" "git.g3e.fr/syonad/two/pkg/logger" @@ -31,6 +32,13 @@ func main() { db := kv.InitDB(kv.Config{Path: cfg.Database.Path}, false) defer 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. + if err := migration.MigrateStates(db, log.With(slog.String("component", "migration"))); err != nil { + log.Error("failed to migrate states", "error", err) + return + } + q := worker.New(cfg.Worker.BufferSize) q.Start(cfg.Worker.Count) diff --git a/internal/migration/state.go b/internal/migration/state.go new file mode 100644 index 0000000..99b5517 --- /dev/null +++ b/internal/migration/state.go @@ -0,0 +1,84 @@ +// Package migration met la DB au format attendu par la version courante de +// l'agent. Les migrations sont idempotentes et jouées au démarrage. +package migration + +import ( + "fmt" + "log/slog" + "strings" + + "git.g3e.fr/syonad/two/internal/state" + "git.g3e.fr/syonad/two/pkg/db/kv" + + "github.com/dgraph-io/badger/v4" +) + +// prefixes énumère les familles de ressources portant une clé /state. +var prefixes = []string{"vpc/", "subnet/", "vm/"} + +// legacyStates traduit l'ancien vocabulaire (états VPC/subnet d'avant +// l'unification, et états VM start/stop) vers l'enum courant. +var legacyStates = map[string]state.State{ + "created": state.Running, + "started": state.Running, + "starting": state.Creating, + "stopping": state.Deleting, + "stopped": state.Deleted, +} + +// MigrateStates convertit les valeurs de state héritées vers l'enum courant, +// puis marque en erreur les ressources restées dans un état transitoire. +// +// La réconciliation est sûre au démarrage : la queue worker est en mémoire et +// vide à ce moment, donc aucune commande n'est en cours. Une ressource en +// creating/deleting est nécessairement orpheline d'un arrêt de l'agent — sans +// ce passage en error, elle resterait indéfiniment non supprimable. +func MigrateStates(db *badger.DB, log *slog.Logger) error { + for _, prefix := range prefixes { + entries, err := kv.ListByPrefix(db, prefix) + if err != nil { + return fmt.Errorf("list %s: %w", prefix, err) + } + for key, value := range entries { + if !strings.HasSuffix(key, "/state") { + continue + } + resource := strings.TrimSuffix(key, "/state") + target, reason := targetState(value) + if target == "" { + continue + } + if err := state.Set(db, resource, target); err != nil { + return fmt.Errorf("migrate %s: %w", key, err) + } + log.Info("state migrated", + "resource", resource, "from", value, "to", string(target), "reason", reason) + } + } + return nil +} + +// targetState retourne l'état vers lequel migrer une valeur brute, ou "" si +// elle doit rester inchangée. +func targetState(value string) (state.State, string) { + if legacy, ok := legacyStates[value]; ok { + // Un état hérité transitoire est orphelin au même titre qu'un état + // transitoire courant : on applique la réconciliation directement. + if state.IsTransient(legacy) { + return state.Error, "orphaned legacy state" + } + return legacy, "legacy state" + } + + current, err := state.Parse(value) + if err != nil { + // Valeur inconnue : ni l'ancien vocabulaire, ni le nouveau. On la + // bascule en error plutôt que de la laisser bloquer l'agent — la + // ressource reste visible et supprimable. + return state.Error, "unknown state" + } + if state.IsTransient(current) { + return state.Error, "orphaned transient state" + } + return "", "" +} diff --git a/internal/migration/state_test.go b/internal/migration/state_test.go new file mode 100644 index 0000000..85137a4 --- /dev/null +++ b/internal/migration/state_test.go @@ -0,0 +1,175 @@ +package migration + +import ( + "io" + "log/slog" + "testing" + + "git.g3e.fr/syonad/two/internal/state" + "git.g3e.fr/syonad/two/pkg/db/kv" + + "github.com/dgraph-io/badger/v4" +) + +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 discardLogger() *slog.Logger { + return slog.New(slog.NewTextHandler(io.Discard, nil)) +} + +func TestMigrateStates_LegacyMapping(t *testing.T) { + db := newTestDB(t) + cases := map[string]struct { + resource string + from string + want state.State + }{ + "vpc created": {"vpc/vpc-1", "created", state.Running}, + "subnet created": {"subnet/sn-1", "created", state.Running}, + "vm started": {"vm/vm-1", "started", state.Running}, + "vm stopped": {"vm/vm-2", "stopped", state.Deleted}, + // Les états transitoires hérités sont orphelins après un redémarrage. + "vm starting": {"vm/vm-3", "starting", state.Error}, + "vm stopping": {"vm/vm-4", "stopping", state.Error}, + "vpc creating": {"vpc/vpc-2", "creating", state.Error}, + "subnet deleting": {"subnet/sn-2", "deleting", state.Error}, + } + for _, c := range cases { + kv.AddInDB(db, c.resource+"/state", c.from) + } + + if err := MigrateStates(db, discardLogger()); err != nil { + t.Fatalf("MigrateStates a échoué : %v", err) + } + + for name, c := range cases { + got, err := state.Get(db, c.resource) + if err != nil { + t.Errorf("%s : Get a échoué : %v", name, err) + continue + } + if got != c.want { + t.Errorf("%s : %q → %q, attendu %q", name, c.from, got, c.want) + } + } +} + +func TestMigrateStates_StableStatesUntouched(t *testing.T) { + db := newTestDB(t) + stable := map[string]state.State{ + "vpc/vpc-1": state.Running, + "subnet/sn-1": state.Deleted, + "vm/vm-1": state.Error, + } + for resource, s := range stable { + if err := state.Set(db, resource, s); err != nil { + t.Fatalf("préparation du test : %v", err) + } + } + + if err := MigrateStates(db, discardLogger()); err != nil { + t.Fatalf("MigrateStates a échoué : %v", err) + } + + for resource, want := range stable { + got, err := state.Get(db, resource) + if err != nil { + t.Fatalf("Get(%s) a échoué : %v", resource, err) + } + if got != want { + t.Errorf("%s : état modifié %q, attendu %q", resource, got, want) + } + } +} + +func TestMigrateStates_Idempotent(t *testing.T) { + db := newTestDB(t) + kv.AddInDB(db, "vpc/vpc-1/state", "created") + kv.AddInDB(db, "vm/vm-1/state", "starting") + + if err := MigrateStates(db, discardLogger()); err != nil { + t.Fatalf("premier passage : %v", err) + } + first := map[string]state.State{} + for _, r := range []string{"vpc/vpc-1", "vm/vm-1"} { + s, err := state.Get(db, r) + if err != nil { + t.Fatalf("Get(%s) : %v", r, err) + } + first[r] = s + } + + if err := MigrateStates(db, discardLogger()); err != nil { + t.Fatalf("second passage : %v", err) + } + for r, want := range first { + got, err := state.Get(db, r) + if err != nil { + t.Fatalf("Get(%s) : %v", r, err) + } + if got != want { + t.Errorf("%s : le second passage a modifié l'état (%q → %q)", r, want, got) + } + } +} + +func TestMigrateStates_UnknownValueBecomesError(t *testing.T) { + db := newTestDB(t) + kv.AddInDB(db, "vpc/vpc-corrompu/state", "n'importe quoi") + + if err := MigrateStates(db, discardLogger()); err != nil { + t.Fatalf("MigrateStates a échoué : %v", err) + } + + got, err := state.Get(db, "vpc/vpc-corrompu") + if err != nil { + t.Fatalf("Get a échoué : %v", err) + } + if got != state.Error { + t.Errorf("une valeur inconnue devrait devenir %q, obtenu %q", state.Error, got) + } +} + +func TestMigrateStates_IgnoresNonStateKeys(t *testing.T) { + db := newTestDB(t) + kv.AddInDB(db, "vpc/vpc-1/state", "created") + kv.AddInDB(db, "vpc/vpc-1/cidr", "10.0.0.0/16") + kv.AddInDB(db, "subnet/sn-1/local_iface", "br-vms") + kv.AddInDB(db, "vm/vm-1/disk/sda", "/data/root.qcow2") + + if err := MigrateStates(db, discardLogger()); err != nil { + t.Fatalf("MigrateStates a échoué : %v", err) + } + + untouched := map[string]string{ + "vpc/vpc-1/cidr": "10.0.0.0/16", + "subnet/sn-1/local_iface": "br-vms", + "vm/vm-1/disk/sda": "/data/root.qcow2", + } + for key, want := range untouched { + got, err := kv.GetFromDB(db, key) + if err != nil { + t.Errorf("la clé %s devrait exister : %v", key, err) + continue + } + if got != want { + t.Errorf("%s = %q, attendu %q", key, got, want) + } + } + // Aucune clé /state parasite ne doit apparaître sur ces ressources. + if _, err := kv.GetFromDB(db, "vm/vm-1/state"); err == nil { + t.Error("aucune clé state ne devrait être créée pour vm-1") + } +} + +func TestMigrateStates_EmptyDB(t *testing.T) { + db := newTestDB(t) + if err := MigrateStates(db, discardLogger()); err != nil { + t.Fatalf("MigrateStates devrait réussir sur une DB vide : %v", err) + } +} From 2150530ef26b6c135a2048939198b2e477e03912 Mon Sep 17 00:00:00 2001 From: GnomeZworc Date: Thu, 13 Aug 2026 23:55:53 +0200 Subject: [PATCH 2/2] f-37: periph: add final touch #37 Signed-off-by: GnomeZworc --- api/agent.yaml | 28 +++-- internal/api/agent/subnet.go | 8 +- internal/api/agent/subnet_test.go | 11 ++ internal/prometheus/agent/collector.go | 18 ++- internal/prometheus/agent/collector_test.go | 126 ++++++++++++++++++++ 5 files changed, 176 insertions(+), 15 deletions(-) create mode 100644 internal/prometheus/agent/collector_test.go diff --git a/api/agent.yaml b/api/agent.yaml index 7e9c461..0861353 100644 --- a/api/agent.yaml +++ b/api/agent.yaml @@ -91,7 +91,7 @@ paths: "404": $ref: "#/components/responses/NotFound" "409": - description: VPC not in a deletable state + description: VPC not deletable — only running or error states can be deleted, and all its subnets must be deleted first content: application/json: schema: @@ -146,7 +146,7 @@ paths: schema: $ref: "#/components/schemas/Error" "422": - description: Subnet not found or not in created state + description: Subnet not found, or not in creating/running state content: application/json: schema: @@ -185,6 +185,12 @@ paths: $ref: "#/components/schemas/VM" "404": $ref: "#/components/responses/NotFound" + "409": + description: VM not stoppable — only running or error states can be stopped + content: + application/json: + schema: + $ref: "#/components/schemas/Error" "500": $ref: "#/components/responses/InternalError" @@ -274,6 +280,12 @@ paths: $ref: "#/components/schemas/Subnet" "404": $ref: "#/components/responses/NotFound" + "409": + description: Subnet not deletable — only running or error states can be deleted + content: + application/json: + schema: + $ref: "#/components/schemas/Error" "500": $ref: "#/components/responses/InternalError" @@ -314,8 +326,8 @@ components: example: vp-00001 state: type: string - enum: [creating, created, deleting, deleted] - example: created + enum: [creating, running, error, deleting, deleted] + example: running cidr: type: string example: "10.0.0.0/16" @@ -373,8 +385,8 @@ components: example: sn-00001 state: type: string - enum: [creating, created, deleting, deleted] - example: created + enum: [creating, running, error, deleting, deleted] + example: running vpc: type: string example: vpc1 @@ -472,8 +484,8 @@ components: example: vm-00001 state: type: string - enum: [starting, started, stopping, stopped] - example: started + enum: [creating, running, error, deleting, deleted] + example: running metadata_port: type: string example: "80" diff --git a/internal/api/agent/subnet.go b/internal/api/agent/subnet.go index 36ee9d7..f0671f0 100644 --- a/internal/api/agent/subnet.go +++ b/internal/api/agent/subnet.go @@ -68,7 +68,13 @@ func (s *Server) getSubnet(w http.ResponseWriter, _ *http.Request, name string) func (s *Server) deleteSubnet(w http.ResponseWriter, _ *http.Request, name string) { cmd := dispatcher.DeleteSubnetCommand{Name: name} if err := s.dispatcher.Prepare(cmd); err != nil { - w.WriteHeader(http.StatusNotFound) + // 404 si la ressource n'existe pas, 409 si elle existe mais n'est pas + // dans un état supprimable — même convention que /vpcs et /vms. + if _, dbErr := kv.GetFromDB(s.db, "subnet/"+name+"/state"); dbErr != nil { + w.WriteHeader(http.StatusNotFound) + } else { + w.WriteHeader(http.StatusConflict) + } json.NewEncoder(w).Encode(ErrorResponse{Error: err.Error()}) return } diff --git a/internal/api/agent/subnet_test.go b/internal/api/agent/subnet_test.go index 1ff067a..77ab5e4 100644 --- a/internal/api/agent/subnet_test.go +++ b/internal/api/agent/subnet_test.go @@ -275,6 +275,17 @@ func TestDeleteSubnet_Success(t *testing.T) { } } +func TestDeleteSubnet_ConflictWhileCreating(t *testing.T) { + s, db := newTestServer(t) + kv.AddInDB(db, "subnet/sn-wip/state", "creating") + req := httptest.NewRequest(http.MethodDelete, "/subnets/sn-wip", nil) + w := httptest.NewRecorder() + s.SubnetByNameHandler(w, req) + if w.Code != http.StatusConflict { + t.Errorf("attendu 409, obtenu %d: %s", w.Code, w.Body.String()) + } +} + func TestDeleteSubnet_NotFound(t *testing.T) { s, _ := newTestServer(t) req := httptest.NewRequest(http.MethodDelete, "/subnets/inexistant", nil) diff --git a/internal/prometheus/agent/collector.go b/internal/prometheus/agent/collector.go index 04712ad..57dc371 100644 --- a/internal/prometheus/agent/collector.go +++ b/internal/prometheus/agent/collector.go @@ -3,12 +3,13 @@ package agentmetrics import ( "strings" + "git.g3e.fr/syonad/two/internal/state" "git.g3e.fr/syonad/two/pkg/db/kv" "github.com/dgraph-io/badger/v4" "github.com/prometheus/client_golang/prometheus" ) -var allStates = []string{"creating", "created", "deleting", "deleted"} +var allStates = state.All() // AgentCollector implements prometheus.Collector and exposes agent metrics // by querying the BadgerDB on each scrape. @@ -47,7 +48,7 @@ func (c *AgentCollector) Collect(ch chan<- prometheus.Metric) { // collectStates counts resources under the given DB prefix by their state value // and emits one gauge per state label. func (c *AgentCollector) collectStates(ch chan<- prometheus.Metric, prefix string, desc *prometheus.Desc) { - counts := make(map[string]float64, len(allStates)) + counts := make(map[state.State]float64, len(allStates)) for _, s := range allStates { counts[s] = 0 } @@ -55,13 +56,18 @@ func (c *AgentCollector) collectStates(ch chan<- prometheus.Metric, prefix strin items, err := kv.ListByPrefix(c.db, prefix) if err == nil { for key, val := range items { - if strings.HasSuffix(key, "/state") { - counts[val]++ + if !strings.HasSuffix(key, "/state") { + continue + } + // Une valeur hors enum n'est pas comptée : la migration au + // démarrage les a toutes ramenées dans l'enum. + if s, err := state.Parse(val); err == nil { + counts[s]++ } } } - for _, state := range allStates { - ch <- prometheus.MustNewConstMetric(desc, prometheus.GaugeValue, counts[state], state) + for _, s := range allStates { + ch <- prometheus.MustNewConstMetric(desc, prometheus.GaugeValue, counts[s], string(s)) } } diff --git a/internal/prometheus/agent/collector_test.go b/internal/prometheus/agent/collector_test.go new file mode 100644 index 0000000..1eed284 --- /dev/null +++ b/internal/prometheus/agent/collector_test.go @@ -0,0 +1,126 @@ +package agentmetrics + +import ( + "strings" + "testing" + + "git.g3e.fr/syonad/two/internal/state" + "git.g3e.fr/syonad/two/pkg/db/kv" + + "github.com/dgraph-io/badger/v4" + "github.com/prometheus/client_golang/prometheus" + dto "github.com/prometheus/client_model/go" +) + +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 +} + +// collect exécute une collecte et retourne, par nom de métrique, la valeur de +// chaque série indexée par son label state. +func collect(t *testing.T, c *AgentCollector) map[string]map[string]float64 { + t.Helper() + ch := make(chan prometheus.Metric, 64) + c.Collect(ch) + close(ch) + + out := map[string]map[string]float64{} + for m := range ch { + var pb dto.Metric + if err := m.Write(&pb); err != nil { + t.Fatalf("write metric: %v", err) + } + // Le fqName n'est pas exposé directement : il figure dans la + // représentation textuelle du Desc. + desc := m.Desc().String() + var name string + switch { + case strings.Contains(desc, "syonad_vpcs_total"): + name = "vpcs" + case strings.Contains(desc, "syonad_subnets_total"): + name = "subnets" + default: + t.Fatalf("métrique inattendue : %s", desc) + } + if out[name] == nil { + out[name] = map[string]float64{} + } + for _, l := range pb.GetLabel() { + if l.GetName() == "state" { + out[name][l.GetValue()] = pb.GetGauge().GetValue() + } + } + } + return out +} + +func TestCollector_CountsByState(t *testing.T) { + db := newTestDB(t) + for resource, s := range map[string]state.State{ + "vpc/vpc-1": state.Running, + "vpc/vpc-2": state.Running, + "vpc/vpc-3": state.Error, + "subnet/sn-1": state.Creating, + } { + if err := state.Set(db, resource, s); err != nil { + t.Fatalf("préparation du test : %v", err) + } + } + + got := collect(t, NewAgentCollector(db)) + wantVPC := map[string]float64{"creating": 0, "running": 2, "error": 1, "deleting": 0, "deleted": 0} + for s, want := range wantVPC { + if got["vpcs"][s] != want { + t.Errorf("vpcs{state=%q} = %v, attendu %v", s, got["vpcs"][s], want) + } + } + if got["subnets"]["creating"] != 1 { + t.Errorf("subnets{state=creating} = %v, attendu 1", got["subnets"]["creating"]) + } +} + +// Une série par état doit être émise même à zéro : une métrique qui disparaît +// côté Prometheus casse les alertes qui s'en servent. +func TestCollector_EmitsEveryState(t *testing.T) { + got := collect(t, NewAgentCollector(newTestDB(t))) + for _, family := range []string{"vpcs", "subnets"} { + if len(got[family]) != len(state.All()) { + t.Errorf("%s : %d séries, attendu %d", family, len(got[family]), len(state.All())) + } + for _, s := range state.All() { + if _, ok := got[family][string(s)]; !ok { + t.Errorf("%s : série manquante pour l'état %q", family, s) + } + } + } +} + +func TestCollector_IgnoresNonStateKeys(t *testing.T) { + db := newTestDB(t) + if err := state.Set(db, "vpc/vpc-1", state.Running); err != nil { + t.Fatalf("préparation du test : %v", err) + } + kv.AddInDB(db, "vpc/vpc-1/cidr", "10.0.0.0/16") + + got := collect(t, NewAgentCollector(db)) + if got["vpcs"]["running"] != 1 { + t.Errorf("vpcs{state=running} = %v, attendu 1", got["vpcs"]["running"]) + } +} + +func TestCollector_Describe(t *testing.T) { + ch := make(chan *prometheus.Desc, 10) + NewAgentCollector(newTestDB(t)).Describe(ch) + close(ch) + + var n int + for range ch { + n++ + } + if n != 2 { + t.Errorf("%d descripteurs, attendu 2", n) + } +}