f-37: migration: add state change at boot #37
Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
This commit is contained in:
parent
aa6611249b
commit
fef04b10fc
3 changed files with 267 additions and 0 deletions
|
|
@ -8,6 +8,7 @@ import (
|
||||||
agentapi "git.g3e.fr/syonad/two/internal/api/agent"
|
agentapi "git.g3e.fr/syonad/two/internal/api/agent"
|
||||||
configuration "git.g3e.fr/syonad/two/internal/config/agent"
|
configuration "git.g3e.fr/syonad/two/internal/config/agent"
|
||||||
dispatcher "git.g3e.fr/syonad/two/internal/dispatcher/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"
|
agentmetrics "git.g3e.fr/syonad/two/internal/prometheus/agent"
|
||||||
"git.g3e.fr/syonad/two/pkg/db/kv"
|
"git.g3e.fr/syonad/two/pkg/db/kv"
|
||||||
"git.g3e.fr/syonad/two/pkg/logger"
|
"git.g3e.fr/syonad/two/pkg/logger"
|
||||||
|
|
@ -31,6 +32,13 @@ func main() {
|
||||||
db := kv.InitDB(kv.Config{Path: cfg.Database.Path}, false)
|
db := kv.InitDB(kv.Config{Path: cfg.Database.Path}, false)
|
||||||
defer db.Close()
|
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 := worker.New(cfg.Worker.BufferSize)
|
||||||
q.Start(cfg.Worker.Count)
|
q.Start(cfg.Worker.Count)
|
||||||
|
|
||||||
|
|
|
||||||
84
internal/migration/state.go
Normal file
84
internal/migration/state.go
Normal file
|
|
@ -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 "", ""
|
||||||
|
}
|
||||||
175
internal/migration/state_test.go
Normal file
175
internal/migration/state_test.go
Normal file
|
|
@ -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)
|
||||||
|
}
|
||||||
|
}
|
||||||
Loading…
Add table
Add a link
Reference in a new issue