Compare commits

..

2 commits

Author SHA1 Message Date
2150530ef2
f-37: periph: add final touch #37
All checks were successful
Pre Release Workflow / set-release-target (push) Successful in 1s
Pre Release Workflow / build (agent, amd64, linux) (push) Successful in 1m35s
Pre Release Workflow / build (metadata, amd64, linux) (push) Successful in 1m31s
Pre Release Workflow / upload-scripts (run-dnsmasq-in-netns.sh) (push) Successful in 7s
Pre Release Workflow / prerelease (push) Successful in 10s
Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
2026-08-13 23:55:53 +02:00
fef04b10fc
f-37: migration: add state change at boot #37
Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
2026-08-13 23:43:16 +02:00
8 changed files with 443 additions and 15 deletions

View file

@ -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"

View file

@ -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)

View file

@ -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
}

View file

@ -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)

View 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 "", ""
}

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

View file

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

View file

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