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) + } +}